{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Arbiter.Core.Sql.Stats
( getQueueStatsSQL
, allQueueStatsSQL
, countChildrenBatchSQL
) where
import Data.Int (Int64)
import Data.Maybe (fromMaybe)
import Data.Text (Text)
import Data.Text qualified as T
import NeatInterpolation (text)
import Arbiter.Core.Codec (RowCodec)
import Arbiter.Core.Job.Schema (SchemaName, TableName, jobQueueDLQTable, jobQueueTable)
import Arbiter.Core.Queues (arbiterQueuesTable)
import Arbiter.Core.Sql.Jobs (jobStatusCaseSQL, unionAllOverQueueTables)
import Arbiter.Core.Sql.QQ (sql)
import Arbiter.Core.Sql.Query (Query, rawRows, rows)
import Arbiter.Core.SqlLiterals (textLiteral)
import Arbiter.Core.Worker (arbiterWorkersTable)
getQueueStatsSQL :: RowCodec a -> SchemaName -> TableName -> [Text] -> Query a
getQueueStatsSQL :: forall a. RowCodec a -> Text -> Text -> [Text] -> Query a
getQueueStatsSQL RowCodec a
codec Text
schema Text
tableName [Text]
kinds =
let stats :: Text
stats = Text -> Text -> [Text] -> Text
queueStatsSelect Text
schema Text
tableName [Text]
kinds
in RowCodec a -> Query () -> Query a
forall a. RowCodec a -> Query () -> Query a
rows RowCodec a
codec [sql|${stats}|]
queueStatsSelect :: SchemaName -> TableName -> [Text] -> Text
queueStatsSelect :: Text -> Text -> [Text] -> Text
queueStatsSelect Text
schema Text
tableName [Text]
kinds =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
dlqTbl :: Text
dlqTbl = Text -> Text -> Text
jobQueueDLQTable Text
schema Text
tableName
kindExpr :: Text
kindExpr = [Text] -> Text
declaredKindSQL [Text]
kinds
in [text|
SELECT MAX(total_jobs) FILTER (WHERE total_row = 1) AS total_jobs,
MAX(ready_jobs) FILTER (WHERE total_row = 1) AS ready_jobs,
MAX(in_flight_jobs) FILTER (WHERE total_row = 1) AS in_flight_jobs,
MAX(scheduled_jobs) FILTER (WHERE total_row = 1) AS scheduled_jobs,
MAX(backoff_jobs) FILTER (WHERE total_row = 1) AS backoff_jobs,
MAX(throttled_jobs) FILTER (WHERE total_row = 1) AS throttled_jobs,
MAX(suspended_jobs) FILTER (WHERE total_row = 1) AS suspended_jobs,
MAX(cancelled_jobs) FILTER (WHERE total_row = 1) AS cancelled_jobs,
MAX(oldest_ready_age_seconds) FILTER (WHERE total_row = 1) AS oldest_ready_age_seconds,
MAX(oldest_in_flight_age_seconds) FILTER (WHERE total_row = 1) AS oldest_in_flight_age_seconds,
(SELECT COUNT(*)::int8 FROM ${dlqTbl}) AS dlq_jobs,
jsonb_object_agg(kind, total_jobs) FILTER (WHERE total_row = 0 AND kind IS NOT NULL) AS kind_counts
FROM (
SELECT kind, GROUPING(kind) AS total_row,
COUNT(*)::int8 AS total_jobs,
COUNT(*) FILTER (WHERE status = 'ready') AS ready_jobs,
COUNT(*) FILTER (WHERE status = 'in_flight') AS in_flight_jobs,
COUNT(*) FILTER (WHERE status = 'scheduled') AS scheduled_jobs,
COUNT(*) FILTER (WHERE status = 'backoff') AS backoff_jobs,
COUNT(*) FILTER (WHERE status = 'throttled') AS throttled_jobs,
COUNT(*) FILTER (WHERE status = 'suspended') AS suspended_jobs,
COUNT(*) FILTER (WHERE status = 'cancelled') AS cancelled_jobs,
EXTRACT(EPOCH FROM (
clock_timestamp() - MIN(inserted_at) FILTER (WHERE status = 'ready')
))::float8 AS oldest_ready_age_seconds,
EXTRACT(EPOCH FROM (
clock_timestamp() - MIN(last_attempted_at) FILTER (WHERE status = 'in_flight')
))::float8 AS oldest_in_flight_age_seconds
FROM (
SELECT inserted_at, last_attempted_at, ${kindExpr} AS kind, ${jobStatusCaseSQL} AS status
FROM ${tbl}
) classified
GROUP BY GROUPING SETS ((), (kind))
) rollup
|]
declaredKindSQL :: [Text] -> Text
declaredKindSQL :: [Text] -> Text
declaredKindSQL [] = Text
"NULL::text"
declaredKindSQL [Text]
kinds =
let literals :: Text
literals = Text -> [Text] -> Text
T.intercalate Text
", " ((Text -> Text) -> [Text] -> [Text]
forall a b. (a -> b) -> [a] -> [b]
map Text -> Text
textLiteral [Text]
kinds)
in [text|CASE WHEN kind IN (${literals}) THEN kind END|]
allQueueStatsSQL :: RowCodec a -> SchemaName -> [(TableName, [Text])] -> Query a
allQueueStatsSQL :: forall a. RowCodec a -> Text -> [(Text, [Text])] -> Query a
allQueueStatsSQL RowCodec a
codec Text
schema [(Text, [Text])]
queueKinds =
let queuesTbl :: Text
queuesTbl = Text -> Text
arbiterQueuesTable Text
schema
workersTbl :: Text
workersTbl = Text -> Text
arbiterWorkersTable Text
schema
in RowCodec a -> Text -> Query a
forall a. RowCodec a -> Text -> Query a
rawRows RowCodec a
codec (Text -> Query a) -> Text -> Query a
forall a b. (a -> b) -> a -> b
$ Text -> [Text] -> (Text -> Text -> Text) -> Text
unionAllOverQueueTables Text
schema (((Text, [Text]) -> Text) -> [(Text, [Text])] -> [Text]
forall a b. (a -> b) -> [a] -> [b]
map (Text, [Text]) -> Text
forall a b. (a, b) -> a
fst [(Text, [Text])]
queueKinds) ((Text -> Text -> Text) -> Text) -> (Text -> Text -> Text) -> Text
forall a b. (a -> b) -> a -> b
$ \Text
tableName Text
_ ->
let stats :: Text
stats = Text -> Text -> [Text] -> Text
queueStatsSelect Text
schema Text
tableName ([Text] -> Maybe [Text] -> [Text]
forall a. a -> Maybe a -> a
fromMaybe [] (Text -> [(Text, [Text])] -> Maybe [Text]
forall a b. Eq a => a -> [(a, b)] -> Maybe b
lookup Text
tableName [(Text, [Text])]
queueKinds))
in [text|
SELECT '${tableName}' AS queue, stats.*,
COALESCE((SELECT paused FROM ${queuesTbl} WHERE queue_name = '${tableName}'), FALSE) AS queue_paused,
worker_counts.workers_live, worker_counts.workers_paused
FROM (${stats}) stats
CROSS JOIN (
-- Live matches the worker health CASE: a fresh heartbeat and not draining.
SELECT COUNT(*)::int8 AS workers_live, COUNT(*) FILTER (WHERE paused)::int8 AS workers_paused
FROM ${workersTbl}
WHERE queue_name = '${tableName}'
AND last_heartbeat >= NOW() - stale_threshold_secs * interval '1 second'
AND NOT shutting_down
) worker_counts
|]
countChildrenBatchSQL :: Text -> Text -> [Int64] -> Query ()
countChildrenBatchSQL :: Text -> Text -> [Int64] -> Query ()
countChildrenBatchSQL Text
schema Text
tableName [Int64]
jobIds =
let tbl :: Text
tbl = Text -> Text -> Text
jobQueueTable Text
schema Text
tableName
in [sql|
SELECT parent_id, COUNT(*),
COUNT(*) FILTER (WHERE suspended)
FROM ${tbl}
WHERE parent_id = ANY(#{jobIds :: [CInt8]})
GROUP BY parent_id
|]