{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}

-- | Jobs SQL templates.
module Arbiter.Core.Sql.Jobs
  ( JobFilter (..)
  , JobSortColumn (..)
  , jobSortColumnName
  , DLQSortColumn (..)
  , dlqSortColumnName
  , ArchiveSortColumn (..)
  , archiveSortColumnName
  , SortDir (..)
  , sortDirSql
  , throttledPredicateSQL
  , jobStatusCaseSQL
  , jobsWithStatusSubquery
  , listJobsFilteredSQL
  , listJobsWithStatusSQL
  , countJobsFilteredSQL
  , getJobByIdWithStatusSQL
  , listDLQFilteredSQL
  , NullsBehavior (..)
  , nullsClause
  , jobColumnNulls
  , dlqColumnNulls
  , buildJobsOrderBy
  , buildDLQOrderBy
  , buildArchiveOrderBy
  , countDLQFilteredSQL
  , allDLQColumns
  , jobColsExceptId
  , dlqCarriedCols
  , requeuedCols
  , enqueuedAgainCols
  , jobColumns
  , dedupUpdateSet
  , insertJobSQL
  , insertJobReplaceSQL
  , insertJobsBatchSQL
  , insertJobsBatchSQL_
  , insertJobsBatchBase
  , getJobByIdSQL
  , getJobByDedupKeySQL
  , cancelJobSQL
  , unionAllOverQueueTables
  ) where

import Data.Aeson (Value)
import Data.Int (Int64)
import Data.Maybe (fromMaybe)
import Data.Text (Text)
import Data.Text qualified as T
import Data.Time (UTCTime)
import Data.UUID.Types (UUID)
import NeatInterpolation (text)

import Arbiter.Core.Codec (dlqRowCodec, jobRowCodec)
import Arbiter.Core.Job.Schema
  ( SchemaName
  , TableName
  , jobQueueDLQTable
  , jobQueueTable
  )
import Arbiter.Core.Job.Types (JobRead, JobStatus)
import Arbiter.Core.Sql.QQ (sql)
import Arbiter.Core.Sql.Query (Query, mwhen, rows)

-- | A narrowing predicate on a job listing. @FilterJobId@ and the @completed_at@
-- range name DLQ and archive columns.
data JobFilter
  = FilterGroupKey Text
  | FilterParentId Int64
  | FilterRootsOnly
  | FilterStatus JobStatus
  | FilterId Int64
  | FilterJobId Int64
  | FilterClaimedBy UUID
  | FilterKind Text
  | FilterPayloadText Text
  | FilterRateLimitPrefix Text
  | FilterConcurrencyPrefix Text
  | FilterInsertedAfter UTCTime
  | FilterInsertedBefore UTCTime
  | FilterCompletedAfter UTCTime
  | FilterCompletedBefore UTCTime
  deriving stock (JobFilter -> JobFilter -> Bool
(JobFilter -> JobFilter -> Bool)
-> (JobFilter -> JobFilter -> Bool) -> Eq JobFilter
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: JobFilter -> JobFilter -> Bool
== :: JobFilter -> JobFilter -> Bool
$c/= :: JobFilter -> JobFilter -> Bool
/= :: JobFilter -> JobFilter -> Bool
Eq, Int -> JobFilter -> ShowS
[JobFilter] -> ShowS
JobFilter -> String
(Int -> JobFilter -> ShowS)
-> (JobFilter -> String)
-> ([JobFilter] -> ShowS)
-> Show JobFilter
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> JobFilter -> ShowS
showsPrec :: Int -> JobFilter -> ShowS
$cshow :: JobFilter -> String
show :: JobFilter -> String
$cshowList :: [JobFilter] -> ShowS
showList :: [JobFilter] -> ShowS
Show)

-- | Sortable columns on the main jobs table.
data JobSortColumn
  = JsId
  | JsPriority
  | JsAttempts
  | JsInsertedAt
  | JsNotVisibleUntil
  | JsGroupKey
  | JsParentId
  | JsLastAttemptedAt
  deriving stock (JobSortColumn
JobSortColumn -> JobSortColumn -> Bounded JobSortColumn
forall a. a -> a -> Bounded a
$cminBound :: JobSortColumn
minBound :: JobSortColumn
$cmaxBound :: JobSortColumn
maxBound :: JobSortColumn
Bounded, Int -> JobSortColumn
JobSortColumn -> Int
JobSortColumn -> [JobSortColumn]
JobSortColumn -> JobSortColumn
JobSortColumn -> JobSortColumn -> [JobSortColumn]
JobSortColumn -> JobSortColumn -> JobSortColumn -> [JobSortColumn]
(JobSortColumn -> JobSortColumn)
-> (JobSortColumn -> JobSortColumn)
-> (Int -> JobSortColumn)
-> (JobSortColumn -> Int)
-> (JobSortColumn -> [JobSortColumn])
-> (JobSortColumn -> JobSortColumn -> [JobSortColumn])
-> (JobSortColumn -> JobSortColumn -> [JobSortColumn])
-> (JobSortColumn
    -> JobSortColumn -> JobSortColumn -> [JobSortColumn])
-> Enum JobSortColumn
forall a.
(a -> a)
-> (a -> a)
-> (Int -> a)
-> (a -> Int)
-> (a -> [a])
-> (a -> a -> [a])
-> (a -> a -> [a])
-> (a -> a -> a -> [a])
-> Enum a
$csucc :: JobSortColumn -> JobSortColumn
succ :: JobSortColumn -> JobSortColumn
$cpred :: JobSortColumn -> JobSortColumn
pred :: JobSortColumn -> JobSortColumn
$ctoEnum :: Int -> JobSortColumn
toEnum :: Int -> JobSortColumn
$cfromEnum :: JobSortColumn -> Int
fromEnum :: JobSortColumn -> Int
$cenumFrom :: JobSortColumn -> [JobSortColumn]
enumFrom :: JobSortColumn -> [JobSortColumn]
$cenumFromThen :: JobSortColumn -> JobSortColumn -> [JobSortColumn]
enumFromThen :: JobSortColumn -> JobSortColumn -> [JobSortColumn]
$cenumFromTo :: JobSortColumn -> JobSortColumn -> [JobSortColumn]
enumFromTo :: JobSortColumn -> JobSortColumn -> [JobSortColumn]
$cenumFromThenTo :: JobSortColumn -> JobSortColumn -> JobSortColumn -> [JobSortColumn]
enumFromThenTo :: JobSortColumn -> JobSortColumn -> JobSortColumn -> [JobSortColumn]
Enum, JobSortColumn -> JobSortColumn -> Bool
(JobSortColumn -> JobSortColumn -> Bool)
-> (JobSortColumn -> JobSortColumn -> Bool) -> Eq JobSortColumn
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: JobSortColumn -> JobSortColumn -> Bool
== :: JobSortColumn -> JobSortColumn -> Bool
$c/= :: JobSortColumn -> JobSortColumn -> Bool
/= :: JobSortColumn -> JobSortColumn -> Bool
Eq, Int -> JobSortColumn -> ShowS
[JobSortColumn] -> ShowS
JobSortColumn -> String
(Int -> JobSortColumn -> ShowS)
-> (JobSortColumn -> String)
-> ([JobSortColumn] -> ShowS)
-> Show JobSortColumn
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> JobSortColumn -> ShowS
showsPrec :: Int -> JobSortColumn -> ShowS
$cshow :: JobSortColumn -> String
show :: JobSortColumn -> String
$cshowList :: [JobSortColumn] -> ShowS
showList :: [JobSortColumn] -> ShowS
Show)

