{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE QuasiQuotes #-}
module Arbiter.Core.Sql.Queues
( queueColumnList
, ensureQueueSQL
, setQueuePausedSQL
, getQueueSQL
, listQueuesSQL
) where
import Data.Int (Int64)
import Data.Text (Text)
import Arbiter.Core.Codec (queueRowCodec)
import Arbiter.Core.Job.Schema (SchemaName, pauseNotifyChannelPrefix)
import Arbiter.Core.Queues (QueueRow, arbiterQueuesTable)
import Arbiter.Core.Sql.QQ (sql)
import Arbiter.Core.Sql.Query (Query, rows)
import Arbiter.Core.SqlLiterals (textLiteral)
import Arbiter.Core.Worker (arbiterWorkersTable)
queueColumnList :: Text
queueColumnList :: Text
queueColumnList = Text
"queue_name, paused, paused_at, metadata, created_at, updated_at"
ensureQueueSQL :: SchemaName -> Text -> Query ()
ensureQueueSQL :: Text -> Text -> Query ()
ensureQueueSQL Text
schemaName Text
queue =
let tbl :: Text
tbl = Text -> Text
arbiterQueuesTable Text
schemaName
in [sql|
INSERT INTO ${tbl} (queue_name) VALUES (#{queue :: CText})
ON CONFLICT (queue_name) DO NOTHING
|]
setQueuePausedSQL :: SchemaName -> Text -> Bool -> Query Int64
setQueuePausedSQL :: Text -> Text -> Bool -> Query Int64
setQueuePausedSQL Text
schemaName Text
queue Bool
paused =
let tbl :: Text
tbl = Text -> Text
arbiterQueuesTable Text
schemaName
workersTbl :: Text
workersTbl = Text -> Text
arbiterWorkersTable Text
schemaName
chanPrefix :: Text
chanPrefix = Text -> Text
textLiteral (Text -> Text
pauseNotifyChannelPrefix Text
schemaName)
in [sql|
WITH upsert AS (
INSERT INTO ${tbl} (queue_name, paused, paused_at)
VALUES (#{queue :: CText}, #{paused :: CBool},
CASE WHEN #{paused :: CBool}::boolean THEN NOW() ELSE NULL END)
ON CONFLICT (queue_name) DO UPDATE
SET paused = EXCLUDED.paused,
paused_at = CASE
WHEN EXCLUDED.paused AND NOT arbiter_queues.paused THEN NOW()
WHEN NOT EXCLUDED.paused THEN NULL
ELSE arbiter_queues.paused_at
END,
updated_at = NOW()
RETURNING queue_name, paused
),
notif AS (
SELECT pg_notify(
LEFT(${chanPrefix} || worker.queue_name, 63),
json_build_object(
'worker_id', worker.worker_id,
'paused', worker.paused OR upserted.paused
)::text
)
FROM upsert upserted
JOIN ${workersTbl} worker ON worker.queue_name = upserted.queue_name
)
SELECT count(*)::int8 AS @{count :: CInt8} FROM upsert
WHERE (SELECT count(*) FROM notif) >= 0
|]
getQueueSQL :: SchemaName -> Text -> Query QueueRow
getQueueSQL :: Text -> Text -> Query QueueRow
getQueueSQL Text
schemaName Text
queue =
let tbl :: Text
tbl = Text -> Text
arbiterQueuesTable Text
schemaName
in RowCodec QueueRow -> Query () -> Query QueueRow
forall a. RowCodec a -> Query () -> Query a
rows RowCodec QueueRow
queueRowCodec [sql|SELECT ${queueColumnList} FROM ${tbl} WHERE queue_name = #{queue :: CText}|]
listQueuesSQL :: SchemaName -> Query QueueRow
listQueuesSQL :: Text -> Query QueueRow
listQueuesSQL Text
schemaName =
let tbl :: Text
tbl = Text -> Text
arbiterQueuesTable Text
schemaName
in RowCodec QueueRow -> Query () -> Query QueueRow
forall a. RowCodec a -> Query () -> Query a
rows RowCodec QueueRow
queueRowCodec [sql|SELECT ${queueColumnList} FROM ${tbl} ORDER BY queue_name|]