{-# LANGUAGE AllowAmbiguousTypes #-}
{-# LANGUAGE DataKinds #-}
{-# LANGUAGE DuplicateRecordFields #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE TypeFamilies #-}
{-# LANGUAGE UndecidableInstances #-}

-- | REST API server for the Arbiter job queue (SimpleDb backend).
--
-- __Security:__ No built-in authentication. All endpoints are publicly
-- accessible. Add auth middleware before exposing to untrusted networks.
module Arbiter.Servant.Server
  ( -- * Server handlers
    arbiterServer
  , arbiterServerHoisted
  , arbiterApp
  , runArbiterAPI
  , ArbiterServerConfig (..)
  , initArbiterServer
  , defaultQueueStatsCacheTtl
  , defaultMaintenanceInterval
  , defaultMaintenanceBucketIdle
  , defaultMaintenanceSparseInterval
  , defaultMaintenanceTimeout
  , BuildServer (..)
  ) where

import Arbiter.Core.CronSchedule qualified as CS
import Arbiter.Core.Health qualified as Health
import Arbiter.Core.HighLevel qualified as HL
import Arbiter.Core.Job.Schema qualified as Schema
import Arbiter.Core.Job.Types (DedupKey (..), JobPayload, JobStatus, isRollup, kindsFor)
import Arbiter.Core.Job.Types qualified as Job
import Arbiter.Core.JobResult (EncodeJobResult, encodeJobResult)
import Arbiter.Core.MonadArbiter (withDbTransaction)
import Arbiter.Core.Operations qualified as Ops
import Arbiter.Core.PoolConfig (PoolConfig (..))
import Arbiter.Core.QueueRegistry (JobPayloadRegistry, RegistryTables (..), SpecName, SpecPayload, SpecResult)
import Arbiter.Core.Queues qualified as Queues
import Arbiter.Core.Sql.Jobs (ArchiveSortColumn, DLQSortColumn, JobFilter (..), JobSortColumn, SortDir)
import Arbiter.Core.Trace (withPublishSpan)
import Arbiter.Simple (SimpleConnectionPool (..), SimpleDb, SimpleEnv (..), createSimpleEnvWithConfig, runSimpleDb)
import Arbiter.Worker (MaintenancePace (..), runMaintenancePass, storeEncodedResult)
import Arbiter.Worker.Config (maintenanceOpName)
import Arbiter.Worker.Cron (nextRunFromExpression, updateCronScheduleChecked)
import Arbiter.Worker.Logger (defaultLogConfig)
import Control.Concurrent (forkIOWithUnmask, threadDelay)
import Control.Concurrent.Async (race_)
import Control.Concurrent.MVar (MVar, modifyMVar, modifyMVar_, newMVar)
import Control.Concurrent.STM
  ( TChan
  , TVar
  , atomically
  , check
  , dupTChan
  , modifyTVar'
  , newBroadcastTChanIO
  , newTVarIO
  , readTChan
  , readTVar
  , readTVarIO
  , writeTChan
  )
import Control.Exception (SomeAsyncException, SomeException, bracket, bracket_, fromException, handle, throwIO, try)
import Control.Monad (forever, guard, join, mfilter, unless, void, when)
import Control.Monad.IO.Class (MonadIO, liftIO)
import Data.Aeson (encode)
import Data.ByteString (ByteString)
import Data.ByteString.Builder qualified as Builder
import Data.ByteString.Char8 qualified as BS8
import Data.ByteString.Lazy qualified as LBS
import Data.IORef (modifyIORef', newIORef, readIORef, writeIORef)
import Data.Int (Int64)
import Data.Map.Strict qualified as Map
import Data.Maybe (catMaybes, fromMaybe, isJust)
import Data.Pool qualified as Pool
import Data.Set qualified as Set
import Data.String (fromString)
import Data.Text (Text)
import Data.Text qualified as T
import Data.Text.Encoding (encodeUtf8)
import Data.Time (NominalDiffTime, UTCTime, diffUTCTime, getCurrentTime)
import Data.Time.Format (defaultTimeLocale, formatTime)
import Data.UUID.Types (UUID)
import Data.UUID.V4 qualified as UUID
import Database.PostgreSQL.Simple qualified as PG
import Database.PostgreSQL.Simple.Notification (Notification (..), getNotification)
import GHC.TypeLits (KnownSymbol, symbolVal)
import Network.HTTP.Types (status200)
import Network.Wai (responseStream)
import Network.Wai.Handler.Warp (Port, defaultSettings, runSettings, setPort)
import Servant
import Servant.Server.Generic (AsServerT)
import System.IO (stderr)
import System.Timeout (timeout)

import Arbiter.Servant.API
  ( ArbiterAPI
  , ArchiveAPI (..)
  , ConcurrencyAPI (..)
  , CronAPI (..)
  , DLQAPI (..)
  , HealthAPI (..)
  , JobsAPI (..)
  , MaintenanceAPI (..)
  , QueuesAPI (..)
  , RateLimitsAPI (..)
  , RegistryToAPI
  , SharedAPI
  , StatsAPI (..)
  , TableAPI (..)
  , WorkersAPI (..)
  )
import Arbiter.Servant.Types

-- | Configuration for the API server.
data ArbiterServerConfig (registry :: JobPayloadRegistry) = ArbiterServerConfig
  { forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> SimpleEnv registry
serverEnv :: SimpleEnv registry
  -- ^ The SimpleEnv containing schema and connection pool
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Bool
enableSSE :: Bool
  -- ^ Enable the Server-Sent Events streaming endpoint. When 'False', the
  -- @\/events\/stream@ endpoint returns one \"disabled\" event and closes.
  -- The admin UI then polls. Default: 'True'.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> MVar (Maybe SSEHub)
sseHub :: MVar (Maybe SSEHub)
  -- ^ Lazily started SSE broadcast hub, shared by all clients. The first
  -- subscriber starts it. The last disconnect tears it down and releases its
  -- @LISTEN@ connection.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell RateLimitPoliciesResponse
rateLimitPoliciesCache :: CacheCell RateLimitPoliciesResponse
  -- ^ Short-TTL cache for the rate-limit policy list.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> CacheCell ConcurrencyPoliciesResponse
concurrencyPoliciesCache :: CacheCell ConcurrencyPoliciesResponse
  -- ^ Short-TTL cache for the concurrency policy list.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell AllStatsResponse
allQueueStatsCache :: CacheCell AllStatsResponse
  -- ^ Short-TTL cache for the all-queues overview aggregate.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell StatsResponse
queueStatsCache :: CacheCell StatsResponse
  -- ^ Per-queue stats cache.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> NominalDiffTime
queueStatsCacheTtl :: NominalDiffTime
  -- ^ Per-queue stats staleness, or zero to always hit the database.
  -- Default: 'defaultQueueStatsCacheTtl'.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell HealthResponse
healthCache :: CacheCell HealthResponse
  -- ^ Short-TTL cache for the readiness probe.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> NominalDiffTime
maintenanceInterval :: NominalDiffTime
  -- ^ Minimum gap between runs of one maintenance operation. Zero runs every
  -- operation on every call. Default: 'defaultMaintenanceInterval'.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> NominalDiffTime
maintenanceSparseInterval :: NominalDiffTime
  -- ^ Gap between runs of one whole-schema operation, independent of
  -- 'maintenanceInterval'. Default: 'defaultMaintenanceSparseInterval'.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> NominalDiffTime
maintenanceBucketIdle :: NominalDiffTime
  -- ^ Idle age at which a pass prunes a rate-limit bucket.
  -- Default: 'defaultMaintenanceBucketIdle'.
  , forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> NominalDiffTime
maintenanceTimeout :: NominalDiffTime
  -- ^ Abort any single maintenance statement that runs longer than this.
  -- Default: 'defaultMaintenanceTimeout'.
  }

-- | A running SSE broadcast hub: the channel every client duplicates, and the
-- live subscriber count the listener watches to release itself at zero.
data SSEHub = SSEHub
  { SSEHub -> TChan ByteString
hubChan :: TChan ByteString
  , SSEHub -> TVar Int
hubRefs :: TVar Int
  }

-- | The schema every handler's statements run against.
serverSchema :: ArbiterServerConfig registry -> Text
serverSchema :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema = SimpleEnv registry -> Text
forall (registry :: JobPayloadRegistry). SimpleEnv registry -> Text
schema (SimpleEnv registry -> Text)
-> (ArbiterServerConfig registry -> SimpleEnv registry)
-> ArbiterServerConfig registry
-> Text
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ArbiterServerConfig registry -> SimpleEnv registry
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> SimpleEnv registry
serverEnv

-- | Run a statement on the server's own pool.
runDb :: (MonadIO n) => ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb :: forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config = IO a -> n a
forall a. IO a -> n a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO a -> n a)
-> (SimpleDb registry IO a -> IO a)
-> SimpleDb registry IO a
-> n a
forall b c a. (b -> c) -> (a -> b) -> a -> c
. SimpleEnv registry -> SimpleDb registry IO a -> IO a
forall (registry :: JobPayloadRegistry) (m :: * -> *) a.
SimpleEnv registry -> SimpleDb registry m a -> m a
runSimpleDb (ArbiterServerConfig registry -> SimpleEnv registry
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> SimpleEnv registry
serverEnv ArbiterServerConfig registry
config)

-- | 'NoContent' when a statement touched a row, 404 otherwise.
rowsOr404 :: LBS.ByteString -> Int64 -> Handler NoContent
rowsOr404 :: ByteString -> Int64 -> Handler NoContent
rowsOr404 ByteString
missing Int64
rowsAffected
  | Int64
rowsAffected Int64 -> Int64 -> Bool
forall a. Ord a => a -> a -> Bool
> Int64
0 = NoContent -> Handler NoContent
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure NoContent
NoContent
  | Bool
otherwise = ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err404 {errBody = missing}

-- | Answer a handler that decided its own error.
noContentOr :: Either ServerError () -> Handler NoContent
noContentOr :: Either ServerError () -> Handler NoContent
noContentOr = (ServerError -> Handler NoContent)
-> (() -> Handler NoContent)
-> Either ServerError ()
-> Handler NoContent
forall a c b. (a -> c) -> (b -> c) -> Either a b -> c
either ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError (Handler NoContent -> () -> Handler NoContent
forall a b. a -> b -> a
const (NoContent -> Handler NoContent
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure NoContent
NoContent))

-- | Run a job mutation. When it touches no row, re-read the job and answer 404, or
-- the 409 that @refuse@ derives from the job's state.
mutateJob
  :: forall payload registry
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> (Text -> SimpleDb registry IO Int64)
  -> (Job.JobRead payload -> LBS.ByteString)
  -> Handler NoContent
mutateJob :: forall payload (registry :: JobPayloadRegistry).
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> (Text -> SimpleDb registry IO Int64)
-> (JobRead payload -> ByteString)
-> Handler NoContent
mutateJob Text
tableName ArbiterServerConfig registry
config Int64
jobId Text -> SimpleDb registry IO Int64
mutate JobRead payload -> ByteString
refuse =
  Either ServerError () -> Handler NoContent
noContentOr (Either ServerError () -> Handler NoContent)
-> Handler (Either ServerError ()) -> Handler NoContent
forall (m :: * -> *) a b. Monad m => (a -> m b) -> m a -> m b
=<< ArbiterServerConfig registry
-> SimpleDb registry IO (Either ServerError ())
-> Handler (Either ServerError ())
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (Text -> SimpleDb registry IO Int64
mutate Text
schemaName SimpleDb registry IO Int64
-> (Int64 -> SimpleDb registry IO (Either ServerError ()))
-> SimpleDb registry IO (Either ServerError ())
forall a b.
SimpleDb registry IO a
-> (a -> SimpleDb registry IO b) -> SimpleDb registry IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= Int64 -> SimpleDb registry IO (Either ServerError ())
diagnose)
  where
    schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
    diagnose :: Int64 -> SimpleDb registry IO (Either ServerError ())
diagnose Int64
rowsAffected
      | Int64
rowsAffected Int64 -> Int64 -> Bool
forall a. Ord a => a -> a -> Bool
> Int64
0 = Either ServerError ()
-> SimpleDb registry IO (Either ServerError ())
forall a. a -> SimpleDb registry IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (() -> Either ServerError ()
forall a b. b -> Either a b
Right ())
      | Bool
otherwise =
          Either ServerError ()
-> (JobRead payload -> Either ServerError ())
-> Maybe (JobRead payload)
-> Either ServerError ()
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (ServerError -> Either ServerError ()
forall a b. a -> Either a b
Left ServerError
err404 {errBody = "Job not found"}) (\JobRead payload
job -> ServerError -> Either ServerError ()
forall a b. a -> Either a b
Left ServerError
err409 {errBody = refuse job})
            (Maybe (JobRead payload) -> Either ServerError ())
-> SimpleDb registry IO (Maybe (JobRead payload))
-> SimpleDb registry IO (Either ServerError ())
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text -> Text -> Int64 -> m (Maybe (JobRead payload))
Ops.getJobById @_ @payload Text
schemaName Text
tableName Int64
jobId

-- | Small pool configuration for admin API traffic.
serverPoolConfig :: PoolConfig
serverPoolConfig :: PoolConfig
serverPoolConfig =
  PoolConfig
    { poolSize :: Int
poolSize = Int
10
    , poolIdleTimeout :: Int
poolIdleTimeout = Int
60
    , poolStripes :: Maybe Int
poolStripes = Int -> Maybe Int
forall a. a -> Maybe a
Just Int
1
    }

-- | Create an 'ArbiterServerConfig' with a connection pool of its own. SSE needs the
-- event-streaming triggers, which 'Arbiter.Migrations.runMigrationsForRegistry' installs
-- when @enableEventStreaming@ is set.
initArbiterServer
  :: forall registry
   . Proxy registry
  -> ByteString
  -> Text
  -> IO (ArbiterServerConfig registry)
initArbiterServer :: forall (registry :: JobPayloadRegistry).
Proxy registry
-> ByteString -> Text -> IO (ArbiterServerConfig registry)
initArbiterServer Proxy registry
_proxy ByteString
connStr Text
schemaName = do
  env <- Proxy registry
-> ByteString -> Text -> PoolConfig -> IO (SimpleEnv registry)
forall (registry :: JobPayloadRegistry) (m :: * -> *).
MonadIO m =>
Proxy registry
-> ByteString -> Text -> PoolConfig -> m (SimpleEnv registry)
createSimpleEnvWithConfig (forall (t :: JobPayloadRegistry). Proxy t
forall {k} (t :: k). Proxy t
Proxy @registry) ByteString
connStr Text
schemaName PoolConfig
serverPoolConfig
  hub <- newMVar Nothing
  rlCache <- newCacheCell
  ccCache <- newCacheCell
  statsCache <- newCacheCell
  perQueueCache <- newCacheCell
  healthCell <- newCacheCell
  pure
    ArbiterServerConfig
      { serverEnv = env
      , enableSSE = True
      , sseHub = hub
      , rateLimitPoliciesCache = rlCache
      , concurrencyPoliciesCache = ccCache
      , allQueueStatsCache = statsCache
      , queueStatsCache = perQueueCache
      , queueStatsCacheTtl = defaultQueueStatsCacheTtl
      , healthCache = healthCell
      , maintenanceInterval = defaultMaintenanceInterval
      , maintenanceSparseInterval = defaultMaintenanceSparseInterval
      , maintenanceBucketIdle = defaultMaintenanceBucketIdle
      , maintenanceTimeout = defaultMaintenanceTimeout
      }

-- | Jobs API handlers for a specific table.
jobsServer
  :: forall registry payload result
   . (EncodeJobResult result, JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> JobsAPI payload result (AsServerT Handler)
jobsServer :: forall (registry :: JobPayloadRegistry) payload result.
(EncodeJobResult result, JobPayload payload) =>
Text
-> ArbiterServerConfig registry
-> JobsAPI payload result (AsServerT Handler)
jobsServer Text
table ArbiterServerConfig registry
config =
  JobsAPI
    { listJobs :: AsServerT Handler
:- (QueryParam "limit" Int
    :> (QueryParam "offset" Int
        :> (QueryParam "group_key" Text
            :> (QueryParam "parent_id" Int64
                :> (QueryParam "job_id" Int64
                    :> (QueryFlag "roots_only"
                        :> (QueryParam "status" JobStatus
                            :> (QueryParam "claimed_by" UUID
                                :> (QueryParam "kind" Text
                                    :> (QueryParam "payload" Text
                                        :> (QueryParam "rate_limit_prefix" Text
                                            :> (QueryParam "concurrency_prefix" Text
                                                :> (QueryParam "sort_by" JobSortColumn
                                                    :> (QueryParam "sort_dir" SortDir
                                                        :> Get
                                                             '[JSON]
                                                             (JobsResponse payload)))))))))))))))
listJobs = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Maybe Int
-> Maybe Int
-> Maybe Text
-> Maybe Int64
-> Maybe Int64
-> Bool
-> Maybe JobStatus
-> Maybe UUID
-> Maybe Text
-> Maybe Text
-> Maybe Text
-> Maybe Text
-> Maybe JobSortColumn
-> Maybe SortDir
-> Handler (JobsResponse payload)
listJobsHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , insertJob :: AsServerT Handler
:- (ReqBody '[JSON] (ApiJobWrite payload)
    :> Post '[JSON] (JobResponse (ApiJob payload)))
insertJob = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> ApiJobWrite payload
-> Handler (JobResponse (ApiJob payload))
insertJobHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , insertJobsBatch :: AsServerT Handler
:- ("batch"
    :> (ReqBody '[JSON] (BatchInsertRequest payload)
        :> Post '[JSON] (BatchInsertResponse payload)))
insertJobsBatch = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> BatchInsertRequest payload
-> Handler (BatchInsertResponse payload)
insertJobsBatchHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , getJob :: AsServerT Handler
:- (Capture "id" Int64
    :> Get '[JSON] (JobResponse (ApiJobWithStatus payload)))
getJob = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> Handler (JobResponse (ApiJobWithStatus payload))
getJobHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , cancelJob :: AsServerT Handler :- (Capture "id" Int64 :> DeleteNoContent)
cancelJob = forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
cancelJobHandler @registry Text
table ArbiterServerConfig registry
config
    , forceCancelJob :: AsServerT Handler
:- (Capture "id" Int64 :> ("force-cancel" :> PostNoContent))
forceCancelJob = forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
forceCancelJobHandler @registry Text
table ArbiterServerConfig registry
config
    , promoteJob :: AsServerT Handler
:- (Capture "id" Int64 :> ("promote" :> PostNoContent))
promoteJob = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
promoteJobHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , moveToDLQ :: AsServerT Handler
:- (Capture "id" Int64 :> ("move-to-dlq" :> PostNoContent))
moveToDLQ = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
moveToDLQHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , pauseChildren :: AsServerT Handler
:- (Capture "id" Int64 :> ("pause-children" :> PostNoContent))
pauseChildren = forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
pauseChildrenHandler @registry Text
table ArbiterServerConfig registry
config
    , resumeChildren :: AsServerT Handler
:- (Capture "id" Int64 :> ("resume-children" :> PostNoContent))
resumeChildren = forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
resumeChildrenHandler @registry Text
table ArbiterServerConfig registry
config
    , suspendJob :: AsServerT Handler
:- (Capture "id" Int64 :> ("suspend" :> PostNoContent))
suspendJob = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
suspendJobHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , resumeJob :: AsServerT Handler
:- (Capture "id" Int64 :> ("resume" :> PostNoContent))
resumeJob = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
resumeJobHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , ackClaimedJob :: AsServerT Handler
:- (Capture "id" Int64
    :> ("ack"
        :> (ReqBody '[JSON] (AckRequest result) :> PostNoContent)))
ackClaimedJob = forall (registry :: JobPayloadRegistry) payload result.
(EncodeJobResult result, JobPayload payload) =>
Text
-> ArbiterServerConfig registry
-> Int64
-> AckRequest result
-> Handler NoContent
ackClaimedJobHandler @registry @payload @result Text
table ArbiterServerConfig registry
config
    , nackClaimedJob :: AsServerT Handler
:- (Capture "id" Int64
    :> ("nack" :> (ReqBody '[JSON] JobLease :> PostNoContent)))
nackClaimedJob = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> JobLease
-> Handler NoContent
nackClaimedJobHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , extendClaimedJob :: AsServerT Handler
:- (Capture "id" Int64
    :> ("extend" :> (ReqBody '[JSON] ExtendRequest :> PostNoContent)))
extendClaimedJob = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> ExtendRequest
-> Handler NoContent
extendClaimedJobHandler @registry @payload Text
table ArbiterServerConfig registry
config
    }

-- | List jobs with pagination and composable filters.
listJobsHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Maybe Int
  -> Maybe Int
  -> Maybe Text
  -> Maybe Int64
  -> Maybe Int64
  -> Bool
  -> Maybe JobStatus
  -> Maybe UUID
  -> Maybe Text
  -> Maybe Text
  -> Maybe Text
  -> Maybe Text
  -> Maybe JobSortColumn
  -> Maybe SortDir
  -> Handler (JobsResponse payload)
listJobsHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Maybe Int
-> Maybe Int
-> Maybe Text
-> Maybe Int64
-> Maybe Int64
-> Bool
-> Maybe JobStatus
-> Maybe UUID
-> Maybe Text
-> Maybe Text
-> Maybe Text
-> Maybe Text
-> Maybe JobSortColumn
-> Maybe SortDir
-> Handler (JobsResponse payload)
listJobsHandler Text
tableName ArbiterServerConfig registry
config Maybe Int
mLimit Maybe Int
mOffset Maybe Text
mGroupKey Maybe Int64
mParentId Maybe Int64
mJobId Bool
rootsOnly Maybe JobStatus
mStatus Maybe UUID
mClaimedBy Maybe Text
mKind Maybe Text
mPayload Maybe Text
mRatePrefix Maybe Text
mConcPrefix Maybe JobSortColumn
mSortBy Maybe SortDir
mSortDir = IO (JobsResponse payload) -> Handler (JobsResponse payload)
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO (JobsResponse payload) -> Handler (JobsResponse payload))
-> IO (JobsResponse payload) -> Handler (JobsResponse payload)
forall a b. (a -> b) -> a -> b
$ do
  let (Int
limit, Int
offset) = Int -> Maybe Int -> Maybe Int -> (Int, Int)
validatePagination Int
50 Maybe Int
mLimit Maybe Int
mOffset
      schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
      filters :: [JobFilter]
filters =
        [Maybe JobFilter] -> [JobFilter]
forall a. [Maybe a] -> [a]
catMaybes
          [ Text -> JobFilter
FilterGroupKey (Text -> JobFilter) -> Maybe Text -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text
mGroupKey
          , Int64 -> JobFilter
FilterParentId (Int64 -> JobFilter) -> Maybe Int64 -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Int64
mParentId
          , Int64 -> JobFilter
FilterId (Int64 -> JobFilter) -> Maybe Int64 -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Int64
mJobId
          , JobFilter
FilterRootsOnly JobFilter -> Maybe () -> Maybe JobFilter
forall a b. a -> Maybe b -> Maybe a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Bool -> Maybe ()
forall (f :: * -> *). Alternative f => Bool -> f ()
guard Bool
rootsOnly
          , JobStatus -> JobFilter
FilterStatus (JobStatus -> JobFilter) -> Maybe JobStatus -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe JobStatus
mStatus
          , UUID -> JobFilter
FilterClaimedBy (UUID -> JobFilter) -> Maybe UUID -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe UUID
mClaimedBy
          , Text -> JobFilter
FilterKind (Text -> JobFilter) -> Maybe Text -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text -> Maybe Text
nonBlank Maybe Text
mKind
          , Text -> JobFilter
FilterPayloadText (Text -> JobFilter) -> Maybe Text -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text -> Maybe Text
nonBlank Maybe Text
mPayload
          , Text -> JobFilter
FilterRateLimitPrefix (Text -> JobFilter) -> Maybe Text -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text
mRatePrefix
          , Text -> JobFilter
FilterConcurrencyPrefix (Text -> JobFilter) -> Maybe Text -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text
mConcPrefix
          ]

  (jobs, total, combined, dlqCounts) <- ArbiterServerConfig registry
-> SimpleDb
     registry
     IO
     ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
      Map Int64 Int64)
-> IO
     ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
      Map Int64 Int64)
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb
   registry
   IO
   ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
    Map Int64 Int64)
 -> IO
      ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
       Map Int64 Int64))
-> SimpleDb
     registry
     IO
     ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
      Map Int64 Int64)
-> IO
     ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
      Map Int64 Int64)
forall a b. (a -> b) -> a -> b
$ SimpleDb
  registry
  IO
  ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
   Map Int64 Int64)
-> SimpleDb
     registry
     IO
     ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
      Map Int64 Int64)
forall a. SimpleDb registry IO a -> SimpleDb registry IO a
forall (m :: * -> *) a. MonadArbiter m => m a -> m a
withDbTransaction (SimpleDb
   registry
   IO
   ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
    Map Int64 Int64)
 -> SimpleDb
      registry
      IO
      ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
       Map Int64 Int64))
-> SimpleDb
     registry
     IO
     ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
      Map Int64 Int64)
-> SimpleDb
     registry
     IO
     ([(JobRead payload, JobStatus)], Int64, Map Int64 (Int64, Int64),
      Map Int64 Int64)
forall a b. (a -> b) -> a -> b
$ do
    page <- Text
-> Text
-> [JobFilter]
-> Maybe JobSortColumn
-> Maybe SortDir
-> Int
-> Int
-> SimpleDb registry IO [(JobRead payload, JobStatus)]
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text
-> Text
-> [JobFilter]
-> Maybe JobSortColumn
-> Maybe SortDir
-> Int
-> Int
-> m [(JobRead payload, JobStatus)]
Ops.listJobsWithStatus Text
schemaName Text
tableName [JobFilter]
filters Maybe JobSortColumn
mSortBy Maybe SortDir
mSortDir Int
limit Int
offset
    matching <- Ops.countJobsFiltered schemaName tableName filters
    -- Every parent is a rollup finalizer. A page without one skips the count queries.
    let jobIds = ((JobRead payload, JobStatus) -> Int64)
-> [(JobRead payload, JobStatus)] -> [Int64]
forall a b. (a -> b) -> [a] -> [b]
map (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
Job.primaryKey (JobRead payload -> Int64)
-> ((JobRead payload, JobStatus) -> JobRead payload)
-> (JobRead payload, JobStatus)
-> Int64
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (JobRead payload, JobStatus) -> JobRead payload
forall a b. (a, b) -> a
fst) [(JobRead payload, JobStatus)]
page
        hasParents = ((JobRead payload, JobStatus) -> Bool)
-> [(JobRead payload, JobStatus)] -> Bool
forall (t :: * -> *) a. Foldable t => (a -> Bool) -> t a -> Bool
any (JobRead payload -> Bool
forall p q t adm. Job p Int64 q t adm -> Bool
isRollup (JobRead payload -> Bool)
-> ((JobRead payload, JobStatus) -> JobRead payload)
-> (JobRead payload, JobStatus)
-> Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (JobRead payload, JobStatus) -> JobRead payload
forall a b. (a, b) -> a
fst) [(JobRead payload, JobStatus)]
page
    if null page || not hasParents
      then pure (page, matching, Map.empty, Map.empty)
      else do
        children <- Ops.countChildrenBatch schemaName tableName jobIds
        dlqChildren <- Ops.countDLQChildrenBatch schemaName tableName jobIds
        pure (page, matching, children, dlqChildren)

  let childCounts = ((Int64, Int64) -> Int64)
-> Map Int64 (Int64, Int64) -> Map Int64 Int64
forall a b. (a -> b) -> Map Int64 a -> Map Int64 b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap (Int64, Int64) -> Int64
forall a b. (a, b) -> a
fst Map Int64 (Int64, Int64)
combined
      pausedParents = Map Int64 (Int64, Int64) -> [Int64]
forall k a. Map k a -> [k]
Map.keys (Map Int64 (Int64, Int64) -> [Int64])
-> Map Int64 (Int64, Int64) -> [Int64]
forall a b. (a -> b) -> a -> b
$ ((Int64, Int64) -> Bool)
-> Map Int64 (Int64, Int64) -> Map Int64 (Int64, Int64)
forall a k. (a -> Bool) -> Map k a -> Map k a
Map.filter (\(Int64
childTotal, Int64
childPaused) -> Int64
childPaused Int64 -> Int64 -> Bool
forall a. Eq a => a -> a -> Bool
== Int64
childTotal) Map Int64 (Int64, Int64)
combined
      apiJobs = ((JobRead payload, JobStatus) -> ApiJobWithStatus payload)
-> [(JobRead payload, JobStatus)] -> [ApiJobWithStatus payload]
forall a b. (a -> b) -> [a] -> [b]
map ((JobRead payload -> JobStatus -> ApiJobWithStatus payload)
-> (JobRead payload, JobStatus) -> ApiJobWithStatus payload
forall a b c. (a -> b -> c) -> (a, b) -> c
uncurry JobRead payload -> JobStatus -> ApiJobWithStatus payload
forall payload.
JobRead payload -> JobStatus -> ApiJobWithStatus payload
ApiJobWithStatus) [(JobRead payload, JobStatus)]
jobs
  pure $
    JobsResponse
      { jobs = apiJobs
      , jobsTotal = fromIntegral total
      , jobsOffset = offset
      , jobsLimit = limit
      , childCounts = childCounts
      , pausedParents = pausedParents
      , dlqChildCounts = dlqCounts
      }

-- | Insert a new job into the queue.
insertJobHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> ApiJobWrite payload
  -> Handler (JobResponse (ApiJob payload))
insertJobHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> ApiJobWrite payload
-> Handler (JobResponse (ApiJob payload))
insertJobHandler Text
tableName ArbiterServerConfig registry
config (ApiJobWrite JobWrite payload
jobWrite) = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  mJob <- ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe (JobRead payload))
-> Handler (Maybe (JobRead payload))
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO (Maybe (JobRead payload))
 -> Handler (Maybe (JobRead payload)))
-> SimpleDb registry IO (Maybe (JobRead payload))
-> Handler (Maybe (JobRead payload))
forall a b. (a -> b) -> a -> b
$ Text
-> [JobWrite payload]
-> SimpleDb registry IO (Maybe (JobRead payload))
-> SimpleDb registry IO (Maybe (JobRead payload))
forall payload (m :: * -> *) a.
(HasKind payload, MonadUnliftIO m) =>
Text -> [JobWrite payload] -> m a -> m a
withPublishSpan Text
tableName [JobWrite payload
jobWrite] (SimpleDb registry IO (Maybe (JobRead payload))
 -> SimpleDb registry IO (Maybe (JobRead payload)))
-> SimpleDb registry IO (Maybe (JobRead payload))
-> SimpleDb registry IO (Maybe (JobRead payload))
forall a b. (a -> b) -> a -> b
$ do
    inserted <- Text
-> Text
-> JobWrite payload
-> SimpleDb registry IO (Maybe (JobRead payload))
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text -> Text -> JobWrite payload -> m (Maybe (JobRead payload))
Ops.insertJob Text
schemaName Text
tableName JobWrite payload
jobWrite
    case (inserted, Job.dedupKey jobWrite) of
      (Just JobRead payload
fresh, Maybe DedupKey
_) -> Maybe (JobRead payload)
-> SimpleDb registry IO (Maybe (JobRead payload))
forall a. a -> SimpleDb registry IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (JobRead payload -> Maybe (JobRead payload)
forall a. a -> Maybe a
Just JobRead payload
fresh)
      (Maybe (JobRead payload)
Nothing, Just (IgnoreDuplicate Text
duplicateKey)) -> Text
-> Text -> Text -> SimpleDb registry IO (Maybe (JobRead payload))
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text -> Text -> Text -> m (Maybe (JobRead payload))
Ops.getJobByDedupKey Text
schemaName Text
tableName Text
duplicateKey
      (Maybe (JobRead payload), Maybe DedupKey)
_ -> Maybe (JobRead payload)
-> SimpleDb registry IO (Maybe (JobRead payload))
forall a. a -> SimpleDb registry IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe (JobRead payload)
forall a. Maybe a
Nothing
  case mJob of
    Just JobRead payload
found -> JobResponse (ApiJob payload)
-> Handler (JobResponse (ApiJob payload))
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (JobResponse (ApiJob payload)
 -> Handler (JobResponse (ApiJob payload)))
-> JobResponse (ApiJob payload)
-> Handler (JobResponse (ApiJob payload))
forall a b. (a -> b) -> a -> b
$ ApiJob payload -> JobResponse (ApiJob payload)
forall a. a -> JobResponse a
JobResponse (JobRead payload -> ApiJob payload
forall payload. JobRead payload -> ApiJob payload
ApiJob JobRead payload
found)
    Maybe (JobRead payload)
Nothing ->
      ServerError -> Handler (JobResponse (ApiJob payload))
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err409 {errBody = "Replace blocked: existing job is actively claimed, force-cancel flagged, or has children"}

-- | Insert multiple jobs in a single batch operation.
insertJobsBatchHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> BatchInsertRequest payload
  -> Handler (BatchInsertResponse payload)
insertJobsBatchHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> BatchInsertRequest payload
-> Handler (BatchInsertResponse payload)
insertJobsBatchHandler Text
tableName ArbiterServerConfig registry
config (BatchInsertRequest [ApiJobWrite payload]
jobWrites) = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
      writes :: [JobWrite payload]
writes = (ApiJobWrite payload -> JobWrite payload)
-> [ApiJobWrite payload] -> [JobWrite payload]
forall a b. (a -> b) -> [a] -> [b]
map ApiJobWrite payload -> JobWrite payload
forall payload. ApiJobWrite payload -> JobWrite payload
unApiJobWrite [ApiJobWrite payload]
jobWrites

  inserted <-
    ArbiterServerConfig registry
-> SimpleDb registry IO [JobRead payload]
-> Handler [JobRead payload]
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config
      (SimpleDb registry IO [JobRead payload]
 -> Handler [JobRead payload])
-> SimpleDb registry IO [JobRead payload]
-> Handler [JobRead payload]
forall a b. (a -> b) -> a -> b
$ Text
-> [JobWrite payload]
-> SimpleDb registry IO [JobRead payload]
-> SimpleDb registry IO [JobRead payload]
forall payload (m :: * -> *) a.
(HasKind payload, MonadUnliftIO m) =>
Text -> [JobWrite payload] -> m a -> m a
withPublishSpan Text
tableName [JobWrite payload]
writes
      (SimpleDb registry IO [JobRead payload]
 -> SimpleDb registry IO [JobRead payload])
-> SimpleDb registry IO [JobRead payload]
-> SimpleDb registry IO [JobRead payload]
forall a b. (a -> b) -> a -> b
$ Text
-> Text
-> [JobWrite payload]
-> SimpleDb registry IO [JobRead payload]
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text -> Text -> [JobWrite payload] -> m [JobRead payload]
Ops.insertJobsBatch Text
schemaName Text
tableName [JobWrite payload]
writes
  let apiJobs = (JobRead payload -> ApiJob payload)
-> [JobRead payload] -> [ApiJob payload]
forall a b. (a -> b) -> [a] -> [b]
map JobRead payload -> ApiJob payload
forall payload. JobRead payload -> ApiJob payload
ApiJob [JobRead payload]
inserted
  pure $ BatchInsertResponse {inserted = apiJobs, insertedCount = length apiJobs}

-- | Fetch a job by id.
getJobHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler (JobResponse (ApiJobWithStatus payload))
getJobHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> Handler (JobResponse (ApiJobWithStatus payload))
getJobHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  mJob <- ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe (JobRead payload, JobStatus))
-> Handler (Maybe (JobRead payload, JobStatus))
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO (Maybe (JobRead payload, JobStatus))
 -> Handler (Maybe (JobRead payload, JobStatus)))
-> SimpleDb registry IO (Maybe (JobRead payload, JobStatus))
-> Handler (Maybe (JobRead payload, JobStatus))
forall a b. (a -> b) -> a -> b
$ Text
-> Text
-> Int64
-> SimpleDb registry IO (Maybe (JobRead payload, JobStatus))
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text -> Text -> Int64 -> m (Maybe (JobRead payload, JobStatus))
Ops.getJobByIdWithStatus Text
schemaName Text
tableName Int64
jobId
  case mJob of
    Maybe (JobRead payload, JobStatus)
Nothing -> ServerError -> Handler (JobResponse (ApiJobWithStatus payload))
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err404 {errBody = "Job not found"}
    Just (JobRead payload
found, JobStatus
jobStatus) -> JobResponse (ApiJobWithStatus payload)
-> Handler (JobResponse (ApiJobWithStatus payload))
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (JobResponse (ApiJobWithStatus payload)
 -> Handler (JobResponse (ApiJobWithStatus payload)))
-> JobResponse (ApiJobWithStatus payload)
-> Handler (JobResponse (ApiJobWithStatus payload))
forall a b. (a -> b) -> a -> b
$ JobResponse {job :: ApiJobWithStatus payload
job = JobRead payload -> JobStatus -> ApiJobWithStatus payload
forall payload.
JobRead payload -> JobStatus -> ApiJobWithStatus payload
ApiJobWithStatus JobRead payload
found JobStatus
jobStatus}

-- | Cancel a job (delete it from the queue).
cancelJobHandler
  :: forall registry
   . Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
cancelJobHandler :: forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
cancelJobHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (Text -> Text -> Int64 -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.cancelJobCascade Text
schemaName Text
tableName Int64
jobId) Handler Int64 -> (Int64 -> Handler NoContent) -> Handler NoContent
forall a b. Handler a -> (a -> Handler b) -> Handler b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= ByteString -> Int64 -> Handler NoContent
rowsOr404 ByteString
"Job not found"

-- | Cascade-cancel a job and async-cancel any in-flight handlers via NOTIFY.
forceCancelJobHandler
  :: forall registry
   . Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
forceCancelJobHandler :: forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
forceCancelJobHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (Text -> Text -> Int64 -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.forceCancelJob Text
schemaName Text
tableName Int64
jobId) Handler Int64 -> (Int64 -> Handler NoContent) -> Handler NoContent
forall a b. Handler a -> (a -> Handler b) -> Handler b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= ByteString -> Int64 -> Handler NoContent
rowsOr404 ByteString
"Job not found"

-- | Promote a job (make it immediately visible).
promoteJobHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
promoteJobHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
promoteJobHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId =
  forall payload (registry :: JobPayloadRegistry).
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> (Text -> SimpleDb registry IO Int64)
-> (JobRead payload -> ByteString)
-> Handler NoContent
mutateJob @payload Text
tableName ArbiterServerConfig registry
config Int64
jobId (\Text
schemaName -> Text -> Text -> Int64 -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.promoteJob Text
schemaName Text
tableName Int64
jobId) JobRead payload -> ByteString
forall {a} {payload} {q} {insertedAt} {adm}.
IsString a =>
JobRecord payload Int64 q insertedAt adm -> a
refuse
  where
    refuse :: JobRecord payload Int64 q insertedAt adm -> a
refuse JobRecord payload Int64 q insertedAt adm
job
      | JobRecord payload Int64 q insertedAt adm -> Bool
forall p q t adm. Job p Int64 q t adm -> Bool
Job.suspended JobRecord payload Int64 q insertedAt adm
job = a
"Job is suspended - use resume endpoint"
      | Maybe UUID -> Bool
forall a. Maybe a -> Bool
isJust (JobRecord payload Int64 q insertedAt adm -> Maybe UUID
forall payload q insertedAt adm.
JobRecord payload Int64 q insertedAt adm -> Maybe UUID
Job.claimedBy JobRecord payload Int64 q insertedAt adm
job) = a
"Job is in flight - wait for its lease to lapse"
      | Bool
otherwise = a
"Job is already visible"

-- | Move a job to the dead letter queue.
moveToDLQHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
moveToDLQHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
moveToDLQHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  result <- ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe Int64) -> Handler (Maybe Int64)
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO (Maybe Int64) -> Handler (Maybe Int64))
-> SimpleDb registry IO (Maybe Int64) -> Handler (Maybe Int64)
forall a b. (a -> b) -> a -> b
$ SimpleDb registry IO (Maybe Int64)
-> SimpleDb registry IO (Maybe Int64)
forall a. SimpleDb registry IO a -> SimpleDb registry IO a
forall (m :: * -> *) a. MonadArbiter m => m a -> m a
withDbTransaction (SimpleDb registry IO (Maybe Int64)
 -> SimpleDb registry IO (Maybe Int64))
