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

{- | The mirror-queue handle: the durable hand-off from the request path to the mirror worker.

Mirroring is demand-driven, so the serve path 'enqueue's a 'MirrorJob' and answers at once
while a worker 'receive's it, verifies the artifact, publishes it and 'ack's. At-least-once
delivery is safe because publishing is idempotent, which is why retry is "do not 'ack'" and
there is no @nack@. This cloud surface is the one whose API differs materially per provider,
so it is a record of functions and 'ReceiptHandle' is opaque. See
@docs\/architecture\/cloud-backends.md@, "Mirror Queue".
-}
module Ecluse.Core.Queue (
    -- * Queue handle
    MirrorQueue (..),
    noMirrorQueue,

    -- * Payloads
    MirrorJob (..),
    RemoteSpanContext (..),
    QueueMessage (..),

    -- * The payload's wire mapping
    encodeJob,
    decodeJob,

    -- * Opaque receipt
    ReceiptHandle,
    mkReceiptHandle,
    unReceiptHandle,

    -- * Durations and the receipt lease
    Seconds (..),
    ReceiptLease (..),

    -- * Dead-letter terminus and the redelivery budget
    DeadLetterTerminus (..),
    DeliveryBudget (..),
    defaultDeliveryBudget,
    effectiveDeliveryBudget,
    retiringDelivery,
    deliveryBudgetSpent,
) where

