packages feed

shibuya-kiroku-adapter 0.5.1.5 → 0.6.0.0

raw patch · 6 files changed

+636/−125 lines, 6 filesdep +contravariant-extrasdep ~kiroku-storedep ~kiroku-test-supportPVP ok

version bump matches the API change (PVP)

Dependencies added: contravariant-extras

Dependency ranges changed: kiroku-store, kiroku-test-support

API changes (from Hackage documentation)

- Shibuya.Adapter.Kiroku: ConsumerGroup :: !Int32 -> !Int32 -> ConsumerGroup
- Shibuya.Adapter.Kiroku: [member] :: ConsumerGroup -> !Int32
- Shibuya.Adapter.Kiroku: [size] :: ConsumerGroup -> !Int32
+ Shibuya.Adapter.Kiroku: [handlerStallWarnAfter] :: KirokuConsumerGroupConfig -> !Maybe NominalDiffTime
+ Shibuya.Adapter.Kiroku: [retryPolicy] :: KirokuConsumerGroupConfig -> !RetryPolicy
+ Shibuya.Adapter.Kiroku: consumerGroupSizeValue :: ConsumerGroupSize -> Int32
+ Shibuya.Adapter.Kiroku: data ConsumerGroupSize
+ Shibuya.Adapter.Kiroku: kirokuProcessor :: forall (es :: [Effect]). Adapter es RecordedEvent -> Handler es RecordedEvent -> QueueProcessor es
+ Shibuya.Adapter.Kiroku: member :: ConsumerGroup -> Int32
+ Shibuya.Adapter.Kiroku: mkConsumerGroup :: Int32 -> ConsumerGroupSize -> Either InvalidConsumerGroup ConsumerGroup
+ Shibuya.Adapter.Kiroku: mkConsumerGroupSize :: Int32 -> Either InvalidConsumerGroup ConsumerGroupSize
+ Shibuya.Adapter.Kiroku: size :: ConsumerGroup -> Int32
- Shibuya.Adapter.Kiroku: KirokuAdapterConfig :: !SubscriptionName -> !SubscriptionTarget -> !Int32 -> !Natural -> !Natural -> !Maybe ConsumerGroup -> !MissingCheckpointPolicy -> !EventTypeFilter -> !Maybe (RecordedEvent -> Bool) -> KirokuAdapterConfig
+ Shibuya.Adapter.Kiroku: KirokuAdapterConfig :: !SubscriptionName -> !SubscriptionTarget -> !BatchSize -> !StreamBufferSize -> !RetryPolicy -> !Maybe NominalDiffTime -> !Natural -> !Maybe ConsumerGroup -> !MissingCheckpointPolicy -> !EventTypeFilter -> !Maybe (RecordedEvent -> Bool) -> KirokuAdapterConfig
- Shibuya.Adapter.Kiroku: KirokuConsumerGroupConfig :: !SubscriptionName -> !SubscriptionTarget -> !Int32 -> !Int32 -> !Natural -> !Natural -> !Concurrency -> !MissingCheckpointPolicy -> !EventTypeFilter -> !Maybe (RecordedEvent -> Bool) -> KirokuConsumerGroupConfig
+ Shibuya.Adapter.Kiroku: KirokuConsumerGroupConfig :: !SubscriptionName -> !SubscriptionTarget -> !ConsumerGroupSize -> !BatchSize -> !StreamBufferSize -> !RetryPolicy -> !Maybe NominalDiffTime -> !Natural -> !Concurrency -> !MissingCheckpointPolicy -> !EventTypeFilter -> !Maybe (RecordedEvent -> Bool) -> KirokuConsumerGroupConfig
- Shibuya.Adapter.Kiroku: [batchSize] :: KirokuConsumerGroupConfig -> !Int32
+ Shibuya.Adapter.Kiroku: [batchSize] :: KirokuConsumerGroupConfig -> !BatchSize
- Shibuya.Adapter.Kiroku: [bufferSize] :: KirokuConsumerGroupConfig -> !Natural
+ Shibuya.Adapter.Kiroku: [bufferSize] :: KirokuConsumerGroupConfig -> !StreamBufferSize
- Shibuya.Adapter.Kiroku: [groupSize] :: KirokuConsumerGroupConfig -> !Int32
+ Shibuya.Adapter.Kiroku: [groupSize] :: KirokuConsumerGroupConfig -> !ConsumerGroupSize
- Shibuya.Adapter.Kiroku: defaultConsumerGroupConfig :: SubscriptionName -> SubscriptionTarget -> Int32 -> KirokuConsumerGroupConfig
+ Shibuya.Adapter.Kiroku: defaultConsumerGroupConfig :: SubscriptionName -> SubscriptionTarget -> ConsumerGroupSize -> KirokuConsumerGroupConfig

Files

