{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Arbiter.Core.Sql.DLQ
( DLQMove (..)
, moveToDLQSQL
, selectExhaustedJobsSQL
, retryFromDLQSQL
, dlqJobExistsSQL
, deleteDLQJobSQL
, moveToDLQBatchSQL
, deleteDLQJobsBatchSQL
, cascadeChildrenToDLQSQL
, countDLQChildrenBatchSQL
) where
import Data.Aeson (Value)
import Data.Int (Int64)
import Data.Text (Text)
import Data.Text qualified as T
import NeatInterpolation (text)
import Arbiter.Core.Codec (jobRowCodec)
import Arbiter.Core.Job.Schema (jobQueueDLQTable, jobQueueTable)
import Arbiter.Core.Job.Types (JobRead, defaultMaxAttempts)
import Arbiter.Core.Sql.Jobs (dlqCarriedCols, jobColumns, requeuedCols)
import Arbiter.Core.Sql.QQ (sql)
import Arbiter.Core.Sql.Query (Query, mwhen, rows)
import Arbiter.Core.Sql.Tree (lockedByIdsCte)
data DLQMove = MoveNow | MoveIfExhausted
deriving stock (DLQMove -> DLQMove -> Bool
(DLQMove -> DLQMove -> Bool)
-> (DLQMove -> DLQMove -> Bool) -> Eq DLQMove
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: DLQMove -> DLQMove -> Bool
== :: DLQMove -> DLQMove -> Bool
$c/= :: DLQMove -> DLQMove -> Bool
/= :: DLQMove -> DLQMove -> Bool
Eq, Int -> DLQMove -> ShowS
[DLQMove] -> ShowS
DLQMove -> String
(Int -> DLQMove -> ShowS)
-> (DLQMove -> String) -> ([DLQMove] -> ShowS) -> Show DLQMove
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> DLQMove -> ShowS
showsPrec :: Int -> DLQMove -> ShowS
$cshow :: DLQMove -> String
show :: DLQMove -> String
$cshowList :: [DLQMove] -> ShowS
showList :: [DLQMove] -> ShowS
Show)
sweepableGuard :: Text
sweepableGuard :: Text
sweepableGuard =
let dma :: Text
dma = String -> Text
T.pack (Int32 -> String
forall a. Show a => a -> String
show Int32
defaultMaxAttempts)
in [text|
NOT suspended AND cancel_requested_at IS NULL
AND attempts >= COALESCE(max_attempts, ${dma})
AND (not_visible_until IS NULL OR not_visible_until <= NOW())
|]
moveToDLQSQL :: DLQMove -> Text -> Text -> Int64 -> Int64 -> Text -> Query Int64
moveToDLQSQL :: DLQMove -> Text -> Text -> Int64 -> Int64 -> Text -> Query Int64
moveToDLQSQL DLQMove
move Text
schema Text
tableName Int64
jobId Int64
cseq Text
errorMsg =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
exhausted :: Text
exhausted = Bool -> Text -> Text
forall m. Monoid m => Bool -> m -> m
mwhen (DLQMove
move DLQMove -> DLQMove -> Bool
forall a. Eq a => a -> a -> Bool
== DLQMove
MoveIfExhausted) [text|AND ${sweepableGuard}|]
in [sql|
WITH deleted_job AS (
DELETE FROM ${tbl}
WHERE id = #{jobId :: CInt8} AND claim_seq = #{cseq :: CInt8} ${exhausted}
RETURNING *
),
inserted_dlq AS (
INSERT INTO ${dlqTbl} (job_id, ${dlqCarriedCols}, last_error)
SELECT id, ${dlqCarriedCols}, #{errorMsg :: CText}
FROM deleted_job
)
SELECT count(*) AS @{count :: CInt8} FROM deleted_job
|]
selectExhaustedJobsSQL :: Text -> Text -> Int -> Query (Int64, Int64, Maybe Int64, Bool)
selectExhaustedJobsSQL :: Text -> Text -> Int -> Query (Int64, Int64, Maybe Int64, Bool)
selectExhaustedJobsSQL Text
schema Text
tableName Int
limit =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
lim :: Text
lim = String -> Text
T.pack (Int -> String
forall a. Show a => a -> String
show Int
limit)
in [sql|
SELECT @{id :: CInt8}, @{claim_seq :: CInt8}, @{parent_id :: Maybe CInt8},
(parent_state IS NOT NULL) AS @{is_rollup :: CBool}
FROM ${tbl}
WHERE ${sweepableGuard}
ORDER BY id ASC
LIMIT ${lim}
|]
retryFromDLQSQL :: Text -> Text -> Int64 -> Query (JobRead Value)
retryFromDLQSQL :: Text -> Text -> Int64 -> Query (JobRead Value)
retryFromDLQSQL Text
schema Text
tableName Int64
dlqId =
let dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable 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|
WITH RECURSIVE
target AS (
SELECT * FROM ${dlqTbl} WHERE id = #{dlqId :: CInt8}
),
-- Walk up through DLQ ancestors to find the root of the tree.
-- Stops when parent_id IS NULL, parent is in main queue, or
-- parent is not found in DLQ (orphaned).
ancestors AS (
SELECT dead.job_id, dead.parent_id, 0 AS depth
FROM ${dlqTbl} dead
WHERE dead.job_id = (SELECT parent_id FROM target)
AND (SELECT parent_id FROM target) IS NOT NULL
AND NOT EXISTS (SELECT 1 FROM ${tbl} WHERE id = (SELECT parent_id FROM target))
UNION ALL
SELECT dead.job_id, dead.parent_id, ancestor.depth + 1
FROM ${dlqTbl} dead
JOIN ancestors ancestor ON dead.job_id = ancestor.parent_id
WHERE ancestor.parent_id IS NOT NULL
AND NOT EXISTS (SELECT 1 FROM ${tbl} WHERE id = ancestor.parent_id)
),
-- Root is the topmost DLQ ancestor, or the target itself
root_job_id AS (
SELECT COALESCE(
(SELECT job_id FROM ancestors ORDER BY depth DESC LIMIT 1),
(SELECT job_id FROM target)
) AS job_id
),
-- The root's parent is NULL or exists in the main queue
can_retry AS (
SELECT EXISTS (
SELECT 1
FROM root_job_id root
JOIN ${dlqTbl} dead ON dead.job_id = root.job_id
WHERE dead.parent_id IS NULL
OR EXISTS (SELECT 1 FROM ${tbl} WHERE id = dead.parent_id)
) AS val
),
-- Walk down from root to collect all DLQ tree members
tree AS (
SELECT dead.id AS dlq_id, dead.job_id
FROM ${dlqTbl} dead
WHERE dead.job_id = (SELECT job_id FROM root_job_id)
UNION ALL
SELECT dead.id AS dlq_id, dead.job_id
FROM ${dlqTbl} dead
JOIN tree member ON dead.parent_id = member.job_id
),
-- Delete all tree members from DLQ (guarded by can_retry)
deleted AS (
DELETE FROM ${dlqTbl}
WHERE id IN (SELECT dlq_id FROM tree)
AND (SELECT val FROM can_retry)
RETURNING job_id, claim_seq, ${requeuedCols}
),
-- Re-insert into the main queue. A rollup finalizer is suspended when
-- it has children in this retry batch or in the main queue.
inserted AS (
INSERT INTO ${tbl} (id, attempts, claim_seq, suspended, ${requeuedCols})
SELECT dead.job_id, 0, dead.claim_seq + 1,
CASE WHEN dead.parent_state IS NOT NULL
THEN EXISTS (SELECT 1 FROM deleted child WHERE child.parent_id = dead.job_id)
OR EXISTS (SELECT 1 FROM ${tbl} WHERE parent_id = dead.job_id)
ELSE FALSE
END,
${requeuedCols}
FROM deleted dead
RETURNING *
),
-- Re-suspend any rollup parents already in the main queue that just
-- got fresh children re-inserted under them.
parent_resuspend AS (
UPDATE ${tbl}
SET suspended = TRUE, updated_at = NOW()
WHERE parent_state IS NOT NULL
AND id IN (SELECT DISTINCT parent_id FROM inserted WHERE parent_id IS NOT NULL)
AND NOT suspended
AND NOT (attempts > 0 AND not_visible_until IS NOT NULL AND not_visible_until > NOW())
)
SELECT ${jobColumns} FROM inserted WHERE id = (SELECT job_id FROM target)
|]
dlqJobExistsSQL :: Text -> Text -> Int64 -> Query Bool
dlqJobExistsSQL :: Text -> Text -> Int64 -> Query Bool
dlqJobExistsSQL Text
schema Text
tableName Int64
dlqId =
let dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
in [sql|SELECT EXISTS (SELECT 1 FROM ${dlqTbl} WHERE id = #{dlqId :: CInt8}) AS @{result :: CBool}|]
deleteDLQJobSQL :: Text -> Text -> Int64 -> Query (Maybe Int64)
deleteDLQJobSQL :: Text -> Text -> Int64 -> Query (Maybe Int64)
deleteDLQJobSQL Text
schema Text
tableName Int64
dlqId =
let dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
in [sql|DELETE FROM ${dlqTbl} WHERE id = #{dlqId :: CInt8} RETURNING @{parent_id :: Maybe CInt8}|]
moveToDLQBatchSQL :: Text -> Text -> [Int64] -> [Int64] -> [Text] -> Query Int64
moveToDLQBatchSQL :: Text -> Text -> [Int64] -> [Int64] -> [Text] -> Query Int64
moveToDLQBatchSQL Text
schema Text
tableName [Int64]
ids [Int64]
cseqs [Text]
errs =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
locked :: Query ()
locked = Text -> [Int64] -> Query ()
lockedByIdsCte Text
tbl [Int64]
ids
in [sql|
WITH input_jobs AS (
SELECT unnest(#{ids :: [CInt8]}::bigint[]) AS id,
unnest(#{cseqs :: [CInt8]}::bigint[]) AS expected_claim_seq,
unnest(#{errs :: [CText]}::text[]) AS error_msg
),
${locked},
deleted_jobs AS (
DELETE FROM ${tbl} job
USING input_jobs input_job
WHERE job.id = input_job.id AND job.claim_seq = input_job.expected_claim_seq
AND job.id IN (SELECT id FROM locked)
RETURNING job.*, input_job.error_msg AS new_error
),
inserted_dlq AS (
INSERT INTO ${dlqTbl} (job_id, failed_at, ${dlqCarriedCols}, last_error)
SELECT id, NOW(), ${dlqCarriedCols}, new_error
FROM deleted_jobs
)
SELECT id AS @{result :: CInt8} FROM deleted_jobs
|]
deleteDLQJobsBatchSQL :: Text -> Text -> [Int64] -> Query (Int64, Maybe Int64)
deleteDLQJobsBatchSQL :: Text -> Text -> [Int64] -> Query (Int64, Maybe Int64)
deleteDLQJobsBatchSQL Text
schema Text
tableName [Int64]
dlqIds =
let dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
in [sql|
DELETE FROM ${dlqTbl} WHERE id = ANY(#{dlqIds :: [CInt8]})
RETURNING @{id :: CInt8}, @{parent_id :: Maybe CInt8}
|]
cascadeChildrenToDLQSQL :: Text -> Text -> Int64 -> Text -> Query Int64
cascadeChildrenToDLQSQL :: Text -> Text -> Int64 -> Text -> Query Int64
cascadeChildrenToDLQSQL Text
schema Text
tableName Int64
parentId Text
errorMsg =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
in [sql|
WITH RECURSIVE descendants AS (
SELECT id FROM ${tbl} WHERE parent_id = #{parentId :: CInt8}
UNION ALL
SELECT job.id FROM ${tbl} job JOIN descendants descendant ON job.parent_id = descendant.id
),
deleted AS (
DELETE FROM ${tbl}
WHERE id IN (SELECT id FROM descendants)
RETURNING id, ${dlqCarriedCols}
),
inserted_dlq AS (
INSERT INTO ${dlqTbl} (job_id, ${dlqCarriedCols}, last_error)
SELECT id, ${dlqCarriedCols}, #{errorMsg :: CText}
FROM deleted
)
SELECT count(*) AS @{count :: CInt8} FROM deleted
|]
countDLQChildrenBatchSQL :: Text -> Text -> [Int64] -> Query (Int64, Int64)
countDLQChildrenBatchSQL :: Text -> Text -> [Int64] -> Query (Int64, Int64)
countDLQChildrenBatchSQL Text
schema Text
tableName [Int64]
jobIds =
let dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
in [sql|
SELECT @{parent_id :: CInt8}, COUNT(*) AS @{count :: CInt8}
FROM ${dlqTbl}
WHERE parent_id = ANY(#{jobIds :: [CInt8]})
GROUP BY parent_id
|]