module Ecluse.Runtime.Queue.Sqs.Internal (
SqsConfig (..),
defaultSqsConfig,
newSqsQueue,
ReceivedMessage (..),
liftReceivedMessages,
deadLetterTerminusOf,
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)
data SqsConfig = SqsConfig
{ SqsConfig -> Text
sqsQueueUrl :: Text
, SqsConfig -> Text
sqsRegion :: Text
, SqsConfig -> Maybe AwsEndpoint
sqsEndpoint :: Maybe AwsEndpoint
, SqsConfig -> Int
sqsBatchSize :: Int
, SqsConfig -> Int
sqsWaitSeconds :: Int
, SqsConfig -> Seconds
sqsVisibilityTimeout :: Seconds
, SqsConfig -> Seconds
sqsTerminalBackoff :: Seconds
, SqsConfig -> DeliveryBudget
sqsMaxReceiveCount :: DeliveryBudget
}
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)
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
}
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
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
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
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
,
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
}
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
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
((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
sqsInFlightMaximum :: Seconds
sqsInFlightMaximum :: Seconds
sqsInFlightMaximum = Int -> Seconds
Seconds Int
43_200
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))
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
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)
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
data ReceivedMessage = ReceivedMessage
{ ReceivedMessage -> Maybe Text
rmBody :: Maybe Text
, ReceivedMessage -> Maybe Text
rmReceipt :: Maybe Text
, ReceivedMessage -> Maybe Text
rmMessageId :: Maybe Text
, ReceivedMessage -> Maybe Text
rmReceiveCount :: Maybe Text
}
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)
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)
}
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)
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)
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
}
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))
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)
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)
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
dropReasonLabel :: SqsDropReason -> Text
dropReasonLabel :: SqsDropReason -> Text
dropReasonLabel = \case
SqsDropReason
MissingBody -> Text
"missing body"
SqsDropReason
MissingReceipt -> Text
"missing receipt"
SqsDropReason
UndecodableBody -> Text
"undecodable body"
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