-- | Underlying SQL column name for a 'JobSortColumn'.
jobSortColumnName :: JobSortColumn -> Text
jobSortColumnName :: JobSortColumn -> Text
jobSortColumnName = \case
  JobSortColumn
JsId -> Text
"id"
  JobSortColumn
JsPriority -> Text
"priority"
  JobSortColumn
JsAttempts -> Text
"attempts"
  JobSortColumn
JsInsertedAt -> Text
"inserted_at"
  JobSortColumn
JsNotVisibleUntil -> Text
"not_visible_until"
  JobSortColumn
JsGroupKey -> Text
"group_key"
  JobSortColumn
JsParentId -> Text
"parent_id"
  JobSortColumn
JsLastAttemptedAt -> Text
"last_attempted_at"

-- | Sortable columns on the DLQ table. @DlqId@ is the DLQ primary key and
-- @DlqJobId@ the original job id.
data DLQSortColumn
  = DlqId
  | DlqFailedAt
  | DlqJobId
  | DlqPriority
  | DlqAttempts
  | DlqInsertedAt
  | DlqGroupKey
  | DlqParentId
  | DlqLastAttemptedAt
  deriving stock (DLQSortColumn
DLQSortColumn -> DLQSortColumn -> Bounded DLQSortColumn
forall a. a -> a -> Bounded a
$cminBound :: DLQSortColumn
minBound :: DLQSortColumn
$cmaxBound :: DLQSortColumn
maxBound :: DLQSortColumn
Bounded, Int -> DLQSortColumn
DLQSortColumn -> Int
DLQSortColumn -> [DLQSortColumn]
DLQSortColumn -> DLQSortColumn
DLQSortColumn -> DLQSortColumn -> [DLQSortColumn]
DLQSortColumn -> DLQSortColumn -> DLQSortColumn -> [DLQSortColumn]
(DLQSortColumn -> DLQSortColumn)
-> (DLQSortColumn -> DLQSortColumn)
-> (Int -> DLQSortColumn)
-> (DLQSortColumn -> Int)
-> (DLQSortColumn -> [DLQSortColumn])
-> (DLQSortColumn -> DLQSortColumn -> [DLQSortColumn])
-> (DLQSortColumn -> DLQSortColumn -> [DLQSortColumn])
-> (DLQSortColumn
    -> DLQSortColumn -> DLQSortColumn -> [DLQSortColumn])
-> Enum DLQSortColumn
forall a.
(a -> a)
-> (a -> a)
-> (Int -> a)
-> (a -> Int)
-> (a -> [a])
-> (a -> a -> [a])
-> (a -> a -> [a])
-> (a -> a -> a -> [a])
-> Enum a
$csucc :: DLQSortColumn -> DLQSortColumn
succ :: DLQSortColumn -> DLQSortColumn
$cpred :: DLQSortColumn -> DLQSortColumn
pred :: DLQSortColumn -> DLQSortColumn
$ctoEnum :: Int -> DLQSortColumn
toEnum :: Int -> DLQSortColumn
$cfromEnum :: DLQSortColumn -> Int
fromEnum :: DLQSortColumn -> Int
$cenumFrom :: DLQSortColumn -> [DLQSortColumn]
enumFrom :: DLQSortColumn -> [DLQSortColumn]
$cenumFromThen :: DLQSortColumn -> DLQSortColumn -> [DLQSortColumn]
enumFromThen :: DLQSortColumn -> DLQSortColumn -> [DLQSortColumn]
$cenumFromTo :: DLQSortColumn -> DLQSortColumn -> [DLQSortColumn]
enumFromTo :: DLQSortColumn -> DLQSortColumn -> [DLQSortColumn]
$cenumFromThenTo :: DLQSortColumn -> DLQSortColumn -> DLQSortColumn -> [DLQSortColumn]
enumFromThenTo :: DLQSortColumn -> DLQSortColumn -> DLQSortColumn -> [DLQSortColumn]
Enum, DLQSortColumn -> DLQSortColumn -> Bool
(DLQSortColumn -> DLQSortColumn -> Bool)
-> (DLQSortColumn -> DLQSortColumn -> Bool) -> Eq DLQSortColumn
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: DLQSortColumn -> DLQSortColumn -> Bool
== :: DLQSortColumn -> DLQSortColumn -> Bool
$c/= :: DLQSortColumn -> DLQSortColumn -> Bool
/= :: DLQSortColumn -> DLQSortColumn -> Bool
Eq, Int -> DLQSortColumn -> ShowS
[DLQSortColumn] -> ShowS
DLQSortColumn -> String
(Int -> DLQSortColumn -> ShowS)
-> (DLQSortColumn -> String)
-> ([DLQSortColumn] -> ShowS)
-> Show DLQSortColumn
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> DLQSortColumn -> ShowS
showsPrec :: Int -> DLQSortColumn -> ShowS
$cshow :: DLQSortColumn -> String
show :: DLQSortColumn -> String
$cshowList :: [DLQSortColumn] -> ShowS
showList :: [DLQSortColumn] -> ShowS
Show)

