{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Arbiter.Core.Sql.Tree
( pauseChildrenSQL
, resumeChildrenSQL
, descendantsCte
, lockedByIdsCte
, lockJobTreesSQL
, lockJobTreesFromRootSQL
, cancelJobCascadeSQL
, forceCancelJobSQL
, deleteCancelledJobsSQL
, selectCancelledReapableJobsSQL
, cancelJobTreeSQL
, tryWakeAncestorSQL
, treeRollupIdsSQL
, suspendJobSQL
, resumeJobSQL
, jobExistsSQL
, getParentIdSQL
, getParentIdsSQL
, insertResultSQL
, insertResultsBatchSQL
, getResultsByParentSQL
, getDLQChildErrorsByParentSQL
, persistParentStateSQL
, getParentStateSnapshotSQL
, readChildResultsSQL
) where
import Data.Aeson (Value)
import Data.Int (Int64)
import Data.Text (Text)
import Data.Text qualified as T
import Data.UUID.Types (UUID)
import Arbiter.Core.Job.Schema
( SchemaName
, TableName
, cancelNotifyChannel
, jobQueueDLQTable
, jobQueueResultsTable
, jobQueueTable
)
import Arbiter.Core.Sql.QQ (sql)
import Arbiter.Core.Sql.Query (Query)
import Arbiter.Core.SqlLiterals (textLiteral)
pauseChildrenSQL :: Text -> Text -> Int64 -> Query ()
pauseChildrenSQL :: Text -> Text -> Int64 -> Query ()
pauseChildrenSQL Text
schema Text
tableName Int64
parentId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
cte :: Query ()
cte = Text -> Int64 -> Query ()
childDescendantsCte Text
tbl Int64
parentId
locked :: Query ()
locked = Text -> Query ()
lockDescendantsCte Text
tbl
in [sql|
${cte},
${locked}
UPDATE ${tbl}
SET suspended = TRUE, updated_at = NOW()
WHERE id IN (SELECT id FROM locked)
AND NOT suspended
AND (not_visible_until IS NULL OR not_visible_until <= NOW())
|]
resumeChildrenSQL :: Text -> Text -> Int64 -> Query ()
resumeChildrenSQL :: Text -> Text -> Int64 -> Query ()
resumeChildrenSQL Text
schema Text
tableName Int64
parentId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
cte :: Query ()
cte = Text -> Int64 -> Query ()
childDescendantsCte Text
tbl Int64
parentId
locked :: Query ()
locked = Text -> Query ()
lockDescendantsCte Text
tbl
in [sql|
${cte},
${locked}
UPDATE ${tbl} job
SET suspended = FALSE, updated_at = NOW()
WHERE job.id IN (SELECT id FROM locked)
AND job.suspended = TRUE
AND NOT (
job.parent_state IS NOT NULL
AND EXISTS (SELECT 1 FROM ${tbl} child WHERE child.parent_id = job.id)
)
|]
descendantsFromCte :: Text -> Query () -> Query ()
descendantsFromCte :: Text -> Query () -> Query ()
descendantsFromCte Text
tbl Query ()
seed =
[sql|
WITH RECURSIVE descendants AS (
SELECT id FROM ${tbl} WHERE ${seed}
UNION ALL
SELECT job.id FROM ${tbl} job JOIN descendants descendant ON job.parent_id = descendant.id
)
|]
childDescendantsCte :: Text -> Int64 -> Query ()
childDescendantsCte :: Text -> Int64 -> Query ()
childDescendantsCte Text
tbl Int64
parentId = Text -> Query () -> Query ()
descendantsFromCte Text
tbl [sql|parent_id = #{parentId :: CInt8}|]
descendantsCte :: Text -> Int64 -> Query ()
descendantsCte :: Text -> Int64 -> Query ()
descendantsCte Text
tbl Int64
jobId = Text -> Query () -> Query ()
descendantsFromCte Text
tbl [sql|id = #{jobId :: CInt8}|]
lockDescendantsCte :: Text -> Query ()
lockDescendantsCte :: Text -> Query ()
lockDescendantsCte Text
tbl =
[sql|
locked AS (
SELECT id FROM ${tbl} WHERE id IN (SELECT id FROM descendants)
ORDER BY id DESC
FOR UPDATE
)
|]
lockedByIdsCte :: Text -> [Int64] -> Query ()
lockedByIdsCte :: Text -> [Int64] -> Query ()
lockedByIdsCte Text
tbl [Int64]
ids =
[sql|
locked AS (
SELECT id FROM ${tbl} WHERE id = ANY(#{ids :: [CInt8]})
ORDER BY id DESC
FOR UPDATE
)
|]
rootsFromCte :: Text -> Query () -> Query ()
rootsFromCte :: Text -> Query () -> Query ()
rootsFromCte Text
tbl Query ()
seed =
[sql|
WITH RECURSIVE ancestors AS (
SELECT id, parent_id FROM ${tbl} WHERE ${seed}
UNION
SELECT job.id, job.parent_id FROM ${tbl} job JOIN ancestors ancestor ON job.id = ancestor.parent_id
),
roots AS (
SELECT id FROM ancestors WHERE parent_id IS NULL
)
|]
descendantsOfCte :: Text -> [Int64] -> Query ()
descendantsOfCte :: Text -> [Int64] -> Query ()
descendantsOfCte Text
tbl [Int64]
jobIds =
[sql|
WITH RECURSIVE descendants AS (
SELECT id FROM ${tbl} WHERE id = ANY(#{jobIds :: [CInt8]})
UNION
SELECT job.id FROM ${tbl} job JOIN descendants descendant ON job.parent_id = descendant.id
)
|]
lockJobTreesSQL :: Text -> Text -> [Int64] -> Query Int64
lockJobTreesSQL :: Text -> Text -> [Int64] -> Query Int64
lockJobTreesSQL Text
schema Text
tableName [Int64]
jobIds =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
cte :: Query ()
cte = Text -> [Int64] -> Query ()
descendantsOfCte Text
tbl [Int64]
jobIds
locked :: Query ()
locked = Text -> Query ()
lockDescendantsCte Text
tbl
in [sql|
${cte},
${locked}
SELECT count(*) AS @{count :: CInt8} FROM locked
|]
lockJobTreesFromRootSQL :: Text -> Text -> [Int64] -> Query Int64
lockJobTreesFromRootSQL :: Text -> Text -> [Int64] -> Query Int64
lockJobTreesFromRootSQL Text
schema Text
tableName [Int64]
jobIds =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
cte :: Query ()
cte = Text -> Query () -> Query ()
rootsFromCte Text
tbl [sql|id = ANY(#{jobIds :: [CInt8]})|]
locked :: Query ()
locked = Text -> Query ()
lockDescendantsCte Text
tbl
in [sql|
${cte},
descendants AS (
SELECT id FROM ${tbl} WHERE id = ANY(#{jobIds :: [CInt8]}) OR id IN (SELECT id FROM roots)
UNION
SELECT job.id FROM ${tbl} job JOIN descendants descendant ON job.parent_id = descendant.id
),
${locked}
SELECT count(*) AS @{count :: CInt8} FROM locked
|]
cancelJobCascadeSQL :: Text -> Text -> Int64 -> Query Int64
cancelJobCascadeSQL :: Text -> Text -> Int64 -> Query Int64
cancelJobCascadeSQL Text
schema Text
tableName Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
cte :: Query ()
cte = Text -> Int64 -> Query ()
descendantsCte Text
tbl Int64
jobId
locked :: Query ()
locked = Text -> Query ()
lockDescendantsCte Text
tbl
in [sql|
${cte},
${locked},
deleted AS (
DELETE FROM ${tbl} WHERE id IN (SELECT id FROM locked)
RETURNING id
)
SELECT count(*) AS @{count :: CInt8} FROM deleted
|]
forceCancelJobSQL :: SchemaName -> TableName -> Int64 -> Query Int64
forceCancelJobSQL :: Text -> Text -> Int64 -> Query Int64
forceCancelJobSQL Text
schema Text
tableName Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
cte :: Query ()
cte = Text -> Int64 -> Query ()
descendantsCte Text
tbl Int64
jobId
chan :: Text
chan = Text -> Text
textLiteral (Text -> Text -> Text
cancelNotifyChannel Text
schema Text
tableName)
in [sql|
${cte},
locked AS (
SELECT id, claimed_by, not_visible_until
FROM ${tbl}
WHERE id IN (SELECT id FROM descendants)
ORDER BY id DESC
FOR UPDATE
),
cancelled AS (
UPDATE ${tbl} job SET cancel_requested_at = NOW(), claim_seq = job.claim_seq + 1
FROM locked held
WHERE job.id = held.id
AND held.claimed_by IS NOT NULL
AND held.not_visible_until IS NOT NULL AND held.not_visible_until > NOW()
RETURNING job.id, job.claimed_by
),
deleted AS (
DELETE FROM ${tbl} job
USING locked held
WHERE job.id = held.id
AND (held.claimed_by IS NULL OR held.not_visible_until IS NULL OR held.not_visible_until <= NOW())
RETURNING job.id, job.claimed_by
),
notif AS (
SELECT pg_notify(
${chan},
json_build_object('worker_id', notified.claimed_by, 'job_id', notified.id)::text
)
FROM (
SELECT id, claimed_by FROM cancelled
UNION ALL
SELECT id, claimed_by FROM deleted WHERE claimed_by IS NOT NULL
) notified
)
SELECT ((SELECT count(*) FROM cancelled) + (SELECT count(*) FROM deleted))::int8 AS @{count :: CInt8}
WHERE (SELECT count(*) FROM notif) >= 0
|]
deleteCancelledJobsSQL :: SchemaName -> TableName -> Maybe UUID -> [Int64] -> Query (Int64, Maybe Int64)
deleteCancelledJobsSQL :: Text -> Text -> Maybe UUID -> [Int64] -> Query (Int64, Maybe Int64)
deleteCancelledJobsSQL Text
schema Text
tableName Maybe UUID
owner [Int64]
jobIds =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|
WITH locked AS (
SELECT id FROM ${tbl}
WHERE id = ANY(#{jobIds :: [CInt8]})
AND cancel_requested_at IS NOT NULL
AND (claimed_by = #{owner :: Maybe CUuid} OR not_visible_until IS NULL OR not_visible_until <= NOW())
ORDER BY id DESC
FOR UPDATE
)
DELETE FROM ${tbl} WHERE id IN (SELECT id FROM locked)
RETURNING @{id :: CInt8}, @{parent_id :: Maybe CInt8}
|]
selectCancelledReapableJobsSQL :: SchemaName -> TableName -> Int -> Query Int64
selectCancelledReapableJobsSQL :: Text -> Text -> Int -> Query Int64
selectCancelledReapableJobsSQL 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}
FROM ${tbl}
WHERE cancel_requested_at IS NOT NULL
AND (not_visible_until IS NULL OR not_visible_until <= NOW())
ORDER BY id ASC
LIMIT ${lim}
|]
cancelJobTreeSQL :: Text -> Text -> Int64 -> Query Int64
cancelJobTreeSQL :: Text -> Text -> Int64 -> Query Int64
cancelJobTreeSQL Text
schema Text
tableName Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
cte :: Query ()
cte = Text -> Query () -> Query ()
rootsFromCte Text
tbl [sql|id = #{jobId :: CInt8}|]
locked :: Query ()
locked = Text -> Query ()
lockDescendantsCte Text
tbl
in [sql|
${cte},
descendants AS (
SELECT id FROM ${tbl} WHERE id IN (SELECT id FROM roots)
UNION ALL
SELECT job.id FROM ${tbl} job JOIN descendants descendant ON job.parent_id = descendant.id
),
${locked},
deleted AS (
DELETE FROM ${tbl} WHERE id IN (SELECT id FROM locked)
RETURNING id
)
SELECT count(*) AS @{count :: CInt8} FROM deleted
|]
tryWakeAncestorSQL :: Text -> Text -> Int64 -> Query ()
tryWakeAncestorSQL :: Text -> Text -> Int64 -> Query ()
tryWakeAncestorSQL Text
schema Text
tableName Int64
ancestorId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|
UPDATE ${tbl}
SET suspended = FALSE, updated_at = NOW()
WHERE id = #{ancestorId :: CInt8}
AND suspended = TRUE
AND NOT EXISTS (SELECT 1 FROM ${tbl} child WHERE child.parent_id = #{ancestorId :: CInt8})
|]
treeRollupIdsSQL :: Text -> Text -> Int64 -> Query Int64
treeRollupIdsSQL :: Text -> Text -> Int64 -> Query Int64
treeRollupIdsSQL Text
schema Text
tableName Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|
WITH RECURSIVE descendants AS (
SELECT id, parent_state FROM ${tbl} WHERE id = #{jobId :: CInt8}
UNION ALL
SELECT job.id, job.parent_state FROM ${tbl} job JOIN descendants descendant ON job.parent_id = descendant.id
)
SELECT id AS @{result :: CInt8} FROM descendants WHERE parent_state IS NOT NULL
|]
suspendJobSQL :: Text -> Text -> Int64 -> Query ()
suspendJobSQL :: Text -> Text -> Int64 -> Query ()
suspendJobSQL Text
schema Text
tableName Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|
UPDATE ${tbl}
SET suspended = TRUE, updated_at = NOW()
WHERE id = #{jobId :: CInt8}
AND NOT suspended
AND NOT (attempts > 0 AND not_visible_until IS NOT NULL AND not_visible_until > NOW())
|]
resumeJobSQL :: Text -> Text -> Int64 -> Query ()
resumeJobSQL :: Text -> Text -> Int64 -> Query ()
resumeJobSQL Text
schema Text
tableName Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|
UPDATE ${tbl}
SET suspended = FALSE, updated_at = NOW()
WHERE id = #{jobId :: CInt8} AND suspended = TRUE
AND NOT (
parent_state IS NOT NULL
AND EXISTS (SELECT 1 FROM ${tbl} child WHERE child.parent_id = ${tbl}.id)
)
|]
jobExistsSQL :: Text -> Text -> Int64 -> Query Bool
jobExistsSQL :: Text -> Text -> Int64 -> Query Bool
jobExistsSQL Text
schema Text
tableName Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|SELECT EXISTS (SELECT 1 FROM ${tbl} WHERE id = #{jobId :: CInt8}) AS @{result :: CBool}|]
getParentIdSQL :: Text -> Text -> Int64 -> Query (Maybe Int64)
getParentIdSQL :: Text -> Text -> Int64 -> Query (Maybe Int64)
getParentIdSQL Text
schema Text
tableName Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|SELECT @{parent_id :: Maybe CInt8} FROM ${tbl} WHERE id = #{jobId :: CInt8}|]
getParentIdsSQL :: Text -> Text -> [Int64] -> Query (Maybe Int64)
getParentIdsSQL :: Text -> Text -> [Int64] -> Query (Maybe Int64)
getParentIdsSQL Text
schema Text
tableName [Int64]
jobIds =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|SELECT @{parent_id :: Maybe CInt8} FROM ${tbl} WHERE id = ANY(#{jobIds :: [CInt8]})|]
insertResultSQL :: Text -> Text -> Int64 -> Int64 -> Value -> Query ()
insertResultSQL :: Text -> Text -> Int64 -> Int64 -> Value -> Query ()
insertResultSQL Text
schema Text
tableName Int64
parentId Int64
childId Value
result =
let resultsTbl :: Text
resultsTbl = Text -> Text -> Text
jobQueueResultsTable Text
schema Text
tableName
in [sql|
INSERT INTO ${resultsTbl} (parent_id, child_id, result)
VALUES (#{parentId :: CInt8}, #{childId :: CInt8}, #{result :: CJsonb})
ON CONFLICT (parent_id, child_id) DO UPDATE SET result = EXCLUDED.result
|]
insertResultsBatchSQL :: Text -> Text -> [Int64] -> [Int64] -> [Value] -> Query ()
insertResultsBatchSQL :: Text -> Text -> [Int64] -> [Int64] -> [Value] -> Query ()
insertResultsBatchSQL Text
schema Text
tableName [Int64]
parentIds [Int64]
childIds [Value]
results =
let resultsTbl :: Text
resultsTbl = Text -> Text -> Text
jobQueueResultsTable Text
schema Text
tableName
in [sql|
INSERT INTO ${resultsTbl} (parent_id, child_id, result)
SELECT parent_id, child_id, result FROM (
SELECT unnest(#{parentIds :: [CInt8]}::bigint[]) AS parent_id,
unnest(#{childIds :: [CInt8]}::bigint[]) AS child_id,
unnest(#{results :: [CJsonb]}::jsonb[]) AS result
) src
ON CONFLICT (parent_id, child_id) DO UPDATE SET result = EXCLUDED.result
|]
getResultsByParentSQL :: Text -> Text -> Int64 -> Query (Int64, Value)
getResultsByParentSQL :: Text -> Text -> Int64 -> Query (Int64, Value)
getResultsByParentSQL Text
schema Text
tableName Int64
parentId =
let resultsTbl :: Text
resultsTbl = Text -> Text -> Text
jobQueueResultsTable Text
schema Text
tableName
in [sql|
SELECT @{child_id :: CInt8}, @{result :: CJsonb} FROM ${resultsTbl} WHERE parent_id = #{parentId :: CInt8}
|]
getDLQChildErrorsByParentSQL :: Text -> Text -> Int64 -> Query (Int64, Maybe Text)
getDLQChildErrorsByParentSQL :: Text -> Text -> Int64 -> Query (Int64, Maybe Text)
getDLQChildErrorsByParentSQL Text
schema Text
tableName Int64
parentId =
let dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
in [sql|
SELECT @{job_id :: CInt8}, @{last_error :: Maybe CText} FROM ${dlqTbl} WHERE parent_id = #{parentId :: CInt8}
|]
persistParentStateSQL :: Text -> Text -> Value -> Int64 -> Query ()
persistParentStateSQL :: Text -> Text -> Value -> Int64 -> Query ()
persistParentStateSQL Text
schema Text
tableName Value
parentState Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|
UPDATE ${tbl} SET parent_state = #{parentState :: CJsonb}, updated_at = NOW() WHERE id = #{jobId :: CInt8}
|]
getParentStateSnapshotSQL :: Text -> Text -> Int64 -> Query (Maybe Value)
getParentStateSnapshotSQL :: Text -> Text -> Int64 -> Query (Maybe Value)
getParentStateSnapshotSQL Text
schema Text
tableName Int64
jobId =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|SELECT @{parent_state :: Maybe CJsonb} FROM ${tbl} WHERE id = #{jobId :: CInt8}|]
readChildResultsSQL :: Text -> Text -> Int64 -> Query (Text, Maybe Int64, Maybe Value, Maybe Text, Maybe Int64)
readChildResultsSQL :: Text
-> Text
-> Int64
-> Query (Text, Maybe Int64, Maybe Value, Maybe Text, Maybe Int64)
readChildResultsSQL Text
schema Text
tableName Int64
parentId =
let resultsTbl :: Text
resultsTbl = Text -> Text -> Text
jobQueueResultsTable Text
schema Text
tableName
dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|
SELECT 'r'::text AS @{source :: CText}, @{child_id :: Maybe CInt8}, @{result :: Maybe CJsonb},
NULL::text AS @{error :: Maybe CText}, NULL::bigint AS @{dlq_pk :: Maybe CInt8}
FROM ${resultsTbl} WHERE parent_id = #{parentId :: CInt8}
UNION ALL
SELECT 'e' AS source, job_id AS child_id, NULL::jsonb AS result, last_error AS error, id AS dlq_pk
FROM ${dlqTbl} WHERE parent_id = #{parentId :: CInt8}
UNION ALL
SELECT 's' AS source, NULL::bigint AS child_id, parent_state AS result, NULL::text AS error,
NULL::bigint AS dlq_pk
FROM ${tbl} WHERE id = #{parentId :: CInt8} AND parent_state IS NOT NULL
|]