import Data.Aeson (eitherDecodeStrict', object, withObject, (.:), (.:?), (.=))
import Data.Aeson qualified as Aeson
import Data.Aeson.Types (Parser, parseEither)

import Ecluse.Core.Ecosystem (Ecosystem, ecosystemName, parseEcosystem)
import Ecluse.Core.Fault (TransportCause (TransportProtocol), TransportFault, transportFault)
import Ecluse.Core.Package (PackageName, pkgEcosystem, pkgNamespace, unScope, unscopedName)
import Ecluse.Core.Queue.Lease (ReceiptLease (..), Seconds (..))
import Ecluse.Core.Security.Egress (RegistryUrl, registryUrlText)
import Ecluse.Core.Server.Path (Filename, mkFilename, unFilename)
import Ecluse.Core.Version (Version, mkVersion, renderVersion)

{- | Everything the worker needs to back-fill one artifact into the mirror target. The payload
is a __trust boundary__: it carries selection keys and never authority, so no digest, no size.
-}
data MirrorJob = MirrorJob
    { MirrorJob -> PackageName
jobPackage :: PackageName
    -- ^ The package whose artifact the worker mirrors.
    , MirrorJob -> Version
jobVersion :: Version
    -- ^ The specific version to mirror.
    , MirrorJob -> RegistryUrl
jobArtifactUrl :: RegistryUrl
    {- ^ Where the worker fetches the artifact bytes from. A wire decode re-forms the validated
    https egress witness rather than trusting the payload's text.
    -}
    , MirrorJob -> Filename
jobArtifactFilename :: Filename
    {- ^ The serve-time-admitted artifact's filename: a selection key the admission gate
    cross-checks against current metadata, not authority.
    -}
    , MirrorJob -> Maybe RemoteSpanContext
jobTraceContext :: Maybe RemoteSpanContext
    {- ^ The trace context of the span that enqueued the job, so the worker's per-job span links
    back across the asynchronous hop. 'Nothing' when the producer carried none.
    -}
    }
    deriving stock (MirrorJob -> MirrorJob -> Bool
(MirrorJob -> MirrorJob -> Bool)
-> (MirrorJob -> MirrorJob -> Bool) -> Eq MirrorJob
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: MirrorJob -> MirrorJob -> Bool
== :: MirrorJob -> MirrorJob -> Bool
$c/= :: MirrorJob -> MirrorJob -> Bool
/= :: MirrorJob -> MirrorJob -> Bool
Eq, Int -> MirrorJob -> ShowS
[MirrorJob] -> ShowS
MirrorJob -> String
(Int -> MirrorJob -> ShowS)
-> (MirrorJob -> String)
-> ([MirrorJob] -> ShowS)
-> Show MirrorJob
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> MirrorJob -> ShowS
showsPrec :: Int -> MirrorJob -> ShowS
$cshow :: MirrorJob -> String
show :: MirrorJob -> String
$cshowList :: [MirrorJob] -> ShowS
showList :: [MirrorJob] -> ShowS
Show)

{- | The @traceparent@ and @tracestate@ of the enqueueing span, verbatim. The queue never parses
them, so an unparseable carrier yields no span link rather than a decode failure.
-}
data RemoteSpanContext = RemoteSpanContext
    { RemoteSpanContext -> Text
rscTraceparent :: Text
    -- ^ The W3C @traceparent@ header value of the enqueueing span.
    , RemoteSpanContext -> Text
rscTracestate :: Text
    -- ^ The W3C @tracestate@ value (possibly empty), so vendor trace state survives the hop.
    }
    deriving stock (RemoteSpanContext -> RemoteSpanContext -> Bool
(RemoteSpanContext -> RemoteSpanContext -> Bool)
-> (RemoteSpanContext -> RemoteSpanContext -> Bool)
-> Eq RemoteSpanContext
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: RemoteSpanContext -> RemoteSpanContext -> Bool
== :: RemoteSpanContext -> RemoteSpanContext -> Bool
$c/= :: RemoteSpanContext -> RemoteSpanContext -> Bool
/= :: RemoteSpanContext -> RemoteSpanContext -> Bool
Eq, Int -> RemoteSpanContext -> ShowS
[RemoteSpanContext] -> ShowS
RemoteSpanContext -> String
(Int -> RemoteSpanContext -> ShowS)
-> (RemoteSpanContext -> String)
-> ([RemoteSpanContext] -> ShowS)
-> Show RemoteSpanContext
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> RemoteSpanContext -> ShowS
showsPrec :: Int -> RemoteSpanContext -> ShowS
$cshow :: RemoteSpanContext -> String
show :: RemoteSpanContext -> String
$cshowList :: [RemoteSpanContext] -> ShowS
showList :: [RemoteSpanContext] -> ShowS
Show)

{- | Encode a 'MirrorJob' as the JSON text of a queue message body, the inverse of 'decodeJob'. The
identity rides as a namespace and a base name, so a namespaced name round-trips on any ecosystem.
-}
encodeJob :: MirrorJob -> Text
encodeJob :: MirrorJob -> Text
encodeJob MirrorJob
job =
    ByteString -> Text
forall a b. ConvertUtf8 a b => b -> a
decodeUtf8 (ByteString -> Text) -> (Value -> ByteString) -> Value -> Text
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Value -> ByteString
forall a. ToJSON a => a -> ByteString
Aeson.encode (Value -> Text) -> Value -> Text
forall a b. (a -> b) -> a -> b
$
        [Pair] -> Value
object
            [ Key
"ecosystem" Key -> Text -> Pair
forall v. ToJSON v => Key -> v -> Pair
forall e kv v. (KeyValue e kv, ToJSON v) => Key -> v -> kv
.= Ecosystem -> Text
ecosystemName (PackageName -> Ecosystem
pkgEcosystem (MirrorJob -> PackageName
jobPackage MirrorJob
job))
            , Key
"namespace" Key -> Maybe Text -> Pair
forall v. ToJSON v => Key -> v -> Pair
forall e kv v. (KeyValue e kv, ToJSON v) => Key -> v -> kv
.= (Scope -> Text
unScope (Scope -> Text) -> Maybe Scope -> Maybe Text
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> PackageName -> Maybe Scope
pkgNamespace (MirrorJob -> PackageName
jobPackage MirrorJob
job))
            , Key
"name" Key -> Text -> Pair
forall v. ToJSON v => Key -> v -> Pair
forall e kv v. (KeyValue e kv, ToJSON v) => Key -> v -> kv
.= PackageName -> Text
unscopedName (MirrorJob -> PackageName
jobPackage MirrorJob
job)
            , Key
"version" Key -> Text -> Pair
forall v. ToJSON v => Key -> v -> Pair
forall e kv v. (KeyValue e kv, ToJSON v) => Key -> v -> kv
.= Version -> Text
renderVersion (MirrorJob -> Version
jobVersion MirrorJob
job)
            , Key
"artifactUrl" Key -> Text -> Pair
forall v. ToJSON v => Key -> v -> Pair
forall e kv v. (KeyValue e kv, ToJSON v) => Key -> v -> kv
.= RegistryUrl -> Text
registryUrlText (MirrorJob -> RegistryUrl
jobArtifactUrl MirrorJob
job)
            , Key
"filename" Key -> Text -> Pair
forall v. ToJSON v => Key -> v -> Pair
forall e kv v. (KeyValue e kv, ToJSON v) => Key -> v -> kv
.= Filename -> Text
unFilename (MirrorJob -> Filename
jobArtifactFilename MirrorJob
job)
            , Key
"traceContext" Key -> Maybe Value -> Pair
forall v. ToJSON v => Key -> v -> Pair
forall e kv v. (KeyValue e kv, ToJSON v) => Key -> v -> kv
.= (RemoteSpanContext -> Value
encodeTraceContext (RemoteSpanContext -> Value)
-> Maybe RemoteSpanContext -> Maybe Value
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> MirrorJob -> Maybe RemoteSpanContext
jobTraceContext MirrorJob
job)
            ]

-- The W3C traceparent and tracestate verbatim, so the worker can re-establish the
-- cross-async span link. A 'Nothing' carrier round-trips through a JSON null.
encodeTraceContext :: RemoteSpanContext -> Aeson.Value
encodeTraceContext :: RemoteSpanContext -> Value
encodeTraceContext RemoteSpanContext
rsc =
    [Pair] -> Value
object
        [ Key
"traceparent" Key -> Text -> Pair
forall v. ToJSON v => Key -> v -> Pair
forall e kv v. (KeyValue e kv, ToJSON v) => Key -> v -> kv
.= RemoteSpanContext -> Text
rscTraceparent RemoteSpanContext
rsc
        , Key
"tracestate" Key -> Text -> Pair
forall v. ToJSON v => Key -> v -> Pair
forall e kv v. (KeyValue e kv, ToJSON v) => Key -> v -> kv
.= RemoteSpanContext -> Text
rscTracestate RemoteSpanContext
rsc
        ]

{- | Decode a queue message body back into a 'MirrorJob'. The payload is a __trust boundary__, so
the name, the filename, and the artifact URL each go back through their own gate.
-}
decodeJob ::
    -- | Read a wire namespace and base name through the ecosystem's own grammar.
    (Ecosystem -> Maybe Text -> Text -> Either Text PackageName) ->
    -- | Re-form the artifact URL's validated https egress witness.
    (Text -> Either Text RegistryUrl) ->
    Text ->
    Either Text MirrorJob
decodeJob :: (Ecosystem -> Maybe Text -> Text -> Either Text PackageName)
-> (Text -> Either Text RegistryUrl)
-> Text
-> Either Text MirrorJob
decodeJob Ecosystem -> Maybe Text -> Text -> Either Text PackageName
packageName Text -> Either Text RegistryUrl
egressUrl Text
body =
    (String -> Text) -> Either String Value -> Either Text Value
forall a b c. (a -> b) -> Either a c -> Either b c
forall (p :: * -> * -> *) a b c.
Bifunctor p =>
(a -> b) -> p a c -> p b c
first String -> Text
forall a. ToText a => a -> Text
toText (ByteString -> Either String Value
forall a. FromJSON a => ByteString -> Either String a
eitherDecodeStrict' (Text -> ByteString
forall a b. ConvertUtf8 a b => a -> b
encodeUtf8 Text
body))
        Either Text Value
