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

{- | The tuning fields, the message lifting and the dead-letter probe behind
"Ecluse.Runtime.Queue.Sqs", which documents the backend and re-exports the curated surface.
Importing this module opts out of that stability promise, the convention @text@ and @bytestring@
use, so production code imports the public one.
-}
module Ecluse.Runtime.Queue.Sqs.Internal (
    -- * Configuration
    SqsConfig (..),
    defaultSqsConfig,

    -- * The backend
    newSqsQueue,

    -- * Received-message lifting
    ReceivedMessage (..),
    liftReceivedMessages,

    -- * Dead-letter probe
    deadLetterTerminusOf,

    -- * The queue payload's wire-name gate
    mirrorJobPackage,
) where

import Amazonka qualified as AWS

import Amazonka.SQS.ChangeMessageVisibility qualified as SQS
import Amazonka.SQS.DeleteMessage qualified as SQS
import Amazonka.SQS.GetQueueAttributes qualified as SQS
import Amazonka.SQS.ReceiveMessage qualified as SQS
import Amazonka.SQS.SendMessage qualified as SQS
import Amazonka.SQS.Types qualified as SQS
import Data.Aeson (eitherDecodeStrict', withObject, (.:))
import Data.Aeson qualified as Aeson
import Data.Aeson.Types (parseMaybe)
import Katip (LogEnv, Severity (DebugS), sl)
import Lens.Micro ((?~), (^.))

import Ecluse.Core.Ecosystem (Ecosystem (Npm, PyPI, RubyGems))
import Ecluse.Core.Fault (TransportFault)
import Ecluse.Core.Package (PackageName, mkPackageName, mkScope)
import Ecluse.Core.Queue (
    DeadLetterTerminus (TerminusAbsent, TerminusAttached),
    DeliveryBudget (DeliveryBudget),
    MirrorQueue (..),
    QueueMessage (..),
    Seconds (..),
    decodeJob,
    defaultDeliveryBudget,
    effectiveDeliveryBudget,
    encodeJob,
    mkReceiptHandle,
    unReceiptHandle,
 )
import Ecluse.Core.Queue.Lease (ReceiptLease, monotonicNow, receiptLease)
import Ecluse.Core.Registry (parseErrorMessage)
import Ecluse.Core.Registry.Npm.Project (projectName)
import Ecluse.Core.Security.Egress (RegistryUrl)
import Ecluse.Core.Text (nonBlank, readDecimalText)
import Ecluse.Runtime.Aws.Env (AwsEndpoint, newAwsEnv)
import Ecluse.Runtime.Aws.Fault (classifyAwsTransport, sendClassified)
import Ecluse.Runtime.Log (logLine, moduleField)

{- | What the SQS backend needs. The provider knobs take their defaults from
'defaultSqsConfig' (see "Ecluse.Core.Queue").
-}
data SqsConfig = SqsConfig
    { SqsConfig -> Text
sqsQueueUrl :: Text
    -- ^ The fully-qualified SQS queue URL mirror jobs are sent to and received from.
    , SqsConfig -> Text
sqsRegion :: Text
    -- ^ The AWS region the queue lives in (e.g. @"us-east-1"@).
    , SqsConfig -> Maybe AwsEndpoint
sqsEndpoint :: Maybe AwsEndpoint
    {- ^ An endpoint override for an emulator or VPC endpoint. 'Nothing' uses
    @amazonka@'s default resolution and the ambient credential chain.
    -}
    , SqsConfig -> Int
sqsBatchSize :: Int
    -- ^ Maximum messages to pull per 'receive' (SQS caps this at 10).
    , SqsConfig -> Int
sqsWaitSeconds :: Int
    {- ^ The long-poll window in seconds (SQS caps this at 20). A 'receive' waits this long
    for a message before returning @[]@, so an idle worker does not hot-loop on empty polls.
    -}
    , SqsConfig -> Seconds
sqsVisibilityTimeout :: Seconds
    {- ^ How long a received message stays hidden before SQS redelivers it: the budget for
    processing-then-'ack', extendable per message via 'extendVisibility'.
    -}
    , SqsConfig -> Seconds
sqsTerminalBackoff :: Seconds
    {- ^ The visibility timeout 'deadLetter' applies to a __terminal__ message. It exceeds the
    processing window, so the worker does not re-fetch an unmirrorable artifact in a hot loop.
    -}
    , SqsConfig -> DeliveryBudget
sqsMaxReceiveCount :: DeliveryBudget
    {- ^ The configured __floor__ on deliveries before the worker retires a message. 'newSqsQueue'
    raises the budget past a redrive policy's count, so it never pre-empts a dead-letter capture.
    -}
    }
    deriving stock (SqsConfig -> SqsConfig -> Bool
(SqsConfig -> SqsConfig -> Bool)
-> (SqsConfig -> SqsConfig -> Bool) -> Eq SqsConfig
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: SqsConfig -> SqsConfig -> Bool
== :: SqsConfig -> SqsConfig -> Bool
$c/= :: SqsConfig -> SqsConfig -> Bool
/= :: SqsConfig -> SqsConfig -> Bool
Eq, Int -> SqsConfig -> ShowS
[SqsConfig] -> ShowS
SqsConfig -> String
(Int -> SqsConfig -> ShowS)
-> (SqsConfig -> String)
-> ([SqsConfig] -> ShowS)
-> Show SqsConfig
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> SqsConfig -> ShowS
showsPrec :: Int -> SqsConfig -> ShowS
$cshow :: SqsConfig -> String
show :: SqsConfig -> String
$cshowList :: [SqsConfig] -> ShowS
showList :: [SqsConfig] -> ShowS
Show)

-- | A 'SqsConfig' for a queue URL and region, with the provider knobs at their defaults.
defaultSqsConfig :: Text -> Text -> SqsConfig
defaultSqsConfig :: Text -> Text -> SqsConfig
defaultSqsConfig Text
queueUrl Text
region =
    SqsConfig
        { sqsQueueUrl :: Text
sqsQueueUrl = Text
queueUrl
        , sqsRegion :: Text
sqsRegion = Text
region
        , sqsEndpoint :: Maybe AwsEndpoint
sqsEndpoint = Maybe AwsEndpoint
forall a. Maybe a
Nothing
        , sqsBatchSize :: Int
sqsBatchSize = Int
10
        , sqsWaitSeconds :: Int
sqsWaitSeconds = Int
20
        , sqsVisibilityTimeout :: Seconds
sqsVisibilityTimeout = Int -> Seconds
Seconds Int
30
        , sqsTerminalBackoff :: Seconds
sqsTerminalBackoff = Int -> Seconds
Seconds Int
300
        , sqsMaxReceiveCount :: DeliveryBudget
sqsMaxReceiveCount = DeliveryBudget
defaultDeliveryBudget
        }

{- | Build an SQS-backed 'MirrorQueue'. The @amazonka@ 'AWS.Env' is built once here and
the returned handle's closures capture it.
-}
newSqsQueue :: LogEnv -> (Text -> Either Text RegistryUrl) -> SqsConfig -> IO MirrorQueue
newSqsQueue :: LogEnv
-> (Text -> Either Text RegistryUrl) -> SqsConfig -> IO MirrorQueue
newSqsQueue LogEnv
logEnv Text -> Either Text RegistryUrl
egressUrl SqsConfig
cfg = do
    env <- SqsConfig -> IO Env
mkEnv SqsConfig
cfg
    -- Every operation reports its failure as the handle's 'TransportFault' value, never
    -- through the exception channel.
    let run :: (AWS.AWSRequest a) => a -> IO (Either TransportFault (AWS.AWSResponse a))
        run = (Error -> TransportFault)
-> Env -> a -> IO (Either TransportFault (AWSResponse a))
forall a e.
AWSRequest a =>
(Error -> e) -> Env -> a -> IO (Either e (AWSResponse a))
sendClassified Error -> TransportFault
classifyAwsTransport Env
env
        queueUrl = SqsConfig -> Text
sqsQueueUrl SqsConfig
cfg
        Seconds terminalBackoffSecs = sqsTerminalBackoff cfg
    -- SQS keeps the redrive configuration on the queue, not on a message, so the backend
    -- probes it once at boot and falls back to the configured floor when the probe faults.
    terminus <- fmap terminusOfResponse <$> run (terminusRequest queueUrl)
    let budget = (TransportFault -> DeliveryBudget)
-> (DeadLetterTerminus -> DeliveryBudget)
-> Either TransportFault DeadLetterTerminus
-> DeliveryBudget
forall a c b. (a -> c) -> (b -> c) -> Either a b -> c
either (DeliveryBudget -> TransportFault -> DeliveryBudget
forall a b. a -> b -> a
const (SqsConfig -> DeliveryBudget
sqsMaxReceiveCount SqsConfig
cfg)) (DeliveryBudget -> DeadLetterTerminus -> DeliveryBudget
effectiveDeliveryBudget (SqsConfig -> DeliveryBudget
sqsMaxReceiveCount SqsConfig
cfg)) Either TransportFault DeadLetterTerminus
terminus
    pure
        MirrorQueue
            { enqueue = fmap void . run . SQS.newSendMessage queueUrl . encodeJob
            , receive = do
                -- Stamped before the request, so a lease never outlives the window SQS granted.
                receivedAt <- monotonicNow
                let lease = MonoTime -> Seconds -> Seconds -> ReceiptLease
receiptLease MonoTime
receivedAt (SqsConfig -> Seconds
sqsVisibilityTimeout SqsConfig
cfg) Seconds
sqsInFlightMaximum
                outcome <- run (receiveRequest cfg)
                traverse (liftReceivedMessages logEnv egressUrl lease . receivedMessages) outcome
            , ack = fmap void . run . SQS.newDeleteMessage queueUrl . unReceiptHandle
            , extendVisibility = \ReceiptHandle
receipt (Seconds Int
secs) ->
                (Either TransportFault ChangeMessageVisibilityResponse
 -> Either TransportFault ())
-> IO (Either TransportFault ChangeMessageVisibilityResponse)
-> IO (Either TransportFault ())
forall a b. (a -> b) -> IO a -> IO b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap Either TransportFault ChangeMessageVisibilityResponse
-> Either TransportFault ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO (Either TransportFault ChangeMessageVisibilityResponse)
 -> IO (Either TransportFault ()))
-> (ChangeMessageVisibility
    -> IO (Either TransportFault ChangeMessageVisibilityResponse))
-> ChangeMessageVisibility
-> IO (Either TransportFault ())
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ChangeMessageVisibility
-> IO (Either TransportFault (AWSResponse ChangeMessageVisibility))
ChangeMessageVisibility
-> IO (Either TransportFault ChangeMessageVisibilityResponse)
forall a.
AWSRequest a =>
a -> IO (Either TransportFault (AWSResponse a))
run (ChangeMessageVisibility -> IO (Either TransportFault ()))
-> ChangeMessageVisibility -> IO (Either TransportFault ())
forall a b. (a -> b) -> a -> b
$
                    Text -> Text -> Int -> ChangeMessageVisibility
SQS.newChangeMessageVisibility Text
queueUrl (ReceiptHandle -> Text
unReceiptHandle ReceiptHandle
receipt) Int
secs
            , -- A terminal fault: return the message with the backoff visibility timeout
              -- (@ChangeMessageVisibility@), __never__ @DeleteMessage@, so the message is not
              -- discarded but rides the operator's redrive policy to the dead-letter queue.
              deadLetter = \ReceiptHandle
receipt ->
                (Either TransportFault ChangeMessageVisibilityResponse
 -> Either TransportFault ())
-> IO (Either TransportFault ChangeMessageVisibilityResponse)
-> IO (Either TransportFault ())
forall a b. (a -> b) -> IO a -> IO b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap Either TransportFault ChangeMessageVisibilityResponse
-> Either TransportFault ()
forall (f :: * -> *) a. Functor f => f a -> f ()
void (IO (Either TransportFault ChangeMessageVisibilityResponse)
 -> IO (Either TransportFault ()))
-> (ChangeMessageVisibility
    -> IO (Either TransportFault ChangeMessageVisibilityResponse))
-> ChangeMessageVisibility
-> IO (Either TransportFault ())
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ChangeMessageVisibility
-> IO (Either TransportFault (AWSResponse ChangeMessageVisibility))
ChangeMessageVisibility
-> IO (Either TransportFault ChangeMessageVisibilityResponse)
forall a.
AWSRequest a =>
a -> IO (Either TransportFault (AWSResponse a))
run (ChangeMessageVisibility -> IO (Either TransportFault ()))
-> ChangeMessageVisibility -> IO (Either TransportFault ())
forall a b. (a -> b) -> a -> b
$
                    Text -> Text -> Int -> ChangeMessageVisibility
SQS.newChangeMessageVisibility Text
queueUrl (ReceiptHandle -> Text
unReceiptHandle ReceiptHandle
receipt) Int
terminalBackoffSecs
            , deliveryBudget = budget
            , deadLetterTerminus = terminus
            }

