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

{- | The STM-backed in-memory 'MirrorQueue': the bounded, best-effort production backend
mirroring rolls over to when no @ECLUSE_QUEUE__URL@ is set.

It honours the handle's contract ("Ecluse.Core.Queue") with two departures. 'enqueue' drops the
newest job past the depth cap, and a 'receive' is final, so nothing redelivers and no delivery
carries a lease. Both are safe because the next demand re-enqueues the job.
-}
module Ecluse.Core.Queue.Memory (
    -- * Bounded in-memory production backend
    MemoryQueueConfig (..),
    defaultMemoryQueueConfig,
    newBoundedInMemoryQueue,
    memoryQueueDropReportInterval,
) where

import Control.Concurrent.STM.TBQueue (TBQueue, newTBQueueIO, readTBQueue, tryReadTBQueue)
import System.Timeout (timeout)

import Ecluse.Core.Queue (
    DeadLetterTerminus (TerminusAbsent),
    MirrorJob,
    MirrorQueue (..),
    QueueMessage (..),
    defaultDeliveryBudget,
    mkReceiptHandle,
 )
import Ecluse.Core.Queue.Buffer (
    reportWorthy,
    writeOrDrop,
 )

{- | The bounded in-memory backend's depth cap and idle-poll window. Build it with
'defaultMemoryQueueConfig' for the production poll window.
-}
data MemoryQueueConfig = MemoryQueueConfig
    { MemoryQueueConfig -> Int
memQueueMaxDepth :: Int
    {- ^ The maximum number of jobs the queue holds. The config layer enforces a positive cap. An
    'enqueue' past it drops the newest job, a safe loss because the next demand re-enqueues.
    -}
    , MemoryQueueConfig -> Int
memQueuePollWaitMicros :: Int
    {- ^ The idle long-poll window in microseconds: how long a 'receive' waits for a job before
    returning @[]@. The bound keeps the worker's liveness heartbeat advancing on an idle queue.
    -}
    }
    deriving stock (MemoryQueueConfig -> MemoryQueueConfig -> Bool
(MemoryQueueConfig -> MemoryQueueConfig -> Bool)
-> (MemoryQueueConfig -> MemoryQueueConfig -> Bool)
-> Eq MemoryQueueConfig
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: MemoryQueueConfig -> MemoryQueueConfig -> Bool
== :: MemoryQueueConfig -> MemoryQueueConfig -> Bool
$c/= :: MemoryQueueConfig -> MemoryQueueConfig -> Bool
/= :: MemoryQueueConfig -> MemoryQueueConfig -> Bool
Eq, Int -> MemoryQueueConfig -> ShowS
[MemoryQueueConfig] -> ShowS
MemoryQueueConfig -> String
(Int -> MemoryQueueConfig -> ShowS)
-> (MemoryQueueConfig -> String)
-> ([MemoryQueueConfig] -> ShowS)
-> Show MemoryQueueConfig
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> MemoryQueueConfig -> ShowS
showsPrec :: Int -> MemoryQueueConfig -> ShowS
$cshow :: MemoryQueueConfig -> String
show :: MemoryQueueConfig -> String
$cshowList :: [MemoryQueueConfig] -> ShowS
showList :: [MemoryQueueConfig] -> ShowS
Show)

{- | A 'MemoryQueueConfig' at the production @20s@ idle-poll window, which sits under
'Ecluse.Core.Worker.workerHeartbeatStaleAfter' so an idle poll cannot stall @\/livez@.
-}
defaultMemoryQueueConfig :: Int -> MemoryQueueConfig
defaultMemoryQueueConfig :: Int -> MemoryQueueConfig
defaultMemoryQueueConfig Int
maxDepth =
    MemoryQueueConfig
        { memQueueMaxDepth :: Int
memQueueMaxDepth = Int
maxDepth
        , memQueuePollWaitMicros :: Int
memQueuePollWaitMicros = Int
20_000_000
        }

-- Held at the SQS batch cap, so the worker sees one bounded batch shape whatever the backend.
memoryQueueBatchSize :: Int
memoryQueueBatchSize :: Int
memoryQueueBatchSize = Int
10

{- | How many cap-overflow drops the bounded in-memory backend absorbs between warning reports.
It reports the first drop, then every multiple of this, so a sustained flood cannot spam.
-}
memoryQueueDropReportInterval :: Int
memoryQueueDropReportInterval :: Int
memoryQueueDropReportInterval = Int
1000

{- | Build the bounded, best-effort in-memory 'MirrorQueue'. A cold-cache @npm ci@ enqueues
thousands of jobs at once, so 'enqueue' sheds past 'memQueueMaxDepth' rather than throwing.
-}
newBoundedInMemoryQueue ::
    -- | The depth cap and the idle-poll window.
    MemoryQueueConfig ->
    -- | Invoked on each report-worthy cap-overflow drop with the running total, for the log.
    (Int -> IO ()) ->
    IO MirrorQueue