-- | The DLQ column a sort key names.
dlqSortColumnName :: DLQSortColumn -> Text
dlqSortColumnName :: DLQSortColumn -> Text
dlqSortColumnName = \case
  DLQSortColumn
DlqId -> Text
"id"
  DLQSortColumn
DlqFailedAt -> Text
"failed_at"
  DLQSortColumn
DlqJobId -> Text
"job_id"
  DLQSortColumn
DlqPriority -> Text
"priority"
  DLQSortColumn
DlqAttempts -> Text
"attempts"
  DLQSortColumn
DlqInsertedAt -> Text
"inserted_at"
  DLQSortColumn
DlqGroupKey -> Text
"group_key"
  DLQSortColumn
DlqParentId -> Text
"parent_id"
  DLQSortColumn
DlqLastAttemptedAt -> Text
"last_attempted_at"

-- | Sortable columns on the archive table. @ArchiveId@ is the archive primary
-- key and @ArchiveJobId@ the original job id.
data ArchiveSortColumn
  = ArchiveId
  | ArchiveCompletedAt
  | ArchiveInsertedAt
  | ArchiveJobId
  | ArchiveAttempts
  | ArchiveGroupKey
  | ArchiveParentId
  deriving stock (ArchiveSortColumn
ArchiveSortColumn -> ArchiveSortColumn -> Bounded ArchiveSortColumn
forall a. a -> a -> Bounded a
$cminBound :: ArchiveSortColumn
minBound :: ArchiveSortColumn
$cmaxBound :: ArchiveSortColumn
maxBound :: ArchiveSortColumn
Bounded, Int -> ArchiveSortColumn
ArchiveSortColumn -> Int
ArchiveSortColumn -> [ArchiveSortColumn]
ArchiveSortColumn -> ArchiveSortColumn
ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn]
ArchiveSortColumn
-> ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn]
(ArchiveSortColumn -> ArchiveSortColumn)
-> (ArchiveSortColumn -> ArchiveSortColumn)
-> (Int -> ArchiveSortColumn)
-> (ArchiveSortColumn -> Int)
-> (ArchiveSortColumn -> [ArchiveSortColumn])
-> (ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn])
-> (ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn])
-> (ArchiveSortColumn
    -> ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn])
-> Enum ArchiveSortColumn
forall a.
(a -> a)
-> (a -> a)
-> (Int -> a)
-> (a -> Int)
-> (a -> [a])
-> (a -> a -> [a])
-> (a -> a -> [a])
-> (a -> a -> a -> [a])
-> Enum a
$csucc :: ArchiveSortColumn -> ArchiveSortColumn
succ :: ArchiveSortColumn -> ArchiveSortColumn
$cpred :: ArchiveSortColumn -> ArchiveSortColumn
pred :: ArchiveSortColumn -> ArchiveSortColumn
$ctoEnum :: Int -> ArchiveSortColumn
toEnum :: Int -> ArchiveSortColumn
$cfromEnum :: ArchiveSortColumn -> Int
fromEnum :: ArchiveSortColumn -> Int
$cenumFrom :: ArchiveSortColumn -> [ArchiveSortColumn]
enumFrom :: ArchiveSortColumn -> [ArchiveSortColumn]
$cenumFromThen :: ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn]
enumFromThen :: ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn]
$cenumFromTo :: ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn]
enumFromTo :: ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn]
$cenumFromThenTo :: ArchiveSortColumn
-> ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn]
enumFromThenTo :: ArchiveSortColumn
-> ArchiveSortColumn -> ArchiveSortColumn -> [ArchiveSortColumn]
Enum, ArchiveSortColumn -> ArchiveSortColumn -> Bool
(ArchiveSortColumn -> ArchiveSortColumn -> Bool)
-> (ArchiveSortColumn -> ArchiveSortColumn -> Bool)
-> Eq ArchiveSortColumn
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: ArchiveSortColumn -> ArchiveSortColumn -> Bool
== :: ArchiveSortColumn -> ArchiveSortColumn -> Bool
$c/= :: ArchiveSortColumn -> ArchiveSortColumn -> Bool
/= :: ArchiveSortColumn -> ArchiveSortColumn -> Bool
Eq, Int -> ArchiveSortColumn -> ShowS
[ArchiveSortColumn] -> ShowS
ArchiveSortColumn -> String
(Int -> ArchiveSortColumn -> ShowS)
-> (ArchiveSortColumn -> String)
-> ([ArchiveSortColumn] -> ShowS)
-> Show ArchiveSortColumn
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> ArchiveSortColumn -> ShowS
showsPrec :: Int -> ArchiveSortColumn -> ShowS
$cshow :: ArchiveSortColumn -> String
show :: ArchiveSortColumn -> String
$cshowList :: [ArchiveSortColumn] -> ShowS
showList :: [ArchiveSortColumn] -> ShowS
Show)