-- Build the region-scoped, optionally endpoint-overridden amazonka environment.
mkEnv :: SqsConfig -> IO AWS.Env
mkEnv :: SqsConfig -> IO Env
mkEnv SqsConfig
cfg = Maybe Text -> Maybe AwsEndpoint -> Service -> IO Env
newAwsEnv (Text -> Maybe Text
forall a. a -> Maybe a
Just (SqsConfig -> Text
sqsRegion SqsConfig
cfg)) (SqsConfig -> Maybe AwsEndpoint
sqsEndpoint SqsConfig
cfg) Service
SQS.defaultService

-- SQS caps the long poll at 20s and clamps a larger configured wait. That stays inside
-- @amazonka@'s default request timeout, so the client needs no response-timeout override.
receiveRequest :: SqsConfig -> SQS.ReceiveMessage
receiveRequest :: SqsConfig -> ReceiveMessage
receiveRequest SqsConfig
cfg =
    Text -> ReceiveMessage
SQS.newReceiveMessage (SqsConfig -> Text
sqsQueueUrl SqsConfig
cfg)
        ReceiveMessage
-> (ReceiveMessage -> ReceiveMessage) -> ReceiveMessage
forall a b. a -> (a -> b) -> b
& (Maybe Int -> Identity (Maybe Int))
-> ReceiveMessage -> Identity ReceiveMessage
Lens' ReceiveMessage (Maybe Int)
SQS.receiveMessage_maxNumberOfMessages
        ((Maybe Int -> Identity (Maybe Int))
 -> ReceiveMessage -> Identity ReceiveMessage)