-> SimpleDb registry IO (Maybe Int64)
-> SimpleDb registry IO (Maybe Int64)
forall a b. (a -> b) -> a -> b
$ do
    mJob <- forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text -> Text -> Int64 -> m (Maybe (JobRead payload))
Ops.getJobById @_ @payload Text
schemaName Text
tableName Int64
jobId
    case mJob of
      Maybe (JobRead payload)
Nothing -> Maybe Int64 -> SimpleDb registry IO (Maybe Int64)
forall a. a -> SimpleDb registry IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe Int64
forall a. Maybe a
Nothing
      Just JobRead payload
job ->
        Int64 -> Maybe Int64
forall a. a -> Maybe a
Just (Int64 -> Maybe Int64)
-> SimpleDb registry IO Int64 -> SimpleDb registry IO (Maybe Int64)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> TreeLocks
-> Text
-> Text
-> Text
-> JobRead payload
-> SimpleDb registry IO Int64
forall (m :: * -> *) payload.
MonadArbiter m =>
TreeLocks -> Text -> Text -> Text -> JobRead payload -> m Int64
Ops.moveToDLQ TreeLocks
Ops.TakeLocks Text
schemaName Text
tableName Text
"Manually moved to DLQ via admin API" JobRead payload
job

  case result of
    Maybe Int64
Nothing -> ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err404 {errBody = "Job not found"}
    Just Int64
0 -> ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err409 {errBody = "Job was concurrently modified"}
    Just Int64
_ -> NoContent -> Handler NoContent
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure NoContent
NoContent

-- | Pause all children of a parent job.
pauseChildrenHandler
  :: forall registry
   . Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