-- | The archive column a sort key names.
archiveSortColumnName :: ArchiveSortColumn -> Text
archiveSortColumnName :: ArchiveSortColumn -> Text
archiveSortColumnName = \case
  ArchiveSortColumn
ArchiveId -> Text
"id"
  ArchiveSortColumn
ArchiveCompletedAt -> Text
"completed_at"
  ArchiveSortColumn
ArchiveInsertedAt -> Text
"inserted_at"
  ArchiveSortColumn
ArchiveJobId -> Text
"job_id"
  ArchiveSortColumn
ArchiveAttempts -> Text
"attempts"
  ArchiveSortColumn
ArchiveGroupKey -> Text
"group_key"
  ArchiveSortColumn
ArchiveParentId -> Text
"parent_id"

-- | Sort direction.
data SortDir = SortAsc | SortDesc
  deriving stock (SortDir
SortDir -> SortDir -> Bounded SortDir
forall a. a -> a -> Bounded a
$cminBound :: SortDir
minBound :: SortDir
$cmaxBound :: SortDir
maxBound :: SortDir
Bounded, Int -> SortDir
SortDir -> Int
SortDir -> [SortDir]
SortDir -> SortDir
SortDir -> SortDir -> [SortDir]
SortDir -> SortDir -> SortDir -> [SortDir]
(SortDir -> SortDir)
-> (SortDir -> SortDir)
-> (Int -> SortDir)
-> (SortDir -> Int)
-> (SortDir -> [SortDir])
-> (SortDir -> SortDir -> [SortDir])
-> (SortDir -> SortDir -> [SortDir])
-> (SortDir -> SortDir -> SortDir -> [SortDir])
-> Enum SortDir
forall a.
(a -> a)
-> (a -> a)
-> (Int -> a)
-> (a -> Int)
-> (a -> [a])
-> (a -> a -> [a])
-> (a -> a -> [a])
-> (a -> a -> a -> [a])
-> Enum a
$csucc :: SortDir -> SortDir
succ :: SortDir -> SortDir
$cpred :: SortDir -> SortDir
pred :: SortDir -> SortDir
$ctoEnum :: Int -> SortDir
toEnum :: Int -> SortDir
$cfromEnum :: SortDir -> Int
fromEnum :: SortDir -> Int
$cenumFrom :: SortDir -> [SortDir]
enumFrom :: SortDir -> [SortDir]
$cenumFromThen :: SortDir -> SortDir -> [SortDir]
enumFromThen :: SortDir -> SortDir -> [SortDir]
$cenumFromTo :: SortDir -> SortDir -> [SortDir]
enumFromTo :: SortDir -> SortDir -> [SortDir]
$cenumFromThenTo :: SortDir -> SortDir -> SortDir -> [SortDir]
enumFromThenTo :: SortDir -> SortDir -> SortDir -> [SortDir]
Enum, SortDir -> SortDir -> Bool
(SortDir -> SortDir -> Bool)
-> (SortDir -> SortDir -> Bool) -> Eq SortDir
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: SortDir -> SortDir -> Bool
== :: SortDir -> SortDir -> Bool
$c/= :: SortDir -> SortDir -> Bool
/= :: SortDir -> SortDir -> Bool
Eq, Int -> SortDir -> ShowS
[SortDir] -> ShowS
SortDir -> String
(Int -> SortDir -> ShowS)
-> (SortDir -> String) -> ([SortDir] -> ShowS) -> Show SortDir
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> SortDir -> ShowS
showsPrec :: Int -> SortDir -> ShowS
$cshow :: SortDir -> String
show :: SortDir -> String
$cshowList :: [SortDir] -> ShowS
showList :: [SortDir] -> ShowS
Show)

-- | A sort direction as SQL.
sortDirSql :: SortDir -> Text
sortDirSql :: SortDir -> Text
sortDirSql SortDir
SortAsc = Text
"ASC"
sortDirSql SortDir
SortDesc = Text
"DESC"

-- | A throttle-deferred job with a live marker, still parked. Shared by status,
-- count, and wake. Claiming clears the marker.
throttledPredicateSQL :: Text
throttledPredicateSQL :: Text
throttledPredicateSQL =
  Text
"throttled_until > NOW() AND not_visible_until > NOW()"

-- | The derived job status. Its string values match 'Arbiter.Core.Job.Status.jobStatusToText'.
jobStatusCaseSQL :: Text
jobStatusCaseSQL :: Text
jobStatusCaseSQL =
  [text|
    CASE
      WHEN cancel_requested_at IS NOT NULL THEN 'cancelled'
      WHEN suspended THEN 'suspended'
      WHEN ${throttledPredicateSQL} THEN 'throttled'
      WHEN claimed_by IS NULL AND attempts > 0
           AND not_visible_until IS NOT NULL AND not_visible_until > NOW() THEN 'backoff'
      WHEN attempts > 0 AND not_visible_until IS NOT NULL AND not_visible_until > NOW() THEN 'in_flight'
      WHEN not_visible_until IS NOT NULL AND not_visible_until > NOW() THEN 'scheduled'
      ELSE 'ready'
    END
  |]

-- | All job columns plus the derived @status@ column, aliased @job@, for filtering.
jobsWithStatusSubquery :: Text -> Text -> Text
jobsWithStatusSubquery :: Text -> Text -> Text
jobsWithStatusSubquery Text
schema Text
tableName =
  let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
   in [text|(SELECT ${jobColumns}, ${jobStatusCaseSQL} AS status FROM ${tbl}) job|]

-- | List filtered jobs without the derived status.
listJobsFilteredSQL :: Text -> Text -> Query () -> Text -> Int64 -> Int64 -> Query (JobRead Value)
listJobsFilteredSQL :: Text
-> Text
-> Query ()
-> Text
-> Int64
-> Int64
-> Query (JobRead Value)
listJobsFilteredSQL Text
schema Text
tableName Query ()
whereFrag Text
orderBy Int64
limit Int64
offset =
  let 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|
          SELECT ${jobColumns}
          FROM ${tbl}
          ${whereFrag}
          ORDER BY ${orderBy} LIMIT #{limit :: CInt8} OFFSET #{offset :: CInt8}
        |]