-> (Value -> Either Text MirrorJob) -> Either Text MirrorJob
forall a b. Either Text a -> (a -> Either Text b) -> Either Text b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= (String -> Text)
-> Either String MirrorJob -> Either Text MirrorJob
forall a b c. (a -> b) -> Either a c -> Either b c
forall (p :: * -> * -> *) a b c.
Bifunctor p =>
(a -> b) -> p a c -> p b c
first String -> Text
forall a. ToText a => a -> Text
toText (Either String MirrorJob -> Either Text MirrorJob)
-> (Value -> Either String MirrorJob)
-> Value
-> Either Text MirrorJob
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (Value -> Parser MirrorJob) -> Value -> Either String MirrorJob
forall a b. (a -> Parser b) -> a -> Either String b
parseEither ((Ecosystem -> Maybe Text -> Text -> Either Text PackageName)
-> (Text -> Either Text RegistryUrl) -> Value -> Parser MirrorJob
parseMirrorJob Ecosystem -> Maybe Text -> Text -> Either Text PackageName
packageName Text -> Either Text RegistryUrl
egressUrl)

parseMirrorJob ::
    (Ecosystem -> Maybe Text -> Text -> Either Text PackageName) ->
    (Text -> Either Text RegistryUrl) ->
    Aeson.Value ->
    Parser MirrorJob