-> Int -> ReceiveMessage -> ReceiveMessage
forall s t a b. ASetter s t a (Maybe b) -> b -> s -> t
?~ SqsConfig -> Int
sqsBatchSize SqsConfig
cfg
            ReceiveMessage
-> (ReceiveMessage -> ReceiveMessage) -> ReceiveMessage
forall a b. a -> (a -> b) -> b
& (Maybe Int -> Identity (Maybe Int))
-> ReceiveMessage -> Identity ReceiveMessage
Lens' ReceiveMessage (Maybe Int)
SQS.receiveMessage_waitTimeSeconds
        ((Maybe Int -> Identity (Maybe Int))
 -> ReceiveMessage -> Identity ReceiveMessage)
-> Int -> ReceiveMessage -> ReceiveMessage
forall s t a b. ASetter s t a (Maybe b) -> b -> s -> t
?~ SqsConfig -> Int
sqsWaitSeconds SqsConfig
cfg
            ReceiveMessage
-> (ReceiveMessage -> ReceiveMessage) -> ReceiveMessage
forall a b. a -> (a -> b) -> b
& (Maybe Int -> Identity (Maybe Int))
-> ReceiveMessage -> Identity ReceiveMessage
Lens' ReceiveMessage (Maybe Int)
SQS.receiveMessage_visibilityTimeout
        ((Maybe Int -> Identity (Maybe Int))
 -> ReceiveMessage -> Identity ReceiveMessage)
