-- SPDX-FileCopyrightText: 2026 Alexandra de Wit
--
-- SPDX-License-Identifier: MIT

{- | Realising a job's verdict at the queue handle: ack, dead-letter, or release for retry.

Every receipt in a batch is leased for the batch's whole run ("Ecluse.Core.Worker.Lease"), so
a job never races the backend's visibility window and its disposition is the one thing that
ends the lease. Jobs run __sequentially__, one artifact task at a time. A delivery that already
spent the queue's redelivery budget is retired before its job runs, so a message nothing else
captures stops cycling instead of re-fetching its artifact on every redelivery.
-}
module Ecluse.Core.Worker.Realise (
    processBatch,
) where

import Katip (Severity (ErrorS, WarningS), logFM, ls)

import Ecluse.Core.Fault (tfDetail)
import Ecluse.Core.Queue (
    DeliveryBudget,
    MirrorQueue (ack, deadLetter, deliveryBudget, extendVisibility),
    QueueMessage (msgJob, msgReceipt, msgReceiveCount),
    ReceiptHandle,
    Seconds (Seconds),
    deliveryBudgetSpent,
    retiringDelivery,
 )
import Ecluse.Core.Telemetry.Metrics qualified as Metric
import Ecluse.Core.Telemetry.Record (WorkerMetricsPort (..))
import Ecluse.Core.Worker.Job (
    JobOutcome (DeadLettered, Dropped, Retried, SourceUnavailable, Succeeded),
    RetryLeg (AfterPublish, BeforePublish),
    processJob,
 )
import Ecluse.Core.Worker.Lease (LeasedReceipt, disposing, leasedMessage, queueLeaseOps, whileLeased, withLeasedBatch)
import Ecluse.Core.Worker.Types

-- What the worker leaves at the queue handle once a message is decided, realised under the
-- receipt's own lease so no renewal can follow it. 'DisposeLeave' touches the handle not at all.
data Disposition
    = DisposeAck
    | DisposeDeadLetter
    | DisposeRelease
    | DisposeLeave

{- | Process one batch sequentially under a lease on every receipt in it. The heartbeat advances
per job, so 'Ecluse.Core.Worker.Liveness.workerHeartbeatStaleAfter' covers one job.
-}
processBatch :: [QueueMessage] -> WorkerM ()
processBatch :: [QueueMessage] -> WorkerM ()
processBatch [QueueMessage]
messages = do
    queue <- (WorkerRuntime -> MirrorQueue) -> WorkerM MirrorQueue
forall r (m :: * -> *) a. MonadReader r m => (r -> a) -> m a
asks WorkerRuntime -> MirrorQueue
wrQueue
    withLeasedBatch (queueLeaseOps queue) messages (traverse_ processLeased)

{- Decide one leased message and realise it, then beat. A receipt whose lease was dropped is
left unacknowledged, so the backend redelivers it once its window lapses. -}
processLeased :: LeasedReceipt -> WorkerM ()
processLeased :: LeasedReceipt -> WorkerM ()
processLeased LeasedReceipt
leased = do
    decided <- LeasedReceipt -> WorkerM Disposition -> WorkerM (Maybe Disposition)
forall (m :: * -> *) a.
MonadUnliftIO m =>
LeasedReceipt -> m a -> m (Maybe a)
whileLeased LeasedReceipt
leased (QueueMessage -> WorkerM Disposition
decideMessage QueueMessage
message)
    whenJust decided (disposing leased . realiseDisposition (msgReceipt message))
    recordWorkerProgress
  where
    message :: QueueMessage
message = LeasedReceipt -> QueueMessage
leasedMessage LeasedReceipt
leased

{- Check the queue's delivery budget before running the job, so a poison message retires without
re-fetching its artifact, even on a queue with no dead-letter terminus. -}
decideMessage :: QueueMessage -> WorkerM Disposition
decideMessage :: QueueMessage -> WorkerM Disposition
decideMessage QueueMessage
message = do
    budget <- (WorkerRuntime -> DeliveryBudget) -> WorkerM DeliveryBudget