-- | List filtered jobs with @status@ as a trailing column. The caller attaches the row decoder.
listJobsWithStatusSQL :: Text -> Text -> Query () -> Text -> Int64 -> Int64 -> Query ()
listJobsWithStatusSQL :: Text -> Text -> Query () -> Text -> Int64 -> Int64 -> Query ()
listJobsWithStatusSQL Text
schema Text
tableName Query ()
whereFrag Text
orderBy Int64
limit Int64
offset =
  let sub :: Text
sub = Text -> Text -> Text
jobsWithStatusSubquery Text
schema Text
tableName
   in [sql|
        SELECT * FROM ${sub}
        ${whereFrag}
        ORDER BY ${orderBy} LIMIT #{limit :: CInt8} OFFSET #{offset :: CInt8}
      |]

-- | Count filtered jobs through the status subquery.
countJobsFilteredSQL :: Text -> Text -> Query () -> Query Int64
countJobsFilteredSQL :: Text -> Text -> Query () -> Query Int64
countJobsFilteredSQL Text
schema Text
tableName Query ()
whereFrag =
  let sub :: Text
sub = Text -> Text -> Text
jobsWithStatusSubquery Text
schema Text
tableName
   in [sql|SELECT COUNT(*) AS @{count :: CInt8} FROM ${sub} ${whereFrag}|]

-- | Fetch a single job by id with its derived @status@ trailing column. The caller attaches the row decoder.
getJobByIdWithStatusSQL :: Text -> Text -> Int64 -> Query ()
getJobByIdWithStatusSQL :: Text -> Text -> Int64 -> Query ()
getJobByIdWithStatusSQL Text
schema Text
tableName Int64
jobId =
  let sub :: Text