parseMirrorJob :: (Ecosystem -> Maybe Text -> Text -> Either Text PackageName)
-> (Text -> Either Text RegistryUrl) -> Value -> Parser MirrorJob
parseMirrorJob Ecosystem -> Maybe Text -> Text -> Either Text PackageName
packageName Text -> Either Text RegistryUrl
egressUrl = String -> (Object -> Parser MirrorJob) -> Value -> Parser MirrorJob
forall a. String -> (Object -> Parser a) -> Value -> Parser a
withObject String
"MirrorJob" ((Object -> Parser MirrorJob) -> Value -> Parser MirrorJob)
-> (Object -> Parser MirrorJob) -> Value -> Parser MirrorJob
forall a b. (a -> b) -> a -> b
$ \Object
o -> do
    ecoName <- Object
o Object -> Key -> Parser Text
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"ecosystem"
    eco <- maybe (fail (unusable "ecosystem" ecoName)) pure (parseEcosystem ecoName)
    rawNamespace <- o .:? "namespace"
    rawName <- o .: "name"
    -- The payload states the namespace and the base name separately, so the ecosystem re-joins
    -- and re-reads them rather than either field being trusted as given.
    package <- either (fail . toString) pure (packageName eco rawNamespace rawName)
    rawVersion <- o .: "version"
    rawArtifactUrl <- o .: "artifactUrl"
    -- The type the worker's fetch requires cannot be fabricated from an unvalidated string.
    artifactUrl <- either (fail . toString) pure (egressUrl rawArtifactUrl)
    rawFilename <- o .: "filename"
    -- The filename is interpolated into an upstream path, so it is refused here unless it is a
    -- safe path component.
    filename <- maybe (fail (unusable "artifact filename" rawFilename)) pure (mkFilename rawFilename)
    -- A job enqueued with tracing off carries no "traceContext". It yields no span link
    -- rather than a decode failure.
    traceContext <- o .:? "traceContext" >>= traverse parseTraceContext
    pure
        MirrorJob
            { jobPackage = package
            , jobVersion = mkVersion eco rawVersion
            , jobArtifactUrl = artifactUrl
            , jobArtifactFilename = filename
            , jobTraceContext = traceContext
            }

-- The refusal a trust-boundary field reports when its own gate rejects the payload's text.
-- 'String' because that is what aeson's 'fail' takes.
unusable :: String -> Text -> String
unusable :: String -> Text -> String
unusable String
field Text
value = String
"unusable " String -> ShowS
forall a. Semigroup a => a -> a -> a
<> String
field String -> ShowS
forall a. Semigroup a => a -> a -> a
<> String
" " String -> ShowS
forall a. Semigroup a => a -> a -> a
<> Text -> String
forall b a. (Show a, IsString b) => a -> b
show Text
value

-- The carrier is untrusted opaque transport, so both fields are taken as-is. An unparseable
-- W3C value yields no link in the tracing port rather than failing the decode.
parseTraceContext :: Aeson.Value -> Parser RemoteSpanContext
parseTraceContext :: Value -> Parser RemoteSpanContext
parseTraceContext = String
-> (Object -> Parser RemoteSpanContext)
-> Value
-> Parser RemoteSpanContext
forall a. String -> (Object -> Parser a) -> Value -> Parser a
withObject String
"RemoteSpanContext" ((Object -> Parser RemoteSpanContext)
 -> Value -> Parser RemoteSpanContext)
-> (Object -> Parser RemoteSpanContext)
-> Value
-> Parser RemoteSpanContext
forall a b. (a -> b) -> a -> b
$ \Object
t ->
    Text -> Text -> RemoteSpanContext
RemoteSpanContext (Text -> Text -> RemoteSpanContext)
-> Parser Text -> Parser (Text -> RemoteSpanContext)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Object
t Object -> Key -> Parser Text
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"traceparent" Parser (Text -> RemoteSpanContext)
-> Parser Text -> Parser RemoteSpanContext
forall a b. Parser (a -> b) -> Parser a -> Parser b
forall (f :: * -> *) a b. Applicative f => f (a -> b) -> f a -> f b
<*> Object
t Object -> Key -> Parser Text
forall a. FromJSON a => Object -> Key -> Parser a
.: Key
"tracestate"

