{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Arbiter.Core.Sql.Gates
( ensureGateRowSQL
, checkGateSQL
, tryClaimGateSQL
, claimOrReadGateSQL
, releaseGateSQL
, gateNameDigestSQL
, bumpGateSQL
, bumpGateStateSQL
, setGateMetadataSQL
) where
import Data.Aeson (Value)
import Data.Text (Text)
import Data.Time (UTCTime)
import Arbiter.Core.Gates (arbiterGatesTable)
import Arbiter.Core.Job.Schema (SchemaName)
import Arbiter.Core.Sql.QQ (sql)
import Arbiter.Core.Sql.Query (Query)
ensureGateRowSQL :: SchemaName -> Text -> Query ()
ensureGateRowSQL :: Text -> Text -> Query ()
ensureGateRowSQL Text
schemaName Text
task =
let tbl :: Text
tbl = Text -> Text
arbiterGatesTable Text
schemaName
in [sql|
INSERT INTO ${tbl} (task_name) VALUES (#{task :: CText})
ON CONFLICT (task_name) DO NOTHING
|]
checkGateSQL :: SchemaName -> Double -> Text -> Query Bool
checkGateSQL :: Text -> Double -> Text -> Query Bool
checkGateSQL Text
schemaName Double
intervalSecs Text
task =
let tbl :: Text
tbl = Text -> Text
arbiterGatesTable Text
schemaName
in [sql|
SELECT (last_run_at < NOW() - (#{intervalSecs :: CFloat8}::double precision * interval '1 second'))
AS @{result :: CBool}
FROM ${tbl}
WHERE task_name = #{task :: CText}
|]
tryClaimGateSQL :: SchemaName -> Text -> Double -> Query (Maybe Value)
tryClaimGateSQL :: Text -> Text -> Double -> Query (Maybe Value)
tryClaimGateSQL Text
schemaName Text
task Double
intervalSecs =
let tbl :: Text
tbl = Text -> Text
arbiterGatesTable Text
schemaName
in [sql|
SELECT metadata -> 'payload' AS @{payload :: Maybe CJsonb} FROM ${tbl}
WHERE task_name = #{task :: CText}
AND last_run_at < NOW() - (#{intervalSecs :: CFloat8}::double precision * interval '1 second')
FOR UPDATE SKIP LOCKED
|]
gateNameDigestSQL :: Text -> Query Text
gateNameDigestSQL :: Text -> Query Text
gateNameDigestSQL Text
name = [sql|SELECT md5(#{name :: CText}) AS @{digest :: CText}|]
claimOrReadGateSQL
:: SchemaName -> Text -> Double -> Double -> Query (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
claimOrReadGateSQL :: Text
-> Text
-> Double
-> Double
-> Query (Maybe UTCTime, Maybe UTCTime, Maybe Value, Maybe Double)
claimOrReadGateSQL Text
schemaName Text
task Double
intervalSecs Double
maxAgeSecs =
let tbl :: Text
tbl = Text -> Text
arbiterGatesTable Text
schemaName
in [sql|
WITH seeded AS (
INSERT INTO ${tbl} (task_name, last_run_at)
SELECT #{task :: CText}, NOW()
WHERE NOT EXISTS (SELECT 1 FROM ${tbl} WHERE task_name = #{task :: CText})
ON CONFLICT (task_name) DO NOTHING
RETURNING NOW() AS claimed_at, '1970-01-01'::timestamptz AS previous_run_at
),
claimed AS (
UPDATE ${tbl} gate SET last_run_at = NOW()
FROM ${tbl} old
WHERE gate.task_name = #{task :: CText}
AND old.task_name = gate.task_name
AND gate.last_run_at < NOW() - (#{intervalSecs :: CFloat8}::double precision * interval '1 second')
RETURNING NOW() AS claimed_at, old.last_run_at AS previous_run_at
),
claim AS (
SELECT claimed_at, previous_run_at FROM seeded
UNION ALL
SELECT claimed_at, previous_run_at FROM claimed
),
published AS (
SELECT metadata -> 'payload' AS payload,
EXTRACT(EPOCH FROM NOW() - (metadata ->> 'at')::timestamptz)::float8 AS age_seconds
FROM ${tbl}
WHERE task_name = #{task :: CText}
AND NOT EXISTS (SELECT 1 FROM claim)
AND metadata -> 'payload' IS NOT NULL
AND (metadata ->> 'at')::timestamptz
> NOW() - (#{maxAgeSecs :: CFloat8}::double precision * interval '1 second')
)
SELECT (SELECT claimed_at FROM claim) AS @{claimed_at :: Maybe CTimestamptz},
(SELECT previous_run_at FROM claim) AS @{previous_run_at :: Maybe CTimestamptz},
(SELECT payload FROM published) AS @{payload :: Maybe CJsonb},
(SELECT age_seconds FROM published) AS @{age_seconds :: Maybe CFloat8}
|]
releaseGateSQL :: SchemaName -> Text -> UTCTime -> UTCTime -> Query ()
releaseGateSQL :: Text -> Text -> UTCTime -> UTCTime -> Query ()
releaseGateSQL Text
schemaName Text
task UTCTime
claimedAt UTCTime
previous =
let tbl :: Text
tbl = Text -> Text
arbiterGatesTable Text
schemaName
in [sql|
UPDATE ${tbl} SET last_run_at = #{previous :: CTimestamptz}
WHERE task_name = #{task :: CText} AND last_run_at = #{claimedAt :: CTimestamptz}
|]
bumpGateSQL :: SchemaName -> Text -> Query ()
bumpGateSQL :: Text -> Text -> Query ()
bumpGateSQL Text
schemaName Text
task =
let tbl :: Text
tbl = Text -> Text
arbiterGatesTable Text
schemaName
in [sql|UPDATE ${tbl} SET last_run_at = NOW() WHERE task_name = #{task :: CText}|]
bumpGateStateSQL :: SchemaName -> Text -> Value -> Query ()
bumpGateStateSQL :: Text -> Text -> Value -> Query ()
bumpGateStateSQL Text
schemaName Text
task Value
payload =
let tbl :: Text
tbl = Text -> Text
arbiterGatesTable Text
schemaName
in [sql|
UPDATE ${tbl}
SET last_run_at = NOW(),
metadata = jsonb_build_object('at', to_jsonb(NOW()), 'payload', #{payload :: CJsonb}::jsonb)
WHERE task_name = #{task :: CText}
|]
setGateMetadataSQL :: SchemaName -> Value -> UTCTime -> Text -> Query ()
setGateMetadataSQL :: Text -> Value -> UTCTime -> Text -> Query ()
setGateMetadataSQL Text
schemaName Value
metadata UTCTime
claimedAt Text
task =
let tbl :: Text
tbl = Text -> Text
arbiterGatesTable Text
schemaName
in [sql|
UPDATE ${tbl}
SET last_run_at = NOW(),
metadata = jsonb_build_object(
'at', to_jsonb(#{claimedAt :: CTimestamptz}::timestamptz),
'payload', #{metadata :: CJsonb}::jsonb
)
WHERE task_name = #{task :: CText}
AND (metadata ->> 'at' IS NULL OR (metadata ->> 'at')::timestamptz < #{claimedAt :: CTimestamptz})
|]