sub = Text -> Text -> Text
jobsWithStatusSubquery Text
schema Text
tableName
   in [sql|SELECT * FROM ${sub} WHERE id = #{jobId :: CInt8}|]

-- | List DLQ jobs under a dynamic WHERE and an @orderBy@ from 'buildDLQOrderBy'.
listDLQFilteredSQL :: Text -> Text -> Query () -> Text -> Int64 -> Int64 -> Query (Int64, UTCTime, JobRead Value)
listDLQFilteredSQL :: Text
-> Text
-> Query ()
-> Text
-> Int64
-> Int64
-> Query (Int64, UTCTime, JobRead Value)
listDLQFilteredSQL Text
schema Text
tableName Query ()
whereFrag Text
orderBy Int64
limit Int64
offset =
  let dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
   in RowCodec (Int64, UTCTime, JobRead Value)
-> Query () -> Query (Int64, UTCTime, JobRead Value)
forall a. RowCodec a -> Query () -> Query a
rows
        (Text -> RowCodec (Int64, UTCTime, JobRead Value)
dlqRowCodec Text
tableName)
        [sql|
          SELECT ${allDLQColumns}
          FROM ${dlqTbl}
          ${whereFrag}
          ORDER BY ${orderBy}
          LIMIT #{limit :: CInt8} OFFSET #{offset :: CInt8}
        |]

-- | How NULLs in a sort column should order relative to non-NULL values.
data NullsBehavior
  = -- | Column is NOT NULL. Emit no NULLS clause.
    NullsNotApplicable
  | -- | NULL means absent. Always sorts last.
    NullsAsAbsent
  | -- | NULL is the column's minimum (e.g. @not_visible_until IS NULL@ means visible now).
    -- NULLS FIRST when ASC, NULLS LAST when DESC.
    NullsAsMinimum

-- | The @NULLS@ placement for a column's null semantics and direction.
nullsClause :: NullsBehavior -> SortDir -> Text
nullsClause :: NullsBehavior -> SortDir -> Text
nullsClause NullsBehavior
NullsNotApplicable SortDir
_ = Text
""
nullsClause NullsBehavior
NullsAsAbsent SortDir
_ = Text
" NULLS LAST"
nullsClause NullsBehavior
NullsAsMinimum SortDir
SortAsc = Text
" NULLS FIRST"
nullsClause NullsBehavior
NullsAsMinimum SortDir
SortDesc = Text
" NULLS LAST"

-- | How nulls sort for a job column.
jobColumnNulls :: JobSortColumn -> NullsBehavior
jobColumnNulls :: JobSortColumn -> NullsBehavior
jobColumnNulls = \case
  JobSortColumn
JsId -> NullsBehavior
NullsNotApplicable
  JobSortColumn
JsPriority -> NullsBehavior
NullsNotApplicable
  JobSortColumn
JsAttempts -> NullsBehavior
NullsNotApplicable
  JobSortColumn
JsInsertedAt -> NullsBehavior
NullsNotApplicable
  JobSortColumn
JsNotVisibleUntil -> NullsBehavior
NullsAsMinimum
  JobSortColumn
JsGroupKey -> NullsBehavior
NullsAsAbsent
  JobSortColumn
JsParentId -> NullsBehavior
NullsAsAbsent
  JobSortColumn
JsLastAttemptedAt -> NullsBehavior
NullsAsMinimum

-- | How nulls sort for a DLQ column.
dlqColumnNulls :: DLQSortColumn -> NullsBehavior
dlqColumnNulls :: DLQSortColumn -> NullsBehavior
dlqColumnNulls = \case
  DLQSortColumn
DlqId -> NullsBehavior
NullsNotApplicable
  DLQSortColumn
DlqFailedAt -> NullsBehavior
NullsNotApplicable
  DLQSortColumn
DlqJobId -> NullsBehavior
NullsNotApplicable
  DLQSortColumn
DlqInsertedAt -> NullsBehavior
NullsNotApplicable
  DLQSortColumn
DlqPriority -> NullsBehavior
NullsNotApplicable
  DLQSortColumn
DlqAttempts -> NullsBehavior
NullsNotApplicable
  DLQSortColumn
DlqGroupKey -> NullsBehavior
NullsAsAbsent
  DLQSortColumn
DlqParentId -> NullsBehavior
NullsAsAbsent
  DLQSortColumn
DlqLastAttemptedAt -> NullsBehavior
NullsAsMinimum

archiveColumnNulls :: ArchiveSortColumn -> NullsBehavior
archiveColumnNulls :: ArchiveSortColumn -> NullsBehavior
archiveColumnNulls = \case
  ArchiveSortColumn
ArchiveId -> NullsBehavior
NullsNotApplicable
  ArchiveSortColumn
ArchiveCompletedAt -> NullsBehavior
NullsNotApplicable
  ArchiveSortColumn
ArchiveInsertedAt -> NullsBehavior
NullsNotApplicable
  ArchiveSortColumn
ArchiveJobId -> NullsBehavior
NullsNotApplicable
  ArchiveSortColumn
ArchiveAttempts -> NullsBehavior
NullsNotApplicable
  ArchiveSortColumn
ArchiveGroupKey -> NullsBehavior
NullsAsAbsent
  ArchiveSortColumn
ArchiveParentId -> NullsBehavior
NullsAsAbsent

-- | Build an ORDER BY clause from a typed sort spec. @nullsFn@ places NULLs per
-- column. A stable @id@ tie-breaker in the primary sort's direction follows any
-- sort column other than @idCol@.
buildOrderBy :: (Eq a) => (a -> Text) -> (a -> NullsBehavior) -> a -> a -> Maybe a -> Maybe SortDir -> Text
buildOrderBy :: forall a.
Eq a =>
(a -> Text)
-> (a -> NullsBehavior)
-> a
-> a
-> Maybe a
-> Maybe SortDir
-> Text
buildOrderBy a -> Text
nameFn a -> NullsBehavior
nullsFn a
defCol a
idCol Maybe a
mCol Maybe SortDir
mDir =
  let sortCol :: a
sortCol = a -> Maybe a -> a
forall a. a -> Maybe a -> a
fromMaybe a
defCol Maybe a
mCol
      dir :: SortDir
dir = SortDir -> Maybe SortDir -> SortDir
forall a. a -> Maybe a -> a
fromMaybe SortDir
SortDesc Maybe SortDir
mDir
      columnName :: Text
columnName = a -> Text
nameFn a
sortCol
      dirText :: Text
dirText = SortDir -> Text
sortDirSql SortDir
dir
      nulls :: Text
nulls = NullsBehavior -> SortDir -> Text
nullsClause (a -> NullsBehavior
nullsFn a
sortCol) SortDir
dir
      tieBreaker :: Text
tieBreaker = Bool -> Text -> Text
forall m. Monoid m => Bool -> m -> m
mwhen (a
sortCol a -> a -> Bool
forall a. Eq a => a -> a -> Bool
/= a
idCol) [text|, id ${dirText}|]
   in [text|${columnName} ${dirText}${nulls}${tieBreaker}|]

-- | ORDER BY for the jobs table. Defaults to @id DESC@.
buildJobsOrderBy :: Maybe JobSortColumn -> Maybe SortDir -> Text
buildJobsOrderBy :: Maybe JobSortColumn -> Maybe SortDir -> Text
buildJobsOrderBy = (JobSortColumn -> Text)
-> (JobSortColumn -> NullsBehavior)
-> JobSortColumn
-> JobSortColumn
-> Maybe JobSortColumn
-> Maybe SortDir
-> Text
forall a.
Eq a =>
(a -> Text)
-> (a -> NullsBehavior)
-> a
-> a
-> Maybe a
-> Maybe SortDir
-> Text
buildOrderBy JobSortColumn -> Text
jobSortColumnName JobSortColumn -> NullsBehavior
jobColumnNulls JobSortColumn
JsId JobSortColumn
JsId

-- | ORDER BY for the DLQ table. Defaults to @failed_at DESC@.
buildDLQOrderBy :: Maybe DLQSortColumn -> Maybe SortDir -> Text
buildDLQOrderBy :: Maybe DLQSortColumn -> Maybe SortDir -> Text
buildDLQOrderBy = (DLQSortColumn -> Text)
-> (DLQSortColumn -> NullsBehavior)
-> DLQSortColumn
-> DLQSortColumn
-> Maybe DLQSortColumn
-> Maybe SortDir
-> Text
forall a.
Eq a =>
(a -> Text)
-> (a -> NullsBehavior)
-> a
-> a
-> Maybe a
-> Maybe SortDir
-> Text
buildOrderBy DLQSortColumn -> Text
dlqSortColumnName DLQSortColumn -> NullsBehavior
dlqColumnNulls DLQSortColumn
DlqFailedAt DLQSortColumn
DlqId

-- | ORDER BY for the archive table. Defaults to @completed_at DESC@.
buildArchiveOrderBy :: Maybe ArchiveSortColumn -> Maybe SortDir -> Text
buildArchiveOrderBy :: Maybe ArchiveSortColumn -> Maybe SortDir -> Text
buildArchiveOrderBy = (ArchiveSortColumn -> Text)
-> (ArchiveSortColumn -> NullsBehavior)
-> ArchiveSortColumn
-> ArchiveSortColumn
-> Maybe ArchiveSortColumn
-> Maybe SortDir
-> Text
forall a.
Eq a =>
(a -> Text)
-> (a -> NullsBehavior)
-> a
-> a
-> Maybe a
-> Maybe SortDir
-> Text
buildOrderBy ArchiveSortColumn -> Text
archiveSortColumnName ArchiveSortColumn -> NullsBehavior
archiveColumnNulls ArchiveSortColumn
ArchiveCompletedAt ArchiveSortColumn
ArchiveId

-- | Count DLQ jobs under a dynamic WHERE.
countDLQFilteredSQL :: Text -> Text -> Query () -> Query Int64
countDLQFilteredSQL :: Text -> Text -> Query () -> Query Int64
countDLQFilteredSQL Text
schema Text
tableName Query ()
whereFrag =
  let dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
   in [sql|SELECT COUNT(*) AS @{count :: CInt8} FROM ${dlqTbl} ${whereFrag}|]

-- | The job read columns, in codec order, for SELECT and RETURNING.
jobColumns :: Text
jobColumns :: Text
jobColumns =
  [text|
    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
  |]

-- | The DLQ read columns, in codec order. The DLQ uses @job_id@ for the main-table @id@.
allDLQColumns :: Text
allDLQColumns :: Text
allDLQColumns =
  [text|
    id, failed_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
  |]

-- | The job read columns except @id@. The archive INSERT copies them with the main
-- table's @id@ as @job_id@.
jobColsExceptId :: Text
jobColsExceptId :: Text
jobColsExceptId =
  [text|
    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
  |]

-- | Job columns carried through a DLQ round-trip. The read columns except @id@ and
-- @last_error@, plus write-only @rate_limit_cost@.
dlqCarriedCols :: Text
dlqCarriedCols :: Text
dlqCarriedCols =
  [text|
    payload, group_key, inserted_at, updated_at, attempts, 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,
    rate_limit_cost
  |]

-- | Job columns a DLQ retry carries back to the main table. The retry re-arms the rest.
requeuedCols :: Text
requeuedCols :: Text
requeuedCols =
  [text|
    payload, group_key, priority, max_attempts, parent_id, parent_state, traceparent, tracestate,
    archive_for, kind, rate_limit_key, rate_limit_prefix, concurrency_key, concurrency_prefix,
    rate_limit_cost
  |]

-- | 'requeuedCols' for an archive re-enqueue, without the parent link.
enqueuedAgainCols :: Text
enqueuedAgainCols :: Text
enqueuedAgainCols =
  [text|
    payload, group_key, priority, max_attempts, traceparent, tracestate,
    archive_for, kind, rate_limit_key, rate_limit_prefix, concurrency_key, concurrency_prefix,
    rate_limit_cost
  |]

-- | @DO UPDATE SET@ body for a replace-dedup upsert. Copies each writable column
-- from the excluded row, then re-arms the replaced job for a fresh run.
dedupUpdateSet :: Text -> Text
dedupUpdateSet :: Text -> Text
dedupUpdateSet Text
tbl =
  [text|
    payload = EXCLUDED.payload,
    group_key = EXCLUDED.group_key,
    priority = EXCLUDED.priority,
    not_visible_until = EXCLUDED.not_visible_until,
    dedup_strategy = EXCLUDED.dedup_strategy,
    max_attempts = EXCLUDED.max_attempts,
    parent_id = EXCLUDED.parent_id,
    parent_state = EXCLUDED.parent_state,
    traceparent = EXCLUDED.traceparent,
    tracestate = EXCLUDED.tracestate,
    suspended = EXCLUDED.suspended,
    archive_for = EXCLUDED.archive_for,
    kind = EXCLUDED.kind,
    rate_limit_key = EXCLUDED.rate_limit_key,
    rate_limit_prefix = EXCLUDED.rate_limit_prefix,
    concurrency_key = EXCLUDED.concurrency_key,
    concurrency_prefix = EXCLUDED.concurrency_prefix,
    rate_limit_cost = EXCLUDED.rate_limit_cost,
    attempts = 0,
    claim_seq = ${tbl}.claim_seq + 1,
    last_error = NULL,
    updated_at = NOW(),
    throttled_until = NULL,
    last_attempted_at = NULL,
    claimed_by = NULL
  |]

-- | An existing row is replaceable when idle, unflagged and childless.
replaceableGuard :: Text -> Text -> Text
replaceableGuard :: Text -> Text -> Text
replaceableGuard Text
tbl Text
dlqTbl =
  [text|
    (${tbl}.attempts = 0
      OR ${tbl}.claimed_by IS NULL
      OR ${tbl}.not_visible_until IS NULL
      OR ${tbl}.not_visible_until <= NOW())
      AND ${tbl}.cancel_requested_at IS NULL
      AND NOT EXISTS (SELECT 1 FROM ${tbl} child WHERE child.parent_id = ${tbl}.id)
      AND NOT EXISTS (SELECT 1 FROM ${dlqTbl} dlq_child WHERE dlq_child.parent_id = ${tbl}.id)
  |]

-- | Insert a job. The write fragment carries the column list and parameters.
insertJobSQL :: SchemaName -> TableName -> Query () -> Query (JobRead Value)
insertJobSQL :: Text -> Text -> Query () -> Query (JobRead Value)
insertJobSQL Text
schema Text
tableName Query ()
valuesFrag =
  let 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} ${valuesFrag}
          ON CONFLICT (dedup_key) WHERE dedup_key IS NOT NULL DO NOTHING
          RETURNING ${jobColumns}
        |]