-> Int -> ReceiveMessage -> ReceiveMessage
forall s t a b. ASetter s t a (Maybe b) -> b -> s -> t
?~ Int
visibilitySeconds
            ReceiveMessage
-> (ReceiveMessage -> ReceiveMessage) -> ReceiveMessage
forall a b. a -> (a -> b) -> b
& (Maybe [MessageAttribute] -> Identity (Maybe [MessageAttribute]))
-> ReceiveMessage -> Identity ReceiveMessage
Lens' ReceiveMessage (Maybe [MessageAttribute])
SQS.receiveMessage_attributeNames
        -- Asked for explicitly: SQS omits the delivery count unless a request names
        -- it. The redelivery budget judges a delivery by that count.
        ((Maybe [MessageAttribute] -> Identity (Maybe [MessageAttribute]))
 -> ReceiveMessage -> Identity ReceiveMessage)
-> [MessageAttribute] -> ReceiveMessage -> ReceiveMessage
forall s t a b. ASetter s t a (Maybe b) -> b -> s -> t
?~ [MessageAttribute
SQS.MessageAttribute_ApproximateReceiveCount]
  where
    Seconds Int
visibilitySeconds = SqsConfig -> Seconds
sqsVisibilityTimeout SqsConfig
cfg

{- The longest SQS keeps one receipt in flight, however often its visibility is renewed. A
renewal never asks past it, so a lease that reaches it is dropped rather than silently lapsing. -}
sqsInFlightMaximum :: Seconds
sqsInFlightMaximum :: Seconds
sqsInFlightMaximum = Int -> Seconds
Seconds Int
43_200

-- The boot-time redrive probe.
terminusRequest :: Text -> SQS.GetQueueAttributes
terminusRequest :: Text -> GetQueueAttributes
terminusRequest Text
queueUrl =
    Text -> GetQueueAttributes
SQS.newGetQueueAttributes Text
queueUrl
        GetQueueAttributes
-> (GetQueueAttributes -> GetQueueAttributes) -> GetQueueAttributes
forall a b. a -> (a -> b) -> b
& (Maybe [QueueAttributeName]
 -> Identity (Maybe [QueueAttributeName]))
-> GetQueueAttributes -> Identity GetQueueAttributes
Lens' GetQueueAttributes (Maybe [QueueAttributeName])
SQS.getQueueAttributes_attributeNames
        ((Maybe [QueueAttributeName]
  -> Identity (Maybe [QueueAttributeName]))
 -> GetQueueAttributes -> Identity GetQueueAttributes)
-> [QueueAttributeName] -> GetQueueAttributes -> GetQueueAttributes
forall s t a b. ASetter s t a (Maybe b) -> b -> s -> t
?~ [QueueAttributeName
SQS.QueueAttributeName_RedrivePolicy]

terminusOfResponse :: SQS.GetQueueAttributesResponse -> DeadLetterTerminus
terminusOfResponse :: GetQueueAttributesResponse -> DeadLetterTerminus
terminusOfResponse GetQueueAttributesResponse
response =
    Maybe Text -> DeadLetterTerminus