{- | The backend's own delivery token (an SQS receipt handle, a Pub\/Sub @ackId@). The
constructor is hidden, so worker code only takes one from a 'QueueMessage' that 'receive' returned.
-}
newtype ReceiptHandle = ReceiptHandle Text
    deriving stock (ReceiptHandle -> ReceiptHandle -> Bool
(ReceiptHandle -> ReceiptHandle -> Bool)
-> (ReceiptHandle -> ReceiptHandle -> Bool) -> Eq ReceiptHandle
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: ReceiptHandle -> ReceiptHandle -> Bool
== :: ReceiptHandle -> ReceiptHandle -> Bool
$c/= :: ReceiptHandle -> ReceiptHandle -> Bool
/= :: ReceiptHandle -> ReceiptHandle -> Bool
Eq, Eq ReceiptHandle
Eq ReceiptHandle =>
(ReceiptHandle -> ReceiptHandle -> Ordering)
-> (ReceiptHandle -> ReceiptHandle -> Bool)
-> (ReceiptHandle -> ReceiptHandle -> Bool)
-> (ReceiptHandle -> ReceiptHandle -> Bool)
-> (ReceiptHandle -> ReceiptHandle -> Bool)
-> (ReceiptHandle -> ReceiptHandle -> ReceiptHandle)
-> (ReceiptHandle -> ReceiptHandle -> ReceiptHandle)
-> Ord ReceiptHandle
ReceiptHandle -> ReceiptHandle -> Bool
ReceiptHandle -> ReceiptHandle -> Ordering
ReceiptHandle -> ReceiptHandle -> ReceiptHandle
forall a.
Eq a =>
(a -> a -> Ordering)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> a)
-> (a -> a -> a)
-> Ord a
$ccompare :: ReceiptHandle -> ReceiptHandle -> Ordering
compare :: ReceiptHandle -> ReceiptHandle -> Ordering
$c< :: ReceiptHandle -> ReceiptHandle -> Bool
< :: ReceiptHandle -> ReceiptHandle -> Bool
$c<= :: ReceiptHandle -> ReceiptHandle -> Bool
<= :: ReceiptHandle -> ReceiptHandle -> Bool
$c> :: ReceiptHandle -> ReceiptHandle -> Bool
> :: ReceiptHandle -> ReceiptHandle -> Bool
$c>= :: ReceiptHandle -> ReceiptHandle -> Bool
>= :: ReceiptHandle -> ReceiptHandle -> Bool
$cmax :: ReceiptHandle -> ReceiptHandle -> ReceiptHandle
max :: ReceiptHandle -> ReceiptHandle -> ReceiptHandle
$cmin :: ReceiptHandle -> ReceiptHandle -> ReceiptHandle
min :: ReceiptHandle -> ReceiptHandle -> ReceiptHandle
Ord, Int -> ReceiptHandle -> ShowS
[ReceiptHandle] -> ShowS
ReceiptHandle -> String
(Int -> ReceiptHandle -> ShowS)
-> (ReceiptHandle -> String)
-> ([ReceiptHandle] -> ShowS)
-> Show ReceiptHandle
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> ReceiptHandle -> ShowS
showsPrec :: Int -> ReceiptHandle -> ShowS
$cshow :: ReceiptHandle -> String
show :: ReceiptHandle -> String
$cshowList :: [ReceiptHandle] -> ShowS
showList :: [ReceiptHandle] -> ShowS
Show)

-- | Wrap a backend's delivery token. For backend implementations only.
mkReceiptHandle :: Text -> ReceiptHandle
mkReceiptHandle :: Text -> ReceiptHandle
mkReceiptHandle = Text -> ReceiptHandle
ReceiptHandle

-- | Recover a backend's delivery token to pass back to it. For backend implementations only.
unReceiptHandle :: ReceiptHandle -> Text
unReceiptHandle :: ReceiptHandle -> Text
unReceiptHandle (ReceiptHandle Text
t) = Text
t

