{-# LANGUAGE DeriveAnyClass #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE TypeFamilies #-}

-- | Configuration types for the arbiter worker pool.
module Arbiter.Worker.Config
  ( -- * Worker Configuration
    WorkerConfig (..)
  , transactionalWorkerConfig
  , manualWorkerConfig
  , defaultBatchedWorkerConfig
  , withHooks
  , withMaintenance
  , HandlerMode (..)
  , handlerBatchSize
  , MaintenanceOp (..)
  , maintenanceOpName
  , ResultOf
  , WorkerConfigException (..)
  , validateWorkerConfig

    -- * Batch Callbacks
  , BatchCallbacks (..)

    -- * Worker State
  , WorkerState (..)
  , shutdownWorker
  , getWorkerState
  , getListenerReady
  , readEffectiveState
  , writePause
  , writePauseIfCurrent
  , workerStateVar
  , pauseVar
  , pauseEpoch
  , heartbeatSignal
  , listenerReadyVar
  ) where

import Arbiter.Core.Job.Types (JobRead, ObservabilityHooks, andThen, defaultObservabilityHooks)
import Arbiter.Core.MonadArbiter (JobHandler, MonadArbiter, ResultOf)
import Control.Exception (Exception)
import Control.Monad (when)
import Control.Monad.IO.Class (MonadIO, liftIO)
import Data.Aeson (Value, (.=))
import Data.Foldable (toList)
import Data.Int (Int64)
import Data.List.NonEmpty (NonEmpty ((:|)))
import Data.Text (Text)
import Data.Text qualified as T
import Data.Time (NominalDiffTime)
import Data.UUID (UUID, toString)
import Data.UUID.V4 qualified as UUID
import Data.Word (Word64)
import Network.HostName (getHostName)
import System.Directory (getTemporaryDirectory)
import UnliftIO (MonadUnliftIO)
import UnliftIO.STM (TMVar, TVar, newEmptyTMVarIO, newTVarIO)
import UnliftIO.STM qualified as STM

import Arbiter.Worker.BackoffStrategy (BackoffStrategy, Jitter (..), exponentialBackoff)
import Arbiter.Worker.Cron (CronJob)
import Arbiter.Worker.Logger (LogConfig (..), defaultLogConfig)
import Arbiter.Worker.WorkerState (WorkerState (..))

-- | Which reaper op a maintenance report came from.
data MaintenanceOp
  = RefreshGroups
  | SweepStaleWorkers
  | SweepExhaustedJobs
  | SweepCancelledJobs
  | PruneRateLimitBuckets
  | ReconcileConcurrencyStale
  | ReconcilePruneConcurrency
  | PurgeArchives
  deriving stock (MaintenanceOp
MaintenanceOp -> MaintenanceOp -> Bounded MaintenanceOp
forall a. a -> a -> Bounded a
$cminBound :: MaintenanceOp
minBound :: MaintenanceOp
$cmaxBound :: MaintenanceOp
maxBound :: MaintenanceOp
Bounded, Int -> MaintenanceOp
MaintenanceOp -> Int
MaintenanceOp -> [MaintenanceOp]
MaintenanceOp -> MaintenanceOp
MaintenanceOp -> MaintenanceOp -> [MaintenanceOp]
MaintenanceOp -> MaintenanceOp -> MaintenanceOp -> [MaintenanceOp]
(MaintenanceOp -> MaintenanceOp)
-> (MaintenanceOp -> MaintenanceOp)
-> (Int -> MaintenanceOp)
-> (MaintenanceOp -> Int)
-> (MaintenanceOp -> [MaintenanceOp])
-> (MaintenanceOp -> MaintenanceOp -> [MaintenanceOp])
-> (MaintenanceOp -> MaintenanceOp -> [MaintenanceOp])
-> (MaintenanceOp
    -> MaintenanceOp -> MaintenanceOp -> [MaintenanceOp])
-> Enum MaintenanceOp
forall a.
(a -> a)
-> (a -> a)
-> (Int -> a)
-> (a -> Int)
-> (a -> [a])
-> (a -> a -> [a])
-> (a -> a -> [a])
-> (a -> a -> a -> [a])
-> Enum a
$csucc :: MaintenanceOp -> MaintenanceOp
succ :: MaintenanceOp -> MaintenanceOp
$cpred :: MaintenanceOp -> MaintenanceOp
pred :: MaintenanceOp -> MaintenanceOp
$ctoEnum :: Int -> MaintenanceOp
toEnum :: Int -> MaintenanceOp
$cfromEnum :: MaintenanceOp -> Int
fromEnum :: MaintenanceOp -> Int
$cenumFrom :: MaintenanceOp -> [MaintenanceOp]
enumFrom :: MaintenanceOp -> [MaintenanceOp]
$cenumFromThen :: MaintenanceOp -> MaintenanceOp -> [MaintenanceOp]
enumFromThen :: MaintenanceOp -> MaintenanceOp -> [MaintenanceOp]
$cenumFromTo :: MaintenanceOp -> MaintenanceOp -> [MaintenanceOp]
enumFromTo :: MaintenanceOp -> MaintenanceOp -> [MaintenanceOp]
$cenumFromThenTo :: MaintenanceOp -> MaintenanceOp -> MaintenanceOp -> [MaintenanceOp]
enumFromThenTo :: MaintenanceOp -> MaintenanceOp -> MaintenanceOp -> [MaintenanceOp]
Enum, MaintenanceOp -> MaintenanceOp -> Bool
(MaintenanceOp -> MaintenanceOp -> Bool)
-> (MaintenanceOp -> MaintenanceOp -> Bool) -> Eq MaintenanceOp
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: MaintenanceOp -> MaintenanceOp -> Bool
== :: MaintenanceOp -> MaintenanceOp -> Bool
$c/= :: MaintenanceOp -> MaintenanceOp -> Bool
/= :: MaintenanceOp -> MaintenanceOp -> Bool
Eq, Eq MaintenanceOp
Eq MaintenanceOp =>
(MaintenanceOp -> MaintenanceOp -> Ordering)
-> (MaintenanceOp -> MaintenanceOp -> Bool)
-> (MaintenanceOp -> MaintenanceOp -> Bool)
-> (MaintenanceOp -> MaintenanceOp -> Bool)
-> (MaintenanceOp -> MaintenanceOp -> Bool)
-> (MaintenanceOp -> MaintenanceOp -> MaintenanceOp)
-> (MaintenanceOp -> MaintenanceOp -> MaintenanceOp)
-> Ord MaintenanceOp
MaintenanceOp -> MaintenanceOp -> Bool
MaintenanceOp -> MaintenanceOp -> Ordering
MaintenanceOp -> MaintenanceOp -> MaintenanceOp
forall a.
Eq a =>
(a -> a -> Ordering)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> a)
-> (a -> a -> a)
-> Ord a
$ccompare :: MaintenanceOp -> MaintenanceOp -> Ordering
compare :: MaintenanceOp -> MaintenanceOp -> Ordering
$c< :: MaintenanceOp -> MaintenanceOp -> Bool
< :: MaintenanceOp -> MaintenanceOp -> Bool
$c<= :: MaintenanceOp -> MaintenanceOp -> Bool
<= :: MaintenanceOp -> MaintenanceOp -> Bool
$c> :: MaintenanceOp -> MaintenanceOp -> Bool
> :: MaintenanceOp -> MaintenanceOp -> Bool
$c>= :: MaintenanceOp -> MaintenanceOp -> Bool
>= :: MaintenanceOp -> MaintenanceOp -> Bool
$cmax :: MaintenanceOp -> MaintenanceOp -> MaintenanceOp
max :: MaintenanceOp -> MaintenanceOp -> MaintenanceOp
$cmin :: MaintenanceOp -> MaintenanceOp -> MaintenanceOp
min :: MaintenanceOp -> MaintenanceOp -> MaintenanceOp
Ord, Int -> MaintenanceOp -> ShowS
[MaintenanceOp] -> ShowS
MaintenanceOp -> String
(Int -> MaintenanceOp -> ShowS)
-> (MaintenanceOp -> String)
-> ([MaintenanceOp] -> ShowS)
-> Show MaintenanceOp
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> MaintenanceOp -> ShowS
showsPrec :: Int -> MaintenanceOp -> ShowS
$cshow :: MaintenanceOp -> String
show :: MaintenanceOp -> String
$cshowList :: [MaintenanceOp] -> ShowS
showList :: [MaintenanceOp] -> ShowS
Show)