deadLetterTerminusOf (HashMap QueueAttributeName Text -> Maybe Text
forall (t :: * -> *). Foldable t => t Text -> Maybe Text
soleValue (HashMap QueueAttributeName Text -> Maybe Text)
-> Maybe (HashMap QueueAttributeName Text) -> Maybe Text
forall (m :: * -> *) a b. Monad m => (a -> m b) -> m a -> m b
=<< (GetQueueAttributesResponse
response GetQueueAttributesResponse
-> Getting
     (Maybe (HashMap QueueAttributeName Text))
     GetQueueAttributesResponse
     (Maybe (HashMap QueueAttributeName Text))
-> Maybe (HashMap QueueAttributeName Text)
forall s a. s -> Getting a s a -> a
^. Getting
  (Maybe (HashMap QueueAttributeName Text))
  GetQueueAttributesResponse
  (Maybe (HashMap QueueAttributeName Text))
Lens'
  GetQueueAttributesResponse
  (Maybe (HashMap QueueAttributeName Text))
SQS.getQueueAttributesResponse_attributes))

-- Every request here names exactly one attribute and SQS returns only the named
-- attributes that are set, so an empty map means unset, not an ambiguous pick.
soleValue :: (Foldable t) => t Text -> Maybe Text
soleValue :: forall (t :: * -> *). Foldable t => t Text -> Maybe Text
soleValue = [Text] -> Maybe Text
forall a. [a] -> Maybe a
listToMaybe ([Text] -> Maybe Text)
-> (t Text -> [Text]) -> t Text -> Maybe Text
forall b c a. (b -> c) -> (a -> b) -> a -> c
. t Text -> [Text]
forall a. t a -> [a]
forall (t :: * -> *) a. Foldable t => t a -> [a]
toList

{- | Classify a queue's raw @RedrivePolicy@ into its dead-letter terminus. No policy or a
blank one is absent, and any other is attached, with @maxReceiveCount@ only when it reads.
-}
deadLetterTerminusOf :: Maybe Text -> DeadLetterTerminus
deadLetterTerminusOf :: Maybe Text -> DeadLetterTerminus
deadLetterTerminusOf Maybe Text
raw =
    DeadLetterTerminus
-> (Text -> DeadLetterTerminus) -> Maybe Text -> DeadLetterTerminus
forall b a. b -> (a -> b) -> Maybe a -> b
maybe DeadLetterTerminus
TerminusAbsent (Maybe DeliveryBudget -> DeadLetterTerminus
TerminusAttached (Maybe DeliveryBudget -> DeadLetterTerminus)
-> (Text -> Maybe DeliveryBudget) -> Text -> DeadLetterTerminus
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Text -> Maybe DeliveryBudget
captureCountOf) (Text -> Maybe Text
nonBlank (Text -> Maybe Text) -> Maybe Text -> Maybe Text
forall (m :: * -> *) a b. Monad m => (a -> m b) -> m a -> m b
=<< Maybe Text
raw)

-- The @maxReceiveCount@ an attached redrive policy declares. AWS embeds the policy as
-- JSON and renders the count as either a string or a number, so both parse.
captureCountOf :: Text -> Maybe DeliveryBudget
captureCountOf :: Text -> Maybe DeliveryBudget
captureCountOf Text
policy = do
    value <- Either String Value -> Maybe Value
forall l r. Either l r -> Maybe r
rightToMaybe (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
policy))
    raw <- parseMaybe (withObject "RedrivePolicy" (.: "maxReceiveCount")) value
    DeliveryBudget <$> countOf raw

countOf :: Aeson.Value -> Maybe Int
countOf :: Value -> Maybe Int
countOf = \case
    Aeson.String Text
text -> Text -> Maybe Int
forall a. Integral a => Text -> Maybe a
readDecimalText Text
text
    Value
other -> case Value -> Result Int
forall a. FromJSON a => Value -> Result a
Aeson.fromJSON Value
other of
        Aeson.Success Int
count -> Int -> Maybe Int
forall a. a -> Maybe a
Just Int
count
        Aeson.Error String
_ -> Maybe Int
forall a. Maybe a
Nothing

