{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Arbiter.Core.Sql.Cron
( allCronColumns
, upsertCronDefaultSQL
, listCronSchedulesSQL
, getCronScheduleByNameSQL
, updateCronScheduleSQL
, touchCronLastFiredSQL
, touchCronCheckedSQL
, tryFireCronGateSQL
, tryAcquireCronLeaderSQL
, requestCronRunSQL
, claimCronRunSQL
, touchCronManualRunSQL
, pendingCronRunsSQL
) where
import Control.Monad (join)
import Data.Maybe (isJust)
import Data.Text (Text)
import Data.Time (UTCTime)
import NeatInterpolation (text)
import Arbiter.Core.Codec (cronScheduleRowCodec)
import Arbiter.Core.CronSchedule (CronScheduleRow, CronScheduleUpdate (..), cronSchedulesTable)
import Arbiter.Core.Job.Schema (cronRunNotifyChannel)
import Arbiter.Core.Sql.QQ (sql)
import Arbiter.Core.Sql.Query (Query, rows)
import Arbiter.Core.SqlLiterals (textLiteral)
allCronColumns :: Text
allCronColumns :: Text
allCronColumns =
[text|
name, queue_name, default_expression, default_overlap, default_timezone,
override_expression, override_overlap, override_timezone, enabled,
last_fired_at, last_checked_at, run_requested_at, last_manual_run_at, created_at, updated_at
|]
cronReadColumns :: Text
cronReadColumns :: Text
cronReadColumns =
[text|
name, queue_name, default_expression, default_overlap, default_timezone,
override_expression, override_overlap, override_timezone, enabled,
last_fired_at, last_checked_at,
(CASE WHEN ${cronRunPending} THEN run_requested_at END) AS run_requested_at,
last_manual_run_at, created_at, updated_at
|]
upsertCronDefaultSQL :: Text -> Text -> Text -> Text -> Text -> Maybe Text -> Query ()
upsertCronDefaultSQL :: Text -> Text -> Text -> Text -> Text -> Maybe Text -> Query ()
upsertCronDefaultSQL Text
schemaName Text
name Text
queueName Text
defaultExpr Text
defaultOv Maybe Text
defaultTz =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
in [sql|
INSERT INTO ${tbl} (name, queue_name, default_expression, default_overlap, default_timezone)
VALUES (#{name :: CText}, #{queueName :: CText}, #{defaultExpr :: CText},
#{defaultOv :: CText}, #{defaultTz :: Maybe CText})
ON CONFLICT (name) DO UPDATE SET
queue_name = EXCLUDED.queue_name,
default_expression = EXCLUDED.default_expression,
default_overlap = EXCLUDED.default_overlap,
default_timezone = EXCLUDED.default_timezone,
updated_at = NOW()
WHERE (${tbl}.queue_name, ${tbl}.default_expression, ${tbl}.default_overlap, ${tbl}.default_timezone)
IS DISTINCT FROM (EXCLUDED.queue_name, EXCLUDED.default_expression,
EXCLUDED.default_overlap, EXCLUDED.default_timezone)
|]
listCronSchedulesSQL :: Text -> Maybe Text -> Query CronScheduleRow
listCronSchedulesSQL :: Text -> Maybe Text -> Query CronScheduleRow
listCronSchedulesSQL Text
schemaName Maybe Text
queue =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
in RowCodec CronScheduleRow -> Query () -> Query CronScheduleRow
forall a. RowCodec a -> Query () -> Query a
rows
RowCodec CronScheduleRow
cronScheduleRowCodec
[sql|
SELECT ${cronReadColumns} FROM ${tbl}
WHERE #{queue :: Maybe CText}::text IS NULL OR queue_name = #{queue :: Maybe CText}::text
ORDER BY name
|]
getCronScheduleByNameSQL :: Text -> Text -> Query CronScheduleRow
getCronScheduleByNameSQL :: Text -> Text -> Query CronScheduleRow
getCronScheduleByNameSQL Text
schemaName Text
name =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
in RowCodec CronScheduleRow -> Query () -> Query CronScheduleRow
forall a. RowCodec a -> Query () -> Query a
rows RowCodec CronScheduleRow
cronScheduleRowCodec [sql|SELECT ${cronReadColumns} FROM ${tbl} WHERE name = #{name :: CText}|]
updateCronScheduleSQL :: Text -> Text -> CronScheduleUpdate -> Maybe (Query ())
updateCronScheduleSQL :: Text -> Text -> CronScheduleUpdate -> Maybe (Query ())
updateCronScheduleSQL Text
_ Text
_ (CronScheduleUpdate Maybe (Maybe Text)
Nothing Maybe (Maybe Text)
Nothing Maybe (Maybe Text)
Nothing Maybe Bool
Nothing) = Maybe (Query ())
forall a. Maybe a
Nothing
updateCronScheduleSQL Text
schemaName Text
name (CronScheduleUpdate Maybe (Maybe Text)
mExpression Maybe (Maybe Text)
mOverlap Maybe (Maybe Text)
mTimezone Maybe Bool
mEnabled) =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
setExpression :: Bool
setExpression = Maybe (Maybe Text) -> Bool
forall a. Maybe a -> Bool
isJust Maybe (Maybe Text)
mExpression
expression :: Maybe Text
expression = Maybe (Maybe Text) -> Maybe Text
forall (m :: * -> *) a. Monad m => m (m a) -> m a
join Maybe (Maybe Text)
mExpression
setOverlap :: Bool
setOverlap = Maybe (Maybe Text) -> Bool
forall a. Maybe a -> Bool
isJust Maybe (Maybe Text)
mOverlap
overlap :: Maybe Text
overlap = Maybe (Maybe Text) -> Maybe Text
forall (m :: * -> *) a. Monad m => m (m a) -> m a
join Maybe (Maybe Text)
mOverlap
setTimezone :: Bool
setTimezone = Maybe (Maybe Text) -> Bool
forall a. Maybe a -> Bool
isJust Maybe (Maybe Text)
mTimezone
timezone :: Maybe Text
timezone = Maybe (Maybe Text) -> Maybe Text
forall (m :: * -> *) a. Monad m => m (m a) -> m a
join Maybe (Maybe Text)
mTimezone
in Query () -> Maybe (Query ())
forall a. a -> Maybe a
Just
[sql|
UPDATE ${tbl}
SET override_expression = CASE WHEN #{setExpression :: CBool}::boolean
THEN #{expression :: Maybe CText}::text
ELSE override_expression END,
override_overlap = CASE WHEN #{setOverlap :: CBool}::boolean
THEN #{overlap :: Maybe CText}::text
ELSE override_overlap END,
override_timezone = CASE WHEN #{setTimezone :: CBool}::boolean
THEN #{timezone :: Maybe CText}::text
ELSE override_timezone END,
enabled = COALESCE(#{mEnabled :: Maybe CBool}::boolean, enabled),
run_requested_at = CASE WHEN #{mEnabled :: Maybe CBool}::boolean IS FALSE
THEN NULL
ELSE run_requested_at END,
updated_at = NOW()
WHERE name = #{name :: CText}
|]
touchCronLastFiredSQL :: Text -> Text -> Query ()
touchCronLastFiredSQL :: Text -> Text -> Query ()
touchCronLastFiredSQL Text
schemaName Text
name =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
in [sql|UPDATE ${tbl} SET last_fired_at = NOW() WHERE name = #{name :: CText}|]
touchCronCheckedSQL :: Text -> UTCTime -> [Text] -> Query ()
touchCronCheckedSQL :: Text -> UTCTime -> [Text] -> Query ()
touchCronCheckedSQL Text
schemaName UTCTime
watermark [Text]
names =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
in [sql|
UPDATE ${tbl}
SET last_checked_at = GREATEST(last_checked_at, #{watermark :: CTimestamptz})
WHERE name = ANY(#{names :: [CText]})
|]
tryFireCronGateSQL :: Text -> UTCTime -> Text -> Query ()
tryFireCronGateSQL :: Text -> UTCTime -> Text -> Query ()
tryFireCronGateSQL Text
schemaName UTCTime
minuteFloor Text
name =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
in [sql|
UPDATE ${tbl}
SET last_fired_at = GREATEST(last_fired_at, #{minuteFloor :: CTimestamptz})
WHERE name = #{name :: CText}
AND (last_fired_at IS NULL OR last_fired_at < #{minuteFloor :: CTimestamptz})
|]
tryAcquireCronLeaderSQL :: Text -> Text -> Text -> Query Bool
tryAcquireCronLeaderSQL :: Text -> Text -> Text -> Query Bool
tryAcquireCronLeaderSQL Text
schema Text
queue Text
name =
[sql|
SELECT pg_try_advisory_xact_lock(
hashtextextended(#{schema :: CText} || ':' || #{queue :: CText} || ':' || #{name :: CText}, 0)
) AS @{result :: CBool}
|]
cronRunRequestTtl :: Text
cronRunRequestTtl :: Text
cronRunRequestTtl = Text
"INTERVAL '5 minutes'"
cronRunPending :: Text
cronRunPending :: Text
cronRunPending = [text|(run_requested_at IS NOT NULL AND run_requested_at > NOW() - ${cronRunRequestTtl})|]
requestCronRunSQL :: Text -> Text -> Query Text
requestCronRunSQL :: Text -> Text -> Query Text
requestCronRunSQL Text
schemaName Text
name =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
chan :: Text
chan = Text -> Text
textLiteral (Text -> Text
cronRunNotifyChannel Text
schemaName)
in [sql|
WITH found AS (SELECT enabled, run_requested_at FROM ${tbl} WHERE name = #{name :: CText} FOR UPDATE),
stamped AS (
UPDATE ${tbl} SET run_requested_at = NOW(), updated_at = NOW()
WHERE name = #{name :: CText} AND enabled AND NOT ${cronRunPending}
RETURNING pg_notify(${chan}, name))
SELECT CASE
WHEN EXISTS (SELECT 1 FROM stamped) THEN 'stamped'
WHEN EXISTS (SELECT 1 FROM found WHERE enabled AND ${cronRunPending}) THEN 'pending'
WHEN EXISTS (SELECT 1 FROM found) THEN 'disabled'
ELSE 'not_found' END AS @{status :: CText}
|]
claimCronRunSQL :: Text -> Text -> Query CronScheduleRow
claimCronRunSQL :: Text -> Text -> Query CronScheduleRow
claimCronRunSQL Text
schemaName Text
name =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
in RowCodec CronScheduleRow -> Query () -> Query CronScheduleRow
forall a. RowCodec a -> Query () -> Query a
rows
RowCodec CronScheduleRow
cronScheduleRowCodec
[sql|
UPDATE ${tbl} SET run_requested_at = NULL, updated_at = NOW()
WHERE name = #{name :: CText} AND enabled AND ${cronRunPending}
RETURNING ${allCronColumns}
|]
touchCronManualRunSQL :: Text -> UTCTime -> Text -> Query ()
touchCronManualRunSQL :: Text -> UTCTime -> Text -> Query ()
touchCronManualRunSQL Text
schemaName UTCTime
firedAt Text
name =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
in [sql|
UPDATE ${tbl}
SET last_manual_run_at = GREATEST(last_manual_run_at, #{firedAt :: CTimestamptz}), updated_at = NOW()
WHERE name = #{name :: CText}
|]
pendingCronRunsSQL :: Text -> [Text] -> Query Text
pendingCronRunsSQL :: Text -> [Text] -> Query Text
pendingCronRunsSQL Text
schemaName [Text]
names =
let tbl :: Text
tbl = Text -> Text
cronSchedulesTable Text
schemaName
in [sql|SELECT @{name :: CText} FROM ${tbl} WHERE name = ANY(#{names :: [CText]}) AND enabled AND ${cronRunPending}|]