packages feed

shibuya-pgmq-adapter 0.8.0.0 → 0.11.0.0

raw patch · 9 files changed

+1128/−343 lines, 9 filesdep +hasql-transactiondep ~shibuya-corePVP ok

version bump matches the API change (PVP)

Dependencies added: hasql-transaction

Dependency ranges changed: shibuya-core

API changes (from Hackage documentation)

+ Shibuya.Adapter.Pgmq: InvalidAckRetryBackoff :: !NominalDiffTime -> !NominalDiffTime -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidAckRetryMaxAttempts :: !Int -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidBatchSize :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidHaltVisibilityTimeout :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidLongPollIntervalMs :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidLongPollSeconds :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidMaxRetries :: !Int64 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidPollRetryBackoff :: !NominalDiffTime -> !NominalDiffTime -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidPollRetryMaxAttempts :: !Int -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidPrefetchBufferSize :: !Natural -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidStandardPollInterval :: !NominalDiffTime -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: InvalidVisibilityTimeout :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq: PgmqAdapterEnv :: !Pool -> (Message -> IO ()) -> (Message -> PgmqRuntimeError -> IO ()) -> PgmqAdapterEnv
+ Shibuya.Adapter.Pgmq: [ackRetry] :: PgmqAdapterConfig -> !PollRetryConfig
+ Shibuya.Adapter.Pgmq: [haltVisibilityTimeout] :: PgmqAdapterConfig -> !Maybe Int32
+ Shibuya.Adapter.Pgmq: [onAckFailure] :: PgmqAdapterEnv -> Message -> PgmqRuntimeError -> IO ()
+ Shibuya.Adapter.Pgmq: [onAutoDeadLetter] :: PgmqAdapterEnv -> Message -> IO ()
+ Shibuya.Adapter.Pgmq: [pool] :: PgmqAdapterEnv -> !Pool
+ Shibuya.Adapter.Pgmq: data PgmqAdapterEnv
+ Shibuya.Adapter.Pgmq: data PgmqConfigError
+ Shibuya.Adapter.Pgmq: mkPgmqAdapterEnv :: Pool -> PgmqAdapterEnv
+ Shibuya.Adapter.Pgmq: validateConfig :: PgmqAdapterConfig -> Either PgmqConfigError PgmqAdapterConfig
+ Shibuya.Adapter.Pgmq.Config: InvalidAckRetryBackoff :: !NominalDiffTime -> !NominalDiffTime -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidAckRetryMaxAttempts :: !Int -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidBatchSize :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidHaltVisibilityTimeout :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidLongPollIntervalMs :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidLongPollSeconds :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidMaxRetries :: !Int64 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidPollRetryBackoff :: !NominalDiffTime -> !NominalDiffTime -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidPollRetryMaxAttempts :: !Int -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidPrefetchBufferSize :: !Natural -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidStandardPollInterval :: !NominalDiffTime -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: InvalidVisibilityTimeout :: !Int32 -> PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: PgmqAdapterEnv :: !Pool -> (Message -> IO ()) -> (Message -> PgmqRuntimeError -> IO ()) -> PgmqAdapterEnv
+ Shibuya.Adapter.Pgmq.Config: [ackRetry] :: PgmqAdapterConfig -> !PollRetryConfig
+ Shibuya.Adapter.Pgmq.Config: [haltVisibilityTimeout] :: PgmqAdapterConfig -> !Maybe Int32
+ Shibuya.Adapter.Pgmq.Config: [onAckFailure] :: PgmqAdapterEnv -> Message -> PgmqRuntimeError -> IO ()
+ Shibuya.Adapter.Pgmq.Config: [onAutoDeadLetter] :: PgmqAdapterEnv -> Message -> IO ()
+ Shibuya.Adapter.Pgmq.Config: [pool] :: PgmqAdapterEnv -> !Pool
+ Shibuya.Adapter.Pgmq.Config: data PgmqAdapterEnv
+ Shibuya.Adapter.Pgmq.Config: data PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: instance GHC.Classes.Eq Shibuya.Adapter.Pgmq.Config.PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: instance GHC.Internal.Generics.Generic Shibuya.Adapter.Pgmq.Config.PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: instance GHC.Internal.Show.Show Shibuya.Adapter.Pgmq.Config.PgmqConfigError
+ Shibuya.Adapter.Pgmq.Config: mkPgmqAdapterEnv :: Pool -> PgmqAdapterEnv
+ Shibuya.Adapter.Pgmq.Config: validateConfig :: PgmqAdapterConfig -> Either PgmqConfigError PgmqAdapterConfig
- Shibuya.Adapter.Pgmq: PgmqAdapterConfig :: !QueueName -> !Int32 -> !Int32 -> !PollingConfig -> !PollRetryConfig -> !Maybe DeadLetterConfig -> !Int64 -> !Maybe FifoConfig -> !Maybe PrefetchConfig -> PgmqAdapterConfig
+ Shibuya.Adapter.Pgmq: PgmqAdapterConfig :: !QueueName -> !Int32 -> !Int32 -> !PollingConfig -> !PollRetryConfig -> !PollRetryConfig -> !Maybe DeadLetterConfig -> !Maybe Int32 -> !Int64 -> !Maybe FifoConfig -> !Maybe PrefetchConfig -> PgmqAdapterConfig
- Shibuya.Adapter.Pgmq: pgmqAdapter :: forall (es :: [Effect]). (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es, Tracing :> es) => PgmqAdapterConfig -> Eff es (Adapter es Value)
+ Shibuya.Adapter.Pgmq: pgmqAdapter :: forall (es :: [Effect]). (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es, Tracing :> es) => PgmqAdapterEnv -> PgmqAdapterConfig -> Eff es (Either PgmqConfigError (Adapter es Value))
- Shibuya.Adapter.Pgmq.Config: PgmqAdapterConfig :: !QueueName -> !Int32 -> !Int32 -> !PollingConfig -> !PollRetryConfig -> !Maybe DeadLetterConfig -> !Int64 -> !Maybe FifoConfig -> !Maybe PrefetchConfig -> PgmqAdapterConfig
+ Shibuya.Adapter.Pgmq.Config: PgmqAdapterConfig :: !QueueName -> !Int32 -> !Int32 -> !PollingConfig -> !PollRetryConfig -> !PollRetryConfig -> !Maybe DeadLetterConfig -> !Maybe Int32 -> !Int64 -> !Maybe FifoConfig -> !Maybe PrefetchConfig -> PgmqAdapterConfig

Files