forall r (m :: * -> *) a. MonadReader r m => (r -> a) -> m a
asks (MirrorQueue -> DeliveryBudget
deliveryBudget (MirrorQueue -> DeliveryBudget)
-> (WorkerRuntime -> MirrorQueue)
-> WorkerRuntime
-> DeliveryBudget
forall b c a. (b -> c) -> (a -> b) -> a -> c
. WorkerRuntime -> MirrorQueue
wrQueue)
    if deliveryBudgetSpent budget message
        then do
            metrics <- asks wrMetrics
            liftIO (wmpMirrorJobProcessed metrics Metric.Discarded)
            -- On a queue with no dead-letter terminus this line is the only record it leaves.
            DisposeAck <$ logFM ErrorS (ls (budgetSpentReason budget message))
        else decideDelivery message

-- Run the job and read its outcome as the disposition for a delivery still within the budget.
decideDelivery :: QueueMessage -> WorkerM Disposition
decideDelivery :: QueueMessage -> WorkerM Disposition
decideDelivery QueueMessage
message = do
    metrics <- (WorkerRuntime -> WorkerMetricsPort) -> WorkerM WorkerMetricsPort
forall r (m :: * -> *) a. MonadReader r m => (r -> a) -> m a
asks WorkerRuntime -> WorkerMetricsPort
wrMetrics
    outcome <- processJob (msgJob message)
    liftIO (wmpMirrorJobProcessed metrics (jobResultMetric outcome))
    case outcome of
        JobOutcome
Succeeded -> Disposition -> WorkerM Disposition
forall a. a -> WorkerM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Disposition
DisposeAck
        Dropped Text
reason ->
            -- Non-retryable, and not worth a dead-letter forensic trail, so retire it instead.
            Disposition
DisposeAck Disposition -> WorkerM () -> WorkerM Disposition
forall a b. a -> WorkerM b -> WorkerM a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Severity -> LogStr -> WorkerM ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
ErrorS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"dropping unrecoverable mirror job: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
reason))
        SourceUnavailable Text
reason ->
            Disposition
DisposeAck Disposition -> WorkerM () -> WorkerM Disposition
forall a b. a -> WorkerM b -> WorkerM a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Severity -> LogStr -> WorkerM ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
ErrorS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"refusing to mirror without the source version object: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
reason))
        DeadLettered Text
reason ->
            -- Alarm first: on the in-memory backend the log and metric are the only record.
            Disposition
DisposeDeadLetter Disposition -> WorkerM () -> WorkerM Disposition
forall a b. a -> WorkerM b -> WorkerM a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Severity -> LogStr -> WorkerM ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
ErrorS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"dead-lettering unmirrorable mirror job (rides the backend's dead-letter terminus): " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
reason))
        Retried RetryLeg
leg Text
reason ->
            RetryLeg -> Disposition
retryDisposition RetryLeg
leg Disposition -> WorkerM () -> WorkerM Disposition
forall a b. a -> WorkerM b -> WorkerM a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ Severity -> LogStr -> WorkerM ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
WarningS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"leaving mirror job un-acked for retry (redelivered by a durable queue, re-mirrored on next demand by the in-memory one): " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
reason))

