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

-- | Job completion, failure, cancellation, and retry settlement.
module Arbiter.Worker.Settlement
  ( ackOrGone
  , reportSuccess
  , batchCallbacks
  , finalizeForceCancelled
  , reportBatchOutcome
  , jobLog
  , batchLog
  ) where

import Arbiter.Core.Exceptions
  ( BranchCancelException (..)
  , JobDeadlineExceeded (..)
  , JobException (..)
  , JobGoneException (..)
  , JobNackException (..)
  , JobPermanentException (..)
  , JobRetryableException (..)
  , ParsingException (..)
  , TreeCancelException (..)
  , namedJobIds
  , throwJobGoneIds
  )
import Arbiter.Core.HighLevel (JobOperation)
import Arbiter.Core.HighLevel qualified as Arb
import Arbiter.Core.Job.Schema (SchemaName)
import Arbiter.Core.Job.Types qualified as Job
import Arbiter.Core.JobResult
import Arbiter.Core.MonadArbiter (MonadArbiter (..))
import Arbiter.Core.Operations qualified as Ops
import Arbiter.Core.Trace
  ( ConsumeShape (..)
  , markSpanError
  , recordJobCancelled
  , recordJobFailure
  )
import Control.Exception (SomeException, fromException, toException)
import Control.Monad (unless, void, when)
import Control.Monad.IO.Class (liftIO)
import Data.Bifunctor (second)
import Data.Bool (bool)
import Data.Either (fromRight)
import Data.Foldable (traverse_)
import Data.Int (Int32, Int64)
import Data.List (partition)
import Data.List.NonEmpty (NonEmpty (..))
import Data.Map.Strict qualified as Map
import Data.Maybe (fromMaybe)
import Data.Set qualified as Set
import Data.Text (Text)
import Data.Text qualified as T
import Data.Time (UTCTime, getCurrentTime)
import Data.UUID (UUID)

import Arbiter.Worker.BackoffStrategy
import Arbiter.Worker.Config
import Arbiter.Worker.Logger
import Arbiter.Worker.Logger.Internal
  ( runHook
  , tryWarn
  , tryWarnWith
  , withJobContext
  , withJobContextList
  , withJobContextOne
  )
import Arbiter.Worker.Results (storeEncodedResult, storeEncodedResults)
import Arbiter.Worker.Settle
  ( CancelHandoff
  , byIdDesc
  , disowned
  , finalized
  , hasIdIn
  , pendingJobs
  , record
  , recordCancelled
  , settle
  , settleBy
  , settleInterruptibly
  , unownedJobs
  )

-- | Add one job to the pool log context.
jobLog :: WorkerConfig m payload -> Job.JobRead payload -> LogConfig
jobLog :: forall (m :: * -> *) payload.
WorkerConfig m payload -> JobRead payload -> LogConfig
jobLog WorkerConfig m payload
config = LogConfig -> JobRead payload -> LogConfig
forall payload. LogConfig -> JobRead payload -> LogConfig
withJobContextOne (WorkerConfig m payload -> LogConfig
forall (m :: * -> *) payload. WorkerConfig m payload -> LogConfig
logConfig WorkerConfig m payload
config)

-- | Add jobs to the pool log context.
jobsLog :: WorkerConfig m payload -> [Job.JobRead payload] -> LogConfig
jobsLog :: forall (m :: * -> *) payload.
WorkerConfig m payload -> [JobRead payload] -> LogConfig
jobsLog WorkerConfig m payload
config = LogConfig -> [JobRead payload] -> LogConfig
forall payload. LogConfig -> [JobRead payload] -> LogConfig
withJobContextList (WorkerConfig m payload -> LogConfig
forall (m :: * -> *) payload. WorkerConfig m payload -> LogConfig
logConfig WorkerConfig m payload
config)

-- | Add a claimed batch to the pool log context.
batchLog :: WorkerConfig m payload -> NonEmpty (Job.JobRead payload) -> LogConfig
batchLog :: forall (m :: * -> *) payload.
WorkerConfig m payload -> NonEmpty (JobRead payload) -> LogConfig
batchLog WorkerConfig m payload
config = LogConfig -> NonEmpty (JobRead payload) -> LogConfig
forall payload.
LogConfig -> NonEmpty (JobRead payload) -> LogConfig
withJobContext (WorkerConfig m payload -> LogConfig
forall (m :: * -> *) payload. WorkerConfig m payload -> LogConfig
logConfig WorkerConfig m payload
config)

unownedReason :: Text
unownedReason :: Text
unownedReason = Text
"no longer claimed by this worker"

-- | Ack a job inside the caller's transaction, throwing if another worker reclaimed it mid-flight.
ackOrGone :: (JobOperation m payload) => Job.JobRead payload -> m ()
ackOrGone :: forall (m :: * -> *) payload.
JobOperation m payload =>
JobRead payload -> m ()
ackOrGone JobRead payload
job = do
  schemaName <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
  rowsAffected <- Ops.ackJobInner schemaName (Job.queueName job) job
  when (rowsAffected == 0) $
    throwJobGoneIds "reclaimed by another worker during processing" [Job.primaryKey job]

-- | Report a completed job, once the handoff already records it.
reportSuccess
  :: (JobOperation m payload)
  => WorkerConfig m payload
  -> UTCTime
  -> Job.JobRead payload
  -> m ()
reportSuccess :: forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload -> UTCTime -> JobRead payload -> m ()
reportSuccess WorkerConfig m payload
config UTCTime
startTime JobRead payload
job = do
  endT <- IO UTCTime -> m UTCTime
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO IO UTCTime
getCurrentTime
  runHook (jobLog config job) "onJobSuccess" $ Job.onJobSuccess (observabilityHooks config) job startTime endT

-- | The settle operations a batch handler drives its jobs through.
batchCallbacks
  :: forall payload m
   . ( EncodeJobResult (ResultOf m payload)
     , JobOperation m payload
     )
  => WorkerConfig m payload
  -> CancelHandoff
  -> NonEmpty (Job.JobRead payload)
  -> UTCTime
  -> Text
  -> BatchCallbacks m payload (ResultOf m payload)