CHANGELOG.md view
@@ -1,5 +1,86 @@ # Changelog +## 0.11.0.0 — 2026-07-04++Paired with `shibuya-core 0.8.0.1`.++### Breaking Changes++- Requires `shibuya-core ^>=0.8.0.1` (up from `^>=0.7.0.0`) in the library and+  test stanzas. `shibuya-core 0.8.0.0` is a breaking release, so adapter+  consumers must migrate along with it:+  - Handlers now receive `Message es msg` (envelope + optional lease, no ack+    finalizer) instead of `Ingested`. Handlers written against the `Handler`+    type alias that read `msg.envelope` / `msg.lease` compile unchanged;+    handlers with an explicit `Ingested es msg -> …` signature must switch to+    `Message es msg -> …`.+  - `runApp` now takes a validated `AppConfig` record instead of positional+    supervision-strategy and inbox-size arguments. `defaultAppConfig`+    (`AppConfig { strategy = IgnoreFailures, inboxSize = 100 }`) is the drop-in+    replacement for the old `runApp IgnoreFailures 100 …`.+  - The runner internals moved under `Shibuya.Internal.*`; metrics types such as+    `ProcessorId` are public via `Shibuya.Core.Metrics` (and re-exported from+    `Shibuya.App`).++  See the `shibuya-core 0.8.0.0` migration guide for the full list.++### Notes++- The `^>=0.8.0.1` lower bound (rather than `0.8.0.0`) pulls in the+  `shibuya-core 0.8.0.1` patch, which cuts per-message allocation on the `Async`+  and `Ahead` concurrency dispatch paths. No API or behavior change; it benefits+  the adapter's throughput on those paths for free.+- The adapter's own public API is unchanged. `pgmqSource` still yields+  `Ingested es Value`; the framework projects each `Ingested` to the+  handler-facing `Message` itself.+- The bundled example (`shibuya-pgmq-example`) and benchmark+  (`shibuya-pgmq-adapter-bench`), plus the README and getting-started guide,+  were updated to the `shibuya-core 0.8.0.0` API (`runApp defaultAppConfig`, the+  `Message` handler pattern, `ProcessorId` from `Shibuya.App`). Neither the+  example nor the benchmark is published to Hackage.++## 0.10.0.0 — 2026-07-04++### Features++- Reintroduced opt-in concurrent prefetch via `prefetchConfig :: Maybe PrefetchConfig`+  (default `Nothing`). When enabled, the polling stage reads the next batches on a+  background worker, overlapping database latency with handler work. The historical+  `parBuffered` deadlock (`thread blocked indefinitely in an STM transaction`) is fixed+  by running only the prefetch stage under effectful's `ConcUnlift` strategy (scoped via+  `morphInner`), so the non-prefetch path is unchanged (still `SeqUnlift`, no overhead).++### Breaking Changes++- `PgmqAdapterConfig` gained a `prefetchConfig :: Maybe PrefetchConfig` field. Callers that+  construct the config by full record literal must add it; `defaultConfig` sets it to `Nothing`.+- `PgmqConfigError` gained an `InvalidPrefetchBufferSize` constructor; `validateConfig` now+  rejects a prefetch `bufferSize` of `0`.++### Notes++- Shutdown with prefetch enabled can leave up to `bufferSize * batchSize` already-read+  messages invisible until their visibility timeout expires. No messages are lost — they are+  redelivered after the visibility timeout; only redelivery is delayed. This bounded,+  at-least-once-safe behaviour is documented on `PrefetchConfig` and in the adapter+  architecture docs.++## 0.9.0.0 — 2026-07-02++### Breaking Changes++- `pgmqAdapter` now takes `PgmqAdapterEnv` and returns `Either PgmqConfigError (Adapter es Value)`.+- `PgmqAdapterConfig` gained `ackRetry` and `haltVisibilityTimeout`, and removed the known-deadlocking concurrent lookahead configuration.++### Reliability++- DLQ send and source delete now run in one PostgreSQL transaction.+- Message finalizers are idempotent after success.+- Ack operations and lease extension use bounded transient retry.+- `AckHalt` uses a configured visibility timeout instead of a hardcoded hour.+- Lease extension uses absolute visibility deadlines so later extension calls do not shorten the tracked lease.+- Trace header merging on the DLQ path tolerates non-UTF8 bytes.+ ## 0.8.0.0 — 2026-06-15  ### Breaking Changes
shibuya-pgmq-adapter.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.12 name: shibuya-pgmq-adapter-version: 0.8.0.0+version: 0.11.0.0 synopsis: PGMQ adapter for the Shibuya queue processing framework description:   A Shibuya adapter that integrates with pgmq (PostgreSQL Message Queue)@@ -43,10 +43,13 @@     base ^>=4.21.0.0,     bytestring ^>=0.12,     effectful-core ^>=2.6.1.0,+    hasql ^>=1.10,+    hasql-pool ^>=1.4,+    hasql-transaction ^>=1.2,     pgmq-core ^>=0.3,     pgmq-effectful ^>=0.3,     pgmq-hasql ^>=0.3,-    shibuya-core ^>=0.7.0.0,+    shibuya-core ^>=0.8.0.1,     stm ^>=2.5,     streamly ^>=0.11,     streamly-core ^>=0.3,@@ -107,6 +110,7 @@     ephemeral-pg,     hasql ^>=1.10,     hasql-pool ^>=1.4,+    hasql-transaction ^>=1.2,     hspec ^>=2.11,     pgmq-core ^>=0.3,     pgmq-effectful ^>=0.3,@@ -114,7 +118,7 @@     pgmq-migration ^>=0.3,     quickcheck-instances ^>=0.3,     random,-    shibuya-core ^>=0.7.0.0,+    shibuya-core ^>=0.8.0.1,     shibuya-pgmq-adapter,     stm,     streamly ^>=0.11,
src/Shibuya/Adapter/Pgmq.hs view
@@ -6,7 +6,7 @@ -- == Example Usage -- -- @--- import Shibuya.App (runApp, QueueProcessor (..))+-- import Shibuya.App (runApp, defaultAppConfig, QueueProcessor (..)) -- import Shibuya.Adapter.Pgmq -- import Pgmq.Effectful (runPgmq) -- import Hasql.Pool qualified as Pool@@ -21,8 +21,8 @@ --       runEff --         . runPgmq pool --         $ do---             adapter <- pgmqAdapter config---             result <- runApp IgnoreFailures 100+--             Right adapter <- pgmqAdapter (mkPgmqAdapterEnv pool) config+--             result <- runApp defaultAppConfig --               [ (ProcessorId "orders", QueueProcessor adapter handleOrder) --               ] --             -- ...@@ -74,13 +74,17 @@      -- * Configuration     PgmqAdapterConfig (..),+    PgmqAdapterEnv (..),+    mkPgmqAdapterEnv,+    PgmqConfigError (..),     PollingConfig (..),     PollRetryConfig (..),+    PrefetchConfig (..),     DeadLetterConfig (..),     DeadLetterTarget (..),     FifoConfig (..),     FifoReadStrategy (..),-    PrefetchConfig (..),+    validateConfig,      -- * Smart Constructors     directDeadLetter,@@ -119,6 +123,8 @@ import Control.Monad (forM_) import Control.Monad.IO.Class (liftIO) import Data.Aeson (Value)+import Data.Function ((&))+import Data.Vector qualified as Vector import Effectful (Eff, IOE, (:>)) import Effectful.Error.Static (Error) import Pgmq.Effectful (PgmqRuntimeError)@@ -145,6 +151,8 @@     FifoConfig (..),     FifoReadStrategy (..),     PgmqAdapterConfig (..),+    PgmqAdapterEnv (..),+    PgmqConfigError (..),     PollRetryConfig (..),     PollingConfig (..),     PrefetchConfig (..),@@ -153,13 +161,17 @@     defaultPollingConfig,     defaultPrefetchConfig,     directDeadLetter,+    mkPgmqAdapterEnv,     topicDeadLetter,+    validateConfig,   )-import Shibuya.Adapter.Pgmq.Internal (pgmqSource, pgmqSourceWithPrefetch)+import Shibuya.Adapter.Pgmq.Internal (mkIngested, pgmqChunks, pgmqChunksPrefetch, releaseMessages)+import Shibuya.Core.Ingested (Ingested) import Shibuya.Telemetry.Effect (Tracing) import Streamly.Data.Stream (Stream) import Streamly.Data.Stream qualified as Stream import Streamly.Data.Stream.Prelude qualified as StreamP+import Streamly.Data.Unfold qualified as Unfold  -- | Create a PGMQ adapter with the given configuration. --@@ -170,7 +182,7 @@ -- * Lease extension capability for long-running handlers -- * Dead-letter queue support (optional) -- * FIFO ordering support (optional)--- * Concurrent prefetching (optional, via 'prefetchConfig')+-- * Typed configuration validation -- -- == Effect Requirements --@@ -181,55 +193,61 @@ -- == Example -- -- @--- adapter <- pgmqAdapter config--- runApp IgnoreFailures 100+-- Right adapter <- pgmqAdapter env config+-- runApp defaultAppConfig --   [ (ProcessorId "my-processor", QueueProcessor adapter myHandler) --   ] -- @------ == Prefetching------ To enable concurrent prefetching (polls next batch while processing current):------ @--- let config = (defaultConfig queueName) { prefetchConfig = Just defaultPrefetchConfig }--- adapter <- pgmqAdapter config--- @ pgmqAdapter ::   (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es, Tracing :> es) =>+  PgmqAdapterEnv ->   PgmqAdapterConfig ->-  Eff es (Adapter es Value)-pgmqAdapter config = do-  -- Create shutdown signal-  shutdownVar <- liftIO $ newTVarIO False--  -- Select source based on prefetch configuration-  let messageSource = case config.prefetchConfig of-        Nothing ->-          -- No prefetching - simple sequential polling-          pgmqSource config-        Just prefetch ->-          -- Concurrent prefetching enabled-          let prefetchSettings = StreamP.maxBuffer (fromIntegral prefetch.bufferSize)-           in pgmqSourceWithPrefetch prefetchSettings config--  pure-    Adapter-      { adapterName = "pgmq:" <> queueNameToText config.queueName,-        source = takeUntilShutdown shutdownVar messageSource,-        shutdown = liftIO $ atomically $ writeTVar shutdownVar True-      }+  Eff es (Either PgmqConfigError (Adapter es Value))+pgmqAdapter env config =+  case validateConfig config of+    Left err -> pure (Left err)+    Right validConfig -> do+      shutdownVar <- liftIO $ newTVarIO False+      let messageSource = pgmqSourceWithShutdown env validConfig shutdownVar+      pure $+        Right+          Adapter+            { adapterName = "pgmq:" <> queueNameToText validConfig.queueName,+              source = messageSource,+              shutdown = liftIO $ atomically $ writeTVar shutdownVar True+            } --- | Take from stream until shutdown signal is set.-takeUntilShutdown ::-  (IOE :> es) =>+-- | Poll in chunks so a shutdown can release messages that were read from pgmq+-- but not yet handed to Shibuya's bounded inbox.+pgmqSourceWithShutdown ::+  (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es, Tracing :> es) =>+  PgmqAdapterEnv ->+  PgmqAdapterConfig ->   TVar Bool ->-  Stream (Eff es) a ->-  Stream (Eff es) a-takeUntilShutdown shutdownVar =-  Stream.takeWhileM $ \_ -> do-    isShutdown <- liftIO $ readTVarIO shutdownVar-    pure (not isShutdown)+  Stream (Eff es) (Ingested es Value)+pgmqSourceWithShutdown env config shutdownVar =+  chunkStream+    & Stream.filter (not . Vector.null)+    & Stream.takeWhileM keepChunk+    & Stream.unfoldEach (Unfold.unfoldr Vector.uncons)+    & Stream.mapMaybeM (mkIngested env config)+  where+    -- Only the polling stage is prefetched (and only under a locally-scoped+    -- ConcUnlift, see 'pgmqChunksPrefetch'). The shutdown gate, flatten, and+    -- 'mkIngested'/finalization stages stay on the consumer thread. When+    -- prefetch is disabled the stream is exactly 'pgmqChunks config', running+    -- under the default SeqUnlift strategy with no added overhead.+    chunkStream = case config.prefetchConfig of+      Nothing -> pgmqChunks config+      Just prefetch ->+        pgmqChunksPrefetch (StreamP.maxBuffer (fromIntegral prefetch.bufferSize)) config+    keepChunk chunk = do+      isShutdown <- liftIO $ readTVarIO shutdownVar+      if isShutdown+        then do+          releaseMessages env config chunk+          pure False+        else pure True  -------------------------------------------------------------------------------- -- Topic Management (pgmq 1.11.0+)
src/Shibuya/Adapter/Pgmq/Config.hs view
@@ -2,11 +2,19 @@ module Shibuya.Adapter.Pgmq.Config   ( -- * Main Configuration     PgmqAdapterConfig (..),+    PgmqAdapterEnv (..),+    mkPgmqAdapterEnv,+    PgmqConfigError (..),+    validateConfig,      -- * Polling Configuration     PollingConfig (..),     PollRetryConfig (..), +    -- * Prefetch Configuration+    PrefetchConfig (..),+    defaultPrefetchConfig,+     -- * Dead-Letter Queue Configuration     DeadLetterConfig (..),     DeadLetterTarget (..),@@ -19,94 +27,202 @@     FifoConfig (..),     FifoReadStrategy (..), -    -- * Prefetch Configuration-    PrefetchConfig (..),-     -- * Defaults     defaultConfig,     defaultPollingConfig,     defaultPollRetryConfig,-    defaultPrefetchConfig,   ) where  import Data.Int (Int32, Int64) import Data.Time (NominalDiffTime) import GHC.Generics (Generic)+import Hasql.Pool qualified as Pool import Numeric.Natural (Natural)+import Pgmq.Effectful (PgmqRuntimeError) import Pgmq.Types (QueueName, RoutingKey)+import Pgmq.Types qualified as Pgmq +-- | Runtime resources and callbacks used by the adapter.+data PgmqAdapterEnv = PgmqAdapterEnv+  { -- | Connection pool used for operations that need one transaction.+    pool :: !Pool.Pool,+    -- | Called when the source stream dead-letters an over-retried message.+    onAutoDeadLetter :: Pgmq.Message -> IO (),+    -- | Called when an ack path fails after retry or fails permanently.+    onAckFailure :: Pgmq.Message -> PgmqRuntimeError -> IO ()+  }++-- | Build an adapter environment with no-op callbacks.+mkPgmqAdapterEnv :: Pool.Pool -> PgmqAdapterEnv+mkPgmqAdapterEnv pool =+  PgmqAdapterEnv+    { pool = pool,+      onAutoDeadLetter = const (pure ()),+      onAckFailure = \_ _ -> pure ()+    }+ -- | Configuration for the PGMQ adapter. data PgmqAdapterConfig = PgmqAdapterConfig-  { -- | Name of the queue to consume from+  { -- | Name of the queue to consume from.     queueName :: !QueueName,-    -- | Visibility timeout in seconds (default: 30)-    -- Messages become invisible for this duration after being read+    -- | Visibility timeout in seconds. Messages become invisible for this duration after being read.     visibilityTimeout :: !Int32,-    -- | Maximum number of messages to read per poll (default: 1)+    -- | Maximum number of messages to read per poll.     batchSize :: !Int32,-    -- | Polling configuration+    -- | Polling configuration.     polling :: !PollingConfig,-    -- | Retry policy for transient errors during queue polling+    -- | Retry policy for transient errors during queue polling.     pollRetry :: !PollRetryConfig,-    -- | Optional dead-letter queue configuration+    -- | Retry policy for transient errors during acknowledgement operations.+    ackRetry :: !PollRetryConfig,+    -- | Optional dead-letter queue configuration.     deadLetterConfig :: !(Maybe DeadLetterConfig),-    -- | Maximum retries before dead-lettering (default: 3)-    -- Based on pgmq's readCount field+    -- | Optional visibility timeout used for 'AckHalt'. Falls back to 'visibilityTimeout'.+    haltVisibilityTimeout :: !(Maybe Int32),+    -- | Maximum deliveries before dead-lettering.+    --+    -- This is based on pgmq's readCount field, which counts deliveries, not+    -- handler failures. A value of 0 is valid and auto-dead-letters every+    -- message before processing.     maxRetries :: !Int64,-    -- | Optional FIFO mode configuration+    -- | Optional FIFO mode configuration.     fifoConfig :: !(Maybe FifoConfig),-    -- | Optional concurrent prefetch configuration-    -- When enabled, polls ahead while processing current messages+    -- | Optional concurrent prefetch configuration.+    --+    -- When 'Just', the polling stage runs on a background worker under+    -- effectful's 'ConcUnlift' strategy so streamly's @parBuffered@ may unlift+    -- 'Eff' off-thread (see "Shibuya.Adapter.Pgmq.Internal".@pgmqChunksPrefetch@).+    -- When 'Nothing' (the default), the source runs on the default 'SeqUnlift'+    -- strategy with no added overhead.+    --+    -- See 'PrefetchConfig' for the visibility-timeout trade-off and the bounded,+    -- loss-free shutdown behaviour (buffered messages are redelivered after+    -- their visibility timeout rather than released immediately).     prefetchConfig :: !(Maybe PrefetchConfig)   }   deriving stock (Show, Eq, Generic) --- | Retry policy for transient database errors during queue polling.+-- | Typed configuration validation errors.+data PgmqConfigError+  = InvalidBatchSize !Int32+  | InvalidVisibilityTimeout !Int32+  | InvalidStandardPollInterval !NominalDiffTime+  | InvalidLongPollSeconds !Int32+  | InvalidLongPollIntervalMs !Int32+  | InvalidMaxRetries !Int64+  | InvalidHaltVisibilityTimeout !Int32+  | InvalidPollRetryMaxAttempts !Int+  | InvalidPollRetryBackoff !NominalDiffTime !NominalDiffTime+  | InvalidAckRetryMaxAttempts !Int+  | InvalidAckRetryBackoff !NominalDiffTime !NominalDiffTime+  | InvalidPrefetchBufferSize !Natural+  deriving stock (Show, Eq, Generic)++-- | Validate adapter configuration before starting a source stream.+validateConfig :: PgmqAdapterConfig -> Either PgmqConfigError PgmqAdapterConfig+validateConfig config+  | config.batchSize < 1 = Left (InvalidBatchSize config.batchSize)+  | config.visibilityTimeout < 1 = Left (InvalidVisibilityTimeout config.visibilityTimeout)+  | Just halt <- config.haltVisibilityTimeout, halt < 1 = Left (InvalidHaltVisibilityTimeout halt)+  | config.maxRetries < 0 = Left (InvalidMaxRetries config.maxRetries)+  | Just prefetch <- config.prefetchConfig,+    prefetch.bufferSize == 0 =+      Left (InvalidPrefetchBufferSize prefetch.bufferSize)+  | otherwise = do+      validatePolling config.polling+      validateRetry InvalidPollRetryMaxAttempts InvalidPollRetryBackoff config.pollRetry+      validateRetry InvalidAckRetryMaxAttempts InvalidAckRetryBackoff config.ackRetry+      pure config+  where+    validatePolling = \case+      StandardPolling interval+        | interval <= 0 -> Left (InvalidStandardPollInterval interval)+        | otherwise -> Right ()+      LongPolling maxPollSeconds pollIntervalMs+        | maxPollSeconds < 1 -> Left (InvalidLongPollSeconds maxPollSeconds)+        | pollIntervalMs < 1 -> Left (InvalidLongPollIntervalMs pollIntervalMs)+        | otherwise -> Right ()++    validateRetry maxAttemptsErr backoffErr retry+      | retry.maxAttempts < 1 = Left (maxAttemptsErr retry.maxAttempts)+      | retry.initialBackoff < 0 || retry.maxBackoff < 0 =+          Left (backoffErr retry.initialBackoff retry.maxBackoff)+      | otherwise = Right ()++-- | Retry policy for transient database errors. data PollRetryConfig = PollRetryConfig-  { -- | Total attempts per poll, including the first attempt+  { -- | Total attempts, including the first attempt.     maxAttempts :: !Int,-    -- | Delay before the first retry+    -- | Delay before the first retry.     initialBackoff :: !NominalDiffTime,-    -- | Maximum delay between retry attempts+    -- | Maximum delay between retry attempts.     maxBackoff :: !NominalDiffTime   }   deriving stock (Show, Eq, Generic) +-- | Configuration for concurrent prefetch.+--+-- When enabled, the polling stage buffers the next batches on a background+-- worker while the current messages are being processed, overlapping database+-- latency with handler work.+--+-- Trade-off: prefetched messages have their visibility timeout ticking while+-- buffered, so ensure+-- @bufferSize * batchSize * avgProcessingTime < visibilityTimeout@.+--+-- === Shutdown behaviour (no data loss)+--+-- Because prefetch reads message batches /ahead/ of processing, a shutdown can+-- leave up to @bufferSize * batchSize@ already-read messages sitting in the+-- prefetch buffer without being handed to a handler. Unlike the non-prefetch+-- path — which releases just-read, undispatched messages immediately on+-- shutdown — these buffered messages are not released promptly.+--+-- __No messages are lost.__ A pgmq read only sets a message's visibility+-- timeout; it never deletes the message. Any message that was read but not+-- acknowledged stays in the queue and pgmq redelivers it once its visibility+-- timeout expires. The only effect is that, after a shutdown, redelivery of+-- those buffered messages is /delayed/ by up to 'visibilityTimeout' seconds.+-- This is a bounded, at-least-once-safe edge case (delayed processing, not data+-- loss), inherent to reading ahead; it is verified by the adapter's test suite.+data PrefetchConfig = PrefetchConfig+  { -- | Number of batches to buffer ahead of consumption (default: 4).+    -- Higher values reduce latency but increase visibility-timeout pressure.+    -- Must be greater than zero (rejected by 'validateConfig' otherwise).+    bufferSize :: !Natural+  }+  deriving stock (Show, Eq, Generic)+ -- | Polling strategy for reading messages. data PollingConfig-  = -- | Standard polling with sleep between reads when queue is empty+  = -- | Standard polling with sleep between reads when queue is empty.     StandardPolling-      { -- | Interval between polls when no messages are available+      { -- | Interval between polls when no messages are available.         pollInterval :: !NominalDiffTime       }-  | -- | Long polling - blocks in database until messages available-    -- More efficient when queue is often empty+  | -- | Long polling blocks in PostgreSQL until messages are available or the wait expires.     LongPolling-      { -- | Maximum seconds to wait for messages (e.g., 10)+      { -- | Maximum seconds to wait for messages.         maxPollSeconds :: !Int32,-        -- | Interval between database checks in milliseconds (e.g., 100)+        -- | Interval between database checks in milliseconds.         pollIntervalMs :: !Int32       }   deriving stock (Show, Eq, Generic)  -- | Target for dead-lettered messages. data DeadLetterTarget-  = -- | Send directly to a specific queue+  = -- | Send directly to a specific queue.     DirectQueue !QueueName-  | -- | Route via topic pattern matching (pgmq 1.11.0+).-    -- Messages are sent using @pgmq.send_topic@ with the given routing key,-    -- allowing fan-out to multiple DLQ consumers based on their topic bindings.+  | -- | Route via topic pattern matching.     TopicRoute !RoutingKey   deriving stock (Show, Eq, Generic)  -- | Configuration for dead-letter queue handling.--- When a message exceeds maxRetries or receives AckDeadLetter,--- it will be sent to the configured target. data DeadLetterConfig = DeadLetterConfig-  { -- | Where to send dead-lettered messages+  { -- | Where to send dead-lettered messages.     dlqTarget :: !DeadLetterTarget,-    -- | Whether to include original message metadata in DLQ message+    -- | Whether to include original message metadata in DLQ message.     includeMetadata :: !Bool   }   deriving stock (Show, Eq, Generic)@@ -119,9 +235,7 @@       includeMetadata = metadata     } --- | Create a dead-letter config using topic-based routing (pgmq 1.11.0+).--- Messages are sent via @pgmq.send_topic@ and delivered to all queues--- whose topic bindings match the routing key.+-- | Create a dead-letter config using topic-based routing. topicDeadLetter :: RoutingKey -> Bool -> DeadLetterConfig topicDeadLetter routingKey metadata =   DeadLetterConfig@@ -130,41 +244,25 @@     }  -- | FIFO queue configuration for ordered message processing.--- Requires pgmq 1.8.0+ with FIFO indexes. data FifoConfig = FifoConfig-  { -- | Strategy for reading messages from FIFO queue+  { -- | Strategy for reading messages from FIFO queue.     readStrategy :: !FifoReadStrategy   }   deriving stock (Show, Eq, Generic)  -- | Strategy for reading messages from FIFO queues. data FifoReadStrategy-  = -- | Fill batch from same message group first (SQS-like behavior)-    -- Good for: order processing, document workflows+  = -- | Fill batch from same message group first.     ThroughputOptimized-  | -- | Fair round-robin distribution across message groups-    -- Good for: multi-tenant systems, load balancing+  | -- | Fair round-robin distribution across groups.     RoundRobin   deriving stock (Show, Eq, Generic) --- | Configuration for concurrent prefetching.--- When enabled, polls the next batch while current messages are being processed.------ Trade-off: Lower latency at the cost of visibility timeout pressure.--- Prefetched messages have their visibility timeout ticking, so ensure:--- @bufferSize * batchSize * avgProcessingTime < visibilityTimeout@-data PrefetchConfig = PrefetchConfig-  { -- | Number of batches to buffer ahead of consumption (default: 4)-    -- Higher values reduce latency but increase visibility timeout pressure-    bufferSize :: !Natural-  }-  deriving stock (Show, Eq, Generic)- -- | Default polling configuration using standard polling with 1 second interval. defaultPollingConfig :: PollingConfig defaultPollingConfig = StandardPolling {pollInterval = 1} --- | Default retry policy for transient poll errors.+-- | Default retry policy for transient database errors. defaultPollRetryConfig :: PollRetryConfig defaultPollRetryConfig =   PollRetryConfig@@ -173,23 +271,11 @@       maxBackoff = 5     } --- | Default prefetch configuration.--- Buffers 4 batches ahead, balancing latency with visibility timeout safety.+-- | Default prefetch configuration: buffers 4 batches ahead. defaultPrefetchConfig :: PrefetchConfig defaultPrefetchConfig = PrefetchConfig {bufferSize = 4}  -- | Default adapter configuration.--- Note: You must set 'queueName' before using.------ @--- let config = defaultConfig { queueName = myQueueName }--- @------ To enable prefetching:------ @--- let config = (defaultConfig myQueueName) { prefetchConfig = Just defaultPrefetchConfig }--- @ defaultConfig :: QueueName -> PgmqAdapterConfig defaultConfig name =   PgmqAdapterConfig@@ -198,8 +284,10 @@       batchSize = 1,       polling = defaultPollingConfig,       pollRetry = defaultPollRetryConfig,+      ackRetry = defaultPollRetryConfig,       deadLetterConfig = Nothing,+      haltVisibilityTimeout = Nothing,       maxRetries = 3,       fifoConfig = Nothing,-      prefetchConfig = Nothing -- Disabled by default+      prefetchConfig = Nothing     }
src/Shibuya/Adapter/Pgmq/Internal.hs view
@@ -3,14 +3,14 @@ module Shibuya.Adapter.Pgmq.Internal   ( -- * Stream Construction     pgmqSource,-    pgmqSourceWithPrefetch,     pgmqChunks,     pgmqChunksPrefetch,     pgmqMessages,-    pgmqMessagesPrefetch,+    releaseMessages,      -- * Ingested Construction     mkIngested,+    finalizeAutoDeadLetter,      -- * AckHandle Construction     mkAckHandle,@@ -27,28 +27,43 @@      -- * Utilities     nominalToSeconds,+    retryingTransient,   ) where  import Control.Concurrent (threadDelay)-import Control.Monad (when)+import Control.Monad (void, when) import Control.Monad.IO.Class (liftIO) import Data.Aeson (Value (..)) import Data.Aeson.Key qualified as Key import Data.Aeson.KeyMap qualified as KeyMap+import Data.Foldable (traverse_) import Data.Function ((&))+import Data.IORef (IORef, newIORef, readIORef, writeIORef) import Data.Int (Int32) import Data.Text qualified as Text import Data.Text.Encoding qualified as TE-import Data.Time (NominalDiffTime, nominalDiffTimeToSeconds)+import Data.Time (NominalDiffTime, addUTCTime, getCurrentTime, nominalDiffTimeToSeconds) import Data.Vector (Vector) import Data.Vector qualified as Vector-import Effectful (Eff, IOE, (:>))+import Effectful+  ( Eff,+    IOE,+    Limit (..),+    Persistence (..),+    UnliftStrategy (..),+    withUnliftStrategy,+    (:>),+  ) import Effectful.Error.Static (Error, catchError, throwError)-import Pgmq.Effectful (PgmqRuntimeError, isTransient)+import Hasql.Pool qualified as Pool+import Hasql.Transaction qualified as Transaction+import Hasql.Transaction.Sessions qualified as Transaction.Sessions+import Pgmq.Effectful (PgmqRuntimeError, fromUsageError, isTransient) import Pgmq.Effectful.Effect   ( Pgmq,     archiveMessage,+    batchChangeVisibilityTimeout,     changeVisibilityTimeout,     deleteMessage,     readGrouped,@@ -57,13 +72,12 @@     readGroupedWithPoll,     readMessage,     readWithPoll,-    sendMessage,-    sendMessageWithHeaders,-    sendTopic,-    sendTopicWithHeaders,+    setVisibilityTimeoutAt,   )+import Pgmq.Hasql.Statements.Message qualified as Msg import Pgmq.Hasql.Statements.Types-  ( MessageQuery (..),+  ( BatchVisibilityTimeoutQuery (..),+    MessageQuery (..),     ReadGrouped (..),     ReadGroupedWithPoll (..),     ReadMessage (..),@@ -72,6 +86,7 @@     SendMessageWithHeaders (..),     SendTopic (..),     SendTopicWithHeaders (..),+    VisibilityTimeoutAtQuery (..),     VisibilityTimeoutQuery (..),   ) import Pgmq.Types qualified as Pgmq@@ -81,6 +96,7 @@     FifoConfig (..),     FifoReadStrategy (..),     PgmqAdapterConfig (..),+    PgmqAdapterEnv (..),     PollRetryConfig (..),     PollingConfig (..),   )@@ -165,24 +181,29 @@  -- | Create a Lease for visibility timeout extension. mkLease ::-  (Pgmq :> es) =>-  Pgmq.QueueName ->-  Pgmq.MessageId ->-  Lease es-mkLease queueName msgId =-  Lease-    { leaseId = Text.pack (show (Pgmq.unMessageId msgId)),-      leaseExtend = \duration -> do-        let vtSeconds = nominalToSeconds duration-        _ <--          changeVisibilityTimeout $-            VisibilityTimeoutQuery-              { queueName = queueName,-                messageId = msgId,-                visibilityTimeoutOffset = vtSeconds-              }-        pure ()-    }+  (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es) =>+  PgmqAdapterConfig ->+  Pgmq.Message ->+  Eff es (Lease es)+mkLease config msg = do+  lastVtRef <- liftIO $ newIORef msg.visibilityTime+  pure+    Lease+      { leaseId = Text.pack (show (Pgmq.unMessageId msg.messageId)),+        leaseExtend = \duration -> do+          now <- liftIO getCurrentTime+          lastVt <- liftIO $ readIORef lastVtRef+          let target = max lastVt (addUTCTime duration now)+          updated <-+            retryingTransient config.ackRetry $+              setVisibilityTimeoutAt $+                VisibilityTimeoutAtQuery+                  { queueName = config.queueName,+                    messageId = msg.messageId,+                    visibilityTime = target+                  }+          liftIO $ writeIORef lastVtRef updated.visibilityTime+      }  -- | Create an AckHandle for a message. --@@ -197,98 +218,125 @@ -- pre-0.5.0.0 behavior. See plan 1 / Finding F3 in the parent -- shibuya repo's plan 9. mkAckHandle ::-  (Pgmq :> es, IOE :> es, Tracing :> es) =>+  (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es, Tracing :> es) =>+  PgmqAdapterEnv ->   PgmqAdapterConfig ->+  IORef Bool ->   Pgmq.Message ->   AckHandle es-mkAckHandle config msg = AckHandle $ \decision -> do-  let queueName = config.queueName-      msgId = msg.messageId+mkAckHandle env config finalizedRef msg = AckHandle $ \decision -> do+  alreadyFinalized <- liftIO $ readIORef finalizedRef+  if alreadyFinalized+    then pure ()+    else do+      runDecision decision+      liftIO $ writeIORef finalizedRef True+  where+    queueName = config.queueName+    msgId = msg.messageId -  case decision of-    AckOk ->-      -- Successfully processed - delete from queue-      void $ deleteMessage (MessageQuery queueName msgId)-    AckRetry (RetryDelay delay) -> do-      -- Retry after delay - extend visibility timeout-      let vtSeconds = nominalToSeconds delay-      void $-        changeVisibilityTimeout $-          VisibilityTimeoutQuery-            { queueName = queueName,-              messageId = msgId,-              visibilityTimeoutOffset = vtSeconds-            }-    AckDeadLetter reason -> do-      -- Handle dead-lettering-      case config.deadLetterConfig of-        Nothing ->-          -- No DLQ configured - just archive the message-          void $ archiveMessage (MessageQuery queueName msgId)-        Just dlqConfig -> do-          -- Build DLQ headers: pull the consumer's current trace-          -- context (Nothing if tracing is off or no active span);-          -- merge with the original message's headers (consumer's-          -- traceparent wins, original preserved under the-          -- x-shibuya-upstream-* keys).-          consumerHdrs <- currentTraceHeaders-          let dlqBody = mkDlqPayload msg reason dlqConfig.includeMetadata-              dlqHeaders = mergeDlqHeaders consumerHdrs msg.headers-          case dlqConfig.dlqTarget of-            DirectQueue dlqQueueName ->-              case dlqHeaders of-                Just headers ->-                  void $-                    sendMessageWithHeaders $-                      SendMessageWithHeaders-                        { queueName = dlqQueueName,-                          messageBody = dlqBody,-                          messageHeaders = Pgmq.MessageHeaders headers,-                          delay = Nothing-                        }-                Nothing ->-                  void $-                    sendMessage $-                      SendMessage-                        { queueName = dlqQueueName,-                          messageBody = dlqBody,-                          delay = Nothing-                        }-            TopicRoute routingKey ->-              case dlqHeaders of-                Just headers ->-                  void $-                    sendTopicWithHeaders $-                      SendTopicWithHeaders-                        { routingKey = routingKey,-                          messageBody = dlqBody,-                          messageHeaders = Pgmq.MessageHeaders headers,-                          delay = Nothing-                        }-                Nothing ->-                  void $-                    sendTopic $-                      SendTopic-                        { routingKey = routingKey,-                          messageBody = dlqBody,-                          delay = Nothing-                        }-          -- Delete from original queue-          void $ deleteMessage (MessageQuery queueName msgId)-    AckHalt _reason -> do-      -- Halt processing - extend VT far into future-      -- Message becomes visible again after processor restarts-      let vtSeconds = 3600 :: Int32 -- 1 hour-      void $-        changeVisibilityTimeout $-          VisibilityTimeoutQuery-            { queueName = queueName,-              messageId = msgId,-              visibilityTimeoutOffset = vtSeconds-            }+    runDecision decision = retryingTransient config.ackRetry $ case decision of+      AckOk -> do+        -- Successfully processed - delete from queue+        void $ deleteMessage (MessageQuery queueName msgId)+      AckRetry (RetryDelay delay) -> do+        -- Retry after delay - extend visibility timeout+        let vtSeconds = nominalToSeconds delay+        void $+          changeVisibilityTimeout $+            VisibilityTimeoutQuery+              { queueName = queueName,+                messageId = msgId,+                visibilityTimeoutOffset = vtSeconds+              }+      AckDeadLetter reason -> do+        -- Handle dead-lettering+        case config.deadLetterConfig of+          Nothing ->+            -- No DLQ configured - just archive the message+            void $ archiveMessage (MessageQuery queueName msgId)+          Just dlqConfig -> do+            consumerHdrs <- currentTraceHeaders+            deadLetterTransactionally env config dlqConfig msg reason (mergeDlqHeaders consumerHdrs msg.headers)+      AckHalt _reason -> do+        -- Park the message by reassigning its visibility timeout. It becomes+        -- visible to any consumer after this many seconds, independent of restart.+        let vtSeconds = maybe config.visibilityTimeout id config.haltVisibilityTimeout+        void $+          changeVisibilityTimeout $+            VisibilityTimeoutQuery+              { queueName = queueName,+                messageId = msgId,+                visibilityTimeoutOffset = vtSeconds+              }++deadLetterTransactionally ::+  (Error PgmqRuntimeError :> es, IOE :> es) =>+  PgmqAdapterEnv ->+  PgmqAdapterConfig ->+  DeadLetterConfig ->+  Pgmq.Message ->+  DeadLetterReason ->+  Maybe Value ->+  Eff es ()+deadLetterTransactionally env config dlqConfig msg reason dlqHeaders = do+  result <- liftIO $ Pool.use env.pool session+  case result of+    Left err -> throwError (fromUsageError err)+    Right () -> pure ()   where-    void :: (Functor f) => f a -> f ()-    void = fmap (const ())+    dlqBody = mkDlqPayload msg reason dlqConfig.includeMetadata+    sourceQuery = MessageQuery config.queueName msg.messageId+    session =+      Transaction.Sessions.transaction+        Transaction.Sessions.ReadCommitted+        Transaction.Sessions.Write+        tx+    tx = do+      case dlqConfig.dlqTarget of+        DirectQueue dlqQueueName ->+          case dlqHeaders of+            Just headers ->+              void $+                Transaction.statement+                  SendMessageWithHeaders+                    { queueName = dlqQueueName,+                      messageBody = dlqBody,+                      messageHeaders = Pgmq.MessageHeaders headers,+                      delay = Nothing+                    }+                  Msg.sendMessageWithHeaders+            Nothing ->+              void $+                Transaction.statement+                  SendMessage+                    { queueName = dlqQueueName,+                      messageBody = dlqBody,+                      delay = Nothing+                    }+                  Msg.sendMessage+        TopicRoute routingKey ->+          case dlqHeaders of+            Just headers ->+              void $+                Transaction.statement+                  SendTopicWithHeaders+                    { routingKey = routingKey,+                      messageBody = dlqBody,+                      messageHeaders = Pgmq.MessageHeaders headers,+                      delay = Nothing+                    }+                  Msg.sendTopicWithHeaders+            Nothing ->+              void $+                Transaction.statement+                  SendTopic+                    { routingKey = routingKey,+                      messageBody = dlqBody,+                      delay = Nothing+                    }+                  Msg.sendTopic+      void $ Transaction.statement sourceQuery Msg.deleteMessage  -- | Merge the consumer's current trace headers with the original -- message's headers JSON for the DLQ-write path.@@ -338,24 +386,31 @@     -- Convert TraceHeaders ([(ByteString, ByteString)]) to a JSON object.     traceHeadersToKeyMap hdrs =       KeyMap.fromList-        [ (Key.fromText (TE.decodeUtf8 k), String (TE.decodeUtf8 v))+        [ (Key.fromText (TE.decodeUtf8Lenient k), String (TE.decodeUtf8Lenient v))         | (k, v) <- hdrs         ]  -- | Create an Ingested from a pgmq Message. -- Handles auto dead-lettering when maxRetries is exceeded. mkIngested ::-  (Pgmq :> es, IOE :> es, Tracing :> es) =>+  (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es, Tracing :> es) =>+  PgmqAdapterEnv ->   PgmqAdapterConfig ->   Pgmq.Message ->   Eff es (Maybe (Ingested es Value))-mkIngested config msg = do+mkIngested env config msg = do+  finalizedRef <- liftIO $ newIORef False+  lease <- mkLease config msg+  let ackHandle = mkAckHandle env config finalizedRef msg   -- Check if max retries exceeded   if msg.readCount > config.maxRetries     then do       -- Auto dead-letter messages that exceed retry limit-      let ackHandle = mkAckHandle config msg-      ackHandle.finalize (AckDeadLetter MaxRetriesExceeded)+      finalizeAutoDeadLetter+        msg+        env.onAutoDeadLetter+        env.onAckFailure+        (ackHandle.finalize (AckDeadLetter MaxRetriesExceeded))       -- Return Nothing - this message won't be processed by handler       pure Nothing     else@@ -363,10 +418,21 @@         Just           Ingested             { envelope = pgmqMessageToEnvelope msg,-              ack = mkAckHandle config msg,-              lease = Just (mkLease config.queueName msg.messageId)+              ack = ackHandle,+              lease = Just lease             } +finalizeAutoDeadLetter ::+  (Error PgmqRuntimeError :> es, IOE :> es) =>+  Pgmq.Message ->+  (Pgmq.Message -> IO ()) ->+  (Pgmq.Message -> PgmqRuntimeError -> IO ()) ->+  Eff es () ->+  Eff es ()+finalizeAutoDeadLetter msg onAutoDeadLetter onAckFailure finalizeAction =+  (finalizeAction >> liftIO (onAutoDeadLetter msg))+    `catchError` \_callStack err -> liftIO (onAckFailure msg err)+ -- | Stream of message batches from pgmq. -- Each element is a Vector of messages from a single poll. -- This is the lowest-level stream that handles polling logic.@@ -374,27 +440,8 @@   (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es) =>   PgmqAdapterConfig ->   Stream (Eff es) (Vector Pgmq.Message)-pgmqChunks config = Stream.repeatM (pollRetrying 1 initialBackoff)+pgmqChunks config = Stream.repeatM (retryingTransient config.pollRetry poll)   where-    PollRetryConfig-      { maxAttempts = retryMaxAttempts,-        initialBackoff = initialBackoff,-        maxBackoff = retryMaxBackoff-      } = config.pollRetry--    pollRetrying ::-      (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es) =>-      Int ->-      NominalDiffTime ->-      Eff es (Vector Pgmq.Message)-    pollRetrying attempt backoff =-      poll `catchError` \_callStack err ->-        if isTransient err && attempt < retryMaxAttempts-          then do-            liftIO $ threadDelay (nominalToMicros backoff)-            pollRetrying (attempt + 1) (min (backoff * 2) retryMaxBackoff)-          else throwError err-     poll :: (Pgmq :> es, IOE :> es) => Eff es (Vector Pgmq.Message)     poll = case config.fifoConfig of       Nothing -> pollNonFifo@@ -431,6 +478,49 @@     nominalToMicros :: NominalDiffTime -> Int     nominalToMicros t = floor (nominalDiffTimeToSeconds t * 1_000_000) +-- | Chunk polling with concurrent prefetch.+--+-- Wraps 'pgmqChunks' in streamly's @parBuffered@ so the next batches are polled+-- on a background worker while the current messages are processed. @parBuffered@+-- forks worker threads that must unlift @Eff es@ to @IO@; effectful's default+-- 'SeqUnlift' strategy throws when its unlift runs off-thread, which is the+-- historical prefetch deadlock. We therefore run the concurrent portion under+-- 'ConcUnlift' (which clones the effect environment per worker thread), scoped+-- locally with 'withUnliftStrategy' via 'Stream.morphInner' so the override is+-- in force at the moment @parBuffered@ forks — and so it does /not/ leak to the+-- non-prefetch path, which keeps running under 'SeqUnlift'.+--+-- 'Stream.morphInner' is applied /after/ @parBuffered@ so it wraps @parBuffered@'s+-- own forking step. Scoping the strategy at stream-construction time instead+-- would not work: the fork happens when the stream is run, not when it is built.+pgmqChunksPrefetch ::+  (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es) =>+  (StreamP.Config -> StreamP.Config) ->+  PgmqAdapterConfig ->+  Stream (Eff es) (Vector Pgmq.Message)+pgmqChunksPrefetch prefetchSettings config =+  pgmqChunks config+    & StreamP.parBuffered prefetchSettings+    & Stream.morphInner (withUnliftStrategy (ConcUnlift Ephemeral Unlimited))++retryingTransient ::+  (Error PgmqRuntimeError :> es, IOE :> es) =>+  PollRetryConfig ->+  Eff es a ->+  Eff es a+retryingTransient retry action = go 1 retry.initialBackoff+  where+    go attempt backoff =+      action `catchError` \_callStack err ->+        if isTransient err && attempt < retry.maxAttempts+          then do+            liftIO $ threadDelay (nominalToMicros backoff)+            go (attempt + 1) (min (backoff * 2) retry.maxBackoff)+          else throwError err++    nominalToMicros :: NominalDiffTime -> Int+    nominalToMicros t = floor (nominalDiffTimeToSeconds t * 1_000_000)+ -- | Flatten message chunks into individual messages. -- Uses Streamly's unfoldEach to expand each Vector into individual elements, -- ensuring ALL messages from each batch are processed (not just the first).@@ -451,63 +541,30 @@ -- Uses unfoldEach to process ALL messages from each batch, not just the first. pgmqSource ::   (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es, Tracing :> es) =>+  PgmqAdapterEnv ->   PgmqAdapterConfig ->   Stream (Eff es) (Ingested es Value)-pgmqSource config =+pgmqSource env config =   pgmqMessages config-    & Stream.mapMaybeM (mkIngested config) -- Convert + filter auto-DLQ'd messages---- | Stream of message batches with concurrent prefetching.--- Uses parBuffered to poll the next batch while current batch is being processed.--- This reduces latency by overlapping polling with message processing.------ Note: Prefetched messages have their visibility timeout ticking. Ensure--- bufferSize * batchSize * avgProcessingTime < visibilityTimeout to avoid--- messages re-appearing before they're processed.-pgmqChunksPrefetch ::-  (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es) =>-  (StreamP.Config -> StreamP.Config) ->-  PgmqAdapterConfig ->-  Stream (Eff es) (Vector Pgmq.Message)-pgmqChunksPrefetch prefetchConfig config =-  pgmqChunks config-    & StreamP.parBuffered prefetchConfig+    & Stream.mapMaybeM (mkIngested env config) -- Convert + filter auto-DLQ'd messages --- | Flatten prefetched message chunks into individual messages.--- Like pgmqMessages but with concurrent prefetching of batches.-pgmqMessagesPrefetch ::+releaseMessages ::   (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es) =>-  (StreamP.Config -> StreamP.Config) ->-  PgmqAdapterConfig ->-  Stream (Eff es) Pgmq.Message-pgmqMessagesPrefetch prefetchConfig config =-  pgmqChunksPrefetch prefetchConfig config-    & Stream.filter (not . Vector.null) -- Skip empty batches-    & Stream.unfoldEach vectorUnfold -- Flatten Vector to individual elements-  where-    vectorUnfold = Unfold.unfoldr Vector.uncons---- | Create message source stream with concurrent prefetching.--- Polls the next batch while current messages are being processed.------ This provides lower latency than pgmqSource by keeping messages ready--- in a buffer for immediate consumption. The trade-off is that prefetched--- messages have their visibility timeout ticking.------ Usage:------ @--- -- With default prefetch settings (4 batches ahead)--- source = pgmqSourceWithPrefetch defaultPrefetchConfig config------ -- With custom buffer size--- source = pgmqSourceWithPrefetch (StreamP.maxBuffer 2) config--- @-pgmqSourceWithPrefetch ::-  (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es, Tracing :> es) =>-  (StreamP.Config -> StreamP.Config) ->+  PgmqAdapterEnv ->   PgmqAdapterConfig ->-  Stream (Eff es) (Ingested es Value)-pgmqSourceWithPrefetch prefetchConfig config =-  pgmqMessagesPrefetch prefetchConfig config-    & Stream.mapMaybeM (mkIngested config)+  Vector Pgmq.Message ->+  Eff es ()+releaseMessages env config messages = do+  let msgIds = fmap (.messageId) (Vector.toList messages)+  retryingTransient+    config.ackRetry+    ( void $+        batchChangeVisibilityTimeout $+          BatchVisibilityTimeoutQuery+            { queueName = config.queueName,+              messageIds = msgIds,+              visibilityTimeoutOffset = 0+            }+    )+    `catchError` \_callStack err ->+      liftIO $ traverse_ (\message -> env.onAckFailure message err) messages
test/Shibuya/Adapter/Pgmq/ChaosSpec.hs view
@@ -15,6 +15,7 @@ import Data.Aeson (Value (..), object, (.=)) import Data.Aeson.KeyMap qualified as KeyMap import Data.IORef (IORef, atomicModifyIORef', newIORef, readIORef)+import Data.Int (Int32) import Data.Text qualified as Text import Data.Vector qualified as Vector import Effectful (Eff, IOE, liftIO, runEff, (:>))@@ -24,26 +25,36 @@ import Pgmq.Effectful qualified as PgmqEff import Pgmq.Hasql.Sessions qualified as Sessions import Pgmq.Hasql.Statements.Types (ReadMessage (..), SendMessage (..), SendMessageWithHeaders (..))-import Pgmq.Types (MessageBody (..), MessageHeaders (..))+import Pgmq.Types (MessageBody (..), MessageHeaders (..), QueueName)+import Shibuya.Adapter (Adapter) import Shibuya.Adapter.Pgmq   ( PgmqAdapterConfig (..),     PollingConfig (..),+    PrefetchConfig (..),     defaultConfig,+    defaultPrefetchConfig,     directDeadLetter,+    mkPgmqAdapterEnv,     pgmqAdapter,   )+import Shibuya.Adapter.Pgmq.Internal (mkIngested) import Shibuya.App-  ( ShutdownConfig (..),+  ( AppConfig (..),+    ProcessorId (..),+    ShutdownConfig (..),     SupervisionStrategy (..),+    defaultAppConfig,     mkProcessor,     runApp,     stopAppGracefully,   )-import Shibuya.Core.Ack (AckDecision (..), DeadLetterReason (..))+import Shibuya.Core.Ack (AckDecision (..), DeadLetterReason (..), HaltReason (..))+import Shibuya.Core.AckHandle (AckHandle (..))+import Shibuya.Core.Ingested (Ingested (..)) import Shibuya.Handler (Handler)-import Shibuya.Runner.Metrics (ProcessorId (..))-import Shibuya.Telemetry.Effect (runTracingNoop)+import Shibuya.Telemetry.Effect (Tracing, runTracingNoop) import System.Environment (lookupEnv)+import System.Timeout (timeout) import Test.Hspec import TmpPostgres (TestFixture (..), runPgmqSession, withPgmqDb, withTestFixture) @@ -54,6 +65,7 @@       poisonMessageSpec       longHandlerSpec       gracefulShutdownSpec+      prefetchSpec  -- | Wrapper to run tests with a temporary database and fixture withTempDbFixture :: (TestFixture -> IO ()) -> IO ()@@ -153,11 +165,11 @@      -- Run processor that dead-letters the message     runAdapterIO pool $ runTracingNoop $ do-      adapter <- pgmqAdapter config+      adapter <- requireAdapter pool config       let handler = deadLetterHandler processedRef           processor = mkProcessor adapter handler -      result <- runApp IgnoreFailures 100 [(ProcessorId "dlq-test", processor)]+      result <- runApp defaultAppConfig [(ProcessorId "dlq-test", processor)]       case result of         Left err -> liftIO $ expectationFailure $ "Failed to start app: " <> show err         Right appHandle -> do@@ -223,11 +235,11 @@      -- Run processor that dead-letters the message     runAdapterIO pool $ runTracingNoop $ do-      adapter <- pgmqAdapter config+      adapter <- requireAdapter pool config       let handler = deadLetterHandler processedRef           processor = mkProcessor adapter handler -      result <- runApp IgnoreFailures 100 [(ProcessorId "dlq-trace-test", processor)]+      result <- runApp defaultAppConfig [(ProcessorId "dlq-trace-test", processor)]       case result of         Left err -> liftIO $ expectationFailure $ "Failed to start app: " <> show err         Right appHandle -> do@@ -265,6 +277,155 @@           _ -> expectationFailure "tracestate header should be a string"       Just _ -> expectationFailure "DLQ message headers should be an object" +  it "AckDeadLetter is idempotent after a successful finalize" $ \TestFixture {pool, queueName, dlqName} -> do+    runPgmqSession pool $ do+      _ <-+        Sessions.sendMessage $+          SendMessage+            { queueName = queueName,+              messageBody = MessageBody (String "idempotent-dlq"),+              delay = Just 0+            }+      pure ()++    let config =+          (defaultConfig queueName)+            { visibilityTimeout = 5,+              batchSize = 1,+              polling = StandardPolling {pollInterval = 0.1},+              deadLetterConfig = Just $ directDeadLetter dlqName True+            }++    runAdapterIO pool $ runTracingNoop $ do+      msgs <-+        PgmqEff.readMessage $+          ReadMessage+            { queueName = queueName,+              delay = 30,+              batchSize = Just 1,+              conditional = Nothing+            }+      case Vector.uncons msgs of+        Nothing -> liftIO $ expectationFailure "expected one source message"+        Just (msg, _) -> do+          ingestedResult <- mkIngested (mkPgmqAdapterEnv pool) config msg+          case ingestedResult of+            Nothing -> liftIO $ expectationFailure "message should not auto-DLQ"+            Just Ingested {ack = AckHandle finalize} -> do+              finalize (AckDeadLetter (PoisonPill "first"))+              finalize (AckDeadLetter (PoisonPill "second"))++    dlqMsgs <-+      runPgmqSession pool $+        Sessions.readMessage $+          ReadMessage+            { queueName = dlqName,+              delay = 30,+              batchSize = Just 10,+              conditional = Nothing+            }+    Vector.length dlqMsgs `shouldBe` 1++    sourceMsgs <-+      runPgmqSession pool $+        Sessions.readMessage $+          ReadMessage+            { queueName = queueName,+              delay = 30,+              batchSize = Just 10,+              conditional = Nothing+            }+    Vector.length sourceMsgs `shouldBe` 0++  it "AckOk is idempotent after a successful finalize" $ \TestFixture {pool, queueName, dlqName = _} -> do+    runPgmqSession pool $ do+      _ <-+        Sessions.sendMessage $+          SendMessage+            { queueName = queueName,+              messageBody = MessageBody (String "idempotent-ok"),+              delay = Just 0+            }+      pure ()++    let config = (defaultConfig queueName) {visibilityTimeout = 5, batchSize = 1}++    runAdapterIO pool $ runTracingNoop $ do+      msgs <-+        PgmqEff.readMessage $+          ReadMessage+            { queueName = queueName,+              delay = 30,+              batchSize = Just 1,+              conditional = Nothing+            }+      case Vector.uncons msgs of+        Nothing -> liftIO $ expectationFailure "expected one source message"+        Just (msg, _) -> do+          ingestedResult <- mkIngested (mkPgmqAdapterEnv pool) config msg+          case ingestedResult of+            Nothing -> liftIO $ expectationFailure "message should not auto-DLQ"+            Just Ingested {ack = AckHandle finalize} -> do+              finalize AckOk+              finalize AckOk++    sourceMsgs <-+      runPgmqSession pool $+        Sessions.readMessage $+          ReadMessage+            { queueName = queueName,+              delay = 30,+              batchSize = Just 10,+              conditional = Nothing+            }+    Vector.length sourceMsgs `shouldBe` 0++  it "AckHalt uses haltVisibilityTimeout when configured" $ \TestFixture {pool, queueName, dlqName = _} -> do+    runPgmqSession pool $ do+      _ <-+        Sessions.sendMessage $+          SendMessage+            { queueName = queueName,+              messageBody = MessageBody (String "halt-vt"),+              delay = Just 0+            }+      pure ()++    let config =+          (defaultConfig queueName)+            { visibilityTimeout = 30,+              haltVisibilityTimeout = Just 1+            }++    runAdapterIO pool $ runTracingNoop $ do+      msgs <-+        PgmqEff.readMessage $+          ReadMessage+            { queueName = queueName,+              delay = 30,+              batchSize = Just 1,+              conditional = Nothing+            }+      case Vector.uncons msgs of+        Nothing -> liftIO $ expectationFailure "expected one source message"+        Just (msg, _) -> do+          ingestedResult <- mkIngested (mkPgmqAdapterEnv pool) config msg+          case ingestedResult of+            Nothing -> liftIO $ expectationFailure "message should not auto-DLQ"+            Just Ingested {ack = AckHandle finalize} -> finalize (AckHalt (HaltFatal "pause"))++    threadDelay 1500000+    visible <-+      runPgmqSession pool $+        Sessions.readMessage $+          ReadMessage+            { queueName = queueName,+              delay = 30,+              batchSize = Just 1,+              conditional = Nothing+            }+    Vector.length visible `shouldBe` 1+ -------------------------------------------------------------------------------- -- Long Handler Tests --------------------------------------------------------------------------------@@ -296,11 +457,11 @@      -- Start processor with slow handler     processorAsync <- async $ runAdapterIO pool $ runTracingNoop $ do-      adapter <- pgmqAdapter config+      adapter <- requireAdapter pool config       let handler = slowHandler processedRef 2000000 -- 2 second delay           processor = mkProcessor adapter handler -      result <- runApp IgnoreFailures 100 [(ProcessorId "slow-test", processor)]+      result <- runApp defaultAppConfig [(ProcessorId "slow-test", processor)]       case result of         Left err -> liftIO $ expectationFailure $ "Failed to start app: " <> show err         Right appHandle -> do@@ -365,11 +526,11 @@      -- Run processor with slow handler     runAdapterIO pool $ runTracingNoop $ do-      adapter <- pgmqAdapter config+      adapter <- requireAdapter pool config       let handler = slowHandler processedRef 50000 -- 0.05 second delay per message           processor = mkProcessor adapter handler -      appResult <- runApp IgnoreFailures 100 [(ProcessorId "drain-test", processor)]+      appResult <- runApp defaultAppConfig [(ProcessorId "drain-test", processor)]       case appResult of         Left err -> liftIO $ expectationFailure $ "Failed to start app: " <> show err         Right appHandle -> do@@ -408,11 +569,11 @@             }      runAdapterIO pool $ runTracingNoop $ do-      adapter <- pgmqAdapter config+      adapter <- requireAdapter pool config       let handler = countingHandler processedRef           processor = mkProcessor adapter handler -      appResult <- runApp IgnoreFailures 100 [(ProcessorId "stop-test", processor)]+      appResult <- runApp defaultAppConfig [(ProcessorId "stop-test", processor)]       case appResult of         Left err -> liftIO $ expectationFailure $ "Failed to start app: " <> show err         Right appHandle -> do@@ -454,6 +615,219 @@     remaining `shouldSatisfy` (>= 0) -- Just ensure no crash  --------------------------------------------------------------------------------+-- Prefetch Tests (scoped ConcUnlift)+--------------------------------------------------------------------------------++-- | Regression test for the historical prefetch deadlock.+--+-- With @prefetchConfig = Just ...@ the polling stage runs on a streamly+-- @parBuffered@ worker thread that must unlift @Eff@. Under effectful's default+-- 'SeqUnlift' this deadlocks ("thread blocked indefinitely in an STM+-- transaction"); the adapter now scopes the concurrent stage to 'ConcUnlift'+-- (see 'Shibuya.Adapter.Pgmq.Internal.pgmqChunksPrefetch').+--+-- Acceptance is behavioral: the queue drains to completion within a timeout. If+-- the deadlock were present the 'timeout' fires and the test fails instead of+-- hanging the suite. Stripping the @morphInner (withUnliftStrategy ...)@ wrapper+-- from 'pgmqChunksPrefetch' makes this test fail (verified during M2).+prefetchSpec :: SpecWith TestFixture+prefetchSpec = describe "Prefetch" $ do+  it "drains the queue under concurrent prefetch without deadlocking" $ \TestFixture {pool, queueName, dlqName = _} -> do+    let total = 50 :: Int+    forM_ [1 .. total] $ \i ->+      runPgmqSession pool $ do+        _ <-+          Sessions.sendMessage $+            SendMessage+              { queueName = queueName,+                messageBody = MessageBody (String $ Text.pack ("prefetch-" <> show i)),+                delay = Just 0+              }+        pure ()++    processedRef <- newIORef (0 :: Int)++    let config =+          (defaultConfig queueName)+            { visibilityTimeout = 30,+              batchSize = 5,+              polling = StandardPolling {pollInterval = 0.05},+              prefetchConfig = Just defaultPrefetchConfig+            }++    -- Guard against a hang: a deadlock makes 'timeout' return Nothing and the+    -- test fails, rather than blocking the whole suite indefinitely.+    result <- timeout 30_000_000 $ runAdapterIO pool $ runTracingNoop $ do+      adapter <- requireAdapter pool config+      let handler = countingHandler processedRef+          processor = mkProcessor adapter handler+      appResult <- runApp defaultAppConfig [(ProcessorId "prefetch-test", processor)]+      case appResult of+        Left err -> liftIO $ expectationFailure $ "Failed to start app: " <> show err+        Right appHandle -> do+          liftIO $ waitForProcessed processedRef total 20_000_000+          let shutdownConfig = ShutdownConfig {drainTimeout = 5}+          _ <- stopAppGracefully shutdownConfig appHandle+          pure ()++    case result of+      Nothing ->+        expectationFailure+          "Prefetch run timed out — likely the parBuffered/effectful STM deadlock"+      Just () -> pure ()++    processed <- readIORef processedRef+    processed `shouldBe` total++    -- All processed messages were AckOk'd, so the source queue is drained.+    remaining <-+      runPgmqSession pool $ do+        msgs <-+          Sessions.readMessage $+            ReadMessage+              { queueName = queueName,+                delay = 30,+                batchSize = Just 100,+                conditional = Nothing+              }+        pure $ Vector.length msgs+    remaining `shouldBe` 0++  -- Differential shutdown-release measurement: does prefetch hold messages+  -- invisible at shutdown that the non-prefetch path would release?+  --+  -- 'parBuffered' buffers chunks upstream of the EP-27 shutdown gate+  -- ('takeWhileM'/'releaseMessages'). On shutdown the gate releases only the+  -- chunk it pulls next, so chunks still sitting in the parBuffered buffer may+  -- never reach 'releaseMessages' and stay invisible until their VT expires.+  --+  -- The metric @held = total - processed - visibleImmediatelyAfterShutdown@+  -- counts messages left invisible in the queue (inbox residual + any buffer+  -- leak); it cancels out never-read and deleted messages. We run the identical+  -- scenario with prefetch OFF (on the main queue) and ON (on the DLQ queue,+  -- reused here as an independent second queue) and compare.+  it "strands only a bounded (<= one prefetch buffer) number of extra messages invisible on shutdown" $ \TestFixture {pool, queueName, dlqName} -> do+    (totalOff, procOff, visOff) <- liftIO $ measureShutdownRelease pool queueName Nothing+    (totalOn, procOn, visOn) <- liftIO $ measureShutdownRelease pool dlqName (Just defaultPrefetchConfig)++    let heldOff = totalOff - procOff - visOff+        heldOn = totalOn - procOn - visOn++    liftIO $+      putStrLn $+        "PREFETCH-SHUTDOWN-DIAG: off(total="+          <> show totalOff+          <> " processed="+          <> show procOff+          <> " visible="+          <> show visOff+          <> " held="+          <> show heldOff+          <> ") on(total="+          <> show totalOn+          <> " processed="+          <> show procOn+          <> " visible="+          <> show visOn+          <> " held="+          <> show heldOn+          <> ")"++    -- ACCEPTED, DOCUMENTED behaviour (M3 decision, 2026-07-04): prefetch strands+    -- up to bufferSize*batchSize extra messages invisible on shutdown, because+    -- 'parBuffered' buffers chunks upstream of the shutdown gate and only the+    -- chunk the gate pulls next is released. This is bounded and loss-free (the+    -- messages are redelivered after their VT — proven by the "loses no messages+    -- under prefetch" test), and is documented in the 'PrefetchConfig' Haddock+    -- and the adapter ARCHITECTURE docs.+    --+    -- This assertion is a regression guard on the *bound*, not the exact count:+    -- the messages held invisible with prefetch must not exceed one full prefetch+    -- buffer (bufferSize * batchSize) plus the ordinary core-inbox residual. An+    -- unbounded leak (stranding most of the queue) would fail here. We bound+    -- heldOn absolutely rather than (heldOn - heldOff), because the baseline's+    -- inbox residual varies run to run (heldOff has been observed at 1–3). Config+    -- under test: batchSize=2, defaultPrefetchConfig bufferSize=4, runApp inbox=2.+    -- Measured numbers are on the PREFETCH-SHUTDOWN-DIAG line above.+    let bufferBound = fromIntegral defaultPrefetchConfig.bufferSize * 2 -- bufferSize * batchSize+        inboxResidualSlack = 6 -- generous allowance for the inbox=2 + in-flight residual+    heldOn `shouldSatisfy` (<= bufferBound + inboxResidualSlack)++  -- Proves the "accept and document" decision is safe: the shutdown strand+  -- above delays redelivery but LOSES NOTHING. A blocking handler forces+  -- messages to be read ahead into the prefetch buffer without ever being+  -- acked; after shutdown and one visibility-timeout window, every message is+  -- recoverable (processed + still-in-queue == total). This is the empirical+  -- basis for the no-data-loss claim in the docs/Haddocks.+  it "loses no messages under prefetch: stranded messages are all redelivered after the VT" $ \TestFixture {pool, queueName, dlqName = _} -> do+    let total = 20 :: Int+        vt = 4 :: Int32+    forM_ [1 .. total] $ \i ->+      runPgmqSession pool $ do+        _ <-+          Sessions.sendMessage $+            SendMessage+              { queueName = queueName,+                messageBody = MessageBody (String $ Text.pack ("noloss-" <> show i)),+                delay = Just 0+              }+        pure ()++    processedRef <- newIORef (0 :: Int)++    let config =+          (defaultConfig queueName)+            { visibilityTimeout = vt,+              batchSize = 2,+              polling = StandardPolling {pollInterval = 0.02},+              prefetchConfig = Just defaultPrefetchConfig+            }++    -- Handler blocks for far longer than the test window, so messages are+    -- read-ahead into the buffer/inbox but never acked before shutdown.+    _ <- timeout 30_000_000 $ runAdapterIO pool $ runTracingNoop $ do+      adapter <- requireAdapter pool config+      let handler = slowHandler processedRef 10_000_000+          processor = mkProcessor adapter handler+      appResult <- runApp (AppConfig {strategy = IgnoreFailures, inboxSize = 2}) [(ProcessorId "noloss-test", processor)]+      case appResult of+        Left err -> liftIO $ expectationFailure $ "Failed to start app: " <> show err+        Right appHandle -> do+          liftIO $ threadDelay 1_000_000 -- let the adapter read a few batches ahead+          _ <- stopAppGracefully (ShutdownConfig {drainTimeout = 1}) appHandle+          pure ()++    processed <- readIORef processedRef++    -- Wait past the visibility timeout so every read-but-stranded message+    -- becomes visible again.+    threadDelay (fromIntegral vt * 1_000_000 + 2_000_000)++    recoverable <-+      runPgmqSession pool $ do+        msgs <-+          Sessions.readMessage $+            ReadMessage+              { queueName = queueName,+                delay = 30,+                batchSize = Just 200,+                conditional = Nothing+              }+        pure $ Vector.length msgs++    liftIO $+      putStrLn $+        "PREFETCH-NOLOSS-DIAG: total="+          <> show total+          <> " processed="+          <> show processed+          <> " recoverableAfterVT="+          <> show recoverable++    -- No message is lost: everything not AckOk-deleted is back in the queue.+    (processed + recoverable) `shouldBe` total++-------------------------------------------------------------------------------- -- Test Helpers -------------------------------------------------------------------------------- @@ -490,3 +864,71 @@             else do               threadDelay pollInterval               go (elapsed + pollInterval)++requireAdapter ::+  (Pgmq :> es, Error PgmqRuntimeError :> es, IOE :> es, Tracing :> es) =>+  Pool.Pool ->+  PgmqAdapterConfig ->+  Eff es (Adapter es Value)+requireAdapter pool config = do+  result <- pgmqAdapter (mkPgmqAdapterEnv pool) config+  case result of+    Left err -> liftIO $ error $ "Invalid PGMQ adapter config: " <> show err+    Right adapter -> pure adapter++-- | Send a batch of messages, run a slow-handler processor against @q@ until a+-- couple are processed, then shut down and measure how many messages remain in+-- the queue but invisible. Returns @(total, processed, visibleAfterShutdown)@.+--+-- Parameters are chosen so 'parBuffered' can buffer several small chunks ahead+-- of the shutdown gate: small @batchSize@, default @bufferSize@ (4), a slow+-- handler, and a small inbox — so the buffer is non-empty at shutdown.+measureShutdownRelease :: Pool.Pool -> QueueName -> Maybe PrefetchConfig -> IO (Int, Int, Int)+measureShutdownRelease pool q mprefetch = do+  let total = 40 :: Int+  forM_ [1 .. total] $ \i ->+    runPgmqSession pool $ do+      _ <-+        Sessions.sendMessage $+          SendMessage+            { queueName = q,+              messageBody = MessageBody (String $ Text.pack ("shutdown-" <> show i)),+              delay = Just 0+            }+      pure ()++  processedRef <- newIORef (0 :: Int)++  let config =+        (defaultConfig q)+          { visibilityTimeout = 30,+            batchSize = 2,+            polling = StandardPolling {pollInterval = 0.02},+            prefetchConfig = mprefetch+          }++  _ <- timeout 30_000_000 $ runAdapterIO pool $ runTracingNoop $ do+    adapter <- requireAdapter pool config+    let handler = slowHandler processedRef 400000 -- 0.4s per message+        processor = mkProcessor adapter handler+    appResult <- runApp (AppConfig {strategy = IgnoreFailures, inboxSize = 2}) [(ProcessorId "shutdown-measure", processor)]+    case appResult of+      Left err -> liftIO $ expectationFailure $ "Failed to start app: " <> show err+      Right appHandle -> do+        liftIO $ waitForProcessed processedRef 1 5_000_000+        _ <- stopAppGracefully (ShutdownConfig {drainTimeout = 1}) appHandle+        pure ()++  processed <- readIORef processedRef+  visibleNow <-+    runPgmqSession pool $ do+      msgs <-+        Sessions.readMessage $+          ReadMessage+            { queueName = q,+              delay = 30,+              batchSize = Just 100,+              conditional = Nothing+            }+      pure $ Vector.length msgs+  pure (total, processed, visibleNow)
test/Shibuya/Adapter/Pgmq/ConfigSpec.hs view
@@ -9,7 +9,7 @@   defaultConfigSpec   defaultPollingConfigSpec   defaultPollRetryConfigSpec-  defaultPrefetchConfigSpec+  validateConfigSpec   deadLetterTargetSpec   smartConstructorSpec @@ -38,18 +38,21 @@   it "uses the default poll retry policy" $ do     config.pollRetry `shouldBe` defaultPollRetryConfig +  it "uses the default ack retry policy" $ do+    config.ackRetry `shouldBe` defaultPollRetryConfig+   it "sets deadLetterConfig to Nothing" $ do     config.deadLetterConfig `shouldBe` Nothing +  it "sets haltVisibilityTimeout to Nothing" $ do+    config.haltVisibilityTimeout `shouldBe` Nothing+   it "sets maxRetries to 3" $ do     config.maxRetries `shouldBe` 3    it "sets fifoConfig to Nothing" $ do     config.fifoConfig `shouldBe` Nothing -  it "sets prefetchConfig to Nothing" $ do-    config.prefetchConfig `shouldBe` Nothing- -- | Tests for defaultPollingConfig defaultPollingConfigSpec :: Spec defaultPollingConfigSpec = describe "defaultPollingConfig" $ do@@ -80,11 +83,43 @@   it "caps backoff at five seconds" $ do     retryMaxBackoff `shouldBe` 5 --- | Tests for defaultPrefetchConfig-defaultPrefetchConfigSpec :: Spec-defaultPrefetchConfigSpec = describe "defaultPrefetchConfig" $ do-  it "has bufferSize of 4" $ do-    defaultPrefetchConfig.bufferSize `shouldBe` 4+validateConfigSpec :: Spec+validateConfigSpec = describe "validateConfig" $ do+  let queueName = case parseQueueName "validate_queue" of+        Right q -> q+        Left e -> error $ "Unexpected: " <> show e+      base = defaultConfig queueName++  it "accepts the default config" $ do+    validateConfig base `shouldBe` Right base++  it "rejects batchSize below one" $ do+    validateConfig base {batchSize = 0} `shouldBe` Left (InvalidBatchSize 0)++  it "rejects visibilityTimeout below one" $ do+    validateConfig base {visibilityTimeout = 0} `shouldBe` Left (InvalidVisibilityTimeout 0)++  it "rejects non-positive standard polling interval" $ do+    validateConfig base {polling = StandardPolling 0} `shouldBe` Left (InvalidStandardPollInterval 0)++  it "rejects invalid long polling fields" $ do+    validateConfig base {polling = LongPolling 0 100} `shouldBe` Left (InvalidLongPollSeconds 0)+    validateConfig base {polling = LongPolling 10 0} `shouldBe` Left (InvalidLongPollIntervalMs 0)++  it "rejects negative maxRetries but accepts zero" $ do+    validateConfig base {maxRetries = -1} `shouldBe` Left (InvalidMaxRetries (-1))+    validateConfig base {maxRetries = 0} `shouldBe` Right base {maxRetries = 0}++  it "rejects invalid halt visibility timeout" $ do+    validateConfig base {haltVisibilityTimeout = Just 0} `shouldBe` Left (InvalidHaltVisibilityTimeout 0)++  it "rejects invalid retry policies" $ do+    let invalidAttempts = defaultPollRetryConfig {maxAttempts = 0}+        invalidBackoff = defaultPollRetryConfig {initialBackoff = -1}+    validateConfig base {pollRetry = invalidAttempts} `shouldBe` Left (InvalidPollRetryMaxAttempts 0)+    validateConfig base {ackRetry = invalidAttempts} `shouldBe` Left (InvalidAckRetryMaxAttempts 0)+    validateConfig base {pollRetry = invalidBackoff} `shouldBe` Left (InvalidPollRetryBackoff (-1) 5)+    validateConfig base {ackRetry = invalidBackoff} `shouldBe` Left (InvalidAckRetryBackoff (-1) 5)  -- | Tests for DeadLetterTarget deadLetterTargetSpec :: Spec
test/Shibuya/Adapter/Pgmq/IntegrationSpec.hs view
@@ -312,7 +312,9 @@       batchSize = 1,       polling = StandardPolling {pollInterval = 0.1}, -- Fast polling for tests       pollRetry = defaultPollRetryConfig,+      ackRetry = defaultPollRetryConfig,       deadLetterConfig = Nothing,+      haltVisibilityTimeout = Nothing,       maxRetries = 3,       fifoConfig = Nothing,       prefetchConfig = Nothing@@ -327,7 +329,9 @@       batchSize = bs,       polling = StandardPolling {pollInterval = 0.1},       pollRetry = defaultPollRetryConfig,+      ackRetry = defaultPollRetryConfig,       deadLetterConfig = Nothing,+      haltVisibilityTimeout = Nothing,       maxRetries = 3,       fifoConfig = Nothing,       prefetchConfig = Nothing
test/Shibuya/Adapter/Pgmq/InternalSpec.hs view
@@ -4,6 +4,7 @@  import Data.Aeson (Value (..), object, (.=)) import Data.Aeson.KeyMap qualified as KeyMap+import Data.ByteString qualified as BS import Data.IORef (IORef, atomicModifyIORef', newIORef, readIORef) import Data.Int (Int32) import Data.Time (NominalDiffTime, UTCTime (..), fromGregorian)@@ -24,7 +25,8 @@     defaultPollingConfig,   ) import Shibuya.Adapter.Pgmq.Internal-  ( mergeDlqHeaders,+  ( finalizeAutoDeadLetter,+    mergeDlqHeaders,     mkReadGrouped,     mkReadMessage,     mkReadWithPoll,@@ -41,6 +43,7 @@   mkReadWithPollSpec   mkReadGroupedSpec   pollRetrySpec+  autoDeadLetterHookSpec   mergeDlqHeadersSpec  -- | Tests for nominalToSeconds@@ -89,7 +92,9 @@             batchSize = 10,             polling = defaultPollingConfig,             pollRetry = defaultPollRetryConfig,+            ackRetry = defaultPollRetryConfig,             deadLetterConfig = Nothing,+            haltVisibilityTimeout = Nothing,             maxRetries = 3,             fifoConfig = Nothing,             prefetchConfig = Nothing@@ -126,7 +131,9 @@             batchSize = 10,             polling = defaultPollingConfig,             pollRetry = defaultPollRetryConfig,+            ackRetry = defaultPollRetryConfig,             deadLetterConfig = Nothing,+            haltVisibilityTimeout = Nothing,             maxRetries = 3,             fifoConfig = Nothing,             prefetchConfig = Nothing@@ -171,7 +178,9 @@             batchSize = 20,             polling = defaultPollingConfig,             pollRetry = defaultPollRetryConfig,+            ackRetry = defaultPollRetryConfig,             deadLetterConfig = Nothing,+            haltVisibilityTimeout = Nothing,             maxRetries = 3,             fifoConfig = Nothing,             prefetchConfig = Nothing@@ -220,6 +229,42 @@     result `shouldBe` Left transientError     readIORef calls `shouldReturn` 2 +autoDeadLetterHookSpec :: Spec+autoDeadLetterHookSpec = describe "finalizeAutoDeadLetter" $ do+  it "calls the auto-DLQ hook only after finalize succeeds" $ do+    autoCalls <- newIORef (0 :: Int)+    failureCalls <- newIORef (0 :: Int)+    let action :: Eff '[Error PgmqRuntimeError, IOE] ()+        action =+          finalizeAutoDeadLetter+            testMessage+            (\_ -> bump autoCalls)+            (\_ _ -> bump failureCalls)+            (pure ())++    result <- runEff $ runErrorNoCallStack action++    result `shouldBe` Right ()+    readIORef autoCalls `shouldReturn` 1+    readIORef failureCalls `shouldReturn` 0++  it "calls the ack-failure hook without reporting auto-DLQ success when finalize fails" $ do+    autoCalls <- newIORef (0 :: Int)+    failureCalls <- newIORef (0 :: Int)+    let action :: Eff '[Error PgmqRuntimeError, IOE] ()+        action =+          finalizeAutoDeadLetter+            testMessage+            (\_ -> bump autoCalls)+            (\_ _ -> bump failureCalls)+            (throwError permanentError)++    result <- runEff $ runErrorNoCallStack action++    result `shouldBe` Right ()+    readIORef autoCalls `shouldReturn` 0+    readIORef failureCalls `shouldReturn` 1+ -- | Tests for mergeDlqHeaders, the helper that injects the failing -- consumer's trace context onto a DLQ message while preserving the -- original producer's trace under x-shibuya-upstream-* keys.@@ -300,6 +345,12 @@           `shouldBe` Just (String "vendor=opaque")       other -> expectationFailure $ "expected merged Object, got " <> show other +  it "leniently decodes non-UTF8 trace headers" $ do+    let consumerHdrs = Just [(BS.pack [0xff, 0xfe], BS.pack [0xff])]+    mergeDlqHeaders consumerHdrs Nothing `shouldSatisfy` \case+      Just (Object obj) -> not (KeyMap.null obj)+      _ -> False+ retryTestConfig :: Int -> PgmqAdapterConfig retryTestConfig attempts =   let queueName = case parseQueueName "retry_test" of@@ -316,7 +367,9 @@                 initialBackoff = 0,                 maxBackoff = 0               },+          ackRetry = defaultPollRetryConfig,           deadLetterConfig = Nothing,+          haltVisibilityTimeout = Nothing,           maxRetries = 3,           fifoConfig = Nothing,           prefetchConfig = Nothing@@ -359,6 +412,9 @@ permanentError :: PgmqRuntimeError permanentError =   PgmqConnectionError (HasqlErrors.AuthenticationConnectionError "bad password")++bump :: IORef Int -> IO ()+bump ref = atomicModifyIORef' ref (\n -> (n + 1, ()))  testMessage :: Message testMessage =