{- | The fields of a received SQS message the backend reads. Lifting them out of the
@amazonka@ 'SQS.Message' keeps the 'QueueMessage' mapping free of the AWS type.
-}
data ReceivedMessage = ReceivedMessage
    { ReceivedMessage -> Maybe Text
rmBody :: Maybe Text
    -- ^ The message body carrying the encoded job (SQS always supplies one).
    , ReceivedMessage -> Maybe Text
rmReceipt :: Maybe Text
    -- ^ The receipt handle a later 'ack' deletes the message by (SQS always supplies one).
    , ReceivedMessage -> Maybe Text
rmMessageId :: Maybe Text
    -- ^ The SQS-assigned message id, for the drop log. Not part of the untrusted body.
    , ReceivedMessage -> Maybe Text
rmReceiveCount :: Maybe Text
    {- ^ The raw @ApproximateReceiveCount@ system attribute, which the poll asks for
    explicitly. 'Nothing' when SQS supplied none, which reads as a first delivery.
    -}
    }
    deriving stock (ReceivedMessage -> ReceivedMessage -> Bool
(ReceivedMessage -> ReceivedMessage -> Bool)
-> (ReceivedMessage -> ReceivedMessage -> Bool)
-> Eq ReceivedMessage
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: ReceivedMessage -> ReceivedMessage -> Bool
== :: ReceivedMessage -> ReceivedMessage -> Bool
$c/= :: ReceivedMessage -> ReceivedMessage -> Bool
/= :: ReceivedMessage -> ReceivedMessage -> Bool
Eq, Int -> ReceivedMessage -> ShowS
[ReceivedMessage] -> ShowS
ReceivedMessage -> String
(Int -> ReceivedMessage -> ShowS)
-> (ReceivedMessage -> String)
-> ([ReceivedMessage] -> ShowS)
-> Show ReceivedMessage
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> ReceivedMessage -> ShowS
showsPrec :: Int -> ReceivedMessage -> ShowS
$cshow :: ReceivedMessage -> String
show :: ReceivedMessage -> String
$cshowList :: [ReceivedMessage] -> ShowS
showList :: [ReceivedMessage] -> ShowS
Show)

-- The read fields of an amazonka Message, lifted at the effectful edge ('receive').
receivedFields :: SQS.Message -> ReceivedMessage
receivedFields :: Message -> ReceivedMessage
receivedFields Message
message =
    ReceivedMessage
        { rmBody :: Maybe Text
rmBody = Message
message Message -> Getting (Maybe Text) Message (Maybe Text) -> Maybe Text
forall s a. s -> Getting a s a -> a
^. Getting (Maybe Text) Message (Maybe Text)
Lens' Message (Maybe Text)
SQS.message_body
        , rmReceipt :: Maybe Text
rmReceipt = Message
message Message -> Getting (Maybe Text) Message (Maybe Text) -> Maybe Text
forall s a. s -> Getting a s a -> a
^. Getting (Maybe Text) Message (Maybe Text)
Lens' Message (Maybe Text)
SQS.message_receiptHandle
        , rmMessageId :: Maybe Text
rmMessageId = Message
message Message -> Getting (Maybe Text) Message (Maybe Text) -> Maybe Text
forall s a. s -> Getting a s a -> a
^. Getting (Maybe Text) Message (Maybe Text)
Lens' Message (Maybe Text)
SQS.message_messageId
        , rmReceiveCount :: Maybe Text
rmReceiveCount = HashMap MessageAttribute Text -> Maybe Text
forall (t :: * -> *). Foldable t => t Text -> Maybe Text
soleValue (HashMap MessageAttribute Text -> Maybe Text)
-> Maybe (HashMap MessageAttribute Text) -> Maybe Text
forall (m :: * -> *) a b. Monad m => (a -> m b) -> m a -> m b
=<< (Message
message Message
-> Getting
     (Maybe (HashMap MessageAttribute Text))
     Message
     (Maybe (HashMap MessageAttribute Text))
-> Maybe (HashMap MessageAttribute Text)
forall s a. s -> Getting a s a -> a
^. Getting
  (Maybe (HashMap MessageAttribute Text))
  Message
  (Maybe (HashMap MessageAttribute Text))
Lens' Message (Maybe (HashMap MessageAttribute Text))
SQS.message_attributes)
        }

-- The received batch's messages, each reduced to the fields the backend reads.
receivedMessages :: SQS.ReceiveMessageResponse -> [ReceivedMessage]
receivedMessages :: ReceiveMessageResponse -> [ReceivedMessage]
receivedMessages ReceiveMessageResponse
response =
    [ReceivedMessage]
-> ([Message] -> [ReceivedMessage])
-> Maybe [Message]
-> [ReceivedMessage]
forall b a. b -> (a -> b) -> Maybe a -> b
maybe [] ((Message -> ReceivedMessage) -> [Message] -> [ReceivedMessage]
forall a b. (a -> b) -> [a] -> [b]
map Message -> ReceivedMessage
receivedFields) (ReceiveMessageResponse
response ReceiveMessageResponse
-> Getting
     (Maybe [Message]) ReceiveMessageResponse (Maybe [Message])
-> Maybe [Message]
forall s a. s -> Getting a s a -> a
^. Getting (Maybe [Message]) ReceiveMessageResponse (Maybe [Message])
Lens' ReceiveMessageResponse (Maybe [Message])
SQS.receiveMessageResponse_messages)