pauseChildrenHandler :: forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
pauseChildrenHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  Handler Int64 -> Handler ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (Handler Int64 -> Handler ())
-> (SimpleDb registry IO Int64 -> Handler Int64)
-> SimpleDb registry IO Int64
-> Handler ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO Int64 -> Handler ())
-> SimpleDb registry IO Int64 -> Handler ()
forall a b. (a -> b) -> a -> b
$
    Text -> Text -> Int64 -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.pauseChildren Text
schemaName Text
tableName Int64
jobId

  -- Pausing nothing is a success. The children may be in flight, suspended or done.
  NoContent -> Handler NoContent
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure NoContent
NoContent

-- | Resume all suspended children of a parent job.
resumeChildrenHandler
  :: forall registry
   . Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
resumeChildrenHandler :: forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
resumeChildrenHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  Handler Int64 -> Handler ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (Handler Int64 -> Handler ())
-> (SimpleDb registry IO Int64 -> Handler Int64)
-> SimpleDb registry IO Int64
-> Handler ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO Int64 -> Handler ())
-> SimpleDb registry IO Int64 -> Handler ()
forall a b. (a -> b) -> a -> b
$
    Text -> Text -> Int64 -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.resumeChildren Text
schemaName Text
tableName Int64
jobId

  -- Resuming nothing is a success. The children may be unsuspended or done.
  NoContent -> Handler NoContent
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure NoContent
NoContent

-- | Suspend a job (make it unclaimable).
suspendJobHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
suspendJobHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
suspendJobHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId =
  forall payload (registry :: JobPayloadRegistry).
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> (Text -> SimpleDb registry IO Int64)
-> (JobRead payload -> ByteString)
-> Handler NoContent
mutateJob @payload Text
tableName ArbiterServerConfig registry
config Int64
jobId (\Text
schemaName -> Text -> Text -> Int64 -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.suspendJob Text
schemaName Text
tableName Int64
jobId) JobRead payload -> ByteString
forall {a} {payload} {q} {insertedAt} {adm}.
IsString a =>
JobRecord payload Int64 q insertedAt adm -> a
refuse
  where
    refuse :: JobRecord payload Int64 q insertedAt adm -> a
refuse JobRecord payload Int64 q insertedAt adm
job
      | JobRecord payload Int64 q insertedAt adm -> Bool
forall p q t adm. Job p Int64 q t adm -> Bool
Job.suspended JobRecord payload Int64 q insertedAt adm
job = a
"Job is already suspended"
      | Bool
otherwise = a
"Job is in-flight - cannot suspend"

-- | Resume a suspended job, making it claimable again. Refuses a finalizer with children
-- still running.
resumeJobHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
resumeJobHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
resumeJobHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId =
  forall payload (registry :: JobPayloadRegistry).
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> (Text -> SimpleDb registry IO Int64)
-> (JobRead payload -> ByteString)
-> Handler NoContent
mutateJob @payload Text
tableName ArbiterServerConfig registry
config Int64
jobId (\Text
schemaName -> Text -> Text -> Int64 -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.resumeJob Text
schemaName Text
tableName Int64
jobId) JobRead payload -> ByteString
forall {a} {payload} {q} {insertedAt} {adm}.
IsString a =>
JobRecord payload Int64 q insertedAt adm -> a
refuse
  where
    refuse :: JobRecord p Int64 q t adm -> a
refuse JobRecord p Int64 q t adm
job
      | Bool -> Bool
not (JobRecord p Int64 q t adm -> Bool
forall p q t adm. Job p Int64 q t adm -> Bool
Job.suspended JobRecord p Int64 q t adm
job) = a
"Job is not suspended"
      | JobRecord p Int64 q t adm -> Bool
forall p q t adm. Job p Int64 q t adm -> Bool
isRollup JobRecord p Int64 q t adm
job = a
"Cannot resume a rollup finalizer with active children"
      | Bool
otherwise = a
"Job could not be resumed (concurrent modification)"

