antiope-sqs 7.0.2 → 7.0.3
raw patch · 2 files changed
+18/−12 lines, 2 filesdep +splitPVP ok
version bump matches the API change (PVP)
Dependencies added: split
API changes (from Hackage documentation)
Files
- antiope-sqs.cabal +2/−1
- src/Antiope/SQS.hs +16/−11
antiope-sqs.cabal view
@@ -1,7 +1,7 @@ cabal-version: 2.2 name: antiope-sqs-version: 7.0.2+version: 7.0.3 synopsis: Please see the README on Github at <https://github.com/arbor/antiope#readme> description: Please see the README on Github at <https://github.com/arbor/antiope#readme>. category: Services@@ -39,6 +39,7 @@ , monad-loops , mtl , network-uri+ , split , text , unliftio-core , unordered-containers
src/Antiope/SQS.hs view
@@ -17,13 +17,14 @@ ) where import Control.Lens-import Control.Monad (forM_, join, void)+import Control.Monad (forM, forM_, join, void) import Control.Monad.IO.Unlift (MonadUnliftIO) import Control.Monad.Loops (unfoldWhileM) import Control.Monad.Trans (lift) import Data.Coerce (coerce) import Data.Conduit import Data.Conduit.Combinators (yieldMany)+import Data.List.Split (chunksOf) import Data.Maybe (catMaybes) import Data.Text (pack) import Network.AWS (HasEnv, MonadAWS, runAWS, runResourceT)@@ -61,16 +62,20 @@ -> [msg] -> m (Either SQSError ()) ackMessages (QueueUrl queueUrl) msgs = do- let receipts = msgs ^.. each . to getReceiptHandle & catMaybes- -- each dmbr needs an ID. just use the list index.- let dmbres = (\(r, i) -> deleteMessageBatchRequestEntry (pack (show i)) r) <$> zip (coerce receipts) ([0..] :: [Int])- resp <- AWS.send $ deleteMessageBatch queueUrl & dmbEntries .~ dmbres- -- only acceptable if no errors.- if resp ^. dmbrsResponseStatus == 200- then case resp ^. dmbrsFailed of- [] -> return $ Right ()- _ -> return $ Left DeleteMessageBatchError- else return $ Left DeleteMessageBatchError+ let receipts' = msgs ^.. each . to getReceiptHandle & catMaybes+ let maxBatchSize = 10 -- Amazon enforces this+ results <- forM (chunksOf maxBatchSize receipts') $ \receipts -> do+ -- each dmbr needs an ID. just use the list index.+ let dmbres = (\(r, i) -> deleteMessageBatchRequestEntry (pack (show i)) r) <$> zip (coerce receipts) ([0..] :: [Int])+ resp <- AWS.send $ deleteMessageBatch queueUrl & dmbEntries .~ dmbres+ -- only acceptable if no errors.+ if resp ^. dmbrsResponseStatus == 200+ then case resp ^. dmbrsFailed of+ [] -> return $ Right ()+ _ -> return $ Left DeleteMessageBatchError+ else return $ Left DeleteMessageBatchError+ pure $ sequence_ results+ -- | Reads from an SQS indefinitely, producing messages into a conduit queueSource :: MonadAWS m => QueueUrl -> ConduitT () Message m ()