-- | The op's stable name, used to coordinate replicas and to label its metrics.
maintenanceOpName :: MaintenanceOp -> Text
maintenanceOpName :: MaintenanceOp -> Text
maintenanceOpName MaintenanceOp
operation = case MaintenanceOp
operation of
  MaintenanceOp
RefreshGroups -> Text
"refresh-all-groups"
  MaintenanceOp
SweepStaleWorkers -> Text
"sweep-stale-workers"
  MaintenanceOp
SweepExhaustedJobs -> Text
"sweep-exhausted-jobs"
  MaintenanceOp
SweepCancelledJobs -> Text
"sweep-cancelled-jobs"
  MaintenanceOp
PruneRateLimitBuckets -> Text
"prune-rate-limit-buckets"
  MaintenanceOp
ReconcileConcurrencyStale -> Text
"reconcile-concurrency-stale"
  MaintenanceOp
ReconcilePruneConcurrency -> Text
"reconcile-prune-concurrency"
  MaintenanceOp
PurgeArchives -> Text
"purge-archives"

-- | Invalid worker configuration detected before a pool starts.
newtype WorkerConfigException = WorkerConfigException Text
  deriving stock (WorkerConfigException -> WorkerConfigException -> Bool
(WorkerConfigException -> WorkerConfigException -> Bool)
-> (WorkerConfigException -> WorkerConfigException -> Bool)
-> Eq WorkerConfigException
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: WorkerConfigException -> WorkerConfigException -> Bool
== :: WorkerConfigException -> WorkerConfigException -> Bool
$c/= :: WorkerConfigException -> WorkerConfigException -> Bool
/= :: WorkerConfigException -> WorkerConfigException -> Bool
Eq, Int -> WorkerConfigException -> ShowS
[WorkerConfigException] -> ShowS
WorkerConfigException -> String
(Int -> WorkerConfigException -> ShowS)
-> (WorkerConfigException -> String)
-> ([WorkerConfigException] -> ShowS)
-> Show WorkerConfigException
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> WorkerConfigException -> ShowS
showsPrec :: Int -> WorkerConfigException -> ShowS
$cshow :: WorkerConfigException -> String
show :: WorkerConfigException -> String
$cshowList :: [WorkerConfigException] -> ShowS
showList :: [WorkerConfigException] -> ShowS
Show)
  deriving anyclass (Show WorkerConfigException
Typeable WorkerConfigException
(Typeable WorkerConfigException, Show WorkerConfigException) =>
(WorkerConfigException -> SomeException)
-> (SomeException -> Maybe WorkerConfigException)
-> (WorkerConfigException -> String)
-> (WorkerConfigException -> Bool)
-> Exception WorkerConfigException
SomeException -> Maybe WorkerConfigException
WorkerConfigException -> Bool
WorkerConfigException -> String
WorkerConfigException -> SomeException
forall e.
(Typeable e, Show e) =>
(e -> SomeException)
-> (SomeException -> Maybe e)
-> (e -> String)
-> (e -> Bool)
-> Exception e
$ctoException :: WorkerConfigException -> SomeException
toException :: WorkerConfigException -> SomeException
$cfromException :: SomeException -> Maybe WorkerConfigException
fromException :: SomeException -> Maybe WorkerConfigException
$cdisplayException :: WorkerConfigException -> String
displayException :: WorkerConfigException -> String
$cbacktraceDesired :: WorkerConfigException -> Bool
backtraceDesired :: WorkerConfigException -> Bool
Exception)

-- | Mutable state owned by one worker pool.
data WorkerRuntime = WorkerRuntime
  { WorkerRuntime -> TVar WorkerState
runtimeStateVar :: TVar WorkerState
  , WorkerRuntime -> TVar Bool
runtimePauseVar :: TVar Bool
  , WorkerRuntime -> TVar Word64
runtimePauseEpoch :: TVar Word64
  , WorkerRuntime -> TMVar ()
runtimeHeartbeatSignal :: TMVar ()
  , WorkerRuntime -> TVar Bool
runtimeListenerReadyVar :: TVar Bool
  }

-- | Configuration for a worker pool.
data WorkerConfig m payload = WorkerConfig
  { forall (m :: * -> *) payload. WorkerConfig m payload -> Int
workerCount :: Int
  -- ^ Number of concurrent worker threads.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> HandlerMode m payload
handlerMode :: HandlerMode m payload
  -- ^ Job handler and claiming strategy. Set by this module's config constructors.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
pollInterval :: NominalDiffTime
  -- ^ Cadence floor in seconds for the dispatcher poll.
  -- Default: 5.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
visibilityTimeout :: NominalDiffTime
  -- ^ How long a claimed job stays invisible to other workers.
  -- Must be greater than 'jobHeartbeatInterval'. Default: 60.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
jobHeartbeatInterval :: NominalDiffTime
  -- ^ Interval for extending a job's visibility timeout during processing.
  -- Must be less than 'visibilityTimeout'. Default: 30.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> Maybe NominalDiffTime
maxJobDuration :: Maybe NominalDiffTime
  -- ^ Interrupt a handler that runs longer than this. Default: 'Nothing', no bound.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
workerHeartbeatInterval :: NominalDiffTime
  -- ^ Cadence for bumping @arbiter_workers.last_heartbeat@, the optional
  -- liveness file, and reconciling pause state from the DB. Must be well below
  -- 'workerStaleThreshold'. Default: 10.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> BackoffStrategy
backoffStrategy :: BackoffStrategy
  -- ^ Retry backoff strategy. Default: exponential with base 2, max 1048576 seconds.
  , forall (m :: * -> *) payload. WorkerConfig m payload -> Jitter
jitter :: Jitter
  -- ^ Jitter strategy for retry delays. Default: 'EqualJitter'.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> ObservabilityHooks m payload
observabilityHooks :: ObservabilityHooks m payload
  -- ^ Callbacks for metrics or tracing. Default: no-op hooks.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> MaintenanceOp -> Int64 -> m ()
onMaintenance :: MaintenanceOp -> Int64 -> m ()
  -- ^ Called after a reaper op this pool won the gate for, with the rows it touched.
  -- Reaper work is schema-wide and carries no queue. Default: no-op.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> WorkerRuntime
workerRuntime :: WorkerRuntime
  -- ^ Mutable lifecycle state allocated for this pool.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> Maybe String
livenessFile :: Maybe FilePath
  -- ^ When set, the heartbeat loop touches this file at the
  -- 'workerHeartbeatInterval' cadence, for a file-based liveness probe.
  -- Default: @arbiter-worker-\<workerId\>@ in the system temporary directory.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> Maybe NominalDiffTime
gracefulShutdownTimeout :: Maybe NominalDiffTime
  -- ^ Seconds a graceful shutdown waits on in-flight jobs before force-exiting.
  -- 'Nothing' waits indefinitely. Default: @Just 30@.
  , forall (m :: * -> *) payload. WorkerConfig m payload -> LogConfig
logConfig :: LogConfig
  -- ^ Where the pool's structured JSON logs go, at what level, and what context they
  -- carry beyond the job's own. Default: Info to stdout.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> [CronJob payload]
cronJobs :: [CronJob payload]
  -- ^ Cron schedules. A non-empty list gives the pool a scheduler thread, which reads
  -- the @cron_schedules@ table each tick for runtime overrides. Default: @[]@.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
reaperInterval :: NominalDiffTime
  -- ^ How often the reaper runs. Default: @300@ (5 minutes).
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
reaperSparseInterval :: NominalDiffTime
  -- ^ How often the reaper runs an operation that scans the whole schema. Default: @3600@ (1 hour).
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
reaperBucketIdle :: NominalDiffTime
  -- ^ Idle age at which the reaper prunes a rate-limit bucket. Default: @300@ (5 minutes).
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
reaperTimeout :: NominalDiffTime
  -- ^ Abort any single reaper statement that runs longer than this. Default: @300@ (5 minutes).
  , forall (m :: * -> *) payload. WorkerConfig m payload -> UUID
workerId :: UUID
  -- ^ Identity for this pool. A fresh one is generated by default.
  , forall (m :: * -> *) payload. WorkerConfig m payload -> Maybe Text
workerHost :: Maybe Text
  -- ^ Hostname recorded in the worker registry. Default: auto-generated.
  , forall (m :: * -> *) payload. WorkerConfig m payload -> Maybe Value
workerMetadata :: Maybe Value
  -- ^ Arbitrary JSONB metadata for the worker registry row (image tag,
  -- git SHA, deploy id, etc.). Default: 'Nothing'.
  , forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
workerStaleThreshold :: NominalDiffTime
  -- ^ Workers whose @last_heartbeat@ is older than this are swept from the
  -- runtime registry by 'reaperInterval'. Must be well above the heartbeat cadence
  -- ('workerHeartbeatInterval', or 'jobHeartbeatInterval' while busy).
  -- Default: @300@ (5 minutes).
  }