newBoundedInMemoryQueue :: MemoryQueueConfig -> (Int -> IO ()) -> IO MirrorQueue
newBoundedInMemoryQueue MemoryQueueConfig
cfg Int -> IO ()
onDrop = do
    -- A capacity of at least one. The config layer enforces a positive cap, but a
    -- directly-constructed queue must never be the degenerate always-full zero.
    queue <- Natural -> IO (TBQueue MirrorJob)
forall a. Natural -> IO (TBQueue a)
newTBQueueIO (Int -> Natural
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
1 (MemoryQueueConfig -> Int
memQueueMaxDepth MemoryQueueConfig
cfg)))
    dropCount <- newTVarIO (0 :: Int)
    nextReceipt <- newTVarIO (0 :: Word64)
    pure
        MirrorQueue
            { enqueue = \MirrorJob
job -> do
                dropped <- STM (Maybe Int) -> IO (Maybe Int)
forall (m :: * -> *) a. MonadIO m => STM a -> m a
atomically (TBQueue MirrorJob -> TVar Int -> MirrorJob -> STM (Maybe Int)
writeOrDrop TBQueue MirrorJob
queue TVar Int
dropCount MirrorJob
job)
                whenJust dropped (\Int
n -> Bool -> IO () -> IO ()
forall (f :: * -> *). Applicative f => Bool -> f () -> f ()
when (Int -> Int -> Bool
reportWorthy Int
n Int
memoryQueueDropReportInterval) (Int -> IO ()
onDrop Int
n))
                -- A cap overflow is the documented drop-newest shed (reported through
                -- the callback), not a backend fault: the enqueue itself worked.
                pure (Right ())
            , -- @timeout@ over @atomically@, not @registerDelay@, so the poll bound holds on the
              -- non-threaded RTS too. A fired timeout aborts the transaction, consuming nothing.
              receive = Right . fromMaybe [] <$> timeout (memQueuePollWaitMicros cfg) (atomically (receiveBatch queue nextReceipt))
            , -- A delivered job is already gone from the queue, so there is nothing to
              -- retire and a failed job redelivers via the next demand, not here.
              ack = const (pure (Right ()))
            , extendVisibility = \ReceiptHandle
_ Seconds
_ -> Either TransportFault () -> IO (Either TransportFault ())
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (() -> Either TransportFault ()
forall a b. b -> Either a b
Right ())
            , -- A delivered job is already dropped, so a terminal fault has nowhere further to
              -- go. Its observability is the worker's error log and metric.
              deadLetter = const (pure (Right ()))
            , -- Nothing here captures a poison message, and nothing redelivers one, so
              -- the backend holds the budget inert at the shipped default.
              deliveryBudget = defaultDeliveryBudget
            , deadLetterTerminus = Right TerminusAbsent
            }

-- One transaction, so a timeout fired by the caller during the initial block aborts the whole
-- batch and consumes nothing.
receiveBatch :: TBQueue MirrorJob -> TVar Word64 -> STM [QueueMessage]
receiveBatch :: TBQueue MirrorJob -> TVar Word64 -> STM [QueueMessage]
receiveBatch TBQueue MirrorJob
queue TVar Word64
nextReceipt = do
    headJob <- TBQueue MirrorJob -> STM MirrorJob
forall a. TBQueue a -> STM a
readTBQueue TBQueue MirrorJob
queue
    rest <- drainUpTo (memoryQueueBatchSize - 1)
    traverse assignReceipt (headJob : rest)
  where
    drainUpTo :: Int -> STM [MirrorJob]
    drainUpTo :: Int -> STM [MirrorJob]
drainUpTo Int
budget
        | Int
budget Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
<= Int
0 = [MirrorJob] -> STM [MirrorJob]
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure []
        | Bool
otherwise =
            TBQueue MirrorJob -> STM (Maybe MirrorJob)
forall a. TBQueue a -> STM (Maybe a)
tryReadTBQueue TBQueue MirrorJob
queue STM (Maybe MirrorJob)
-> (Maybe MirrorJob -> STM [MirrorJob]) -> STM [MirrorJob]
forall a b. STM a -> (a -> STM b) -> STM b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \case
                Maybe MirrorJob
Nothing -> [MirrorJob] -> STM [MirrorJob]
forall a. a -> STM a
forall (f :: * -> *) a. Applicative f => a -> f a
pure []
                Just MirrorJob
job -> (MirrorJob
job MirrorJob -> [MirrorJob] -> [MirrorJob]
forall a. a -> [a] -> [a]
:) ([MirrorJob] -> [MirrorJob]) -> STM [MirrorJob] -> STM [MirrorJob]
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Int -> STM [MirrorJob]
drainUpTo (Int
budget Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
1)

    assignReceipt :: MirrorJob -> STM QueueMessage
    assignReceipt :: MirrorJob -> STM QueueMessage
assignReceipt MirrorJob
job = do
        n <- TVar Word64 -> STM Word64
forall a. TVar a -> STM a
readTVar TVar Word64
nextReceipt
        writeTVar nextReceipt (n + 1)
        -- Every delivery is a first delivery and none expires: a received job leaves the
        -- queue for good, so this backend never redelivers one and grants no lease.
        pure QueueMessage{msgJob = job, msgReceipt = mkReceiptHandle (show n), msgReceiveCount = 1, msgLease = Nothing}