CHANGELOG.md view
@@ -1,5 +1,26 @@ # Changelog +## 0.6.0.0 — 2026-10-10++### Breaking Changes++* Both configurations use validated `BatchSize` and `StreamBufferSize` capacities.+  `KirokuConsumerGroupConfig.groupSize` and `defaultConsumerGroupConfig` take+  validated `ConsumerGroupSize`; smart constructors and read-only accessors are+  re-exported.+* Both configurations add `retryPolicy` (five total deliveries) and+  `handlerStallWarnAfter` (default `Nothing`), forwarded to each underlying worker.+  Pending raw acknowledgements remain pending; enabled warnings are advisory.++### New Features++* Add `kirokuProcessor`, composing single-processor defaults (`Unordered`,+  `Serial`) with the one-second synchronous exception retry guard.++### Other Changes++* Require `kiroku-store ^>=0.10.0.0`.+ ## 0.5.1.5 — 2026-09-25  ### Other Changes
+ LICENSE view
@@ -0,0 +1,11 @@+Copyright (c) 2026 Nadeem Bitar.++Redistribution and use in source and binary forms, with or without modification, are permitted provided that the following conditions are met:++1. Redistributions of source code must retain the above copyright notice, this list of conditions and the following disclaimer.++2. Redistributions in binary form must reproduce the above copyright notice, this list of conditions and the following disclaimer in the documentation and/or other materials provided with the distribution.++3. Neither the name of the copyright holder nor the names of its contributors may be used to endorse or promote products derived from this software without specific prior written permission.++THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
+ bench/WriteProbe.hs view
@@ -0,0 +1,262 @@+{-# LANGUAGE CPP #-}+{-# LANGUAGE OverloadedRecordDot #-}++{- | MP-12's durable append/subscription comparison probe. Build this identical+workload against the original control and candidate; only the legacy group+constructor needs a compatibility branch. No production instrumentation.+-}+module Main (main) where++import Contravariant.Extras (contrazip2)+import Control.Concurrent (threadDelay)+import Control.Concurrent.Async qualified as Async+import Control.Exception (bracket)+import Control.Monad (forM_, unless, void, when)+import Data.Aeson (encode, object, (.=))+import Data.ByteString.Lazy.Char8 qualified as Bytes+import Data.IORef+import Data.Int (Int32, Int64)+import Data.List (sort)+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Vector qualified as Vector+import Data.Word (Word64)+import Effectful (liftIO, runEff)+import EphemeralPg qualified as Pg+import GHC.Clock (getMonotonicTimeNSec)+import GHC.Stats+import Hasql.Decoders qualified as D+import Hasql.Encoders qualified as E+import Hasql.Pool qualified as Pool+import Hasql.Session qualified as Session+import Hasql.Statement (Statement, preparable)+import Kiroku.Store+import Kiroku.Test.Postgres (ephemeralConfig, migrateTestDatabase)+import Shibuya.Adapter.Kiroku (defaultKirokuAdapterConfig, kirokuAdapter)+import Shibuya.Adapter.Kiroku qualified as Adapter+import Shibuya.App (ProcessorId (..), defaultAppConfig, mkProcessor, runApp, stopApp)+import Shibuya.Core.Ack (AckDecision (..))+import Shibuya.Telemetry.Effect (runTracingNoop)+import System.Environment (getArgs)+import System.Mem (performMajorGC)+import System.Timeout (timeout)+import Text.Read (readMaybe)++data Workload = Workload+    { seconds :: !Int+    , mode :: !String+    , width :: !Int+    , appendBatch :: !Int+    , checkpointBatch :: !Int32+    , fresh :: !Bool+    , offered :: !Int+    }++main :: IO ()+main = do+    args <- getArgs+    workload <- case args of+        [duration, mode, width, batch, checkpoint, fresh, offered] ->+            Workload <$> number duration <*> pure mode <*> number width <*> number batch <*> number checkpoint <*> number fresh <*> number offered+        _ -> fail "usage: write-probe SECONDS MODE WIDTH APPEND_BATCH CHECKPOINT_BATCH FRESH OFFERED_CALLS_PER_SEC (MODE=none|all|category|group|adapter)"+    unless (workload.seconds > 0 && workload.width > 0 && workload.appendBatch > 0 && workload.checkpointBatch > 0 && workload.offered >= 0 && workload.mode `elem` ["none", "all", "category", "group", "adapter"]) (fail "invalid workload")+    original <- ephemeralConfig+    let overridden = ["fsync", "full_page_writes", "synchronous_commit", "shared_buffers", "wal_level"]+        durable = original{Pg.postgresSettings = filter (\(key, _) -> key `notElem` overridden) original.postgresSettings ++ [("fsync", "on"), ("synchronous_commit", "on"), ("full_page_writes", "on"), ("shared_buffers", "128MB"), ("wal_level", "replica")]}+    result <- Pg.withCachedConfig durable Pg.defaultCacheConfig $ \database -> do+        migrateTestDatabase (Pg.connectionString database)+        live <- newIORef (0 :: Int)+        failures <- newIORef (0 :: Int)+        batches <- newIORef (0 :: Int)+        let observe = \case+                KirokuEventSubscriptionCaughtUp{} -> atomicModifyIORef' live (\n -> (n + 1, ()))+                KirokuEventSubscriptionDbError{} -> atomicModifyIORef' failures (\n -> (n + 1, ()))+                KirokuEventSubscriptionDelivered{} -> atomicModifyIORef' batches (\n -> (n + 1, ()))+                _ -> pure ()+            settings = (defaultConnectionSettings (Pg.connectionString database)){poolSize = 10, eventHandler = Just observe}+        withStore settings $ \store -> do+            let event = EventData Nothing (EventType "Probe") (object ["body" .= Text.replicate 512 "x"]) Nothing Nothing Nothing+            -- Identical existing-stream fixtures and explicit future-only live entry.+            forM_ [0 .. 3] $ \writer -> forM_ [0 .. workload.width - 1] $ \slot ->+                void $ append store [(stream writer slot 0 False, AnyVersion, [event])]+            delivered <- newIORef (0 :: Int)+            let handler _ = atomicModifyIORef' delivered (\n -> (n + 1, ())) >> pure Continue+                config m =+                    (defaultSubscriptionConfig (SubscriptionName "probe") (workloadTarget workload) handler)+                        { missingCheckpointPolicy = FromCurrentHead+                        , batchSize = either (error . show) Prelude.id (mkBatchSize workload.checkpointBatch)+                        , consumerGroup = if workload.mode == "group" then Just (membership m) else Nothing+                        }+                members = if workload.mode == "group" then [0, 1, 2, 3] else [0]+                native = bracket (mapM (subscribe store . config) members) (mapM_ cancel) $ \_ -> measure workload store event delivered live failures batches (length members)+            case workload.mode of+                "none" -> measure workload store event delivered live failures batches 0+                "adapter" -> runEff $ runTracingNoop $ do+                    adapter <-+                        kirokuAdapter store $+                            (defaultKirokuAdapterConfig (SubscriptionName "probe") AllStreams)+                                { Adapter.missingCheckpointPolicy = FromCurrentHead+                                , Adapter.batchSize = either (error . show) Prelude.id (mkBatchSize workload.checkpointBatch)+                                }+                    app <- runApp defaultAppConfig [(ProcessorId "probe", mkProcessor adapter (\_ -> liftIO (atomicModifyIORef' delivered (\n -> (n + 1, ()))) >> pure AckOk))]+                    case app of+                        Left err -> liftIO $ fail (show err)+                        Right handle -> do+                            liftIO $ measure workload store event delivered live failures batches 1+                            stopApp handle+                _ -> native+    either (fail . show) pure result+  where+    number text = maybe (fail ("invalid argument " <> text)) pure (readMaybe text)++membership :: Int32 -> ConsumerGroup+#ifdef LEGACY_TOPOLOGY+membership m = ConsumerGroup m 4+#else+membership m = either (error . show) Prelude.id (mkConsumerGroupSize 4 >>= mkConsumerGroup m)+#endif++workloadTarget :: Workload -> SubscriptionTarget+workloadTarget workload = if workload.mode == "category" || workload.mode == "group" then Category (CategoryName "probe") else AllStreams++stream :: Int -> Int -> Int -> Bool -> StreamName+stream writer slot iteration fresh = StreamName ("probe-" <> Text.pack (show writer <> "-" <> show slot <> if fresh then "-" <> show iteration else ""))++append :: KirokuStore -> [(StreamName, ExpectedVersion, [EventData])] -> IO [AppendResult]+append store operations = runStoreIO store (appendMultiStream operations) >>= either (fail . show) pure++measure :: Workload -> KirokuStore -> EventData -> IORef Int -> IORef Int -> IORef Int -> IORef Int -> Int -> IO ()+measure workload store event delivered live failures batches members = do+    when (members > 0) $ await "live entry" (fmap (>= members) (readIORef live))+    -- Warm up the same append shape for two seconds, then fully drain before+    -- measuring. The output contains only the declared measurement interval.+    void $ writers 2+    when (members > 0) $ await "warmup checkpoint drain" durable+    flushStats store+    writeIORef delivered 0+    writeIORef batches 0+    performMajorGC+    stats0 <- getRTSStats+    wal0 <- scalar store "SELECT pg_current_wal_insert_lsn()::text"+    tables0 <- tableStats store+    start <- getMonotonicTimeNSec+    samples <- writers workload.seconds+    finish <- getMonotonicTimeNSec+    let calls = sum [length latencies | (latencies, _) <- samples]+        events = calls * workload.width * workload.appendBatch+        elapsed = secondsBetween start finish+        sorted = sort (concat [latencies | (latencies, _) <- samples])+    when (members > 0) $ do+        await "delivery drain" (fmap (== events) (readIORef delivered))+        await "durable progress drain" durable+    performMajorGC+    stats1 <- getRTSStats+    errors <- readIORef failures+    unless (errors == 0) (fail "database errors during workload")+    flushStats store+    wal1 <- scalar store "SELECT pg_current_wal_insert_lsn()::text"+    walBytes <- scalarInt store ("SELECT pg_wal_lsn_diff('" <> wal1 <> "'::pg_lsn, '" <> wal0 <> "'::pg_lsn)::bigint")+    tables1 <- tableStats store+    count <- readIORef delivered+    batchCount <- readIORef batches+    server <- scalar store "SELECT version()"+    durability <- scalar store "SELECT current_setting('fsync') || ',' || current_setting('synchronous_commit') || ',' || current_setting('full_page_writes')"+    let (saves0, hot0) = tables0+        (saves1, hot1) = tables1+    Bytes.putStrLn $+        encode $+            object+                [ "workload" .= object ["mode" .= workload.mode, "width" .= workload.width, "append_batch" .= workload.appendBatch, "checkpoint_batch" .= workload.checkpointBatch, "fresh" .= workload.fresh, "offered" .= workload.offered, "seconds" .= workload.seconds]+                , "server" .= server+                , "durability" .= durability+                , "calls" .= calls+                , "events" .= events+                , "elapsed" .= elapsed+                , "events_per_second" .= (fromIntegral events / elapsed :: Double)+                , "append_p50_ms" .= percentile sorted 0.50+                , "append_p95_ms" .= percentile sorted 0.95+                , "append_p99_ms" .= percentile sorted 0.99+                , "delivered" .= count+                , "delivery_batches" .= batchCount+                , "durable_drained" .= True+                , "wal_bytes" .= walBytes+                , "checkpoint_updates" .= (saves1 - saves0)+                , "checkpoint_hot_updates" .= (hot1 - hot0)+                , "allocated_bytes" .= (stats1.allocated_bytes - stats0.allocated_bytes)+                , "gc_cpu_ns" .= (stats1.gc_cpu_ns - stats0.gc_cpu_ns)+                , "gc_elapsed_ns" .= (stats1.gc_elapsed_ns - stats0.gc_elapsed_ns)+                , "max_live_bytes" .= stats1.max_live_bytes+                ]+  where+    durable = do+        Right inventory <- runStoreIO store subscriptionCheckpointInventory+        -- Group members have sparse partitions, so their own final event may+        -- precede the global head; verify against the last matching fetch below.+        if workload.mode /= "group"+            then pure (all (\row -> row.checkpointPosition >= inventory.storePosition) (Vector.toList inventory.checkpoints))+            else+                and+                    <$> mapM+                        ( \row -> do+                            n <- runSession store $ Session.statement (let { GlobalPosition p = row.checkpointPosition } in p, row.consumerGroupMember) groupRemaining+                            pure (n == 0)+                        )+                        (Vector.toList inventory.checkpoints)+    writers duration = do+        begin <- getMonotonicTimeNSec+        let end = begin + fromIntegral duration * 1_000_000_000+        Async.mapConcurrently (writer begin end) [0 .. 3]+    writer begin end writerId = loop 0 []+      where+        spacing = if workload.offered == 0 then 0 else 4_000_000_000 `div` fromIntegral workload.offered+        loop i collected = do+            now <- getMonotonicTimeNSec+            let scheduled = if spacing == 0 then now else begin + fromIntegral i * spacing+            if now >= end || scheduled >= end+                then pure (collected, ())+                else do+                    when (scheduled > now) (threadDelay (fromIntegral ((scheduled - now) `div` 1000)))+                    void $ append store [(stream writerId slot (i + fromIntegral begin) workload.fresh, AnyVersion, replicate workload.appendBatch event) | slot <- [0 .. workload.width - 1]]+                    completed <- getMonotonicTimeNSec+                    loop (i + 1) (fromIntegral (completed - scheduled) / 1_000_000 : collected)++secondsBetween :: Word64 -> Word64 -> Double+secondsBetween start finish = fromIntegral (finish - start) / 1_000_000_000++percentile :: [Double] -> Double -> Double+percentile [] _ = 0+percentile values quantile = values !! min (length values - 1) (floor (quantile * fromIntegral (length values - 1)))++await :: String -> IO Bool -> IO ()+await label predicate = timeout 30_000_000 loop >>= maybe (fail (label <> " timed out")) pure+  where+    loop = predicate >>= \ok -> unless ok (threadDelay 1_000 >> loop)++runSession :: KirokuStore -> Session.Session a -> IO a+runSession store session = Pool.use store.pool session >>= either (fail . show) pure++scalar :: KirokuStore -> Text -> IO Text+scalar store sql = runSession store (Session.statement () (preparable sql E.noParams (D.singleRow (D.column (D.nonNullable D.text)))))++scalarInt :: KirokuStore -> Text -> IO Int64+scalarInt store sql = runSession store (Session.statement () (preparable sql E.noParams (D.singleRow (D.column (D.nonNullable D.int8)))))++tableStats :: KirokuStore -> IO (Int64, Int64)+tableStats store = runSession store $ Session.statement () (preparable "SELECT n_tup_upd, n_tup_hot_upd FROM pg_stat_user_tables WHERE schemaname = 'kiroku' AND relname = 'subscriptions'" E.noParams (D.singleRow ((,) <$> D.column (D.nonNullable D.int8) <*> D.column (D.nonNullable D.int8))))++groupRemaining :: Statement (Int64, Int32) Int64+groupRemaining =+    preparable+        "SELECT count(*) FROM kiroku.stream_events WHERE stream_id = 0 AND category = 'probe' AND stream_version > $1 AND (((hashtextextended(original_stream_id::text, 0) % 4) + 4) % 4) = $2"+        (contrazip2 (E.param (E.nonNullable E.int8)) (E.param (E.nonNullable E.int4)))+        (D.singleRow (D.column (D.nonNullable D.int8)))++-- Flush backend-local table statistics outside the timed interval. Ten+-- concurrent sessions occupy the ten pool slots, including idle writer and+-- checkpoint connections; otherwise an idle backend's warmup updates can leak+-- into the measured HOT-update deltas.+flushStats :: KirokuStore -> IO ()+flushStats store =+    Async.replicateConcurrently_ 10 $+        runSession store $+            Session.script "SELECT pg_stat_force_next_flush(); SELECT pg_sleep(0.2)"
shibuya-kiroku-adapter.cabal view
@@ -1,6 +1,6 @@ cabal-version:   3.0 name:            shibuya-kiroku-adapter-version:         0.5.1.5+version:         0.6.0.0 synopsis:   Kiroku event store adapter for the Shibuya queue processing framework @@ -18,6 +18,7 @@ author:          Nadeem Bitar maintainer:      nadeem@gmail.com license:         BSD-3-Clause+license-file:    LICENSE build-type:      Simple category:        Concurrency, Database, Eventing extra-doc-files: CHANGELOG.md@@ -59,19 +60,20 @@     Shibuya.Adapter.Kiroku.Convert    build-depends:-    , aeson                                  >=2.1      && <2.3-    , base                                   >=4.18     && <5-    , effectful-core                         >=2.6.1    && <2.7  || >=2.7.1.1 && <2.8+    , aeson                                  >=2.1       && <2.3+    , base                                   >=4.18      && <5+    , effectful-core                         >=2.6.1     && <2.7  || >=2.7.1.1 && <2.8     , hs-opentelemetry-api                   ^>=1.0     , hs-opentelemetry-semantic-conventions  ^>=1.40-    , kiroku-store                           ^>=0.9.0.1-    , shibuya-core                           >=0.10     && <0.11+    , kiroku-store                           ^>=0.10.0.0+    , shibuya-core                           >=0.10      && <0.11     , shibuya-kiroku-adapter-internal-    , stm                                    >=2.5      && <2.6-    , streamly-core                          >=0.3      && <0.4-    , text                                   >=2.0      && <2.2-    , unordered-containers                   >=0.2      && <0.3-    , uuid                                   >=1.3      && <1.4+    , stm                                    >=2.5       && <2.6+    , streamly-core                          >=0.3       && <0.4+    , text                                   >=2.0       && <2.2+    , time                                   >=1.12      && <1.15+    , unordered-containers                   >=0.2       && <0.3+    , uuid                                   >=1.3       && <1.4    hs-source-dirs:  src @@ -82,30 +84,30 @@   hs-source-dirs: test   ghc-options:    -threaded -rtsopts -with-rtsopts=-N   build-depends:-    , aeson                            >=2.1      && <2.3-    , base                             >=4.18     && <5-    , containers                       >=0.6      && <0.8+    , aeson                            >=2.1       && <2.3+    , base                             >=4.18      && <5+    , containers                       >=0.6       && <0.8     , directory-    , effectful                        >=2.6.1    && <2.8-    , effectful-core                   >=2.6.1    && <2.7  || >=2.7.1.1 && <2.8-    , ephemeral-pg                     >=0.3.1    && <0.4-    , generic-lens                     >=2.2      && <2.4-    , hasql                            >=1.10     && <1.11-    , hasql-pool                       >=1.2      && <1.5+    , effectful                        >=2.6.1     && <2.8+    , effectful-core                   >=2.6.1     && <2.7  || >=2.7.1.1 && <2.8+    , ephemeral-pg                     >=0.3.1     && <0.4+    , generic-lens                     >=2.2       && <2.4+    , hasql                            >=1.10      && <1.11+    , hasql-pool                       >=1.2       && <1.5     , hs-opentelemetry-api             ^>=1.0-    , hspec                            >=2.10     && <2.12-    , kiroku-store                     ^>=0.9.0.1+    , hspec                            >=2.10      && <2.12+    , kiroku-store                     ^>=0.10.0.0     , kiroku-test-support              ^>=0.1.0.0-    , lens                             >=5.2      && <5.4-    , shibuya-core                     >=0.10     && <0.11+    , lens                             >=5.2       && <5.4+    , shibuya-core                     >=0.10      && <0.11     , shibuya-kiroku-adapter     , shibuya-kiroku-adapter-internal-    , stm                              >=2.5      && <2.6-    , streamly-core                    >=0.3      && <0.4-    , text                             >=2.0      && <2.2-    , time                             >=1.12     && <1.15-    , unordered-containers             >=0.2      && <0.3-    , uuid                             >=1.3      && <1.4+    , stm                              >=2.5       && <2.6+    , streamly-core                    >=0.3       && <0.4+    , text                             >=2.0       && <2.2+    , time                             >=1.12      && <1.15+    , unordered-containers             >=0.2       && <0.3+    , uuid                             >=1.3       && <1.4  -- EP-45 live-store performance, restart, backlog, and retained-memory fixture. -- The executable owns an ephemeral PostgreSQL instance and emits the common@@ -121,18 +123,43 @@   ghc-options:        -threaded -rtsopts "-with-rtsopts=-N4 -T -A32m" -O2   default-extensions: OverloadedRecordDot   build-depends:-    , aeson                   >=2.1      && <2.3-    , async                   >=2.2      && <2.3-    , base                    >=4.18     && <5-    , bytestring              >=0.11     && <0.13-    , containers              >=0.6      && <0.8-    , effectful-core          >=2.6.1    && <2.7  || >=2.7.1.1 && <2.8-    , ephemeral-pg            >=0.3.1    && <0.4-    , kiroku-store            ^>=0.9.0.1+    , aeson                   >=2.1       && <2.3+    , async                   >=2.2       && <2.3+    , base                    >=4.18      && <5+    , bytestring              >=0.11      && <0.13+    , containers              >=0.6       && <0.8+    , effectful-core          >=2.6.1     && <2.7  || >=2.7.1.1 && <2.8+    , ephemeral-pg            >=0.3.1     && <0.4+    , kiroku-store            ^>=0.10.0.0     , kiroku-test-support     ^>=0.1.0.0-    , shibuya-core            >=0.10     && <0.11+    , shibuya-core            >=0.10      && <0.11     , shibuya-kiroku-adapter-    , stm                     >=2.5      && <2.6-    , text                    >=2.0      && <2.2-    , time                    >=1.12     && <1.15-    , vector                  >=0.13     && <0.14+    , stm                     >=2.5       && <2.6+    , text                    >=2.0       && <2.2+    , time                    >=1.12      && <1.15+    , vector                  >=0.13      && <0.14++-- MP-12: identical durable mixed-workload probe for control/candidate builds.+benchmark kiroku-mp12-write-probe+  import:             common+  type:               exitcode-stdio-1.0+  main-is:            WriteProbe.hs+  hs-source-dirs:     bench+  default-extensions: OverloadedRecordDot+  ghc-options:        -threaded -rtsopts "-with-rtsopts=-N4 -T -A32m" -O2+  build-depends:+    , aeson                   >=2.1   && <2.3+    , async                   >=2.2   && <2.3+    , base                    >=4.18  && <5+    , bytestring              >=0.11  && <0.13+    , contravariant-extras    >=0.3+    , effectful-core          >=2.6.1 && <2.7  || >=2.7.1.1 && <2.8+    , ephemeral-pg            >=0.3.1 && <0.4+    , hasql                   >=1.10  && <1.11+    , hasql-pool              >=1.2   && <1.5+    , kiroku-store+    , kiroku-test-support+    , shibuya-core            >=0.10  && <0.11+    , shibuya-kiroku-adapter+    , text                    >=2.0   && <2.2+    , vector                  >=0.13  && <0.14
src/Shibuya/Adapter/Kiroku.hs view
@@ -26,7 +26,7 @@                 pure AckOk          Right appHandle <- runApp defaultAppConfig-            [(ProcessorId \"my-projection\", mkProcessor adapter handler)]+            [(ProcessorId \"my-projection\", kirokuProcessor adapter handler)]          waitApp appHandle @@@ -45,10 +45,11 @@ main :: IO () main = withStore settings $ \\store ->     runEff $ runTracingNoop $ do+        groupSize <- either (fail . show) pure (mkConsumerGroupSize 4)         let cfg = defaultConsumerGroupConfig                 (SubscriptionName \"orders-projection\")                 (Category (CategoryName \"orders\"))-                4   -- group size+                groupSize          Right processors <- kirokuConsumerGroupProcessors store cfg handler         Right appHandle <- runApp defaultAppConfig processors@@ -105,14 +106,22 @@  Shibuya's supervised runner converts a synchronous handler exception to an immediate 'AckRetry' and finalizes it, so the ack-coupled Kiroku worker cannot be-left blocked by an abandoned reply. 'guardKirokuHandlerWith' remains useful when+left blocked by an abandoned reply. 'kirokuProcessor' applies the one-second+paced guard for a single processor. 'guardKirokuHandlerWith' remains useful when the application wants a different exception disposition, and 'kirokuConsumerGroupProcessors' applies the adapter's default guard automatically. Asynchronous cancellation is never converted into an ack.++Raw consumers of @adapter.source@ must finalize every item. Leaving one pending+blocks delivery and checkpoint advancement. Opt in with a positive+@handlerStallWarnAfter@ to receive advisory store handler-stall events; warnings+never finalize, retry or checkpoint the item. @retryPolicy@ controls total+deliveries independently of the delay chosen by 'AckRetry'. -} module Shibuya.Adapter.Kiroku (     -- * Adapter     kirokuAdapter,+    kirokuProcessor,     guardKirokuHandlerWith,     guardKirokuHandler, @@ -130,7 +139,13 @@     -- * Re-exports from kiroku-store     SubscriptionName (..),     SubscriptionTarget (..),-    ConsumerGroup (..),+    ConsumerGroup,+    ConsumerGroupSize,+    consumerGroupSizeValue,+    mkConsumerGroupSize,+    mkConsumerGroup,+    member,+    size,     EventTypeFilter (..),     MissingCheckpointPolicy (..), ) where@@ -138,21 +153,27 @@ import Control.Exception (SomeException) import Data.Int (Int32) import Data.Text qualified as T+import Data.Time (NominalDiffTime) import Effectful (Eff, IOE, liftIO, (:>))-import Effectful.Exception (catchSync, throwIO)+import Effectful.Exception (catchSync) import GHC.Generics (Generic) import Kiroku.Store.Connection (KirokuStore)-import Kiroku.Store.Subscription.Stream (subscriptionAckStream)+import Kiroku.Store.Subscription.Stream (StreamBufferSize, defaultStreamBufferSize, subscriptionAckStream) import Kiroku.Store.Subscription.Types (-    ConsumerGroup (..),+    ConsumerGroup,+    ConsumerGroupSize,     EventTypeFilter (..),-    InvalidConsumerGroup (..),     MissingCheckpointPolicy (..),     SubscriptionConfig,     SubscriptionName (..),     SubscriptionResult (..),     SubscriptionTarget (..),+    consumerGroupSizeValue,     defaultSubscriptionConfig,+    member,+    mkConsumerGroup,+    mkConsumerGroupSize,+    size,  ) import Kiroku.Store.Subscription.Types qualified as Sub import Kiroku.Store.Types (RecordedEvent)@@ -160,7 +181,7 @@ import Shibuya.Adapter (Adapter (..)) import Shibuya.Adapter.Kiroku.Convert (kirokuEnvelopeAttrs, toIngestedAck) import Shibuya.Adapter.Kiroku.Internal (acquireAllAndTransfer)-import Shibuya.App (ProcessorId (..), QueueProcessor (..))+import Shibuya.App (ProcessorId (..), QueueProcessor (..), mkProcessor) import Shibuya.Core.Ack (AckDecision (..), RetryDelay (..)) import Shibuya.Core.Error (PolicyError (..)) import Shibuya.Handler (Handler)@@ -185,10 +206,14 @@     -- ^ Unique subscription identifier (checkpoint key)     , subscriptionTarget :: !SubscriptionTarget     -- ^ 'AllStreams' or @'Category' categoryName@-    , batchSize :: !Int32+    , batchSize :: !Sub.BatchSize     -- ^ Events per database fetch during catch-up-    , bufferSize :: !Natural+    , bufferSize :: !StreamBufferSize     -- ^ Bridge 'TBQueue' capacity; must be at least 1.+    , retryPolicy :: !Sub.RetryPolicy+    -- ^ Total delivery attempts, default five; each AckRetry chooses its delay.+    , handlerStallWarnAfter :: !(Maybe NominalDiffTime)+    -- ^ Optional positive advisory warning interval, disabled by default.     , queueCapacity :: !Natural     {- ^ Publisher-side capacity in batches, where each batch contains up to     Kiroku's publisher batch size (currently 1000 events). When this fills,@@ -197,15 +222,14 @@     , consumerGroup :: !(Maybe ConsumerGroup)     {- ^ Optional consumer-group membership for this adapter instance.     'Nothing' (the default) = ordinary single-consumer subscription.-    @'Just' ('ConsumerGroup' { member = m, size = n })@ = this adapter is+    @'Just' cg@ (built with 'mkConsumerGroup') = this adapter is     member @m@ of a group of size @n@, receiving only the events whose     originating stream hashes to slot @m@ (in global-position order). To run a     full size-@n@ group, create @n@ adapters with the same 'subscriptionName'     and distinct 'member' indices, each backed by its own Shibuya processor. -    The validity invariant (@size >= 1@, @0 <= member < size@) is enforced by-    the underlying 'Kiroku.Store.Subscription.subscribe' call, which throws-    'Kiroku.Store.Subscription.Types.InvalidConsumerGroup' on violation.+    'mkConsumerGroupSize' and 'mkConsumerGroup' validate this invariant before+    an adapter can be configured.     -}     , missingCheckpointPolicy :: !MissingCheckpointPolicy     {- ^ What the underlying Kiroku worker does when this adapter's exact@@ -238,8 +262,8 @@     }     deriving stock (Generic) -{- | A 'KirokuAdapterConfig' with sensible defaults: @batchSize = 100@,-@bufferSize = 256@, @queueCapacity = 16@, @consumerGroup = 'Nothing'@+{- | A 'KirokuAdapterConfig' with sensible defaults: @batchSize = Sub.defaultBatchSize@,+@bufferSize = defaultStreamBufferSize@, @queueCapacity = 16@, @consumerGroup = 'Nothing'@ (ordinary single-consumer subscription), @missingCheckpointPolicy = 'FromBeginning'@, @eventTypeFilter = 'AllEventTypes'@ (deliver every type), and @selector = 'Nothing'@ (no extra predicate@@ -262,8 +286,10 @@     KirokuAdapterConfig         { subscriptionName = name         , subscriptionTarget = target-        , batchSize = 100-        , bufferSize = 256+        , batchSize = Sub.defaultBatchSize+        , bufferSize = defaultStreamBufferSize+        , retryPolicy = Sub.defaultRetryPolicy+        , handlerStallWarnAfter = Nothing         , queueCapacity = 16         , consumerGroup = Nothing         , missingCheckpointPolicy = FromBeginning@@ -293,6 +319,12 @@ guardKirokuHandler :: Handler es msg -> Handler es msg guardKirokuHandler = guardKirokuHandlerWith (const (AckRetry (RetryDelay 1))) +{- | Recommended single-processor constructor. Like mkProcessor, defaults to+unordered policy and serial concurrency, with the paced exception guard.+-}+kirokuProcessor :: Adapter es RecordedEvent -> Handler es RecordedEvent -> QueueProcessor es+kirokuProcessor adapter handler = mkProcessor adapter (guardKirokuHandler handler)+ {- | Create a Shibuya 'Adapter' backed by a Kiroku subscription.  The adapter:@@ -315,7 +347,7 @@     KirokuStore ->     KirokuAdapterConfig ->     Eff es (Adapter es RecordedEvent)-kirokuAdapter store KirokuAdapterConfig{subscriptionName = subName, subscriptionTarget = subTarget, batchSize = bs, bufferSize = buf, queueCapacity = qCap, consumerGroup = cg, missingCheckpointPolicy = checkpointPolicy, eventTypeFilter = etf, selector = sel} = do+kirokuAdapter store KirokuAdapterConfig{subscriptionName = subName, subscriptionTarget = subTarget, batchSize = bs, bufferSize = buf, queueCapacity = qCap, retryPolicy = attempts, handlerStallWarnAfter = stallWarn, consumerGroup = cg, missingCheckpointPolicy = checkpointPolicy, eventTypeFilter = etf, selector = sel} = do     -- Build from 'defaultSubscriptionConfig' and override only the non-default     -- fields. Using the smart constructor (rather than a full record literal)     -- means any future field added to 'SubscriptionConfigM' is inherited at its@@ -324,6 +356,8 @@         subConfig =             (defaultSubscriptionConfig subName subTarget (\_ -> pure Continue))                 { Sub.batchSize = bs+                , Sub.retryPolicy = attempts+                , Sub.handlerStallWarnAfter = stallWarn                 , Sub.queueCapacity = qCap                 , Sub.consumerGroup = cg                 , Sub.missingCheckpointPolicy = checkpointPolicy@@ -343,7 +377,7 @@         envAttrs =             kirokuEnvelopeAttrs                 subNameText-                (fmap (\ConsumerGroup{member = m} -> fromIntegral m) cg)+                (fmap (fromIntegral . member) cg)         ingestedStream = fmap (toIngestedAck envAttrs cancelAction) (Stream.morphInner liftIO ioStream)      pure@@ -378,15 +412,18 @@     {- ^ 'AllStreams' or @'Category' categoryName@ — the same source for every     member; kiroku partitions it across members in SQL.     -}-    , groupSize :: !Int32+    , groupSize :: !ConsumerGroupSize     {- ^ @N@ members; must be @>= 1@ (enforced by the underlying-    'Kiroku.Store.Subscription.subscribe', which throws-    'Kiroku.Store.Subscription.Types.InvalidConsumerGroup' otherwise).+    'mkConsumerGroupSize' at construction).     -}-    , batchSize :: !Int32+    , batchSize :: !Sub.BatchSize     -- ^ Events per database fetch during catch-up (per member).-    , bufferSize :: !Natural+    , bufferSize :: !StreamBufferSize     -- ^ Per-member bridge 'TBQueue' capacity; must be at least 1.+    , retryPolicy :: !Sub.RetryPolicy+    -- ^ Total delivery attempts, default five; each AckRetry chooses its delay.+    , handlerStallWarnAfter :: !(Maybe NominalDiffTime)+    -- ^ Optional positive advisory warning interval, disabled by default.     , queueCapacity :: !Natural     {- ^ Per-member publisher-side capacity in batches. When this fills, Kiroku     pauses and later resumes the member losslessly.@@ -421,21 +458,23 @@     deriving stock (Generic)  {- | A 'KirokuConsumerGroupConfig' with sensible defaults: @memberConcurrency =-'Serial'@ (the only legal per-member concurrency), @batchSize = 100@,-@bufferSize = 256@, @queueCapacity = 16@, @missingCheckpointPolicy =+'Serial'@ (the only legal per-member concurrency), @batchSize = Sub.defaultBatchSize@,+@bufferSize = defaultStreamBufferSize@, @queueCapacity = 16@, @missingCheckpointPolicy = 'FromBeginning'@, @eventTypeFilter = 'AllEventTypes'@ (deliver every type), @selector = 'Nothing'@ (no extra predicate filtering). Supply the subscription name, target, and group size. -} defaultConsumerGroupConfig ::-    SubscriptionName -> SubscriptionTarget -> Int32 -> KirokuConsumerGroupConfig+    SubscriptionName -> SubscriptionTarget -> ConsumerGroupSize -> KirokuConsumerGroupConfig defaultConsumerGroupConfig name target n =     KirokuConsumerGroupConfig         { subscriptionName = name         , subscriptionTarget = target         , groupSize = n-        , batchSize = 100-        , bufferSize = 256+        , batchSize = Sub.defaultBatchSize+        , bufferSize = defaultStreamBufferSize+        , retryPolicy = Sub.defaultRetryPolicy+        , handlerStallWarnAfter = Nothing         , queueCapacity = 16         , memberConcurrency = Serial         , missingCheckpointPolicy = FromBeginning@@ -476,9 +515,7 @@  Each processor's 'ProcessorId' is @\"\<subscriptionName\>-member-\<m\>\"@ so member identity is readable off the id-and two members never collide. @groupSize >= 1@ is validated before any-subscription opens, throwing 'InvalidConsumerGroup' with member 0 when it is-violated.+and two members never collide. The group size is validated by 'mkConsumerGroupSize' at construction. -} kirokuConsumerGroupProcessors ::     (IOE :> es) =>@@ -486,7 +523,7 @@     KirokuConsumerGroupConfig ->     Handler es RecordedEvent ->     Eff es (Either PolicyError [(ProcessorId, QueueProcessor es)])-kirokuConsumerGroupProcessors store cfg@KirokuConsumerGroupConfig{subscriptionName = subName, subscriptionTarget = subTarget, groupSize = n, batchSize = bs, bufferSize = buf, queueCapacity = qCap, missingCheckpointPolicy = checkpointPolicy, eventTypeFilter = etf, selector = sel} handler =+kirokuConsumerGroupProcessors store cfg@KirokuConsumerGroupConfig{subscriptionName = subName, subscriptionTarget = subTarget, groupSize = n, batchSize = bs, bufferSize = buf, queueCapacity = qCap, retryPolicy = attempts, handlerStallWarnAfter = stallWarn, missingCheckpointPolicy = checkpointPolicy, eventTypeFilter = etf, selector = sel} handler =     kirokuConsumerGroupProcessorsWith mkMemberAdapter cfg handler   where     mkMemberAdapter m =@@ -497,8 +534,10 @@                 , subscriptionTarget = subTarget                 , batchSize = bs                 , bufferSize = buf+                , retryPolicy = attempts+                , handlerStallWarnAfter = stallWarn                 , queueCapacity = qCap-                , consumerGroup = Just (ConsumerGroup{member = m, size = n})+                , consumerGroup = Just (either (error . show) Prelude.id (mkConsumerGroup m n))                 , missingCheckpointPolicy = checkpointPolicy                 , eventTypeFilter = etf                 , selector = sel@@ -525,21 +564,19 @@         , memberConcurrency = mc         }     handler =-        if n < 1-            then throwIO (InvalidConsumerGroup 0 n)-            else case consumerGroupPolicy mc of-                Left e -> pure (Left e)-                Right (ordering, conc) -> do-                    let SubscriptionName name = subName-                    acquireAllAndTransfer-                        n-                        mkMemberAdapter-                        (\Adapter{shutdown = shutdownAction} -> shutdownAction)-                        (\_ _ -> pure ())-                        ( \adapters ->-                            pure . Right $-                                [ let pid = ProcessorId (name <> "-member-" <> T.pack (show m))-                                   in (pid, QueueProcessor adapter (guardKirokuHandler handler) ordering conc)-                                | (m, adapter) <- zip [0 .. n - 1] adapters-                                ]-                        )+        case consumerGroupPolicy mc of+            Left e -> pure (Left e)+            Right (ordering, conc) -> do+                let SubscriptionName name = subName+                acquireAllAndTransfer+                    (consumerGroupSizeValue n)+                    mkMemberAdapter+                    (\Adapter{shutdown = shutdownAction} -> shutdownAction)+                    (\_ _ -> pure ())+                    ( \adapters ->+                        pure . Right $+                            [ let pid = ProcessorId (name <> "-member-" <> T.pack (show m))+                               in (pid, QueueProcessor adapter (guardKirokuHandler handler) ordering conc)+                            | (m, adapter) <- zip [0 .. consumerGroupSizeValue n - 1] adapters+                            ]+                    )
test/Main.hs view
@@ -6,23 +6,27 @@ import Control.Concurrent.STM qualified as STM import Control.Exception qualified as E import Control.Lens ((&), (.~), (^.))-import Control.Monad (unless)+import Control.Monad (forM_, unless) import Control.Monad.IO.Class (liftIO) import Data.Aeson (Value) import Data.Aeson qualified as Aeson import Data.Foldable (toList) import Data.Generics.Labels () import Data.HashMap.Strict qualified as HashMap-import Data.IORef (modifyIORef', newIORef, readIORef)+import Data.IORef (atomicModifyIORef', modifyIORef', newIORef, readIORef) import Data.Int (Int32, Int64) import Data.List (nub, sort) import Data.Map.Strict qualified as Map+import Data.Maybe (isNothing) import Data.Set qualified as Set import Data.Text (Text) import Data.Text qualified as T import Data.Time (UTCTime (..), fromGregorian) import Data.UUID qualified as UUID+import Data.Word (Word64) import Effectful (runEff)+import Effectful.Timeout qualified as Timeout+import GHC.Clock (getMonotonicTimeNSec) import Hasql.Pool qualified as Pool import Hasql.Session qualified as Session import Kiroku.Store@@ -41,6 +45,7 @@     kirokuAdapter,     kirokuConsumerGroupProcessors,     kirokuConsumerGroupProcessorsWith,+    kirokuProcessor,  ) import Shibuya.Adapter.Kiroku.Convert (     KirokuEnvelopeAttrs,@@ -189,6 +194,7 @@             consumerGroupPolicy (Async 4)                 `shouldBe` Left (InvalidPolicyCombo "StrictInOrder requires Serial concurrency") +    acknowledgementSpec     around withTestStore $ do         describe "kirokuAdapter" $ do             it "delivers only matching event types when an eventTypeFilter is set (EP-43)" $ \store -> do@@ -246,7 +252,7 @@                 countVar <- newTVarIO (0 :: Int)                 runEff $ runTracingNoop $ do                     let cfg =-                            defaultConsumerGroupConfig (SubscriptionName "etfg") (Category (CategoryName "etfg")) 4+                            defaultConsumerGroupConfig (SubscriptionName "etfg") (Category (CategoryName "etfg")) (validGroupSize 4)                                 & #eventTypeFilter .~ OnlyEventTypes (Set.fromList [EventType "A"])                         handler ingested = do                             liftIO $ do@@ -1027,7 +1033,7 @@                             ( \m ->                                 kirokuAdapter store $                                     defaultKirokuAdapterConfig (SubscriptionName "cg-shibuya-group") (Category (CategoryName "cg"))-                                        & #consumerGroup .~ Just (ConsumerGroup{member = m, size = 4})+                                        & #consumerGroup .~ Just (validConsumerGroup m 4)                             )                             [0, 1, 2, 3] @@ -1070,7 +1076,7 @@                         runTracingNoop $                             kirokuConsumerGroupProcessors                                 store-                                ( defaultConsumerGroupConfig (SubscriptionName "cgp-reject") (Category (CategoryName "cgp-reject")) 4+                                ( defaultConsumerGroupConfig (SubscriptionName "cgp-reject") (Category (CategoryName "cgp-reject")) (validGroupSize 4)                                     & #memberConcurrency .~ Async 4                                 )                                 ( \ingested -> do@@ -1089,7 +1095,7 @@                                     liftIO $ expectationFailure "factory should not be called for an invalid policy"                                     pure stubAdapter                                 )-                                ( defaultConsumerGroupConfig (SubscriptionName "cgp-reject-with") (Category (CategoryName "cgp-reject-with")) 4+                                ( defaultConsumerGroupConfig (SubscriptionName "cgp-reject-with") (Category (CategoryName "cgp-reject-with")) (validGroupSize 4)                                     & #memberConcurrency .~ Async 4                                 )                                 ( \ingested -> do@@ -1100,23 +1106,9 @@                     Left err -> err `shouldBe` InvalidPolicyCombo "StrictInOrder requires Serial concurrency"                     Right _ -> expectationFailure "expected Left PolicyError for Async member concurrency" -            it "throws InvalidConsumerGroup for non-positive group sizes" $ \store -> do-                let handler ingested = do-                        let _ = envelopePayload ingested-                        pure AckOk-                    throwsSize n =-                        runEff-                            ( runTracingNoop $-                                kirokuConsumerGroupProcessors-                                    store-                                    (defaultConsumerGroupConfig (SubscriptionName ("cgp-invalid-" <> T.pack (show n))) AllStreams n)-                                    handler-                            )-                            `shouldThrow` ( \InvalidConsumerGroup{invalidMember = member, invalidSize = size} ->-                                                member == 0 && size == n-                                          )-                throwsSize 0-                throwsSize (-1)+            it "rejects non-positive group sizes at construction" $ \_store -> do+                mkConsumerGroupSize 0 `shouldBe` Left (InvalidConsumerGroup 0 0)+                mkConsumerGroupSize (-1) `shouldBe` Left (InvalidConsumerGroup 0 (-1))              it "shuts down every created member and preserves the factory failure when cleanup throws" $ \_store -> do                 shutdowns <- newIORef ([] :: [Int32])@@ -1143,7 +1135,7 @@                             runTracingNoop $                                 kirokuConsumerGroupProcessorsWith                                     factory-                                    (defaultConsumerGroupConfig (SubscriptionName "cgp-partial-cleanup") AllStreams 3)+                                    (defaultConsumerGroupConfig (SubscriptionName "cgp-partial-cleanup") AllStreams (validGroupSize 3))                                     handler                 case result of                     Left (e :: E.SomeException) -> show e `shouldContain` "member 2 failed"@@ -1195,7 +1187,7 @@                         result <-                             kirokuConsumerGroupProcessors                                 store-                                (defaultConsumerGroupConfig (SubscriptionName "cgp-guard") AllStreams 1)+                                (defaultConsumerGroupConfig (SubscriptionName "cgp-guard") AllStreams (validGroupSize 1))                                 ( \_ -> do                                     n <- liftIO $ atomically $ do                                         c <- readTVar countVar@@ -1237,7 +1229,7 @@                         defaultConsumerGroupConfig                             (SubscriptionName "cgp-shibuya-group")                             (Category (CategoryName "cgp"))-                            4+                            (validGroupSize 4)                  runEff $ runTracingNoop $ do                     let handler ingested = do@@ -1305,8 +1297,11 @@  -- | Read the dead letters recorded for a non-group subscription (member 0). readDeadLetters :: KirokuStore -> Text -> IO [SQL.DeadLetterRecord]-readDeadLetters store subName = do-    result <- Pool.use (store ^. #pool) (Session.statement (subName, 0 :: Int32) SQL.readDeadLettersStmt)+readDeadLetters store subName = readMemberDeadLetters store subName 0++readMemberDeadLetters :: KirokuStore -> Text -> Int32 -> IO [SQL.DeadLetterRecord]+readMemberDeadLetters store subName memberIndex = do+    result <- Pool.use (store ^. #pool) (Session.statement (subName, memberIndex) SQL.readDeadLettersStmt)     case result of         Left err -> error ("readDeadLetters failed: " <> show err)         Right v -> pure (toList v)@@ -1428,3 +1423,161 @@ withTestStore action =     withMigratedTestDatabase $ \connStr ->         withStore (defaultConnectionSettings connStr) action++validGroupSize :: Int32 -> ConsumerGroupSize+validGroupSize = either (error . show) Prelude.id . mkConsumerGroupSize++validConsumerGroup :: Int32 -> Int32 -> ConsumerGroup+validConsumerGroup m n = either (error . show) Prelude.id (mkConsumerGroup m (validGroupSize n))++acknowledgementSpec :: Spec+acknowledgementSpec = do+    describe "acknowledgement liveness" $ do+        forM_ [False, True] $ \guarded ->+            it (if guarded then "guarded single processor retries with pacing and continues" else "standard Shibuya handler exceptions are finalized") $+                withTestStore $ \store -> do+                    Right _ <-+                        runStoreIO store $+                            appendToStream+                                (StreamName "ack-liveness")+                                NoStream+                                [makeEvent "First" (Aeson.object []), makeEvent "Second" (Aeson.object [])]+                    attempts <- newIORef ([] :: [(Int64, Maybe Attempt, Word64)])+                    count <- newTVarIO (0 :: Int)+                    within "runner finalization" $ runEff $ runTracingNoop $ do+                        adapter <- kirokuAdapter store (defaultKirokuAdapterConfig (SubscriptionName "ack-liveness") AllStreams)+                        let handler ingested = do+                                now <- liftIO getMonotonicTimeNSec+                                let pos = globalPos (envelopePayload ingested)+                                    Message{envelope = Envelope{attempt = attempt'}} = ingested+                                n <- liftIO $ do+                                    modifyIORef' attempts ((pos, attempt', now) :)+                                    atomically $ do+                                        old <- readTVar count+                                        writeTVar count (old + 1)+                                        pure old+                                if n == 0 then liftIO (E.throwIO (userError "one-shot handler failure")) else pure AckOk+                            processor = if guarded then kirokuProcessor adapter handler else mkProcessor adapter handler+                        app <- runApp defaultAppConfig [(ProcessorId "ack-liveness", processor)] >>= either (liftIO . fail . show) pure+                        liftIO $ waitForCount count 3 5_000_000+                        liftIO $ waitForCheckpointPosition store (SubscriptionCheckpointKey (SubscriptionName "ack-liveness") 0) (GlobalPosition 2)+                        stopApp app+                    recorded <- reverse <$> readIORef attempts+                    map (\(pos, attempt', _) -> (pos, attempt')) recorded+                        `shouldBe` [(1, Just (Attempt 0)), (1, Just (Attempt 1)), (2, Just (Attempt 0))]+                    case recorded of+                        (_, _, firstAt) : (_, _, retriedAt) : _ ->+                            if guarded+                                then retriedAt - firstAt `shouldSatisfy` (>= 900_000_000)+                                else retriedAt - firstAt `shouldSatisfy` (< 900_000_000)+                        _ -> expectationFailure "expected three finalized deliveries"++        forM_ [Nothing, Just 0.05] $ \warnAfter ->+            it ("keeps a raw unfinalized item pending and forwards warning policy " <> show warnAfter) $ do+                warnings <- newTVarIO (0 :: Int)+                warning <- STM.newEmptyTMVarIO+                let observe event = case event of+                        KirokuEventSubscriptionHandlerStalled _ pos _ elapsed _ -> atomically $ do+                            STM.modifyTVar' warnings (+ 1)+                            _ <- STM.tryPutTMVar warning (pos, elapsed)+                            pure ()+                        _ -> pure ()+                withMigratedTestDatabase $ \connStr ->+                    withStore (defaultConnectionSettings connStr & #eventHandler .~ Just observe) $ \store -> do+                        Right _ <-+                            runStoreIO store $+                                appendToStream+                                    (StreamName "ack-raw")+                                    NoStream+                                    [makeEvent "First" (Aeson.object []), makeEvent "Second" (Aeson.object [])]+                        within "raw acknowledgement ownership" $ runEff $ runTracingNoop $ Timeout.runTimeout $ do+                            adapter <-+                                kirokuAdapter+                                    store+                                    ( (defaultKirokuAdapterConfig (SubscriptionName "ack-raw") AllStreams)+                                        & #handlerStallWarnAfter .~ warnAfter+                                    )+                            let Adapter{source = sourceStream, shutdown = shutdownAction} = adapter+                            first <- Stream.uncons sourceStream+                            case first of+                                Nothing -> liftIO (expectationFailure "expected first item")+                                Just (item, rest) -> do+                                    -- The raw consumer intentionally leaves item unfinalized.+                                    secondRead <- Timeout.timeout 150_000 (Stream.uncons rest)+                                    liftIO (isNothing secondRead `shouldBe` True)+                                    liftIO $ readCheckpointPosition store (SubscriptionCheckpointKey (SubscriptionName "ack-raw") 0) `shouldReturn` Just (GlobalPosition 0)+                                    case warnAfter of+                                        Nothing -> liftIO $ STM.readTVarIO warnings `shouldReturn` 0+                                        Just threshold -> liftIO $ do+                                            (pos, elapsed) <- within "raw ack warning" (atomically (STM.readTMVar warning))+                                            pos `shouldBe` GlobalPosition 1+                                            elapsed `shouldSatisfy` (>= threshold)+                                    -- Finalization remains explicitly owned by this consumer.+                                    let Ingested{ack = AckHandle{finalize = finalizeFirst}} = item+                                    finalizeFirst AckOk+                                    (second, _) <- Stream.uncons rest >>= maybe (liftIO (fail "expected second item")) pure+                                    let Ingested{envelope = Envelope{payload = secondEvent}} = second+                                    liftIO $ globalPos secondEvent `shouldBe` 2+                                    let Ingested{ack = AckHandle{finalize = finalizeSecond}} = second+                                    finalizeSecond AckOk+                                    liftIO $ waitForCheckpointPosition store (SubscriptionCheckpointKey (SubscriptionName "ack-raw") 0) (GlobalPosition 2)+                            shutdownAction+                        stoppedCount <- STM.readTVarIO warnings+                        threadDelay 100_000+                        STM.readTVarIO warnings `shouldReturn` stoppedCount++    describe "retry policy" $ do+        forM_ [Nothing, Just 2] $ \attemptLimit ->+            it ("honors single-adapter total delivery limit " <> show attemptLimit) $ withTestStore $ \store -> do+                Right _ <- runStoreIO store $ appendToStream (StreamName "ack-policy") NoStream [makeEvent "Retry" (Aeson.object [])]+                count <- newTVarIO (0 :: Int)+                runEff $ runTracingNoop $ do+                    let config =+                            (defaultKirokuAdapterConfig (SubscriptionName "ack-policy") AllStreams)+                                & #retryPolicy .~ maybe defaultRetryPolicy RetryPolicy attemptLimit+                    adapter <- kirokuAdapter store config+                    let handler _ = liftIO (atomically (STM.modifyTVar' count (+ 1))) >> pure (AckRetry (Ack.RetryDelay 0))+                    app <- runApp defaultAppConfig [(ProcessorId "ack-policy", kirokuProcessor adapter handler)] >>= either (liftIO . fail . show) pure+                    liftIO $ waitForCheckpointPosition store (SubscriptionCheckpointKey (SubscriptionName "ack-policy") 0) (GlobalPosition 1)+                    stopApp app+                STM.readTVarIO count `shouldReturn` maybe 5 Prelude.id attemptLimit+                letters <- readDeadLetters store "ack-policy"+                map SQL.deadLetterAttemptCount letters `shouldBe` [fromIntegral (maybe 5 Prelude.id attemptLimit)]++        it "forwards two attempts and stall policy to both consumer-group members" $ do+            warned <- newTVarIO Set.empty+            release <- STM.newEmptyTMVarIO+            let observe event = case event of+                    KirokuEventSubscriptionHandlerStalled _ _ _ _ (GroupMember member' 2) -> atomically (STM.modifyTVar' warned (Set.insert member'))+                    _ -> pure ()+            withMigratedTestDatabase $ \connStr ->+                withStore (defaultConnectionSettings connStr & #eventHandler .~ Just observe) $ \store -> do+                    -- Enough distinct streams to exercise both PostgreSQL hash slots;+                    -- require both warnings before release so an empty slot fails.+                    forM_ [1 .. 20 :: Int] $ \i -> do+                        Right _ <- runStoreIO store $ appendToStream (StreamName ("ack-group-" <> T.pack (show i))) NoStream [makeEvent "Retry" (Aeson.object [])]+                        pure ()+                    seen <- newIORef Map.empty+                    runEff $ runTracingNoop $ do+                        let config =+                                (defaultConsumerGroupConfig (SubscriptionName "ack-policy-group") AllStreams (validGroupSize 2))+                                    & #retryPolicy .~ RetryPolicy 2+                                    & #handlerStallWarnAfter .~ Just 0.05+                            handler ingested = do+                                liftIO $ atomically (STM.readTMVar release)+                                liftIO $ atomicModifyIORef' seen (\old -> (Map.insertWith (+) (globalPos (envelopePayload ingested)) (1 :: Int) old, ()))+                                pure (AckRetry (Ack.RetryDelay 0))+                        processors <- kirokuConsumerGroupProcessors store config handler >>= either (liftIO . fail . show) pure+                        app <- runApp defaultAppConfig processors >>= either (liftIO . fail . show) pure+                        liftIO $ within "both group members stalled" $ atomically $ readTVar warned >>= STM.check . (== Set.fromList [0, 1])+                        liftIO $ atomically (STM.putTMVar release ())+                        liftIO $ within "all group dead letters" $ do+                            let loop = do+                                    letters <- concat <$> mapM (readMemberDeadLetters store "ack-policy-group") [0, 1]+                                    if length letters == 20 then pure () else threadDelay 10_000 >> loop+                            loop+                        stopApp app+                    Map.elems <$> readIORef seen `shouldReturn` replicate 20 2+                    letters <- concat <$> mapM (readMemberDeadLetters store "ack-policy-group") [0, 1]+                    sort (map SQL.deadLetterGlobalPosition letters) `shouldBe` [1 .. 20]+                    map SQL.deadLetterAttemptCount letters `shouldBe` replicate 20 2