-- Why a received message could not become a QueueMessage. A closed set with no
-- payload, so a drop log never echoes any of the (untrusted) message contents.
data SqsDropReason = MissingBody | MissingReceipt | UndecodableBody
    deriving stock (SqsDropReason -> SqsDropReason -> Bool
(SqsDropReason -> SqsDropReason -> Bool)
-> (SqsDropReason -> SqsDropReason -> Bool) -> Eq SqsDropReason
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: SqsDropReason -> SqsDropReason -> Bool
== :: SqsDropReason -> SqsDropReason -> Bool
$c/= :: SqsDropReason -> SqsDropReason -> Bool
/= :: SqsDropReason -> SqsDropReason -> Bool
Eq, Int -> SqsDropReason -> ShowS
[SqsDropReason] -> ShowS
SqsDropReason -> String
(Int -> SqsDropReason -> ShowS)
-> (SqsDropReason -> String)
-> ([SqsDropReason] -> ShowS)
-> Show SqsDropReason
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> SqsDropReason -> ShowS
showsPrec :: Int -> SqsDropReason -> ShowS
$cshow :: SqsDropReason -> String
show :: SqsDropReason -> String
$cshowList :: [SqsDropReason] -> ShowS
showList :: [SqsDropReason] -> ShowS
Show)

{- A message missing its body or receipt, which SQS always supplies, or one whose body
does not decode, is dropped rather than crashing the poll. -}
toQueueMessage :: (Text -> Either Text RegistryUrl) -> ReceiptLease -> ReceivedMessage -> Either SqsDropReason QueueMessage
toQueueMessage :: (Text -> Either Text RegistryUrl)
-> ReceiptLease
-> ReceivedMessage
-> Either SqsDropReason QueueMessage
toQueueMessage Text -> Either Text RegistryUrl
egressUrl ReceiptLease
lease ReceivedMessage
received = do
    body <- SqsDropReason -> Maybe Text -> Either SqsDropReason Text
forall l r. l -> Maybe r -> Either l r
maybeToRight SqsDropReason
MissingBody (ReceivedMessage -> Maybe Text
rmBody ReceivedMessage
received)
    receipt <- maybeToRight MissingReceipt (rmReceipt received)
    job <- first (const UndecodableBody) (decodeJob mirrorJobPackage egressUrl body)
    pure
        QueueMessage
            { msgJob = job
            , msgReceipt = mkReceiptHandle receipt
            , msgReceiveCount = receiveCountOf (rmReceiveCount received)
            , msgLease = Just lease
            }

-- The delivery count SQS reported, or a first delivery when it reported none, an
-- unreadable one, or a count below one. Only evidence ever puts a message past its budget.
receiveCountOf :: Maybe Text -> Int
receiveCountOf :: Maybe Text -> Int
receiveCountOf Maybe Text
raw = Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
1 (Int -> Maybe Int -> Int
forall a. a -> Maybe a -> a
fromMaybe Int
1 (Text -> Maybe Int
forall a. Integral a => Text -> Maybe a
readDecimalText (Text -> Maybe Int) -> Maybe Text -> Maybe Int
forall (m :: * -> *) a b. Monad m => (a -> m b) -> m a -> m b
=<< Maybe Text
raw))

{- | Lift a received batch into deliverable 'QueueMessage's under one poll's lease. A drop is
logged, omitted, and left un-'ack'ed, so redelivery and dead-letter behaviour are unchanged.
-}
liftReceivedMessages :: LogEnv -> (Text -> Either Text RegistryUrl) -> ReceiptLease -> [ReceivedMessage] -> IO [QueueMessage]
liftReceivedMessages :: LogEnv
-> (Text -> Either Text RegistryUrl)
-> ReceiptLease
-> [ReceivedMessage]
-> IO [QueueMessage]
liftReceivedMessages LogEnv
logEnv Text -> Either Text RegistryUrl
egressUrl ReceiptLease
lease =
    ([Maybe QueueMessage] -> [QueueMessage])
-> IO [Maybe QueueMessage] -> IO [QueueMessage]
forall a b. (a -> b) -> IO a -> IO b
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
fmap [Maybe QueueMessage] -> [QueueMessage]
forall a. [Maybe a] -> [a]
catMaybes (IO [Maybe QueueMessage] -> IO [QueueMessage])
-> ([ReceivedMessage] -> IO [Maybe QueueMessage])
-> [ReceivedMessage]
-> IO [QueueMessage]
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (ReceivedMessage -> IO (Maybe QueueMessage))
-> [ReceivedMessage] -> IO [Maybe QueueMessage]
forall (t :: * -> *) (f :: * -> *) a b.
(Traversable t, Applicative f) =>
(a -> f b) -> t a -> f (t b)
forall (f :: * -> *) a b.
Applicative f =>
(a -> f b) -> [a] -> f [b]
traverse (LogEnv
-> (Text -> Either Text RegistryUrl)
-> ReceiptLease
-> ReceivedMessage
-> IO (Maybe QueueMessage)
liftReceivedMessage LogEnv
logEnv Text -> Either Text RegistryUrl
egressUrl ReceiptLease
lease)

-- Deliver a received message, or log the drop at DebugS and yield Nothing.
liftReceivedMessage :: LogEnv -> (Text -> Either Text RegistryUrl) -> ReceiptLease -> ReceivedMessage -> IO (Maybe QueueMessage)
liftReceivedMessage :: LogEnv
-> (Text -> Either Text RegistryUrl)
-> ReceiptLease
-> ReceivedMessage
-> IO (Maybe QueueMessage)
liftReceivedMessage LogEnv
logEnv Text -> Either Text RegistryUrl
egressUrl ReceiptLease
lease ReceivedMessage
received =
    case (Text -> Either Text RegistryUrl)