{- Only a failed publish resets the window. A transient failure before it keeps its lease, so an
upstream outage cannot spend the queue's whole redelivery budget in seconds. -}
retryDisposition :: RetryLeg -> Disposition
retryDisposition :: RetryLeg -> Disposition
retryDisposition = \case
    RetryLeg
AfterPublish -> Disposition
DisposeRelease
    RetryLeg
BeforePublish -> Disposition
DisposeLeave

realiseDisposition :: ReceiptHandle -> Disposition -> WorkerM ()
realiseDisposition :: ReceiptHandle -> Disposition -> WorkerM ()
realiseDisposition ReceiptHandle
receipt = \case
    Disposition
DisposeAck -> ReceiptHandle -> WorkerM ()
ackMessage ReceiptHandle
receipt
    Disposition
DisposeDeadLetter -> ReceiptHandle -> WorkerM ()
deadLetterMessage ReceiptHandle
receipt
    Disposition
DisposeRelease -> ReceiptHandle -> WorkerM ()
releaseForRetry ReceiptHandle
receipt
    Disposition
DisposeLeave -> WorkerM ()
forall (f :: * -> *). Applicative f => f ()
pass

budgetSpentReason :: DeliveryBudget -> QueueMessage -> Text
budgetSpentReason :: DeliveryBudget -> QueueMessage -> Text
budgetSpentReason DeliveryBudget
budget QueueMessage
message =
    Text
"discarding a mirror job after "
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show (QueueMessage -> Int
msgReceiveCount QueueMessage
message)
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" deliveries (this queue retires one on delivery "
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show (DeliveryBudget -> Int
retiringDelivery DeliveryBudget
budget)
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"): "
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> MirrorJob -> Text
renderJob (QueueMessage -> MirrorJob
msgJob QueueMessage
message)
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
". No dead-letter queue captured it, so it is retired here rather than left to"
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" cycle until the queue's retention window drops it unseen. Attach a redrive"
        Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" policy to retain it for inspection."

-- Classify a job outcome for the @ecluse.mirror.jobs.processed@ metric. 'Metric.Discarded' is
-- absent here on purpose: the worker counts a budget-spent delivery at its retirement.
jobResultMetric :: JobOutcome -> Metric.MirrorResult
jobResultMetric :: JobOutcome -> MirrorResult
jobResultMetric = \case
    JobOutcome
Succeeded -> MirrorResult
Metric.Published
    Dropped Text
_ -> MirrorResult
Metric.Failed
    SourceUnavailable Text
_ -> MirrorResult
Metric.Failed
    DeadLettered Text
_ -> MirrorResult
Metric.Failed
    Retried RetryLeg
_ Text
_ -> MirrorResult
Metric.Failed

ackMessage :: ReceiptHandle -> WorkerM ()
ackMessage :: ReceiptHandle -> WorkerM ()
ackMessage ReceiptHandle
receipt =
    (MirrorQueue -> IO (Either TransportFault ()))
-> (TransportFault -> WorkerM ()) -> WorkerM ()
forall a.
(MirrorQueue -> IO (Either TransportFault a))
-> (TransportFault -> WorkerM ()) -> WorkerM ()
queueOp (MirrorQueue -> ReceiptHandle -> IO (Either TransportFault ())
`ack` ReceiptHandle
receipt) ((TransportFault -> WorkerM ()) -> WorkerM ())
-> (TransportFault -> WorkerM ()) -> WorkerM ()
forall a b. (a -> b) -> a -> b
$ \TransportFault
fault ->
        Severity -> LogStr -> WorkerM ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
WarningS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"ack failed; the processed message will redeliver (harmless, publishing is idempotent): " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> TransportFault -> Text
tfDetail TransportFault
fault))

-- Hand the message to the queue's dead-letter terminus, never a plain delete, which would
-- silently discard it on a durable queue.
deadLetterMessage :: ReceiptHandle -> WorkerM ()
deadLetterMessage :: ReceiptHandle -> WorkerM ()
deadLetterMessage ReceiptHandle
receipt =
    (MirrorQueue -> IO (Either TransportFault ()))
-> (TransportFault -> WorkerM ()) -> WorkerM ()
forall a.
(MirrorQueue -> IO (Either TransportFault a))
-> (TransportFault -> WorkerM ()) -> WorkerM ()
queueOp (MirrorQueue -> ReceiptHandle -> IO (Either TransportFault ())
`deadLetter` ReceiptHandle
receipt) ((TransportFault -> WorkerM ()) -> WorkerM ())
-> (TransportFault -> WorkerM ()) -> WorkerM ()
forall a b. (a -> b) -> a -> b
$ \TransportFault
fault ->
        Severity -> LogStr -> WorkerM ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
WarningS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"dead-letter realisation failed; the message redelivers and re-fails terminally (harmless): " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> TransportFault -> Text
tfDetail TransportFault
fault))

-- Reset the message to visible, so a failed publish redelivers at once instead of waiting out the
-- lease the worker held. Best effort: a missed reset only delays the redelivery.
releaseForRetry :: ReceiptHandle -> WorkerM ()
releaseForRetry :: ReceiptHandle -> WorkerM ()
releaseForRetry ReceiptHandle
receipt =
    (MirrorQueue -> IO (Either TransportFault ()))
-> (TransportFault -> WorkerM ()) -> WorkerM ()
forall a.
(MirrorQueue -> IO (Either TransportFault a))
-> (TransportFault -> WorkerM ()) -> WorkerM ()
queueOp (\MirrorQueue
queue -> MirrorQueue
-> ReceiptHandle -> Seconds -> IO (Either TransportFault ())
extendVisibility MirrorQueue
queue ReceiptHandle
receipt (Int -> Seconds
Seconds Int
0)) (WorkerM () -> TransportFault -> WorkerM ()
forall a b. a -> b -> a
const WorkerM ()
forall (f :: * -> *). Applicative f => f ()
pass)