{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Arbiter.Core.Sql.Groups
( groupsWindowSQL
, emptiedWindowSQL
, lockGroupsSQL
, refreshGroupsSQL
, insertMissingGroupsSQL
) where
import Data.Int (Int64)
import Data.Text (Text)
import NeatInterpolation (text)
import Arbiter.Core.Job.Schema (jobQueueGroupsTable, jobQueueTable)
import Arbiter.Core.Job.Schema.Groups (groupAggregates, inFlightPredicate)
import Arbiter.Core.Sql.QQ (sql)
import Arbiter.Core.Sql.Query (Query)
groupsWindowSQL :: Text -> Text -> Int -> Maybe Text -> Query Text
groupsWindowSQL :: Text -> Text -> Int -> Maybe Text -> Query Text
groupsWindowSQL Text
schema Text
tableName Int
limit Maybe Text
cursor =
let groupsTbl :: Text
groupsTbl = Text -> Text -> Text
jobQueueGroupsTable Text
schema Text
tableName
lim :: Int64
lim = Int -> Int64
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int
limit :: Int64
after :: Query ()
after = (Text -> Query ()) -> Maybe Text -> Query ()
forall m a. Monoid m => (a -> m) -> Maybe a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap (\Text
key -> [sql|WHERE group_key > #{key :: CText}|]) Maybe Text
cursor
in [sql|SELECT @{group_key :: CText} FROM ${groupsTbl} ${after} ORDER BY group_key LIMIT #{lim :: CInt8}|]
emptiedWindowSQL :: Text -> Text -> Int -> Maybe Text -> Query Text
emptiedWindowSQL :: Text -> Text -> Int -> Maybe Text -> Query Text
emptiedWindowSQL Text
schema Text
tableName Int
limit Maybe Text
cursor =
let groupsTbl :: Text
groupsTbl = Text -> Text -> Text
jobQueueGroupsTable Text
schema Text
tableName
lim :: Int64
lim = Int -> Int64
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int
limit :: Int64
after :: Query ()
after = (Text -> Query ()) -> Maybe Text -> Query ()
forall m a. Monoid m => (a -> m) -> Maybe a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap (\Text
key -> [sql|AND group_key > #{key :: CText}|]) Maybe Text
cursor
in [sql|
SELECT @{group_key :: CText} FROM ${groupsTbl}
WHERE job_count = 0 ${after}
ORDER BY group_key LIMIT #{lim :: CInt8}
|]
lockGroupsSQL :: Text -> Text -> Maybe Text -> Maybe Text -> [Text] -> Query Text
lockGroupsSQL :: Text -> Text -> Maybe Text -> Maybe Text -> [Text] -> Query Text
lockGroupsSQL Text
schema Text
tableName Maybe Text
cursor Maybe Text
upper [Text]
emptied =
let groupsTbl :: Text
groupsTbl = Text -> Text -> Text
jobQueueGroupsTable Text
schema Text
tableName
after :: Query ()
after = (Text -> Query ()) -> Maybe Text -> Query ()
forall m a. Monoid m => (a -> m) -> Maybe a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap (\Text
key -> [sql|AND group_key > #{key :: CText}|]) Maybe Text
cursor
window :: Query ()
window =
(Text -> Query ()) -> Maybe Text -> Query ()
forall m a. Monoid m => (a -> m) -> Maybe a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap
( \Text
upperKey ->
[sql|SELECT group_key FROM ${groupsTbl} WHERE group_key <= #{upperKey :: CText} ${after} UNION |]
)
Maybe Text
upper
in [sql|
WITH targets AS (
${window}SELECT unnest(#{emptied :: [CText]}::text[]) AS group_key
)
SELECT @{group_key :: CText} FROM ${groupsTbl}
WHERE group_key IN (SELECT group_key FROM targets)
ORDER BY group_key
FOR UPDATE SKIP LOCKED
|]
summaryAggregates :: Text
summaryAggregates :: Text
summaryAggregates =
let aggs :: Text
aggs = Text -> Text
groupAggregates Text
""
inFlight :: Text
inFlight = Text -> Text
inFlightPredicate Text
""
in [text|${aggs}, MAX(not_visible_until) FILTER (WHERE ${inFlight}) AS in_flight_until|]
refreshGroupsSQL :: Text -> Text -> [Text] -> Query Int64
refreshGroupsSQL :: Text -> Text -> [Text] -> Query Int64
refreshGroupsSQL Text
schema Text
tableName [Text]
keys =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
groupsTbl :: Text
groupsTbl = Text -> Text -> Text
jobQueueGroupsTable Text
schema Text
tableName
aggs :: Text
aggs = Text
summaryAggregates
in [sql|
WITH params AS (
SELECT unnest(#{keys :: [CText]}::text[]) AS group_key
),
current AS (
SELECT group_key, ${aggs}
FROM ${tbl}
WHERE group_key IN (SELECT group_key FROM params)
GROUP BY group_key
),
deleted AS (
DELETE FROM ${groupsTbl} summary
WHERE summary.group_key IN (SELECT group_key FROM params)
AND NOT EXISTS (SELECT 1 FROM current fresh WHERE fresh.group_key = summary.group_key)
RETURNING 1
),
updated AS (
UPDATE ${groupsTbl} summary
SET min_priority = fresh.min_priority,
min_id = fresh.min_id,
job_count = fresh.job_count,
ready_count = fresh.ready_count,
next_due = fresh.next_due,
in_flight_until = fresh.in_flight_until
FROM current fresh
WHERE summary.group_key = fresh.group_key
AND (summary.min_priority <> fresh.min_priority OR summary.min_id <> fresh.min_id
OR summary.job_count <> fresh.job_count
OR summary.ready_count <> fresh.ready_count
OR summary.next_due IS DISTINCT FROM fresh.next_due
OR summary.in_flight_until IS DISTINCT FROM fresh.in_flight_until)
RETURNING 1
)
SELECT (SELECT count(*) FROM deleted) + (SELECT count(*) FROM updated) AS @{rewritten :: CInt8}
|]
insertMissingGroupsSQL :: Text -> Text -> Int -> Maybe Text -> Maybe Text -> Query (Text, Bool)
insertMissingGroupsSQL :: Text
-> Text -> Int -> Maybe Text -> Maybe Text -> Query (Text, Bool)
insertMissingGroupsSQL Text
schema Text
tableName Int
limit Maybe Text
lower Maybe Text
upper =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
groupsTbl :: Text
groupsTbl = Text -> Text -> Text
jobQueueGroupsTable Text
schema Text
tableName
aggs :: Text
aggs = Text
summaryAggregates
lim :: Int64
lim = Int -> Int64
forall a b. (Integral a, Num b) => a -> b
fromIntegral Int
limit :: Int64
after :: Query ()
after = (Text -> Query ()) -> Maybe Text -> Query ()
forall m a. Monoid m => (a -> m) -> Maybe a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap (\Text
key -> [sql|AND job.group_key > #{key :: CText}|]) Maybe Text
lower
upTo :: Query ()
upTo = (Text -> Query ()) -> Maybe Text -> Query ()
forall m a. Monoid m => (a -> m) -> Maybe a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap (\Text
key -> [sql|AND job.group_key <= #{key :: CText}|]) Maybe Text
upper
in [sql|
WITH missing_keys AS (
SELECT DISTINCT job.group_key
FROM ${tbl} job
WHERE job.group_key IS NOT NULL
${after}
${upTo}
AND NOT EXISTS (SELECT 1 FROM ${groupsTbl} summary WHERE summary.group_key = job.group_key)
ORDER BY job.group_key
LIMIT #{lim :: CInt8}
),
missing AS (
SELECT group_key, ${aggs}
FROM ${tbl}
WHERE group_key IN (SELECT group_key FROM missing_keys)
GROUP BY group_key
),
inserted AS (
INSERT INTO ${groupsTbl} (group_key, min_priority, min_id, job_count, ready_count, next_due, in_flight_until)
SELECT group_key, min_priority, min_id, job_count, ready_count, next_due, in_flight_until
FROM missing
ORDER BY group_key
ON CONFLICT (group_key) DO NOTHING
RETURNING group_key
)
SELECT missing_key.group_key AS @{group_key :: CText}, (landed.group_key IS NOT NULL) AS @{inserted :: CBool}
FROM missing_keys missing_key
LEFT JOIN inserted landed ON landed.group_key = missing_key.group_key
ORDER BY missing_key.group_key
|]