-- | Per-job finalizers handed to a batched handler. Untouched jobs are
-- reprocessed. The @With@ variants store a result for the job's parent rollup
-- or its archive entry.
--
-- Each callback runs in its own transaction and commits on return. Call them at
-- the top level of the handler. Wrapping one in your own 'Arbiter.Core.MonadArbiter.withDbTransaction'
-- enlists the callback into that transaction as a savepoint, committing atomically
-- with your writes. The success hook then fires at savepoint release. An outer
-- rollback reprocesses the job after the visibility timeout.
data BatchCallbacks m payload result = BatchCallbacks
  { forall (m :: * -> *) payload result.
BatchCallbacks m payload result -> JobRead payload -> m ()
ack :: JobRead payload -> m ()
  -- ^ Ack and fire onJobSuccess, storing no result. Available on any queue. A
  -- job acked this way is absent from its parent rollup's child results and
  -- leaves its archive entry's result @NULL@. Aborts the handler if another
  -- worker holds the job. The abort nacks the siblings still unfinalized.
  , forall (m :: * -> *) payload result.
BatchCallbacks m payload result
-> JobRead payload -> result -> m ()
ackWith :: JobRead payload -> result -> m ()
  -- ^ Ack, store the result for the parent rollup or the job's archive entry,
  -- fire onJobSuccess.
  , forall (m :: * -> *) payload result.
BatchCallbacks m payload result -> [JobRead payload] -> m ()
ackAll :: [JobRead payload] -> m ()
  -- ^ Bulk-'ack' in one parent-aware transaction, storing no results. Fires
  -- onJobSuccess per acked job. A job another worker holds is reported and
  -- skipped. The handler continues.
  , forall (m :: * -> *) payload result.
BatchCallbacks m payload result
-> [(JobRead payload, result)] -> m ()
ackAllWith :: [(JobRead payload, result)] -> m ()
  -- ^ 'ackAll' storing each job's result for its parent rollup or archive entry.
  , forall (m :: * -> *) payload result.
BatchCallbacks m payload result -> JobRead payload -> Text -> m ()
failRetry :: JobRead payload -> Text -> m ()
  -- ^ Retry with backoff, then DLQ at the job's maxAttempts.
  , forall (m :: * -> *) payload result.
BatchCallbacks m payload result -> JobRead payload -> Text -> m ()
failPermanent :: JobRead payload -> Text -> m ()
  -- ^ Straight to the DLQ.
  , forall (m :: * -> *) payload result.
BatchCallbacks m payload result -> JobRead payload -> Text -> m ()
cancelBranch :: JobRead payload -> Text -> m ()
  -- ^ Cancel this job's branch (its parent and all siblings).
  , forall (m :: * -> *) payload result.
BatchCallbacks m payload result -> JobRead payload -> Text -> m ()
cancelTree :: JobRead payload -> Text -> m ()
  -- ^ Cancel the whole tree from the root down.
  , forall (m :: * -> *) payload result.
BatchCallbacks m payload result -> JobRead payload -> m ()
nack :: JobRead payload -> m ()
  -- ^ Reprocess after the visibility timeout. Records no failure and consumes
  -- no attempt.
  }

-- | How the worker claims and runs jobs. Set by this module's config constructors.
data HandlerMode m payload
  = -- | Automatic single-job mode: claim one job per group and run the handler
    -- in a worker transaction, storing its result and acking atomically.
    SingleJobMode (JobHandler m payload (ResultOf m payload))
  | -- | Batched callback mode: claim up to N jobs per group and hand the batch to
    -- the handler with a 'BatchCallbacks' record to finalize each job. No worker
    -- transaction. Batch size 1 is the manual single-job case.
    BatchedJobsMode
      Int
      (NonEmpty (JobRead payload) -> BatchCallbacks m payload (ResultOf m payload) -> m ())

-- | How many jobs a pool claims per group.
handlerBatchSize :: WorkerConfig m payload -> Int
handlerBatchSize :: forall (m :: * -> *) payload. WorkerConfig m payload -> Int
handlerBatchSize WorkerConfig m payload
config = case WorkerConfig m payload -> HandlerMode m payload
forall (m :: * -> *) payload.
WorkerConfig m payload -> HandlerMode m payload
handlerMode WorkerConfig m payload
config of
  SingleJobMode JobHandler m payload (ResultOf m payload)
_ -> Int
1
  BatchedJobsMode Int
batchSize NonEmpty (JobRead payload)
-> BatchCallbacks m payload (ResultOf m payload) -> m ()
_ -> Int
batchSize

-- | Validate all invariants required for safe worker execution. Return all
-- violations.
validateWorkerConfig :: WorkerConfig m payload -> Either Text ()
validateWorkerConfig :: forall (m :: * -> *) payload.
WorkerConfig m payload -> Either Text ()
validateWorkerConfig WorkerConfig m payload
config =
  case WorkerConfig m payload -> Validation (WorkerConfig m payload)
forall (m :: * -> *) payload.
WorkerConfig m payload -> Validation (WorkerConfig m payload)
fieldInvariants WorkerConfig m payload
config Validation (WorkerConfig m payload)
-> Validation () -> Validation ()
forall a b. Validation a -> Validation b -> Validation b
forall (f :: * -> *) a b. Applicative f => f a -> f b -> f b
*> WorkerConfig m payload -> Validation ()
forall (m :: * -> *) payload.
WorkerConfig m payload -> Validation ()
crossFieldInvariants WorkerConfig m payload
config of
    Validation (Left NonEmpty Text
messages) -> Text -> Either Text ()
forall a b. a -> Either a b
Left (Text -> [Text] -> Text
T.intercalate Text
"; " (NonEmpty Text -> [Text]
forall a. NonEmpty a -> [a]
forall (t :: * -> *) a. Foldable t => t a -> [a]
toList NonEmpty Text
messages))
    Validation (Right ()) -> () -> Either Text ()