-- | A received message: the job to process and the handle that settles this delivery.
data QueueMessage = QueueMessage
    { QueueMessage -> MirrorJob
msgJob :: MirrorJob
    -- ^ The job to process.
    , QueueMessage -> ReceiptHandle
msgReceipt :: ReceiptHandle
    -- ^ The handle identifying this delivery, for 'ack' \/ 'extendVisibility'.
    , QueueMessage -> Int
msgReceiveCount :: Int
    {- ^ Deliveries of this message including this one, @1@ on a first delivery. A backend that
    cannot report a count reports @1@, so only evidence puts a delivery past the 'deliveryBudget'.
    -}
    , QueueMessage -> Maybe ReceiptLease
msgLease :: Maybe ReceiptLease
    {- ^ How long this delivery stays hidden, for the worker's renewal controller.
    'Nothing' from a backend that never expires a delivery.
    -}
    }
    deriving stock (QueueMessage -> QueueMessage -> Bool
(QueueMessage -> QueueMessage -> Bool)
-> (QueueMessage -> QueueMessage -> Bool) -> Eq QueueMessage
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: QueueMessage -> QueueMessage -> Bool
== :: QueueMessage -> QueueMessage -> Bool
$c/= :: QueueMessage -> QueueMessage -> Bool
/= :: QueueMessage -> QueueMessage -> Bool
Eq, Int -> QueueMessage -> ShowS
[QueueMessage] -> ShowS
QueueMessage -> String
(Int -> QueueMessage -> ShowS)
-> (QueueMessage -> String)
-> ([QueueMessage] -> ShowS)
-> Show QueueMessage
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> QueueMessage -> ShowS
showsPrec :: Int -> QueueMessage -> ShowS
$cshow :: QueueMessage -> String
show :: QueueMessage -> String
$cshowList :: [QueueMessage] -> ShowS
showList :: [QueueMessage] -> ShowS
Show)

{- | How many deliveries of one message a queue grants before the worker itself retires it.
A 'newtype', so no caller confuses a count of receives with some other 'Int'.
-}
newtype DeliveryBudget = DeliveryBudget Int
    deriving stock (DeliveryBudget -> DeliveryBudget -> Bool
(DeliveryBudget -> DeliveryBudget -> Bool)
-> (DeliveryBudget -> DeliveryBudget -> Bool) -> Eq DeliveryBudget
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: DeliveryBudget -> DeliveryBudget -> Bool
== :: DeliveryBudget -> DeliveryBudget -> Bool
$c/= :: DeliveryBudget -> DeliveryBudget -> Bool
/= :: DeliveryBudget -> DeliveryBudget -> Bool
Eq, Eq DeliveryBudget
Eq DeliveryBudget =>
(DeliveryBudget -> DeliveryBudget -> Ordering)
-> (DeliveryBudget -> DeliveryBudget -> Bool)
-> (DeliveryBudget -> DeliveryBudget -> Bool)
-> (DeliveryBudget -> DeliveryBudget -> Bool)
-> (DeliveryBudget -> DeliveryBudget -> Bool)
-> (DeliveryBudget -> DeliveryBudget -> DeliveryBudget)
-> (DeliveryBudget -> DeliveryBudget -> DeliveryBudget)
-> Ord DeliveryBudget
DeliveryBudget -> DeliveryBudget -> Bool
DeliveryBudget -> DeliveryBudget -> Ordering
DeliveryBudget -> DeliveryBudget -> DeliveryBudget
forall a.
Eq a =>
(a -> a -> Ordering)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> Bool)
-> (a -> a -> a)
-> (a -> a -> a)
-> Ord a
$ccompare :: DeliveryBudget -> DeliveryBudget -> Ordering
compare :: DeliveryBudget -> DeliveryBudget -> Ordering
$c< :: DeliveryBudget -> DeliveryBudget -> Bool
< :: DeliveryBudget -> DeliveryBudget -> Bool
$c<= :: DeliveryBudget -> DeliveryBudget -> Bool
<= :: DeliveryBudget -> DeliveryBudget -> Bool
$c> :: DeliveryBudget -> DeliveryBudget -> Bool
> :: DeliveryBudget -> DeliveryBudget -> Bool
$c>= :: DeliveryBudget -> DeliveryBudget -> Bool
>= :: DeliveryBudget -> DeliveryBudget -> Bool
$cmax :: DeliveryBudget -> DeliveryBudget -> DeliveryBudget
max :: DeliveryBudget -> DeliveryBudget -> DeliveryBudget
$cmin :: DeliveryBudget -> DeliveryBudget -> DeliveryBudget
min :: DeliveryBudget -> DeliveryBudget -> DeliveryBudget
Ord, Int -> DeliveryBudget -> ShowS
[DeliveryBudget] -> ShowS
DeliveryBudget -> String
(Int -> DeliveryBudget -> ShowS)
-> (DeliveryBudget -> String)
-> ([DeliveryBudget] -> ShowS)
-> Show DeliveryBudget
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> DeliveryBudget -> ShowS
showsPrec :: Int -> DeliveryBudget -> ShowS
$cshow :: DeliveryBudget -> String
show :: DeliveryBudget -> String
$cshowList :: [DeliveryBudget] -> ShowS
showList :: [DeliveryBudget] -> ShowS
Show)

