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

{- | Supervision for the worker's consume loop. A failed @receive@ arrives as the queue
handle's typed fault value, which the step logs and backs off from at its own pacing. Residue,
an exception escaping a dependency's typed contract, is 'superviseLoop''s concern under the
caller's policy.

Shutdown cancels the loop thread, and an un-acked in-flight message simply redelivers, which
is safe because publishing is idempotent.
-}
module Ecluse.Core.Worker.Loop (
    workerLoop,
) where

import Katip (Severity (DebugS, WarningS), logFM, ls)
import UnliftIO.Concurrent (threadDelay)

import Ecluse.Core.Fault (TransportFault, tfDetail)
import Ecluse.Core.Queue (MirrorQueue (receive), QueueMessage)
import Ecluse.Core.Supervision (SupervisionPolicy, superviseLoop)
import Ecluse.Core.Worker.Realise (processBatch)
import Ecluse.Core.Worker.Types

{- | The continuous consume loop: long-poll, process, repeat, under the supervision policy. The
heartbeat advances only on progress, so a persistently faulting @receive@ goes stale on @\/livez@.
-}
workerLoop :: SupervisionPolicy -> WorkerM Void
workerLoop :: SupervisionPolicy -> WorkerM Void
workerLoop SupervisionPolicy
policy = SupervisionPolicy -> WorkerM () -> WorkerM Void
forall (m :: * -> *).
(MonadUnliftIO m, KatipContext m) =>
SupervisionPolicy -> m () -> m Void
superviseLoop SupervisionPolicy
policy WorkerM ()
pollAndProcess

pollAndProcess :: WorkerM ()
pollAndProcess :: WorkerM ()
pollAndProcess = do
    queue <- (WorkerRuntime -> MirrorQueue) -> WorkerM MirrorQueue
forall r (m :: * -> *) a. MonadReader r m => (r -> a) -> m a
asks WorkerRuntime -> MirrorQueue
wrQueue
    liftIO (receive queue) >>= either backOffFrom processPolled

-- No heartbeat advance: the loop is retrying, not healthy-idle, so a persistent fault
-- escalates on @\/livez@ rather than here.
backOffFrom :: TransportFault -> WorkerM ()
backOffFrom :: TransportFault -> WorkerM ()
backOffFrom TransportFault
fault = do
    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
"worker receive failed, backing off: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> TransportFault -> Text
tfDetail TransportFault
fault))
    WorkerM ()
backoff

-- Beat on every successful poll: an empty long-poll is a healthy idle. 'processBatch' beats
-- again after each job, so a long batch cannot starve it.
processPolled :: [QueueMessage] -> WorkerM ()
processPolled :: [QueueMessage] -> WorkerM ()
processPolled [QueueMessage]
messages = do
    Bool -> WorkerM () -> WorkerM ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
unless ([QueueMessage] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null [QueueMessage]
messages) (WorkerM () -> WorkerM ()) -> WorkerM () -> WorkerM ()
forall a b. (a -> b) -> a -> b
$
        Severity -> LogStr -> WorkerM ()
forall (m :: * -> *).
(Applicative m, KatipContext m) =>
Severity -> LogStr -> m ()
logFM Severity
DebugS (Text -> LogStr
forall a. StringConv a Text => a -> LogStr
ls (Text
"worker received " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Int -> Text
forall b a. (Show a, IsString b) => a -> b
show ([QueueMessage] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [QueueMessage]
messages) Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
" messages" :: Text))
    WorkerM ()
recordWorkerProgress
    [QueueMessage] -> WorkerM ()
processBatch [QueueMessage]
messages

-- The fixed pause after a faulted poll, so the loop retries a persistently failing
-- queue backend at a bounded rate rather than hot-looping.
backoff :: WorkerM ()
backoff :: WorkerM ()
backoff = Int -> WorkerM ()
forall (m :: * -> *). MonadIO m => Int -> m ()
threadDelay Int
1_000_000