-- | Insert under the replace dedup strategy. Replaces an existing job that is idle
-- and childless in the main queue and the DLQ. @ON CONFLICT DO UPDATE@ fires the
-- groups UPDATE trigger, which maintains a cross-group move.
insertJobReplaceSQL :: SchemaName -> TableName -> Query () -> Query (JobRead Value)
insertJobReplaceSQL :: Text -> Text -> Query () -> Query (JobRead Value)
insertJobReplaceSQL Text
schema Text
tableName Query ()
valuesFrag =
  let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
      dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
      guard :: Text
guard = Text -> Text -> Text
replaceableGuard Text
tbl Text
dlqTbl
      dedupSet :: Text
dedupSet = Text -> Text
dedupUpdateSet Text
tbl
   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} ${valuesFrag}
          ON CONFLICT (dedup_key) WHERE dedup_key IS NOT NULL DO UPDATE SET
            ${dedupSet}
          WHERE ${guard}
          RETURNING ${jobColumns}
        |]

-- | Batch insert over @unnest@ed parallel arrays. An ignore-dedup job is skipped on
-- conflict. A replace-dedup job updates an idle existing row.
insertJobsBatchSQL :: SchemaName -> TableName -> Query () -> Query (JobRead Value)
insertJobsBatchSQL :: Text -> Text -> Query () -> Query (JobRead Value)
insertJobsBatchSQL Text
schema Text
tableName Query ()
batchSrc =
  RowCodec (JobRead Value) -> Query () -> Query (JobRead Value)