{- | The redelivery budget a backend holds when the operator configures none: five deliveries,
SQS's own redrive convention, pinned to this same value in @config\/default.yaml@.
-}
defaultDeliveryBudget :: DeliveryBudget
defaultDeliveryBudget :: DeliveryBudget
defaultDeliveryBudget = Int -> DeliveryBudget
DeliveryBudget Int
5

{- | Whether a queue has somewhere that captures a message the worker can never mirror. Without
one the message cycles until the queue's retention window discards it unseen.
-}
data DeadLetterTerminus
    = -- | A terminus captures poison messages, at this capture count when the backend reports one.
      TerminusAttached (Maybe DeliveryBudget)
    | -- | Nothing captures poison messages: the worker's budget is the only terminus.
      TerminusAbsent
    deriving stock (DeadLetterTerminus -> DeadLetterTerminus -> Bool
(DeadLetterTerminus -> DeadLetterTerminus -> Bool)
-> (DeadLetterTerminus -> DeadLetterTerminus -> Bool)
-> Eq DeadLetterTerminus
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: DeadLetterTerminus -> DeadLetterTerminus -> Bool
== :: DeadLetterTerminus -> DeadLetterTerminus -> Bool
$c/= :: DeadLetterTerminus -> DeadLetterTerminus -> Bool
/= :: DeadLetterTerminus -> DeadLetterTerminus -> Bool
Eq, Int -> DeadLetterTerminus -> ShowS
[DeadLetterTerminus] -> ShowS
DeadLetterTerminus -> String
(Int -> DeadLetterTerminus -> ShowS)
-> (DeadLetterTerminus -> String)
-> ([DeadLetterTerminus] -> ShowS)
-> Show DeadLetterTerminus
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> DeadLetterTerminus -> ShowS
showsPrec :: Int -> DeadLetterTerminus -> ShowS
$cshow :: DeadLetterTerminus -> String
show :: DeadLetterTerminus -> String
$cshowList :: [DeadLetterTerminus] -> ShowS
showList :: [DeadLetterTerminus] -> ShowS
Show)

{- | The budget the worker enforces: the configured floor, raised past an attached terminus's
capture count, so the dead-letter queue always captures a poison message first.
-}
effectiveDeliveryBudget :: DeliveryBudget -> DeadLetterTerminus -> DeliveryBudget
effectiveDeliveryBudget :: DeliveryBudget -> DeadLetterTerminus -> DeliveryBudget
effectiveDeliveryBudget DeliveryBudget
configured = \case
    TerminusAttached (Just (DeliveryBudget Int
captureAt)) -> DeliveryBudget -> DeliveryBudget -> DeliveryBudget
forall a. Ord a => a -> a -> a
max DeliveryBudget
configured (Int -> DeliveryBudget
DeliveryBudget (Int
captureAt Int -> Int -> Int
forall a. Num a => a -> a -> a
+ Int
1))
    TerminusAttached Maybe DeliveryBudget
Nothing -> DeliveryBudget
configured
    DeadLetterTerminus
TerminusAbsent -> DeliveryBudget
configured

{- | The delivery a budget retires on: the configured value, floored at two, so a first delivery
never spends it. The verdict and the worker's alarm read this one number, so they cannot disagree.
-}
retiringDelivery :: DeliveryBudget -> Int
retiringDelivery :: DeliveryBudget -> Int
retiringDelivery (DeliveryBudget Int
budget) = Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
2 Int
budget

{- | Whether this delivery spends the queue's redelivery budget. A backend supplies the count
('msgReceiveCount'), never the verdict.
-}
deliveryBudgetSpent :: DeliveryBudget -> QueueMessage -> Bool
deliveryBudgetSpent :: DeliveryBudget -> QueueMessage -> Bool
deliveryBudgetSpent DeliveryBudget
budget QueueMessage
message = QueueMessage -> Int
msgReceiveCount QueueMessage
message Int -> Int -> Bool
forall a. Ord a => a -> a -> Bool
>= DeliveryBudget -> Int
retiringDelivery DeliveryBudget
budget