-- | DLQ API handlers for a specific table.
dlqServer
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> DLQAPI payload (AsServerT Handler)
dlqServer :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> DLQAPI payload (AsServerT Handler)
dlqServer Text
table ArbiterServerConfig registry
config =
  DLQAPI
    { listDLQ :: AsServerT Handler
:- (QueryParam "limit" Int
    :> (QueryParam "offset" Int
        :> (QueryParam "parent_id" Int64
            :> (QueryParam "job_id" Int64
                :> (QueryParam "group_key" Text
                    :> (QueryParam "kind" Text
                        :> (QueryParam "sort_by" DLQSortColumn
                            :> (QueryParam "sort_dir" SortDir
                                :> Get '[JSON] (DLQResponse payload)))))))))
listDLQ = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Maybe Int
-> Maybe Int
-> Maybe Int64
-> Maybe Int64
-> Maybe Text
-> Maybe Text
-> Maybe DLQSortColumn
-> Maybe SortDir
-> Handler (DLQResponse payload)
listDLQHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , retryFromDLQ :: AsServerT Handler
:- (Capture "id" Int64 :> ("retry" :> PostNoContent))
retryFromDLQ = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
retryFromDLQHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , deleteDLQ :: AsServerT Handler :- (Capture "id" Int64 :> DeleteNoContent)
deleteDLQ = forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
deleteDLQHandler @registry Text
table ArbiterServerConfig registry
config
    , deleteDLQBatch :: AsServerT Handler
:- ("batch-delete"
    :> (ReqBody '[JSON] BatchDeleteRequest
        :> Post '[JSON] BatchDeleteResponse))
deleteDLQBatch = forall (registry :: JobPayloadRegistry).
Text
-> ArbiterServerConfig registry
-> BatchDeleteRequest
-> Handler BatchDeleteResponse
deleteDLQBatchHandler @registry Text
table ArbiterServerConfig registry
config
    }

-- | List DLQ jobs with pagination and composable filters.
listDLQHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Maybe Int
  -> Maybe Int
  -> Maybe Int64
  -> Maybe Int64
  -> Maybe Text
  -> Maybe Text
  -> Maybe DLQSortColumn
  -> Maybe SortDir
  -> Handler (DLQResponse payload)
listDLQHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Maybe Int
-> Maybe Int
-> Maybe Int64
-> Maybe Int64
-> Maybe Text
-> Maybe Text
-> Maybe DLQSortColumn
-> Maybe SortDir
-> Handler (DLQResponse payload)
listDLQHandler Text
tableName ArbiterServerConfig registry
config Maybe Int
mLimit Maybe Int
mOffset Maybe Int64
mParentId Maybe Int64
mJobId Maybe Text
mGroupKey Maybe Text
mKind Maybe DLQSortColumn
mSortBy Maybe SortDir
mSortDir = do
  let (Int
limit, Int
offset) = Int -> Maybe Int -> Maybe Int -> (Int, Int)
validatePagination Int
50 Maybe Int
mLimit Maybe Int
mOffset
      schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
      filters :: [JobFilter]
filters =
        [Maybe JobFilter] -> [JobFilter]
forall a. [Maybe a] -> [a]
catMaybes
          [ Int64 -> JobFilter
FilterParentId (Int64 -> JobFilter) -> Maybe Int64 -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Int64
mParentId
          , Int64 -> JobFilter
FilterJobId (Int64 -> JobFilter) -> Maybe Int64 -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Int64
mJobId
          , Text -> JobFilter
FilterGroupKey (Text -> JobFilter) -> Maybe Text -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text
mGroupKey
          , Text -> JobFilter
FilterKind (Text -> JobFilter) -> Maybe Text -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text -> Maybe Text
nonBlank Maybe Text
mKind
          ]

  (dlqJobs, total) <- ArbiterServerConfig registry
-> SimpleDb registry IO ([DLQJob payload], Int64)
-> Handler ([DLQJob payload], Int64)
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO ([DLQJob payload], Int64)
 -> Handler ([DLQJob payload], Int64))
-> SimpleDb registry IO ([DLQJob payload], Int64)
-> Handler ([DLQJob payload], Int64)
forall a b. (a -> b) -> a -> b
$ SimpleDb registry IO ([DLQJob payload], Int64)
-> SimpleDb registry IO ([DLQJob payload], Int64)
forall a. SimpleDb registry IO a -> SimpleDb registry IO a
forall (m :: * -> *) a. MonadArbiter m => m a -> m a
withDbTransaction (SimpleDb registry IO ([DLQJob payload], Int64)
 -> SimpleDb registry IO ([DLQJob payload], Int64))
-> SimpleDb registry IO ([DLQJob payload], Int64)
-> SimpleDb registry IO ([DLQJob payload], Int64)
forall a b. (a -> b) -> a -> b
$ do
    page <- Text
-> Text
-> [JobFilter]
-> Maybe DLQSortColumn
-> Maybe SortDir
-> Int
-> Int
-> SimpleDb registry IO [DLQJob payload]
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text
-> Text
-> [JobFilter]
-> Maybe DLQSortColumn
-> Maybe SortDir
-> Int
-> Int
-> m [DLQJob payload]
Ops.listDLQFilteredOrdered Text
schemaName Text
tableName [JobFilter]
filters Maybe DLQSortColumn
mSortBy Maybe SortDir
mSortDir Int
limit Int
offset
    matching <- Ops.countDLQFiltered schemaName tableName filters
    pure (page, matching)

  let apiDlqJobs = (DLQJob payload -> ApiDLQJob payload)
-> [DLQJob payload] -> [ApiDLQJob payload]
forall a b. (a -> b) -> [a] -> [b]
map DLQJob payload -> ApiDLQJob payload
forall payload. DLQJob payload -> ApiDLQJob payload
ApiDLQJob [DLQJob payload]
dlqJobs
  pure $
    DLQResponse
      { dlqJobs = apiDlqJobs
      , dlqTotal = fromIntegral total
      , dlqOffset = offset
      , dlqLimit = limit
      }

-- | Retry a DLQ job back into the main queue. 409 when its parent is gone.
retryFromDLQHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
retryFromDLQHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
retryFromDLQHandler Text
tableName ArbiterServerConfig registry
config Int64
dlqId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  result <- ArbiterServerConfig registry
-> SimpleDb registry IO (Either Bool ())
-> Handler (Either Bool ())
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO (Either Bool ()) -> Handler (Either Bool ()))
-> SimpleDb registry IO (Either Bool ())
-> Handler (Either Bool ())
forall a b. (a -> b) -> a -> b
$ SimpleDb registry IO (Either Bool ())
-> SimpleDb registry IO (Either Bool ())
forall a. SimpleDb registry IO a -> SimpleDb registry IO a
forall (m :: * -> *) a. MonadArbiter m => m a -> m a
withDbTransaction (SimpleDb registry IO (Either Bool ())
 -> SimpleDb registry IO (Either Bool ()))
-> SimpleDb registry IO (Either Bool ())
-> SimpleDb registry IO (Either Bool ())
forall a b. (a -> b) -> a -> b
$ do
    mJob <- forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text -> Text -> Int64 -> m (Maybe (JobRead payload))
Ops.retryFromDLQ @_ @payload Text
schemaName Text
tableName Int64
dlqId
    case mJob of
      Just JobRead payload
_ -> Either Bool () -> SimpleDb registry IO (Either Bool ())
forall a. a -> SimpleDb registry IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (() -> Either Bool ()
forall a b. b -> Either a b
Right ())
      Maybe (JobRead payload)
Nothing -> Bool -> Either Bool ()
forall a b. a -> Either a b
Left (Bool -> Either Bool ())
-> SimpleDb registry IO Bool
-> SimpleDb registry IO (Either Bool ())
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Text -> Text -> Int64 -> SimpleDb registry IO Bool
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Bool
Ops.dlqJobExists Text
schemaName Text
tableName Int64
dlqId

  case result of
    Right () -> NoContent -> Handler NoContent
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure NoContent
NoContent
    Left Bool
True -> ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err409 {errBody = "Cannot retry: parent job no longer exists (not in queue or DLQ)"}
    Left Bool
False -> ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err404 {errBody = "DLQ job not found"}

-- | Delete a job from DLQ permanently.
deleteDLQHandler
  :: forall registry
   . Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
deleteDLQHandler :: forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
deleteDLQHandler Text
tableName ArbiterServerConfig registry
config Int64
dlqId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (Text -> Text -> Int64 -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.deleteDLQJob Text
schemaName Text
tableName Int64
dlqId) Handler Int64 -> (Int64 -> Handler NoContent) -> Handler NoContent
forall a b. Handler a -> (a -> Handler b) -> Handler b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= ByteString -> Int64 -> Handler NoContent
rowsOr404 ByteString
"DLQ job not found"

-- | Batch delete jobs from DLQ permanently.
deleteDLQBatchHandler
  :: forall registry
   . Text
  -> ArbiterServerConfig registry
  -> BatchDeleteRequest
  -> Handler BatchDeleteResponse
deleteDLQBatchHandler :: forall (registry :: JobPayloadRegistry).
Text
-> ArbiterServerConfig registry
-> BatchDeleteRequest
-> Handler BatchDeleteResponse
deleteDLQBatchHandler Text
tableName ArbiterServerConfig registry
config (BatchDeleteRequest [Int64]
dlqIds) = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  rowsDeleted <- ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO Int64 -> Handler Int64)
-> SimpleDb registry IO Int64 -> Handler Int64
forall a b. (a -> b) -> a -> b
$ Text -> Text -> [Int64] -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> [Int64] -> m Int64
Ops.deleteDLQJobsBatch Text
schemaName Text
tableName [Int64]
dlqIds
  pure $ BatchDeleteResponse {deleted = rowsDeleted}

-- | Archive API handler for a specific table.
archiveServer
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> ArchiveAPI payload (AsServerT Handler)
archiveServer :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> ArchiveAPI payload (AsServerT Handler)
archiveServer Text
table ArbiterServerConfig registry
config =
  ArchiveAPI
    { listArchive :: AsServerT Handler
:- (QueryParam "limit" Int
    :> (QueryParam "offset" Int
        :> (QueryParam "parent_id" Int64
            :> (QueryParam "job_id" Int64
                :> (QueryParam "group_key" Text
                    :> (QueryParam "kind" Text
                        :> (QueryParam "completed_after" UTCTime
                            :> (QueryParam "completed_before" UTCTime
                                :> (QueryParam "sort_by" ArchiveSortColumn
                                    :> (QueryParam "sort_dir" SortDir
                                        :> Get '[JSON] (ArchiveResponse payload)))))))))))
listArchive = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Maybe Int
-> Maybe Int
-> Maybe Int64
-> Maybe Int64
-> Maybe Text
-> Maybe Text
-> Maybe UTCTime
-> Maybe UTCTime
-> Maybe ArchiveSortColumn
-> Maybe SortDir
-> Handler (ArchiveResponse payload)
listArchiveHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , reEnqueueArchive :: AsServerT Handler
:- (Capture "id" Int64 :> ("reenqueue" :> PostNoContent))
reEnqueueArchive = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
reEnqueueArchiveHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , deleteArchive :: AsServerT Handler :- (Capture "id" Int64 :> DeleteNoContent)
deleteArchive = forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
deleteArchiveHandler @registry Text
table ArbiterServerConfig registry
config
    , deleteArchiveBatch :: AsServerT Handler
:- ("batch-delete"
    :> (ReqBody '[JSON] BatchDeleteRequest
        :> Post '[JSON] BatchDeleteResponse))
deleteArchiveBatch = forall (registry :: JobPayloadRegistry).
Text
-> ArbiterServerConfig registry
-> BatchDeleteRequest
-> Handler BatchDeleteResponse
deleteArchiveBatchHandler @registry Text
table ArbiterServerConfig registry
config
    }

-- | List archived jobs with pagination and composable filters.
listArchiveHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Maybe Int
  -> Maybe Int
  -> Maybe Int64
  -> Maybe Int64
  -> Maybe Text
  -> Maybe Text
  -> Maybe UTCTime
  -> Maybe UTCTime
  -> Maybe ArchiveSortColumn
  -> Maybe SortDir
  -> Handler (ArchiveResponse payload)
listArchiveHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Maybe Int
-> Maybe Int
-> Maybe Int64
-> Maybe Int64
-> Maybe Text
-> Maybe Text
-> Maybe UTCTime
-> Maybe UTCTime
-> Maybe ArchiveSortColumn
-> Maybe SortDir
-> Handler (ArchiveResponse payload)
listArchiveHandler Text
tableName ArbiterServerConfig registry
config Maybe Int
mLimit Maybe Int
mOffset Maybe Int64
mParentId Maybe Int64
mJobId Maybe Text
mGroupKey Maybe Text
mKind Maybe UTCTime
mCompletedAfter Maybe UTCTime
mCompletedBefore Maybe ArchiveSortColumn
mSortBy Maybe SortDir
mSortDir = do
  let (Int
limit, Int
offset) = Int -> Maybe Int -> Maybe Int -> (Int, Int)
validatePagination Int
50 Maybe Int
mLimit Maybe Int
mOffset
      schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
      filters :: [JobFilter]
filters =
        [Maybe JobFilter] -> [JobFilter]
forall a. [Maybe a] -> [a]
catMaybes
          [ Int64 -> JobFilter
FilterParentId (Int64 -> JobFilter) -> Maybe Int64 -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Int64
mParentId
          , Int64 -> JobFilter
FilterJobId (Int64 -> JobFilter) -> Maybe Int64 -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Int64
mJobId
          , Text -> JobFilter
FilterGroupKey (Text -> JobFilter) -> Maybe Text -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text
mGroupKey
          , Text -> JobFilter
FilterKind (Text -> JobFilter) -> Maybe Text -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text -> Maybe Text
nonBlank Maybe Text
mKind
          , UTCTime -> JobFilter
FilterCompletedAfter (UTCTime -> JobFilter) -> Maybe UTCTime -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe UTCTime
mCompletedAfter
          , UTCTime -> JobFilter
FilterCompletedBefore (UTCTime -> JobFilter) -> Maybe UTCTime -> Maybe JobFilter
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe UTCTime
mCompletedBefore
          ]

  (archived, total) <- ArbiterServerConfig registry
-> SimpleDb registry IO ([ArchiveJob payload], Int64)
-> Handler ([ArchiveJob payload], Int64)
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO ([ArchiveJob payload], Int64)
 -> Handler ([ArchiveJob payload], Int64))
-> SimpleDb registry IO ([ArchiveJob payload], Int64)
-> Handler ([ArchiveJob payload], Int64)
forall a b. (a -> b) -> a -> b
$ SimpleDb registry IO ([ArchiveJob payload], Int64)
-> SimpleDb registry IO ([ArchiveJob payload], Int64)
forall a. SimpleDb registry IO a -> SimpleDb registry IO a
forall (m :: * -> *) a. MonadArbiter m => m a -> m a
withDbTransaction (SimpleDb registry IO ([ArchiveJob payload], Int64)
 -> SimpleDb registry IO ([ArchiveJob payload], Int64))
-> SimpleDb registry IO ([ArchiveJob payload], Int64)
-> SimpleDb registry IO ([ArchiveJob payload], Int64)
forall a b. (a -> b) -> a -> b
$ do
    page <- Text
-> Text
-> [JobFilter]
-> Maybe ArchiveSortColumn
-> Maybe SortDir
-> Int
-> Int
-> SimpleDb registry IO [ArchiveJob payload]
forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text
-> Text
-> [JobFilter]
-> Maybe ArchiveSortColumn
-> Maybe SortDir
-> Int
-> Int
-> m [ArchiveJob payload]
Ops.listArchiveFiltered Text
schemaName Text
tableName [JobFilter]
filters Maybe ArchiveSortColumn
mSortBy Maybe SortDir
mSortDir Int
limit Int
offset
    matching <- Ops.countArchiveFiltered schemaName tableName filters
    pure (page, matching)

  pure $
    ArchiveResponse
      { archiveJobs = map ApiArchiveJob archived
      , archiveTotal = fromIntegral total
      , archiveOffset = offset
      , archiveLimit = limit
      }

-- | Re-enqueue an archived job as a fresh job. 404 if the archive row is gone.
reEnqueueArchiveHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
reEnqueueArchiveHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
reEnqueueArchiveHandler Text
tableName ArbiterServerConfig registry
config Int64
archiveId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  mJob <- ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe (JobRead payload))
-> Handler (Maybe (JobRead payload))
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO (Maybe (JobRead payload))
 -> Handler (Maybe (JobRead payload)))
-> SimpleDb registry IO (Maybe (JobRead payload))
-> Handler (Maybe (JobRead payload))
forall a b. (a -> b) -> a -> b
$ forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text -> Text -> Int64 -> m (Maybe (JobRead payload))
Ops.reEnqueueFromArchive @_ @payload Text
schemaName Text
tableName Int64
archiveId
  case mJob of
    Just JobRead payload
_ -> NoContent -> Handler NoContent
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure NoContent
NoContent
    Maybe (JobRead payload)
Nothing -> ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err404 {errBody = "Archived job not found"}

-- | Purge one archived job by its archive primary key.
deleteArchiveHandler
  :: forall registry
   . Text
  -> ArbiterServerConfig registry
  -> Int64
  -> Handler NoContent
deleteArchiveHandler :: forall (registry :: JobPayloadRegistry).
Text -> ArbiterServerConfig registry -> Int64 -> Handler NoContent
deleteArchiveHandler Text
tableName ArbiterServerConfig registry
config Int64
archiveId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (Text -> Text -> Int64 -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.deleteArchiveJob Text
schemaName Text
tableName Int64
archiveId) Handler Int64 -> (Int64 -> Handler NoContent) -> Handler NoContent
forall a b. Handler a -> (a -> Handler b) -> Handler b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= ByteString -> Int64 -> Handler NoContent
rowsOr404 ByteString
"Archived job not found"

-- | Bulk-purge archived jobs by archive primary key.
deleteArchiveBatchHandler
  :: forall registry
   . Text
  -> ArbiterServerConfig registry
  -> BatchDeleteRequest
  -> Handler BatchDeleteResponse
deleteArchiveBatchHandler :: forall (registry :: JobPayloadRegistry).
Text
-> ArbiterServerConfig registry
-> BatchDeleteRequest
-> Handler BatchDeleteResponse
deleteArchiveBatchHandler Text
tableName ArbiterServerConfig registry
config (BatchDeleteRequest [Int64]
archiveIds) = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  rowsDeleted <- ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO Int64 -> Handler Int64)
-> SimpleDb registry IO Int64 -> Handler Int64
forall a b. (a -> b) -> a -> b
$ Text -> Text -> [Int64] -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> [Int64] -> m Int64
Ops.deleteArchiveJobsBatch Text
schemaName Text
tableName [Int64]
archiveIds
  pure $ BatchDeleteResponse {deleted = rowsDeleted}

-- | Stats API handler for a specific table.
statsServer
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> StatsAPI (AsServerT Handler)
statsServer :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry -> StatsAPI (AsServerT Handler)
statsServer Text
tableName ArbiterServerConfig registry
config =
  StatsAPI
    { getStats :: AsServerT Handler :- Get '[JSON] StatsResponse
getStats = forall (registry :: JobPayloadRegistry).
Text
-> [Text] -> ArbiterServerConfig registry -> Handler StatsResponse
getStatsHandler @registry Text
tableName (forall payload. HasKind payload => [Text]
kindsFor @payload) ArbiterServerConfig registry
config
    }

-- | Get queue statistics.
getStatsHandler
  :: forall registry
   . Text
  -> [Text]
  -> ArbiterServerConfig registry
  -> Handler StatsResponse
getStatsHandler :: forall (registry :: JobPayloadRegistry).
Text
-> [Text] -> ArbiterServerConfig registry -> Handler StatsResponse
getStatsHandler Text
tableName [Text]
kinds ArbiterServerConfig registry
config =
  IO StatsResponse -> Handler StatsResponse
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO StatsResponse -> Handler StatsResponse)
-> IO StatsResponse -> Handler StatsResponse
forall a b. (a -> b) -> a -> b
$ NominalDiffTime
-> CacheCell StatsResponse
-> Text
-> IO StatsResponse
-> IO StatsResponse
forall a. NominalDiffTime -> CacheCell a -> Text -> IO a -> IO a
cachedForKey (ArbiterServerConfig registry -> NominalDiffTime
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> NominalDiffTime
queueStatsCacheTtl ArbiterServerConfig registry
config) (ArbiterServerConfig registry -> CacheCell StatsResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell StatsResponse
queueStatsCache ArbiterServerConfig registry
config) Text
tableName (IO StatsResponse -> IO StatsResponse)
-> IO StatsResponse -> IO StatsResponse
forall a b. (a -> b) -> a -> b
$ do
    let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config

    queueStats <- ArbiterServerConfig registry
-> SimpleDb registry IO QueueStats -> IO QueueStats
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO QueueStats -> IO QueueStats)
-> SimpleDb registry IO QueueStats -> IO QueueStats
forall a b. (a -> b) -> a -> b
$ Text -> Text -> [Text] -> SimpleDb registry IO QueueStats
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> [Text] -> m QueueStats
Ops.getQueueStats Text
schemaName Text
tableName [Text]
kinds
    now <- getCurrentTime
    let timestamp = String -> Text
T.pack (String -> Text) -> String -> Text
forall a b. (a -> b) -> a -> b
$ TimeLocale -> String -> UTCTime -> String
forall t. FormatTime t => TimeLocale -> String -> t -> String
formatTime TimeLocale
defaultTimeLocale String
"%Y-%m-%dT%H:%M:%S%z" UTCTime
now

    pure $ StatsResponse {stats = queueStats, timestamp = timestamp}

-- | Every queue's stats in one request, for the landing overview.
getAllStatsHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> [(Text, [Text])]
  -> Handler AllStatsResponse
getAllStatsHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> [(Text, [Text])] -> Handler AllStatsResponse
getAllStatsHandler ArbiterServerConfig registry
config [(Text, [Text])]
queueKinds =
  IO AllStatsResponse -> Handler AllStatsResponse
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO AllStatsResponse -> Handler AllStatsResponse)
-> IO AllStatsResponse -> Handler AllStatsResponse
forall a b. (a -> b) -> a -> b
$ NominalDiffTime
-> CacheCell AllStatsResponse
-> IO AllStatsResponse
-> IO AllStatsResponse
forall a. NominalDiffTime -> CacheCell a -> IO a -> IO a
cachedFor NominalDiffTime
overviewStatsCacheTtl (ArbiterServerConfig registry -> CacheCell AllStatsResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell AllStatsResponse
allQueueStatsCache ArbiterServerConfig registry
config) (IO AllStatsResponse -> IO AllStatsResponse)
-> IO AllStatsResponse -> IO AllStatsResponse
forall a b. (a -> b) -> a -> b
$ do
    let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
    [QueueOverview] -> AllStatsResponse
AllStatsResponse ([QueueOverview] -> AllStatsResponse)
-> IO [QueueOverview] -> IO AllStatsResponse
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> ArbiterServerConfig registry
-> SimpleDb registry IO [QueueOverview] -> IO [QueueOverview]
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (Text -> [(Text, [Text])] -> SimpleDb registry IO [QueueOverview]
forall (m :: * -> *).
MonadArbiter m =>
Text -> [(Text, [Text])] -> m [QueueOverview]
Ops.getAllQueueStats Text
schemaName [(Text, [Text])]
queueKinds)

-- | Table API handlers for a specific table.
tableServer
  :: forall registry payload result
   . (EncodeJobResult result, JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> TableAPI payload result (AsServerT Handler)
tableServer :: forall (registry :: JobPayloadRegistry) payload result.
(EncodeJobResult result, JobPayload payload) =>
Text
-> ArbiterServerConfig registry
-> TableAPI payload result (AsServerT Handler)
tableServer Text
table ArbiterServerConfig registry
config =
  TableAPI
    { jobs :: AsServerT Handler
:- ("jobs" :> NamedRoutes (JobsAPI payload result))
jobs = forall (registry :: JobPayloadRegistry) payload result.
(EncodeJobResult result, JobPayload payload) =>
Text
-> ArbiterServerConfig registry
-> JobsAPI payload result (AsServerT Handler)
jobsServer @registry @payload @result Text
table ArbiterServerConfig registry
config
    , claimJobs :: AsServerT Handler
:- ("claim"
    :> (ReqBody '[JSON] ClaimRequest
        :> Post '[JSON] (ClaimResponse payload)))
claimJobs = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> ClaimRequest
-> Handler (ClaimResponse payload)
claimJobsHandler @registry @payload Text
table ArbiterServerConfig registry
config
    , dlq :: AsServerT Handler :- ("dlq" :> NamedRoutes (DLQAPI payload))
dlq = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> DLQAPI payload (AsServerT Handler)
dlqServer @registry @payload Text
table ArbiterServerConfig registry
config
    , archive :: AsServerT Handler
:- ("archive" :> NamedRoutes (ArchiveAPI payload))
archive = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> ArchiveAPI payload (AsServerT Handler)
archiveServer @registry @payload Text
table ArbiterServerConfig registry
config
    , stats :: AsServerT Handler :- ("stats" :> NamedRoutes StatsAPI)
stats = forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry -> StatsAPI (AsServerT Handler)
statsServer @registry @payload Text
table ArbiterServerConfig registry
config
    , listKinds :: AsServerT Handler :- ("kinds" :> Get '[JSON] [Text])
listKinds = [Text] -> Handler [Text]
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (forall payload. HasKind payload => [Text]
kindsFor @payload)
    }

-- | Lease visible jobs to a consumer outside a worker pool. Each returned job
-- contains the claim sequence and claimant required for finalization.
claimJobsHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> ClaimRequest
  -> Handler (ClaimResponse payload)
claimJobsHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> ClaimRequest
-> Handler (ClaimResponse payload)
claimJobsHandler Text
tableName ArbiterServerConfig registry
config ClaimRequest
req = IO (ClaimResponse payload) -> Handler (ClaimResponse payload)
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO (ClaimResponse payload) -> Handler (ClaimResponse payload))
-> IO (ClaimResponse payload) -> Handler (ClaimResponse payload)
forall a b. (a -> b) -> a -> b
$ do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
      maxJobs :: Int
maxJobs = Int -> Int -> Int -> Int
forall a. Ord a => a -> a -> a -> a
clamp Int
1 Int
1000 (Int -> Maybe Int -> Int
forall a. a -> Maybe a -> a
fromMaybe Int
1 (ClaimRequest -> Maybe Int
crMaxJobs ClaimRequest
req))
      leaseSecs :: NominalDiffTime
leaseSecs = Double -> NominalDiffTime
forall a b. (Real a, Fractional b) => a -> b
realToFrac (Double -> Double -> Double -> Double
forall a. Ord a => a -> a -> a -> a
clamp Double
1 Double
3600 (Double -> Maybe Double -> Double
forall a. a -> Maybe a -> a
fromMaybe Double
60 (ClaimRequest -> Maybe Double
crLeaseSeconds ClaimRequest
req)))
  claimant <- IO UUID
UUID.nextRandom
  jobs <- runDb config $ do
    mQueue <- Ops.getQueue schemaName tableName
    if maybe False Queues.paused mQueue
      then pure []
      else Ops.claimNextVisibleJobsAs @_ @payload schemaName tableName maxJobs leaseSecs claimant
  pure $ ClaimResponse (map ApiJob jobs)

-- | Complete a job that the caller holds. Store an optional result in the
-- parent rollup or archive entry, as worker @ackWith@ does.
ackClaimedJobHandler
  :: forall registry payload result
   . (EncodeJobResult result, JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> AckRequest result
  -> Handler NoContent
ackClaimedJobHandler :: forall (registry :: JobPayloadRegistry) payload result.
(EncodeJobResult result, JobPayload payload) =>
Text
-> ArbiterServerConfig registry
-> Int64
-> AckRequest result
-> Handler NoContent
ackClaimedJobHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId AckRequest result
req =
  forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> JobLease
-> (Text -> JobRead payload -> SimpleDb registry IO Int64)
-> Handler NoContent
withHeldJob @registry @payload Text
tableName ArbiterServerConfig registry
config Int64
jobId (AckRequest result -> JobLease
forall result. AckRequest result -> JobLease
arLease AckRequest result
req) ((Text -> JobRead payload -> SimpleDb registry IO Int64)
 -> Handler NoContent)
-> (Text -> JobRead payload -> SimpleDb registry IO Int64)
-> Handler NoContent
forall a b. (a -> b) -> a -> b
$ \Text
schemaName JobRead payload
job ->
    SimpleDb registry IO Int64 -> SimpleDb registry IO Int64
forall a. SimpleDb registry IO a -> SimpleDb registry IO a
forall (m :: * -> *) a. MonadArbiter m => m a -> m a
withDbTransaction (SimpleDb registry IO Int64 -> SimpleDb registry IO Int64)
-> SimpleDb registry IO Int64 -> SimpleDb registry IO Int64
forall a b. (a -> b) -> a -> b
$ do
      rows <- Text -> Text -> JobRead payload -> SimpleDb registry IO Int64
forall (m :: * -> *) payload.
MonadArbiter m =>
Text -> Text -> JobRead payload -> m Int64
Ops.ackJob Text
schemaName Text
tableName JobRead payload
job
      when (rows > 0) $ storeEncodedResult schemaName job (arResult req >>= encodeJobResult)
      pure rows

-- | Restore the attempt used by a claim. The job becomes available when its
-- lease expires.
nackClaimedJobHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> JobLease
  -> Handler NoContent
nackClaimedJobHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> JobLease
-> Handler NoContent
nackClaimedJobHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId JobLease
lease =
  forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> JobLease
-> (Text -> JobRead payload -> SimpleDb registry IO Int64)
-> Handler NoContent
withHeldJob @registry @payload Text
tableName ArbiterServerConfig registry
config Int64
jobId JobLease
lease ((Text -> JobRead payload -> SimpleDb registry IO Int64)
 -> Handler NoContent)
-> (Text -> JobRead payload -> SimpleDb registry IO Int64)
-> Handler NoContent
forall a b. (a -> b) -> a -> b
$ \Text
schemaName JobRead payload
job ->
    Text -> Text -> JobRead payload -> SimpleDb registry IO Int64
forall (m :: * -> *) payload.
MonadArbiter m =>
Text -> Text -> JobRead payload -> m Int64
Ops.nackJob Text
schemaName Text
tableName JobRead payload
job

-- | Extend a held lease. This is the HTTP consumer equivalent of a worker heartbeat.
extendClaimedJobHandler
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> ExtendRequest
  -> Handler NoContent
extendClaimedJobHandler :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> ExtendRequest
-> Handler NoContent
extendClaimedJobHandler Text
tableName ArbiterServerConfig registry
config Int64
jobId ExtendRequest
req =
  forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> JobLease
-> (Text -> JobRead payload -> SimpleDb registry IO Int64)
-> Handler NoContent
withHeldJob @registry @payload Text
tableName ArbiterServerConfig registry
config Int64
jobId (ExtendRequest -> JobLease
erLease ExtendRequest
req) ((Text -> JobRead payload -> SimpleDb registry IO Int64)
 -> Handler NoContent)
-> (Text -> JobRead payload -> SimpleDb registry IO Int64)
-> Handler NoContent
forall a b. (a -> b) -> a -> b
$ \Text
schemaName JobRead payload
job ->
    Text
-> Text
-> NominalDiffTime
-> JobRead payload
-> SimpleDb registry IO Int64
forall (m :: * -> *) payload.
MonadArbiter m =>
Text -> Text -> NominalDiffTime -> JobRead payload -> m Int64
Ops.setVisibilityTimeout Text
schemaName Text
tableName (Double -> NominalDiffTime
forall a b. (Real a, Fractional b) => a -> b
realToFrac (Double -> Double -> Double -> Double
forall a. Ord a => a -> a -> a -> a
clamp Double
1 Double
3600 (ExtendRequest -> Double
erSeconds ExtendRequest
req))) JobRead payload
job

-- | Finalize the job identified by a lease. Refuse a lease that the caller no
-- longer holds or a lease held by a worker pool. Each statement checks the claim
-- sequence and writes no change after a lease is lost.
withHeldJob
  :: forall registry payload
   . (JobPayload payload)
  => Text
  -> ArbiterServerConfig registry
  -> Int64
  -> JobLease
  -> (Text -> Job.JobRead payload -> SimpleDb registry IO Int64)
  -> Handler NoContent
withHeldJob :: forall (registry :: JobPayloadRegistry) payload.
JobPayload payload =>
Text
-> ArbiterServerConfig registry
-> Int64
-> JobLease
-> (Text -> JobRead payload -> SimpleDb registry IO Int64)
-> Handler NoContent
withHeldJob Text
tableName ArbiterServerConfig registry
config Int64
jobId JobLease
lease Text -> JobRead payload -> SimpleDb registry IO Int64
finalize = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
      held :: JobRead payload -> Bool
held JobRead payload
job = JobRead payload -> Maybe UUID
forall payload q insertedAt adm.
JobRecord payload Int64 q insertedAt adm -> Maybe UUID
Job.claimedBy JobRead payload
job Maybe UUID -> Maybe UUID -> Bool
forall a. Eq a => a -> a -> Bool
== UUID -> Maybe UUID
forall a. a -> Maybe a
Just (JobLease -> UUID
jlClaimedBy JobLease
lease) Bool -> Bool -> Bool
&& JobRead payload -> Int64
forall payload q insertedAt adm.
JobRecord payload Int64 q insertedAt adm -> Int64
Job.claimSeq JobRead payload
job Int64 -> Int64 -> Bool
forall a. Eq a => a -> a -> Bool
== JobLease -> Int64
jlClaimSeq JobLease
lease
      refuse :: ByteString -> Either ServerError b
refuse ByteString
body = ServerError -> Either ServerError b
forall a b. a -> Either a b
Left ServerError
err409 {errBody = body}

  result <- ArbiterServerConfig registry
-> SimpleDb registry IO (Either ServerError ())
-> Handler (Either ServerError ())
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO (Either ServerError ())
 -> Handler (Either ServerError ()))
-> SimpleDb registry IO (Either ServerError ())
-> Handler (Either ServerError ())
forall a b. (a -> b) -> a -> b
$ do
    mJob <- forall (m :: * -> *) payload.
(JobPayload payload, MonadArbiter m) =>
Text -> Text -> Int64 -> m (Maybe (JobRead payload))
Ops.getJobById @_ @payload Text
schemaName Text
tableName Int64
jobId
    case mJob of
      Maybe (JobRead payload)
Nothing -> Either ServerError ()
-> SimpleDb registry IO (Either ServerError ())
forall a. a -> SimpleDb registry IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Either ServerError ()
 -> SimpleDb registry IO (Either ServerError ()))
-> Either ServerError ()
-> SimpleDb registry IO (Either ServerError ())
forall a b. (a -> b) -> a -> b
$ ServerError -> Either ServerError ()
forall a b. a -> Either a b
Left ServerError
err404 {errBody = "Job not found"}
      Just JobRead payload
job
        | Bool -> Bool
not (JobRead payload -> Bool
held JobRead payload
job) -> Either ServerError ()
-> SimpleDb registry IO (Either ServerError ())
forall a. a -> SimpleDb registry IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Either ServerError ()
 -> SimpleDb registry IO (Either ServerError ()))
-> Either ServerError ()
-> SimpleDb registry IO (Either ServerError ())
forall a b. (a -> b) -> a -> b
$ ByteString -> Either ServerError ()
forall {b}. ByteString -> Either ServerError b
refuse ByteString
"Job is not held by this lease"
        | Bool
otherwise -> do
            pooled <- Text -> UUID -> SimpleDb registry IO Bool
forall (m :: * -> *). MonadArbiter m => Text -> UUID -> m Bool
Ops.workerRegistered Text
schemaName (JobLease -> UUID
jlClaimedBy JobLease
lease)
            if pooled
              then pure $ refuse "Job is held by a worker pool"
              else do
                rowsAffected <- finalize schemaName job
                pure $
                  if rowsAffected > 0
                    then Right ()
                    else refuse (if Job.suspended job then "Job is suspended" else "Lease no longer held")

  noContentOr result

-- | Maintenance API handler.
maintenanceServer
  :: forall registry
   . (HL.RegistryAdmissionPolicies registry, RegistryTables registry)
  => ArbiterServerConfig registry
  -> MaintenanceAPI (AsServerT Handler)
maintenanceServer :: forall (registry :: JobPayloadRegistry).
(RegistryAdmissionPolicies registry, RegistryTables registry) =>
ArbiterServerConfig registry -> MaintenanceAPI (AsServerT Handler)
maintenanceServer ArbiterServerConfig registry
config = MaintenanceAPI {runMaintenance :: AsServerT Handler :- Post '[JSON] MaintenanceResponse
runMaintenance = forall (registry :: JobPayloadRegistry).
(RegistryAdmissionPolicies registry, RegistryTables registry) =>
ArbiterServerConfig registry -> Handler MaintenanceResponse
maintenanceHandler @registry ArbiterServerConfig registry
config}

-- | Run one maintenance pass. A worker pool's reaper does the same work. Operations
-- exclude each other across callers. An operation another caller is running is
-- skipped and absent from the response.
maintenanceHandler
  :: forall registry
   . (HL.RegistryAdmissionPolicies registry, RegistryTables registry)
  => ArbiterServerConfig registry
  -> Handler MaintenanceResponse
maintenanceHandler :: forall (registry :: JobPayloadRegistry).
(RegistryAdmissionPolicies registry, RegistryTables registry) =>
ArbiterServerConfig registry -> Handler MaintenanceResponse
maintenanceHandler ArbiterServerConfig registry
config = IO MaintenanceResponse -> Handler MaintenanceResponse
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO MaintenanceResponse -> Handler MaintenanceResponse)
-> IO MaintenanceResponse -> Handler MaintenanceResponse
forall a b. (a -> b) -> a -> b
$ do
  touched <- Map Text Int64 -> IO (IORef (Map Text Int64))
forall a. a -> IO (IORef a)
newIORef Map Text Int64
forall k a. Map k a
Map.empty
  let report MaintenanceOp
operation Int64
rows = IO () -> SimpleDb registry IO ()
forall a. IO a -> SimpleDb registry IO a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> SimpleDb registry IO ())
-> IO () -> SimpleDb registry IO ()
forall a b. (a -> b) -> a -> b
$ IORef (Map Text Int64)
-> (Map Text Int64 -> Map Text Int64) -> IO ()
forall a. IORef a -> (a -> a) -> IO ()
modifyIORef' IORef (Map Text Int64)
touched ((Int64 -> Int64 -> Int64)
-> Text -> Int64 -> Map Text Int64 -> Map Text Int64
forall k a. Ord k => (a -> a -> a) -> k -> a -> Map k a -> Map k a
Map.insertWith Int64 -> Int64 -> Int64
forall a. Num a => a -> a -> a
(+) (MaintenanceOp -> Text
maintenanceOpName MaintenanceOp
operation) Int64
rows)
      pace =
        MaintenancePace
          { paceWindow :: NominalDiffTime
paceWindow = ArbiterServerConfig registry -> NominalDiffTime
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> NominalDiffTime
maintenanceInterval ArbiterServerConfig registry
config
          , paceSparseWindow :: NominalDiffTime
paceSparseWindow = ArbiterServerConfig registry -> NominalDiffTime
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> NominalDiffTime
maintenanceSparseInterval ArbiterServerConfig registry
config
          , paceBucketIdle :: NominalDiffTime
paceBucketIdle = ArbiterServerConfig registry -> NominalDiffTime
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> NominalDiffTime
maintenanceBucketIdle ArbiterServerConfig registry
config
          }
  failed <-
    runDb config $
      runMaintenancePass defaultLogConfig report pace (maintenanceTimeout config)
  ops <- readIORef touched
  pure $ MaintenanceResponse ops (map maintenanceOpName failed)

-- | Hold a caller-supplied value inside the range an endpoint accepts.
clamp :: (Ord a) => a -> a -> a -> a
clamp :: forall a. Ord a => a -> a -> a -> a
clamp a
lower a
upper = a -> a -> a
forall a. Ord a => a -> a -> a
max a
lower (a -> a) -> (a -> a) -> a -> a
forall b c a. (b -> c) -> (a -> b) -> a -> c
. a -> a -> a
forall a. Ord a => a -> a -> a
min a
upper

-- | Queues API handler.
queuesServer
  :: forall registry
   . (RegistryTables registry)
  => Proxy registry
  -> ArbiterServerConfig registry
  -> QueuesAPI (AsServerT Handler)
queuesServer :: forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
Proxy registry
-> ArbiterServerConfig registry -> QueuesAPI (AsServerT Handler)
queuesServer Proxy registry
registryProxy ArbiterServerConfig registry
config =
  let known :: [Text]
known = Proxy registry -> [Text]
forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
Proxy registry -> [Text]
registryTableNames Proxy registry
registryProxy
   in QueuesAPI
        { listQueues :: AsServerT Handler :- Get '[JSON] QueuesResponse
listQueues = QueuesResponse -> Handler QueuesResponse
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (QueuesResponse -> Handler QueuesResponse)
-> QueuesResponse -> Handler QueuesResponse
forall a b. (a -> b) -> a -> b
$ QueuesResponse {queues :: [Text]
queues = [Text]
known}
        , getAllStats :: AsServerT Handler :- ("stats" :> Get '[JSON] AllStatsResponse)
getAllStats = ArbiterServerConfig registry
-> [(Text, [Text])] -> Handler AllStatsResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> [(Text, [Text])] -> Handler AllStatsResponse
getAllStatsHandler ArbiterServerConfig registry
config (Proxy registry -> [(Text, [Text])]
forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
Proxy registry -> [(Text, [Text])]
registryQueueKinds Proxy registry
registryProxy)
        , getDetails :: AsServerT Handler
:- (Capture "queue" Text
    :> ("details" :> Get '[JSON] (Maybe QueueRow)))
getDetails = ArbiterServerConfig registry -> Text -> Handler (Maybe QueueRow)
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text -> Handler (Maybe QueueRow)
getQueueDetailsHandler ArbiterServerConfig registry
config
        , pauseQueue :: AsServerT Handler
:- (Capture "queue" Text :> ("pause" :> PostNoContent))
pauseQueue = ArbiterServerConfig registry
-> [Text] -> Bool -> Text -> Handler NoContent
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> [Text] -> Bool -> Text -> Handler NoContent
setQueuePausedHandler ArbiterServerConfig registry
config [Text]
known Bool
True
        , resumeQueue :: AsServerT Handler
:- (Capture "queue" Text :> ("resume" :> PostNoContent))
resumeQueue = ArbiterServerConfig registry
-> [Text] -> Bool -> Text -> Handler NoContent
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> [Text] -> Bool -> Text -> Handler NoContent
setQueuePausedHandler ArbiterServerConfig registry
config [Text]
known Bool
False
        }

-- | Get a queue's operator config.
getQueueDetailsHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Text
  -> Handler (Maybe QueueRow)
getQueueDetailsHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text -> Handler (Maybe QueueRow)
getQueueDetailsHandler ArbiterServerConfig registry
config Text
queue = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe QueueRow)
-> Handler (Maybe QueueRow)
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO (Maybe QueueRow) -> Handler (Maybe QueueRow))
-> SimpleDb registry IO (Maybe QueueRow)
-> Handler (Maybe QueueRow)
forall a b. (a -> b) -> a -> b
$ Text -> Text -> SimpleDb registry IO (Maybe QueueRow)
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> m (Maybe QueueRow)
Ops.getQueue Text
schemaName Text
queue

-- | Flip the @paused@ flag for a queue, validated against the registry. The
-- @arbiter_queues@ row is created lazily on first pause.
setQueuePausedHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> [Text]
  -> Bool
  -> Text
  -> Handler NoContent
setQueuePausedHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> [Text] -> Bool -> Text -> Handler NoContent
setQueuePausedHandler ArbiterServerConfig registry
config [Text]
knownQueues Bool
pauseFlag Text
queue = do
  Bool -> Handler () -> Handler ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless (Text
queue Text -> [Text] -> Bool
forall a. Eq a => a -> [a] -> Bool
forall (t :: * -> *) a. (Foldable t, Eq a) => a -> t a -> Bool
`elem` [Text]
knownQueues) (Handler () -> Handler ()) -> Handler () -> Handler ()
forall a b. (a -> b) -> a -> b
$
    ServerError -> Handler ()
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err404 {errBody = "Unknown queue"}
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  Handler Int64 -> Handler ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (Handler Int64 -> Handler ())
-> (SimpleDb registry IO Int64 -> Handler Int64)
-> SimpleDb registry IO Int64
-> Handler ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO Int64 -> Handler ())
-> SimpleDb registry IO Int64 -> Handler ()
forall a b. (a -> b) -> a -> b
$ Text -> Text -> Bool -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Bool -> m Int64
Ops.setQueuePaused Text
schemaName Text
queue Bool
pauseFlag
  -- The landing overview shows each queue's paused flag.
  CacheCell AllStatsResponse -> Handler ()
forall a. CacheCell a -> Handler ()
invalidate (ArbiterServerConfig registry -> CacheCell AllStatsResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell AllStatsResponse
allQueueStatsCache ArbiterServerConfig registry
config)
  CacheCell StatsResponse -> Handler ()
forall a. CacheCell a -> Handler ()
invalidate (ArbiterServerConfig registry -> CacheCell StatsResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell StatsResponse
queueStatsCache ArbiterServerConfig registry
config)
  NoContent -> Handler NoContent
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure NoContent
NoContent

-- | Serve the SSE stream as a raw WAI application. Flush after each event and
-- send a keepalive comment every 15 seconds. Each client reads a duplicate of
-- the shared broadcast hub and does not use a PostgreSQL pool connection. If
-- 'enableSSE' is false, send one @disabled@ event and close the stream. The
-- admin UI then stops reconnection attempts.
eventsServer
  :: forall registry
   . ArbiterServerConfig registry
  -> Tagged Handler Application
eventsServer :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Tagged Handler Application
eventsServer ArbiterServerConfig registry
config = Application -> Tagged Handler Application
forall {k} (s :: k) b. b -> Tagged s b
Tagged (Application -> Tagged Handler Application)
-> Application -> Tagged Handler Application
forall a b. (a -> b) -> a -> b
$ \Request
_req Response -> IO ResponseReceived
sendResponse ->
  if Bool -> Bool
not (ArbiterServerConfig registry -> Bool
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Bool
enableSSE ArbiterServerConfig registry
config)
    then Response -> IO ResponseReceived
sendResponse (Response -> IO ResponseReceived)
-> Response -> IO ResponseReceived
forall a b. (a -> b) -> a -> b
$ Status -> ResponseHeaders -> StreamingBody -> Response
responseStream Status
status200 ResponseHeaders
sseHeaders (StreamingBody -> Response) -> StreamingBody -> Response
forall a b. (a -> b) -> a -> b
$ \Builder -> IO ()
write IO ()
flush -> do
      Builder -> IO ()
write Builder
"data: {\"event\":\"disabled\"}\n\n"
      IO ()
flush
    else
      -- 'bracket' pairs the refcount increment with its decrement around the
      -- whole response. The hub is released when the streaming body never runs.
      IO (Maybe (TChan ByteString))
-> (Maybe (TChan ByteString) -> IO ())
-> (Maybe (TChan ByteString) -> IO ResponseReceived)
-> IO ResponseReceived
forall a b c. IO a -> (a -> IO b) -> (a -> IO c) -> IO c
bracket (ArbiterServerConfig registry -> IO (Maybe (TChan ByteString))
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> IO (Maybe (TChan ByteString))
subscribeSSE ArbiterServerConfig registry
config) (IO ()
-> (TChan ByteString -> IO ()) -> Maybe (TChan ByteString) -> IO ()
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (() -> IO ()
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()) (IO () -> TChan ByteString -> IO ()
forall a b. a -> b -> a
const (ArbiterServerConfig registry -> IO ()
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> IO ()
unsubscribeSSE ArbiterServerConfig registry
config))) ((Maybe (TChan ByteString) -> IO ResponseReceived)
 -> IO ResponseReceived)
-> (Maybe (TChan ByteString) -> IO ResponseReceived)
-> IO ResponseReceived
forall a b. (a -> b) -> a -> b
$ \Maybe (TChan ByteString)
mSub ->
        case Maybe (TChan ByteString)
mSub of
          Maybe (TChan ByteString)
Nothing -> Response -> IO ResponseReceived
sendResponse (Response -> IO ResponseReceived)
-> Response -> IO ResponseReceived
forall a b. (a -> b) -> a -> b
$ Status -> ResponseHeaders -> StreamingBody -> Response
responseStream Status
status200 ResponseHeaders
sseHeaders (StreamingBody -> Response) -> StreamingBody -> Response
forall a b. (a -> b) -> a -> b
$ \Builder -> IO ()
write IO ()
flush -> do
            Builder -> IO ()
write Builder
"data: {\"event\":\"error\",\"message\":\"No connection pool available\"}\n\n"
            IO ()
flush
          Just TChan ByteString
sub -> Response -> IO ResponseReceived
sendResponse (Response -> IO ResponseReceived)
-> Response -> IO ResponseReceived
forall a b. (a -> b) -> a -> b
$ Status -> ResponseHeaders -> StreamingBody -> Response
responseStream Status
status200 ResponseHeaders
sseHeaders (StreamingBody -> Response) -> StreamingBody -> Response
forall a b. (a -> b) -> a -> b
$ \Builder -> IO ()
write IO ()
flush ->
            -- A failed write (client gone) ends the stream. The enclosing bracket
            -- then drops this subscriber's refcount.
            (SomeException -> IO ()) -> IO () -> IO ()
forall e a. Exception e => (e -> IO a) -> IO a -> IO a
handle SomeException -> IO ()
swallowSync (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
              Builder -> IO ()
write Builder
"data: {\"event\":\"connected\",\"message\":\"Stream connected\"}\n\n"
              IO ()
flush
              -- Read this client's channel with a 15s keepalive heartbeat.
              let go :: IO ()
go = do
                    mPayload <- Int -> IO ByteString -> IO (Maybe ByteString)
forall a. Int -> IO a -> IO (Maybe a)
timeout Int
15_000_000 (STM ByteString -> IO ByteString
forall a. STM a -> IO a
atomically (TChan ByteString -> STM ByteString
forall a. TChan a -> STM a
readTChan TChan ByteString
sub))
                    case mPayload of
                      Just ByteString
payload -> do
                        Builder -> IO ()
write (Builder
"data: " Builder -> Builder -> Builder
forall a. Semigroup a => a -> a -> a
<> ByteString -> Builder
Builder.byteString ByteString
payload Builder -> Builder -> Builder
forall a. Semigroup a => a -> a -> a
<> Builder
"\n\n")
                        IO ()
flush
                        IO ()
go
                      Maybe ByteString
Nothing -> do
                        Builder -> IO ()
write Builder
": keepalive\n\n"
                        IO ()
flush
                        IO ()
go
              IO ()
go
  where
    sseHeaders :: ResponseHeaders
sseHeaders =
      [ (HeaderName
"Content-Type", ByteString
"text/event-stream")
      , (HeaderName
"Cache-Control", ByteString
"no-cache")
      , (HeaderName
"Connection", ByteString
"keep-alive")
      , (HeaderName
"X-Accel-Buffering", ByteString
"no")
      ]

swallowSync :: SomeException -> IO ()
swallowSync :: SomeException -> IO ()
swallowSync SomeException
exception = case SomeException -> Maybe SomeAsyncException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception of
  Just (SomeAsyncException
_ :: SomeAsyncException) -> SomeException -> IO ()
forall e a. (HasCallStack, Exception e) => e -> IO a
throwIO SomeException
exception
  Maybe SomeAsyncException
Nothing -> () -> IO ()
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()

-- | Subscribe to the shared SSE hub, returning a duplicated channel to stream
-- from. The first subscriber starts the hub (one @LISTEN@ connection). Later
-- subscribers bump the refcount. 'Nothing' when there is no pool.
subscribeSSE :: ArbiterServerConfig registry -> IO (Maybe (TChan ByteString))
subscribeSSE :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> IO (Maybe (TChan ByteString))
subscribeSSE ArbiterServerConfig registry
config =
  MVar (Maybe SSEHub)
-> (Maybe SSEHub -> IO (Maybe SSEHub, Maybe (TChan ByteString)))
-> IO (Maybe (TChan ByteString))
forall a b. MVar a -> (a -> IO (a, b)) -> IO b
modifyMVar (ArbiterServerConfig registry -> MVar (Maybe SSEHub)
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> MVar (Maybe SSEHub)
sseHub ArbiterServerConfig registry
config) ((Maybe SSEHub -> IO (Maybe SSEHub, Maybe (TChan ByteString)))
 -> IO (Maybe (TChan ByteString)))
-> (Maybe SSEHub -> IO (Maybe SSEHub, Maybe (TChan ByteString)))
-> IO (Maybe (TChan ByteString))
forall a b. (a -> b) -> a -> b
$ \Maybe SSEHub
mhub -> case Maybe SSEHub
mhub of
    Just SSEHub
hub -> do
      sub <- STM (TChan ByteString) -> IO (TChan ByteString)
forall a. STM a -> IO a
atomically (STM (TChan ByteString) -> IO (TChan ByteString))
-> STM (TChan ByteString) -> IO (TChan ByteString)
forall a b. (a -> b) -> a -> b
$ do
        TVar Int -> (Int -> Int) -> STM ()
forall a. TVar a -> (a -> a) -> STM ()
modifyTVar' (SSEHub -> TVar Int
hubRefs SSEHub
hub) (Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
1)
        TChan ByteString -> STM (TChan ByteString)
forall a. TChan a -> STM (TChan a)
dupTChan (SSEHub -> TChan ByteString
hubChan SSEHub
hub)
      pure (Just hub, Just sub)
    Maybe SSEHub
Nothing -> case SimpleConnectionPool -> Maybe (Pool Connection)
connectionPool (SimpleEnv registry -> SimpleConnectionPool
forall (registry :: JobPayloadRegistry).
SimpleEnv registry -> SimpleConnectionPool
simplePool (ArbiterServerConfig registry -> SimpleEnv registry
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> SimpleEnv registry
serverEnv ArbiterServerConfig registry
config)) of
      Maybe (Pool Connection)
Nothing -> (Maybe SSEHub, Maybe (TChan ByteString))
-> IO (Maybe SSEHub, Maybe (TChan ByteString))
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (Maybe SSEHub
forall a. Maybe a
Nothing, Maybe (TChan ByteString)
forall a. Maybe a
Nothing)
      Just Pool Connection
pool -> do
        broadcast <- IO (TChan ByteString)
forall a. IO (TChan a)
newBroadcastTChanIO
        refs <- newTVarIO 1
        sub <- atomically (dupTChan broadcast)
        -- subscribeSSE runs masked as a 'bracket' acquire. The listener is
        -- forked with an explicit unmask.
        void $ forkIOWithUnmask $ \forall a. IO a -> IO a
unmask -> IO () -> IO ()
forall a. IO a -> IO a
unmask (Pool Connection -> TChan ByteString -> TVar Int -> IO ()
sseListenerLoop Pool Connection
pool TChan ByteString
broadcast TVar Int
refs)
        pure (Just (SSEHub broadcast refs), Just sub)

-- | Drop one subscriber. When the count reaches zero the hub is removed from
-- the config. The listener observes the same count and releases its connection
-- (see 'sseListenerLoop').
unsubscribeSSE :: ArbiterServerConfig registry -> IO ()
unsubscribeSSE :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> IO ()
unsubscribeSSE ArbiterServerConfig registry
config =
  MVar (Maybe SSEHub) -> (Maybe SSEHub -> IO (Maybe SSEHub)) -> IO ()
forall a. MVar a -> (a -> IO a) -> IO ()
modifyMVar_ (ArbiterServerConfig registry -> MVar (Maybe SSEHub)
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> MVar (Maybe SSEHub)
sseHub ArbiterServerConfig registry
config) ((Maybe SSEHub -> IO (Maybe SSEHub)) -> IO ())
-> (Maybe SSEHub -> IO (Maybe SSEHub)) -> IO ()
forall a b. (a -> b) -> a -> b
$ \Maybe SSEHub
mhub -> case Maybe SSEHub
mhub of
    Maybe SSEHub
Nothing -> Maybe SSEHub -> IO (Maybe SSEHub)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe SSEHub
forall a. Maybe a
Nothing
    Just SSEHub
hub -> do
      remaining <- STM Int -> IO Int
forall a. STM a -> IO a
atomically (STM Int -> IO Int) -> STM Int -> IO Int
forall a b. (a -> b) -> a -> b
$ do
        TVar Int -> (Int -> Int) -> STM ()
forall a. TVar a -> (a -> a) -> STM ()
modifyTVar' (SSEHub -> TVar Int
hubRefs SSEHub
hub) (Int -> Int -> Int
forall a. Num a => a -> a -> a
subtract Int
1)
        TVar Int -> STM Int
forall a. TVar a -> STM a
readTVar (SSEHub -> TVar Int
hubRefs SSEHub
hub)
      pure $ if remaining <= 0 then Nothing else Just hub

-- | The single listener. It runs @LISTEN@ on one borrowed connection and fans
-- every notification into the broadcast channel, raced against the subscriber
-- count reaching zero. On connection loss it destroys the dead resource and
-- reconnects. When the count hits zero the race ends, the 'bracket' releases
-- the connection, and the thread exits.
sseListenerLoop :: Pool.Pool PG.Connection -> TChan ByteString -> TVar Int -> IO ()
sseListenerLoop :: Pool Connection -> TChan ByteString -> TVar Int -> IO ()
sseListenerLoop Pool Connection
pool TChan ByteString
broadcast TVar Int
refs = do
  backoff <- Int -> IO (IORef Int)
forall a. a -> IO (IORef a)
newIORef Int
baseBackoff
  race_ waitForIdle (pump backoff)
  where
    baseBackoff :: Int
baseBackoff = Int
1_000_000 -- 1s
    maxBackoff :: Int
maxBackoff = Int
30_000_000 -- 30s
    waitForIdle :: IO ()
waitForIdle = STM () -> IO ()
forall a. STM a -> IO a
atomically (STM () -> IO ()) -> STM () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
      subscribers <- TVar Int -> STM Int
forall a. TVar a -> STM a
readTVar TVar Int
refs
      check (subscribers == 0)
    pump :: IORef Int -> IO (ZonkAny 0)
pump IORef Int
backoff =
      IO () -> IO (ZonkAny 0)
forall (f :: * -> *) a b. Applicative f => f a -> f b
forever
        (IO () -> IO (ZonkAny 0)) -> IO () -> IO (ZonkAny 0)
forall a b. (a -> b) -> a -> b
$ (SomeException -> IO ()) -> IO () -> IO ()
forall e a. Exception e => (e -> IO a) -> IO a -> IO a
handle (IORef Int -> SomeException -> IO ()
onError IORef Int
backoff)
        (IO () -> IO ()) -> IO () -> IO ()
forall a b. (a -> b) -> a -> b
$ IO (Connection, LocalPool Connection)
-> ((Connection, LocalPool Connection) -> IO ())
-> ((Connection, LocalPool Connection) -> IO ())
-> IO ()
forall a b c. IO a -> (a -> IO b) -> (a -> IO c) -> IO c
bracket
          (Pool Connection -> IO (Connection, LocalPool Connection)
forall a. Pool a -> IO (a, LocalPool a)
Pool.takeResource Pool Connection
pool)
          (\(Connection
conn, LocalPool Connection
localPool) -> Pool Connection -> LocalPool Connection -> Connection -> IO ()
forall a. Pool a -> LocalPool a -> a -> IO ()
Pool.destroyResource Pool Connection
pool LocalPool Connection
localPool Connection
conn)
        (((Connection, LocalPool Connection) -> IO ()) -> IO ())
-> ((Connection, LocalPool Connection) -> IO ()) -> IO ()
forall a b. (a -> b) -> a -> b
$ \(Connection
conn, LocalPool Connection
_) -> do
          _ <- Connection -> Query -> IO Int64
PG.execute_ Connection
conn (Query -> IO Int64) -> Query -> IO Int64
forall a b. (a -> b) -> a -> b
$ Query
"LISTEN " Query -> Query -> Query
forall a. Semigroup a => a -> a -> a
<> String -> Query
forall a. IsString a => String -> a
fromString (Text -> String
T.unpack Text
Schema.eventStreamingChannel)
          writeIORef backoff baseBackoff -- reset on connect
          forever $ do
            notification <- getNotification conn
            atomically $ writeTChan broadcast (notificationData notification)
    -- Re-raise async exceptions. On a sync error log it and retry with capped
    -- exponential backoff.
    onError :: IORef Int -> SomeException -> IO ()
onError IORef Int
backoff SomeException
exception = case SomeException -> Maybe SomeAsyncException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception of
      Just (SomeAsyncException
_ :: SomeAsyncException) -> SomeException -> IO ()
forall e a. (HasCallStack, Exception e) => e -> IO a
throwIO SomeException
exception
      Maybe SomeAsyncException
Nothing -> do
        delay <- IORef Int -> IO Int
forall a. IORef a -> IO a
readIORef IORef Int
backoff
        BS8.hPutStr stderr . encodeUtf8 $
          "[arbiter:sse] listener error, retrying in "
            <> T.pack (show (delay `div` 1_000_000))
            <> "s: "
            <> T.pack (show exception)
            <> "\n"
        threadDelay delay
        writeIORef backoff (min maxBackoff (delay * 2))

-- | Cron API handlers.
cronServer
  :: forall registry
   . ArbiterServerConfig registry
  -> CronAPI (AsServerT Handler)
cronServer :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CronAPI (AsServerT Handler)
cronServer ArbiterServerConfig registry
config =
  CronAPI
    { listSchedules :: AsServerT Handler
:- ("schedules"
    :> (QueryParam "queue" Text :> Get '[JSON] CronSchedulesResponse))
listSchedules = ArbiterServerConfig registry
-> Maybe Text -> Handler CronSchedulesResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Maybe Text -> Handler CronSchedulesResponse
listCronSchedulesHandler ArbiterServerConfig registry
config
    , updateSchedule :: AsServerT Handler
:- ("schedules"
    :> (Capture "name" Text
        :> (ReqBody '[JSON] CronScheduleUpdate
            :> Patch '[JSON] CronScheduleView)))
updateSchedule = ArbiterServerConfig registry
-> Text -> CronScheduleUpdate -> Handler CronScheduleView
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Text -> CronScheduleUpdate -> Handler CronScheduleView
updateCronScheduleHandler ArbiterServerConfig registry
config
    , runSchedule :: AsServerT Handler
:- ("schedules"
    :> (Capture "name" Text :> ("run" :> PostNoContent)))
runSchedule = ArbiterServerConfig registry -> Text -> Handler NoContent
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text -> Handler NoContent
runCronScheduleHandler ArbiterServerConfig registry
config
    }

-- | List cron schedules, optionally scoped to a queue.
listCronSchedulesHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Maybe Text
  -> Handler CronSchedulesResponse
listCronSchedulesHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Maybe Text -> Handler CronSchedulesResponse
listCronSchedulesHandler ArbiterServerConfig registry
config Maybe Text
mQueue = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  rows <- ArbiterServerConfig registry
-> SimpleDb registry IO [CronScheduleRow]
-> Handler [CronScheduleRow]
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO [CronScheduleRow]
 -> Handler [CronScheduleRow])
-> SimpleDb registry IO [CronScheduleRow]
-> Handler [CronScheduleRow]
forall a b. (a -> b) -> a -> b
$ Text -> Maybe Text -> SimpleDb registry IO [CronScheduleRow]
forall (m :: * -> *).
MonadArbiter m =>
Text -> Maybe Text -> m [CronScheduleRow]
Ops.listCronSchedules Text
schemaName Maybe Text
mQueue
  now <- liftIO getCurrentTime
  pure $ CronSchedulesResponse {cronSchedules = map (cronScheduleView now) rows}

-- | A schedule row with the next tick it fires at. A disabled schedule has none.
cronScheduleView :: UTCTime -> CronScheduleRow -> CronScheduleView
cronScheduleView :: UTCTime -> CronScheduleRow -> CronScheduleView
cronScheduleView UTCTime
now row :: CronScheduleRow
row@CS.CronScheduleRow {enabled :: CronScheduleRow -> Bool
CS.enabled = Bool
isEnabled} =
  CronScheduleView
    { schedule :: CronScheduleRow
schedule = CronScheduleRow
row
    , nextRunAt :: Maybe UTCTime
nextRunAt = do
        Bool -> Maybe ()
forall (f :: * -> *). Alternative f => Bool -> f ()
guard Bool
isEnabled
        Maybe Text -> Text -> UTCTime -> Maybe UTCTime
nextRunFromExpression (CronScheduleRow -> Maybe Text
CS.effectiveTimezone CronScheduleRow
row) (CronScheduleRow -> Text
CS.effectiveExpression CronScheduleRow
row) UTCTime
now
    }

-- | Update a cron schedule.
updateCronScheduleHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Text
  -> CronScheduleUpdate
  -> Handler CronScheduleView
updateCronScheduleHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Text -> CronScheduleUpdate -> Handler CronScheduleView
updateCronScheduleHandler ArbiterServerConfig registry
config Text
name CronScheduleUpdate
update = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  result <- ArbiterServerConfig registry
-> SimpleDb registry IO (Either Text (Maybe CronScheduleRow))
-> Handler (Either Text (Maybe CronScheduleRow))
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO (Either Text (Maybe CronScheduleRow))
 -> Handler (Either Text (Maybe CronScheduleRow)))
-> SimpleDb registry IO (Either Text (Maybe CronScheduleRow))
-> Handler (Either Text (Maybe CronScheduleRow))
forall a b. (a -> b) -> a -> b
$ SimpleDb registry IO (Either Text (Maybe CronScheduleRow))
-> SimpleDb registry IO (Either Text (Maybe CronScheduleRow))
forall a. SimpleDb registry IO a -> SimpleDb registry IO a
forall (m :: * -> *) a. MonadArbiter m => m a -> m a
withDbTransaction (SimpleDb registry IO (Either Text (Maybe CronScheduleRow))
 -> SimpleDb registry IO (Either Text (Maybe CronScheduleRow)))
-> SimpleDb registry IO (Either Text (Maybe CronScheduleRow))
-> SimpleDb registry IO (Either Text (Maybe CronScheduleRow))
forall a b. (a -> b) -> a -> b
$ do
    outcome <- Text
-> CronScheduleUpdate -> SimpleDb registry IO (Either Text Int64)
forall (m :: * -> *).
MonadArbiter m =>
Text -> CronScheduleUpdate -> m (Either Text Int64)
updateCronScheduleChecked Text
name CronScheduleUpdate
update
    traverse (const (Ops.getCronScheduleByName schemaName name)) outcome

  case result of
    Left Text
err -> ServerError -> Handler CronScheduleView
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err400 {errBody = LBS.fromStrict (encodeUtf8 err)}
    Right Maybe CronScheduleRow
Nothing -> ServerError -> Handler CronScheduleView
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err404 {errBody = "Cron schedule not found"}
    Right (Just CronScheduleRow
row) -> (UTCTime -> CronScheduleRow -> CronScheduleView)
-> CronScheduleRow -> UTCTime -> CronScheduleView
forall a b c. (a -> b -> c) -> b -> a -> c
flip UTCTime -> CronScheduleRow -> CronScheduleView
cronScheduleView CronScheduleRow
row (UTCTime -> CronScheduleView)
-> Handler UTCTime -> Handler CronScheduleView
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> IO UTCTime -> Handler UTCTime
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO IO UTCTime
getCurrentTime

-- | Request an out-of-band run of a cron schedule. A disabled schedule is
-- refused. A schedule with a run already pending is refused.
runCronScheduleHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Text
  -> Handler NoContent
runCronScheduleHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text -> Handler NoContent
runCronScheduleHandler ArbiterServerConfig registry
config Text
name = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  outcome <- ArbiterServerConfig registry
-> SimpleDb registry IO RunRequestOutcome
-> Handler RunRequestOutcome
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO RunRequestOutcome
 -> Handler RunRequestOutcome)
-> SimpleDb registry IO RunRequestOutcome
-> Handler RunRequestOutcome
forall a b. (a -> b) -> a -> b
$ Text -> Text -> SimpleDb registry IO RunRequestOutcome
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> m RunRequestOutcome
Ops.requestCronRun Text
schemaName Text
name
  case outcome of
    RunRequestOutcome
Ops.RunReqNotFound -> ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err404 {errBody = "Cron schedule not found"}
    RunRequestOutcome
Ops.RunReqDisabled -> ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err409 {errBody = "Cron schedule is disabled"}
    RunRequestOutcome
Ops.RunReqPending -> ServerError -> Handler NoContent
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err409 {errBody = "Cron schedule already has a run pending"}
    RunRequestOutcome
Ops.RunReqStamped -> NoContent -> Handler NoContent
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure NoContent
NoContent

-- | Workers API handlers.
workersServer
  :: forall registry
   . ArbiterServerConfig registry
  -> WorkersAPI (AsServerT Handler)
workersServer :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> WorkersAPI (AsServerT Handler)
workersServer ArbiterServerConfig registry
config =
  WorkersAPI
    { listWorkers :: AsServerT Handler
:- (QueryParam "queue" Text
    :> (QueryParam "live" Double :> Get '[JSON] WorkersResponse))
listWorkers = ArbiterServerConfig registry
-> Maybe Text -> Maybe Double -> Handler WorkersResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Maybe Text -> Maybe Double -> Handler WorkersResponse
listWorkersHandler ArbiterServerConfig registry
config
    , pauseWorker :: AsServerT Handler
:- (Capture "id" UUID :> ("pause" :> PostNoContent))
pauseWorker = ArbiterServerConfig registry -> Bool -> UUID -> Handler NoContent
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Bool -> UUID -> Handler NoContent
setWorkerPausedHandler ArbiterServerConfig registry
config Bool
True
    , resumeWorker :: AsServerT Handler
:- (Capture "id" UUID :> ("resume" :> PostNoContent))
resumeWorker = ArbiterServerConfig registry -> Bool -> UUID -> Handler NoContent
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Bool -> UUID -> Handler NoContent
setWorkerPausedHandler ArbiterServerConfig registry
config Bool
False
    }

-- | List workers, optionally scoped to a queue and/or to recent heartbeats.
listWorkersHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Maybe Text
  -> Maybe Double
  -> Handler WorkersResponse
listWorkersHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Maybe Text -> Maybe Double -> Handler WorkersResponse
listWorkersHandler ArbiterServerConfig registry
config Maybe Text
mQueue Maybe Double
mLiveSecs = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  rows <- ArbiterServerConfig registry
-> SimpleDb registry IO [WorkerRow] -> Handler [WorkerRow]
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (SimpleDb registry IO [WorkerRow] -> Handler [WorkerRow])
-> SimpleDb registry IO [WorkerRow] -> Handler [WorkerRow]
forall a b. (a -> b) -> a -> b
$ Text
-> Maybe Text
-> Maybe NominalDiffTime
-> SimpleDb registry IO [WorkerRow]
forall (m :: * -> *).
MonadArbiter m =>
Text -> Maybe Text -> Maybe NominalDiffTime -> m [WorkerRow]
Ops.listWorkers Text
schemaName Maybe Text
mQueue (Double -> NominalDiffTime
forall a b. (Real a, Fractional b) => a -> b
realToFrac (Double -> NominalDiffTime)
-> Maybe Double -> Maybe NominalDiffTime
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Double
mLiveSecs)
  pure $ WorkersResponse {workers = rows}

-- | Set a worker's @paused@ flag. The worker reconciles its local state on
-- the next heartbeat.
setWorkerPausedHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Bool
  -> UUID
  -> Handler NoContent
setWorkerPausedHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Bool -> UUID -> Handler NoContent
setWorkerPausedHandler ArbiterServerConfig registry
config Bool
pauseFlag UUID
workerId = do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (Text -> UUID -> Bool -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> UUID -> Bool -> m Int64
Ops.setWorkerPaused Text
schemaName UUID
workerId Bool
pauseFlag) Handler Int64 -> (Int64 -> Handler NoContent) -> Handler NoContent
forall a b. Handler a -> (a -> Handler b) -> Handler b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= ByteString -> Int64 -> Handler NoContent
rowsOr404 ByteString
"Worker not found"

-- | Rate-limit management/observability handlers.
rateLimitsServer
  :: forall registry
   . (RegistryTables registry)
  => ArbiterServerConfig registry
  -> RateLimitsAPI (AsServerT Handler)
rateLimitsServer :: forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry -> RateLimitsAPI (AsServerT Handler)
rateLimitsServer ArbiterServerConfig registry
config =
  RateLimitsAPI
    { listRateLimits :: AsServerT Handler :- Get '[JSON] RateLimitPoliciesResponse
listRateLimits = ArbiterServerConfig registry -> Handler RateLimitPoliciesResponse
forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry -> Handler RateLimitPoliciesResponse
listRateLimitsHandler ArbiterServerConfig registry
config
    , listRateLimitBuckets :: AsServerT Handler
:- (Capture "prefix" Text
    :> ("buckets"
        :> (QueryParam "limit" Int
            :> (QueryParam "offset" Int
                :> Get '[JSON] RateLimitBucketsResponse))))
listRateLimitBuckets = ArbiterServerConfig registry
-> Text
-> Maybe Int
-> Maybe Int
-> Handler RateLimitBucketsResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Text
-> Maybe Int
-> Maybe Int
-> Handler RateLimitBucketsResponse
listRateLimitBucketsHandler ArbiterServerConfig registry
config
    , updateRateLimitPolicy :: AsServerT Handler
:- (Capture "prefix" Text
    :> (ReqBody '[JSON] RateLimitPolicyUpdate
        :> Patch '[JSON] RateLimitPolicyView))
updateRateLimitPolicy = ArbiterServerConfig registry
-> Text -> RateLimitPolicyUpdate -> Handler RateLimitPolicyView
forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry
-> Text -> RateLimitPolicyUpdate -> Handler RateLimitPolicyView
updateRateLimitPolicyHandler ArbiterServerConfig registry
config
    , resetRateLimitBuckets :: AsServerT Handler
:- (Capture "prefix" Text
    :> ("reset" :> Post '[JSON] RateLimitResetResponse))
resetRateLimitBuckets = ArbiterServerConfig registry
-> Text -> Handler RateLimitResetResponse
forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry
-> Text -> Handler RateLimitResetResponse
resetRateLimitBucketsHandler ArbiterServerConfig registry
config
    }

-- | Liveness and readiness handlers.
healthServer
  :: forall registry
   . ArbiterServerConfig registry
  -> HealthAPI (AsServerT Handler)
healthServer :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> HealthAPI (AsServerT Handler)
healthServer ArbiterServerConfig registry
config =
  HealthAPI
    { getHealth :: AsServerT Handler :- Get '[JSON] HealthResponse
getHealth = ArbiterServerConfig registry -> Handler HealthResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Handler HealthResponse
healthHandler ArbiterServerConfig registry
config
    , getLiveness :: AsServerT Handler :- ("live" :> Get '[JSON] LivenessResponse)
getLiveness = LivenessResponse -> Handler LivenessResponse
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure LivenessResponse {alive :: Bool
alive = Bool
True}
    }

-- | Readiness for probes and the dashboard. An unreachable database is a 503.
-- Both answers carry the same body.
healthHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Handler HealthResponse
healthHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Handler HealthResponse
healthHandler ArbiterServerConfig registry
config = do
  report <- IO HealthResponse -> Handler HealthResponse
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (ArbiterServerConfig registry -> IO HealthResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> IO HealthResponse
probeHealth ArbiterServerConfig registry
config)
  case status report of
    HealthStatus
Ok -> HealthResponse -> Handler HealthResponse
forall a. a -> Handler a
forall (f :: * -> *) a. Applicative f => a -> f a
pure HealthResponse
report
    HealthStatus
Down ->
      ServerError -> Handler HealthResponse
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError
        ServerError
err503
          { errBody = encode report
          , errHeaders = [("Content-Type", "application/json;charset=utf-8")]
          }

-- | Time a database round-trip and report what it says about itself. Cancellation
-- still propagates.
probeHealth
  :: forall registry
   . ArbiterServerConfig registry
  -> IO HealthResponse
probeHealth :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> IO HealthResponse
probeHealth ArbiterServerConfig registry
config = NominalDiffTime
-> CacheCell HealthResponse
-> IO HealthResponse
-> IO HealthResponse
forall a. NominalDiffTime -> CacheCell a -> IO a -> IO a
cachedFor NominalDiffTime
healthCacheTtl (ArbiterServerConfig registry -> CacheCell HealthResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell HealthResponse
healthCache ArbiterServerConfig registry
config) (IO HealthResponse -> IO HealthResponse)
-> IO HealthResponse -> IO HealthResponse
forall a b. (a -> b) -> a -> b
$ do
  let schemaName :: Text
schemaName = ArbiterServerConfig registry -> Text
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Text
serverSchema ArbiterServerConfig registry
config
  started <- IO UTCTime
getCurrentTime
  probed <- try @SomeException $ timeout healthProbeMicros (runDb config Health.getPgDbHealth)
  either swallowSync (const (pure ())) probed
  finished <- getCurrentTime
  let elapsedMs = NominalDiffTime -> Double
forall a b. (Real a, Fractional b) => a -> b
realToFrac (UTCTime -> UTCTime -> NominalDiffTime
diffUTCTime UTCTime
finished UTCTime
started) Double -> Double -> Double
forall a. Num a => a -> a -> a
* Double
1000
  let reached = Maybe (Maybe (Maybe PgDbHealth)) -> Maybe (Maybe PgDbHealth)
forall (m :: * -> *) a. Monad m => m (m a) -> m a
join (Either SomeException (Maybe (Maybe PgDbHealth))
-> Maybe (Maybe (Maybe PgDbHealth))
forall {a} {a}. Either a a -> Maybe a
eitherToMaybe Either SomeException (Maybe (Maybe PgDbHealth))
probed)
  pure
    HealthResponse
      { status = maybe Down (const Ok) reached
      , schemaName = schemaName
      , checkedAt = finished
      , dbLatencyMs = elapsedMs <$ reached
      , db = join reached
      }
  where
    eitherToMaybe :: Either a a -> Maybe a
eitherToMaybe = (a -> Maybe a) -> (a -> Maybe a) -> Either a a -> Maybe a
forall a c b. (a -> c) -> (b -> c) -> Either a b -> c
either (Maybe a -> a -> Maybe a
forall a b. a -> b -> a
const Maybe a
forall a. Maybe a
Nothing) a -> Maybe a
forall a. a -> Maybe a
Just

-- | Poll-collapsing TTL for the readiness probe.
healthCacheTtl :: NominalDiffTime
healthCacheTtl :: NominalDiffTime
healthCacheTtl = NominalDiffTime
2

-- | Probe time limit. A pool with no available connection reports @down@ when it
-- lapses. A connect blocked in the driver is not interruptible. The connection
-- string needs its own @connect_timeout@.
healthProbeMicros :: Int
healthProbeMicros :: Int
healthProbeMicros = Int
5_000_000

-- | Poll-collapsing TTL for the dashboard list-policy stats.
policyStatsCacheTtl :: NominalDiffTime
policyStatsCacheTtl :: NominalDiffTime
policyStatsCacheTtl = NominalDiffTime
10

-- | Shorter TTL for the faster-polling all-queues overview.
overviewStatsCacheTtl :: NominalDiffTime
overviewStatsCacheTtl :: NominalDiffTime
overviewStatsCacheTtl = NominalDiffTime
5

-- | Default floor between per-queue stats scans.
defaultQueueStatsCacheTtl :: NominalDiffTime
defaultQueueStatsCacheTtl :: NominalDiffTime
defaultQueueStatsCacheTtl = NominalDiffTime
2

-- | No minimum gap. An explicit maintenance call runs every operation.
-- Concurrent callers exclude each other on the gate.
defaultMaintenanceInterval :: NominalDiffTime
defaultMaintenanceInterval :: NominalDiffTime
defaultMaintenanceInterval = NominalDiffTime
0

-- | Gap the whole-schema operations keep, matching a worker pool's reaper.
defaultMaintenanceSparseInterval :: NominalDiffTime
defaultMaintenanceSparseInterval :: NominalDiffTime
defaultMaintenanceSparseInterval = NominalDiffTime
3600

-- | Bucket idle age, matching a worker pool's reaper.
defaultMaintenanceBucketIdle :: NominalDiffTime
defaultMaintenanceBucketIdle :: NominalDiffTime
defaultMaintenanceBucketIdle = NominalDiffTime
300

-- | Statement timeout for one maintenance operation.
defaultMaintenanceTimeout :: NominalDiffTime
defaultMaintenanceTimeout :: NominalDiffTime
defaultMaintenanceTimeout = NominalDiffTime
300

-- | Keyed TTL cache under an epoch bumped by 'invalidate'.
data CacheCell a = CacheCell
  { forall a. CacheCell a -> TVar (Word, Map Text (UTCTime, a))
cacheEntries :: TVar (Word, Map.Map Text (UTCTime, a))
  , forall a. CacheCell a -> TVar (Set Text)
cacheFilling :: TVar (Set.Set Text)
  }

newCacheCell :: IO (CacheCell a)
newCacheCell :: forall a. IO (CacheCell a)
newCacheCell = TVar (Word, Map Text (UTCTime, a))
-> TVar (Set Text) -> CacheCell a
forall a.
TVar (Word, Map Text (UTCTime, a))
-> TVar (Set Text) -> CacheCell a
CacheCell (TVar (Word, Map Text (UTCTime, a))
 -> TVar (Set Text) -> CacheCell a)
-> IO (TVar (Word, Map Text (UTCTime, a)))
-> IO (TVar (Set Text) -> CacheCell a)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> (Word, Map Text (UTCTime, a))
-> IO (TVar (Word, Map Text (UTCTime, a)))
forall a. a -> IO (TVar a)
newTVarIO (Word
0, Map Text (UTCTime, a)
forall k a. Map k a
Map.empty) IO (TVar (Set Text) -> CacheCell a)
-> IO (TVar (Set Text)) -> IO (CacheCell a)
forall a b. IO (a -> b) -> IO a -> IO b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Set Text -> IO (TVar (Set Text))
forall a. a -> IO (TVar a)
newTVarIO Set Text
forall a. Set a
Set.empty

-- | Serve the sole entry of a single-key cell.
cachedFor :: NominalDiffTime -> CacheCell a -> IO a -> IO a
cachedFor :: forall a. NominalDiffTime -> CacheCell a -> IO a -> IO a
cachedFor NominalDiffTime
ttl CacheCell a
cell = NominalDiffTime -> CacheCell a -> Text -> IO a -> IO a
forall a. NominalDiffTime -> CacheCell a -> Text -> IO a -> IO a
cachedForKey NominalDiffTime
ttl CacheCell a
cell Text
""

-- | Serve one key. Concurrent misses on a key collapse onto one @produce@. Its
-- write is skipped when 'invalidate' bumped the epoch meanwhile.
cachedForKey :: NominalDiffTime -> CacheCell a -> Text -> IO a -> IO a
cachedForKey :: forall a. NominalDiffTime -> CacheCell a -> Text -> IO a -> IO a
cachedForKey NominalDiffTime
ttl CacheCell a
cell Text
key IO a
produce
  | NominalDiffTime
ttl NominalDiffTime -> NominalDiffTime -> Bool
forall a. Ord a => a -> a -> Bool
<= NominalDiffTime
0 = IO a
produce
  | Bool
otherwise = IO (Maybe a)
fresh IO (Maybe a) -> (Maybe a -> IO a) -> IO a
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= IO a -> (a -> IO a) -> Maybe a -> IO a
forall b a. b -> (a -> b) -> Maybe a -> b
maybe IO a
fill a -> IO a
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure
  where
    fresh :: IO (Maybe a)
fresh = do
      now <- IO UTCTime
getCurrentTime
      (_, entries) <- readTVarIO (cacheEntries cell)
      pure $ do
        (storedAt, value) <- Map.lookup key entries
        guard (diffUTCTime now storedAt < ttl)
        pure value
    fill :: IO a
fill = IO () -> IO () -> IO a -> IO a
forall a b c. IO a -> IO b -> IO c -> IO c
bracket_ IO ()
acquire IO ()
release (IO (Maybe a)
fresh IO (Maybe a) -> (Maybe a -> IO a) -> IO a
forall a b. IO a -> (a -> IO b) -> IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= IO a -> (a -> IO a) -> Maybe a -> IO a
forall b a. b -> (a -> b) -> Maybe a -> b
maybe IO a
store a -> IO a
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure)
    acquire :: IO ()
acquire = STM () -> IO ()
forall a. STM a -> IO a
atomically (STM () -> IO ()) -> STM () -> IO ()
forall a b. (a -> b) -> a -> b
$ do
      inflight <- TVar (Set Text) -> STM (Set Text)
forall a. TVar a -> STM a
readTVar (CacheCell a -> TVar (Set Text)
forall a. CacheCell a -> TVar (Set Text)
cacheFilling CacheCell a
cell)
      check (Set.notMember key inflight)
      modifyTVar' (cacheFilling cell) (Set.insert key)
    release :: IO ()
release = STM () -> IO ()
forall a. STM a -> IO a
atomically (STM () -> IO ()) -> STM () -> IO ()
forall a b. (a -> b) -> a -> b
$ TVar (Set Text) -> (Set Text -> Set Text) -> STM ()
forall a. TVar a -> (a -> a) -> STM ()
modifyTVar' (CacheCell a -> TVar (Set Text)
forall a. CacheCell a -> TVar (Set Text)
cacheFilling CacheCell a
cell) (Text -> Set Text -> Set Text
forall a. Ord a => a -> Set a -> Set a
Set.delete Text
key)
    store :: IO a
store = do
      (epoch, _) <- TVar (Word, Map Text (UTCTime, a))
-> IO (Word, Map Text (UTCTime, a))
forall a. TVar a -> IO a
readTVarIO (CacheCell a -> TVar (Word, Map Text (UTCTime, a))
forall a. CacheCell a -> TVar (Word, Map Text (UTCTime, a))
cacheEntries CacheCell a
cell)
      value <- produce
      now <- getCurrentTime
      atomically $ modifyTVar' (cacheEntries cell) $ \(Word
current, Map Text (UTCTime, a)
cached) ->
        if Word
current Word -> Word -> Bool
forall a. Eq a => a -> a -> Bool
== Word
epoch then (Word
current, Text
-> (UTCTime, a) -> Map Text (UTCTime, a) -> Map Text (UTCTime, a)
forall k a. Ord k => k -> a -> Map k a -> Map k a
Map.insert Text
key (UTCTime
now, a
value) Map Text (UTCTime, a)
cached) else (Word
current, Map Text (UTCTime, a)
cached)
      pure value

-- | Bump a cache cell's epoch and drop its entries, after an operator mutation.
invalidate :: CacheCell a -> Handler ()
invalidate :: forall a. CacheCell a -> Handler ()
invalidate CacheCell a
cell = IO () -> Handler ()
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> Handler ()) -> IO () -> Handler ()
forall a b. (a -> b) -> a -> b
$ STM () -> IO ()
forall a. STM a -> IO a
atomically (STM () -> IO ()) -> STM () -> IO ()
forall a b. (a -> b) -> a -> b
$ TVar (Word, Map Text (UTCTime, a))
-> ((Word, Map Text (UTCTime, a)) -> (Word, Map Text (UTCTime, a)))
-> STM ()
forall a. TVar a -> (a -> a) -> STM ()
modifyTVar' (CacheCell a -> TVar (Word, Map Text (UTCTime, a))
forall a. CacheCell a -> TVar (Word, Map Text (UTCTime, a))
cacheEntries CacheCell a
cell) (((Word, Map Text (UTCTime, a)) -> (Word, Map Text (UTCTime, a)))
 -> STM ())
-> ((Word, Map Text (UTCTime, a)) -> (Word, Map Text (UTCTime, a)))
-> STM ()
forall a b. (a -> b) -> a -> b
$ \(Word
epoch, Map Text (UTCTime, a)
_) -> (Word
epoch Word -> Word -> Word
forall a. Num a => a -> a -> a
+ Word
1, Map Text (UTCTime, a)
forall k a. Map k a
Map.empty)

-- | List policies with bucket stats and currently-throttled job counts.
listRateLimitsHandler
  :: forall registry
   . (RegistryTables registry)
  => ArbiterServerConfig registry
  -> Handler RateLimitPoliciesResponse
listRateLimitsHandler :: forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry -> Handler RateLimitPoliciesResponse
listRateLimitsHandler ArbiterServerConfig registry
config =
  IO RateLimitPoliciesResponse -> Handler RateLimitPoliciesResponse
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO RateLimitPoliciesResponse -> Handler RateLimitPoliciesResponse)
-> IO RateLimitPoliciesResponse
-> Handler RateLimitPoliciesResponse
forall a b. (a -> b) -> a -> b
$ NominalDiffTime
-> CacheCell RateLimitPoliciesResponse
-> IO RateLimitPoliciesResponse
-> IO RateLimitPoliciesResponse
forall a. NominalDiffTime -> CacheCell a -> IO a -> IO a
cachedFor NominalDiffTime
policyStatsCacheTtl (ArbiterServerConfig registry -> CacheCell RateLimitPoliciesResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CacheCell RateLimitPoliciesResponse
rateLimitPoliciesCache ArbiterServerConfig registry
config) (IO RateLimitPoliciesResponse -> IO RateLimitPoliciesResponse)
-> IO RateLimitPoliciesResponse -> IO RateLimitPoliciesResponse
forall a b. (a -> b) -> a -> b
$ do
    views <- ArbiterServerConfig registry
-> SimpleDb registry IO [RateLimitPolicyView]
-> IO [RateLimitPolicyView]
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config SimpleDb registry IO [RateLimitPolicyView]
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
m [RateLimitPolicyView]
HL.listRateLimitPolicies
    pure $ RateLimitPoliciesResponse {policies = views}

-- | List a prefix's buckets with fill levels, paginated (default 100, max 1000).
listRateLimitBucketsHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Text
  -> Maybe Int
  -> Maybe Int
  -> Handler RateLimitBucketsResponse
listRateLimitBucketsHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Text
-> Maybe Int
-> Maybe Int
-> Handler RateLimitBucketsResponse
listRateLimitBucketsHandler ArbiterServerConfig registry
config Text
prefix Maybe Int
mLimit Maybe Int
mOffset = do
  let (Int
limit, Int
offset) = Int -> Maybe Int -> Maybe Int -> (Int, Int)
validatePagination Int
100 Maybe Int
mLimit Maybe Int
mOffset
  rows <- ArbiterServerConfig registry
-> SimpleDb registry IO [RateLimitBucketView]
-> Handler [RateLimitBucketView]
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (Text -> Int -> Int -> SimpleDb registry IO [RateLimitBucketView]
forall (m :: * -> *).
MonadArbiter m =>
Text -> Int -> Int -> m [RateLimitBucketView]
HL.listRateLimitBuckets Text
prefix Int
limit Int
offset)
  pure $ RateLimitBucketsResponse {buckets = rows}

updateThenView
  :: ArbiterServerConfig registry
  -> SimpleDb registry IO (Maybe a)
  -> LBS.ByteString
  -> Handler a
updateThenView :: forall (registry :: JobPayloadRegistry) a.
ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe a) -> ByteString -> Handler a
updateThenView ArbiterServerConfig registry
config SimpleDb registry IO (Maybe a)
action ByteString
notFound = do
  mView <- ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe a) -> Handler (Maybe a)
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config SimpleDb registry IO (Maybe a)
action
  maybe (throwError err404 {errBody = notFound}) pure mView

-- | Set or clear a policy's override params, then return the updated view.
updateRateLimitPolicyHandler
  :: forall registry
   . (RegistryTables registry)
  => ArbiterServerConfig registry
  -> Text
  -> RateLimitPolicyUpdate
  -> Handler RateLimitPolicyView
updateRateLimitPolicyHandler :: forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry
-> Text -> RateLimitPolicyUpdate -> Handler RateLimitPolicyView
updateRateLimitPolicyHandler ArbiterServerConfig registry
config Text
prefix upd :: RateLimitPolicyUpdate
upd@(RateLimitPolicyUpdate Maybe (Maybe Double)
mMax Maybe (Maybe Double)
mRefill Maybe (Maybe Double)
mInterval) = do
  let invalid :: Maybe ByteString
invalid
        | Bool -> (Double -> Bool) -> Maybe Double -> Bool
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Bool
False (Double -> Double -> Bool
forall a. Ord a => a -> a -> Bool
< Double
0) (Maybe (Maybe Double) -> Maybe Double
forall (m :: * -> *) a. Monad m => m (m a) -> m a
join Maybe (Maybe Double)
mMax) = ByteString -> Maybe ByteString
forall a. a -> Maybe a
Just ByteString
"override max tokens must be >= 0"
        | Bool -> (Double -> Bool) -> Maybe Double -> Bool
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Bool
False (Double -> Double -> Bool
forall a. Ord a => a -> a -> Bool
< Double
0) (Maybe (Maybe Double) -> Maybe Double
forall (m :: * -> *) a. Monad m => m (m a) -> m a
join Maybe (Maybe Double)
mRefill) = ByteString -> Maybe ByteString
forall a. a -> Maybe a
Just ByteString
"override refill amount must be >= 0"
        | Bool -> (Double -> Bool) -> Maybe Double -> Bool
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Bool
False (Double -> Double -> Bool
forall a. Ord a => a -> a -> Bool
<= Double
0) (Maybe (Maybe Double) -> Maybe Double
forall (m :: * -> *) a. Monad m => m (m a) -> m a
join Maybe (Maybe Double)
mInterval) = ByteString -> Maybe ByteString
forall a. a -> Maybe a
Just ByteString
"override interval must be > 0"
        | Bool
otherwise = Maybe ByteString
forall a. Maybe a
Nothing
  case Maybe ByteString
invalid of
    Just ByteString
msg -> ServerError -> Handler RateLimitPolicyView
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err400 {errBody = msg}
    Maybe ByteString
Nothing -> do
      -- An all-absent patch reads the view without rewriting the row.
      let update :: SimpleDb registry IO ()
update = case (Maybe (Maybe Double)
mMax, Maybe (Maybe Double)
mRefill, Maybe (Maybe Double)
mInterval) of
            (Maybe (Maybe Double)
Nothing, Maybe (Maybe Double)
Nothing, Maybe (Maybe Double)
Nothing) -> () -> SimpleDb registry IO ()
forall a. a -> SimpleDb registry IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
            (Maybe (Maybe Double), Maybe (Maybe Double), Maybe (Maybe Double))
_ -> SimpleDb registry IO Int64 -> SimpleDb registry IO ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (SimpleDb registry IO Int64 -> SimpleDb registry IO ())
-> SimpleDb registry IO Int64 -> SimpleDb registry IO ()
forall a b. (a -> b) -> a -> b
$ Text -> RateLimitPolicyUpdate -> SimpleDb registry IO Int64
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Text -> RateLimitPolicyUpdate -> m Int64
HL.updateRateLimitPolicyOverrides Text
prefix RateLimitPolicyUpdate
upd
      view <- ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe RateLimitPolicyView)
-> ByteString
-> Handler RateLimitPolicyView
forall (registry :: JobPayloadRegistry) a.
ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe a) -> ByteString -> Handler a
updateThenView ArbiterServerConfig registry
config (SimpleDb registry IO ()
update SimpleDb registry IO ()
-> SimpleDb registry IO (Maybe RateLimitPolicyView)
-> SimpleDb registry IO (Maybe RateLimitPolicyView)
forall a b.
SimpleDb registry IO a
-> SimpleDb registry IO b -> SimpleDb registry IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> Text -> SimpleDb registry IO (Maybe RateLimitPolicyView)
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Text -> m (Maybe RateLimitPolicyView)
HL.getRateLimitPolicy Text
prefix) ByteString
"Rate-limit policy not found"
      invalidate (rateLimitPoliciesCache config)
      pure view

-- | Clear every bucket for a prefix. Returns the number reset. 404s an unknown prefix.
resetRateLimitBucketsHandler
  :: forall registry
   . (RegistryTables registry)
  => ArbiterServerConfig registry
  -> Text
  -> Handler RateLimitResetResponse
resetRateLimitBucketsHandler :: forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry
-> Text -> Handler RateLimitResetResponse
resetRateLimitBucketsHandler ArbiterServerConfig registry
config Text
prefix = do
  let action :: SimpleDb registry IO (Maybe Int64)
action =
        Text -> SimpleDb registry IO Bool
forall (m :: * -> *). MonadArbiter m => Text -> m Bool
HL.rateLimitPolicyExists Text
prefix SimpleDb registry IO Bool
-> (Bool -> SimpleDb registry IO (Maybe Int64))
-> SimpleDb registry IO (Maybe Int64)
forall a b.
SimpleDb registry IO a
-> (a -> SimpleDb registry IO b) -> SimpleDb registry IO b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Bool
exists -> if Bool
exists then Int64 -> Maybe Int64
forall a. a -> Maybe a
Just (Int64 -> Maybe Int64)
-> SimpleDb registry IO Int64 -> SimpleDb registry IO (Maybe Int64)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Text -> SimpleDb registry IO Int64
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Text -> m Int64
HL.resetRateLimitBuckets Text
prefix else Maybe Int64 -> SimpleDb registry IO (Maybe Int64)
forall a. a -> SimpleDb registry IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe Int64
forall a. Maybe a
Nothing
  count <- ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe Int64)
-> ByteString
-> Handler Int64
forall (registry :: JobPayloadRegistry) a.
ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe a) -> ByteString -> Handler a
updateThenView ArbiterServerConfig registry
config SimpleDb registry IO (Maybe Int64)
action ByteString
"Rate-limit policy not found"
  invalidate (rateLimitPoliciesCache config)
  pure $ RateLimitResetResponse {reset = count}

-- | Concurrency management/observability handlers.
concurrencyServer
  :: forall registry
   . (RegistryTables registry)
  => ArbiterServerConfig registry
  -> ConcurrencyAPI (AsServerT Handler)
concurrencyServer :: forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry -> ConcurrencyAPI (AsServerT Handler)
concurrencyServer ArbiterServerConfig registry
config =
  ConcurrencyAPI
    { listConcurrency :: AsServerT Handler :- Get '[JSON] ConcurrencyPoliciesResponse
listConcurrency = ArbiterServerConfig registry -> Handler ConcurrencyPoliciesResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Handler ConcurrencyPoliciesResponse
listConcurrencyHandler ArbiterServerConfig registry
config
    , listConcurrencyKeys :: AsServerT Handler
:- (Capture "prefix" Text
    :> ("keys"
        :> (QueryParam "limit" Int
            :> (QueryParam "offset" Int
                :> Get '[JSON] ConcurrencyKeysResponse))))
listConcurrencyKeys = ArbiterServerConfig registry
-> Text
-> Maybe Int
-> Maybe Int
-> Handler ConcurrencyKeysResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Text
-> Maybe Int
-> Maybe Int
-> Handler ConcurrencyKeysResponse
listConcurrencyKeysHandler ArbiterServerConfig registry
config
    , updateConcurrencyPolicy :: AsServerT Handler
:- (Capture "prefix" Text
    :> (ReqBody '[JSON] ConcurrencyPolicyUpdate
        :> Patch '[JSON] ConcurrencyPolicyView))
updateConcurrencyPolicy = ArbiterServerConfig registry
-> Text -> ConcurrencyPolicyUpdate -> Handler ConcurrencyPolicyView
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Text -> ConcurrencyPolicyUpdate -> Handler ConcurrencyPolicyView
updateConcurrencyPolicyHandler ArbiterServerConfig registry
config
    , reconcileConcurrency :: AsServerT Handler
:- ("reconcile" :> Post '[JSON] ConcurrencyReconcileResponse)
reconcileConcurrency = ArbiterServerConfig registry
-> Handler ConcurrencyReconcileResponse
forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry
-> Handler ConcurrencyReconcileResponse
reconcileConcurrencyHandler ArbiterServerConfig registry
config
    }

-- | List pools with their default/override limit and live key/in-flight stats.
listConcurrencyHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Handler ConcurrencyPoliciesResponse
listConcurrencyHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Handler ConcurrencyPoliciesResponse
listConcurrencyHandler ArbiterServerConfig registry
config =
  IO ConcurrencyPoliciesResponse
-> Handler ConcurrencyPoliciesResponse
forall a. IO a -> Handler a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO ConcurrencyPoliciesResponse
 -> Handler ConcurrencyPoliciesResponse)
-> IO ConcurrencyPoliciesResponse
-> Handler ConcurrencyPoliciesResponse
forall a b. (a -> b) -> a -> b
$ NominalDiffTime
-> CacheCell ConcurrencyPoliciesResponse
-> IO ConcurrencyPoliciesResponse
-> IO ConcurrencyPoliciesResponse
forall a. NominalDiffTime -> CacheCell a -> IO a -> IO a
cachedFor NominalDiffTime
policyStatsCacheTtl (ArbiterServerConfig registry
-> CacheCell ConcurrencyPoliciesResponse
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> CacheCell ConcurrencyPoliciesResponse
concurrencyPoliciesCache ArbiterServerConfig registry
config) (IO ConcurrencyPoliciesResponse -> IO ConcurrencyPoliciesResponse)
-> IO ConcurrencyPoliciesResponse -> IO ConcurrencyPoliciesResponse
forall a b. (a -> b) -> a -> b
$ do
    views <- ArbiterServerConfig registry
-> SimpleDb registry IO [ConcurrencyPolicyView]
-> IO [ConcurrencyPolicyView]
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config SimpleDb registry IO [ConcurrencyPolicyView]
forall (m :: * -> *). MonadArbiter m => m [ConcurrencyPolicyView]
HL.listConcurrencyPolicies
    pure $ ConcurrencyPoliciesResponse {policies = views}

-- | List a prefix's keys with in-flight fill levels, paginated (default 100, max 1000).
listConcurrencyKeysHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Text
  -> Maybe Int
  -> Maybe Int
  -> Handler ConcurrencyKeysResponse
listConcurrencyKeysHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Text
-> Maybe Int
-> Maybe Int
-> Handler ConcurrencyKeysResponse
listConcurrencyKeysHandler ArbiterServerConfig registry
config Text
prefix Maybe Int
mLimit Maybe Int
mOffset = do
  let (Int
limit, Int
offset) = Int -> Maybe Int -> Maybe Int -> (Int, Int)
validatePagination Int
100 Maybe Int
mLimit Maybe Int
mOffset
  rows <- ArbiterServerConfig registry
-> SimpleDb registry IO [ConcurrencyKeyView]
-> Handler [ConcurrencyKeyView]
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config (Text -> Int -> Int -> SimpleDb registry IO [ConcurrencyKeyView]
forall (m :: * -> *).
MonadArbiter m =>
Text -> Int -> Int -> m [ConcurrencyKeyView]
HL.listConcurrencyKeys Text
prefix Int
limit Int
offset)
  pure $ ConcurrencyKeysResponse {keys = rows}

-- | Set or clear a pool's override limit, then return the updated view.
updateConcurrencyPolicyHandler
  :: forall registry
   . ArbiterServerConfig registry
  -> Text
  -> ConcurrencyPolicyUpdate
  -> Handler ConcurrencyPolicyView
updateConcurrencyPolicyHandler :: forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry
-> Text -> ConcurrencyPolicyUpdate -> Handler ConcurrencyPolicyView
updateConcurrencyPolicyHandler ArbiterServerConfig registry
config Text
prefix upd :: ConcurrencyPolicyUpdate
upd@(ConcurrencyPolicyUpdate Maybe (Maybe Int32)
mLimit) = do
  let invalid :: Maybe ByteString
invalid
        | Bool -> (Int32 -> Bool) -> Maybe Int32 -> Bool
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Bool
False (Int32 -> Int32 -> Bool
forall a. Ord a => a -> a -> Bool
< Int32
0) (Maybe (Maybe Int32) -> Maybe Int32
forall (m :: * -> *) a. Monad m => m (m a) -> m a
join Maybe (Maybe Int32)
mLimit) = ByteString -> Maybe ByteString
forall a. a -> Maybe a
Just ByteString
"override limit must be >= 0"
        | Bool
otherwise = Maybe ByteString
forall a. Maybe a
Nothing
  case Maybe ByteString
invalid of
    Just ByteString
msg -> ServerError -> Handler ConcurrencyPolicyView
forall a. ServerError -> Handler a
forall e (m :: * -> *) a. MonadError e m => e -> m a
throwError ServerError
err400 {errBody = msg}
    Maybe ByteString
Nothing -> do
      -- An absent overrideLimit reads the view without rewriting the row.
      let action :: SimpleDb registry IO (Maybe ConcurrencyPolicyView)
action = case Maybe (Maybe Int32)
mLimit of
            Maybe (Maybe Int32)
Nothing -> Text -> SimpleDb registry IO (Maybe ConcurrencyPolicyView)
forall (m :: * -> *).
MonadArbiter m =>
Text -> m (Maybe ConcurrencyPolicyView)
HL.getConcurrencyPolicy Text
prefix
            Just Maybe Int32
_ -> Text -> ConcurrencyPolicyUpdate -> SimpleDb registry IO Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> ConcurrencyPolicyUpdate -> m Int64
HL.updateConcurrencyPolicyOverrides Text
prefix ConcurrencyPolicyUpdate
upd SimpleDb registry IO Int64
-> SimpleDb registry IO (Maybe ConcurrencyPolicyView)
-> SimpleDb registry IO (Maybe ConcurrencyPolicyView)
forall a b.
SimpleDb registry IO a
-> SimpleDb registry IO b -> SimpleDb registry IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> Text -> SimpleDb registry IO (Maybe ConcurrencyPolicyView)
forall (m :: * -> *).
MonadArbiter m =>
Text -> m (Maybe ConcurrencyPolicyView)
HL.getConcurrencyPolicy Text
prefix
      view <- ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe ConcurrencyPolicyView)
-> ByteString
-> Handler ConcurrencyPolicyView
forall (registry :: JobPayloadRegistry) a.
ArbiterServerConfig registry
-> SimpleDb registry IO (Maybe a) -> ByteString -> Handler a
updateThenView ArbiterServerConfig registry
config SimpleDb registry IO (Maybe ConcurrencyPolicyView)
action ByteString
"Concurrency pool not found"
      invalidate (concurrencyPoliciesCache config)
      pure view

-- | Recompute every key's in-flight count from live jobs. Returns rows repaired.
reconcileConcurrencyHandler
  :: forall registry
   . (RegistryTables registry)
  => ArbiterServerConfig registry
  -> Handler ConcurrencyReconcileResponse
reconcileConcurrencyHandler :: forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry
-> Handler ConcurrencyReconcileResponse
reconcileConcurrencyHandler ArbiterServerConfig registry
config = do
  repaired <- ArbiterServerConfig registry
-> SimpleDb registry IO Int64 -> Handler Int64
forall (n :: * -> *) (registry :: JobPayloadRegistry) a.
MonadIO n =>
ArbiterServerConfig registry -> SimpleDb registry IO a -> n a
runDb ArbiterServerConfig registry
config SimpleDb registry IO Int64
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
m Int64
HL.reconcileConcurrencyCounts
  invalidate (concurrencyPoliciesCache config)
  pure $ ConcurrencyReconcileResponse {reconciled = repaired}

-- | Server for the shared top-level routes.
sharedServer
  :: forall registry
   . (HL.RegistryAdmissionPolicies registry, RegistryTables registry)
  => ArbiterServerConfig registry
  -> ServerT SharedAPI Handler
sharedServer :: forall (registry :: JobPayloadRegistry).
(RegistryAdmissionPolicies registry, RegistryTables registry) =>
ArbiterServerConfig registry -> ServerT SharedAPI Handler
sharedServer ArbiterServerConfig registry
config =
  forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
Proxy registry
-> ArbiterServerConfig registry -> QueuesAPI (AsServerT Handler)
queuesServer @registry (forall (t :: JobPayloadRegistry). Proxy t
forall {k} (t :: k). Proxy t
Proxy @registry) ArbiterServerConfig registry
config
    QueuesAPI (AsServerT Handler)
-> (MaintenanceAPI (AsServerT Handler)
    :<|> (Tagged Handler Application
          :<|> (CronAPI (AsServerT Handler)
                :<|> (WorkersAPI (AsServerT Handler)
                      :<|> (RateLimitsAPI (AsServerT Handler)
                            :<|> (ConcurrencyAPI (AsServerT Handler)
                                  :<|> HealthAPI (AsServerT Handler)))))))
-> QueuesAPI (AsServerT Handler)
   :<|> (MaintenanceAPI (AsServerT Handler)
         :<|> (Tagged Handler Application
               :<|> (CronAPI (AsServerT Handler)
                     :<|> (WorkersAPI (AsServerT Handler)
                           :<|> (RateLimitsAPI (AsServerT Handler)
                                 :<|> (ConcurrencyAPI (AsServerT Handler)
                                       :<|> HealthAPI (AsServerT Handler)))))))
forall a b. a -> b -> a :<|> b
:<|> forall (registry :: JobPayloadRegistry).
(RegistryAdmissionPolicies registry, RegistryTables registry) =>
ArbiterServerConfig registry -> MaintenanceAPI (AsServerT Handler)
maintenanceServer @registry ArbiterServerConfig registry
config
    MaintenanceAPI (AsServerT Handler)
-> (Tagged Handler Application
    :<|> (CronAPI (AsServerT Handler)
          :<|> (WorkersAPI (AsServerT Handler)
                :<|> (RateLimitsAPI (AsServerT Handler)
                      :<|> (ConcurrencyAPI (AsServerT Handler)
                            :<|> HealthAPI (AsServerT Handler))))))
-> MaintenanceAPI (AsServerT Handler)
   :<|> (Tagged Handler Application
         :<|> (CronAPI (AsServerT Handler)
               :<|> (WorkersAPI (AsServerT Handler)
                     :<|> (RateLimitsAPI (AsServerT Handler)
                           :<|> (ConcurrencyAPI (AsServerT Handler)
                                 :<|> HealthAPI (AsServerT Handler))))))
forall a b. a -> b -> a :<|> b
:<|> ArbiterServerConfig registry -> Tagged Handler Application
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> Tagged Handler Application
eventsServer ArbiterServerConfig registry
config
    Tagged Handler Application
-> (CronAPI (AsServerT Handler)
    :<|> (WorkersAPI (AsServerT Handler)
          :<|> (RateLimitsAPI (AsServerT Handler)
                :<|> (ConcurrencyAPI (AsServerT Handler)
                      :<|> HealthAPI (AsServerT Handler)))))
-> Tagged Handler Application
   :<|> (CronAPI (AsServerT Handler)
         :<|> (WorkersAPI (AsServerT Handler)
               :<|> (RateLimitsAPI (AsServerT Handler)
                     :<|> (ConcurrencyAPI (AsServerT Handler)
                           :<|> HealthAPI (AsServerT Handler)))))
forall a b. a -> b -> a :<|> b
:<|> ArbiterServerConfig registry -> CronAPI (AsServerT Handler)
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> CronAPI (AsServerT Handler)
cronServer ArbiterServerConfig registry
config
    CronAPI (AsServerT Handler)
-> (WorkersAPI (AsServerT Handler)
    :<|> (RateLimitsAPI (AsServerT Handler)
          :<|> (ConcurrencyAPI (AsServerT Handler)
                :<|> HealthAPI (AsServerT Handler))))
-> CronAPI (AsServerT Handler)
   :<|> (WorkersAPI (AsServerT Handler)
         :<|> (RateLimitsAPI (AsServerT Handler)
               :<|> (ConcurrencyAPI (AsServerT Handler)
                     :<|> HealthAPI (AsServerT Handler))))
forall a b. a -> b -> a :<|> b
:<|> ArbiterServerConfig registry -> WorkersAPI (AsServerT Handler)
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> WorkersAPI (AsServerT Handler)
workersServer ArbiterServerConfig registry
config
    WorkersAPI (AsServerT Handler)
-> (RateLimitsAPI (AsServerT Handler)
    :<|> (ConcurrencyAPI (AsServerT Handler)
          :<|> HealthAPI (AsServerT Handler)))
-> WorkersAPI (AsServerT Handler)
   :<|> (RateLimitsAPI (AsServerT Handler)
         :<|> (ConcurrencyAPI (AsServerT Handler)
               :<|> HealthAPI (AsServerT Handler)))
forall a b. a -> b -> a :<|> b
:<|> ArbiterServerConfig registry -> RateLimitsAPI (AsServerT Handler)
forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry -> RateLimitsAPI (AsServerT Handler)
rateLimitsServer ArbiterServerConfig registry
config
    RateLimitsAPI (AsServerT Handler)
-> (ConcurrencyAPI (AsServerT Handler)
    :<|> HealthAPI (AsServerT Handler))
-> RateLimitsAPI (AsServerT Handler)
   :<|> (ConcurrencyAPI (AsServerT Handler)
         :<|> HealthAPI (AsServerT Handler))
forall a b. a -> b -> a :<|> b
:<|> ArbiterServerConfig registry -> ConcurrencyAPI (AsServerT Handler)
forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
ArbiterServerConfig registry -> ConcurrencyAPI (AsServerT Handler)
concurrencyServer ArbiterServerConfig registry
config
    ConcurrencyAPI (AsServerT Handler)
-> HealthAPI (AsServerT Handler)
-> ConcurrencyAPI (AsServerT Handler)
   :<|> HealthAPI (AsServerT Handler)
forall a b. a -> b -> a :<|> b
:<|> ArbiterServerConfig registry -> HealthAPI (AsServerT Handler)
forall (registry :: JobPayloadRegistry).
ArbiterServerConfig registry -> HealthAPI (AsServerT Handler)
healthServer ArbiterServerConfig registry
config

-- | Builds a registry's per-queue server implementations.
class BuildServer registry (reg :: JobPayloadRegistry) where
  buildServer :: ArbiterServerConfig registry -> ServerT (RegistryToAPI reg) Handler

-- The empty registry builds the shared top-level routes alone.
instance
  (HL.RegistryAdmissionPolicies registry, RegistryTables registry)
  => BuildServer registry '[]
  where
  buildServer :: ArbiterServerConfig registry -> ServerT (RegistryToAPI '[]) Handler
buildServer = ArbiterServerConfig registry -> ServerT SharedAPI Handler
ArbiterServerConfig registry -> ServerT (RegistryToAPI '[]) Handler
forall (registry :: JobPayloadRegistry).
(RegistryAdmissionPolicies registry, RegistryTables registry) =>
ArbiterServerConfig registry -> ServerT SharedAPI Handler
sharedServer

-- One table builds its endpoints, then the rest of the registry.
instance
  ( BuildServer registry rest
  , EncodeJobResult (SpecResult spec)
  , JobPayload (SpecPayload spec)
  , KnownSymbol (SpecName spec)
  )
  => BuildServer registry (spec ': rest)
  where
  buildServer :: ArbiterServerConfig registry
-> ServerT (RegistryToAPI (spec : rest)) Handler
buildServer ArbiterServerConfig registry
config =
    let tableName :: Text
tableName = String -> Text
T.pack (String -> Text) -> String -> Text
forall a b. (a -> b) -> a -> b
$ Proxy (SpecName spec) -> String
forall (n :: Symbol) (proxy :: Symbol -> *).
KnownSymbol n =>
proxy n -> String
symbolVal (forall {k} (t :: k). Proxy t
forall (t :: Symbol). Proxy t
Proxy @(SpecName spec))
     in forall (registry :: JobPayloadRegistry) payload result.
(EncodeJobResult result, JobPayload payload) =>
Text
-> ArbiterServerConfig registry
-> TableAPI payload result (AsServerT Handler)
tableServer @registry @(SpecPayload spec) @(SpecResult spec) Text
tableName ArbiterServerConfig registry
config
          TableAPI (SpecPayload spec) (SpecResult spec) (AsServerT Handler)
-> ServerT (RegistryToAPI rest) Handler
-> TableAPI
     (SpecPayload spec) (SpecResult spec) (AsServerT Handler)
   :<|> ServerT (RegistryToAPI rest) Handler
forall a b. a -> b -> a :<|> b
:<|> forall (registry :: JobPayloadRegistry)
       (reg :: JobPayloadRegistry).
BuildServer registry reg =>
ArbiterServerConfig registry -> ServerT (RegistryToAPI reg) Handler
buildServer @registry @rest ArbiterServerConfig registry
config

-- | Complete Arbiter server at @\/api\/v1\/...@
arbiterServer
  :: forall registry
   . (BuildServer registry registry)
  => ArbiterServerConfig registry
  -> ServerT (ArbiterAPI registry) Handler
arbiterServer :: forall (registry :: JobPayloadRegistry).
BuildServer registry registry =>
ArbiterServerConfig registry
-> ServerT (ArbiterAPI registry) Handler
arbiterServer = forall (registry :: JobPayloadRegistry)
       (reg :: JobPayloadRegistry).
BuildServer registry reg =>
ArbiterServerConfig registry -> ServerT (RegistryToAPI reg) Handler
buildServer @registry @registry

-- | Hoisted server for integration into a route tree using a custom monad.
arbiterServerHoisted
  :: forall registry m
   . ( BuildServer registry registry
     , HasServer (ArbiterAPI registry) '[]
     )
  => (forall x. Handler x -> m x)
  -> ArbiterServerConfig registry
  -> ServerT (ArbiterAPI registry) m
arbiterServerHoisted :: forall (registry :: JobPayloadRegistry) (m :: * -> *).
(BuildServer registry registry,
 HasServer (ArbiterAPI registry) '[]) =>
(forall x. Handler x -> m x)
-> ArbiterServerConfig registry -> ServerT (ArbiterAPI registry) m
arbiterServerHoisted forall x. Handler x -> m x
natTrans ArbiterServerConfig registry
config =
  Proxy (ArbiterAPI registry)
-> (forall x. Handler x -> m x)
-> ServerT (ArbiterAPI registry) Handler
-> ServerT (ArbiterAPI registry) m
forall {k} (api :: k) (m :: * -> *) (n :: * -> *).
HasServer api '[] =>
Proxy api
-> (forall x. m x -> n x) -> ServerT api m -> ServerT api n
hoistServer (forall t. Proxy t
forall {k} (t :: k). Proxy t
Proxy @(ArbiterAPI registry)) Handler x -> m x
forall x. Handler x -> m x
natTrans (ArbiterServerConfig registry
-> ServerT (ArbiterAPI registry) Handler
forall (registry :: JobPayloadRegistry).
BuildServer registry registry =>
ArbiterServerConfig registry
-> ServerT (ArbiterAPI registry) Handler
arbiterServer ArbiterServerConfig registry
config)

-- | Convert to WAI Application. Each 'QueueWithResult' result type needs
-- @FromJSON@ and @ToJSON@.
arbiterApp
  :: forall registry
   . ( BuildServer registry registry
     , HasServer (ArbiterAPI registry) '[]
     )
  => ArbiterServerConfig registry
  -> Application
arbiterApp :: forall (registry :: JobPayloadRegistry).
(BuildServer registry registry,
 HasServer (ArbiterAPI registry) '[]) =>
ArbiterServerConfig registry -> Application
arbiterApp ArbiterServerConfig registry
config =
  Proxy (ArbiterAPI registry)
-> Server (ArbiterAPI registry) -> Application
forall {k} (api :: k).
HasServer api '[] =>
Proxy api -> Server api -> Application
serve (forall t. Proxy t
forall {k} (t :: k). Proxy t
Proxy @(ArbiterAPI registry)) (ArbiterServerConfig registry -> Server (ArbiterAPI registry)
forall (registry :: JobPayloadRegistry).
BuildServer registry registry =>
ArbiterServerConfig registry
-> ServerT (ArbiterAPI registry) Handler
arbiterServer ArbiterServerConfig registry
config)

-- | Run the API server on a port.
runArbiterAPI
  :: forall registry
   . ( BuildServer registry registry
     , HasServer (ArbiterAPI registry) '[]
     )
  => Port
  -> ArbiterServerConfig registry
  -> IO ()
runArbiterAPI :: forall (registry :: JobPayloadRegistry).
(BuildServer registry registry,
 HasServer (ArbiterAPI registry) '[]) =>
Int -> ArbiterServerConfig registry -> IO ()
runArbiterAPI Int
port ArbiterServerConfig registry
config = do
  String -> IO ()
putStrLn (String -> IO ()) -> String -> IO ()
forall a b. (a -> b) -> a -> b
$ String
"Starting Arbiter API server on port " String -> String -> String
forall a. Semigroup a => a -> a -> a
<> Int -> String
forall a. Show a => a -> String
show Int
port
  let settings :: Settings
settings = Int -> Settings -> Settings
setPort Int
port Settings
defaultSettings
  Settings -> Application -> IO ()
runSettings Settings
settings (ArbiterServerConfig registry -> Application
forall (registry :: JobPayloadRegistry).
(BuildServer registry registry,
 HasServer (ArbiterAPI registry) '[]) =>
ArbiterServerConfig registry -> Application
arbiterApp ArbiterServerConfig registry
config)

-- | Remove an empty search parameter.
nonBlank :: Maybe Text -> Maybe Text
nonBlank :: Maybe Text -> Maybe Text
nonBlank = (Text -> Bool) -> Maybe Text -> Maybe Text
forall (m :: * -> *) a. MonadPlus m => (a -> Bool) -> m a -> m a
mfilter (Bool -> Bool
not (Bool -> Bool) -> (Text -> Bool) -> Text -> Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> Bool
T.null (Text -> Bool) -> (Text -> Text) -> Text -> Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> Text
T.strip)

-- | Clamp pagination parameters to a limit of 1 to 1000 and a non-negative offset.
validatePagination :: Int -> Maybe Int -> Maybe Int -> (Int, Int)
validatePagination :: Int -> Maybe Int -> Maybe Int -> (Int, Int)
validatePagination Int
defLimit Maybe Int
mLimit Maybe Int
mOffset =
  let limit :: Int
limit = Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
1 (Int -> Int) -> Int -> Int
forall a b. (a -> b) -> a -> b
$ Int -> Int -> Int
forall a. Ord a => a -> a -> a
min Int
1000 (Int -> Int) -> Int -> Int
forall a b. (a -> b) -> a -> b
$ Int -> Maybe Int -> Int
forall a. a -> Maybe a -> a
fromMaybe Int
defLimit Maybe Int
mLimit
      offset :: Int
offset = Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
0 (Int -> Int) -> Int -> Int
forall a b. (a -> b) -> a -> b
$ Int -> Maybe Int -> Int
forall a. a -> Maybe a -> a
fromMaybe Int
0 Maybe Int
mOffset
   in (Int
limit, Int
offset)