forall a b. b -> Either a b
Right ()

-- | Rebuilds the config with every field checked or waived. A field added
-- to 'WorkerConfig' is a type error here until it is classified.
fieldInvariants :: WorkerConfig m payload -> Validation (WorkerConfig m payload)
fieldInvariants :: forall (m :: * -> *) payload.
WorkerConfig m payload -> Validation (WorkerConfig m payload)
fieldInvariants WorkerConfig m payload
config =
  Int
-> HandlerMode m payload
-> NominalDiffTime
-> NominalDiffTime
-> NominalDiffTime
-> Maybe NominalDiffTime
-> NominalDiffTime
-> BackoffStrategy
-> Jitter
-> ObservabilityHooks m payload
-> (MaintenanceOp -> Int64 -> m ())
-> WorkerRuntime
-> Maybe String
-> Maybe NominalDiffTime
-> LogConfig
-> [CronJob payload]
-> NominalDiffTime
-> NominalDiffTime
-> NominalDiffTime
-> NominalDiffTime
-> UUID
-> Maybe Text
-> Maybe Value
-> NominalDiffTime
-> WorkerConfig m payload
forall (m :: * -> *) payload.
Int
-> HandlerMode m payload
-> NominalDiffTime
-> NominalDiffTime
-> NominalDiffTime
-> Maybe NominalDiffTime
-> NominalDiffTime
-> BackoffStrategy
-> Jitter
-> ObservabilityHooks m payload
-> (MaintenanceOp -> Int64 -> m ())
-> WorkerRuntime
-> Maybe String
-> Maybe NominalDiffTime
-> LogConfig
-> [CronJob payload]
-> NominalDiffTime
-> NominalDiffTime
-> NominalDiffTime
-> NominalDiffTime
-> UUID
-> Maybe Text
-> Maybe Value
-> NominalDiffTime
-> WorkerConfig m payload
WorkerConfig
    (Int
 -> HandlerMode m payload
 -> NominalDiffTime
 -> NominalDiffTime
 -> NominalDiffTime
 -> Maybe NominalDiffTime
 -> NominalDiffTime
 -> BackoffStrategy
 -> Jitter
 -> ObservabilityHooks m payload
 -> (MaintenanceOp -> Int64 -> m ())
 -> WorkerRuntime
 -> Maybe String
 -> Maybe NominalDiffTime
 -> LogConfig
 -> [CronJob payload]
 -> NominalDiffTime
 -> NominalDiffTime
 -> NominalDiffTime
 -> NominalDiffTime
 -> UUID
 -> Maybe Text
 -> Maybe Value
 -> NominalDiffTime
 -> WorkerConfig m payload)