forall a. RowCodec a -> Query () -> Query a
rows (Text -> RowCodec (JobRead Value)
jobRowCodec Text
tableName) (Text -> Text -> Query () -> Text -> Query ()
insertJobsBatchBase Text
schema Text
tableName Query ()
batchSrc [text|RETURNING ${jobColumns}|])

-- | 'insertJobsBatchBase' with no @RETURNING@.
insertJobsBatchSQL_ :: SchemaName -> TableName -> Query () -> Query ()
insertJobsBatchSQL_ :: Text -> Text -> Query () -> Query ()
insertJobsBatchSQL_ Text
schema Text
tableName Query ()
batchSrc =
  Text -> Text -> Query () -> Text -> Query ()
insertJobsBatchBase Text
schema Text
tableName Query ()
batchSrc Text
""

-- | Batch insert with dedup and replaceable-job handling, plus a caller's @RETURNING@.
insertJobsBatchBase :: SchemaName -> TableName -> Query () -> Text -> Query ()
insertJobsBatchBase :: Text -> Text -> Query () -> Text -> Query ()
insertJobsBatchBase Text
schema Text
tableName Query ()
batchSrc Text
returning =
  let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
      dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
      guard :: Text
guard = Text -> Text -> Text
replaceableGuard Text
tbl Text
dlqTbl
      dedupSet :: Text
dedupSet = Text -> Text
dedupUpdateSet Text
tbl
   in [sql|
        INSERT INTO ${tbl} ${batchSrc}
        WHERE (src.parent_id IS NULL
            OR EXISTS (SELECT 1 FROM ${tbl} parent WHERE parent.id = src.parent_id))
        ON CONFLICT (dedup_key) WHERE dedup_key IS NOT NULL DO UPDATE SET
          ${dedupSet}
        WHERE EXCLUDED.dedup_strategy = 'replace'
          AND ${guard}
        ${returning}
      |]

-- ---------------------------------------------------------------------------
-- Admin Operations
-- ---------------------------------------------------------------------------

-- | Fetch a job by id.
getJobByIdSQL :: Text -> Text -> Int64 -> Query (JobRead Value)
getJobByIdSQL :: Text -> Text -> Int64 -> Query (JobRead Value)
getJobByIdSQL Text
schema Text
tableName Int64
jobId =
  let 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|
          SELECT ${jobColumns}
          FROM ${tbl}
          WHERE id = #{jobId :: CInt8}
        |]

-- | Fetch a job by its dedup key. The partial unique index guarantees at most one row.
getJobByDedupKeySQL :: Text -> Text -> Text -> Query (JobRead Value)
getJobByDedupKeySQL :: Text -> Text -> Text -> Query (JobRead Value)
getJobByDedupKeySQL Text
schema Text
tableName Text
key =
  let 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|
          SELECT ${jobColumns}
          FROM ${tbl}
          WHERE dedup_key = #{key :: CText}
        |]

-- | Delete a job by id. Refuses one with children, which 'cancelJobCascadeSQL' takes.
-- A deleted child with no siblings left resumes its parent for a completion round.
cancelJobSQL :: Text -> Text -> Int64 -> Query Int64
cancelJobSQL :: Text -> Text -> Int64 -> Query Int64
cancelJobSQL Text
schema Text
tableName Int64
jobId =
  let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
   in [sql|
        WITH cancel AS (
          DELETE FROM ${tbl}
          WHERE id = #{jobId :: CInt8}
            AND NOT EXISTS (SELECT 1 FROM ${tbl} child WHERE child.parent_id = #{jobId :: CInt8})
          RETURNING id, parent_id
        ),
        wake_parent AS (
          UPDATE ${tbl}
          SET suspended = FALSE, updated_at = NOW()
          WHERE id = (SELECT parent_id FROM cancel WHERE parent_id IS NOT NULL)
            AND suspended = TRUE
            AND NOT EXISTS (
              SELECT 1 FROM ${tbl} child
              WHERE child.parent_id = (SELECT parent_id FROM cancel WHERE parent_id IS NOT NULL)
                AND child.id NOT IN (SELECT id FROM cancel)
            )
          RETURNING id
        )
        SELECT (SELECT count(*) FROM cancel) AS @{result :: CInt8}
      |]

-- | @UNION ALL@ of @body@ over each job table, passing its raw name and schema-qualified reference.
unionAllOverQueueTables :: SchemaName -> [TableName] -> (TableName -> Text -> Text) -> Text
unionAllOverQueueTables :: Text -> [Text] -> (Text -> Text -> Text) -> Text
unionAllOverQueueTables Text
schema [Text]
tableNames Text -> Text -> Text
body =
  Text -> [Text] -> Text
T.intercalate Text
" UNION ALL " ((Text -> Text) -> [Text] -> [Text]
forall a b. (a -> b) -> [a] -> [b]
map (\Text
tableName -> Text -> Text -> Text
body Text
tableName (Text -> Text -> Text
jobQueueTable Text
schema Text
tableName)) [Text]
tableNames)