{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
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)
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)
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)
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"
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)
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"
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)
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"
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)
sortDirSql :: SortDir -> Text
sortDirSql :: SortDir -> Text
sortDirSql SortDir
SortAsc = Text
"ASC"
sortDirSql SortDir
SortDesc = Text
"DESC"
throttledPredicateSQL :: Text
throttledPredicateSQL :: Text
throttledPredicateSQL =
Text
"throttled_until > NOW() AND not_visible_until > NOW()"
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
|]
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|]
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}
|]
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}
|]
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}|]
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}|]
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}
|]
data NullsBehavior
=
NullsNotApplicable
|
NullsAsAbsent
|
NullsAsMinimum
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"
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
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
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}|]
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
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
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
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}|]
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
|]
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
|]
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
|]
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
|]
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
|]
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
|]
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
|]
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)
|]
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}
|]
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}
|]
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}|])
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
""
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}
|]
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}
|]
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}
|]
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}
|]
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)