{-# LANGUAGE FlexibleInstances #-}
{-# LANGUAGE OverloadedStrings #-}

-- | Job records, the enqueue setters, and the lifecycle hook types.
module Arbiter.Core.Job.Types
  ( -- * Core Job Type
    Job
  , HasKind (..)
  , constructorKind
  , constructorKinds
  , PayloadKeys (..)
  , PayloadColumns (..)
  , JobRead
  , JobWrite
  , primaryKey
  , payload
  , queueName
  , groupKey
  , insertedAt
  , updatedAt
  , attempts
  , lastError
  , priority
  , lastAttemptedAt
  , notVisibleUntil
  , dedupKey
  , maxAttempts
  , parentId
  , parentState
  , traceContext
  , suspended
  , claimedBy
  , claimSeq
  , archiveFor
  , payloadKeys
  , defaultJob
  , defaultGroupedJob
  , setPayload
  , setGroupKey
  , setPriority
  , setNotVisibleUntil
  , setDedupKey
  , setMaxAttempts
  , setTraceContext
  , setArchiveFor
  , mapPayload
  , defaultMaxAttempts
  , dayRetention
  , isRollup

    -- * Derived status
  , JobStatus (..)
  , jobStatusToText
  , jobStatusFromText

    -- * Type Constraints
  , JobPayload
  , RegistryAdmissionPolicies

    -- * Deduplication
  , DedupKey (..)
  , dedupParts

    -- * Trace context
  , TraceContext (..)
  , toTraceContext

    -- * Observability
  , ObservabilityHooks (..)
  , defaultObservabilityHooks
  , andThen
  , JobId
  , ClaimSeq
  , ClaimTime
  , CurrentTime
  , StartTime
  , EndTime
  , ErrorMsg
  , BackoffDelay
  ) where

import Control.Exception qualified as E
import Data.Aeson (FromJSON (..), ToJSON (..), withObject, (.!=), (.:), (.:?))
import Data.Int (Int32, Int64)
import Data.Maybe (isJust)
import Data.Text (Text)
import Data.Time (NominalDiffTime, UTCTime)
import GHC.Generics (Generic)
import UnliftIO (MonadUnliftIO, withRunInIO)

import Arbiter.Core.Concurrency.Spec (ConcurrencyKey, HasConcurrency, RegistryConcurrencyPolicies)
import Arbiter.Core.Job.Dedup (DedupKey (..), dedupParts)
import Arbiter.Core.Job.Kind (HasKind (..), constructorKind, constructorKinds)
import Arbiter.Core.Job.Status (JobStatus (..), jobStatusFromText, jobStatusToText)
import Arbiter.Core.Job.TraceContext (TraceContext (..), toTraceContext)
import Arbiter.Core.Job.Types.Internal
  ( JobRecord (..)
  , archiveFor
  , attempts
  , claimSeq
  , claimedBy
  , dedupKey
  , groupKey
  , insertedAt
  , lastAttemptedAt
  , lastError
  , maxAttempts
  , notVisibleUntil
  , parentId
  , parentState
  , payload
  , payloadKeys
  , primaryKey
  , priority
  , queueName
  , suspended
  , traceContext
  , updatedAt
  )
import Arbiter.Core.RateLimit.Spec (HasRateLimit, RateLimitKey, RegistryRateLimitPolicies)

-- | A job parametrized over payload, primary key, queue name, insertion
-- timestamp, and the columns derived from the payload. The constructor is internal.
type Job payload key q insertedAt adm =
  JobRecord payload key q insertedAt adm

-- | The labels and keys a stored job carries from its payload, one field each.
data PayloadKeys = PayloadKeys
  { PayloadKeys -> Maybe Text
jobKind :: Maybe Text
  -- ^ From the payload's 'Arbiter.Core.Job.Kind.HasKind' instance.
  , PayloadKeys -> Maybe RateLimitKey
jobRateLimitKey :: Maybe RateLimitKey
  -- ^ From the payload's 'Arbiter.Core.RateLimit.Spec.HasRateLimit' instance.
  , PayloadKeys -> Maybe ConcurrencyKey
jobConcurrencyKey :: Maybe ConcurrencyKey
  -- ^ From the payload's 'Arbiter.Core.Concurrency.Spec.HasConcurrency' instance.
  }
  deriving stock (PayloadKeys -> PayloadKeys -> Bool
(PayloadKeys -> PayloadKeys -> Bool)
-> (PayloadKeys -> PayloadKeys -> Bool) -> Eq PayloadKeys
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: PayloadKeys -> PayloadKeys -> Bool
== :: PayloadKeys -> PayloadKeys -> Bool
$c/= :: PayloadKeys -> PayloadKeys -> Bool
/= :: PayloadKeys -> PayloadKeys -> Bool
Eq, (forall x. PayloadKeys -> Rep PayloadKeys x)
-> (forall x. Rep PayloadKeys x -> PayloadKeys)
-> Generic PayloadKeys
forall x. Rep PayloadKeys x -> PayloadKeys
forall x. PayloadKeys -> Rep PayloadKeys x
forall a.
(forall x. a -> Rep a x) -> (forall x. Rep a x -> a) -> Generic a
$cfrom :: forall x. PayloadKeys -> Rep PayloadKeys x
from :: forall x. PayloadKeys -> Rep PayloadKeys x
$cto :: forall x. Rep PayloadKeys x -> PayloadKeys
to :: forall x. Rep PayloadKeys x -> PayloadKeys
Generic, Int -> PayloadKeys -> ShowS
[PayloadKeys] -> ShowS
PayloadKeys -> String
(Int -> PayloadKeys -> ShowS)
-> (PayloadKeys -> String)
-> ([PayloadKeys] -> ShowS)
-> Show PayloadKeys
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> PayloadKeys -> ShowS
showsPrec :: Int -> PayloadKeys -> ShowS
$cshow :: PayloadKeys -> String
show :: PayloadKeys -> String
$cshowList :: [PayloadKeys] -> ShowS
showList :: [PayloadKeys] -> ShowS
Show)

-- | The writable columns resolved from a payload at enqueue. The @kind@, @key@
-- and @prefix@ columns round-trip via 'PayloadKeys'. @cost@ is write-only.
data PayloadColumns = PayloadColumns
  { PayloadColumns -> Maybe Text
pcKind :: Maybe Text
  , PayloadColumns -> Maybe Text
pcRateLimitKey :: Maybe Text
  , PayloadColumns -> Maybe Text
pcRateLimitPrefix :: Maybe Text
  , PayloadColumns -> Double
pcRateLimitCost :: Double
  , PayloadColumns -> Maybe Text
pcConcurrencyKey :: Maybe Text
  , PayloadColumns -> Maybe Text
pcConcurrencyPrefix :: Maybe Text
  }
  deriving stock (PayloadColumns -> PayloadColumns -> Bool
(PayloadColumns -> PayloadColumns -> Bool)
-> (PayloadColumns -> PayloadColumns -> Bool) -> Eq PayloadColumns
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: PayloadColumns -> PayloadColumns -> Bool
== :: PayloadColumns -> PayloadColumns -> Bool
$c/= :: PayloadColumns -> PayloadColumns -> Bool
/= :: PayloadColumns -> PayloadColumns -> Bool
Eq, (forall x. PayloadColumns -> Rep PayloadColumns x)
-> (forall x. Rep PayloadColumns x -> PayloadColumns)
-> Generic PayloadColumns
forall x. Rep PayloadColumns x -> PayloadColumns
forall x. PayloadColumns -> Rep PayloadColumns x
forall a.
(forall x. a -> Rep a x) -> (forall x. Rep a x -> a) -> Generic a
$cfrom :: forall x. PayloadColumns -> Rep PayloadColumns x
from :: forall x. PayloadColumns -> Rep PayloadColumns x
$cto :: forall x. Rep PayloadColumns x -> PayloadColumns
to :: forall x. Rep PayloadColumns x -> PayloadColumns
Generic, Int -> PayloadColumns -> ShowS
[PayloadColumns] -> ShowS
PayloadColumns -> String
(Int -> PayloadColumns -> ShowS)
-> (PayloadColumns -> String)
-> ([PayloadColumns] -> ShowS)
-> Show PayloadColumns
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> PayloadColumns -> ShowS
showsPrec :: Int -> PayloadColumns -> ShowS
$cshow :: PayloadColumns -> String
show :: PayloadColumns -> String
$cshowList :: [PayloadColumns] -> ShowS
showList :: [PayloadColumns] -> ShowS
Show)

-- | Default attempt limit stamped onto jobs whose 'maxAttempts' is unset.
defaultMaxAttempts :: Int32
defaultMaxAttempts :: Int32
defaultMaxAttempts = Int32
10

-- | 24h in seconds, a convenience value for 'archiveFor'.
dayRetention :: Int32
dayRetention :: Int32
dayRetention = Int32
86400

-- | A rollup finalizer is any job whose 'parentState' snapshot is present
-- (an empty object on insert, the merged child results before a DLQ move).
isRollup :: Job p Int64 q t adm -> Bool
isRollup :: forall p q t adm. Job p Int64 q t adm -> Bool
isRollup = Maybe Value -> Bool
forall a. Maybe a -> Bool
isJust (Maybe Value -> Bool)
-> (Job p Int64 q t adm -> Maybe Value)
-> Job p Int64 q t adm
-> Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Job p Int64 q t adm -> Maybe Value
forall payload q insertedAt adm.
JobRecord payload Int64 q insertedAt adm -> Maybe Value
parentState

-- | A job read from the database.
type JobRead payload = Job payload Int64 Text UTCTime PayloadKeys

-- | A job ready to enqueue. Arbiter owns claim, retry, parent, rollup, and
-- suspension state. Use the exported setters to configure enqueue fields.
type JobWrite payload = Job payload () () () ()

-- | Decode the complete persisted representation of a job.
instance (FromJSON payload) => FromJSON (JobRead payload) where
  parseJSON :: Value -> Parser (JobRead payload)
parseJSON = String
-> (Object -> Parser (JobRead payload))
-> Value
-> Parser (JobRead payload)
forall a. String -> (Object -> Parser a) -> Value -> Parser a
withObject String
"Job" ((Object -> Parser (JobRead payload))
 -> Value -> Parser (JobRead payload))
-> (Object -> Parser (JobRead payload))
-> Value
-> Parser (JobRead payload)
forall a b. (a -> b) -> a -> b
$ \Object
obj ->
    Int64
-> payload
-> Text
-> Maybe Text
-> UTCTime
-> Maybe UTCTime
-> Int32
-> Maybe Text
-> Int32
-> Maybe UTCTime
-> Maybe UTCTime
-> Maybe DedupKey
-> Maybe Int32
-> Maybe Int64
-> Maybe Value
-> Maybe TraceContext
-> Bool
-> Maybe UUID
-> Int64
-> Maybe Int32
-> PayloadKeys
-> JobRead payload
forall payload key q insertedAt adm.
key
-> payload
-> q
-> Maybe Text
-> insertedAt
-> Maybe UTCTime
-> Int32
-> Maybe Text
-> Int32
-> Maybe UTCTime
-> Maybe UTCTime
-> Maybe DedupKey
-> Maybe Int32
-> Maybe Int64
-> Maybe Value
-> Maybe TraceContext
-> Bool
-> Maybe UUID
-> Int64
-> Maybe Int32
-> adm
-> JobRecord payload key q insertedAt adm
Job
      (Int64
 -> payload
 -> Text
 -> Maybe Text
 -> UTCTime
 -> Maybe UTCTime
 -> Int32
 -> Maybe Text
 -> Int32
 -> Maybe UTCTime
 -> Maybe UTCTime
 -> Maybe DedupKey
 -> Maybe Int32
 -> Maybe Int64
 -> Maybe Value
 -> Maybe TraceContext
 -> Bool
 -> Maybe UUID
 -> Int64
 -> Maybe Int32
 -> PayloadKeys
 -> JobRead payload)
-> Parser Int64
-> Parser
     (payload
      -> Text
      -> Maybe Text
      -> UTCTime
      -> Maybe UTCTime
      -> Int32
      -> Maybe Text
      -> Int32
      -> Maybe UTCTime
      -> Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Object
obj Object -> Key -> Parser Int64
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"primaryKey"
      Parser
  (payload
   -> Text
   -> Maybe Text
   -> UTCTime
   -> Maybe UTCTime
   -> Int32
   -> Maybe Text
   -> Int32
   -> Maybe UTCTime
   -> Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser payload
-> Parser
     (Text
      -> Maybe Text
      -> UTCTime
      -> Maybe UTCTime
      -> Int32
      -> Maybe Text
      -> Int32
      -> Maybe UTCTime
      -> Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser payload
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"payload"
      Parser
  (Text
   -> Maybe Text
   -> UTCTime
   -> Maybe UTCTime
   -> Int32
   -> Maybe Text
   -> Int32
   -> Maybe UTCTime
   -> Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser Text
-> Parser
     (Maybe Text
      -> UTCTime
      -> Maybe UTCTime
      -> Int32
      -> Maybe Text
      -> Int32
      -> Maybe UTCTime
      -> Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser Text
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"queueName"
      Parser
  (Maybe Text
   -> UTCTime
   -> Maybe UTCTime
   -> Int32
   -> Maybe Text
   -> Int32
   -> Maybe UTCTime
   -> Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe Text)
-> Parser
     (UTCTime
      -> Maybe UTCTime
      -> Int32
      -> Maybe Text
      -> Int32
      -> Maybe UTCTime
      -> Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe Text)
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"groupKey"
      Parser
  (UTCTime
   -> Maybe UTCTime
   -> Int32
   -> Maybe Text
   -> Int32
   -> Maybe UTCTime
   -> Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser UTCTime
-> Parser
     (Maybe UTCTime
      -> Int32
      -> Maybe Text
      -> Int32
      -> Maybe UTCTime
      -> Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser UTCTime
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"insertedAt"
      Parser
  (Maybe UTCTime
   -> Int32
   -> Maybe Text
   -> Int32
   -> Maybe UTCTime
   -> Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe UTCTime)
-> Parser
     (Int32
      -> Maybe Text
      -> Int32
      -> Maybe UTCTime
      -> Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe UTCTime)
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"updatedAt"
      Parser
  (Int32
   -> Maybe Text
   -> Int32
   -> Maybe UTCTime
   -> Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser Int32
-> Parser
     (Maybe Text
      -> Int32
      -> Maybe UTCTime
      -> Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser Int32
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"attempts"
      Parser
  (Maybe Text
   -> Int32
   -> Maybe UTCTime
   -> Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe Text)
-> Parser
     (Int32
      -> Maybe UTCTime
      -> Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe Text)
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"lastError"
      Parser
  (Int32
   -> Maybe UTCTime
   -> Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser Int32
-> Parser
     (Maybe UTCTime
      -> Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser Int32
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"priority"
      Parser
  (Maybe UTCTime
   -> Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe UTCTime)
-> Parser
     (Maybe UTCTime
      -> Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe UTCTime)
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"lastAttemptedAt"
      Parser
  (Maybe UTCTime
   -> Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe UTCTime)
-> Parser
     (Maybe DedupKey
      -> Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe UTCTime)
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"notVisibleUntil"
      Parser
  (Maybe DedupKey
   -> Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe DedupKey)
-> Parser
     (Maybe Int32
      -> Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe DedupKey)
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"dedupKey"
      Parser
  (Maybe Int32
   -> Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe Int32)
-> Parser
     (Maybe Int64
      -> Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe Int32)
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"maxAttempts"
      Parser
  (Maybe Int64
   -> Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe Int64)
-> Parser
     (Maybe Value
      -> Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe Int64)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"parentId"
      Parser
  (Maybe Value
   -> Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe Value)
-> Parser
     (Maybe TraceContext
      -> Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe Value)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"parentState"
      Parser
  (Maybe TraceContext
   -> Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser (Maybe TraceContext)
-> Parser
     (Bool
      -> Maybe UUID
      -> Int64
      -> Maybe Int32
      -> PayloadKeys
      -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> (Maybe Text -> Maybe Text -> Maybe TraceContext
toTraceContext (Maybe Text -> Maybe Text -> Maybe TraceContext)
-> Parser (Maybe Text) -> Parser (Maybe Text -> Maybe TraceContext)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Object
obj Object -> Key -> Parser (Maybe Text)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"traceparent" Parser (Maybe Text -> Maybe TraceContext)
-> Parser (Maybe Text) -> Parser (Maybe TraceContext)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe Text)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"tracestate")
      Parser
  (Bool
   -> Maybe UUID
   -> Int64
   -> Maybe Int32
   -> PayloadKeys
   -> JobRead payload)
-> Parser Bool
-> Parser
     (Maybe UUID
      -> Int64 -> Maybe Int32 -> PayloadKeys -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe Bool)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"suspended" Parser (Maybe Bool) -> Bool -> Parser Bool
forall a. Parser (Maybe a) -> a -> Parser a
.!= Bool
False
      Parser
  (Maybe UUID
   -> Int64 -> Maybe Int32 -> PayloadKeys -> JobRead payload)
-> Parser (Maybe UUID)
-> Parser (Int64 -> Maybe Int32 -> PayloadKeys -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe UUID)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"claimedBy"
      Parser (Int64 -> Maybe Int32 -> PayloadKeys -> JobRead payload)
-> Parser Int64
-> Parser (Maybe Int32 -> PayloadKeys -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe Int64)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"claimSeq" Parser (Maybe Int64) -> Int64 -> Parser Int64
forall a. Parser (Maybe a) -> a -> Parser a
.!= Int64
0
      Parser (Maybe Int32 -> PayloadKeys -> JobRead payload)
-> Parser (Maybe Int32) -> Parser (PayloadKeys -> JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe Int32)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"archiveFor"
      Parser (PayloadKeys -> JobRead payload)
-> Parser PayloadKeys -> Parser (JobRead payload)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> (Maybe Text
-> Maybe RateLimitKey -> Maybe ConcurrencyKey -> PayloadKeys
PayloadKeys (Maybe Text
 -> Maybe RateLimitKey -> Maybe ConcurrencyKey -> PayloadKeys)
-> Parser (Maybe Text)
-> Parser
     (Maybe RateLimitKey -> Maybe ConcurrencyKey -> PayloadKeys)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Object
obj Object -> Key -> Parser (Maybe Text)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"kind" Parser (Maybe RateLimitKey -> Maybe ConcurrencyKey -> PayloadKeys)
-> Parser (Maybe RateLimitKey)
-> Parser (Maybe ConcurrencyKey -> PayloadKeys)
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe RateLimitKey)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"rateLimit" Parser (Maybe ConcurrencyKey -> PayloadKeys)
-> Parser (Maybe ConcurrencyKey) -> Parser PayloadKeys
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
obj Object -> Key -> Parser (Maybe ConcurrencyKey)
forall a. FromJSON a => Object -> Key -> Parser (Maybe a)
.:? Key
"concurrency")

-- | Ungrouped 'JobWrite' with default values. For serial processing within a
-- group, use 'defaultGroupedJob'.
defaultJob :: payload -> JobWrite payload
defaultJob :: forall payload. payload -> JobWrite payload
defaultJob payload
value =
  Job
    { primaryKey :: ()
primaryKey = ()
    , payload :: payload
payload = payload
value
    , queueName :: ()
queueName = ()
    , groupKey :: Maybe Text
groupKey = Maybe Text
forall a. Maybe a
Nothing
    , insertedAt :: ()
insertedAt = ()
    , updatedAt :: Maybe UTCTime
updatedAt = Maybe UTCTime
forall a. Maybe a
Nothing
    , attempts :: Int32
attempts = Int32
0
    , lastError :: Maybe Text
lastError = Maybe Text
forall a. Maybe a
Nothing
    , priority :: Int32
priority = Int32
0
    , lastAttemptedAt :: Maybe UTCTime
lastAttemptedAt = Maybe UTCTime
forall a. Maybe a
Nothing
    , notVisibleUntil :: Maybe UTCTime
notVisibleUntil = Maybe UTCTime
forall a. Maybe a
Nothing
    , dedupKey :: Maybe DedupKey
dedupKey = Maybe DedupKey
forall a. Maybe a
Nothing
    , maxAttempts :: Maybe Int32
maxAttempts = Maybe Int32
forall a. Maybe a
Nothing
    , parentId :: Maybe Int64
parentId = Maybe Int64
forall a. Maybe a
Nothing
    , parentState :: Maybe Value
parentState = Maybe Value
forall a. Maybe a
Nothing
    , traceContext :: Maybe TraceContext
traceContext = Maybe TraceContext
forall a. Maybe a
Nothing
    , suspended :: Bool
suspended = Bool
False
    , claimedBy :: Maybe UUID
claimedBy = Maybe UUID
forall a. Maybe a
Nothing
    , claimSeq :: Int64
claimSeq = Int64
0
    , archiveFor :: Maybe Int32
archiveFor = Maybe Int32
forall a. Maybe a
Nothing
    , payloadKeys :: ()
payloadKeys = ()
    }

-- | 'defaultJob' with a group key. Jobs sharing a group key are processed serially.
defaultGroupedJob :: Text -> payload -> JobWrite payload
defaultGroupedJob :: forall payload. Text -> payload -> JobWrite payload
defaultGroupedJob Text
key = Maybe Text -> JobWrite payload -> JobWrite payload
forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
key) (JobWrite payload -> JobWrite payload)
-> (payload -> JobWrite payload) -> payload -> JobWrite payload
forall b c a. (b -> c) -> (a -> b) -> a -> c
. payload -> JobWrite payload
forall payload. payload -> JobWrite payload
defaultJob

-- | Replace the payload.
setPayload :: payload' -> JobWrite payload -> JobWrite payload'
setPayload :: forall payload' payload.
payload' -> JobWrite payload -> JobWrite payload'
setPayload payload'
value JobWrite payload
job = JobWrite payload
job {payload = value}

-- | Set the group key.
setGroupKey :: Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey :: forall payload. Maybe Text -> JobWrite payload -> JobWrite payload
setGroupKey Maybe Text
value JobWrite payload
job = JobWrite payload
job {groupKey = value}

-- | Set the claim priority. Lower numbers claim first.
setPriority :: Int32 -> JobWrite payload -> JobWrite payload
setPriority :: forall payload. Int32 -> JobWrite payload -> JobWrite payload
setPriority Int32
value JobWrite payload
job = JobWrite payload
job {priority = value}

-- | Delay the job's first visibility.
setNotVisibleUntil :: Maybe UTCTime -> JobWrite payload -> JobWrite payload
setNotVisibleUntil :: forall payload.
Maybe UTCTime -> JobWrite payload -> JobWrite payload
setNotVisibleUntil Maybe UTCTime
value JobWrite payload
job = JobWrite payload
job {notVisibleUntil = value}

-- | Set the dedup key.
setDedupKey :: Maybe DedupKey -> JobWrite payload -> JobWrite payload
setDedupKey :: forall payload.
Maybe DedupKey -> JobWrite payload -> JobWrite payload
setDedupKey Maybe DedupKey
value JobWrite payload
job = JobWrite payload
job {dedupKey = value}

-- | Override the queue's attempt limit.
setMaxAttempts :: Maybe Int32 -> JobWrite payload -> JobWrite payload
setMaxAttempts :: forall payload. Maybe Int32 -> JobWrite payload -> JobWrite payload
setMaxAttempts Maybe Int32
value JobWrite payload
job = JobWrite payload
job {maxAttempts = value}

-- | Attach a trace context.
setTraceContext :: Maybe TraceContext -> JobWrite payload -> JobWrite payload
setTraceContext :: forall payload.
Maybe TraceContext -> JobWrite payload -> JobWrite payload
setTraceContext Maybe TraceContext
value JobWrite payload
job = JobWrite payload
job {traceContext = value}

-- | Set the archive retention, in seconds.
setArchiveFor :: Maybe Int32 -> JobWrite payload -> JobWrite payload
setArchiveFor :: forall payload. Maybe Int32 -> JobWrite payload -> JobWrite payload
setArchiveFor Maybe Int32
value JobWrite payload
job = JobWrite payload
job {archiveFor = value}

-- | Transform a job's payload without changing its stored metadata.
mapPayload
  :: (payload -> payload')
  -> Job payload key q insertedAt adm
  -> Job payload' key q insertedAt adm
mapPayload :: forall payload payload' key q insertedAt adm.
(payload -> payload')
-> Job payload key q insertedAt adm
-> Job payload' key q insertedAt adm
mapPayload payload -> payload'
transform Job payload key q insertedAt adm
job = Job payload key q insertedAt adm
job {payload = transform (payload job)}

-- | The full payload contract. JSON round-trip for JSONB storage plus the rate-limit
-- and concurrency declarations. Both default to unlimited.
type JobPayload payload =
  (FromJSON payload, ToJSON payload, HasKind payload, HasRateLimit payload, HasConcurrency payload)

-- | The registry declares both admission policy kinds.
type RegistryAdmissionPolicies registry =
  (RegistryConcurrencyPolicies registry, RegistryRateLimitPolicies registry)

-- | A job's primary key.
type JobId = Int64

-- | The token identifying one claim of a job.
type ClaimSeq = Int64

-- | When a claim was taken.
type ClaimTime = UTCTime

-- | The transaction's clock reading.
type CurrentTime = UTCTime

-- | When a handler began.
type StartTime = UTCTime

-- | When a handler finished.
type EndTime = UTCTime

-- | A failure message.
type ErrorMsg = Text

-- | How long to wait before the next attempt.
type BackoffDelay = NominalDiffTime

-- | Callbacks fired at each point of a job's lifecycle, for metrics, logging or tracing.
-- An exception thrown inside one is caught and dropped.
data ObservabilityHooks m payload = ObservabilityHooks
  { forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> UTCTime -> m ()
onJobClaimed
      :: (JobPayload payload)
      => JobRead payload
      -> ClaimTime
      -> m ()
  -- ^ Called immediately after a job is claimed by a worker.
  , forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
onJobSuccess
      :: (JobPayload payload)
      => JobRead payload
      -> StartTime
      -> EndTime
      -> m ()
  -- ^ Called after a job handler succeeds.
  , forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> Text -> UTCTime -> UTCTime -> m ()
onJobFailure
      :: (JobPayload payload)
      => JobRead payload
      -> ErrorMsg
      -> StartTime
      -> EndTime
      -> m ()
  -- ^ Called after a job handler fails and the job was retried or dead-lettered.
  -- A deliberate cancel reports through 'onJobCancelled'.
  , forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> BackoffDelay -> m ()
onJobRetry
      :: (JobPayload payload)
      => JobRead payload
      -> BackoffDelay
      -> m ()
  -- ^ Called when a failed job is successfully scheduled for retry.
  , forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => Text -> JobRead payload -> m ()
onJobFailedAndMovedToDLQ
      :: (JobPayload payload)
      => ErrorMsg
      -> JobRead payload
      -> m ()
  -- ^ Called when a job is successfully moved to the dead-letter queue.
  , forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
onJobCancelled
      :: (JobPayload payload)
      => JobRead payload
      -> ErrorMsg
      -> m ()
  -- ^ Called when a handler cancelled the job's tree or branch and the rows were deleted.
  , forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
onJobUnavailable
      :: (JobPayload payload)
      => JobRead payload
      -> ErrorMsg
      -> m ()
  -- ^ Called when a claimed job went away mid-flight and will not be retried here.
  , forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
onJobHeartbeat
      :: (JobPayload payload)
      => JobRead payload
      -> CurrentTime
      -> StartTime
      -> m ()
  -- ^ Called periodically for a running job.
  }

-- | No-op hooks. Override fields to add observability:
--
-- @
-- myHooks = defaultObservabilityHooks
--   { onJobSuccess = \\job startTime endTime -> do
--       let duration = diffUTCTime endTime startTime
--       logInfo $ "Job " <> show (primaryKey job) <> " took " <> show duration
--   }
-- @
defaultObservabilityHooks :: (Applicative m) => ObservabilityHooks m payload
defaultObservabilityHooks :: forall (m :: * -> *) payload.
Applicative m =>
ObservabilityHooks m payload
defaultObservabilityHooks =
  ObservabilityHooks
    { onJobClaimed :: JobPayload payload => JobRead payload -> UTCTime -> m ()
onJobClaimed = \JobRead payload
_ UTCTime
_ -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
    , onJobSuccess :: JobPayload payload => JobRead payload -> UTCTime -> UTCTime -> m ()
onJobSuccess = \JobRead payload
_ UTCTime
_ UTCTime
_ -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
    , onJobFailure :: JobPayload payload =>
JobRead payload -> Text -> UTCTime -> UTCTime -> m ()
onJobFailure = \JobRead payload
_ Text
_ UTCTime
_ UTCTime
_ -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
    , onJobRetry :: JobPayload payload => JobRead payload -> BackoffDelay -> m ()
onJobRetry = \JobRead payload
_ BackoffDelay
_ -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
    , onJobFailedAndMovedToDLQ :: JobPayload payload => Text -> JobRead payload -> m ()
onJobFailedAndMovedToDLQ = \Text
_ JobRead payload
_ -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
    , onJobCancelled :: JobPayload payload => JobRead payload -> Text -> m ()
onJobCancelled = \JobRead payload
_ Text
_ -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
    , onJobUnavailable :: JobPayload payload => JobRead payload -> Text -> m ()
onJobUnavailable = \JobRead payload
_ Text
_ -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
    , onJobHeartbeat :: JobPayload payload => JobRead payload -> UTCTime -> UTCTime -> m ()
onJobHeartbeat = \JobRead payload
_ UTCTime
_ UTCTime
_ -> () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
    }

-- | Run both hooks at each lifecycle point, left before right. The right one runs
-- however the left ended. When both throw, the right's failure propagates.
instance (MonadUnliftIO m) => Semigroup (ObservabilityHooks m payload) where
  ObservabilityHooks m payload
left <> :: ObservabilityHooks m payload
-> ObservabilityHooks m payload -> ObservabilityHooks m payload
<> ObservabilityHooks m payload
right =
    ObservabilityHooks
      { onJobClaimed :: JobPayload payload => JobRead payload -> UTCTime -> m ()
onJobClaimed = \JobRead payload
job UTCTime
claimTime -> ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> UTCTime -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> UTCTime -> m ()
onJobClaimed ObservabilityHooks m payload
left JobRead payload
job UTCTime
claimTime m () -> m () -> m ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
`andThen` ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> UTCTime -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> UTCTime -> m ()
onJobClaimed ObservabilityHooks m payload
right JobRead payload
job UTCTime
claimTime
      , onJobSuccess :: JobPayload payload => JobRead payload -> UTCTime -> UTCTime -> m ()
onJobSuccess = \JobRead payload
job UTCTime
start UTCTime
end -> ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
onJobSuccess ObservabilityHooks m payload
left JobRead payload
job UTCTime
start UTCTime
end m () -> m () -> m ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
`andThen` ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
onJobSuccess ObservabilityHooks m payload
right JobRead payload
job UTCTime
start UTCTime
end
      , onJobFailure :: JobPayload payload =>
JobRead payload -> Text -> UTCTime -> UTCTime -> m ()
onJobFailure = \JobRead payload
job Text
msg UTCTime
start UTCTime
end -> ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> Text -> UTCTime -> UTCTime -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> Text -> UTCTime -> UTCTime -> m ()
onJobFailure ObservabilityHooks m payload
left JobRead payload
job Text
msg UTCTime
start UTCTime
end m () -> m () -> m ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
`andThen` ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> Text -> UTCTime -> UTCTime -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> Text -> UTCTime -> UTCTime -> m ()
onJobFailure ObservabilityHooks m payload
right JobRead payload
job Text
msg UTCTime
start UTCTime
end
      , onJobRetry :: JobPayload payload => JobRead payload -> BackoffDelay -> m ()
onJobRetry = \JobRead payload
job BackoffDelay
delay -> ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> BackoffDelay -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> BackoffDelay -> m ()
onJobRetry ObservabilityHooks m payload
left JobRead payload
job BackoffDelay
delay m () -> m () -> m ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
`andThen` ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> BackoffDelay -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> BackoffDelay -> m ()
onJobRetry ObservabilityHooks m payload
right JobRead payload
job BackoffDelay
delay
      , onJobFailedAndMovedToDLQ :: JobPayload payload => Text -> JobRead payload -> m ()
onJobFailedAndMovedToDLQ = \Text
msg JobRead payload
job -> ObservabilityHooks m payload
-> JobPayload payload => Text -> JobRead payload -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => Text -> JobRead payload -> m ()
onJobFailedAndMovedToDLQ ObservabilityHooks m payload
left Text
msg JobRead payload
job m () -> m () -> m ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
`andThen` ObservabilityHooks m payload
-> JobPayload payload => Text -> JobRead payload -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => Text -> JobRead payload -> m ()
onJobFailedAndMovedToDLQ ObservabilityHooks m payload
right Text
msg JobRead payload
job
      , onJobCancelled :: JobPayload payload => JobRead payload -> Text -> m ()
onJobCancelled = \JobRead payload
job Text
msg -> ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
onJobCancelled ObservabilityHooks m payload
left JobRead payload
job Text
msg m () -> m () -> m ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
`andThen` ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
onJobCancelled ObservabilityHooks m payload
right JobRead payload
job Text
msg
      , onJobUnavailable :: JobPayload payload => JobRead payload -> Text -> m ()
onJobUnavailable = \JobRead payload
job Text
msg -> ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
onJobUnavailable ObservabilityHooks m payload
left JobRead payload
job Text
msg m () -> m () -> m ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
`andThen` ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
onJobUnavailable ObservabilityHooks m payload
right JobRead payload
job Text
msg
      , onJobHeartbeat :: JobPayload payload => JobRead payload -> UTCTime -> UTCTime -> m ()
onJobHeartbeat = \JobRead payload
job UTCTime
now UTCTime
start -> ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
onJobHeartbeat ObservabilityHooks m payload
left JobRead payload
job UTCTime
now UTCTime
start m () -> m () -> m ()
forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
`andThen` ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload =>
   JobRead payload -> UTCTime -> UTCTime -> m ()
onJobHeartbeat ObservabilityHooks m payload
right JobRead payload
job UTCTime
now UTCTime
start
      }

instance (MonadUnliftIO m) => Monoid (ObservabilityHooks m payload) where
  mempty :: ObservabilityHooks m payload
mempty = ObservabilityHooks m payload
forall (m :: * -> *) payload.
Applicative m =>
ObservabilityHooks m payload
defaultObservabilityHooks

-- | base's @finally@. The second action stays interruptible.
andThen :: (MonadUnliftIO m) => m () -> m () -> m ()
andThen :: forall (m :: * -> *). MonadUnliftIO m => m () -> m () -> m ()
andThen m ()
first m ()
second = ((forall a. m a -> IO a) -> IO ()) -> m ()
forall b. ((forall a. m a -> IO a) -> IO b) -> m b
forall (m :: * -> *) b.
MonadUnliftIO m =>
((forall a. m a -> IO a) -> IO b) -> m b
withRunInIO (((forall a. m a -> IO a) -> IO ()) -> m ())
-> ((forall a. m a -> IO a) -> IO ()) -> m ()
forall a b. (a -> b) -> a -> b
$ \forall a. m a -> IO a
run -> m () -> IO ()
forall a. m a -> IO a
run m ()
first IO () -> IO () -> IO ()
forall a b. IO a -> IO b -> IO a
`E.finally` m () -> IO ()
forall a. m a -> IO a
run m ()
second