{- | The mirror-queue handle: a record of functions over a backend whose state the closures
capture. Every operation reports failure as an 'Ecluse.Core.Fault.TransportFault' value.
-}
data MirrorQueue = MirrorQueue
    { MirrorQueue -> MirrorJob -> IO (Either TransportFault ())
enqueue :: MirrorJob -> IO (Either TransportFault ())
    {- ^ Producer. Best-effort: the caller logs a 'Left' and never fails the client response,
    since a later pull re-enqueues.
    -}
    , MirrorQueue -> IO (Either TransportFault [QueueMessage])
receive :: IO (Either TransportFault [QueueMessage])
    {- ^ Consumer. One long-poll for a batch, @Right []@ on a healthy empty poll. A 'Left' does
    not advance the liveness heartbeat, so a persistently failing backend surfaces at @\/livez@.
    -}
    , MirrorQueue -> ReceiptHandle -> IO (Either TransportFault ())
ack :: ReceiptHandle -> IO (Either TransportFault ())
    {- ^ Acknowledge a processed message. The caller logs a 'Left' and absorbs it, since
    idempotent publishing makes the repeat harmless.
    -}
    , MirrorQueue
-> ReceiptHandle -> Seconds -> IO (Either TransportFault ())
extendVisibility :: ReceiptHandle -> Seconds -> IO (Either TransportFault ())
    {- ^ Reset a received message's visibility window: renew the worker's lease on it, or at
    zero release it for an immediate redelivery.
    -}
    , MirrorQueue -> ReceiptHandle -> IO (Either TransportFault ())
deadLetter :: ReceiptHandle -> IO (Either TransportFault ())
    {- ^ Realise a terminal fault, routed to the backend's own dead-letter terminus. The caller
    logs a 'Left' and absorbs it, like 'ack'.
    -}
    , MirrorQueue -> DeliveryBudget
deliveryBudget :: DeliveryBudget
    {- ^ The budget 'deliveryBudgetSpent' judges against, settled once at construction with
    'effectiveDeliveryBudget' so judging a delivery costs no per-message work.
    -}
    , MirrorQueue -> Either TransportFault DeadLetterTerminus
deadLetterTerminus :: Either TransportFault DeadLetterTerminus
    {- ^ What the backend's dead-letter probe found at construction, or the fault that stopped
    it. A 'Left' leaves the configured budget standing and never blocks boot.
    -}
    }

{- | The inert queue a deployment with zero mirroring mounts carries, so the composition-root
'Env' keeps its shape. Reached anyway, 'enqueue' refuses with a typed fault rather than crashing.
-}
noMirrorQueue :: MirrorQueue
noMirrorQueue :: MirrorQueue
noMirrorQueue =
    MirrorQueue
        { enqueue :: MirrorJob -> IO (Either TransportFault ())
enqueue = \MirrorJob
_ -> Either TransportFault () -> IO (Either TransportFault ())
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (TransportFault -> Either TransportFault ()
forall a b. a -> Either a b
Left TransportFault
inertFault)
        , receive :: IO (Either TransportFault [QueueMessage])
receive = Either TransportFault [QueueMessage]
-> IO (Either TransportFault [QueueMessage])
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure ([QueueMessage] -> Either TransportFault [QueueMessage]
forall a b. b -> Either a b
Right [])
        , ack :: ReceiptHandle -> IO (Either TransportFault ())
ack = \ReceiptHandle
_ -> 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 ())
        , extendVisibility :: ReceiptHandle -> Seconds -> IO (Either TransportFault ())
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 ())
        , deadLetter :: ReceiptHandle -> IO (Either TransportFault ())
deadLetter = \ReceiptHandle
_ -> 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 ())
        , deliveryBudget :: DeliveryBudget
deliveryBudget = DeliveryBudget
defaultDeliveryBudget
        , deadLetterTerminus :: Either TransportFault DeadLetterTerminus
deadLetterTerminus = DeadLetterTerminus -> Either TransportFault DeadLetterTerminus
forall a b. b -> Either a b
Right DeadLetterTerminus
TerminusAbsent
        }
  where
    inertFault :: TransportFault
inertFault = TransportCause -> Text -> TransportFault
transportFault TransportCause
TransportProtocol Text
"no mount mirrors, so no mirror queue is built"