{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Arbiter.Core.Sql.Archive
( archiveAckCte
, updateArchiveResultSQL
, updateArchiveResultsBatchSQL
, purgeArchiveSQL
, archivePurgeBatch
, listArchiveFilteredSQL
, countArchiveFilteredSQL
, deleteArchiveJobSQL
, deleteArchiveJobsBatchSQL
, reEnqueueFromArchiveSQL
, allArchiveColumns
) where
import Data.Aeson (Value)
import Data.Int (Int64)
import Data.Text (Text)
import Data.Text qualified as T
import Data.Time (UTCTime)
import NeatInterpolation (text)
import Arbiter.Core.Codec (archiveRowCodec, jobRowCodec)
import Arbiter.Core.Job.Schema (jobQueueArchiveTable, jobQueueTable)
import Arbiter.Core.Job.Types (JobRead)
import Arbiter.Core.Sql.Jobs (enqueuedAgainCols, jobColsExceptId, jobColumns)
import Arbiter.Core.Sql.QQ (sql)
import Arbiter.Core.Sql.Query (Query, rows)
allArchiveColumns :: Text
allArchiveColumns :: Text
allArchiveColumns =
[text|
id, completed_at, job_id, payload, group_key, inserted_at, updated_at, attempts, last_error, priority,
last_attempted_at, not_visible_until, dedup_key, dedup_strategy, max_attempts,
parent_id, parent_state, traceparent, tracestate, suspended, claimed_by, claim_seq,
archive_for, kind, rate_limit_key, rate_limit_prefix, concurrency_key, concurrency_prefix,
result
|]
archiveAckCte :: Text -> Text -> Text -> Text
archiveAckCte :: Text -> Text -> Text -> Text
archiveAckCte Text
schema Text
tableName Text
ackCte =
let archiveTbl :: Text
archiveTbl = Text -> Text -> Text
jobQueueArchiveTable Text
schema Text
tableName
in [text|
archived AS (
INSERT INTO ${archiveTbl} (job_id, ${jobColsExceptId}, rate_limit_cost, completed_at, archive_expires_at)
SELECT id, ${jobColsExceptId}, rate_limit_cost, NOW(), NOW() + (archive_for * interval '1 second')
FROM ${ackCte}
WHERE archive_for > 0
),
|]
updateArchiveResultSQL :: Text -> Text -> Value -> Int64 -> Query ()
updateArchiveResultSQL :: Text -> Text -> Value -> Int64 -> Query ()
updateArchiveResultSQL Text
schema Text
tableName Value
result Int64
jobId =
let archiveTbl :: Text
archiveTbl = Text -> Text -> Text
jobQueueArchiveTable Text
schema Text
tableName
in [sql|UPDATE ${archiveTbl} SET result = #{result :: CJsonb} WHERE job_id = #{jobId :: CInt8}|]
updateArchiveResultsBatchSQL :: Text -> Text -> [Int64] -> [Value] -> Query ()
updateArchiveResultsBatchSQL :: Text -> Text -> [Int64] -> [Value] -> Query ()
updateArchiveResultsBatchSQL Text
schema Text
tableName [Int64]
jobIds [Value]
results =
let archiveTbl :: Text
archiveTbl = Text -> Text -> Text
jobQueueArchiveTable Text
schema Text
tableName
in [sql|
UPDATE ${archiveTbl} archive_row SET result = src.result
FROM (
SELECT unnest(#{jobIds :: [CInt8]}::bigint[]) AS job_id,
unnest(#{results :: [CJsonb]}::jsonb[]) AS result
) src
WHERE archive_row.job_id = src.job_id
|]
archivePurgeBatch :: Int
archivePurgeBatch :: Int
archivePurgeBatch = Int
10000
purgeArchiveSQL :: Text -> Text -> Text
purgeArchiveSQL :: Text -> Text -> Text
purgeArchiveSQL Text
schema Text
tableName =
let archiveTbl :: Text
archiveTbl = Text -> Text -> Text
jobQueueArchiveTable Text
schema Text
tableName
lim :: Text
lim = String -> Text
T.pack (Int -> String
forall a. Show a => a -> String
show Int
archivePurgeBatch)
in [text|
DELETE FROM ${archiveTbl}
WHERE ctid IN (
SELECT ctid FROM ${archiveTbl}
WHERE archive_expires_at < NOW()
LIMIT ${lim}
)
|]
listArchiveFilteredSQL
:: Text -> Text -> Query () -> Text -> Int64 -> Int64 -> Query (Int64, UTCTime, JobRead Value, Maybe Value)
listArchiveFilteredSQL :: Text
-> Text
-> Query ()
-> Text
-> Int64
-> Int64
-> Query (Int64, UTCTime, JobRead Value, Maybe Value)
listArchiveFilteredSQL Text
schema Text
tableName Query ()
whereFrag Text
orderBy Int64
limit Int64
offset =
let archiveTbl :: Text
archiveTbl = Text -> Text -> Text
jobQueueArchiveTable Text
schema Text
tableName
in RowCodec (Int64, UTCTime, JobRead Value, Maybe Value)
-> Query () -> Query (Int64, UTCTime, JobRead Value, Maybe Value)
forall a. RowCodec a -> Query () -> Query a
rows
(Text -> RowCodec (Int64, UTCTime, JobRead Value, Maybe Value)
archiveRowCodec Text
tableName)
[sql|
SELECT ${allArchiveColumns}
FROM ${archiveTbl}
${whereFrag}
ORDER BY ${orderBy}
LIMIT #{limit :: CInt8} OFFSET #{offset :: CInt8}
|]
countArchiveFilteredSQL :: Text -> Text -> Query () -> Query Int64
countArchiveFilteredSQL :: Text -> Text -> Query () -> Query Int64
countArchiveFilteredSQL Text
schema Text
tableName Query ()
whereFrag =
let archiveTbl :: Text
archiveTbl = Text -> Text -> Text
jobQueueArchiveTable Text
schema Text
tableName
in [sql|SELECT COUNT(*) AS @{count :: CInt8} FROM ${archiveTbl} ${whereFrag}|]
deleteArchiveJobSQL :: Text -> Text -> Int64 -> Query ()
deleteArchiveJobSQL :: Text -> Text -> Int64 -> Query ()
deleteArchiveJobSQL Text
schema Text
tableName Int64
archiveId =
let archiveTbl :: Text
archiveTbl = Text -> Text -> Text
jobQueueArchiveTable Text
schema Text
tableName
in [sql|DELETE FROM ${archiveTbl} WHERE id = #{archiveId :: CInt8}|]
deleteArchiveJobsBatchSQL :: Text -> Text -> [Int64] -> Query ()
deleteArchiveJobsBatchSQL :: Text -> Text -> [Int64] -> Query ()
deleteArchiveJobsBatchSQL Text
schema Text
tableName [Int64]
archiveIds =
let archiveTbl :: Text
archiveTbl = Text -> Text -> Text
jobQueueArchiveTable Text
schema Text
tableName
in [sql|DELETE FROM ${archiveTbl} WHERE id = ANY(#{archiveIds :: [CInt8]})|]
reEnqueueFromArchiveSQL :: Text -> Text -> Int64 -> Query (JobRead Value)
reEnqueueFromArchiveSQL :: Text -> Text -> Int64 -> Query (JobRead Value)
reEnqueueFromArchiveSQL Text
schema Text
tableName Int64
archiveId =
let archiveTbl :: Text
archiveTbl = Text -> Text -> Text
jobQueueArchiveTable Text
schema Text
tableName
tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in RowCodec (JobRead Value) -> Query () -> Query (JobRead Value)
forall a. RowCodec a -> Query () -> Query a
rows
(Text -> RowCodec (JobRead Value)
jobRowCodec Text
tableName)
[sql|
INSERT INTO ${tbl} (${enqueuedAgainCols})
SELECT ${enqueuedAgainCols}
FROM ${archiveTbl}
WHERE id = #{archiveId :: CInt8}
RETURNING ${jobColumns}
|]