-> ReceiptLease
-> ReceivedMessage
-> Either SqsDropReason QueueMessage
toQueueMessage Text -> Either Text RegistryUrl
egressUrl ReceiptLease
lease ReceivedMessage
received of
        Right QueueMessage
queueMessage -> Maybe QueueMessage -> IO (Maybe QueueMessage)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure (QueueMessage -> Maybe QueueMessage
forall a. a -> Maybe a
Just QueueMessage
queueMessage)
        Left SqsDropReason
reason -> Maybe QueueMessage
forall a. Maybe a
Nothing Maybe QueueMessage -> IO () -> IO (Maybe QueueMessage)
forall a b. a -> IO b -> IO a
forall (f :: * -> *) a b. Functor f => a -> f b -> f a
<$ LogEnv -> SqsDropReason -> Maybe Text -> IO ()
logSqsDrop LogEnv
logEnv SqsDropReason
reason (ReceivedMessage -> Maybe Text
rmMessageId ReceivedMessage
received)

-- One DebugS line naming why a received message was dropped, and its SQS message id when
-- present. The body is untrusted payload, and @module@ names the public module operators filter on.
logSqsDrop :: LogEnv -> SqsDropReason -> Maybe Text -> IO ()
logSqsDrop :: LogEnv -> SqsDropReason -> Maybe Text -> IO ()
logSqsDrop LogEnv
logEnv SqsDropReason
reason Maybe Text
messageId =
    LogEnv -> SimpleLogPayload -> Severity -> Text -> IO ()
logLine LogEnv
logEnv SimpleLogPayload
payload Severity
DebugS Text
message
  where
    payload :: SimpleLogPayload
payload =
        Text -> SimpleLogPayload
moduleField Text
"Ecluse.Runtime.Queue.Sqs"
            SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"reason" (SqsDropReason -> Text
dropReasonLabel SqsDropReason
reason)
            SimpleLogPayload -> SimpleLogPayload -> SimpleLogPayload
forall a. Semigroup a => a -> a -> a
<> SimpleLogPayload
-> (Text -> SimpleLogPayload) -> Maybe Text -> SimpleLogPayload
forall b a. b -> (a -> b) -> Maybe a -> b
maybe SimpleLogPayload
forall a. Monoid a => a
mempty (Text -> Text -> SimpleLogPayload
forall a. ToJSON a => Text -> a -> SimpleLogPayload
sl Text
"messageId") Maybe Text
messageId
    message :: Text
message = Text
"dropped an unusable SQS message: " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> SqsDropReason -> Text
dropReasonLabel SqsDropReason
reason

-- The operator-facing phrase for each drop reason.
dropReasonLabel :: SqsDropReason -> Text
dropReasonLabel :: SqsDropReason -> Text
dropReasonLabel = \case
    SqsDropReason
MissingBody -> Text
"missing body"
    SqsDropReason
MissingReceipt -> Text
"missing receipt"
    SqsDropReason
UndecodableBody -> Text
"undecodable body"

{- | Read a queue payload's package identity through its ecosystem's own grammar, the gate
'Ecluse.Core.Queue.decodeJob' applies at the trust boundary. Only npm has one to read it through.
-}
mirrorJobPackage :: Ecosystem -> Maybe Text -> Text -> Either Text PackageName
mirrorJobPackage :: Ecosystem -> Maybe Text -> Text -> Either Text PackageName
mirrorJobPackage Ecosystem
eco Maybe Text
namespace Text
rawName = case Ecosystem
eco of
    Ecosystem
Npm -> (ParseError -> Text)
-> Either ParseError PackageName -> Either Text PackageName
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 ParseError -> Text
parseErrorMessage (Text -> Either ParseError PackageName
projectName Text
npmWireName)
    Ecosystem
PyPI -> PackageName -> Either Text PackageName
forall a b. b -> Either a b
Right PackageName
asGiven
    Ecosystem
RubyGems -> PackageName -> Either Text PackageName
forall a b. b -> Either a b
Right PackageName
asGiven
  where
    asGiven :: PackageName
asGiven = Ecosystem -> Maybe Scope -> Text -> PackageName
mkPackageName Ecosystem
eco (Text -> Scope
mkScope (Text -> Scope) -> Maybe Text -> Maybe Scope
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Maybe Text
namespace) Text
rawName
    npmWireName :: Text
npmWireName = Text -> (Text -> Text) -> Maybe Text -> Text
forall b a. b -> (a -> b) -> Maybe a -> b
maybe Text
rawName (\Text
ns -> Text
"@" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
ns Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
"/" Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
rawName) Maybe Text
namespace