batchCallbacks :: forall payload (m :: * -> *).
(EncodeJobResult (ResultOf m payload), JobOperation m payload) =>
WorkerConfig m payload
-> CancelHandoff
-> NonEmpty (JobRead payload)
-> UTCTime
-> Text
-> BatchCallbacks m payload (ResultOf m payload)
batchCallbacks WorkerConfig m payload
config CancelHandoff
handoff NonEmpty (JobRead payload)
jobs UTCTime
startTime Text
schemaName =
  BatchCallbacks
    { ack :: JobRead payload -> m ()
ack = (JobRead payload -> Maybe Value -> m ()
`ackOneStoring` Maybe Value
forall a. Maybe a
Nothing)
    , ackWith :: JobRead payload
-> SpecResult (MatchIn payload '[] (RegistryOf m)) -> m ()
ackWith = \JobRead payload
job SpecResult (MatchIn payload '[] (RegistryOf m))
result -> JobRead payload -> Maybe Value -> m ()
ackOneStoring JobRead payload
job (SpecResult (MatchIn payload '[] (RegistryOf m)) -> Maybe Value
forall a. EncodeJobResult a => a -> Maybe Value
encodeJobResult SpecResult (MatchIn payload '[] (RegistryOf m))
result)
    , ackAll :: [JobRead payload] -> m ()
ackAll = \[JobRead payload]
toAck -> [(JobRead payload, Maybe Value)] -> m ()
ackBatchStoring ((JobRead payload -> (JobRead payload, Maybe Value))
-> [JobRead payload] -> [(JobRead payload, Maybe Value)]
forall a b. (a -> b) -> [a] -> [b]
map (\JobRead payload
job -> (JobRead payload
job, Maybe Value
forall a. Maybe a
Nothing)) [JobRead payload]
toAck)
    , ackAllWith :: [(JobRead payload,
  SpecResult (MatchIn payload '[] (RegistryOf m)))]
-> m ()
ackAllWith = \[(JobRead payload,
  SpecResult (MatchIn payload '[] (RegistryOf m)))]
resultPairs -> [(JobRead payload, Maybe Value)] -> m ()
ackBatchStoring (((JobRead payload, SpecResult (MatchIn payload '[] (RegistryOf m)))
 -> (JobRead payload, Maybe Value))
-> [(JobRead payload,
     SpecResult (MatchIn payload '[] (RegistryOf m)))]
-> [(JobRead payload, Maybe Value)]
forall a b. (a -> b) -> [a] -> [b]
map ((SpecResult (MatchIn payload '[] (RegistryOf m)) -> Maybe Value)
-> (JobRead payload,
    SpecResult (MatchIn payload '[] (RegistryOf m)))
-> (JobRead payload, Maybe Value)
forall b c a. (b -> c) -> (a, b) -> (a, c)
forall (p :: * -> * -> *) b c a.
Bifunctor p =>
(b -> c) -> p a b -> p a c
second SpecResult (MatchIn payload '[] (RegistryOf m)) -> Maybe Value
forall a. EncodeJobResult a => a -> Maybe Value
encodeJobResult) [(JobRead payload,
  SpecResult (MatchIn payload '[] (RegistryOf m)))]
resultPairs)
    , failRetry :: JobRead payload -> Text -> m ()
failRetry = (Text -> JobException) -> JobRead payload -> Text -> m ()
failAs (JobRetryableException -> JobException
Retryable (JobRetryableException -> JobException)
-> (Text -> JobRetryableException) -> Text -> JobException
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> JobRetryableException
JobRetryableException)
    , failPermanent :: JobRead payload -> Text -> m ()
failPermanent = (Text -> JobException) -> JobRead payload -> Text -> m ()
failAs (JobPermanentException -> JobException
Permanent (JobPermanentException -> JobException)
-> (Text -> JobPermanentException) -> Text -> JobException
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> JobPermanentException
JobPermanentException)
    , cancelBranch :: JobRead payload -> Text -> m ()
cancelBranch = (Text -> JobException) -> JobRead payload -> Text -> m ()
failAs (BranchCancelException -> JobException
BranchCancel (BranchCancelException -> JobException)
-> (Text -> BranchCancelException) -> Text -> JobException
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> BranchCancelException
BranchCancelException)
    , cancelTree :: JobRead payload -> Text -> m ()
cancelTree = (Text -> JobException) -> JobRead payload -> Text -> m ()
failAs (TreeCancelException -> JobException
TreeCancel (TreeCancelException -> JobException)
-> (Text -> TreeCancelException) -> Text -> JobException
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> TreeCancelException
TreeCancelException)
    , nack :: JobRead payload -> m ()
nack = JobRead payload -> m ()
nackOne
    }
  where
    (JobRead payload
firstJob :| [JobRead payload]
_) = NonEmpty (JobRead payload)
jobs
    shape :: ConsumeShape
shape = NonEmpty (JobRead payload) -> ConsumeShape
forall payload. NonEmpty (JobRead payload) -> ConsumeShape
batchSpanShape NonEmpty (JobRead payload)
jobs
    nackOne :: JobRead payload -> m ()
nackOne JobRead payload
job = WorkerConfig m payload
-> CancelHandoff -> [JobRead payload] -> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> [JobRead payload] -> m ()
releaseJobs WorkerConfig m payload
config CancelHandoff
handoff [JobRead payload
job]
    failAs :: (Text -> JobException) -> JobRead payload -> Text -> m ()
failAs Text -> JobException
mkExc JobRead payload
job Text
msg = JobRead payload -> SomeException -> m ()
failWith JobRead payload
job (JobException -> SomeException
forall e. Exception e => e -> SomeException
toException (Text -> JobException
mkExc Text
msg))
    failWith :: JobRead payload -> SomeException -> m ()
failWith JobRead payload
job SomeException
exc = do
      endT <- IO UTCTime -> m UTCTime
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO IO UTCTime
getCurrentTime
      settle
        handoff
        (finalized [job])
        ( withDbTransaction $
            handleJobFailure config Ops.TakeLocks shape (classifyException exc) startTime endT job
        )
        $ \FailureOutcome m
outcome -> do
          FailureOutcome m -> m ()
forall (m :: * -> *). Monad m => FailureOutcome m -> m ()
reportWritten FailureOutcome m
outcome
          WorkerConfig m payload
-> CancelHandoff -> [(JobRead payload, FailureOutcome m)] -> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> [(JobRead payload, FailureOutcome m)] -> m ()
settleUnwritten WorkerConfig m payload
config CancelHandoff
handoff [(JobRead payload
job, FailureOutcome m
outcome)]
    ackOneStoring :: JobRead payload -> Maybe Value -> m ()
ackOneStoring JobRead payload
job Maybe Value
mVal =
      CancelHandoff -> Settled payload -> m () -> (() -> m ()) -> m ()
forall (m :: * -> *) payload a b.
MonadUnliftIO m =>
CancelHandoff -> Settled payload -> m a -> (a -> m b) -> m b
settle
        CancelHandoff
handoff
        ([JobRead payload] -> Settled payload
forall payload. [JobRead payload] -> Settled payload
finalized [JobRead payload
job])
        ( m () -> m ()
forall a. m a -> m a
forall (m :: * -> *) a. MonadArbiter m => m a -> m a
withDbTransaction (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$ do
            JobRead payload -> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
JobRead payload -> m ()
ackOrGone JobRead payload
job
            Text -> JobRead payload -> Maybe Value -> m ()
forall (m :: * -> *) payload.
MonadArbiter m =>
Text -> JobRead payload -> Maybe Value -> m ()
storeEncodedResult Text
schemaName JobRead payload
job Maybe Value
mVal
        )
        (m () -> () -> m ()
forall a b. a -> b -> a
const (WorkerConfig m payload -> UTCTime -> JobRead payload -> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload -> UTCTime -> JobRead payload -> m ()
reportSuccess WorkerConfig m payload
config UTCTime
startTime JobRead payload
job))
    ackBatchStoring :: [(JobRead payload, Maybe Value)] -> m ()
ackBatchStoring [(JobRead payload, Maybe Value)]
pairs =
      let jobsToAck :: [JobRead payload]
jobsToAck = ((JobRead payload, Maybe Value) -> JobRead payload)
-> [(JobRead payload, Maybe Value)] -> [JobRead payload]
forall a b. (a -> b) -> [a] -> [b]
map (JobRead payload, Maybe Value) -> JobRead payload
forall a b. (a, b) -> a
fst [(JobRead payload, Maybe Value)]
pairs
       in CancelHandoff
-> m ([JobRead payload], [JobRead payload])
-> (([JobRead payload], [JobRead payload]) -> Settled payload)
-> (([JobRead payload], [JobRead payload]) -> m ())
-> m ()
forall (m :: * -> *) a payload b.
MonadUnliftIO m =>
CancelHandoff -> m a -> (a -> Settled payload) -> (a -> m b) -> m b
settleBy
            CancelHandoff
handoff
            ( m ([JobRead payload], [JobRead payload])
-> m ([JobRead payload], [JobRead payload])
forall a. m a -> m a
forall (m :: * -> *) a. MonadArbiter m => m a -> m a
withDbTransaction (m ([JobRead payload], [JobRead payload])
 -> m ([JobRead payload], [JobRead payload]))
-> m ([JobRead payload], [JobRead payload])
-> m ([JobRead payload], [JobRead payload])
forall a b. (a -> b) -> a -> b
$ do
                acked <- [Int64] -> Set Int64
forall a. Ord a => [a] -> Set a
Set.fromList ([Int64] -> Set Int64) -> m [Int64] -> m (Set Int64)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Text -> Text -> [JobRead payload] -> m [Int64]
forall (m :: * -> *) payload.
MonadArbiter m =>
Text -> Text -> [JobRead payload] -> m [Int64]
Ops.ackJobsBatchInner Text
schemaName (JobRead payload -> Text
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> q
Job.queueName JobRead payload
firstJob) [JobRead payload]
jobsToAck
                storeEncodedResults schemaName (filter (hasIdIn acked . fst) pairs)
                pure (partition (hasIdIn acked) jobsToAck)
            )
            (\([JobRead payload]
done, [JobRead payload]
reclaimed) -> [JobRead payload] -> Settled payload
forall payload. [JobRead payload] -> Settled payload
finalized [JobRead payload]
done Settled payload -> Settled payload -> Settled payload
forall a. Semigroup a => a -> a -> a
<> [JobRead payload] -> Settled payload
forall payload. [JobRead payload] -> Settled payload
disowned [JobRead payload]
reclaimed)
            ((([JobRead payload], [JobRead payload]) -> m ()) -> m ())
-> (([JobRead payload], [JobRead payload]) -> m ()) -> m ()
forall a b. (a -> b) -> a -> b
$ \([JobRead payload]
done, [JobRead payload]
reclaimed) -> do
              (JobRead payload -> m ()) -> [JobRead payload] -> m ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (WorkerConfig m payload -> UTCTime -> JobRead payload -> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload -> UTCTime -> JobRead payload -> m ()
reportSuccess WorkerConfig m payload
config UTCTime
startTime) [JobRead payload]
done
              let reclaimedLog :: LogConfig
reclaimedLog = WorkerConfig m payload -> [JobRead payload] -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> [JobRead payload] -> LogConfig
jobsLog WorkerConfig m payload
config [JobRead payload]
reclaimed
              Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless ([JobRead payload] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null [JobRead payload]
reclaimed) (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$ do
                LogConfig -> LogLevel -> Text -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> LogLevel -> Text -> m ()
tryLog LogConfig
reclaimedLog LogLevel
Info Text
"Jobs no longer claimed by this worker during bulk completion, skipped"
                m (Set Int64) -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (m (Set Int64) -> m ()) -> m (Set Int64) -> m ()
forall a b. (a -> b) -> a -> b
$ WorkerConfig m payload
-> CancelHandoff -> Text -> [JobRead payload] -> m (Set Int64)
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> Text -> [JobRead payload] -> m (Set Int64)
settleGoneJobs WorkerConfig m payload
config CancelHandoff
handoff Text
unownedReason [JobRead payload]
reclaimed

-- | Delete the jobs a force-cancel flagged, report them, and hand back the attempt
-- the claim consumed for the batch siblings it interrupted.
finalizeForceCancelled
  :: (JobOperation m payload)
  => WorkerConfig m payload
  -> NonEmpty (Job.JobRead payload)
  -> [Int64]
  -- ^ The jobs the cancel named.
  -> [Int64]
  -- ^ Jobs that the same signal found unavailable. Report these jobs without a nack.
  -> CancelHandoff
  -> m ()
finalizeForceCancelled :: forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> NonEmpty (JobRead payload)
-> [Int64]
-> [Int64]
-> CancelHandoff
-> m ()
finalizeForceCancelled WorkerConfig m payload
config NonEmpty (JobRead payload)
jobs [Int64]
cancelledIds [Int64]
goneIds CancelHandoff
handoff = do
  LogConfig -> LogLevel -> Text -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> LogLevel -> Text -> m ()
tryLog (WorkerConfig m payload -> NonEmpty (JobRead payload) -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> NonEmpty (JobRead payload) -> LogConfig
batchLog WorkerConfig m payload
config NonEmpty (JobRead payload)
jobs) LogLevel
Info Text
"Job(s) force-cancelled"
  schemaName <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
  pending <- pendingJobs handoff jobs
  -- The cancel can name a job the handler finalized after the cancel took effect.
  let settling = (JobRead payload -> Bool)
-> NonEmpty (JobRead payload) -> [JobRead payload]
forall payload.
(JobRead payload -> Bool)
-> NonEmpty (JobRead payload) -> [JobRead payload]
byIdDesc (Set Int64 -> JobRead payload -> Bool
forall payload. Set Int64 -> JobRead payload -> Bool
hasIdIn ([Int64] -> Set Int64
forall a. Ord a => [a] -> Set a
Set.fromList ([Int64]
cancelledIds [Int64] -> [Int64] -> [Int64]
forall a. Semigroup a => a -> a -> a
<> (JobRead payload -> Int64) -> [JobRead payload] -> [Int64]
forall a b. (a -> b) -> [a] -> [b]
map JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
Job.primaryKey [JobRead payload]
pending))) NonEmpty (JobRead payload)
jobs
      goneSet = [Int64] -> Set Int64
forall a. Ord a => [a] -> Set a
Set.fromList [Int64]
goneIds
      deletable = [JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
Job.primaryKey JobRead payload
job | JobRead payload
job <- [JobRead payload]
settling, Bool -> Bool
not (Set Int64 -> JobRead payload -> Bool
forall payload. Set Int64 -> JobRead payload -> Bool
hasIdIn Set Int64
goneSet JobRead payload
job)]
  deleted <-
    deleteCancelledOrWarn (batchLog config jobs) (workerId config) schemaName (Job.queueName firstJob) deletable
  (fresh, cancelled) <- recordCancelled handoff (deleted <> Set.fromList cancelledIds)
  let (gone, interrupted) = partition (hasIdIn cancelled) settling
      (unreported, alreadyReported) = partition (hasIdIn fresh) gone
      (unavailable, siblings) = partition (hasIdIn goneSet) interrupted
      unavailableLog = WorkerConfig m payload -> [JobRead payload] -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> [JobRead payload] -> LogConfig
jobsLog WorkerConfig m payload
config [JobRead payload]
unavailable
  reportGoneJobs config handoff cancelled "force-cancelled" unreported
  record handoff (finalized alreadyReported)
  unless (null unavailable) $ do
    tryLog unavailableLog Info "Job(s) no longer claimed by this worker, skipping retry"
    reportGoneJobs config handoff mempty unownedReason unavailable
  releaseOrWarn config handoff "Releasing a force-cancel batch sibling failed" siblings
  where
    (JobRead payload
firstJob :| [JobRead payload]
_) = NonEmpty (JobRead payload)
jobs

-- | Hand back the attempt the claim consumed for the jobs left unfinalized, in one
-- statement, and report whichever of them the nack found under another claim.
releaseJobs
  :: (JobOperation m payload)
  => WorkerConfig m payload
  -> CancelHandoff
  -> [Job.JobRead payload]
  -> m ()
releaseJobs :: forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> [JobRead payload] -> m ()
releaseJobs WorkerConfig m payload
_ CancelHandoff
_ [] = () -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()
releaseJobs WorkerConfig m payload
config CancelHandoff
handoff [JobRead payload]
jobs =
  CancelHandoff
-> Settled payload -> m (Set Int64) -> (Set Int64 -> m ()) -> m ()
forall (m :: * -> *) payload a b.
MonadUnliftIO m =>
CancelHandoff -> Settled payload -> m a -> (a -> m b) -> m b
settle
    CancelHandoff
handoff
    ([JobRead payload] -> Settled payload
forall payload. [JobRead payload] -> Settled payload
finalized [JobRead payload]
jobs)
    ([Int64] -> Set Int64
forall a. Ord a => [a] -> Set a
Set.fromList ([Int64] -> Set Int64) -> m [Int64] -> m (Set Int64)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> [JobRead payload] -> m [Int64]
forall payload (m :: * -> *).
JobOperation m payload =>
[JobRead payload] -> m [Int64]
Arb.nackJobsBatch [JobRead payload]
jobs)
    ((Set Int64 -> m ()) -> m ()) -> (Set Int64 -> m ()) -> m ()
forall a b. (a -> b) -> a -> b
$ \Set Int64
released ->
      WorkerConfig m payload
-> CancelHandoff -> [(JobRead payload, FailureOutcome m)] -> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> [(JobRead payload, FailureOutcome m)] -> m ()
settleUnwritten WorkerConfig m payload
config CancelHandoff
handoff [(JobRead payload
job, Text -> FailureOutcome m
forall a b. a -> Either a b
Left Text
unownedReason) | JobRead payload
job <- [JobRead payload]
jobs, Bool -> Bool
not (Set Int64 -> JobRead payload -> Bool
forall payload. Set Int64 -> JobRead payload -> Bool
hasIdIn Set Int64
released JobRead payload
job)]

-- | 'releaseJobs' on an unwinding path. A failure is logged as a warning.
releaseOrWarn
  :: (JobOperation m payload)
  => WorkerConfig m payload
  -> CancelHandoff
  -> Text
  -> [Job.JobRead payload]
  -> m ()
releaseOrWarn :: forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> Text -> [JobRead payload] -> m ()
releaseOrWarn WorkerConfig m payload
config CancelHandoff
handoff Text
warning [JobRead payload]
jobs =
  LogConfig -> Text -> m () -> m ()
forall (m :: * -> *) a.
MonadUnliftIO m =>
LogConfig -> Text -> m a -> m ()
tryWarn (WorkerConfig m payload -> [JobRead payload] -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> [JobRead payload] -> LogConfig
jobsLog WorkerConfig m payload
config [JobRead payload]
jobs) Text
warning (WorkerConfig m payload
-> CancelHandoff -> [JobRead payload] -> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> [JobRead payload] -> m ()
releaseJobs WorkerConfig m payload
config CancelHandoff
handoff [JobRead payload]
jobs)

-- | Settle the jobs a failure could not be written for, each under its own reason.
settleUnwritten
  :: (JobOperation m payload)
  => WorkerConfig m payload
  -> CancelHandoff
  -> [(Job.JobRead payload, FailureOutcome m)]
  -> m ()
settleUnwritten :: forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> [(JobRead payload, FailureOutcome m)] -> m ()
settleUnwritten WorkerConfig m payload
config CancelHandoff
handoff [(JobRead payload, FailureOutcome m)]
pairs =
  ((Text, [JobRead payload]) -> m ())
-> [(Text, [JobRead payload])] -> m ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (Text, [JobRead payload]) -> m ()
settleOne (Map Text [JobRead payload] -> [(Text, [JobRead payload])]
forall k a. Map k a -> [(k, a)]
Map.toList (([JobRead payload] -> [JobRead payload] -> [JobRead payload])
-> [(Text, [JobRead payload])] -> Map Text [JobRead payload]
forall k a. Ord k => (a -> a -> a) -> [(k, a)] -> Map k a
Map.fromListWith [JobRead payload] -> [JobRead payload] -> [JobRead payload]
forall a. Semigroup a => a -> a -> a
(<>) [(Text
reason, [JobRead payload
job]) | (JobRead payload
job, Left Text
reason) <- [(JobRead payload, FailureOutcome m)]
pairs]))
  where
    settleOne :: (Text, [JobRead payload]) -> m ()
settleOne (Text
reason, [JobRead payload]
jobs) = do
      cancelled <- WorkerConfig m payload
-> CancelHandoff -> Text -> [JobRead payload] -> m (Set Int64)
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> Text -> [JobRead payload] -> m (Set Int64)
settleGoneJobs WorkerConfig m payload
config CancelHandoff
handoff Text
reason [JobRead payload]
jobs
      let (forced, unavailable) = partition (hasIdIn cancelled) jobs
      unless (null forced) $ tryLog (jobsLog config forced) Info "Job(s) force-cancelled"
      unless (null unavailable) $ tryLog (jobsLog config unavailable) Warning ("Job(s) " <> reason)

-- | Delete whichever of @jobs@ a force-cancel flagged, report them all, return those ids.
-- The delete is recorded in the handoff.
settleGoneJobs
  :: (JobOperation m payload)
  => WorkerConfig m payload
  -> CancelHandoff
  -> Text
  -> [Job.JobRead payload]
  -> m (Set.Set Int64)
settleGoneJobs :: forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> Text -> [JobRead payload] -> m (Set Int64)
settleGoneJobs WorkerConfig m payload
config CancelHandoff
handoff Text
reason = \case
  [] -> Set Int64 -> m (Set Int64)
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Set Int64
forall a. Monoid a => a
mempty
  jobs :: [JobRead payload]
jobs@(JobRead payload
firstJob : [JobRead payload]
_) -> do
    schemaName <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
    let logCfg = WorkerConfig m payload -> [JobRead payload] -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> [JobRead payload] -> LogConfig
jobsLog WorkerConfig m payload
config [JobRead payload]
jobs
    cancelled <-
      deleteCancelledOrWarn logCfg (workerId config) schemaName (Job.queueName firstJob) (map Job.primaryKey jobs)
    void $ recordCancelled handoff cancelled
    cancelled <$ reportGoneJobs config handoff cancelled reason jobs

-- | Delete whichever of @jobIds@ a force-cancel flagged against this worker's lease,
-- or against none, returning the ids it deleted. One held elsewhere is that worker's.
deleteCancelledOrWarn
  :: (MonadArbiter m)
  => LogConfig
  -> UUID
  -> SchemaName
  -> Text
  -- ^ Queue the jobs belong to.
  -> [Int64]
  -> m (Set.Set Int64)
deleteCancelledOrWarn :: forall (m :: * -> *).
MonadArbiter m =>
LogConfig -> UUID -> Text -> Text -> [Int64] -> m (Set Int64)
deleteCancelledOrWarn LogConfig
logCfg UUID
owner Text
schemaName Text
queue [Int64]
jobIds =
  LogConfig -> Text -> Set Int64 -> m (Set Int64) -> m (Set Int64)
forall (m :: * -> *) a.
MonadUnliftIO m =>
LogConfig -> Text -> a -> m a -> m a
tryWarnWith LogConfig
logCfg Text
"Deleting force-cancelled jobs failed" Set Int64
forall a. Monoid a => a
mempty (m (Set Int64) -> m (Set Int64)) -> m (Set Int64) -> m (Set Int64)
forall a b. (a -> b) -> a -> b
$
    [Int64] -> Set Int64
forall a. Ord a => [a] -> Set a
Set.fromList ([Int64] -> Set Int64) -> m [Int64] -> m (Set Int64)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Text -> Text -> Maybe UUID -> [Int64] -> m [Int64]
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Maybe UUID -> [Int64] -> m [Int64]
Ops.deleteCancelledJobs Text
schemaName Text
queue (UUID -> Maybe UUID
forall a. a -> Maybe a
Just UUID
owner) [Int64]
jobIds

-- | Interpret a finished batch. Warn when the handler left jobs unfinalized. Skip
-- retry for gone or nacked jobs. Fail the rest.
reportBatchOutcome
  :: forall payload m
   . (JobOperation m payload)
  => WorkerConfig m payload
  -> UTCTime
  -> UTCTime
  -> NonEmpty (Job.JobRead payload)
  -> CancelHandoff
  -> Either SomeException ()
  -> m ()
reportBatchOutcome :: forall payload (m :: * -> *).
JobOperation m payload =>
WorkerConfig m payload
-> UTCTime
-> UTCTime
-> NonEmpty (JobRead payload)
-> CancelHandoff
-> Either SomeException ()
-> m ()
reportBatchOutcome WorkerConfig m payload
config UTCTime
startTime UTCTime
endTime NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
jobs CancelHandoff
handoff Either SomeException ()
outcome = do
  unhandled <- CancelHandoff
-> NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall (m :: * -> *) payload.
MonadIO m =>
CancelHandoff -> NonEmpty (JobRead payload) -> m [JobRead payload]
pendingJobs CancelHandoff
handoff NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
jobs
  let splitNamed [Int64]
ids = (JobRecord payload Int64 Text UTCTime PayloadKeys -> Bool)
-> [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> ([JobRecord payload Int64 Text UTCTime PayloadKeys],
    [JobRecord payload Int64 Text UTCTime PayloadKeys])
forall a. (a -> Bool) -> [a] -> ([a], [a])
partition (Set Int64
-> JobRecord payload Int64 Text UTCTime PayloadKeys -> Bool
forall payload. Set Int64 -> JobRead payload -> Bool
hasIdIn ([Int64] -> Set Int64
forall a. Ord a => [a] -> Set a
Set.fromList [Int64]
ids)) [JobRecord payload Int64 Text UTCTime PayloadKeys]
unhandled
      reportUnavailable [Int64]
ids Text
reason = do
        -- An exception naming no job speaks for none of them. The remainder keeps
        -- its attempt.
        let ([JobRecord payload Int64 Text UTCTime PayloadKeys]
jobsGone, [JobRecord payload Int64 Text UTCTime PayloadKeys]
siblings) = [Int64]
-> ([JobRecord payload Int64 Text UTCTime PayloadKeys],
    [JobRecord payload Int64 Text UTCTime PayloadKeys])
splitNamed [Int64]
ids
        m (Set Int64) -> m ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (m (Set Int64) -> m ()) -> m (Set Int64) -> m ()
forall a b. (a -> b) -> a -> b
$ WorkerConfig m payload
-> CancelHandoff
-> Text
-> [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> m (Set Int64)
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> Text -> [JobRead payload] -> m (Set Int64)
settleGoneJobs WorkerConfig m payload
config CancelHandoff
handoff Text
reason [JobRecord payload Int64 Text UTCTime PayloadKeys]
jobsGone
        Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless ([Int64] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null [Int64]
ids) (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
          WorkerConfig m payload
-> CancelHandoff
-> Text
-> [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> Text -> [JobRead payload] -> m ()
releaseOrWarn WorkerConfig m payload
config CancelHandoff
handoff Text
"Releasing an interrupted batch sibling failed" [JobRecord payload Int64 Text UTCTime PayloadKeys]
siblings
      unownedOf = CancelHandoff
-> NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
-> m [JobRecord payload Int64 Text UTCTime PayloadKeys]
forall (m :: * -> *) payload.
MonadIO m =>
CancelHandoff -> NonEmpty (JobRead payload) -> m [JobRead payload]
unownedJobs CancelHandoff
handoff NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
jobs
  case outcome of
    Right () ->
      Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless ([JobRecord payload Int64 Text UTCTime PayloadKeys] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null [JobRecord payload Int64 Text UTCTime PayloadKeys]
unhandled) (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
        LogConfig -> LogLevel -> Text -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> LogLevel -> Text -> m ()
tryLog (WorkerConfig m payload
-> [JobRecord payload Int64 Text UTCTime PayloadKeys] -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> [JobRead payload] -> LogConfig
jobsLog WorkerConfig m payload
config [JobRecord payload Int64 Text UTCTime PayloadKeys]
unhandled) LogLevel
Warning Text
"Handler left jobs unfinalized, will reprocess"
    Left SomeException
exc
      | Just (JobGoneException Text
reason [Int64]
gone) <- SomeException -> Maybe JobGoneException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exc -> do
          LogConfig -> LogLevel -> Text -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> LogLevel -> Text -> m ()
tryLog (WorkerConfig m payload
-> NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
-> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> NonEmpty (JobRead payload) -> LogConfig
batchLog WorkerConfig m payload
config NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
jobs) LogLevel
Info (Text -> m ()) -> Text -> m ()
forall a b. (a -> b) -> a -> b
$ Text
"Job(s) " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
reason Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
", skipping retry" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> [Int64] -> Text
namedJobIds [Int64]
gone
          [Int64] -> Text -> m ()
reportUnavailable [Int64]
gone Text
reason
      | Just JobNackException
JobNackException <- SomeException -> Maybe JobNackException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exc -> do
          -- Hand back the attempt the claim consumed for every job the handler
          -- left unfinalized.
          WorkerConfig m payload
-> CancelHandoff
-> Text
-> [JobRecord payload Int64 Text UTCTime PayloadKeys]
-> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> Text -> [JobRead payload] -> m ()
releaseOrWarn WorkerConfig m payload
config CancelHandoff
handoff Text
"Handing back a nacked batch's attempt failed" [JobRecord payload Int64 Text UTCTime PayloadKeys]
unhandled
          LogConfig -> LogLevel -> Text -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> LogLevel -> Text -> m ()
tryLog (WorkerConfig m payload
-> NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
-> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> NonEmpty (JobRead payload) -> LogConfig
batchLog WorkerConfig m payload
config NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
jobs) LogLevel
Info Text
"Job(s) nacked, will be reprocessed"
      | Bool
otherwise -> do
          -- Fail the jobs the handler did not finalize, in a separate transaction.
          let failure :: (Text, FailureKind)
failure@(Text
reason, FailureKind
kind) = SomeException -> (Text, FailureKind)
classifyException SomeException
exc
              queue :: Text
queue = JobRecord payload Int64 Text UTCTime PayloadKeys -> Text
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> q
Job.queueName JobRecord payload Int64 Text UTCTime PayloadKeys
firstJob
          Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when ([JobRecord payload Int64 Text UTCTime PayloadKeys] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null [JobRecord payload Int64 Text UTCTime PayloadKeys]
unhandled) (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
            LogConfig -> LogLevel -> Text -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> LogLevel -> Text -> m ()
tryLog (WorkerConfig m payload
-> NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
-> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> NonEmpty (JobRead payload) -> LogConfig
batchLog WorkerConfig m payload
config NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
jobs) LogLevel
Warning (Text
"Handler stopped with its jobs already finalized: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
reason)
          schemaName <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
          -- A tree or branch cancel acts on the whole tree.
          unowned <- if cancelsTree kind then unownedOf else pure []
          -- Lock every row the settles below touch, in one pass. A cancel's pass
          -- covers the whole tree.
          let lockTrees
                | FailureKind -> Bool
cancelsTree FailureKind
kind = Text -> Text -> [Int64] -> m ()
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> [Int64] -> m ()
Ops.lockJobTreesFromRoot
                | Bool
otherwise = Text -> Text -> [Int64] -> m ()
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> [Int64] -> m ()
Ops.lockJobTrees
          settle
            handoff
            (finalized unhandled)
            ( withDbTransaction $ do
                Ops.lockJobParents schemaName queue (map Job.parentId unhandled)
                lockTrees schemaName queue (map Job.primaryKey (unhandled <> unowned))
                outcomes <- traverse (\JobRecord payload Int64 Text UTCTime PayloadKeys
job -> (JobRecord payload Int64 Text UTCTime PayloadKeys
job,) (FailureOutcome m
 -> (JobRecord payload Int64 Text UTCTime PayloadKeys,
     FailureOutcome m))
-> m (FailureOutcome m)
-> m (JobRecord payload Int64 Text UTCTime PayloadKeys,
      FailureOutcome m)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> (Text, FailureKind)
-> JobRecord payload Int64 Text UTCTime PayloadKeys
-> m (FailureOutcome m)
failJob (Text, FailureKind)
failure JobRecord payload Int64 Text UTCTime PayloadKeys
job) unhandled
                traverse_ (void . cancelJobFor kind) unowned
                pure outcomes
            )
            $ \[(JobRecord payload Int64 Text UTCTime PayloadKeys,
  FailureOutcome m)]
outcomes -> do
              ((JobRecord payload Int64 Text UTCTime PayloadKeys,
  FailureOutcome m)
 -> m ())
-> [(JobRecord payload Int64 Text UTCTime PayloadKeys,
     FailureOutcome m)]
-> m ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (JobRecord payload Int64 Text UTCTime PayloadKeys,
 FailureOutcome m)
-> m ()
report [(JobRecord payload Int64 Text UTCTime PayloadKeys,
  FailureOutcome m)]
outcomes
              WorkerConfig m payload
-> CancelHandoff
-> [(JobRecord payload Int64 Text UTCTime PayloadKeys,
     FailureOutcome m)]
-> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> [(JobRead payload, FailureOutcome m)] -> m ()
settleUnwritten WorkerConfig m payload
config CancelHandoff
handoff [(JobRecord payload Int64 Text UTCTime PayloadKeys,
  FailureOutcome m)]
outcomes
  where
    (JobRecord payload Int64 Text UTCTime PayloadKeys
firstJob :| [JobRecord payload Int64 Text UTCTime PayloadKeys]
_) = NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
jobs
    shape :: ConsumeShape
shape = NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
-> ConsumeShape
forall payload. NonEmpty (JobRead payload) -> ConsumeShape
batchSpanShape NonEmpty (JobRecord payload Int64 Text UTCTime PayloadKeys)
jobs
    report :: (JobRecord payload Int64 Text UTCTime PayloadKeys,
 FailureOutcome m)
-> m ()
report (JobRecord payload Int64 Text UTCTime PayloadKeys
job, FailureOutcome m
written) =
      LogConfig -> Text -> m () -> m ()
forall (m :: * -> *) a.
MonadUnliftIO m =>
LogConfig -> Text -> m a -> m ()
tryWarn (WorkerConfig m payload
-> JobRecord payload Int64 Text UTCTime PayloadKeys -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> JobRead payload -> LogConfig
jobLog WorkerConfig m payload
config JobRecord payload Int64 Text UTCTime PayloadKeys
job) Text
"Reporting a job's failure failed" (FailureOutcome m -> m ()
forall (m :: * -> *). Monad m => FailureOutcome m -> m ()
reportWritten FailureOutcome m
written)
    failJob :: (Text, FailureKind)
-> JobRecord payload Int64 Text UTCTime PayloadKeys
-> m (FailureOutcome m)
failJob (Text, FailureKind)
failure JobRecord payload Int64 Text UTCTime PayloadKeys
job =
      WorkerConfig m payload
-> TreeLocks
-> ConsumeShape
-> (Text, FailureKind)
-> UTCTime
-> UTCTime
-> JobRecord payload Int64 Text UTCTime PayloadKeys
-> m (FailureOutcome m)
forall payload (m :: * -> *).
JobOperation m payload =>
WorkerConfig m payload
-> TreeLocks
-> ConsumeShape
-> (Text, FailureKind)
-> UTCTime
-> UTCTime
-> JobRead payload
-> m (FailureOutcome m)
handleJobFailure WorkerConfig m payload
config TreeLocks
Ops.LocksHeld ConsumeShape
shape (Text, FailureKind)
failure UTCTime
startTime UTCTime
endTime JobRecord payload Int64 Text UTCTime PayloadKeys
job

data FailureKind = RetryFailure | PermanentFailure | TreeCancelFailure | BranchCancelFailure
  deriving stock (FailureKind -> FailureKind -> Bool
(FailureKind -> FailureKind -> Bool)
-> (FailureKind -> FailureKind -> Bool) -> Eq FailureKind
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: FailureKind -> FailureKind -> Bool
== :: FailureKind -> FailureKind -> Bool
$c/= :: FailureKind -> FailureKind -> Bool
/= :: FailureKind -> FailureKind -> Bool
Eq)

-- | The job's own attempt budget, or the default.
jobMaxAtts :: Job.JobRead payload -> Int32
jobMaxAtts :: forall payload. JobRead payload -> Int32
jobMaxAtts JobRead payload
job = Int32 -> Maybe Int32 -> Int32
forall a. a -> Maybe a -> a
fromMaybe Int32
Job.defaultMaxAttempts (JobRead payload -> Maybe Int32
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> Maybe Int32
Job.maxAttempts JobRead payload
job)

-- | Classify a handler exception into an error message and failure disposition.
-- 'reportBatchOutcome' intercepts 'JobGoneException' before this.
classifyException :: SomeException -> (T.Text, FailureKind)
classifyException :: SomeException -> (Text, FailureKind)
classifyException SomeException
exception
  | Just (Retryable (JobRetryableException Text
msg)) <- SomeException -> Maybe JobException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = (Text
msg, FailureKind
RetryFailure)
  | Just (Permanent (JobPermanentException Text
msg)) <- SomeException -> Maybe JobException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = (Text
msg, FailureKind
PermanentFailure)
  | Just (TreeCancel (TreeCancelException Text
msg)) <- SomeException -> Maybe JobException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = (Text
msg, FailureKind
TreeCancelFailure)
  | Just (BranchCancel (BranchCancelException Text
msg)) <- SomeException -> Maybe JobException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = (Text
msg, FailureKind
BranchCancelFailure)
  | Just (ParsingException Text
msg) <- SomeException -> Maybe ParsingException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = (Text
msg, FailureKind
PermanentFailure)
  | Just (JobDeadlineExceeded Text
msg) <- SomeException -> Maybe JobDeadlineExceeded
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = (Text
msg, FailureKind
RetryFailure)
  | Bool
otherwise = (String -> Text
T.pack (String -> Text) -> String -> Text
forall a b. (a -> b) -> a -> b
$ SomeException -> String
forall a. Show a => a -> String
show SomeException
exception, FailureKind
RetryFailure) -- Unknown exception, treat as retryable

-- | Whether a failure deletes a job tree or updates the job claim.
cancelsTree :: FailureKind -> Bool
cancelsTree :: FailureKind -> Bool
cancelsTree FailureKind
kind = FailureKind
kind FailureKind -> [FailureKind] -> Bool
forall a. Eq a => a -> [a] -> Bool
forall (t :: * -> *) a. (Foldable t, Eq a) => a -> t a -> Bool
`elem` [FailureKind
TreeCancelFailure, FailureKind
BranchCancelFailure]

-- | Report jobs this worker can no longer act on. A job a force-cancel deleted is
-- reported as cancelled. Each is recorded against the handoff.
reportGoneJobs
  :: (JobOperation m payload)
  => WorkerConfig m payload
  -> CancelHandoff
  -> Set.Set Int64
  -- ^ Ids a force-cancel accounted for.
  -> Text
  -> [Job.JobRead payload]
  -> m ()
reportGoneJobs :: forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> CancelHandoff -> Set Int64 -> Text -> [JobRead payload] -> m ()
reportGoneJobs WorkerConfig m payload
config CancelHandoff
handoff Set Int64
cancelled Text
reason [JobRead payload]
jobs =
  CancelHandoff -> Settled payload -> m () -> (() -> m ()) -> m ()
forall (m :: * -> *) payload a b.
MonadIO m =>
CancelHandoff -> Settled payload -> m a -> (a -> m b) -> m b
settleInterruptibly CancelHandoff
handoff ([JobRead payload] -> Settled payload
forall payload. [JobRead payload] -> Settled payload
finalized [JobRead payload]
jobs) (() -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ()) (m () -> () -> m ()
forall a b. a -> b -> a
const ((JobRead payload -> m ()) -> [JobRead payload] -> m ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ JobRead payload -> m ()
report [JobRead payload]
jobs))
  where
    report :: JobRead payload -> m ()
report JobRead payload
job
      | Set Int64 -> JobRead payload -> Bool
forall payload. Set Int64 -> JobRead payload -> Bool
hasIdIn Set Int64
cancelled JobRead payload
job = WorkerConfig m payload -> JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload -> JobRead payload -> Text -> m ()
fireCancelled WorkerConfig m payload
config JobRead payload
job Text
"force-cancelled"
      | Bool
otherwise =
          LogConfig -> Text -> m () -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> Text -> m () -> m ()
runHook (WorkerConfig m payload -> JobRead payload -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> JobRead payload -> LogConfig
jobLog WorkerConfig m payload
config JobRead payload
job) Text
"onJobUnavailable" (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
            ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
Job.onJobUnavailable (WorkerConfig m payload -> ObservabilityHooks m payload
forall (m :: * -> *) payload.
WorkerConfig m payload -> ObservabilityHooks m payload
observabilityHooks WorkerConfig m payload
config) JobRead payload
job Text
reason

-- | What the consumer span over this batch covers. A batch of one narrows it to
-- that job.
batchSpanShape :: NonEmpty (Job.JobRead payload) -> ConsumeShape
batchSpanShape :: forall payload. NonEmpty (JobRead payload) -> ConsumeShape
batchSpanShape = ConsumeShape -> ConsumeShape -> Bool -> ConsumeShape
forall a. a -> a -> Bool -> a
bool ConsumeShape
PerJob ConsumeShape
PerBatch (Bool -> ConsumeShape)
-> (NonEmpty (JobRead payload) -> Bool)
-> NonEmpty (JobRead payload)
-> ConsumeShape
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
> Int
1) (Int -> Bool)
-> (NonEmpty (JobRead payload) -> Int)
-> NonEmpty (JobRead payload)
-> Bool
forall b c a. (b -> c) -> (a -> b) -> a -> c
. NonEmpty (JobRead payload) -> Int
forall a. NonEmpty a -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length

-- | Report a failed job to its hooks and to the consumer span it ran under. Only a
-- single-job span takes the error status.
fireFailure
  :: (JobOperation m payload)
  => WorkerConfig m payload
  -> ConsumeShape
  -- ^ What the span this job ran under covers.
  -> Job.JobRead payload
  -> Text
  -> UTCTime
  -> UTCTime
  -> m ()
fireFailure :: forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> ConsumeShape
-> JobRead payload
-> Text
-> UTCTime
-> UTCTime
-> m ()
fireFailure WorkerConfig m payload
config ConsumeShape
shape JobRead payload
job Text
errorMsg UTCTime
startTime UTCTime
endTime = do
  JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
MonadIO m =>
JobRead payload -> Text -> m ()
recordJobFailure JobRead payload
job Text
errorMsg
  Bool -> m () -> m ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (ConsumeShape
shape ConsumeShape -> ConsumeShape -> Bool
forall a. Eq a => a -> a -> Bool
== ConsumeShape
PerJob) (Text -> m ()
forall (m :: * -> *). MonadIO m => Text -> m ()
markSpanError Text
errorMsg)
  LogConfig -> Text -> m () -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> Text -> m () -> m ()
runHook (WorkerConfig m payload -> JobRead payload -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> JobRead payload -> LogConfig
jobLog WorkerConfig m payload
config JobRead payload
job) Text
"onJobFailure" (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
    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 ()
Job.onJobFailure (WorkerConfig m payload -> ObservabilityHooks m payload
forall (m :: * -> *) payload.
WorkerConfig m payload -> ObservabilityHooks m payload
observabilityHooks WorkerConfig m payload
config) JobRead payload
job Text
errorMsg UTCTime
startTime UTCTime
endTime

-- | Report a cancelled job.
fireCancelled
  :: (JobOperation m payload)
  => WorkerConfig m payload
  -> Job.JobRead payload
  -> Text
  -> m ()
fireCancelled :: forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload -> JobRead payload -> Text -> m ()
fireCancelled WorkerConfig m payload
config JobRead payload
job Text
errorMsg = do
  JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
MonadIO m =>
JobRead payload -> Text -> m ()
recordJobCancelled JobRead payload
job Text
errorMsg
  LogConfig -> Text -> m () -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> Text -> m () -> m ()
runHook (WorkerConfig m payload -> JobRead payload -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> JobRead payload -> LogConfig
jobLog WorkerConfig m payload
config JobRead payload
job) Text
"onJobCancelled" (m () -> m ()) -> m () -> m ()
forall a b. (a -> b) -> a -> b
$
    ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
ObservabilityHooks m payload
-> JobPayload payload => JobRead payload -> Text -> m ()
Job.onJobCancelled (WorkerConfig m payload -> ObservabilityHooks m payload
forall (m :: * -> *) payload.
WorkerConfig m payload -> ObservabilityHooks m payload
observabilityHooks WorkerConfig m payload
config) JobRead payload
job Text
errorMsg

-- | Delete what a tree or branch cancel names, returning the rows deleted.
cancelJobFor
  :: (JobOperation m payload)
  => FailureKind
  -> Job.JobRead payload
  -> m Int64
cancelJobFor :: forall (m :: * -> *) payload.
JobOperation m payload =>
FailureKind -> JobRead payload -> m Int64
cancelJobFor FailureKind
kind JobRead payload
job = do
  schemaName <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
  case kind of
    FailureKind
BranchCancelFailure ->
      Text -> Text -> Int64 -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.cancelJobCascade Text
schemaName (JobRead payload -> Text
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> q
Job.queueName JobRead payload
job) (Int64 -> Maybe Int64 -> Int64
forall a. a -> Maybe a -> a
fromMaybe (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
Job.primaryKey JobRead payload
job) (JobRead payload -> Maybe Int64
forall payload q insertedAt adm.
JobRecord payload Int64 q insertedAt adm -> Maybe Int64
Job.parentId JobRead payload
job))
    FailureKind
_ -> Text -> Text -> Int64 -> m Int64
forall (m :: * -> *).
MonadArbiter m =>
Text -> Text -> Int64 -> m Int64
Ops.cancelJobTree Text
schemaName (JobRead payload -> Text
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> q
Job.queueName JobRead payload
job) (JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
Job.primaryKey JobRead payload
job)

-- | 'Left' why a failure write found no row, 'Right' how to report the one it wrote.
type FailureOutcome m = Either Text (m ())

-- | Report a failure the write landed.
reportWritten :: (Monad m) => FailureOutcome m -> m ()
reportWritten :: forall (m :: * -> *). Monad m => FailureOutcome m -> m ()
reportWritten = m () -> Either Text (m ()) -> m ()
forall b a. b -> Either a b -> b
fromRight (() -> m ()
forall a. a -> m a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ())

-- | Handle failure for a single job (retry or move to DLQ), for the caller to report
-- once it commits.
handleJobFailure
  :: forall payload m
   . (JobOperation m payload)
  => WorkerConfig m payload
  -> Ops.TreeLocks
  -- ^ Whether the caller already holds the parent and tree locks.
  -> ConsumeShape
  -- ^ What the span this job ran under covers.
  -> (Text, FailureKind)
  -- ^ The handler exception, classified once for the whole batch.
  -> UTCTime
  -> UTCTime
  -> Job.JobRead payload
  -> m (FailureOutcome m)
handleJobFailure :: forall payload (m :: * -> *).
JobOperation m payload =>
WorkerConfig m payload
-> TreeLocks
-> ConsumeShape
-> (Text, FailureKind)
-> UTCTime
-> UTCTime
-> JobRead payload
-> m (FailureOutcome m)
handleJobFailure WorkerConfig m payload
config TreeLocks
locks ConsumeShape
shape (Text
errorMsg, FailureKind
failureKind) UTCTime
startTime UTCTime
endTime JobRead payload
job
  -- A batch sibling's cancel takes out the whole tree. Zero rows still means gone.
  | FailureKind -> Bool
cancelsTree FailureKind
failureKind =
      m () -> Either Text (m ())
forall a b. b -> Either a b
Right (WorkerConfig m payload -> JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload -> JobRead payload -> Text -> m ()
fireCancelled WorkerConfig m payload
config JobRead payload
job Text
errorMsg) Either Text (m ()) -> m Int64 -> m (Either Text (m ()))
forall a b. a -> m b -> m a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ FailureKind -> JobRead payload -> m Int64
forall (m :: * -> *) payload.
JobOperation m payload =>
FailureKind -> JobRead payload -> m Int64
cancelJobFor FailureKind
failureKind JobRead payload
job
  | FailureKind
failureKind FailureKind -> FailureKind -> Bool
forall a. Eq a => a -> a -> Bool
== FailureKind
PermanentFailure Bool -> Bool -> Bool
|| JobRead payload -> Int32
forall payload q insertedAt adm.
JobRecord payload Int64 q insertedAt adm -> Int32
Job.attempts JobRead payload
job Int32 -> Int32 -> Bool
forall a. Ord a => a -> a -> Bool
>= JobRead payload -> Int32
forall payload. JobRead payload -> Int32
jobMaxAtts JobRead payload
job = do
      schemaName <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
      wrote
        "no longer available for the dead-letter queue"
        (runHook cfg "onJobFailedAndMovedToDLQ" $ Job.onJobFailedAndMovedToDLQ hooks errorMsg job)
        <$> Ops.moveToDLQ locks schemaName (Job.queueName job) errorMsg job
  | Bool
otherwise = do
      let baseDelay :: NominalDiffTime
baseDelay = BackoffStrategy -> Int32 -> NominalDiffTime
calculateBackoff (WorkerConfig m payload -> BackoffStrategy
forall (m :: * -> *) payload.
WorkerConfig m payload -> BackoffStrategy
backoffStrategy WorkerConfig m payload
config) (JobRead payload -> Int32
forall payload q insertedAt adm.
JobRecord payload Int64 q insertedAt adm -> Int32
Job.attempts JobRead payload
job)
      backoffSecs <- IO NominalDiffTime -> m NominalDiffTime
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO NominalDiffTime -> m NominalDiffTime)
-> IO NominalDiffTime -> m NominalDiffTime
forall a b. (a -> b) -> a -> b
$ Jitter -> NominalDiffTime -> IO NominalDiffTime
applyJitter (WorkerConfig m payload -> Jitter
forall (m :: * -> *) payload. WorkerConfig m payload -> Jitter
jitter WorkerConfig m payload
config) NominalDiffTime
baseDelay
      wrote
        "no longer available for retry"
        (runHook cfg "onJobRetry" $ Job.onJobRetry hooks job backoffSecs)
        <$> Arb.updateJobForRetry backoffSecs errorMsg job
  where
    cfg :: LogConfig
cfg = WorkerConfig m payload -> JobRead payload -> LogConfig
forall (m :: * -> *) payload.
WorkerConfig m payload -> JobRead payload -> LogConfig
jobLog WorkerConfig m payload
config JobRead payload
job
    hooks :: ObservabilityHooks m payload
hooks = WorkerConfig m payload -> ObservabilityHooks m payload
forall (m :: * -> *) payload.
WorkerConfig m payload -> ObservabilityHooks m payload
observabilityHooks WorkerConfig m payload
config
    -- Nothing written means the job went elsewhere.
    wrote :: Text -> m () -> Int64 -> Either Text (m ())
wrote Text
reason m ()
after Int64
rowsAffected
      | Int64
rowsAffected Int64 -> Int64 -> Bool
forall a. Eq a => a -> a -> Bool
== Int64
0 = Text -> Either Text (m ())
forall a b. a -> Either a b
Left Text
reason
      | Bool
otherwise = m () -> Either Text (m ())
forall a b. b -> Either a b
Right (WorkerConfig m payload
-> ConsumeShape
-> JobRead payload
-> Text
-> UTCTime
-> UTCTime
-> m ()
forall (m :: * -> *) payload.
JobOperation m payload =>
WorkerConfig m payload
-> ConsumeShape
-> JobRead payload
-> Text
-> UTCTime
-> UTCTime
-> m ()
fireFailure WorkerConfig m payload
config ConsumeShape
shape JobRead payload
job Text
errorMsg UTCTime
startTime UTCTime
endTime m () -> m () -> m ()
forall a b. m a -> m b -> m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> m ()
after)