-> Validation Int
-> Validation
     (HandlerMode m payload
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> Maybe NominalDiffTime
      -> NominalDiffTime
      -> BackoffStrategy
      -> Jitter
      -> ObservabilityHooks m payload
      -> (MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Text -> Int -> Validation Int
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"workerCount" (WorkerConfig m payload -> Int
forall (m :: * -> *) payload. WorkerConfig m payload -> Int
workerCount WorkerConfig m payload
config)
    Validation
  (HandlerMode m payload
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> Maybe NominalDiffTime
   -> NominalDiffTime
   -> BackoffStrategy
   -> Jitter
   -> ObservabilityHooks m payload
   -> (MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation (HandlerMode m payload)
-> Validation
     (NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> Maybe NominalDiffTime
      -> NominalDiffTime
      -> BackoffStrategy
      -> Jitter
      -> ObservabilityHooks m payload
      -> (MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> (WorkerConfig m payload -> HandlerMode m payload
forall (m :: * -> *) payload.
WorkerConfig m payload -> HandlerMode m payload
handlerMode WorkerConfig m payload
config HandlerMode m payload
-> Validation Int -> Validation (HandlerMode m payload)
forall a b. a -> Validation b -> Validation a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Text -> Int -> Validation Int
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"handler batch size" (WorkerConfig m payload -> Int
forall (m :: * -> *) payload. WorkerConfig m payload -> Int
handlerBatchSize WorkerConfig m payload
config))
    Validation
  (NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> Maybe NominalDiffTime
   -> NominalDiffTime
   -> BackoffStrategy
   -> Jitter
   -> ObservabilityHooks m payload
   -> (MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation NominalDiffTime
-> Validation
     (NominalDiffTime
      -> NominalDiffTime
      -> Maybe NominalDiffTime
      -> NominalDiffTime
      -> BackoffStrategy
      -> Jitter
      -> ObservabilityHooks m payload
      -> (MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"pollInterval" (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
pollInterval WorkerConfig m payload
config)
    Validation
  (NominalDiffTime
   -> NominalDiffTime
   -> Maybe NominalDiffTime
   -> NominalDiffTime
   -> BackoffStrategy
   -> Jitter
   -> ObservabilityHooks m payload
   -> (MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation NominalDiffTime
-> Validation
     (NominalDiffTime
      -> Maybe NominalDiffTime
      -> NominalDiffTime
      -> BackoffStrategy
      -> Jitter
      -> ObservabilityHooks m payload
      -> (MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"visibilityTimeout" (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
visibilityTimeout WorkerConfig m payload
config)
    Validation
  (NominalDiffTime
   -> Maybe NominalDiffTime
   -> NominalDiffTime
   -> BackoffStrategy
   -> Jitter
   -> ObservabilityHooks m payload
   -> (MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation NominalDiffTime
-> Validation
     (Maybe NominalDiffTime
      -> NominalDiffTime
      -> BackoffStrategy
      -> Jitter
      -> ObservabilityHooks m payload
      -> (MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"jobHeartbeatInterval" (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
jobHeartbeatInterval WorkerConfig m payload
config)
    Validation
  (Maybe NominalDiffTime
   -> NominalDiffTime
   -> BackoffStrategy
   -> Jitter
   -> ObservabilityHooks m payload
   -> (MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation (Maybe NominalDiffTime)
-> Validation
     (NominalDiffTime
      -> BackoffStrategy
      -> Jitter
      -> ObservabilityHooks m payload
      -> (MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> (NominalDiffTime -> Validation NominalDiffTime)
-> Maybe NominalDiffTime -> Validation (Maybe NominalDiffTime)
forall (t :: * -> *) (f :: * -> *) a b.
(Traversable t, Applicative f) =>
(a -> f b) -> t a -> f (t b)
forall (f :: * -> *) a b.
Applicative f =>
(a -> f b) -> Maybe a -> f (Maybe b)
traverse (Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"maxJobDuration") (WorkerConfig m payload -> Maybe NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> Maybe NominalDiffTime
maxJobDuration WorkerConfig m payload
config)
    Validation
  (NominalDiffTime
   -> BackoffStrategy
   -> Jitter
   -> ObservabilityHooks m payload
   -> (MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation NominalDiffTime
-> Validation
     (BackoffStrategy
      -> Jitter
      -> ObservabilityHooks m payload
      -> (MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"workerHeartbeatInterval" (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
workerHeartbeatInterval WorkerConfig m payload
config)
    Validation
  (BackoffStrategy
   -> Jitter
   -> ObservabilityHooks m payload
   -> (MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation BackoffStrategy
-> Validation
     (Jitter
      -> ObservabilityHooks m payload
      -> (MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> BackoffStrategy -> Validation BackoffStrategy
forall a. a -> Validation a
waived (WorkerConfig m payload -> BackoffStrategy
forall (m :: * -> *) payload.
WorkerConfig m payload -> BackoffStrategy
backoffStrategy WorkerConfig m payload
config)
    Validation
  (Jitter
   -> ObservabilityHooks m payload
   -> (MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation Jitter
-> Validation
     (ObservabilityHooks m payload
      -> (MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Jitter -> Validation Jitter
forall a. a -> Validation a
waived (WorkerConfig m payload -> Jitter
forall (m :: * -> *) payload. WorkerConfig m payload -> Jitter
jitter WorkerConfig m payload
config)
    Validation
  (ObservabilityHooks m payload
   -> (MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation (ObservabilityHooks m payload)
-> Validation
     ((MaintenanceOp -> Int64 -> m ())
      -> WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> ObservabilityHooks m payload
-> Validation (ObservabilityHooks m payload)
forall a. a -> Validation a
waived (WorkerConfig m payload -> ObservabilityHooks m payload
forall (m :: * -> *) payload.
WorkerConfig m payload -> ObservabilityHooks m payload
observabilityHooks WorkerConfig m payload
config)
    Validation
  ((MaintenanceOp -> Int64 -> m ())
   -> WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation (MaintenanceOp -> Int64 -> m ())
-> Validation
     (WorkerRuntime
      -> Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> (MaintenanceOp -> Int64 -> m ())
-> Validation (MaintenanceOp -> Int64 -> m ())
forall a. a -> Validation a
waived (WorkerConfig m payload -> MaintenanceOp -> Int64 -> m ()
forall (m :: * -> *) payload.
WorkerConfig m payload -> MaintenanceOp -> Int64 -> m ()
onMaintenance WorkerConfig m payload
config)
    Validation
  (WorkerRuntime
   -> Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation WorkerRuntime
-> Validation
     (Maybe String
      -> Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> WorkerRuntime -> Validation WorkerRuntime
forall a. a -> Validation a
waived (WorkerConfig m payload -> WorkerRuntime
forall (m :: * -> *) payload.
WorkerConfig m payload -> WorkerRuntime
workerRuntime WorkerConfig m payload
config)
    Validation
  (Maybe String
   -> Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation (Maybe String)
-> Validation
     (Maybe NominalDiffTime
      -> LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Maybe String -> Validation (Maybe String)
forall a. a -> Validation a
waived (WorkerConfig m payload -> Maybe String
forall (m :: * -> *) payload.
WorkerConfig m payload -> Maybe String
livenessFile WorkerConfig m payload
config)
    Validation
  (Maybe NominalDiffTime
   -> LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation (Maybe NominalDiffTime)
-> Validation
     (LogConfig
      -> [CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> (NominalDiffTime -> Validation NominalDiffTime)
-> Maybe NominalDiffTime -> Validation (Maybe NominalDiffTime)
forall (t :: * -> *) (f :: * -> *) a b.
(Traversable t, Applicative f) =>
(a -> f b) -> t a -> f (t b)
forall (f :: * -> *) a b.
Applicative f =>
(a -> f b) -> Maybe a -> f (Maybe b)
traverse (Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
nonNegative Text
"gracefulShutdownTimeout") (WorkerConfig m payload -> Maybe NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> Maybe NominalDiffTime
gracefulShutdownTimeout WorkerConfig m payload
config)
    Validation
  (LogConfig
   -> [CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation LogConfig
-> Validation
     ([CronJob payload]
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> LogConfig -> Validation LogConfig
forall a. a -> Validation a
waived (WorkerConfig m payload -> LogConfig
forall (m :: * -> *) payload. WorkerConfig m payload -> LogConfig
logConfig WorkerConfig m payload
config)
    Validation
  ([CronJob payload]
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation [CronJob payload]
-> Validation
     (NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> [CronJob payload] -> Validation [CronJob payload]
forall a. a -> Validation a
waived (WorkerConfig m payload -> [CronJob payload]
forall (m :: * -> *) payload.
WorkerConfig m payload -> [CronJob payload]
cronJobs WorkerConfig m payload
config)
    Validation
  (NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation NominalDiffTime
-> Validation
     (NominalDiffTime
      -> NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"reaperInterval" (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
reaperInterval WorkerConfig m payload
config)
    Validation
  (NominalDiffTime
   -> NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation NominalDiffTime
-> Validation
     (NominalDiffTime
      -> NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"reaperSparseInterval" (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
reaperSparseInterval WorkerConfig m payload
config)
    Validation
  (NominalDiffTime
   -> NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation NominalDiffTime
-> Validation
     (NominalDiffTime
      -> UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"reaperBucketIdle" (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
reaperBucketIdle WorkerConfig m payload
config)
    Validation
  (NominalDiffTime
   -> UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation NominalDiffTime
-> Validation
     (UUID
      -> Maybe Text
      -> Maybe Value
      -> NominalDiffTime
      -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"reaperTimeout" (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
reaperTimeout WorkerConfig m payload
config)
    Validation
  (UUID
   -> Maybe Text
   -> Maybe Value
   -> NominalDiffTime
   -> WorkerConfig m payload)
-> Validation UUID
-> Validation
     (Maybe Text
      -> Maybe Value -> NominalDiffTime -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> UUID -> Validation UUID
forall a. a -> Validation a
waived (WorkerConfig m payload -> UUID
forall (m :: * -> *) payload. WorkerConfig m payload -> UUID
workerId WorkerConfig m payload
config)
    Validation
  (Maybe Text
   -> Maybe Value -> NominalDiffTime -> WorkerConfig m payload)
-> Validation (Maybe Text)
-> Validation
     (Maybe Value -> NominalDiffTime -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Maybe Text -> Validation (Maybe Text)
forall a. a -> Validation a
waived (WorkerConfig m payload -> Maybe Text
forall (m :: * -> *) payload. WorkerConfig m payload -> Maybe Text
workerHost WorkerConfig m payload
config)
    Validation
  (Maybe Value -> NominalDiffTime -> WorkerConfig m payload)
-> Validation (Maybe Value)
-> Validation (NominalDiffTime -> WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Maybe Value -> Validation (Maybe Value)
forall a. a -> Validation a
waived (WorkerConfig m payload -> Maybe Value
forall (m :: * -> *) payload. WorkerConfig m payload -> Maybe Value
workerMetadata WorkerConfig m payload
config)
    Validation (NominalDiffTime -> WorkerConfig m payload)
-> Validation NominalDiffTime
-> Validation (WorkerConfig m payload)
forall a b. Validation (a -> b) -> Validation a -> Validation b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Text -> NominalDiffTime -> Validation NominalDiffTime
forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
"workerStaleThreshold" (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
workerStaleThreshold WorkerConfig m payload
config)

-- | Invariants spanning more than one field.
crossFieldInvariants :: WorkerConfig m payload -> Validation ()
crossFieldInvariants :: forall (m :: * -> *) payload.
WorkerConfig m payload -> Validation ()
crossFieldInvariants WorkerConfig m payload
config =
  Bool -> Text -> Validation ()
require
    (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
jobHeartbeatInterval WorkerConfig m payload
config NominalDiffTime -> NominalDiffTime -> Bool
forall a. Ord a => a -> a -> Bool
< WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
visibilityTimeout WorkerConfig m payload
config)
    Text
"jobHeartbeatInterval must be less than visibilityTimeout"
    Validation () -> Validation () -> Validation ()
forall a b. Validation a -> Validation b -> Validation b
forall (f :: * -> *) a b. Applicative f => f a -> f b -> f b
*> Bool -> Text -> Validation ()
require
      (WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
workerHeartbeatInterval WorkerConfig m payload
config NominalDiffTime -> NominalDiffTime -> Bool
forall a. Ord a => a -> a -> Bool
< WorkerConfig m payload -> NominalDiffTime
forall (m :: * -> *) payload.
WorkerConfig m payload -> NominalDiffTime
workerStaleThreshold WorkerConfig m payload
config)
      Text
"workerHeartbeatInterval must be less than workerStaleThreshold"

newtype Validation a = Validation (Either (NonEmpty Text) a)

instance Functor Validation where
  fmap :: forall a b. (a -> b) -> Validation a -> Validation b
fmap a -> b
mapper (Validation Either (NonEmpty Text) a
result) = Either (NonEmpty Text) b -> Validation b
forall a. Either (NonEmpty Text) a -> Validation a
Validation ((a -> b) -> Either (NonEmpty Text) a -> Either (NonEmpty Text) b
forall a b.
(a -> b) -> Either (NonEmpty Text) a -> Either (NonEmpty Text) b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap a -> b
mapper Either (NonEmpty Text) a
result)

instance Applicative Validation where
  pure :: forall a. a -> Validation a
pure = Either (NonEmpty Text) a -> Validation a
forall a. Either (NonEmpty Text) a -> Validation a
Validation (Either (NonEmpty Text) a -> Validation a)
-> (a -> Either (NonEmpty Text) a) -> a -> Validation a
forall b c a. (b -> c) -> (a -> b) -> a -> c
. a -> Either (NonEmpty Text) a
forall a b. b -> Either a b
Right
  Validation (Left NonEmpty Text
left) <*> :: forall a b. Validation (a -> b) -> Validation a -> Validation b
<*> Validation (Left NonEmpty Text
right) = Either (NonEmpty Text) b -> Validation b
forall a. Either (NonEmpty Text) a -> Validation a
Validation (NonEmpty Text -> Either (NonEmpty Text) b
forall a b. a -> Either a b
Left (NonEmpty Text
left NonEmpty Text -> NonEmpty Text -> NonEmpty Text
forall a. Semigroup a => a -> a -> a
<> NonEmpty Text
right))
  Validation Either (NonEmpty Text) (a -> b)
apply <*> Validation Either (NonEmpty Text) a
value = Either (NonEmpty Text) b -> Validation b
forall a. Either (NonEmpty Text) a -> Validation a
Validation (Either (NonEmpty Text) (a -> b)
apply Either (NonEmpty Text) (a -> b)
-> Either (NonEmpty Text) a -> Either (NonEmpty Text) b
forall a b.
Either (NonEmpty Text) (a -> b)
-> Either (NonEmpty Text) a -> Either (NonEmpty Text) b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Either (NonEmpty Text) a
value)

-- | A field carrying no invariant of its own.
waived :: a -> Validation a
waived :: forall a. a -> Validation a
waived = a -> Validation a
forall a. a -> Validation a
forall (f :: * -> *) a. Applicative f => a -> f a
pure

-- | Reject a non-positive field, passing it through unchanged.
positive :: (Num a, Ord a) => Text -> a -> Validation a
positive :: forall a. (Num a, Ord a) => Text -> a -> Validation a
positive Text
label a
value = a
value a -> Validation () -> Validation a
forall a b. a -> Validation b -> Validation a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Bool -> Text -> Validation ()
require (a
value a -> a -> Bool
forall a. Ord a => a -> a -> Bool
> a
0) (Text
label Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" must be greater than zero")

-- | Reject a negative field and return a valid field unchanged. Zero disables
-- the applicable wait.
nonNegative :: (Num a, Ord a) => Text -> a -> Validation a
nonNegative :: forall a. (Num a, Ord a) => Text -> a -> Validation a
nonNegative Text
label a
value = a
value a -> Validation () -> Validation a
forall a b. a -> Validation b -> Validation a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Bool -> Text -> Validation ()
require (a
value a -> a -> Bool
forall a. Ord a => a -> a -> Bool
>= a
0) (Text
label Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" must not be negative")

require :: Bool -> Text -> Validation ()
require :: Bool -> Text -> Validation ()
require Bool
condition Text
message = Either (NonEmpty Text) () -> Validation ()
forall a. Either (NonEmpty Text) a -> Validation a
Validation (if Bool
condition then () -> Either (NonEmpty Text) ()
forall a b. b -> Either a b
Right () else NonEmpty Text -> Either (NonEmpty Text) ()
forall a b. a -> Either a b
Left (Text
message Text -> [Text] -> NonEmpty Text
forall a. a -> [a] -> NonEmpty a
:| []))

-- | Create a t'WorkerConfig' running one job per group in a worker transaction
-- held for the duration of the handler.
--
-- The handler returns the result type @payload@'s registry entry declares.
transactionalWorkerConfig
  :: (MonadArbiter n, MonadIO m)
  => Int
  -- ^ Worker count
  -> JobHandler n payload (ResultOf n payload)
  -> m (WorkerConfig n payload)
transactionalWorkerConfig :: forall (n :: * -> *) (m :: * -> *) payload.
(MonadArbiter n, MonadIO m) =>
Int
-> JobHandler n payload (ResultOf n payload)
-> m (WorkerConfig n payload)
transactionalWorkerConfig Int
workerCnt JobHandler n payload (ResultOf n payload)
handler =
  Int -> HandlerMode n payload -> m (WorkerConfig n payload)
forall (n :: * -> *) (m :: * -> *) payload.
(Applicative n, MonadIO m) =>
Int -> HandlerMode n payload -> m (WorkerConfig n payload)
mkDefaultConfig Int
workerCnt (JobHandler n payload (ResultOf n payload) -> HandlerMode n payload
forall (m :: * -> *) payload.
JobHandler m payload (ResultOf m payload) -> HandlerMode m payload
SingleJobMode JobHandler n payload (ResultOf n payload)
handler)

-- | Create a t'WorkerConfig' for batched job processing, no worker transaction.
-- The handler receives the batch and a 'BatchCallbacks' record to finalize each
-- job (ack, fail, cancel, or nack). Jobs left untouched are reprocessed. To store
-- a result per job, ack with 'ackWith' or 'ackAllWith'.
defaultBatchedWorkerConfig
  :: (MonadArbiter n, MonadIO m)
  => Int
  -- ^ Worker count
  -> Int
  -- ^ Batch size (max jobs per group to claim together)
  -> (NonEmpty (JobRead payload) -> BatchCallbacks n payload (ResultOf n payload) -> n ())
  -> m (WorkerConfig n payload)
defaultBatchedWorkerConfig :: forall (n :: * -> *) (m :: * -> *) payload.
(MonadArbiter n, MonadIO m) =>
Int
-> Int
-> (NonEmpty (JobRead payload)
    -> BatchCallbacks n payload (ResultOf n payload) -> n ())
-> m (WorkerConfig n payload)
defaultBatchedWorkerConfig Int
workerCnt Int
batchSize NonEmpty (JobRead payload)
-> BatchCallbacks n payload (ResultOf n payload) -> n ()
handler =
  Int -> HandlerMode n payload -> m (WorkerConfig n payload)
forall (n :: * -> *) (m :: * -> *) payload.
(Applicative n, MonadIO m) =>
Int -> HandlerMode n payload -> m (WorkerConfig n payload)
mkDefaultConfig Int
workerCnt (Int
-> (NonEmpty (JobRead payload)
    -> BatchCallbacks n payload (ResultOf n payload) -> n ())
-> HandlerMode n payload
forall (m :: * -> *) payload.
Int
-> (NonEmpty (JobRead payload)
    -> BatchCallbacks m payload (ResultOf m payload) -> m ())
-> HandlerMode m payload
BatchedJobsMode Int
batchSize NonEmpty (JobRead payload)
-> BatchCallbacks n payload (ResultOf n payload) -> n ()
handler)

-- | Create a t'WorkerConfig' running one job at a time, no worker transaction.
-- The handler finalizes the job through 'BatchCallbacks'. An unfinalized job is
-- reprocessed.
manualWorkerConfig
  :: (MonadArbiter n, MonadIO m)
  => Int
  -- ^ Worker count
  -> (JobRead payload -> BatchCallbacks n payload (ResultOf n payload) -> n ())
  -> m (WorkerConfig n payload)
manualWorkerConfig :: forall (n :: * -> *) (m :: * -> *) payload.
(MonadArbiter n, MonadIO m) =>
Int
-> (JobRead payload
    -> BatchCallbacks n payload (ResultOf n payload) -> n ())
-> m (WorkerConfig n payload)
manualWorkerConfig Int
workerCnt JobRead payload
-> BatchCallbacks n payload (ResultOf n payload) -> n ()
handler =
  Int
-> Int
-> (NonEmpty (JobRead payload)
    -> BatchCallbacks n payload (ResultOf n payload) -> n ())
-> m (WorkerConfig n payload)
forall (n :: * -> *) (m :: * -> *) payload.
(MonadArbiter n, MonadIO m) =>
Int
-> Int
-> (NonEmpty (JobRead payload)
    -> BatchCallbacks n payload (ResultOf n payload) -> n ())
-> m (WorkerConfig n payload)
defaultBatchedWorkerConfig Int
workerCnt Int
1 (\(JobRead payload
job :| [JobRead payload]
_) -> JobRead payload
-> BatchCallbacks n payload (ResultOf n payload) -> n ()
handler JobRead payload
job)

-- | Rework a pool's observability hooks.
withHooks
  :: (ObservabilityHooks m payload -> ObservabilityHooks m payload)
  -> WorkerConfig m payload
  -> WorkerConfig m payload
withHooks :: forall (m :: * -> *) payload.
(ObservabilityHooks m payload -> ObservabilityHooks m payload)
-> WorkerConfig m payload -> WorkerConfig m payload
withHooks ObservabilityHooks m payload -> ObservabilityHooks m payload
rework WorkerConfig m payload
cfg = WorkerConfig m payload
cfg {observabilityHooks = rework (observabilityHooks cfg)}

-- | Run @report@ before the pool's own maintenance callback. Each runs even when the other fails.
withMaintenance
  :: (MonadUnliftIO m)
  => (MaintenanceOp -> Int64 -> m ())
  -> WorkerConfig m payload
  -> WorkerConfig m payload
withMaintenance :: forall (m :: * -> *) payload.
MonadUnliftIO m =>
(MaintenanceOp -> Int64 -> m ())
-> WorkerConfig m payload -> WorkerConfig m payload
withMaintenance MaintenanceOp -> Int64 -> m ()
report WorkerConfig m payload
cfg =
  WorkerConfig m payload
cfg {onMaintenance = \MaintenanceOp
operation Int64
count -> MaintenanceOp -> Int64 -> m ()
report MaintenanceOp
operation Int64
count m () -> m () -> m ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
`andThen` WorkerConfig m payload -> MaintenanceOp -> Int64 -> m ()
forall (m :: * -> *) payload.
WorkerConfig m payload -> MaintenanceOp -> Int64 -> m ()
onMaintenance WorkerConfig m payload
cfg MaintenanceOp
operation Int64
count}

-- | Internal helper to create a config with the given handler mode.
mkDefaultConfig
  :: (Applicative n, MonadIO m)
  => Int
  -> HandlerMode n payload
  -> m (WorkerConfig n payload)
mkDefaultConfig :: forall (n :: * -> *) (m :: * -> *) payload.
(Applicative n, MonadIO m) =>
Int -> HandlerMode n payload -> m (WorkerConfig n payload)
mkDefaultConfig Int
workerCnt HandlerMode n payload
mode = do
  heartbeatTMVar <- IO (TMVar ()) -> m (TMVar ())
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO IO (TMVar ())
forall (m :: * -> *) a. MonadIO m => m (TMVar a)
newEmptyTMVarIO
  shutdownTVar <- newTVarIO Running
  pauseTVar <- newTVarIO False
  pauseEpochTVar <- newTVarIO 0
  listenerReadyTVar <- newTVarIO False
  uuid <- liftIO UUID.nextRandom
  tmpDir <- liftIO getTemporaryDirectory
  host <- liftIO getHostName
  let livenessPath = String
tmpDir String -> ShowS
forall a. Semigroup a => a -> a -> a
<> String
"/arbiter-worker-" String -> ShowS
forall a. Semigroup a => a -> a -> a
<> UUID -> String
toString UUID
uuid
  pure
    WorkerConfig
      { workerCount = workerCnt
      , handlerMode = mode
      , pollInterval = 5
      , visibilityTimeout = 60
      , jobHeartbeatInterval = 30
      , maxJobDuration = Nothing
      , workerHeartbeatInterval = 10
      , backoffStrategy = exponentialBackoff 2.0 1_048_576
      , jitter = EqualJitter
      , observabilityHooks = defaultObservabilityHooks
      , onMaintenance = \MaintenanceOp
_ Int64
_ -> () -> n ()
forall a. a -> n a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
      , workerRuntime =
          WorkerRuntime
            { runtimeStateVar = shutdownTVar
            , runtimePauseVar = pauseTVar
            , runtimePauseEpoch = pauseEpochTVar
            , runtimeHeartbeatSignal = heartbeatTMVar
            , runtimeListenerReadyVar = listenerReadyTVar
            }
      , livenessFile = Just livenessPath
      , gracefulShutdownTimeout = Just 30
      , logConfig = withWorkerIdContext uuid defaultLogConfig
      , cronJobs = []
      , reaperInterval = 300
      , reaperSparseInterval = 3600
      , reaperBucketIdle = 300
      , reaperTimeout = 300
      , workerId = uuid
      , workerHost = Just (T.pack host)
      , workerMetadata = Nothing
      , workerStaleThreshold = 300
      }

withWorkerIdContext :: UUID -> LogConfig -> LogConfig
withWorkerIdContext :: UUID -> LogConfig -> LogConfig
withWorkerIdContext UUID
workerId LogConfig
logCfg =
  LogConfig
logCfg {identityContext = identityContext logCfg <> ["worker_id" .= workerId]}

-- | Run/shutdown state for a pool.
workerStateVar :: WorkerConfig n payload -> TVar WorkerState
workerStateVar :: forall (n :: * -> *) payload.
WorkerConfig n payload -> TVar WorkerState
workerStateVar = WorkerRuntime -> TVar WorkerState
runtimeStateVar (WorkerRuntime -> TVar WorkerState)
-> (WorkerConfig n payload -> WorkerRuntime)
-> WorkerConfig n payload
-> TVar WorkerState
forall b c a. (b -> c) -> (a -> b) -> a -> c
. WorkerConfig n payload -> WorkerRuntime
forall (m :: * -> *) payload.
WorkerConfig m payload -> WorkerRuntime
workerRuntime

-- | Pause state for a pool.
pauseVar :: WorkerConfig n payload -> TVar Bool
pauseVar :: forall (n :: * -> *) payload. WorkerConfig n payload -> TVar Bool
pauseVar = WorkerRuntime -> TVar Bool
runtimePauseVar (WorkerRuntime -> TVar Bool)
-> (WorkerConfig n payload -> WorkerRuntime)
-> WorkerConfig n payload
-> TVar Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. WorkerConfig n payload -> WorkerRuntime
forall (m :: * -> *) payload.
WorkerConfig m payload -> WorkerRuntime
workerRuntime

-- | Version of the current pause state.
pauseEpoch :: WorkerConfig n payload -> TVar Word64
pauseEpoch :: forall (n :: * -> *) payload. WorkerConfig n payload -> TVar Word64
pauseEpoch = WorkerRuntime -> TVar Word64
runtimePauseEpoch (WorkerRuntime -> TVar Word64)
-> (WorkerConfig n payload -> WorkerRuntime)
-> WorkerConfig n payload
-> TVar Word64
forall b c a. (b -> c) -> (a -> b) -> a -> c
. WorkerConfig n payload -> WorkerRuntime
forall (m :: * -> *) payload.
WorkerConfig m payload -> WorkerRuntime
workerRuntime

-- | Signal used to coordinate worker heartbeats.
heartbeatSignal :: WorkerConfig n payload -> TMVar ()
heartbeatSignal :: forall (n :: * -> *) payload. WorkerConfig n payload -> TMVar ()
heartbeatSignal = WorkerRuntime -> TMVar ()
runtimeHeartbeatSignal (WorkerRuntime -> TMVar ())
-> (WorkerConfig n payload -> WorkerRuntime)
-> WorkerConfig n payload
-> TMVar ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. WorkerConfig n payload -> WorkerRuntime
forall (m :: * -> *) payload.
WorkerConfig m payload -> WorkerRuntime
workerRuntime

-- | State of the pool's notification listener.
listenerReadyVar :: WorkerConfig n payload -> TVar Bool
listenerReadyVar :: forall (n :: * -> *) payload. WorkerConfig n payload -> TVar Bool
listenerReadyVar = WorkerRuntime -> TVar Bool
runtimeListenerReadyVar (WorkerRuntime -> TVar Bool)
-> (WorkerConfig n payload -> WorkerRuntime)
-> WorkerConfig n payload
-> TVar Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. WorkerConfig n payload -> WorkerRuntime
forall (m :: * -> *) payload.
WorkerConfig m payload -> WorkerRuntime
workerRuntime

-- | Initiate graceful shutdown of the worker pool
--
-- Stops claiming new jobs. In-flight jobs will complete, then the pool exits.
shutdownWorker :: (MonadIO m) => WorkerConfig n payload -> m ()
shutdownWorker :: forall (m :: * -> *) (n :: * -> *) payload.
MonadIO m =>
WorkerConfig n payload -> m ()
shutdownWorker WorkerConfig n payload
config = IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> m ()) -> (STM () -> IO ()) -> STM () -> m ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. STM () -> IO ()
forall (m :: * -> *) a. MonadIO m => STM a -> m a
STM.atomically (STM () -> m ()) -> STM () -> m ()
forall a b. (a -> b) -> a -> b
$ TVar WorkerState -> WorkerState -> STM ()
forall a. TVar a -> a -> STM ()
STM.writeTVar (WorkerConfig n payload -> TVar WorkerState
forall (n :: * -> *) payload.
WorkerConfig n payload -> TVar WorkerState
workerStateVar WorkerConfig n payload
config) WorkerState
ShuttingDown

-- | The pool's state, with a pause reported as paused.
getWorkerState :: (MonadIO m) => WorkerConfig n payload -> m WorkerState
getWorkerState :: forall (m :: * -> *) (n :: * -> *) payload.
MonadIO m =>
WorkerConfig n payload -> m WorkerState
getWorkerState WorkerConfig n payload
config = IO WorkerState -> m WorkerState
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO WorkerState -> m WorkerState)
-> (STM WorkerState -> IO WorkerState)
-> STM WorkerState
-> m WorkerState
forall b c a. (b -> c) -> (a -> b) -> a -> c
. STM WorkerState -> IO WorkerState
forall (m :: * -> *) a. MonadIO m => STM a -> m a
STM.atomically (STM WorkerState -> m WorkerState)
-> STM WorkerState -> m WorkerState
forall a b. (a -> b) -> a -> b
$ WorkerConfig n payload -> STM WorkerState
forall (n :: * -> *) payload.
WorkerConfig n payload -> STM WorkerState
readEffectiveState WorkerConfig n payload
config

-- | Whether this pool's LISTEN channels are subscribed (or there is no listener).
getListenerReady :: (MonadIO m) => WorkerConfig n payload -> m Bool
getListenerReady :: forall (m :: * -> *) (n :: * -> *) payload.
MonadIO m =>
WorkerConfig n payload -> m Bool
getListenerReady WorkerConfig n payload
config = IO Bool -> m Bool
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO Bool -> m Bool) -> (STM Bool -> IO Bool) -> STM Bool -> m Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. STM Bool -> IO Bool
forall (m :: * -> *) a. MonadIO m => STM a -> m a
STM.atomically (STM Bool -> m Bool) -> STM Bool -> m Bool
forall a b. (a -> b) -> a -> b
$ TVar Bool -> STM Bool
forall a. TVar a -> STM a
STM.readTVar (WorkerConfig n payload -> TVar Bool
forall (n :: * -> *) payload. WorkerConfig n payload -> TVar Bool
listenerReadyVar WorkerConfig n payload
config)

-- | Set the pause flag, superseding any reading a caller has in flight.
writePause :: WorkerConfig n payload -> Bool -> STM.STM ()
writePause :: forall (n :: * -> *) payload.
WorkerConfig n payload -> Bool -> STM ()
writePause WorkerConfig n payload
config Bool
paused = do
  TVar Bool -> Bool -> STM ()
forall a. TVar a -> a -> STM ()
STM.writeTVar (WorkerConfig n payload -> TVar Bool
forall (n :: * -> *) payload. WorkerConfig n payload -> TVar Bool
pauseVar WorkerConfig n payload
config) Bool
paused
  TVar Word64 -> (Word64 -> Word64) -> STM ()
forall a. TVar a -> (a -> a) -> STM ()
STM.modifyTVar' (WorkerConfig n payload -> TVar Word64
forall (n :: * -> *) payload. WorkerConfig n payload -> TVar Word64
pauseEpoch WorkerConfig n payload
config) (Word64 -> Word64 -> Word64
forall a. Num a => a -> a -> a
+ Word64
1)

-- | Apply a pause reading taken at @epoch@. A newer write wins.
writePauseIfCurrent :: WorkerConfig n payload -> Word64 -> Bool -> STM.STM ()
writePauseIfCurrent :: forall (n :: * -> *) payload.
WorkerConfig n payload -> Word64 -> Bool -> STM ()
writePauseIfCurrent WorkerConfig n payload
config Word64
epoch Bool
paused = do
  current <- TVar Word64 -> STM Word64
forall a. TVar a -> STM a
STM.readTVar (WorkerConfig n payload -> TVar Word64
forall (n :: * -> *) payload. WorkerConfig n payload -> TVar Word64
pauseEpoch WorkerConfig n payload
config)
  when (current == epoch) $ writePause config paused

-- | 'getWorkerState' inside 'STM.STM'.
readEffectiveState :: WorkerConfig n payload -> STM.STM WorkerState
readEffectiveState :: forall (n :: * -> *) payload.
WorkerConfig n payload -> STM WorkerState
readEffectiveState WorkerConfig n payload
config = do
  state <- TVar WorkerState -> STM WorkerState
forall a. TVar a -> STM a
STM.readTVar (WorkerConfig n payload -> TVar WorkerState
forall (n :: * -> *) payload.
WorkerConfig n payload -> TVar WorkerState
workerStateVar WorkerConfig n payload
config)
  case state of
    WorkerState
ShuttingDown -> WorkerState -> STM WorkerState
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure WorkerState
ShuttingDown
    WorkerState
_ -> do
      paused <- TVar Bool -> STM Bool
forall a. TVar a -> STM a
STM.readTVar (WorkerConfig n payload -> TVar Bool
forall (n :: * -> *) payload. WorkerConfig n payload -> TVar Bool
pauseVar WorkerConfig n payload
config)
      pure $ if paused then Paused else state