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
data Disposition
= DisposeAck
| DisposeDeadLetter
| DisposeRelease
| DisposeLeave
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)
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
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)
DisposeAck <$ logFM ErrorS (ls (budgetSpentReason budget message))
else decideDelivery message
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 ->
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 ->
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))
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."
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))
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))
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)