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
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
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
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
backoff :: WorkerM ()
backoff :: WorkerM ()
backoff = Int -> WorkerM ()
forall (m :: * -> *). MonadIO m => Int -> m ()
threadDelay Int
1_000_000