shibuya-pgmq-adapter 0.13.0.0 → 0.14.0.0
raw patch · 6 files changed
+272/−40 lines, 6 filesdep ~shibuya-corePVP ok
version bump matches the API change (PVP)
Dependency ranges changed: shibuya-core
API changes (from Hackage documentation)
Files
- CHANGELOG.md +19/−0
- shibuya-pgmq-adapter.cabal +3/−3
- src/Shibuya/Adapter/Pgmq/Convert.hs +18/−13
- test/Shibuya/Adapter/Pgmq/ChaosSpec.hs +130/−1
- test/Shibuya/Adapter/Pgmq/ConvertSpec.hs +66/−19
- test/Shibuya/Adapter/Pgmq/PropertySpec.hs +36/−4
CHANGELOG.md view
@@ -1,5 +1,24 @@ # Changelog +## 0.14.0.0 — 2026-08-10++This release adopts the application-defined dead-letter reason API from `shibuya-core` 0.9.+The adapter's Haskell API is unchanged, while its PGMQ dead-letter payload gains structured,+machine-queryable reason fields.++### Breaking Changes++- Requires `shibuya-core ^>=0.9.0.0`, up from `^>=0.8.0.1`. The dependency release extends+ `DeadLetterReason` with `ApplicationFailure` and adds total public reason projections.++### Dead-Letter Queues++- Every new DLQ payload retains the compatibility field `dead_letter_reason` and adds+ `dead_letter_reason_code` plus an always-present `dead_letter_reason_detail`. Reasons without+ detail encode that last field as JSON `null`.+- Application-owned dead-letter codes and details are transported verbatim through the public+ Shibuya projections; the adapter no longer duplicates Shibuya's constructor renderer.+ ## 0.13.0.0 — 2026-08-09 Driven by the `pgmq-hs` 0.5 release. The adapter remains paired with
shibuya-pgmq-adapter.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.12 name: shibuya-pgmq-adapter-version: 0.13.0.0+version: 0.14.0.0 synopsis: PGMQ adapter for the Shibuya queue processing framework description: A Shibuya adapter that integrates with pgmq (PostgreSQL Message Queue)@@ -49,7 +49,7 @@ pgmq-core ^>=0.5, pgmq-effectful ^>=0.5, pgmq-hasql ^>=0.5,- shibuya-core ^>=0.8.0.1,+ shibuya-core ^>=0.9.0.0, stm ^>=2.5, streamly ^>=0.11, streamly-core ^>=0.3,@@ -119,7 +119,7 @@ pgmq-migration ^>=0.5, quickcheck-instances ^>=0.3, random,- shibuya-core ^>=0.8.0.1,+ shibuya-core ^>=0.9.0.0, shibuya-pgmq-adapter, stm, streamly ^>=0.11,
src/Shibuya/Adapter/Pgmq/Convert.hs view
@@ -25,7 +25,13 @@ import Data.Text qualified as Text import Data.Text.Encoding qualified as TE import Pgmq.Types qualified as Pgmq-import Shibuya.Core.Ack (DeadLetterReason (..))+import Shibuya.Core.Ack+ ( DeadLetterReason,+ deadLetterCodeText,+ deadLetterReasonCode,+ deadLetterReasonDetail,+ renderDeadLetterReason,+ ) import Shibuya.Core.Types (Attempt (..), Cursor (..), Envelope (..), MessageId (..), TraceHeaders) -- | Convert a pgmq MessageId to a Shibuya MessageId.@@ -142,12 +148,17 @@ -- | DLQ message body Pgmq.MessageBody mkDlqPayload msg reason includeMetadata =- Pgmq.MessageBody $- object $- [ "original_message" .= Pgmq.unMessageBody msg.body,- "dead_letter_reason" .= reasonToText reason- ]- ++ metadataFields+ let rendered = renderDeadLetterReason reason+ code = deadLetterCodeText (deadLetterReasonCode reason)+ detail = deadLetterReasonDetail reason+ in Pgmq.MessageBody $+ object $+ [ "original_message" .= Pgmq.unMessageBody msg.body,+ "dead_letter_reason" .= rendered,+ "dead_letter_reason_code" .= code,+ "dead_letter_reason_detail" .= detail+ ]+ ++ metadataFields where metadataFields | includeMetadata =@@ -158,9 +169,3 @@ "original_headers" .= msg.headers ] | otherwise = []-- reasonToText :: DeadLetterReason -> Text- reasonToText = \case- PoisonPill t -> "poison_pill: " <> t- InvalidPayload t -> "invalid_payload: " <> t- MaxRetriesExceeded -> "max_retries_exceeded"
test/Shibuya/Adapter/Pgmq/ChaosSpec.hs view
@@ -20,9 +20,13 @@ import Data.Vector qualified as Vector import Effectful (Eff, IOE, liftIO, runEff, (:>)) import Effectful.Error.Static (Error, runErrorNoCallStack)+import Hasql.Decoders qualified as D import Hasql.Pool qualified as Pool+import Hasql.Session qualified as Session+import Hasql.Statement qualified as Statement import Pgmq.Effectful (Pgmq, PgmqRuntimeError, runPgmq) import Pgmq.Effectful qualified as PgmqEff+import Pgmq.Hasql.Encoders qualified as Encoders import Pgmq.Hasql.Sessions qualified as Sessions import Pgmq.Hasql.Statements.Types (ReadMessage (..), SendMessage (..), SendMessageWithHeaders (..)) import Pgmq.Types (MessageBody (..), MessageHeaders (..), QueueName)@@ -48,7 +52,7 @@ runApp, stopAppGracefully, )-import Shibuya.Core.Ack (AckDecision (..), DeadLetterReason (..), HaltReason (..))+import Shibuya.Core.Ack (AckDecision (..), DeadLetterCode, DeadLetterReason (..), HaltReason (..), mkDeadLetterCode) import Shibuya.Core.AckHandle (AckHandle (..)) import Shibuya.Core.Ingested (Ingested (..)) import Shibuya.Handler (Handler)@@ -205,6 +209,95 @@ } Vector.length mainCount `shouldBe` 0 + it "preserves an application reason as queryable JSONB fields" $ \TestFixture {pool, queueName, dlqName} -> do+ let detail = "selected 101 recipients; configured limit is 100" :: Text.Text+ rendered = "keiro.router.selection.recipient_overflow: " <> detail++ runPgmqSession pool $ do+ _ <-+ Sessions.sendMessage $+ SendMessage+ { queueName = queueName,+ messageBody = MessageBody (String "application-failure"),+ delay = Just 0+ }+ pure ()++ processedRef <- newIORef (0 :: Int)+ let config =+ (defaultConfig queueName)+ { visibilityTimeout = 5,+ batchSize = 1,+ polling = StandardPolling {pollInterval = 0.1},+ deadLetterConfig = Just $ directDeadLetter dlqName False+ }++ runAdapterIO pool $ runTracingNoop $ do+ adapter <- requireAdapter pool config+ let handler = applicationFailureHandler processedRef+ processor = mkProcessor adapter handler++ result <- runApp defaultAppConfig [(ProcessorId "application-dlq-test", processor)]+ case result of+ Left err -> liftIO $ expectationFailure $ "Failed to start app: " <> show err+ Right appHandle -> do+ liftIO $ waitForProcessed processedRef 1 3000000+ _ <- stopAppGracefully ShutdownConfig {drainTimeout = 5} appHandle+ pure ()++ applicationInspection <- runPgmqSession pool $ inspectNextDlqPayload dlqName+ applicationSize <- case applicationInspection of+ Just (Just storedRendered, Just storedCode, Just storedDetail, storedSize) -> do+ storedRendered `shouldBe` rendered+ storedCode `shouldBe` "keiro.router.selection.recipient_overflow"+ storedDetail `shouldBe` detail+ storedSize `shouldSatisfy` (> 0)+ pure storedSize+ other -> expectationFailure ("expected queryable application DLQ fields, got " <> show other) >> pure 0++ let legacyControl =+ object+ [ "original_message" .= String "application-failure",+ "dead_letter_reason" .= rendered+ ]+ runPgmqSession pool $ do+ _ <-+ Sessions.sendMessage $+ SendMessage+ { queueName = dlqName,+ messageBody = MessageBody legacyControl,+ delay = Just 0+ }+ pure ()++ legacyInspection <- runPgmqSession pool $ inspectNextDlqPayload dlqName+ legacySize <- case legacyInspection of+ Just (Just storedRendered, Nothing, Nothing, storedSize) -> do+ storedRendered `shouldBe` rendered+ storedSize `shouldSatisfy` (> 0)+ pure storedSize+ other -> expectationFailure ("expected legacy-shaped DLQ control, got " <> show other) >> pure 0++ putStrLn $+ "DLQ-PAYLOAD-SIZE-DIAG: dualWriteJsonb="+ <> show applicationSize+ <> "B legacyJsonb="+ <> show legacySize+ <> "B delta="+ <> show (applicationSize - legacySize)+ <> "B"++ sourceMessages <-+ runPgmqSession pool $+ Sessions.readMessage $+ ReadMessage+ { queueName = queueName,+ delay = 30,+ batchSize = Just 1,+ conditional = Nothing+ }+ Vector.length sourceMessages `shouldBe` 0+ it "preserves trace headers when moving to DLQ" $ \TestFixture {pool, queueName, dlqName} -> do -- Send a message with trace headers let traceparent = "00-0af7651916cd43dd8448eb211c80319c-b7ad6b7169203331-01" :: Text.Text@@ -855,6 +948,42 @@ deadLetterHandler processedRef _ = do liftIO $ atomicModifyIORef' processedRef (\n -> (n + 1, ())) pure $ AckDeadLetter (PoisonPill "Test dead-letter")++applicationFailureHandler :: (IOE :> es) => IORef Int -> Handler es Value+applicationFailureHandler processedRef _ = do+ liftIO $ atomicModifyIORef' processedRef (\n -> (n + 1, ()))+ pure $+ AckDeadLetter $+ ApplicationFailure+ representativeDeadLetterCode+ "selected 101 recipients; configured limit is 100"++representativeDeadLetterCode :: DeadLetterCode+representativeDeadLetterCode =+ case mkDeadLetterCode "keiro.router.selection.recipient_overflow" of+ Left err -> error (Text.unpack err)+ Right code -> code++inspectNextDlqPayload :: QueueName -> Session.Session (Maybe (Maybe Text.Text, Maybe Text.Text, Maybe Text.Text, Int32))+inspectNextDlqPayload queueName = Session.statement queueName inspectDlqPayloadStatement++inspectDlqPayloadStatement :: Statement.Statement QueueName (Maybe (Maybe Text.Text, Maybe Text.Text, Maybe Text.Text, Int32))+inspectDlqPayloadStatement =+ Statement.preparable sql Encoders.queueNameEncoder decoder+ where+ sql =+ "SELECT message->>'dead_letter_reason', "+ <> "message->>'dead_letter_reason_code', "+ <> "message->>'dead_letter_reason_detail', "+ <> "pg_column_size(message) "+ <> "FROM pgmq.read($1, 30, 1, '{}'::jsonb)"+ decoder =+ D.rowMaybe $+ (,,,)+ <$> D.column (D.nullable D.text)+ <*> D.column (D.nullable D.text)+ <*> D.column (D.nullable D.text)+ <*> D.column (D.nonNullable D.int4) -- | Handler that takes a specified delay before acking. slowHandler :: (IOE :> es) => IORef Int -> Int -> Handler es Value
test/Shibuya/Adapter/Pgmq/ConvertSpec.hs view
@@ -7,7 +7,7 @@ import Data.Time (UTCTime (..), fromGregorian) import Pgmq.Types qualified as Pgmq import Shibuya.Adapter.Pgmq.Convert-import Shibuya.Core.Ack (DeadLetterReason (..))+import Shibuya.Core.Ack (DeadLetterCode, DeadLetterReason (..), mkDeadLetterCode) import Shibuya.Core.Types (Attempt (..), Cursor (..), Envelope (..), MessageId (..)) import Test.Hspec import Test.QuickCheck@@ -339,6 +339,15 @@ Object obj -> KeyMap.member "dead_letter_reason" obj `shouldBe` True _ -> expectationFailure "Expected Object" + it "always includes structured reason fields" $ do+ let Pgmq.MessageBody payload = mkDlqPayload sampleMessage MaxRetriesExceeded False+ case payload of+ Object obj -> do+ KeyMap.lookup "dead_letter_reason_code" obj+ `shouldBe` Just (String "max_retries_exceeded")+ KeyMap.lookup "dead_letter_reason_detail" obj `shouldBe` Just Null+ _ -> expectationFailure "Expected Object"+ it "does not include original_message_id" $ do let Pgmq.MessageBody payload = mkDlqPayload sampleMessage MaxRetriesExceeded False case payload of@@ -384,27 +393,65 @@ Object obj -> KeyMap.member "last_read_at" obj `shouldBe` True _ -> expectationFailure "Expected Object" - describe "reason formatting" $ do- it "formats MaxRetriesExceeded" $ do- let Pgmq.MessageBody payload = mkDlqPayload sampleMessage MaxRetriesExceeded False+ it "retains all reason fields unchanged" $ do+ let Pgmq.MessageBody payload = mkDlqPayload sampleMessage MaxRetriesExceeded True case payload of- Object obj ->+ Object obj -> do KeyMap.lookup "dead_letter_reason" obj `shouldBe` Just (String "max_retries_exceeded")+ KeyMap.lookup "dead_letter_reason_code" obj+ `shouldBe` Just (String "max_retries_exceeded")+ KeyMap.lookup "dead_letter_reason_detail" obj `shouldBe` Just Null _ -> expectationFailure "Expected Object" - it "formats PoisonPill with message" $ do- let Pgmq.MessageBody payload = mkDlqPayload sampleMessage (PoisonPill "corrupt data") False- case payload of- Object obj ->- KeyMap.lookup "dead_letter_reason" obj- `shouldBe` Just (String "poison_pill: corrupt data")- _ -> expectationFailure "Expected Object"+ describe "reason fields" $ do+ let expectedPayload rendered code detail =+ object+ [ "original_message" .= object ["data" .= ("test" :: Text.Text)],+ "dead_letter_reason" .= (rendered :: Text.Text),+ "dead_letter_reason_code" .= (code :: Text.Text),+ "dead_letter_reason_detail" .= (detail :: Maybe Text.Text)+ ]+ payloadFor reason =+ let Pgmq.MessageBody payload = mkDlqPayload sampleMessage reason False+ in payload - it "formats InvalidPayload with message" $ do- let Pgmq.MessageBody payload = mkDlqPayload sampleMessage (InvalidPayload "parse error") False- case payload of- Object obj ->- KeyMap.lookup "dead_letter_reason" obj- `shouldBe` Just (String "invalid_payload: parse error")- _ -> expectationFailure "Expected Object"+ it "writes the exact PoisonPill object" $ do+ payloadFor (PoisonPill "corrupt data")+ `shouldBe` expectedPayload "poison_pill: corrupt data" "poison_pill" (Just "corrupt data")++ it "writes the exact InvalidPayload object" $ do+ payloadFor (InvalidPayload "parse error")+ `shouldBe` expectedPayload "invalid_payload: parse error" "invalid_payload" (Just "parse error")++ it "writes the exact MaxRetriesExceeded object with null detail" $ do+ payloadFor MaxRetriesExceeded+ `shouldBe` expectedPayload "max_retries_exceeded" "max_retries_exceeded" Nothing++ it "writes the exact ApplicationFailure object" $ do+ payloadFor (ApplicationFailure representativeDeadLetterCode "selected 101 recipients; configured limit is 100")+ `shouldBe` expectedPayload+ "keiro.router.selection.recipient_overflow: selected 101 recipients; configured limit is 100"+ "keiro.router.selection.recipient_overflow"+ (Just "selected 101 recipients; configured limit is 100")++ it "preserves empty application detail as an empty string" $ do+ payloadFor (ApplicationFailure representativeDeadLetterCode "")+ `shouldBe` expectedPayload+ "keiro.router.selection.recipient_overflow: "+ "keiro.router.selection.recipient_overflow"+ (Just "")++ it "preserves Unicode, quotes, backslashes, and colons in structured detail" $ do+ let detail = "選択失敗: recipient=\"north\" \\ retry" :: Text.Text+ payloadFor (ApplicationFailure representativeDeadLetterCode detail)+ `shouldBe` expectedPayload+ ("keiro.router.selection.recipient_overflow: " <> detail)+ "keiro.router.selection.recipient_overflow"+ (Just detail)++representativeDeadLetterCode :: DeadLetterCode+representativeDeadLetterCode =+ case mkDeadLetterCode "keiro.router.selection.recipient_overflow" of+ Left err -> error (Text.unpack err)+ Right code -> code
test/Shibuya/Adapter/Pgmq/PropertySpec.hs view
@@ -1,8 +1,8 @@-{-# OPTIONS_GHC -Wno-orphans #-}+{-# OPTIONS_GHC -Wno-orphans -Werror=incomplete-patterns #-} module Shibuya.Adapter.Pgmq.PropertySpec (spec) where -import Data.Aeson (Value (..))+import Data.Aeson (Value (..), toJSON) import Data.Aeson.KeyMap qualified as KeyMap import Data.Int (Int64) import Data.Maybe (isJust)@@ -17,7 +17,15 @@ pgmqMessageIdToCursor, pgmqMessageToEnvelope, )-import Shibuya.Core.Ack (DeadLetterReason (..))+import Shibuya.Core.Ack+ ( DeadLetterCode,+ DeadLetterReason (..),+ deadLetterCodeText,+ deadLetterReasonCode,+ deadLetterReasonDetail,+ mkDeadLetterCode,+ renderDeadLetterReason,+ ) import Shibuya.Core.Types (Cursor (..), Envelope (..)) import Test.Hspec import Test.QuickCheck@@ -87,6 +95,21 @@ Object obj -> KeyMap.member "dead_letter_reason" obj === True _ -> property False + it "writes all reason fields from Shibuya's total public projections" $ property $ \(reason :: DeadLetterReason) (includeMeta :: Bool) ->+ let msg = mkTestMessage 1+ Pgmq.MessageBody payload = mkDlqPayload msg reason includeMeta+ in case payload of+ Object obj ->+ conjoin+ [ KeyMap.lookup "dead_letter_reason" obj+ === Just (String (renderDeadLetterReason reason)),+ KeyMap.lookup "dead_letter_reason_code" obj+ === Just (String (deadLetterCodeText (deadLetterReasonCode reason))),+ KeyMap.lookup "dead_letter_reason_detail" obj+ === Just (toJSON (deadLetterReasonDetail reason))+ ]+ _ -> property False+ it "metadata keys present iff includeMeta is True" $ property $ \(includeMeta :: Bool) -> let msg = mkTestMessage 1 Pgmq.MessageBody payload = mkDlqPayload msg MaxRetriesExceeded includeMeta@@ -116,7 +139,8 @@ oneof [ pure MaxRetriesExceeded, PoisonPill <$> arbitraryText,- InvalidPayload <$> arbitraryText+ InvalidPayload <$> arbitraryText,+ ApplicationFailure propertyDeadLetterCode <$> arbitraryText ] where arbitraryText :: Gen Text@@ -125,3 +149,11 @@ shrink MaxRetriesExceeded = [] shrink (PoisonPill t) = MaxRetriesExceeded : [PoisonPill (Text.pack t') | t' <- shrink (Text.unpack t)] shrink (InvalidPayload t) = MaxRetriesExceeded : [InvalidPayload (Text.pack t') | t' <- shrink (Text.unpack t)]+ shrink (ApplicationFailure _ t) =+ MaxRetriesExceeded : [ApplicationFailure propertyDeadLetterCode (Text.pack t') | t' <- shrink (Text.unpack t)]++propertyDeadLetterCode :: DeadLetterCode+propertyDeadLetterCode =+ case mkDeadLetterCode "test.property.application_failure" of+ Left err -> error (Text.unpack err)+ Right code -> code