packages feed

kiroku-store 0.10.0.0 → 0.11.0.0

raw patch · 40 files changed

+2046/−24 lines, 40 filesPVP ok

version bump matches the API change (PVP)

API changes (from Hackage documentation)

+ Kiroku.Store.Effect: [GetEvent] :: forall (a :: Type -> Type). EventId -> Store a (Maybe RecordedEvent)
+ Kiroku.Store.Effect: [GetStreamWithHead] :: forall (a :: Type -> Type). StreamName -> Store a (Maybe (StreamInfo, Maybe GlobalPosition))
+ Kiroku.Store.Effect: [ListCategories] :: forall (a :: Type -> Type). Maybe CategoryName -> BrowsePageSize -> Store a (Vector CategoryName)
+ Kiroku.Store.Effect: [ListStreams] :: forall (a :: Type -> Type). Maybe CategoryName -> Maybe Text -> Maybe StreamName -> BrowsePageSize -> Store a (Vector StreamInfo)
+ Kiroku.Store.Effect: [ListSubscriptionDeadLetters] :: forall (a :: Type -> Type). SubscriptionDeadLetterQuery -> Store a SubscriptionDeadLetterPage
+ Kiroku.Store.Read: getEvent :: forall (es :: [Effect]). (HasCallStack, Store :> es) => EventId -> Eff es (Maybe RecordedEvent)
+ Kiroku.Store.Read: getStreamWithHead :: forall (es :: [Effect]). (HasCallStack, Store :> es) => StreamName -> Eff es (Maybe (StreamInfo, Maybe GlobalPosition))
+ Kiroku.Store.Read: listCategories :: forall (es :: [Effect]). (HasCallStack, Store :> es) => Maybe CategoryName -> BrowsePageSize -> Eff es (Vector CategoryName)
+ Kiroku.Store.Read: listStreams :: forall (es :: [Effect]). (HasCallStack, Store :> es) => Maybe CategoryName -> Maybe Text -> Maybe StreamName -> BrowsePageSize -> Eff es (Vector StreamInfo)
+ Kiroku.Store.SQL: ExactName :: StreamRangeBound
+ Kiroku.Store.SQL: ExclusiveLower :: StreamRangeBound
+ Kiroku.Store.SQL: InclusiveLower :: StreamRangeBound
+ Kiroku.Store.SQL: data StreamRangeBound
+ Kiroku.Store.SQL: getEventStmt :: Statement UUID (Maybe RecordedEvent)
+ Kiroku.Store.SQL: getStreamWithHeadStmt :: Statement Text (Maybe (StreamInfo, Maybe GlobalPosition))
+ Kiroku.Store.SQL: instance GHC.Classes.Eq Kiroku.Store.SQL.StreamRangeBound
+ Kiroku.Store.SQL: instance GHC.Internal.Show.Show Kiroku.Store.SQL.StreamRangeBound
+ Kiroku.Store.SQL: listCategoriesStmt :: Bool -> Statement (Maybe Text, Int32) (Vector Text)
+ Kiroku.Store.SQL: listStreamsPairStmt :: StreamRangeBound -> StreamRangeBound -> Statement (Text, Maybe Text, Text, Maybe Text, Int32) (Vector StreamInfo)
+ Kiroku.Store.SQL: listStreamsRangeStmt :: StreamRangeBound -> Statement (Text, Maybe Text, Int32) (Vector StreamInfo)
+ Kiroku.Store.SQL: listStreamsSession :: Maybe CategoryName -> Maybe Text -> Maybe StreamName -> BrowsePageSize -> Session (Vector StreamInfo)
+ Kiroku.Store.SQL: listSubscriptionDeadLettersFromStartStmt :: Statement (Text, Int32) (Vector SubscriptionDeadLetter)
+ Kiroku.Store.SQL: listSubscriptionDeadLettersStmt :: Statement (Text, Int64, Int64, Int32) (Vector SubscriptionDeadLetter)
+ Kiroku.Store.SQL: listSubscriptionMemberDeadLettersFromStartStmt :: Statement (Text, Int32, Int32) (Vector SubscriptionDeadLetter)
+ Kiroku.Store.SQL: listSubscriptionMemberDeadLettersStmt :: Statement (Text, Int32, Int64, Int64, Int32) (Vector SubscriptionDeadLetter)
+ Kiroku.Store.Subscription: subscriptionDeadLetters :: forall (es :: [Effect]). (HasCallStack, Store :> es) => SubscriptionDeadLetterQuery -> Eff es SubscriptionDeadLetterPage
+ Kiroku.Store.Subscription.EventPublisher: PublisherSubscription :: !TBQueue DecodedBatch -> !TVar SubscriberStatus -> !TVar Word64 -> !IO () -> PublisherSubscription
+ Kiroku.Store.Subscription.EventPublisher: [subDropped] :: Subscriber -> !TVar Word64
+ Kiroku.Store.Subscription.EventPublisher: [subscriptionDropped] :: PublisherSubscription -> !TVar Word64
+ Kiroku.Store.Subscription.EventPublisher: [subscriptionQueue] :: PublisherSubscription -> !TBQueue DecodedBatch
+ Kiroku.Store.Subscription.EventPublisher: [subscriptionStatus] :: PublisherSubscription -> !TVar SubscriberStatus
+ Kiroku.Store.Subscription.EventPublisher: [unsubscribe] :: PublisherSubscription -> !IO ()
+ Kiroku.Store.Subscription.EventPublisher: data PublisherSubscription
+ Kiroku.Store.Subscription.EventPublisher: subscribePublisherWith :: EventPublisher -> Natural -> OverflowPolicy -> STM PublisherSubscription
+ Kiroku.Store.Subscription.Types: SubscriptionDeadLetter :: !Int64 -> !SubscriptionName -> !Int32 -> !GlobalPosition -> !EventId -> !Value -> !Text -> !Int32 -> !UTCTime -> SubscriptionDeadLetter
+ Kiroku.Store.Subscription.Types: SubscriptionDeadLetterCursor :: !GlobalPosition -> !Int64 -> SubscriptionDeadLetterCursor
+ Kiroku.Store.Subscription.Types: SubscriptionDeadLetterLimitOutOfRange :: Int32 -> SubscriptionDeadLetterLimitOutOfRange
+ Kiroku.Store.Subscription.Types: SubscriptionDeadLetterPage :: !Vector SubscriptionDeadLetter -> !Maybe SubscriptionDeadLetterCursor -> SubscriptionDeadLetterPage
+ Kiroku.Store.Subscription.Types: SubscriptionDeadLetterQuery :: !SubscriptionName -> !Maybe Int32 -> !Maybe SubscriptionDeadLetterCursor -> !SubscriptionDeadLetterLimit -> SubscriptionDeadLetterQuery
+ Kiroku.Store.Subscription.Types: [after] :: SubscriptionDeadLetterQuery -> !Maybe SubscriptionDeadLetterCursor
+ Kiroku.Store.Subscription.Types: [attemptCount] :: SubscriptionDeadLetter -> !Int32
+ Kiroku.Store.Subscription.Types: [createdAt] :: SubscriptionDeadLetter -> !UTCTime
+ Kiroku.Store.Subscription.Types: [cursorDeadLetterId] :: SubscriptionDeadLetterCursor -> !Int64
+ Kiroku.Store.Subscription.Types: [cursorGlobalPosition] :: SubscriptionDeadLetterCursor -> !GlobalPosition
+ Kiroku.Store.Subscription.Types: [deadLetterId] :: SubscriptionDeadLetter -> !Int64
+ Kiroku.Store.Subscription.Types: [deadLetters] :: SubscriptionDeadLetterPage -> !Vector SubscriptionDeadLetter
+ Kiroku.Store.Subscription.Types: [eventId] :: SubscriptionDeadLetter -> !EventId
+ Kiroku.Store.Subscription.Types: [globalPosition] :: SubscriptionDeadLetter -> !GlobalPosition
+ Kiroku.Store.Subscription.Types: [limit] :: SubscriptionDeadLetterQuery -> !SubscriptionDeadLetterLimit
+ Kiroku.Store.Subscription.Types: [nextCursor] :: SubscriptionDeadLetterPage -> !Maybe SubscriptionDeadLetterCursor
+ Kiroku.Store.Subscription.Types: [reasonSummary] :: SubscriptionDeadLetter -> !Text
+ Kiroku.Store.Subscription.Types: [reason] :: SubscriptionDeadLetter -> !Value
+ Kiroku.Store.Subscription.Types: data SubscriptionDeadLetter
+ Kiroku.Store.Subscription.Types: data SubscriptionDeadLetterCursor
+ Kiroku.Store.Subscription.Types: data SubscriptionDeadLetterLimit
+ Kiroku.Store.Subscription.Types: data SubscriptionDeadLetterPage
+ Kiroku.Store.Subscription.Types: data SubscriptionDeadLetterQuery
+ Kiroku.Store.Subscription.Types: defaultSubscriptionDeadLetterLimit :: SubscriptionDeadLetterLimit
+ Kiroku.Store.Subscription.Types: defaultSubscriptionDeadLetterQuery :: SubscriptionName -> SubscriptionDeadLetterQuery
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.SubscriptionDeadLetter
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.SubscriptionDeadLetterCursor
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.SubscriptionDeadLetterLimit
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.SubscriptionDeadLetterLimitOutOfRange
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.SubscriptionDeadLetterPage
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.SubscriptionDeadLetterQuery
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Ord Kiroku.Store.Subscription.Types.SubscriptionDeadLetterCursor
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Ord Kiroku.Store.Subscription.Types.SubscriptionDeadLetterLimit
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Generics.Generic Kiroku.Store.Subscription.Types.SubscriptionDeadLetter
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Generics.Generic Kiroku.Store.Subscription.Types.SubscriptionDeadLetterCursor
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Generics.Generic Kiroku.Store.Subscription.Types.SubscriptionDeadLetterPage
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Generics.Generic Kiroku.Store.Subscription.Types.SubscriptionDeadLetterQuery
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.SubscriptionDeadLetter
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.SubscriptionDeadLetterCursor
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.SubscriptionDeadLetterLimit
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.SubscriptionDeadLetterLimitOutOfRange
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.SubscriptionDeadLetterPage
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.SubscriptionDeadLetterQuery
+ Kiroku.Store.Subscription.Types: mkSubscriptionDeadLetterLimit :: Int32 -> Either SubscriptionDeadLetterLimitOutOfRange SubscriptionDeadLetterLimit
+ Kiroku.Store.Subscription.Types: newtype SubscriptionDeadLetterLimitOutOfRange
+ Kiroku.Store.Subscription.Types: subscriptionDeadLetterCursor :: SubscriptionDeadLetter -> SubscriptionDeadLetterCursor
+ Kiroku.Store.Subscription.Types: subscriptionDeadLetterLimitValue :: SubscriptionDeadLetterLimit -> Int32
+ Kiroku.Store.Types: InvalidBrowsePageSize :: !Int -> BrowsePageSizeError
+ Kiroku.Store.Types: browsePageSizeValue :: BrowsePageSize -> Int32
+ Kiroku.Store.Types: data BrowsePageSize
+ Kiroku.Store.Types: data BrowsePageSizeError
+ Kiroku.Store.Types: instance GHC.Classes.Eq Kiroku.Store.Types.BrowsePageSize
+ Kiroku.Store.Types: instance GHC.Classes.Eq Kiroku.Store.Types.BrowsePageSizeError
+ Kiroku.Store.Types: instance GHC.Internal.Show.Show Kiroku.Store.Types.BrowsePageSize
+ Kiroku.Store.Types: instance GHC.Internal.Show.Show Kiroku.Store.Types.BrowsePageSizeError
+ Kiroku.Store.Types: mkBrowsePageSize :: Int -> Either BrowsePageSizeError BrowsePageSize
- Kiroku.Store.Subscription.EventPublisher: Subscriber :: !TBQueue DecodedBatch -> !TVar SubscriberStatus -> !OverflowPolicy -> Subscriber
+ Kiroku.Store.Subscription.EventPublisher: Subscriber :: !TBQueue DecodedBatch -> !TVar SubscriberStatus -> !OverflowPolicy -> !TVar Word64 -> Subscriber
- Kiroku.Store.Subscription.Types: [consumerGroupMember] :: SubscriptionCheckpoint -> !Int32
+ Kiroku.Store.Subscription.Types: [consumerGroupMember] :: SubscriptionDeadLetterQuery -> !Maybe Int32

Files

CHANGELOG.md view
@@ -1,5 +1,34 @@ # Changelog +## 0.11.0.0 — 2026-10-11++### Breaking Changes++- The closed `Store` effect gains `GetStreamWithHead`; exhaustive custom interpreters+  must handle it. `StreamInfo` construction is unchanged.++- `Subscriber` gains `subDropped :: TVar Word64`, counting batches discarded only under `DropOldest`; complete record constructors must supply it.++- The closed `Store` effect gains `ListSubscriptionDeadLetters`; exhaustive custom interpreters must handle it.++- The closed `Store` effect gains `ListStreams`, `ListCategories` and `GetEvent`.+  Custom interpreters must handle them. Catalog limits use validated `BrowsePageSize`.++### New Features++- Opt-in `getStreamWithHead` captures metadata and the newest surviving originated+  global position in one statement. Links and `$all` do not originate a head;+  missing streams and existing streams without a head remain distinguishable.+  Existing metadata/event reads and writes acquire no extra work.++- `subscribePublisherWith` returns `PublisherSubscription`, including the dropped-batch counter and idempotent deregistration; `subscribePublisher` retains its original triple.++- Public `subscriptionDeadLetters` with validated query limits, exclusive composite cursors, and newest-first pages across historical members or one selected member. Structured reasons are unchanged.++- Byte-ordered stream browsing with exact category, literal prefix and exclusive+  name cursors; paged category enumeration and canonical event lookup by id.++ ## 0.10.0.0 — 2026-10-10  ### Breaking Changes
bench/ShibuyaOverhead.hs view
@@ -30,7 +30,7 @@ import Kiroku.Store import Kiroku.Store.Subscription.Stream (subscriptionStream) import Kiroku.Store.Subscription.Stream qualified as Buffer-import Kiroku.Test.Postgres (ephemeralConfig)+import Kiroku.Test.Postgres (ephemeralConfig, migrateTestDatabase) import Shibuya.Adapter (Adapter (..)) import Shibuya.App (ProcessorId (..), defaultAppConfig, mkProcessor, runApp, stopApp) import Shibuya.Core.Ack (AckDecision (..))@@ -50,6 +50,7 @@ main = do     config <- ephemeralConfig     result <- Pg.withCachedConfig config Pg.defaultCacheConfig $ \db -> do+        migrateTestDatabase (Pg.connectionString db)         let settings = defaultConnectionSettings (Pg.connectionString db)         withStore settings $ \store -> do             putStrLn "=== Kiroku Shibuya Adapter Overhead Benchmark ==="
+ bench/StreamHeadCost.hs view
@@ -0,0 +1,119 @@+{-# LANGUAGE MultilineStrings #-}++module Main where++import Control.Exception (finally)+import Control.Lens ((^.))+import Control.Monad (replicateM_, unless)+import Data.Generics.Labels ()+import Data.Text qualified as T+import Data.Time.Clock (diffUTCTime, getCurrentTime)+import Effectful+import Effectful.Dispatch.Dynamic (interpret_)+import Effectful.Error.Static (Error, runErrorNoCallStack, throwError)+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.Fixtures.StreamHead (streamHeadFixtureSql)+import Kiroku.Test.Postgres (withMigratedTestDatabase, withSharedMigratedPostgres)+import Test.Tasty.Bench++main :: IO ()+main = do+    start <- getCurrentTime+    withSharedMigratedPostgres $ withMigratedTestDatabase $ \connection ->+        withStore (defaultConnectionSettings connection) $ \store -> do+            sql store streamHeadFixtureSql+            sql store "VACUUM (ANALYZE) stream_events"+            details <-+                Pool.use (store ^. #pool) $+                    Session.statement () $+                        preparable+                            "SELECT json_build_object('server', version(), 'shared_buffers', current_setting('shared_buffers'), 'work_mem', current_setting('work_mem'), 'jit', current_setting('jit'), 'plan_cache_mode', current_setting('plan_cache_mode'), 'streams', (SELECT count(*) FROM streams), 'events', (SELECT count(*) FROM events), 'junctions', (SELECT count(*) FROM stream_events))::text"+                            E.noParams+                            (D.singleRow (D.column (D.nonNullable D.text)))+            either (error . show) (putStrLn . T.unpack) details+            setup <- getCurrentTime+            putStrLn ("setup_seconds=" <> show (diffUTCTime setup start))+            groups <- mapM (sizeGroup store) [("100", StreamName "bench-1"), ("100000", StreamName "long-1")]+            warmed <- getCurrentTime+            putStrLn ("warmup_seconds=" <> show (diffUTCTime warmed setup))+            defaultMain groups `finally` do+                finished <- getCurrentTime+                putStrLn ("measurement_seconds=" <> show (diffUTCTime finished warmed))+                putStrLn ("total_seconds=" <> show (diffUTCTime finished start))++sizeGroup :: KirokuStore -> (String, StreamName) -> IO Benchmark+sizeGroup store (label, name) = do+    Right (Just expected) <- runFrozen store (getStream name)+    unless (expected ^. #version == StreamVersion (read label)) (error "invalid fixture version")+    let validate result = case result of+            Right (Just actual) -> unless (actual == expected) (error "unexpected metadata")+            other -> error ("metadata read failed: " <> show other)+        control = replicateM_ 100 (runFrozen store (getStream name) >>= validate)+        production = replicateM_ 100 (runStoreIO store (getStream name) >>= validate)+    -- Equality checks every StreamInfo field in successful calls in both arms.+    control+    production+    let withHead = replicateM_ 100 $ do+            result <- runStoreIO store (getStreamWithHead name)+            case result of+                Right (Just (actual, Just headPosition)) -> do+                    validate (Right (Just actual))+                    unless+                        (headPosition == GlobalPosition (if label == "100" then 99001 else 200000))+                        (error "unexpected originated head")+                other -> error ("head read failed: " <> show other)+    withHead+    pure $+        bgroup+            label+            [ bench "control-metadata" (whnfIO control)+            , bcompareWithin 0 1.10 ("$(NF-1) == \"" <> label <> "\" && $NF == \"control-metadata\"") $+                bench "production-metadata" (whnfIO production)+            , bench "production-with-head" (whnfIO withHead)+            ]++sql :: KirokuStore -> T.Text -> IO ()+sql store command = Pool.use (store ^. #pool) (Session.script command) >>= either (error . show) pure++-- Frozen from 109d58f57dbd5757ad55792474d046a37cc2e87d before EP-97.+-- Do not substitute production SQL, decoders, handlers or pool helpers here.+runFrozen :: KirokuStore -> Eff '[Store, Error StoreError, IOE] a -> IO (Either StoreError a)+runFrozen store = runEff . runErrorNoCallStack . frozenInterpreter store++frozenInterpreter :: (IOE :> es, Error StoreError :> es) => KirokuStore -> Eff (Store : es) a -> Eff es a+frozenInterpreter store = interpret_ $ \case+    GetStream (StreamName name) -> frozenPool (store ^. #pool) (Session.statement name frozenStatement)+    _ -> error "unexpected operation in frozen metadata control"++frozenPool :: (IOE :> es, Error StoreError :> es) => Pool.Pool -> Session.Session a -> Eff es a+frozenPool pool session = do+    result <- liftIO (Pool.use pool session)+    case result of+        Left usageErr -> throwError (ConnectionError (T.pack (show usageErr)))+        Right a -> pure a++frozenStatement :: Statement T.Text (Maybe StreamInfo)+frozenStatement =+    preparable+        """+        SELECT stream_id, stream_name, stream_version, created_at, deleted_at, truncate_before+        FROM streams+        WHERE stream_name = $1+        """+        (E.param (E.nonNullable E.text))+        (D.rowMaybe frozenRow)++frozenRow :: D.Row StreamInfo+frozenRow =+    StreamInfo+        <$> (StreamId <$> D.column (D.nonNullable D.int8))+        <*> (StreamName <$> D.column (D.nonNullable D.text))+        <*> (StreamVersion <$> D.column (D.nonNullable D.int8))+        <*> D.column (D.nonNullable D.timestamptz)+        <*> D.column (D.nullable D.timestamptz)+        <*> (StreamVersion <$> D.column (D.nonNullable D.int8))
+ bench/StreamHeadPaired.hs view
@@ -0,0 +1,145 @@+{-# LANGUAGE MultilineStrings #-}++module Main where++import Control.Lens ((^.))+import Control.Monad (forM_, replicateM_, unless)+import Data.Generics.Labels ()+import Data.Text qualified as T+import Data.Time.Clock (diffUTCTime, getCurrentTime)+import Effectful+import Effectful.Dispatch.Dynamic (interpret_)+import Effectful.Error.Static (Error, runErrorNoCallStack, throwError)+import GHC.Clock (getMonotonicTimeNSec)+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.Fixtures.StreamHead (streamHeadFixtureSql)+import Kiroku.Test.Postgres (withMigratedTestDatabase, withSharedMigratedPostgres)+import System.Environment (getArgs)+import System.IO (BufferMode (LineBuffering), IOMode (WriteMode), hPutStrLn, hSetBuffering, withFile)++-- Fixed, paired measurements. The cell entry writes a predeclared balanced+-- ABBA/BAAB schedule. Each CSV row is one batch; no adaptive early stopping.+main :: IO ()+main = do+    [schedulePath, outputPath, callsText] <- getArgs+    let calls = read callsText :: Int+    unless (calls > 0 && calls <= 10000 && calls `mod` 5 == 0) (error "invalid batch size")+    schedule <- map (T.splitOn ",") . T.lines . T.pack <$> readFile schedulePath+    start <- getCurrentTime+    withSharedMigratedPostgres $ withMigratedTestDatabase $ \connection ->+        withStore (defaultConnectionSettings connection) $ \store -> do+            sql store streamHeadFixtureSql+            sql store "VACUUM (ANALYZE) stream_events"+            details <-+                Pool.use (store ^. #pool) $+                    Session.statement () $+                        preparable+                            "SELECT json_build_object('server', version(), 'shared_buffers', current_setting('shared_buffers'), 'work_mem', current_setting('work_mem'), 'jit', current_setting('jit'), 'plan_cache_mode', current_setting('plan_cache_mode'), 'streams', (SELECT count(*) FROM streams), 'events', (SELECT count(*) FROM events), 'junctions', (SELECT count(*) FROM stream_events))::text"+                            E.noParams+                            (D.singleRow (D.column (D.nonNullable D.text)))+            either (error . show) (putStrLn . T.unpack) details+            setup <- getCurrentTime+            putStrLn ("setup_seconds=" <> show (diffUTCTime setup start))+            small <- actions store (StreamName "bench-1") (StreamVersion 100)+            large <- actions store (StreamName "long-1") (StreamVersion 100000)+            let (frozenSmall, productionSmall) = small+                (frozenLarge, productionLarge) = large+                groups =+                    [ ("legacy-100", (frozenSmall, productionSmall, calls))+                    , ("legacy-100000", (frozenLarge, productionLarge, calls))+                    , ("aa-100", (productionSmall, productionSmall, calls))+                    , ("aa-100000", (productionLarge, productionLarge, calls))+                    , ("positive-20pct", (productionLarge, productionLarge, calls * 6 `div` 5))+                    ]+            -- Warm both actual implementations and connections before measurement.+            forM_ [frozenSmall, productionSmall, frozenLarge, productionLarge] (replicateM_ 1000)+            warmed <- getCurrentTime+            putStrLn ("warmup_seconds=" <> show (diffUTCTime warmed setup))+            withFile outputPath WriteMode $ \out -> do+                hSetBuffering out LineBuffering+                hPutStrLn out "round,case,order,slot,arm,calls,start_ns,end_ns,duration_ns"+                forM_ schedule $ \row -> case row of+                    [roundText, caseName, order] -> do+                        unless (order == "ABBA" || order == "BAAB") (error "invalid order")+                        (armA, armB, callsB) <- maybe (error "unknown case") pure (lookup caseName groups)+                        forM_ (zip [0 :: Int ..] (T.unpack order)) $ \(slot, arm) -> do+                            let count = if arm == 'A' then calls else callsB+                                action = if arm == 'A' then armA else armB+                            before <- getMonotonicTimeNSec+                            replicateM_ count action+                            after <- getMonotonicTimeNSec+                            hPutStrLn out $+                                T.unpack $+                                    T.intercalate+                                        ","+                                        [ roundText+                                        , caseName+                                        , order+                                        , T.pack (show slot)+                                        , T.singleton arm+                                        , T.pack (show count)+                                        , T.pack (show before)+                                        , T.pack (show after)+                                        , T.pack (show (after - before))+                                        ]+                    _ -> error "invalid schedule row"+            finished <- getCurrentTime+            putStrLn ("measurement_seconds=" <> show (diffUTCTime finished warmed))+            putStrLn ("total_seconds=" <> show (diffUTCTime finished start))++-- The A/A arms receive the same IO action, including validation. The positive+-- control repeats that same action 20% more times, without normalizing it away.+actions :: KirokuStore -> StreamName -> StreamVersion -> IO (IO (), IO ())+actions store name version = do+    Right (Just expected) <- runFrozen store (getStream name)+    unless (expected ^. #version == version) (error "invalid fixture version")+    let validate result = case result of+            Right (Just actual) -> unless (actual == expected) (error "unexpected metadata")+            other -> error ("metadata read failed: " <> show other)+    pure (runFrozen store (getStream name) >>= validate, runStoreIO store (getStream name) >>= validate)++sql :: KirokuStore -> T.Text -> IO ()+sql store command = Pool.use (store ^. #pool) (Session.script command) >>= either (error . show) pure++-- Frozen from 109d58f57dbd5757ad55792474d046a37cc2e87d before EP-97.+-- Do not substitute production SQL, decoders, handlers or pool helpers here.+runFrozen :: KirokuStore -> Eff '[Store, Error StoreError, IOE] a -> IO (Either StoreError a)+runFrozen store = runEff . runErrorNoCallStack . frozenInterpreter store++frozenInterpreter :: (IOE :> es, Error StoreError :> es) => KirokuStore -> Eff (Store : es) a -> Eff es a+frozenInterpreter store = interpret_ $ \case+    GetStream (StreamName name) -> frozenPool (store ^. #pool) (Session.statement name frozenStatement)+    _ -> error "unexpected operation in frozen metadata control"++frozenPool :: (IOE :> es, Error StoreError :> es) => Pool.Pool -> Session.Session a -> Eff es a+frozenPool pool session = do+    result <- liftIO (Pool.use pool session)+    case result of+        Left usageErr -> throwError (ConnectionError (T.pack (show usageErr)))+        Right a -> pure a++frozenStatement :: Statement T.Text (Maybe StreamInfo)+frozenStatement =+    preparable+        """+        SELECT stream_id, stream_name, stream_version, created_at, deleted_at, truncate_before+        FROM streams+        WHERE stream_name = $1+        """+        (E.param (E.nonNullable E.text))+        (D.rowMaybe frozenRow)++frozenRow :: D.Row StreamInfo+frozenRow =+    StreamInfo+        <$> (StreamId <$> D.column (D.nonNullable D.int8))+        <*> (StreamName <$> D.column (D.nonNullable D.text))+        <*> (StreamVersion <$> D.column (D.nonNullable D.int8))+        <*> D.column (D.nonNullable D.timestamptz)+        <*> D.column (D.nullable D.timestamptz)+        <*> (StreamVersion <$> D.column (D.nonNullable D.int8))
kiroku-store.cabal view
@@ -1,6 +1,6 @@ cabal-version:   3.0 name:            kiroku-store-version:         0.10.0.0+version:         0.11.0.0 synopsis:        High-performance PostgreSQL event store description:   Kiroku is a PostgreSQL-backed event store for Haskell applications. It@@ -16,6 +16,7 @@ build-type:      Simple category:        Database, Eventing extra-doc-files: CHANGELOG.md+data-files:      test/fixtures/stream-head-isolation/*.sql  source-repository head   type:     git@@ -66,6 +67,7 @@     Kiroku.Store.HistoryRetention.SQL     Kiroku.Store.Subscription.Checkpoint.SQL     Kiroku.Store.Subscription.CheckpointInventory.SQL+    Kiroku.Store.Subscription.DeadLetter.SQL    build-depends:     , aeson                 >=2.1   && <2.3@@ -94,11 +96,15 @@   hs-source-dirs:  src  test-suite kiroku-store-test-  import:         common-  type:           exitcode-stdio-1.0-  main-is:        Main.hs-  hs-source-dirs: test+  import:          common+  type:            exitcode-stdio-1.0+  main-is:         Main.hs+  hs-source-dirs:  test   other-modules:+    Paths_kiroku_store+    Test.BrowseQueryPlans+    Test.BrowseReads+    Test.BrowseReadsMock     Test.CatchupDbErrorNoPrematureSwitch     Test.Category     Test.CategoryIdleNoSpin@@ -108,6 +114,7 @@     Test.ConsumerGroupEffect     Test.ConsumerGroupResize     Test.ConsumerGroupSql+    Test.DeadLetterQueryPlans     Test.EventTypeFilter     Test.FailureInjection     Test.HandlerStall@@ -119,11 +126,15 @@     Test.PerformanceStructure     Test.Properties     Test.PublisherCallbackResilience+    Test.PublisherDropCounter     Test.PublisherIdleAdvance     Test.PublisherRestartNoRebroadcast     Test.ReadStream     Test.StartupFailureSurfacing     Test.StreamBridgeTermination+    Test.StreamHead+    Test.StreamHeadIsolation+    Test.StreamHeadMock     Test.StreamHistoryGuard     Test.StreamNameLookup     Test.SubscriptionCheckpointInitialization@@ -132,6 +143,8 @@     Test.SubscriptionCheckpointInventoryMock     Test.SubscriptionCheckpointReset     Test.SubscriptionCheckpointWorker+    Test.SubscriptionDeadLetters+    Test.SubscriptionDeadLettersMock     Test.SubscriptionPauseResume     Test.SubscriptionReconnect     Test.SubscriptionRegistry@@ -144,7 +157,8 @@     Test.VisibleGlobalHeadPosition     Test.VisibleGlobalHeadPositionMock -  ghc-options:    -threaded -rtsopts "-with-rtsopts=-N -T"+  autogen-modules: Paths_kiroku_store+  ghc-options:     -threaded -rtsopts "-with-rtsopts=-N -T"   build-depends:     , aeson                 >=2.1   && <2.3     , async                 >=2.2   && <2.3@@ -294,3 +308,53 @@     , kiroku-test-support     , lens     , text++benchmark kiroku-stream-head-cost+  import:         common+  type:           exitcode-stdio-1.0+  main-is:        StreamHeadCost.hs+  hs-source-dirs: bench+  ghc-options:+    -threaded -rtsopts "-with-rtsopts=-N -A32m" -fproc-alignment=64++  build-depends:+    , aeson                >=2.1   && <2.3+    , base                 >=4.18  && <5+    , effectful-core       >=2.6.1 && <2.7  || >=2.7.1.1 && <2.8+    , generic-lens         >=2.2   && <2.4+    , hasql                >=1.10  && <1.11+    , hasql-pool           >=1.2   && <1.5+    , hasql-transaction    >=1.1   && <1.3+    , kiroku-store+    , kiroku-test-support+    , lens                 >=5.2   && <5.4+    , tasty                >=1.4   && <1.6+    , tasty-bench          >=0.4+    , text                 >=2.0   && <2.2+    , time                 >=1.12  && <1.15+    , vector               >=0.13  && <0.14++benchmark kiroku-stream-head-paired+  import:         common+  type:           exitcode-stdio-1.0+  main-is:        StreamHeadPaired.hs+  hs-source-dirs: bench+  ghc-options:+    -threaded -rtsopts "-with-rtsopts=-N -A32m" -fproc-alignment=64++  build-depends:+    , aeson                >=2.1   && <2.3+    , base                 >=4.18  && <5+    , effectful-core       >=2.6.1 && <2.7  || >=2.7.1.1 && <2.8+    , generic-lens         >=2.2   && <2.4+    , hasql                >=1.10  && <1.11+    , hasql-pool           >=1.2   && <1.5+    , hasql-transaction    >=1.1   && <1.3+    , kiroku-store+    , kiroku-test-support+    , lens                 >=5.2   && <5.4+    , tasty                >=1.4   && <1.6+    , tasty-bench          >=0.4+    , text                 >=2.0   && <2.2+    , time                 >=1.12  && <1.15+    , vector               >=0.13  && <0.14
src/Kiroku/Store/Effect.hs view
@@ -59,11 +59,14 @@ import Kiroku.Store.Settings (StoreSettings (..), enrichEvents) import Kiroku.Store.Subscription.Checkpoint.SQL qualified as CheckpointSQL import Kiroku.Store.Subscription.CheckpointInventory.SQL qualified as CheckpointInventorySQL+import Kiroku.Store.Subscription.DeadLetter.SQL qualified as DeadLetterSQL import Kiroku.Store.Subscription.Types (     CheckpointInitialization,     MissingCheckpointPolicy,     SubscriptionCheckpointInventory,     SubscriptionCheckpointMissing,+    SubscriptionDeadLetterPage,+    SubscriptionDeadLetterQuery,     SubscriptionName,  ) import Kiroku.Store.Types@@ -88,6 +91,12 @@     -}     GetVisibleGlobalHeadPosition :: Store m GlobalPosition     GetStream :: StreamName -> Store m (Maybe StreamInfo)+    {- | One-statement metadata and newest surviving originated head observation.+    Outer 'Nothing' means absent; inner 'Nothing' means no originated events,+    including link-only streams and @$all@. Logical lifecycle markers do not+    hide the head. Custom exhaustive interpreters must handle this constructor.+    -}+    GetStreamWithHead :: StreamName -> Store m (Maybe (StreamInfo, Maybe GlobalPosition))     {- | Resolve a 'StreamName' to its surrogate 'StreamId' without     materializing the full 'StreamInfo' row. Mirrors 'GetStream'\'s     soft-delete semantics: returns 'Just' for both live and soft-deleted@@ -117,6 +126,14 @@     'Kiroku.Store.Read.lookupStreamName').     -}     LookupStreamNames :: [StreamId] -> Store m (Map StreamId StreamName)+    {- | Byte-ordered names; optional exact category, literal prefix and exclusive cursor.+    Includes soft-deleted streams, excludes only the reserved stream row.+    -}+    ListStreams :: Maybe CategoryName -> Maybe Text -> Maybe StreamName -> BrowsePageSize -> Store m (Vector StreamInfo)+    -- | Deployment-collation category order with an exclusive cursor.+    ListCategories :: Maybe CategoryName -> BrowsePageSize -> Store m (Vector CategoryName)+    -- | The canonical global-log event; absent after hard deletion.+    GetEvent :: EventId -> Store m (Maybe RecordedEvent)     LinkToStream :: StreamName -> [EventId] -> Store m LinkResult     ReadCategoryForward :: CategoryName -> GlobalPosition -> Int32 -> Store m (Vector RecordedEvent)     AppendMultiStream :: [(StreamName, ExpectedVersion, [EventData])] -> Store m [AppendResult]@@ -148,6 +165,8 @@     'Kiroku.Store.Subscription.subscriptionCheckpointInventory'.     -}     GetSubscriptionCheckpointInventory :: Store m SubscriptionCheckpointInventory+    -- | Read-only keyset page; no event decode hook or checkpoint write.+    ListSubscriptionDeadLetters :: SubscriptionDeadLetterQuery -> Store m SubscriptionDeadLetterPage     {- | Resolve one exact subscription checkpoint key according to its     missing-row policy. Existing rows always take precedence. @@ -258,6 +277,9 @@     GetVisibleGlobalHeadPosition ->         usePool (store ^. #pool) $             Session.statement () SQL.visibleGlobalHeadPositionStmt+    GetStreamWithHead (StreamName name) ->+        usePool (store ^. #pool) $+            Session.statement name SQL.getStreamWithHeadStmt     GetStream (StreamName name) ->         usePool (store ^. #pool) $             Session.statement name SQL.getStreamStmt@@ -276,6 +298,16 @@             ( usePool (store ^. #pool) $                 Session.statement [s | StreamId s <- sids] SQL.lookupStreamNamesStmt             )+    ListStreams category prefix after limit ->+        usePool (store ^. #pool) (SQL.listStreamsSession category prefix after limit)+    ListCategories after limit ->+        fmap (V.map CategoryName) $+            usePool (store ^. #pool) $+                Session.statement (fmap (\(CategoryName value) -> value) after, browsePageSizeValue limit) (SQL.listCategoriesStmt (not (isNothing after)))+    GetEvent (EventId eid) -> do+        found <- usePool (store ^. #pool) $ Session.statement eid SQL.getEventStmt+        decoded <- decodeReadEvents (store ^. #storeSettings) (maybe V.empty V.singleton found)+        pure (decoded V.!? 0)     LinkToStream (StreamName name) eventIds -> do         rejectInvalidApplicationStream name         let uuids = V.fromList [uid | EventId uid <- eventIds]@@ -403,6 +435,8 @@         rejectInvalidApplicationStream name         usePool (store ^. #pool) $             Session.statement (name, v) SQL.setStreamTruncateBeforeStmt+    ListSubscriptionDeadLetters query ->+        usePool (store ^. #pool) (DeadLetterSQL.listSubscriptionDeadLettersSession query)     GetSubscriptionCheckpointInventory ->         usePool (store ^. #pool) $             Session.statement () CheckpointInventorySQL.getSubscriptionCheckpointInventoryStmt
src/Kiroku/Store/Read.hs view
@@ -7,6 +7,10 @@     visibleGlobalHeadPosition,     readCategory,     getStream,+    getStreamWithHead,+    listStreams,+    listCategories,+    getEvent,     lookupStreamId,     eventExistsInStream,     lookupStreamName,@@ -18,6 +22,7 @@ import Data.Int (Int32) import Data.Map.Strict (Map) import Data.Map.Strict qualified as Map+import Data.Text (Text) import Data.Vector (Vector) import Data.Vector qualified as V import Effectful (Eff, (:>))@@ -184,6 +189,32 @@     Eff es (Maybe StreamInfo) getStream name = send (GetStream name) +{- | Read metadata and the newest surviving /originated/ event's global+position in one SQL statement, from one database snapshot.++Outer 'Nothing' means absent or hard-deleted. Inner 'Nothing' means the+stream exists without surviving originated events, including empty and+link-only streams. @$all@ always has an inner 'Nothing': it aggregates events+but originates none. Links advance the stream version without advancing this+head. Soft deletion and logical truncation preserve the originated head;+physical retention can remove it. This is neither the store-wide visible+head nor the monotonically allocated append frontier.++For an origin-only stream with retained required history, capture this pair,+require @version >= N@ and @Just head@, then wait for the projection cursor to+reach @head@. Waiting and timeouts belong to the consumer. The observation does+not lock history or freeze subsequent appends; later physical removal can+require a timeout. A positive linked version gives no such guarantee.++Only this opt-in operation pays for the indexed head probe. 'getStream' and+'StreamInfo' retain their existing representation and cost.+-}+getStreamWithHead ::+    (HasCallStack, Store :> es) =>+    StreamName ->+    Eff es (Maybe (StreamInfo, Maybe GlobalPosition))+getStreamWithHead name = send (GetStreamWithHead name)+ {- | Look up a stream's surrogate id by name.  Returns 'Just' the 'StreamId' for both live and soft-deleted streams (mirroring@@ -252,3 +283,22 @@     StreamId ->     Eff es (Maybe StreamName) lookupStreamName sid = Map.lookup sid <$> lookupStreamNames [sid]++{- | Page streams in stable UTF-8 byte order, independent of database locale.+Category equality includes its bare name; prefix matching is literal, including+@%@ and @_@. Use the last returned name as the exclusive cursor. Soft-deleted+streams remain visible, hard-deleted streams do not. Only the reserved row is+excluded, so @$all-x@ and empty categories are valid. A short page proves current+exhaustion. Concurrent lifecycle changes are not a snapshot-isolation guarantee.+TypeID timestamps give generation order, not event/commit order.+-}+listStreams :: (HasCallStack, Store :> es) => Maybe CategoryName -> Maybe Text -> Maybe StreamName -> BrowsePageSize -> Eff es (Vector StreamInfo)+listStreams category prefix after limit = send (ListStreams category prefix after limit)++-- | Distinct categories, in deployment locale order, after an exclusive cursor.+listCategories :: (HasCallStack, Store :> es) => Maybe CategoryName -> BrowsePageSize -> Eff es (Vector CategoryName)+listCategories after limit = send (ListCategories after limit)++-- | An event as seen in the global log, using the configured typed decode hook.+getEvent :: (HasCallStack, Store :> es) => EventId -> Eff es (Maybe RecordedEvent)+getEvent eid = send (GetEvent eid)
src/Kiroku/Store/SQL.hs view
@@ -21,8 +21,15 @@     readCategoryForwardStmt,     readCategoryEncoder,     getStreamStmt,+    getStreamWithHeadStmt,     eventExistsInStreamStmt,     lookupStreamNamesStmt,+    listStreamsSession,+    StreamRangeBound (..),+    listStreamsRangeStmt,+    listStreamsPairStmt,+    listCategoriesStmt,+    getEventStmt,     currentGlobalPositionStmt,     visibleGlobalHeadPositionStmt, @@ -69,22 +76,30 @@     DeadLetterRecord (..),     insertDeadLetterAndCheckpointStmt,     readDeadLettersStmt,+    listSubscriptionDeadLettersStmt,+    listSubscriptionDeadLettersFromStartStmt,+    listSubscriptionMemberDeadLettersStmt,+    listSubscriptionMemberDeadLettersFromStartStmt, ) where  import Contravariant.Extras (contrazip2, contrazip3, contrazip4, contrazip5, contrazip6) import Control.Lens ((^.)) import Data.Aeson (Value)+import Data.Char (chr, ord) import Data.Functor.Contravariant ((>$<)) import Data.Generics.Labels () import Data.Int (Int32, Int64) import Data.Text (Text)+import Data.Text qualified as T import Data.Time (UTCTime) import Data.UUID (UUID) import Data.Vector (Vector) import GHC.Generics (Generic) import Hasql.Decoders qualified as D import Hasql.Encoders qualified as E+import Hasql.Session qualified as Session import Hasql.Statement (Statement, preparable)+import Kiroku.Store.Subscription.DeadLetter.SQL (listSubscriptionDeadLettersFromStartStmt, listSubscriptionDeadLettersStmt, listSubscriptionMemberDeadLettersFromStartStmt, listSubscriptionMemberDeadLettersStmt) import Kiroku.Store.Types  -- | Parameters for append CTE variants (the 7 parallel arrays + stream name).@@ -483,6 +498,29 @@         E.noParams         (D.singleRow (GlobalPosition <$> D.column (D.nonNullable D.int8))) +-- | Metadata and the newest surviving originated global position in one snapshot.+getStreamWithHeadStmt :: Statement Text (Maybe (StreamInfo, Maybe GlobalPosition))+getStreamWithHeadStmt =+    preparable+        getStreamWithHeadSQL+        (E.param (E.nonNullable E.text))+        (D.rowMaybe ((,) <$> streamInfoRow <*> (fmap GlobalPosition <$> D.column (D.nullable D.int8))))++getStreamWithHeadSQL :: Text+getStreamWithHeadSQL =+    """+    SELECT s.stream_id, s.stream_name, s.stream_version,+           s.created_at, s.deleted_at, s.truncate_before,+           (SELECT se.stream_version+            FROM stream_events AS se+            WHERE se.stream_id = 0+              AND se.original_stream_id = s.stream_id+            ORDER BY se.stream_version DESC+            LIMIT 1) AS head_global_position+    FROM streams AS s+    WHERE s.stream_name = $1+    """+ -- | Get stream metadata by name. getStreamStmt :: Statement Text (Maybe StreamInfo) getStreamStmt =@@ -1416,3 +1454,121 @@       AND consumer_group_member = $2     ORDER BY global_position DESC, dead_letter_id DESC     """++-- Catalog statements have a finite family of prepared shapes. Upper bounds+-- are deliberately checked after the ordered LIMIT, keeping generic plans+-- from bitmap-scanning and sorting an entire matching prefix.+data StreamRangeBound = InclusiveLower | ExclusiveLower | ExactName+    deriving stock (Eq, Show)++streamColumns :: Text+streamColumns = "stream_id, stream_name, stream_version, created_at, deleted_at, truncate_before"++streamRangeSQL :: StreamRangeBound -> Text -> Text -> Text -> Text+streamRangeSQL bound lower upper limit =+    "SELECT * FROM (SELECT "+        <> streamColumns+        <> " FROM streams WHERE stream_id <> 0 AND stream_name COLLATE \"C\" "+        <> operator+        <> " "+        <> lower+        <> " ORDER BY stream_name COLLATE \"C\" LIMIT "+        <> limit+        <> ") bounded WHERE ("+        <> upper+        <> "::text IS NULL OR stream_name COLLATE \"C\" < "+        <> upper+        <> ")"+  where+    operator = case bound of InclusiveLower -> ">="; ExclusiveLower -> ">"; ExactName -> "="++listStreamsRangeStmt :: StreamRangeBound -> Statement (Text, Maybe Text, Int32) (Vector StreamInfo)+listStreamsRangeStmt bound =+    preparable+        (streamRangeSQL bound "$1" "$2" "$3" <> " ORDER BY stream_name COLLATE \"C\"")+        (contrazip3 (E.param (E.nonNullable E.text)) (E.param (E.nullable E.text)) (E.param (E.nonNullable E.int4)))+        (D.rowVector streamInfoRow)++listStreamsPairStmt :: StreamRangeBound -> StreamRangeBound -> Statement (Text, Maybe Text, Text, Maybe Text, Int32) (Vector StreamInfo)+listStreamsPairStmt first second =+    preparable+        ( "SELECT * FROM (("+            <> streamRangeSQL first "$1" "$2" "$5"+            <> ") UNION ALL ("+            <> streamRangeSQL second "$3" "$4" "$5"+            <> ")) matching ORDER BY stream_name COLLATE \"C\" LIMIT $5"+        )+        ( contrazip5+            (E.param (E.nonNullable E.text))+            (E.param (E.nullable E.text))+            (E.param (E.nonNullable E.text))+            (E.param (E.nullable E.text))+            (E.param (E.nonNullable E.int4))+        )+        (D.rowVector streamInfoRow)++{- | One statement and one checkout, including the bare-category/descendant union.+All text is encoded as parameters; only closed operator choices build SQL.+-}+listStreamsSession :: Maybe CategoryName -> Maybe Text -> Maybe StreamName -> BrowsePageSize -> Session.Session (Vector StreamInfo)+listStreamsSession category prefix after page = case ranges of+    [] -> pure mempty+    [(kind, lower, upper)] -> Session.statement (lower, upper, limit) (listStreamsRangeStmt kind)+    [(kind, lower, upper), (kind2, lower2, upper2)] ->+        Session.statement (lower, upper, lower2, upper2, limit) (listStreamsPairStmt kind kind2)+    _ -> error "catalog range invariant: at most two disjoint ranges"+  where+    limit = browsePageSizeValue page+    cursor = fmap (\(StreamName value) -> value) after+    ranges = case category of+        Nothing -> maybe [] pure (interval "" Nothing)+        Just (CategoryName value)+            | T.any (== '-') value -> []+            | otherwise -> bare value <> maybe [] pure (interval (value <> "-") (prefixEnd (value <> "-")))+    bare value+        | maybe True (`T.isPrefixOf` value) prefix && maybe True (< value) cursor = [(ExactName, value, Nothing)]+        | otherwise = []+    interval start end =+        let lower = max start (maybe "" (\value -> value) prefix)+            upper = minimumEnd end (prefix >>= prefixEnd)+            (kind, seek) = case cursor of+                Just value | value >= lower -> (ExclusiveLower, value)+                _ -> (InclusiveLower, lower)+         in if maybe False (<= seek) upper then Nothing else Just (kind, seek, upper)+    minimumEnd Nothing b = b+    minimumEnd a Nothing = a+    minimumEnd (Just a) (Just b) = Just (min a b)++-- Valid UTF-8 byte order follows Unicode scalar order. The all-maximum and+-- empty prefixes have no finite upper bound; PostgreSQL text has no NUL.+prefixEnd :: Text -> Maybe Text+prefixEnd value = case T.unsnoc value of+    Nothing -> Nothing+    Just (initial, lastChar)+        | ord lastChar == 0x10ffff -> prefixEnd initial+        | otherwise -> Just (T.snoc initial (chr (if ord lastChar + 1 == 0xd800 then 0xe000 else ord lastChar + 1)))++-- First and later category pages have separate cursor shapes. Category order+-- retains the deployment collation and its existing category index.+listCategoriesStmt :: Bool -> Statement (Maybe Text, Int32) (Vector Text)+listCategoriesStmt after =+    preparable+        query+        (contrazip2 (E.param (E.nullable E.text)) (E.param (E.nonNullable E.int4)))+        (D.rowVector (D.column (D.nonNullable D.text)))+  where+    first = if after then "AND s.category > $1" else "AND $1::text IS NULL"+    query =+        "WITH RECURSIVE next_category AS (SELECT (SELECT s.category FROM streams s WHERE s.stream_id <> 0 "+            <> first+            <> " ORDER BY s.category LIMIT 1) AS category UNION ALL SELECT (SELECT s.category FROM streams s "+            <> "WHERE s.stream_id <> 0 AND s.category > n.category ORDER BY s.category LIMIT 1) "+            <> "FROM next_category n WHERE n.category IS NOT NULL) SELECT category FROM next_category "+            <> "WHERE category IS NOT NULL LIMIT $2"++getEventStmt :: Statement UUID (Maybe RecordedEvent)+getEventStmt =+    preparable+        "SELECT e.event_id,e.event_type,se.stream_version,se.stream_version AS global_position, se.original_stream_id,se.original_stream_version,e.data,e.metadata,e.causation_id,e.correlation_id,e.created_at FROM events e JOIN stream_events se ON se.event_id=e.event_id AND se.stream_id=0 WHERE e.event_id=$1"+        (E.param (E.nonNullable E.uuid))+        (D.rowMaybe recordedEventRow)
src/Kiroku/Store/Subscription.hs view
@@ -6,6 +6,7 @@     -- * Observability     initializeSubscriptionCheckpoint,     subscriptionCheckpointInventory,+    subscriptionDeadLetters,     subscriptionStates,     SubscriptionStateView (..), @@ -30,7 +31,7 @@ import GHC.Generics (Generic) import GHC.Stack (HasCallStack) import Kiroku.Store.Connection (KirokuStore (..))-import Kiroku.Store.Effect (Store (GetSubscriptionCheckpointInventory, InitializeSubscriptionCheckpoint))+import Kiroku.Store.Effect (Store (GetSubscriptionCheckpointInventory, InitializeSubscriptionCheckpoint, ListSubscriptionDeadLetters)) import Kiroku.Store.Notification qualified as Notifier import Kiroku.Store.Subscription.EventPublisher qualified as Pub import Kiroku.Store.Subscription.Fsm (SubscriptionState (..), stateCursor, stateName)@@ -333,3 +334,18 @@                     }         )         cells++{- | Read one subscription's dead letters, newest first by global position and+then dead-letter id. Echo the page's 'nextCursor' as 'after' for an exclusive+next page. 'consumerGroupMember' Nothing includes all historical members; Just+selects one (ungrouped subscriptions use 0). All-member work is bounded by the+number of historical members times the page size; use member-scoped polling for+large groups. One statement and one pool checkout serve a page.++The reason JSON is returned unchanged, without an event decoding hook. An empty+page means no dead letters were recorded for this query. Hard-deleting source+streams can remove rows, but their cursors remain valid. This operation neither+writes a checkpoint nor changes the worker or cleanup paths.+-}+subscriptionDeadLetters :: (HasCallStack, Store :> es) => SubscriptionDeadLetterQuery -> Eff es SubscriptionDeadLetterPage+subscriptionDeadLetters = send . ListSubscriptionDeadLetters
+ src/Kiroku/Store/Subscription/DeadLetter/SQL.hs view
@@ -0,0 +1,142 @@+{-# LANGUAGE MultilineStrings #-}++module Kiroku.Store.Subscription.DeadLetter.SQL (+    listSubscriptionDeadLettersSession,+    listSubscriptionDeadLettersStmt,+    listSubscriptionDeadLettersFromStartStmt,+    listSubscriptionMemberDeadLettersStmt,+    listSubscriptionMemberDeadLettersFromStartStmt,+) where++import Contravariant.Extras (contrazip2, contrazip3, contrazip4, contrazip5)+import Control.Lens ((^.))+import Data.Generics.Labels ()+import Data.Int (Int32, Int64)+import Data.Text (Text)+import Data.Vector (Vector)+import Data.Vector qualified as V+import Hasql.Decoders qualified as D+import Hasql.Encoders qualified as E+import Hasql.Session (Session)+import Hasql.Session qualified as Session+import Hasql.Statement (Statement, preparable)+import Kiroku.Store.Subscription.Types+import Kiroku.Store.Types (EventId (..), GlobalPosition (..))++listSubscriptionDeadLettersFromStartStmt :: Statement (Text, Int32) (Vector SubscriptionDeadLetter)+listSubscriptionDeadLettersFromStartStmt =+    preparable+        """+        WITH RECURSIVE members(member) AS (+            SELECT (SELECT consumer_group_member FROM dead_letters+                    WHERE subscription_name = $1 ORDER BY consumer_group_member LIMIT 1)+            UNION ALL+            SELECT (SELECT consumer_group_member FROM dead_letters+                    WHERE subscription_name = $1 AND consumer_group_member > m.member+                    ORDER BY consumer_group_member LIMIT 1)+            FROM members m WHERE m.member IS NOT NULL+        )+        SELECT d.dead_letter_id, d.subscription_name, d.consumer_group_member, d.global_position, d.event_id, d.reason, d.reason_summary, d.attempt_count, d.created_at+        FROM members m+        CROSS JOIN LATERAL (+        SELECT dead_letter_id, subscription_name, consumer_group_member, global_position, event_id, reason, reason_summary, attempt_count, created_at+        FROM dead_letters+        WHERE subscription_name = $1 AND consumer_group_member = m.member++        ORDER BY global_position DESC, dead_letter_id DESC LIMIT $2+        ) d+        WHERE m.member IS NOT NULL+        ORDER BY d.global_position DESC, d.dead_letter_id DESC LIMIT $2+        """+        (contrazip2 text int4)+        (D.rowVector deadLetterRow)++listSubscriptionDeadLettersStmt :: Statement (Text, Int64, Int64, Int32) (Vector SubscriptionDeadLetter)+listSubscriptionDeadLettersStmt =+    preparable+        """+        WITH RECURSIVE members(member) AS (+            SELECT (SELECT consumer_group_member FROM dead_letters+                    WHERE subscription_name = $1 ORDER BY consumer_group_member LIMIT 1)+            UNION ALL+            SELECT (SELECT consumer_group_member FROM dead_letters+                    WHERE subscription_name = $1 AND consumer_group_member > m.member+                    ORDER BY consumer_group_member LIMIT 1)+            FROM members m WHERE m.member IS NOT NULL+        )+        SELECT d.dead_letter_id, d.subscription_name, d.consumer_group_member, d.global_position, d.event_id, d.reason, d.reason_summary, d.attempt_count, d.created_at+        FROM members m+        CROSS JOIN LATERAL (+        SELECT dead_letter_id, subscription_name, consumer_group_member, global_position, event_id, reason, reason_summary, attempt_count, created_at+        FROM dead_letters+        WHERE subscription_name = $1 AND consumer_group_member = m.member+        AND (global_position, dead_letter_id) < ($2, $3)+        ORDER BY global_position DESC, dead_letter_id DESC LIMIT $4+        ) d+        WHERE m.member IS NOT NULL+        ORDER BY d.global_position DESC, d.dead_letter_id DESC LIMIT $4+        """+        (contrazip4 text int8 int8 int4)+        (D.rowVector deadLetterRow)++listSubscriptionMemberDeadLettersFromStartStmt :: Statement (Text, Int32, Int32) (Vector SubscriptionDeadLetter)+listSubscriptionMemberDeadLettersFromStartStmt =+    preparable+        """+        SELECT dead_letter_id, subscription_name, consumer_group_member, global_position, event_id, reason, reason_summary, attempt_count, created_at+        FROM dead_letters+        WHERE subscription_name = $1 AND consumer_group_member = $2++        ORDER BY global_position DESC, dead_letter_id DESC LIMIT $3+        """+        (contrazip3 text int4 int4)+        (D.rowVector deadLetterRow)++listSubscriptionMemberDeadLettersStmt :: Statement (Text, Int32, Int64, Int64, Int32) (Vector SubscriptionDeadLetter)+listSubscriptionMemberDeadLettersStmt =+    preparable+        """+        SELECT dead_letter_id, subscription_name, consumer_group_member, global_position, event_id, reason, reason_summary, attempt_count, created_at+        FROM dead_letters+        WHERE subscription_name = $1 AND consumer_group_member = $2+        AND (global_position, dead_letter_id) < ($3, $4)+        ORDER BY global_position DESC, dead_letter_id DESC LIMIT $5+        """+        (contrazip5 text int4 int8 int8 int4)+        (D.rowVector deadLetterRow)++text :: E.Params Text+text = E.param (E.nonNullable E.text)+int4 :: E.Params Int32+int4 = E.param (E.nonNullable E.int4)+int8 :: E.Params Int64+int8 = E.param (E.nonNullable E.int8)++deadLetterRow :: D.Row SubscriptionDeadLetter+deadLetterRow =+    SubscriptionDeadLetter+        <$> D.column (D.nonNullable D.int8)+        <*> (SubscriptionName <$> D.column (D.nonNullable D.text))+        <*> D.column (D.nonNullable D.int4)+        <*> (GlobalPosition <$> D.column (D.nonNullable D.int8))+        <*> (EventId <$> D.column (D.nonNullable D.uuid))+        <*> D.column (D.nonNullable D.jsonb)+        <*> D.column (D.nonNullable D.text)+        <*> D.column (D.nonNullable D.int4)+        <*> D.column (D.nonNullable D.timestamptz)++listSubscriptionDeadLettersSession :: SubscriptionDeadLetterQuery -> Session SubscriptionDeadLetterPage+listSubscriptionDeadLettersSession query = do+    let SubscriptionName name = query ^. #subscriptionName+        pageSize = subscriptionDeadLetterLimitValue (query ^. #limit)+        fetch = pageSize + 1+    rows <- case (query ^. #consumerGroupMember, query ^. #after) of+        (Nothing, Nothing) -> Session.statement (name, fetch) listSubscriptionDeadLettersFromStartStmt+        (Just selectedMember, Nothing) -> Session.statement (name, selectedMember, fetch) listSubscriptionMemberDeadLettersFromStartStmt+        (Nothing, Just (SubscriptionDeadLetterCursor (GlobalPosition position) ident)) ->+            Session.statement (name, position, ident, fetch) listSubscriptionDeadLettersStmt+        (Just selectedMember, Just (SubscriptionDeadLetterCursor (GlobalPosition position) ident)) ->+            Session.statement (name, selectedMember, position, ident, fetch) listSubscriptionMemberDeadLettersStmt+    let page = V.take (fromIntegral pageSize) rows+        cursor = if V.length rows > fromIntegral pageSize then Just (subscriptionDeadLetterCursor (V.last page)) else Nothing+    pure (SubscriptionDeadLetterPage page cursor)
src/Kiroku/Store/Subscription/EventPublisher.hs view
@@ -15,7 +15,7 @@ 'SubscriberStatus': under the default @PauseAndResume@ it marks the subscriber 'Paused' and stops pushing (the worker drains and re-catches-up losslessly), under @DropSubscription@ it marks it 'Overflowed', and under @DropOldest@ it-evicts the oldest batch. The publisher itself never blocks on a slow consumer.+evicts the oldest batch and counts the drop in 'subDropped'. The publisher itself never blocks on a slow consumer. -} module Kiroku.Store.Subscription.EventPublisher (     EventPublisher (..),@@ -24,6 +24,8 @@     startPublisher,     stopPublisher,     subscribePublisher,+    PublisherSubscription (..),+    subscribePublisherWith,     publisherPosition, ) where @@ -59,6 +61,7 @@ import Data.IntMap.Strict (IntMap) import Data.IntMap.Strict qualified as IntMap import Data.Vector qualified as V+import Data.Word (Word64) import Hasql.Pool (Pool) import Hasql.Pool qualified as Pool import Hasql.Session qualified as Session@@ -93,8 +96,19 @@     { subQueue :: !(TBQueue DecodedBatch)     , subStatus :: !(TVar SubscriberStatus)     , subPolicy :: !OverflowPolicy+    , subDropped :: !(TVar Word64)+    -- ^ Dropped batches under DropOldest; compare successive readings (modulo Word64).     } +-- | A bounded publisher subscription with its exact dropped-batch counter.+data PublisherSubscription = PublisherSubscription+    { subscriptionQueue :: !(TBQueue DecodedBatch)+    , subscriptionStatus :: !(TVar SubscriberStatus)+    , subscriptionDropped :: !(TVar Word64)+    , unsubscribe :: !(IO ())+    -- ^ Idempotent deregistration; invoke on every exit.+    }+ -- | Subscriber lifecycle status as observed by its worker. data SubscriberStatus     = -- | Healthy; worker reads from the queue normally.@@ -184,16 +198,23 @@     OverflowPolicy ->     STM (TBQueue DecodedBatch, TVar SubscriberStatus, IO ()) subscribePublisher pub cap policy = do+    sub <- subscribePublisherWith pub cap policy+    pure (subscriptionQueue sub, subscriptionStatus sub, unsubscribe sub)++-- | Register a subscriber, also exposing the DropOldest counter (initially zero).+subscribePublisherWith :: EventPublisher -> Natural -> OverflowPolicy -> STM PublisherSubscription+subscribePublisherWith pub cap policy = do     queue <- newTBQueue cap     status <- newTVar Active+    dropped <- newTVar 0     sid <- readTVar (nextSubscriberId pub)     writeTVar (nextSubscriberId pub) (sid + 1)-    let sub = Subscriber{subQueue = queue, subStatus = status, subPolicy = policy}+    let sub = Subscriber{subQueue = queue, subStatus = status, subPolicy = policy, subDropped = dropped}     modifyTVar' (subscribers pub) (IntMap.insert sid sub)-    let unsubscribe =+    let deregister =             atomically $                 modifyTVar' (subscribers pub) (IntMap.delete sid)-    pure (queue, status, unsubscribe)+    pure (PublisherSubscription queue status dropped deregister)  -- | Read the last-published global position. publisherPosition :: EventPublisher -> STM GlobalPosition@@ -322,6 +343,7 @@                 DropSubscription -> writeTVar (subStatus sub) Overflowed                 DropOldest -> do                     _ <- tryReadTBQueue (subQueue sub)+                    modifyTVar' (subDropped sub) (+ 1)                     writeTBQueue (subQueue sub) events  -- Wait for either a tick or a 30-second timeout (safety poll).
src/Kiroku/Store/Subscription/Types.hs view
@@ -27,6 +27,19 @@     SubscriptionCheckpointMissing (..),     SubscriptionCheckpoint (..),     SubscriptionCheckpointInventory (..),++    -- * Dead-letter inspection+    SubscriptionDeadLetter (..),+    SubscriptionDeadLetterCursor (..),+    subscriptionDeadLetterCursor,+    SubscriptionDeadLetterLimit,+    SubscriptionDeadLetterLimitOutOfRange (..),+    mkSubscriptionDeadLetterLimit,+    subscriptionDeadLetterLimitValue,+    defaultSubscriptionDeadLetterLimit,+    SubscriptionDeadLetterQuery (..),+    defaultSubscriptionDeadLetterQuery,+    SubscriptionDeadLetterPage (..),     SubscriptionTarget (..),     TargetBindingPolicy (..),     SubscriptionTargetMismatch (..),@@ -77,7 +90,8 @@ ) where  import Control.Exception (Exception (..), SomeException)-import Data.Int (Int32)+import Data.Aeson (Value)+import Data.Int (Int32, Int64) import Data.Set (Set) import Data.Set qualified as Set import Data.Text (Text)@@ -95,7 +109,7 @@     deadLetterSummary,     retryDelayMicros,  )-import Kiroku.Store.Types (CategoryName, EventType, GlobalPosition, RecordedEvent (..))+import Kiroku.Store.Types (CategoryName, EventId, EventType, GlobalPosition, RecordedEvent (..)) import Numeric.Natural (Natural)  {- | A declarative, closed filter over event types for a subscription.@@ -235,6 +249,69 @@ data SubscriptionCheckpointInventory = SubscriptionCheckpointInventory     { storePosition :: !GlobalPosition     , checkpoints :: !(Vector SubscriptionCheckpoint)+    }+    deriving stock (Eq, Show, Generic)++-- | A durably parked event. The worker's structured reason JSON is unchanged.+data SubscriptionDeadLetter = SubscriptionDeadLetter+    { deadLetterId :: !Int64+    , subscriptionName :: !SubscriptionName+    , consumerGroupMember :: !Int32+    , globalPosition :: !GlobalPosition+    , eventId :: !EventId+    , reason :: !Value+    , reasonSummary :: !Text+    , attemptCount :: !Int32+    , createdAt :: !UTCTime+    }+    deriving stock (Eq, Show, Generic)++{- | Exclusive cursor in descending (global position, dead-letter id) order.+It remains usable after the originating row is deleted.+-}+data SubscriptionDeadLetterCursor = SubscriptionDeadLetterCursor+    { cursorGlobalPosition :: !GlobalPosition+    , cursorDeadLetterId :: !Int64+    }+    deriving stock (Eq, Ord, Show, Generic)++subscriptionDeadLetterCursor :: SubscriptionDeadLetter -> SubscriptionDeadLetterCursor+subscriptionDeadLetterCursor SubscriptionDeadLetter{globalPosition = position, deadLetterId = ident} =+    SubscriptionDeadLetterCursor position ident++{- | Validated page size, 1 through 1,000. Use the smart constructor.+No Generic or numeric instance exposes a construction bypass.+-}+newtype SubscriptionDeadLetterLimit = SubscriptionDeadLetterLimit Int32+    deriving stock (Eq, Ord, Show)++newtype SubscriptionDeadLetterLimitOutOfRange = SubscriptionDeadLetterLimitOutOfRange Int32+    deriving stock (Eq, Show)+mkSubscriptionDeadLetterLimit :: Int32 -> Either SubscriptionDeadLetterLimitOutOfRange SubscriptionDeadLetterLimit+mkSubscriptionDeadLetterLimit value+    | value < 1 || value > 1000 = Left (SubscriptionDeadLetterLimitOutOfRange value)+    | otherwise = Right (SubscriptionDeadLetterLimit value)+subscriptionDeadLetterLimitValue :: SubscriptionDeadLetterLimit -> Int32+subscriptionDeadLetterLimitValue (SubscriptionDeadLetterLimit value) = value+defaultSubscriptionDeadLetterLimit :: SubscriptionDeadLetterLimit+defaultSubscriptionDeadLetterLimit = SubscriptionDeadLetterLimit 100++-- | Nothing selects all historical members, or starts at the newest row.+data SubscriptionDeadLetterQuery = SubscriptionDeadLetterQuery+    { subscriptionName :: !SubscriptionName+    , consumerGroupMember :: !(Maybe Int32)+    , after :: !(Maybe SubscriptionDeadLetterCursor)+    , limit :: !SubscriptionDeadLetterLimit+    }+    deriving stock (Eq, Show, Generic)++defaultSubscriptionDeadLetterQuery :: SubscriptionName -> SubscriptionDeadLetterQuery+defaultSubscriptionDeadLetterQuery name = SubscriptionDeadLetterQuery name Nothing Nothing defaultSubscriptionDeadLetterLimit++-- | Trimmed page; a cursor is present exactly when another row exists.+data SubscriptionDeadLetterPage = SubscriptionDeadLetterPage+    { deadLetters :: !(Vector SubscriptionDeadLetter)+    , nextCursor :: !(Maybe SubscriptionDeadLetterCursor)     }     deriving stock (Eq, Show, Generic) 
src/Kiroku/Store/Subscription/Worker.hs view
@@ -75,7 +75,7 @@     stateCursor,     step,  )-import Kiroku.Store.Subscription.Types+import Kiroku.Store.Subscription.Types hiding (eventId, globalPosition) import Kiroku.Store.Types (CategoryName (..), EventId (..), GlobalPosition (..), RecordedEvent (..)) import System.IO.Unsafe (unsafePerformIO) 
src/Kiroku/Store/Types.hs view
@@ -15,10 +15,14 @@     categoryName,     streamNameInCategory,     EventFilter (..),+    BrowsePageSize,+    BrowsePageSizeError (..),+    mkBrowsePageSize,+    browsePageSizeValue, ) where  import Data.Aeson (Value)-import Data.Int (Int64)+import Data.Int (Int32, Int64) import Data.Text (Text) import Data.Text qualified as Text import Data.Time (UTCTime)@@ -337,3 +341,20 @@       -}       FilterCausationAncestors !EventId     deriving stock (Eq, Show, Generic)++{- | Validated catalog page size: 1–1001, including one HTTP over-fetch row.+The constructor is private; arithmetic on untrusted limits cannot overflow.+-}+newtype BrowsePageSize = BrowsePageSize Int32+    deriving stock (Eq, Show)++data BrowsePageSizeError = InvalidBrowsePageSize !Int+    deriving stock (Eq, Show)++mkBrowsePageSize :: Int -> Either BrowsePageSizeError BrowsePageSize+mkBrowsePageSize n+    | n >= 1 && n <= 1001 = Right (BrowsePageSize (fromIntegral n))+    | otherwise = Left (InvalidBrowsePageSize n)++browsePageSizeValue :: BrowsePageSize -> Int32+browsePageSizeValue (BrowsePageSize n) = n
test/Main.hs view
@@ -23,7 +23,11 @@ import Kiroku.Store import Kiroku.Store.Subscription.Effect qualified as SubEff import Kiroku.Store.Subscription.EventPublisher (publisherPosition)+import Kiroku.Store.Subscription.Fsm qualified as Fsm import Kiroku.Test.Postgres (withMigratedTestDatabase)+import Test.BrowseQueryPlans qualified as BrowseQueryPlans+import Test.BrowseReads qualified as BrowseReads+import Test.BrowseReadsMock qualified as BrowseReadsMock import Test.CatchupDbErrorNoPrematureSwitch qualified as CatchupDbErrorNoPrematureSwitch import Test.Category qualified as Category import Test.CategoryIdleNoSpin qualified as CategoryIdleNoSpin@@ -45,11 +49,15 @@ import Test.PerformanceStructure qualified as PerformanceStructure import Test.Properties qualified as Properties import Test.PublisherCallbackResilience qualified as PublisherCallbackResilience+import Test.PublisherDropCounter qualified as PublisherDropCounter import Test.PublisherIdleAdvance qualified as PublisherIdleAdvance import Test.PublisherRestartNoRebroadcast qualified as PublisherRestartNoRebroadcast import Test.ReadStream qualified as ReadStream import Test.StartupFailureSurfacing qualified as StartupFailureSurfacing import Test.StreamBridgeTermination qualified as StreamBridgeTermination+import Test.StreamHead qualified as StreamHead+import Test.StreamHeadIsolation qualified as StreamHeadIsolation+import Test.StreamHeadMock qualified as StreamHeadMock import Test.StreamHistoryGuard qualified as StreamHistoryGuard import Test.StreamNameLookup qualified as StreamNameLookup import Test.SubscriptionCheckpointInitialization qualified as SubscriptionCheckpointInitialization@@ -58,6 +66,8 @@ import Test.SubscriptionCheckpointInventoryMock qualified as SubscriptionCheckpointInventoryMock import Test.SubscriptionCheckpointReset qualified as SubscriptionCheckpointReset import Test.SubscriptionCheckpointWorker qualified as SubscriptionCheckpointWorker+import Test.SubscriptionDeadLetters qualified as SubscriptionDeadLetters+import Test.SubscriptionDeadLettersMock qualified as SubscriptionDeadLettersMock import Test.SubscriptionPauseResume qualified as SubscriptionPauseResume import Test.SubscriptionReconnect qualified as SubscriptionReconnect import Test.SubscriptionRegistry qualified as SubscriptionRegistry@@ -72,6 +82,10 @@  main :: IO () main = withSharedMigratedPostgres $ hspec $ do+    BrowseReads.spec+    BrowseReadsMock.spec+    SubscriptionDeadLetters.spec+    SubscriptionDeadLettersMock.spec     UniqueViolationMapping.spec     SubscriptionTarget.spec     Category.spec@@ -93,12 +107,15 @@     ConsumerGroup.spec     ConsumerGroupEffect.spec     describe "performance structure" $ do+        BrowseQueryPlans.spec+        StreamHeadIsolation.spec         PerformanceStructure.spec         NotifyGuard.spec         StreamNameLookup.noOpSpec     CategoryIdleNoSpin.spec     HandlerStall.spec     PublisherCallbackResilience.spec+    PublisherDropCounter.spec     PublisherIdleAdvance.spec     PublisherRestartNoRebroadcast.spec     CatchupDbErrorNoPrematureSwitch.spec@@ -115,6 +132,8 @@     SubscriptionRegistry.spec     SubscriptionRetryDeadLetter.spec     EventTypeFilter.spec+    StreamHead.spec+    StreamHeadMock.spec     VisibleGlobalHeadPosition.spec     VisibleGlobalHeadPositionMock.spec     around withTestStore $ do@@ -1794,6 +1813,14 @@                             , selector = Nothing                             }                 handle <- subscribe store cfg+                -- This fixture requires a live handler; catch-up could otherwise+                -- consume the queued range before observing overflow.+                let awaitLive = do+                        st <- currentState handle+                        case st of+                            Just Fsm.Live{} -> pure ()+                            _ -> threadDelay 1_000 >> awaitLive+                Async.race (threadDelay 5_000_000) awaitLive >>= (`shouldBe` Right ())                 -- First append: triggers handler, which blocks on the release MVar.                 Right _ <- runStoreIO store $ appendToStream (StreamName "f6-1") NoStream [makeEvent "E1" (Aeson.object [])]                 takeMVar firstSeen
+ test/Test/BrowseQueryPlans.hs view
@@ -0,0 +1,109 @@+module Test.BrowseQueryPlans (spec) where++import Control.Exception (bracket_)+import Control.Lens ((^.))+import Control.Monad (forM_)+import Data.Aeson (Value (..))+import Data.Aeson qualified as Aeson+import Data.Aeson.KeyMap qualified as KM+import Data.ByteString (ByteString)+import Data.Generics.Labels ()+import Data.Text (Text)+import Data.Text qualified as T+import Data.Vector qualified as V+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, unpreparable)+import Hasql.Statement qualified as Statement+import Kiroku.Store+import Kiroku.Store.SQL qualified as SQL+import Kiroku.Test.Postgres (migrateTestDatabase, withMigratedTestDatabase)+import Test.Helpers (withTestStore)+import Test.Hspec++spec :: Spec+spec = describe "browse reads prepared query work" $ do+    around withTestStore $ it "bounds production generic and custom index seeks on a 40000-stream catalog" preparedChecks+    it "uses the same bounded production statements in an English ICU database" $ withIcuStore $ \store -> do+        preparedChecks store+        Right _ <- runStoreIO store $ appendToStream (StreamName "orders-éclair") NoStream [EventData Nothing (EventType "Created") (Aeson.object []) Nothing Nothing Nothing]+        Right _ <- runStoreIO store $ appendToStream (StreamName "orders-éclair") NoStream [EventData Nothing (EventType "Created") (Aeson.object []) Nothing Nothing Nothing]+        let page = either (error . show) (\size -> size) (mkBrowsePageSize 11)+        Right rows <- runStoreIO store $ listStreams (Just (CategoryName "orders")) (Just "orders-é") Nothing page+        map (\row -> row ^. #name) (V.toList rows) `shouldBe` [StreamName "orders-éclair"]++preparedChecks :: KirokuStore -> IO ()+preparedChecks store = do+    seeded <-+        Pool.use (store ^. #pool) $+            Session.script+                "INSERT INTO streams(stream_name) SELECT c || '-' || lpad(n::text,10,'0') FROM unnest(ARRAY['orders','noise']) c CROSS JOIN generate_series(1,20000) n; ANALYZE streams;"+    seeded `shouldBe` Right ()+    forM_ ["force_generic_plan", "force_custom_plan"] $ \mode -> do+        forM_ [(SQL.InclusiveLower, "'orders-', 'orders.', 11"), (SQL.ExclusiveLower, "'orders-0000019990', 'orders.', 11"), (SQL.InclusiveLower, "'absent', 'absenu', 11"), (SQL.InclusiveLower, "'', NULL, 11")] $ \(bound, args) -> do+            plan <- explainPrepared store mode "text,text,int4" args (SQL.listStreamsRangeStmt bound)+            expectBounded plan 11+        plan <- explainPrepared store mode "text,text,text,text,int4" "'orders', NULL, 'orders-', 'orders.', 11" (SQL.listStreamsPairStmt SQL.ExactName SQL.InclusiveLower)+        expectBounded plan 22+        forM_ [(False, "NULL, 11"), (True, "'orders', 11")] $ \(hasCursor, args) -> do+            categoryPlan <- explainPrepared store mode "text,int4" args (SQL.listCategoriesStmt hasCursor)+            indexNames categoryPlan `shouldContain` ["ix_streams_category"]+            -- Each loose-index seek returns at most one distinct category.+            scannedRows categoryPlan `shouldSatisfy` (<= 12)++explainPrepared :: KirokuStore -> Text -> Text -> Text -> Statement p r -> IO Value+explainPrepared store mode types args statement = do+    result <- Pool.use (store ^. #pool) $ do+        Session.script ("SET plan_cache_mode=" <> mode <> "; PREPARE mp13_browse(" <> types <> ") AS " <> Statement.toSql statement)+        bytes <- Session.statement () (unpreparable ("EXPLAIN (ANALYZE, BUFFERS, TIMING OFF, FORMAT JSON) EXECUTE mp13_browse(" <> args <> ")") E.noParams (D.singleRow (D.column (D.nonNullable (D.jsonBytes Right)))))+        Session.script "DEALLOCATE mp13_browse; RESET plan_cache_mode"+        pure bytes+    case result of+        Left err -> expectationFailure (show err) >> fail "EXPLAIN failed"+        Right bytes -> decodePlan bytes+  where+    decodePlan :: ByteString -> IO Value+    decodePlan bytes = either fail pure (Aeson.eitherDecodeStrict' bytes)++expectBounded :: Value -> Double -> Expectation+expectBounded plan budget = do+    indexNames plan `shouldContain` ["ix_streams_browse_name"]+    scannedRows plan `shouldSatisfy` (<= budget)+    case plan of+        Array entries+            | Object entry <- V.head entries+            , Just (Object top) <- KM.lookup "Plan" entry ->+                number "Shared Hit Blocks" top + number "Shared Read Blocks" top `shouldSatisfy` (<= 64)+        _ -> expectationFailure (show plan)++-- Count rows examined at relation scans (including filters), not only results.+scannedRows :: Value -> Double+scannedRows (Object fields) = current + sum (map scannedRows (KM.elems fields))+  where+    current = case KM.lookup "Relation Name" fields of+        Just (String "streams") -> (number "Actual Rows" fields + number "Rows Removed by Filter" fields + number "Rows Removed by Index Recheck" fields) * number "Actual Loops" fields+        _ -> 0+scannedRows (Array values) = sum (map scannedRows (foldr (:) [] values))+scannedRows _ = 0+indexNames :: Value -> [Text]+indexNames (Object fields) = [name | Just (String name) <- [KM.lookup "Index Name" fields]] <> concatMap indexNames (KM.elems fields)+indexNames (Array values) = concatMap indexNames (foldr (:) [] values)+indexNames _ = []+number :: Aeson.Key -> Aeson.Object -> Double+number key fields = case KM.lookup key fields of Just (Number n) -> realToFrac n; _ -> 0++-- A separately migrated database verifies production SQL under locale ordering,+-- rather than altering existing columns or using a different hand-written query.+withIcuStore :: (KirokuStore -> IO ()) -> IO ()+withIcuStore action = withMigratedTestDatabase $ \connection ->+    withStore (defaultConnectionSettings connection) $ \admin -> do+        let sql command = Pool.use (admin ^. #pool) (Session.script command) >>= either (fail . show) pure+            newConnection = T.unwords (filter (not . T.isPrefixOf "dbname=") (T.words connection) <> ["dbname=mp13_browse_icu"])+        bracket_+            (sql "CREATE DATABASE mp13_browse_icu TEMPLATE template0 LOCALE_PROVIDER icu ICU_LOCALE 'en'")+            (sql "DROP DATABASE mp13_browse_icu")+            $ do+                migrateTestDatabase newConnection+                withStore (defaultConnectionSettings newConnection) action
+ test/Test/BrowseReads.hs view
@@ -0,0 +1,94 @@+module Test.BrowseReads (spec) where++import Control.Lens ((&), (.~), (^.))+import Control.Monad (forM_, void)+import Data.Aeson qualified as Aeson+import Data.Char (chr)+import Data.Generics.Labels ()+import Data.List (sort)+import Data.Text (Text)+import Data.Text qualified as T+import Data.UUID qualified as UUID+import Data.Vector qualified as V+import Kiroku.Store+import Test.Helpers (makeEvent, withTestStore, withTestStoreSettings)+import Test.Hspec++page :: Int -> BrowsePageSize+page n = either (error . show) (\p -> p) (mkBrowsePageSize n)++names :: V.Vector StreamInfo -> [Text]+names = map (\s -> let StreamName n = s ^. #name in n) . V.toList++seed :: KirokuStore -> [Text] -> IO ()+seed store = mapM_ $ \name -> do+    result <- runStoreIO store $ appendToStream (StreamName name) NoStream [makeEvent "Created" (Aeson.object [])]+    result `shouldSatisfy` either (const False) (const True)++spec :: Spec+spec = describe "browse reads" $ do+    it "validates bounded page sizes" $ do+        forM_ [minBound, -1, 0, 1002, maxBound] $ \n -> mkBrowsePageSize n `shouldBe` Left (InvalidBrowsePageSize n)+        forM_ [1, 1000, 1001] $ \n -> fmap browsePageSizeValue (mkBrowsePageSize n) `shouldBe` Right (fromIntegral n)+    around withTestStore $ do+        it "excludes only the reserved all row and pages in UTF-8 byte order" $ \store -> do+            runStoreIO store (listStreams Nothing Nothing Nothing (page 10)) `shouldReturn` Right V.empty+            let values = ["orders-2", "orders", "orders-1", "$all-x", "singleton", "é-1", "Z-1", "-empty-category"]+            seed store values+            Right first <- runStoreIO store $ listStreams Nothing Nothing Nothing (page 3)+            names first `shouldBe` take 3 (sort values)+            Right rest <- runStoreIO store $ listStreams Nothing Nothing (Just (StreamName (last (names first)))) (page 100)+            names first <> names rest `shouldBe` sort values+            runStoreIO store (listStreams Nothing Nothing (Just (StreamName (last (sort values)))) (page 10)) `shouldReturn` Right V.empty+        it "intersects exact categories, literal prefixes and exclusive cursors" $ \store -> do+            let maxScalar = T.singleton (chr 0x10ffff)+                values = ["orders", "orders-1", "orders-2", "orders:other-1", "orders_-1", "orders%-1", "$all-x", "-a", "é-a", "é-a", maxScalar, maxScalar <> "-a", "\xD7FF-a", "\xE000-a"]+            seed store values+            forM_ [Nothing, Just (CategoryName "orders"), Just (CategoryName ""), Just (CategoryName "$all"), Just (CategoryName "not-a-category"), Just (CategoryName maxScalar)] $ \category ->+                forM_ [Nothing, Just "", Just "orders", Just "orders-", Just "orders%", Just "orders_", Just "absent", Just "é", Just maxScalar, Just "\xD7FF"] $ \prefix ->+                    forM_ [Nothing, Just (StreamName "orders"), Just (StreamName "orders-1"), Just (StreamName "zzz")] $ \cursor -> do+                        let expected = sort [v | v <- values, maybe True (== categoryName (StreamName v)) category, maybe True (`T.isPrefixOf` v) prefix, maybe True (\(StreamName c) -> v > c) cursor]+                        Right actual <- runStoreIO store $ listStreams category prefix cursor (page 100)+                        names actual `shouldBe` expected+        it "enumerates distinct categories including bare and empty categories" $ \store -> do+            seed store ["orders", "orders-1", "orders-2", "$all-x", "-a", "singleton"]+            Right first <- runStoreIO store $ listCategories Nothing (page 2)+            V.toList first `shouldBe` [CategoryName "", CategoryName "$all"]+            Right rest <- runStoreIO store $ listCategories (Just (V.last first)) (page 10)+            V.toList rest `shouldBe` [CategoryName "orders", CategoryName "singleton"]+        it "keeps soft-deleted summaries and removes hard-deleted summaries" $ \store -> do+            seed store ["orders-soft", "orders-hard"]+            void $ runStoreIO store $ softDeleteStream (StreamName "orders-soft")+            void $ runStoreIO store $ setStreamTruncateBefore (StreamName "orders-soft") (StreamVersion 1)+            void $ runStoreIO store $ hardDeleteStream (StreamName "orders-hard")+            Right rows <- runStoreIO store $ listStreams Nothing Nothing Nothing (page 10)+            names rows `shouldBe` ["orders-soft"]+            (V.head rows ^. #deletedAt) `shouldSatisfy` maybe False (const True)+        it "returns the canonical global row once even when linked, then nothing after hard deletion" $ \store -> do+            seed store ["orders-1"]+            Right events <- runStoreIO store $ readAllForward (GlobalPosition 0) 10+            let event = V.head events+            Right _ <- runStoreIO store $ linkToStream (StreamName "links-1") [event ^. #eventId]+            runStoreIO store (getEvent (event ^. #eventId)) `shouldReturn` Right (Just event)+            runStoreIO store (getEvent (EventId UUID.nil)) `shouldReturn` Right Nothing+            Right _ <- runStoreIO store $ hardDeleteStream (StreamName "orders-1")+            runStoreIO store (getEvent (event ^. #eventId)) `shouldReturn` Right Nothing+    it "applies the configured decode hook to event lookup" $ do+        let marker = Aeson.object ["decoded" Aeson..= True]+            hook event = pure (Right (event & #metadata .~ Just marker))+        withTestStoreSettings (\settings -> settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just hook}) $ \store -> do+            seed store ["orders-1"]+            Right events <- runStoreIO store $ readAllForward (GlobalPosition 0) 10+            Right (Just event) <- runStoreIO store $ getEvent (V.head events ^. #eventId)+            (event ^. #metadata) `shouldBe` Just marker++    it "propagates typed decode failure from event lookup" $ do+        let hook event = pure (Left (DecodeFailure (event ^. #eventId) "failure"))+        withTestStoreSettings (\settings -> settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just hook}) $ \store -> do+            seed store ["orders-1"]+            Right rows <- runStoreIO store $ readStreamForward (StreamName "missing") (StreamVersion 0) 1+            rows `shouldBe` V.empty+            -- Use the append-provided id to avoid decoding a fixture through the failing hook.+            let event = makeEvent "Created" (Aeson.object []) & #eventId .~ Just (EventId UUID.nil)+            Right _ <- runStoreIO store $ appendToStream (StreamName "orders-2") NoStream [event]+            runStoreIO store (getEvent (EventId UUID.nil)) `shouldReturn` Left (EventDecodeFailed (DecodeFailure (EventId UUID.nil) "failure"))
+ test/Test/BrowseReadsMock.hs view
@@ -0,0 +1,36 @@+module Test.BrowseReadsMock (spec) where++import Control.Monad.IO.Class (liftIO)+import Data.IORef (modifyIORef', newIORef, readIORef)+import Data.UUID qualified as UUID+import Data.Vector qualified as V+import Effectful (runEff)+import Effectful.Dispatch.Dynamic (interpret_)+import Kiroku.Store.Effect (Store (..))+import Kiroku.Store.Read+import Kiroku.Store.Types+import Test.Hspec++spec :: Spec+spec = describe "browse reads mock" $ it "dispatches each public wrapper exactly once with its parameters" $ do+    calls <- newIORef ([] :: [String])+    let page = either (error . show) (\size -> size) (mkBrowsePageSize 7)+        runner = interpret_ $ \case+            ListStreams category prefix cursor size -> do+                liftIO $ (category, prefix, cursor, size) `shouldBe` (Just (CategoryName "orders"), Just "orders-", Just (StreamName "orders-1"), page)+                liftIO $ modifyIORef' calls (<> ["streams"])+                pure V.empty+            ListCategories cursor size -> do+                liftIO $ (cursor, size) `shouldBe` (Just (CategoryName "orders"), page)+                liftIO $ modifyIORef' calls (<> ["categories"])+                pure V.empty+            GetEvent eid -> do+                liftIO $ eid `shouldBe` EventId UUID.nil+                liftIO $ modifyIORef' calls (<> ["event"])+                pure Nothing+            _ -> error "unexpected browse mock operation"+    _ <- runEff $ runner $ do+        _ <- listStreams (Just (CategoryName "orders")) (Just "orders-") (Just (StreamName "orders-1")) page+        _ <- listCategories (Just (CategoryName "orders")) page+        getEvent (EventId UUID.nil)+    readIORef calls `shouldReturn` ["streams", "categories", "event"]
+ test/Test/DeadLetterQueryPlans.hs view
@@ -0,0 +1,98 @@+{-# LANGUAGE MultilineStrings #-}+{-# LANGUAGE OverloadedLabels #-}++module Test.DeadLetterQueryPlans (spec) where++import Control.Lens ((^.))+import Control.Monad (forM_)+import Data.Aeson (Value (..))+import Data.Aeson qualified as Aeson+import Data.Aeson.KeyMap qualified as KM+import Data.ByteString.Lazy.Char8 qualified as LBS+import Data.Generics.Labels ()+import Data.Text (Text)+import Data.Text qualified+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, unpreparable)+import Hasql.Statement qualified as Statement+import Kiroku.Store+import Kiroku.Store.SQL qualified as SQL+import Test.Helpers (withTestStore)+import Test.Hspec++spec :: Spec+spec = describe "dead-letter production query plans" $ around withTestStore $ it "bounds historical-member enumeration and first/later pages in generic/custom plans as history grows" $ \store -> do+    execute store "INSERT INTO events(event_id,event_type,data,created_at) VALUES ('00000000-0000-0000-0000-000000000001','E','{}',now())"+    forM_ ["1000", "20000"] $ \history -> do+        execute store ("INSERT INTO dead_letters(subscription_name,consumer_group_member,global_position,event_id,reason,reason_summary,attempt_count) SELECT 'plans',m,n,'00000000-0000-0000-0000-000000000001','{}','fixture',1 FROM unnest(ARRAY[0,1,9]) m CROSS JOIN generate_series(1," <> history <> ") n ON CONFLICT DO NOTHING; ANALYZE dead_letters")+        forM_ ["force_generic_plan", "force_custom_plan"] $ \mode -> do+            memberFirst <- explain store mode "text,int4,int4" "'plans',9,6" SQL.listSubscriptionMemberDeadLettersFromStartStmt+            memberLater <- explain store mode "text,int4,int8,int8,int4" "'plans',9,500,9223372036854775807,6" SQL.listSubscriptionMemberDeadLettersStmt+            allFirst <- explain store mode "text,int4" "'plans',6" SQL.listSubscriptionDeadLettersFromStartStmt+            allLater <- explain store mode "text,int8,int8,int4" "'plans',500,9223372036854775807,6" SQL.listSubscriptionDeadLettersStmt+            forM_ [memberFirst, memberLater] $ \plan -> do+                expect ("index" :: Text) plan ("ix_dead_letters_subscription_position" `elem` texts "Index Name" plan)+                expect "member page work" plan (scanned plan <= 6)+                expect "no member sort" plan ("Sort" `notElem` texts "Node Type" plan)+            forM_ [allFirst, allLater] $ \plan -> do+                -- Three loose member probes plus three bounded lateral scans.+                expect "bounded member enumeration/page" plan (scanned plan <= 24)+                expect "bounded merge input" plan (sum (sortInputs plan) <= 18)+                -- The existing natural-key index has the same (name, member)+                -- prefix and can serve the loose member probes too.+                expect "existing ordered indexes only" plan (all (`elem` ["ix_dead_letters_subscription_position", "dead_letters_subscription_name_consumer_group_member_global_key"]) (texts "Index Name" plan))+                expect "ordered historical member advance" plan (any ("consumer_group_member >" `contains`) (texts "Index Cond" plan))+                expect "recency index for lateral pages" plan ("ix_dead_letters_subscription_position" `elem` texts "Index Name" plan)+                expect "no historical sequential scan" plan ("Seq Scan" `notElem` texts "Node Type" plan)+            forM_ [memberLater, allLater] $ \plan ->+                expect "tuple cursor in index condition" plan (any (\condition -> "ROW(global_position, dead_letter_id)" `contains` condition) (texts "Index Cond" plan))+            forM_ [memberFirst, memberLater, allFirst, allLater] $ \plan ->+                expect "bounded shared buffers" plan (buffers plan <= 128)+            -- Preserve the actual production EXPLAIN evidence in the test log.+            mapM_ (LBS.putStrLn . Aeson.encode) [memberFirst, memberLater, allFirst, allLater]+  where+    contains = Data.Text.isInfixOf++execute :: KirokuStore -> Text -> IO ()+execute store sql = Pool.use (store ^. #pool) (Session.script sql) >>= either (fail . show) pure+explain :: KirokuStore -> Text -> Text -> Text -> Statement p r -> IO Value+explain store mode types args statement = do+    Right bytes <- Pool.use (store ^. #pool) $ do+        Session.script ("SET plan_cache_mode=" <> mode <> "; PREPARE mp13_deadletters(" <> types <> ") AS " <> Statement.toSql statement)+        bytes <- Session.statement () (unpreparable ("EXPLAIN (ANALYZE,BUFFERS,TIMING OFF,FORMAT JSON) EXECUTE mp13_deadletters(" <> args <> ")") E.noParams (D.singleRow (D.column (D.nonNullable (D.jsonBytes Right)))))+        Session.script "DEALLOCATE mp13_deadletters; RESET plan_cache_mode"+        pure bytes+    either fail pure (Aeson.eitherDecodeStrict' bytes)+expect :: Text -> Value -> Bool -> Expectation+expect label plan ok = if ok then pure () else expectationFailure (show label <> ": " <> show plan)+texts :: Aeson.Key -> Value -> [Text]+texts key (Object fields) = [s | Just (String s) <- [KM.lookup key fields]] <> concatMap (texts key) (KM.elems fields)+texts key (Array values) = concatMap (texts key) values+texts _ _ = []+number :: Aeson.Key -> Aeson.Object -> Double+number key fields = case KM.lookup key fields of Just (Number n) -> realToFrac n; _ -> 0+scanned :: Value -> Double+scanned (Object fields) = current + sum (map scanned (KM.elems fields))+  where+    current = case KM.lookup "Relation Name" fields of+        Just (String "dead_letters") -> (number "Actual Rows" fields + number "Rows Removed by Filter" fields + number "Rows Removed by Index Recheck" fields) * number "Actual Loops" fields+        _ -> 0+scanned (Array values) = sum (map scanned (foldr (:) [] values))+scanned _ = 0+sortInputs :: Value -> [Double]+sortInputs (Object fields) = current <> concatMap sortInputs (KM.elems fields)+  where+    current = case (KM.lookup "Node Type" fields, KM.lookup "Plans" fields) of+        (Just (String "Sort"), Just (Array children)) -> [number "Actual Rows" child * number "Actual Loops" child | Object child <- foldr (:) [] children]+        _ -> []+sortInputs (Array values) = concatMap sortInputs values+sortInputs _ = []++buffers :: Value -> Double+buffers (Array values) = case foldr (:) [] values of+    Object entry : _ | Just (Object top) <- KM.lookup "Plan" entry -> number "Shared Hit Blocks" top + number "Shared Read Blocks" top+    _ -> 1 / 0+buffers _ = 1 / 0
test/Test/Helpers.hs view
@@ -31,6 +31,7 @@     countEvents,     countDeadLettersForEvents,     insertDeadLetterForEvent,+    insertDeadLetterWith,     insertEventUsingDefaultId,     serverVersionNum,     truncateRejected,@@ -131,21 +132,24 @@  -- | Insert a dead-letter row for a recorded event without running a subscription. insertDeadLetterForEvent :: KirokuStore -> Text -> RecordedEvent -> IO ()-insertDeadLetterForEvent store subscriptionName event = do+insertDeadLetterForEvent store subscriptionName = insertDeadLetterWith store subscriptionName 0 (Aeson.object [("source", Aeson.String "test")]) "test dead letter" 1++insertDeadLetterWith :: KirokuStore -> Text -> Int32 -> Aeson.Value -> Text -> Int32 -> RecordedEvent -> IO ()+insertDeadLetterWith store subscriptionName selectedMember reason summary attempts event = do     let EventId eid = event ^. #eventId         GlobalPosition globalPosition = event ^. #globalPosition         params =             SQL.DeadLetterParams                 { SQL.dlSubscriptionName = subscriptionName-                , SQL.dlMember = 0+                , SQL.dlMember = selectedMember                 , SQL.dlTargetKind = "unbound"                 , SQL.dlTargetCategory = Nothing-                , SQL.dlGroupSize = 1+                , SQL.dlGroupSize = max 1 (selectedMember + 1)                 , SQL.dlGlobalPosition = globalPosition                 , SQL.dlEventId = eid-                , SQL.dlReason = Aeson.object [("source", Aeson.String "test")]-                , SQL.dlReasonSummary = "test dead letter"-                , SQL.dlAttemptCount = 1+                , SQL.dlReason = reason+                , SQL.dlReasonSummary = summary+                , SQL.dlAttemptCount = attempts                 }     result <- Pool.use (store ^. #pool) (Session.statement params SQL.insertDeadLetterAndCheckpointStmt)     case result of
test/Test/PerformanceStructure.hs view
@@ -26,6 +26,8 @@ import Kiroku.Store.Subscription.Stream qualified as Buffer import Kiroku.Store.Subscription.Worker (withLoadCheckpointHookForTest) import Kiroku.Test.Fixtures.CategoryScaling (categoryScalingFixtureSql, categoryScalingHead)+import Kiroku.Test.Fixtures.StreamHead (streamHeadFixtureSql)+import Test.DeadLetterQueryPlans qualified as DeadLetterQueryPlans import Test.Helpers (makeEvent, validConsumerGroup, waitWithTimeout, withTestStore, withTestStoreSettings) import Test.Hspec @@ -33,11 +35,21 @@ spec = do     noOpAppendSpec     queryPlanSpec+    DeadLetterQueryPlans.spec     categoryReadCostSpec+    streamHeadQuerySpec  noOpAppendSpec :: Spec noOpAppendSpec =     describe "no-op paths use no pooled connection" $ do+        it "rejects invalid dead-letter page sizes before pool checkout" $ do+            checkouts <- newIORef (0 :: Int)+            withObservedStore checkouts $ \_ -> do+                before <- readIORef checkouts+                forM_ [0, -1, 1001] $ \n -> mkSubscriptionDeadLetterLimit n `shouldBe` Left (SubscriptionDeadLetterLimitOutOfRange n)+                after <- readIORef checkouts+                after - before `shouldBe` 0+         it "rejects invalid batch and buffer sizes before pool checkout" $ do             checkouts <- newIORef (0 :: Int)             withObservedStore checkouts $ \_ -> do@@ -558,3 +570,81 @@                 <*> D.column (D.nonNullable D.int8)             )         )++streamHeadQuerySpec :: Spec+streamHeadQuerySpec = describe "stream head query work" $+    aroundAll withStreamHeadStore $+        forM_ ["bench-1", "long-1", "empty-1", "missing-1"] $ \name ->+            it ("bounds literal and warmed prepared probes for " <> T.unpack name) $ \store -> do+                let literal = "'" <> name <> "'"+                    prefix = "EXPLAIN (ANALYZE, BUFFERS, TIMING OFF, FORMAT JSON) "+                direct <- explainWith prefix store SQL.getStreamWithHeadStmt [("$1", literal)]+                prepared <- explainStreamHeadPrepared store literal+                forM_ [("literal", direct), ("prepared", prepared)] $ \(kind, plan) -> do+                    -- Emit complete JSON into retained test logs, including buffers and loops.+                    putStrLn ("STREAM_HEAD_PLAN " <> T.unpack name <> " " <> kind <> " " <> show (Aeson.encode plan))+                    expectIndex "ix_stream_events_all_by_origin" plan+                    expectNoNodeType "Sort" plan+                    let nodes = planObjects plan+                        probes = filter ((== Just "stream_events") . textField "Relation Name") nodes+                        limits = filter ((== Just "Limit") . textField "Node Type") nodes+                    length probes `shouldBe` 1+                    length limits `shouldBe` 1+                    forM_ probes $ \probe -> do+                        textField "Node Type" probe `shouldSatisfy` (`elem` [Just "Index Scan", Just "Index Only Scan"])+                        textField "Scan Direction" probe `shouldBe` Just "Backward"+                        planNumber "Actual Rows" probe `shouldSatisfy` (<= 1)+                        if name == "missing-1" then planNumber "Actual Loops" probe `shouldBe` 0 else pure ()+                    forM_ limits $ \limit -> do+                        expectIndex "ix_stream_events_all_by_origin" (Object limit)+                        planNumber "Actual Rows" limit `shouldSatisfy` (<= 1)+                    map (textField "Relation Name") nodes `shouldNotContain` [Just "events"]+                    case plan of+                        Array entries+                            | Object entry : _ <- foldr (:) [] entries+                            , Just (Object top) <- KeyMap.lookup "Plan" entry -> do+                                planNumber "Shared Hit Blocks" top + planNumber "Shared Read Blocks" top `shouldSatisfy` (<= 32)+                                planNumber "Actual Rows" top `shouldBe` (if name == "missing-1" then 0 else 1)+                        _ -> expectationFailure (show plan)+                actual <- runStoreIO store (getStreamWithHead (StreamName name))+                case name of+                    "missing-1" -> actual `shouldBe` Right Nothing+                    "empty-1" -> case actual of+                        Right (Just (_, headPosition)) -> headPosition `shouldBe` Nothing+                        _ -> expectationFailure (show actual)+                    _ -> case actual of+                        Right (Just (_, Just _)) -> pure ()+                        _ -> expectationFailure (show actual)++withStreamHeadStore :: (KirokuStore -> IO ()) -> IO ()+withStreamHeadStore action = withTestStore $ \store -> do+    forM_ [streamHeadFixtureSql, "VACUUM (ANALYZE) stream_events"] $ \command ->+        Pool.use (store ^. #pool) (Session.script command) >>= either (fail . show) pure+    action store++explainStreamHeadPrepared :: KirokuStore -> Text -> IO Value+explainStreamHeadPrepared store literal = do+    result <- Pool.use (store ^. #pool) $ do+        Session.script ("PREPARE ep97_head(text) AS " <> Statement.toSql SQL.getStreamWithHeadStmt)+        -- Default auto policy, past PostgreSQL's first five custom executions.+        Session.script (T.replicate 6 ("EXECUTE ep97_head(" <> literal <> ");"))+        bytes <-+            Session.statement () $+                unpreparable+                    ("EXPLAIN (ANALYZE, BUFFERS, TIMING OFF, FORMAT JSON) EXECUTE ep97_head(" <> literal <> ")")+                    E.noParams+                    (D.singleRow (D.column (D.nonNullable (D.jsonBytes Right))))+        Session.script "DEALLOCATE ep97_head"+        pure bytes+    bytes <- either (fail . show) pure result+    either fail pure (Aeson.eitherDecodeStrict' bytes)++planObjects :: Value -> [Aeson.Object]+planObjects (Object fields) = fields : concatMap planObjects (KeyMap.elems fields)+planObjects (Array values) = concatMap planObjects (foldr (:) [] values)+planObjects _ = []++planNumber :: Aeson.Key -> Aeson.Object -> Double+planNumber key fields = case KeyMap.lookup key fields of+    Just (Number n) -> realToFrac n+    _ -> 0
+ test/Test/PublisherDropCounter.hs view
@@ -0,0 +1,59 @@+{-# LANGUAGE OverloadedRecordDot #-}++module Test.PublisherDropCounter (spec) where++import Control.Concurrent.STM+import Control.Exception (bracket)+import Data.Aeson (object)+import Data.IntMap.Strict qualified as IntMap+import Data.Text (Text)+import Data.Vector qualified as V+import Kiroku.Store+import Kiroku.Store.Subscription.EventPublisher qualified as Pub+import System.Timeout (timeout)+import Test.Helpers (makeEvent, waitForPublisher, withTestStore)+import Test.Hspec++spec :: Spec+spec = describe "publisher drop counter" $ do+    it "counts drops, keeps the newest batch and preserves the idempotent wrapper" $ withTestStore $ \store -> do+        baseline <- IntMap.size <$> readTVarIO (Pub.subscribers store.publisher)+        bracket (atomically (Pub.subscribePublisherWith store.publisher 1 DropOldest)) Pub.unsubscribe $ \sub -> do+            appendAndWait store "drop-a" 1+            readTVarIO sub.subscriptionDropped `shouldReturn` 0+            appendAndWait store "drop-b" 2+            readTVarIO sub.subscriptionDropped `shouldReturn` 1+            readTVarIO sub.subscriptionStatus `shouldReturn` Pub.Active+            newest sub.subscriptionQueue `shouldReturn` [GlobalPosition 2]+            bracket (atomically (Pub.subscribePublisher store.publisher 1 DropOldest)) (\(_, _, stop) -> stop) $ \(queue, _, stop) -> do+                appendAndWait store "drop-c" 3+                appendAndWait store "drop-d" 4+                newest queue `shouldReturn` [GlobalPosition 4]+                stop >> stop+            Pub.unsubscribe sub >> Pub.unsubscribe sub+        (IntMap.size <$> readTVarIO (Pub.subscribers store.publisher)) `shouldReturn` baseline+    it "does not count PauseAndResume or DropSubscription overflow" $ withTestStore $ \store ->+        mapM_+            ( \policy -> bracket (atomically (Pub.subscribePublisherWith store.publisher 1 policy)) Pub.unsubscribe $ \sub -> do+                Right pos <- runStoreIO store (appendToStream (StreamName "policies-a") AnyVersion [makeEvent "E" (object [])])+                bounded (waitForPublisher store pos.globalPosition)+                Right pos2 <- runStoreIO store (appendToStream (StreamName "policies-a") AnyVersion [makeEvent "E" (object [])])+                bounded (waitForPublisher store pos2.globalPosition)+                readTVarIO sub.subscriptionDropped `shouldReturn` 0+                readTVarIO sub.subscriptionStatus `shouldReturn` (if policy == PauseAndResume then Pub.Paused else Pub.Overflowed)+            )+            [PauseAndResume, DropSubscription]++appendAndWait :: KirokuStore -> Text -> Int -> IO ()+appendAndWait store name n = do+    Right _ <- runStoreIO store (appendToStream (StreamName name) NoStream [makeEvent "E" (object [])])+    bounded (waitForPublisher store (GlobalPosition (fromIntegral n)))++bounded :: IO a -> IO a+bounded action = timeout 5_000_000 action >>= maybe (fail "publisher timeout") pure++newest :: TBQueue DecodedBatch -> IO [GlobalPosition]+newest queue = do+    Just (UnchangedBatch events) <- atomically (tryReadTBQueue queue)+    atomically (isEmptyTBQueue queue) `shouldReturn` True+    pure (map (.globalPosition) (V.toList events))
+ test/Test/StreamHead.hs view
@@ -0,0 +1,108 @@+module Test.StreamHead (spec) where++import Control.Concurrent.Async qualified as Async+import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)+import Control.Lens ((&), (.~), (^.))+import Control.Monad (forM_, replicateM_, void)+import Data.Aeson (Value (Null))+import Data.Generics.Labels ()+import Data.IORef (modifyIORef', newIORef, readIORef)+import Data.Text (Text)+import Data.Vector qualified as V+import Effectful (runEff)+import Effectful.Error.Static (runErrorNoCallStack)+import Hasql.Pool qualified as Pool+import Hasql.Session qualified as Session+import Kiroku.Store+import System.Timeout (timeout)+import Test.Helpers (makeEvent, withTestStore, withTestStoreSettings)+import Test.Hspec++spec :: Spec+spec = describe "stream head" $ do+    forM_ [False, True] $ \resource -> describe (if resource then "resource runner" else "direct runner") $ do+        let readHead store name =+                if resource+                    then runEff . runErrorNoCallStack @StoreError . runKirokuStoreWith store . runStoreResource $ getStreamWithHead name+                    else runStoreIO store (getStreamWithHead name)+            check store name version headPosition = do+                Right (Just metadata) <- runStoreIO store (getStream name)+                metadata ^. #version `shouldBe` StreamVersion version+                readHead store name `shouldReturn` Right (Just (metadata, headPosition))+        it "distinguishes absent, empty and reserved $all in empty and populated stores" $ withTestStore $ \store -> do+            readHead store (StreamName "missing") `shouldReturn` Right Nothing+            raw store "INSERT INTO streams(stream_name) VALUES ('empty')"+            check store (StreamName "empty") 0 Nothing+            check store (StreamName "$all") 0 Nothing+            appendOne store (StreamName "origin")+            check store (StreamName "$all") 1 Nothing+        it "captures interleaved heads, logical lifecycle and a fresh identity after hard deletion" $ withTestStore $ \store -> do+            let a = StreamName "a"; b = StreamName "b"+            mapM_ (appendOne store) [a, b, a, b]+            originated <- appendPosition store a+            originated `shouldBe` GlobalPosition 5+            check store a 3 (Just originated)+            appendOne store b+            check store a 3 (Just originated)+            Right (Just _) <- runStoreIO store (setStreamTruncateBefore a (StreamVersion 3))+            check store a 3 (Just originated)+            Right (Just _) <- runStoreIO store (softDeleteStream a)+            check store a 3 (Just originated)+            Right (Just (old, _)) <- readHead store a+            Right (Just _) <- runStoreIO store (hardDeleteStream a)+            readHead store a `shouldReturn` Right Nothing+            raw store "INSERT INTO streams(stream_name) VALUES ('a')"+            check store a 0 Nothing+            Right (Just (fresh, _)) <- readHead store a+            fresh ^. #id `shouldNotBe` (old ^. #id)+            recreated <- appendPosition store a+            recreated `shouldBe` GlobalPosition 7+            check store a 1 (Just recreated)+        it "ignores links while tracking subsequent originated appends" $ withTestStore $ \store -> do+            let mixed = StreamName "mixed"; source = StreamName "source"; linked = StreamName "linked"+            appendOne store mixed+            appendOne store source+            Right events <- runStoreIO store (readAllForward (GlobalPosition 0) 10)+            let eid = (events V.! 1) ^. #eventId+            Right _ <- runStoreIO store (linkToStream linked [eid])+            check store linked 1 Nothing+            Right _ <- runStoreIO store (linkToStream mixed [eid])+            check store mixed 2 (Just ((events V.! 0) ^. #globalPosition))+            next <- appendPosition store mixed+            next `shouldBe` GlobalPosition 3+            check store mixed 3 (Just next)+    it "never invokes the event decode hook" $ do+        calls <- newIORef (0 :: Int)+        let hook _ = modifyIORef' calls (+ 1) >> error "unexpected event decode"+        withTestStoreSettings (\settings -> settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just hook}) $ \store -> do+            expected <- appendPosition store (StreamName "hook")+            Right (Just (_, actual)) <- runStoreIO store (getStreamWithHead (StreamName "hook"))+            actual `shouldBe` Just expected+            readIORef calls `shouldReturn` 0+    it "observes version and head from one snapshot while appends commit" $ withTestStore $ \store -> do+        let name = StreamName "concurrent"+        appendOne store name+        start <- newEmptyMVar+        let writer = takeMVar start >> replicateM_ 100 (appendOne store name)+            reader = replicateM_ 200 $ do+                Right (Just (info, Just (GlobalPosition position))) <- runStoreIO store (getStreamWithHead name)+                info ^. #version `shouldBe` StreamVersion position+        outcome <- timeout 20_000_000 $ Async.withAsync writer $ \worker -> do+            putMVar start ()+            reader+            Async.wait worker+        outcome `shouldBe` Just ()+        Right (Just (info, headPosition)) <- runStoreIO store (getStreamWithHead name)+        info ^. #version `shouldBe` StreamVersion 101+        headPosition `shouldBe` Just (GlobalPosition 101)++appendOne :: KirokuStore -> StreamName -> IO ()+appendOne store name = void (appendPosition store name)++appendPosition :: KirokuStore -> StreamName -> IO GlobalPosition+appendPosition store name = do+    result <- runStoreIO store (appendToStream name AnyVersion [makeEvent "StreamHeadEvent" Null])+    either (fail . show) (pure . (^. #globalPosition)) result++raw :: KirokuStore -> Text -> IO ()+raw store command = Pool.use (store ^. #pool) (Session.script command) >>= either (expectationFailure . show) pure
+ test/Test/StreamHeadIsolation.hs view
@@ -0,0 +1,36 @@+module Test.StreamHeadIsolation (spec) where++import Control.Monad (forM_)+import Data.Text qualified as T+import Data.Text.IO qualified as T+import Hasql.Statement qualified as Statement+import Kiroku.Store.SQL qualified as SQL+import Paths_kiroku_store (getDataFileName)+import Test.Hspec++spec :: Spec+spec = describe "stream head legacy isolation" $ do+    forM_ statements $ \(name, actual) ->+        it ("preserves " <> name <> " byte for byte") $ do+            path <- getDataFileName ("test/fixtures/stream-head-isolation/" <> name <> ".sql")+            expected <- T.readFile path+            actual `shouldBe` expected+    it "keeps the six metadata and eleven event result columns" $ do+        columnCount (Statement.toSql SQL.getStreamStmt) `shouldBe` 6+        columnCount (Statement.toSql SQL.readAllForwardStmt) `shouldBe` 11+  where+    columnCount = length . T.splitOn "," . fst . T.breakOn "FROM" . snd . T.breakOn "SELECT"+    statements =+        [ ("getStreamStmt", Statement.toSql SQL.getStreamStmt)+        , ("readStreamForwardStmt", Statement.toSql SQL.readStreamForwardStmt)+        , ("readStreamBackwardStmt", Statement.toSql SQL.readStreamBackwardStmt)+        , ("readAllForwardStmt", Statement.toSql SQL.readAllForwardStmt)+        , ("readAllBackwardStmt", Statement.toSql SQL.readAllBackwardStmt)+        , ("readCategoryForwardStmt", Statement.toSql SQL.readCategoryForwardStmt)+        , ("readCategoryForwardConsumerGroupStmt", Statement.toSql SQL.readCategoryForwardConsumerGroupStmt)+        , ("appendExpectedVersion", Statement.toSql SQL.appendExpectedVersion)+        , ("appendStreamExists", Statement.toSql SQL.appendStreamExists)+        , ("appendNoStream", Statement.toSql SQL.appendNoStream)+        , ("appendAnyVersion", Statement.toSql SQL.appendAnyVersion)+        , ("linkToStreamStmt", Statement.toSql SQL.linkToStreamStmt)+        ]
+ test/Test/StreamHeadMock.hs view
@@ -0,0 +1,28 @@+module Test.StreamHeadMock (spec) where++import Control.Monad (forM_)+import Data.IORef (IORef, modifyIORef', newIORef, readIORef)+import Data.Time.Clock (getCurrentTime)+import Effectful+import Effectful.Dispatch.Dynamic (interpret_)+import Kiroku.Store+import Test.Hspec++spec :: Spec+spec = describe "stream head mock" $ it "sends exactly one dedicated constructor and preserves all nested-Maybe outcomes" $ do+    now <- getCurrentTime+    let name = StreamName "mock"+        info = StreamInfo (StreamId 42) name (StreamVersion 3) now Nothing (StreamVersion 0)+    forM_ [Nothing, Just (info, Nothing), Just (info, Just (GlobalPosition 5))] $ \expected -> do+        calls <- newIORef (0 :: Int)+        actual <- runEff $ runMock calls name expected (getStreamWithHead name)+        actual `shouldBe` expected+        readIORef calls `shouldReturn` 1++runMock :: (IOE :> es) => IORef Int -> StreamName -> Maybe (StreamInfo, Maybe GlobalPosition) -> Eff (Store : es) a -> Eff es a+runMock calls name expected = interpret_ $ \case+    GetStreamWithHead actual -> do+        liftIO $ actual `shouldBe` name+        liftIO $ modifyIORef' calls (+ 1)+        pure expected+    _ -> error "unexpected Store operation in stream-head mock"
+ test/Test/SubscriptionDeadLetters.hs view
@@ -0,0 +1,126 @@+{-# LANGUAGE OverloadedLabels #-}++module Test.SubscriptionDeadLetters (spec) where++import Control.Lens ((&), (.~), (^.))+import Control.Monad (forM_, void)+import Data.Aeson (object, (.=))+import Data.Generics.Labels ()+import Data.Int (Int32, Int64)+import Data.Text qualified as T+import Data.Time.Clock (getCurrentTime)+import Data.Vector qualified as V+import Effectful (runEff)+import Effectful.Error.Static (runErrorNoCallStack)+import Hasql.Pool qualified as Pool+import Hasql.Session qualified as Session+import Kiroku.Store hiding (HistoryRetentionInventoryQuery (..), member)+import Kiroku.Store.SQL qualified as SQL+import Test.Helpers+import Test.Hspec++spec :: Spec+spec = describe "SubscriptionDeadLetters" $ do+    it "validates page sizes including both boundaries" $ do+        forM_ [0, -1, 1001] $ \n -> mkSubscriptionDeadLetterLimit n `shouldBe` Left (SubscriptionDeadLetterLimitOutOfRange n)+        forM_ [1, 1000] $ \n -> fmap subscriptionDeadLetterLimitValue (mkSubscriptionDeadLetterLimit n) `shouldBe` Right n+        subscriptionDeadLetterLimitValue defaultSubscriptionDeadLetterLimit `shouldBe` 100+    around withTestStore $ do+        it "returns empty pages through both Store interpreters" $ \store -> do+            let query = defaultSubscriptionDeadLetterQuery (SubscriptionName "absent")+            runStoreIO store (subscriptionDeadLetters query) `shouldReturn` Right (SubscriptionDeadLetterPage V.empty Nothing)+            (runEff . runErrorNoCallStack @StoreError . runKirokuStoreWith store . runStoreResource $ subscriptionDeadLetters query)+                `shouldReturn` Right (SubscriptionDeadLetterPage V.empty Nothing)+        it "pages without omissions, isolates names and matches the internal member read" $ \store -> do+            events <- seedEvents store 7+            mapM_ (insertDeadLetterForEvent store "page") (take 5 events)+            mapM_ (insertDeadLetterForEvent store "other") (drop 5 events)+            let query = (defaultSubscriptionDeadLetterQuery (SubscriptionName "page")){limit = validLimit 2}+            first <- page store query+            second <- page store query{after = first ^. #nextCursor}+            third <- page store query{after = second ^. #nextCursor}+            map positions [first, second, third] `shouldBe` [[5, 4], [3, 2], [1]]+            third ^. #nextCursor `shouldBe` Nothing+            Right internal <- Pool.use (store ^. #pool) (Session.statement ("page", 0) SQL.readDeadLettersStmt)+            map SQL.deadLetterGlobalPosition (V.toList internal) `shouldBe` [5, 4, 3, 2, 1]+            forM_ [5, 1000] $ \n -> do+                final <- page store query{limit = validLimit n}+                positions final `shouldBe` [5, 4, 3, 2, 1]+                final ^. #nextCursor `shouldBe` Nothing+        it "merges historical members and resolves same-position ties with an exclusive cursor" $ \store -> do+            [event] <- seedEvents store 1+            forM_ [0, 1, 9] $ \member -> insertDeadLetterWith store "ties" member (object []) "tie" 2 event+            let query = (defaultSubscriptionDeadLetterQuery (SubscriptionName "ties")){limit = validLimit 1}+            first <- page store query+            second <- page store query{after = first ^. #nextCursor}+            third <- page store query{after = second ^. #nextCursor}+            map (\p -> map (\row -> row ^. #consumerGroupMember) (V.toList (p ^. #deadLetters))) [first, second, third] `shouldBe` [[9], [1], [0]]+            third ^. #nextCursor `shouldBe` Nothing+            forM_ [0, 1, 9, 7] $ \member -> do+                selected <- page store (query & #consumerGroupMember .~ Just member)+                positions selected `shouldBe` if member == 7 then [] else [1]+        it "preserves arbitrary and decode-failure JSON unchanged" $ \store -> do+            events <- seedEvents store 2+            let reasons = [object ["kind" .= ("other" :: T.Text), "detail" .= object ["code" .= (42 :: Int)]], object ["kind" .= ("decode_failure" :: T.Text), "detail" .= ("bad version" :: T.Text)]]+            forM_ (zip events reasons) $ \(event, reason) -> insertDeadLetterWith store "json" 0 reason "summary" 3 event+            result <- page store (defaultSubscriptionDeadLetterQuery (SubscriptionName "json"))+            map (\row -> row ^. #reason) (V.toList (result ^. #deadLetters)) `shouldBe` reverse reasons+            now <- getCurrentTime+            forM_ (V.toList (result ^. #deadLetters)) $ \row -> do+                row ^. #attemptCount `shouldBe` 3+                row ^. #createdAt `shouldSatisfy` (<= now)+        it "never invokes a configured event decode hook for dead-letter reads" $ \_ ->+            withTestStoreSettings (\settings -> (settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just (\_ -> fail "dead-letter read invoked decode hook")})) $ \store -> do+                Right _ <- runStoreIO store $ appendToStream (StreamName "hook-1") NoStream [makeEvent "E" (object [])]+                Right () <- Pool.use (store ^. #pool) (Session.script "INSERT INTO dead_letters(subscription_name,global_position,event_id,reason,reason_summary,attempt_count) SELECT 'hook',1,event_id,'{}','fixture',1 FROM events")+                result <- page store (defaultSubscriptionDeadLetterQuery (SubscriptionName "hook"))+                positions result `shouldBe` [1]++        it "keeps a cursor valid after its event and row are hard-deleted" $ \store -> do+            events <- seedEvents store 3+            mapM_ (insertDeadLetterForEvent store "delete") events+            let query = (defaultSubscriptionDeadLetterQuery (SubscriptionName "delete")){limit = validLimit 1}+            first <- page store query+            Right _ <- runStoreIO store $ hardDeleteStream (StreamName "dead-3")+            older <- page store query{after = first ^. #nextCursor, limit = validLimit 10}+            positions older `shouldBe` [2, 1]+        it "includes legal maxBound pairs on first pages" $ \store -> do+            [event] <- seedEvents store 1+            insertDeadLetterForEvent store "boundary" event+            Right () <- Pool.use (store ^. #pool) (Session.script "UPDATE dead_letters SET global_position=9223372036854775807, dead_letter_id=9223372036854775807 WHERE subscription_name='boundary'")+            let query = defaultSubscriptionDeadLetterQuery (SubscriptionName "boundary")+            forM_ [Nothing, Just 0] $ \member -> do+                result <- page store (query & #consumerGroupMember .~ member)+                positions result `shouldBe` [maxBound]+                older <- page store (query & #consumerGroupMember .~ member & #after .~ Just (SubscriptionDeadLetterCursor (GlobalPosition maxBound) maxBound))+                positions older `shouldBe` []+        it "reads a real worker-produced poison reason without changing its checkpoint" $ \store -> do+            events <- seedEvents store 3+            waitForPublisher store (GlobalPosition 3)+            handle <- subscribe store $ defaultSubscriptionConfig (SubscriptionName "worker") AllStreams $ \event -> pure $ case event ^. #globalPosition of+                GlobalPosition 2 -> DeadLetter (DeadLetterPoison "boom")+                GlobalPosition 3 -> Stop+                _ -> Continue+            waitWithTimeout 10_000_000 handle >>= \case+                Right (Right ()) -> pure ()+                other -> expectationFailure (show other)+            result <- page store (defaultSubscriptionDeadLetterQuery (SubscriptionName "worker"))+            case V.toList (result ^. #deadLetters) of+                [row] -> do+                    row ^. #eventId `shouldBe` (events !! 1) ^. #eventId+                    row ^. #reason `shouldBe` object ["kind" .= ("poison" :: T.Text), "detail" .= ("boom" :: T.Text)]+                    row ^. #reasonSummary `shouldBe` "poison: boom"+                _ -> expectationFailure (show result)+            Pool.use (store ^. #pool) (Session.statement ("worker", 0) SQL.getCheckpointMemberStmt) `shouldReturn` Right (Just 3)++validLimit :: Int32 -> SubscriptionDeadLetterLimit+validLimit = either (error . show) (\value -> value) . mkSubscriptionDeadLetterLimit+page :: KirokuStore -> SubscriptionDeadLetterQuery -> IO SubscriptionDeadLetterPage+page store query = runStoreIO store (subscriptionDeadLetters query) >>= either (fail . show) pure+positions :: SubscriptionDeadLetterPage -> [Int64]+positions result = map (\row -> let GlobalPosition p = row ^. #globalPosition in p) (V.toList (result ^. #deadLetters))+seedEvents :: KirokuStore -> Int -> IO [RecordedEvent]+seedEvents store count = do+    forM_ [1 .. count] $ \n -> void $ runStoreIO store $ appendToStream (StreamName ("dead-" <> T.pack (show n))) NoStream [makeEvent "E" (object [])]+    Right events <- runStoreIO store $ readAllForward (GlobalPosition 0) (fromIntegral count)+    pure (V.toList events)
+ test/Test/SubscriptionDeadLettersMock.hs view
@@ -0,0 +1,23 @@+module Test.SubscriptionDeadLettersMock (spec) where++import Control.Monad.IO.Class (liftIO)+import Data.IORef+import Data.Vector qualified as V+import Effectful (runEff)+import Effectful.Dispatch.Dynamic (interpret_)+import Kiroku.Store+import Test.Hspec++spec :: Spec+spec = describe "SubscriptionDeadLetters mock interpreter" $ it "dispatches the public wrapper once with the entire query and returns the interpreter page" $ do+    calls <- newIORef (0 :: Int)+    let query = (defaultSubscriptionDeadLetterQuery (SubscriptionName "mock")){consumerGroupMember = Just 9, after = Just (SubscriptionDeadLetterCursor (GlobalPosition 21) 7)}+        expected = SubscriptionDeadLetterPage V.empty (Just (SubscriptionDeadLetterCursor (GlobalPosition 10) 3))+        runner = interpret_ $ \case+            ListSubscriptionDeadLetters actual -> do+                liftIO $ actual `shouldBe` query+                liftIO $ modifyIORef' calls (+ 1)+                pure expected+            _ -> error "unexpected Store operation"+    runEff (runner (subscriptionDeadLetters query)) `shouldReturn` expected+    readIORef calls `shouldReturn` 1
test/Test/SubscriptionPauseResume.hs view
@@ -30,19 +30,22 @@ -} module Test.SubscriptionPauseResume (spec) where +import Control.Concurrent (threadDelay) import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar) import Control.Concurrent.STM (atomically, newTVarIO, readTVar, writeTVar) import Control.Exception qualified import Control.Lens ((&), (.~), (^.)) import Data.Aeson qualified as Aeson import Data.Generics.Labels ()-import Data.IORef (modifyIORef', newIORef, readIORef)+import Data.IORef (atomicModifyIORef', modifyIORef', newIORef, readIORef) import Data.Int (Int32) import Data.Text qualified as T import Hasql.Pool qualified as Pool import Hasql.Session qualified as Session import Kiroku.Store import Kiroku.Store.SQL qualified as SQL+import Kiroku.Store.Subscription.Fsm qualified as Fsm+import System.Timeout qualified as Timeout import Test.Helpers (makeEvent, waitForPublisher, waitWithTimeout, withTestStoreSettings) import Test.Hspec @@ -59,7 +62,8 @@ spec = describe "subscription FSM — recoverable backpressure (EP-41 M2)" $ do     it "pauses a slow AllStreams consumer and resumes, delivering all events" $ do         evtRef <- newIORef ([] :: [KirokuEvent])-        let evtHandler e = modifyIORef' evtRef (e :)+        -- Publisher and worker callbacks can arrive concurrently.+        let evtHandler e = atomicModifyIORef' evtRef (\events -> (e : events, ()))         withTestStoreSettings (& #eventHandler .~ Just evtHandler) $ \store -> do             firstSeen <- newEmptyMVar             release <- newEmptyMVar@@ -95,6 +99,7 @@                         , selector = Nothing                         }             handle <- subscribe store cfg+            waitForLive handle             -- First append: the worker reads it from the queue and the handler             -- blocks inside it on `release`, so the worker stops draining.             Right _ <- runStoreIO store $ appendToStream (StreamName "pr-1") NoStream [makeEvent "E1" (Aeson.object [])]@@ -163,6 +168,7 @@                         , selector = Nothing                         }             handle <- subscribe store cfg+            waitForLive handle             Right _ <- runStoreIO store $ appendToStream (StreamName "ds-1") NoStream [makeEvent "E1" (Aeson.object [])]             takeMVar firstSeen             let appendOne i = do@@ -182,3 +188,13 @@                     case Control.Exception.fromException e of                         Just (SubscriptionOverflowed sn) -> sn `shouldBe` SubscriptionName "dropsub-overflow-test"                         Nothing -> expectationFailure ("expected SubscriptionOverflowed, got: " <> show e)++-- These fixtures must block live delivery, not the initial catch-up read.+waitForLive :: SubscriptionHandle -> IO ()+waitForLive handle = do+    let go = do+            st <- currentState handle+            case st of+                Just Fsm.Live{} -> pure ()+                _ -> threadDelay 1_000 >> go+    Timeout.timeout 5_000_000 go >>= (`shouldBe` Just ())
+ test/fixtures/stream-head-isolation/appendAnyVersion.sql view
@@ -0,0 +1,46 @@+WITH+  new_events AS (+    SELECT *+    FROM unnest($1::uuid[], $2::text[], $3::uuid[], $4::uuid[], $5::jsonb[], $6::jsonb[], $7::timestamptz[])+    WITH ORDINALITY AS t(event_id, event_type, causation_id, correlation_id, data, metadata, created_at, idx)+  ),+  stream_upsert AS (+    INSERT INTO streams (stream_name, stream_version)+    VALUES ($8, (SELECT count(*) FROM new_events))+    ON CONFLICT (stream_name)+    DO UPDATE SET stream_version = streams.stream_version + (SELECT count(*) FROM new_events)+      WHERE streams.deleted_at IS NULL+    RETURNING stream_id, category, stream_version - (SELECT count(*) FROM new_events) AS initial_version+  ),+  inserted_events AS (+    INSERT INTO events (event_id, event_type, causation_id, correlation_id, data, metadata, created_at)+    SELECT event_id, event_type, causation_id, correlation_id, data, metadata, created_at+    FROM new_events+    WHERE EXISTS (SELECT 1 FROM stream_upsert)+    ORDER BY idx+  ),+  source_links AS (+    INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version)+    SELECT ne.event_id, su.stream_id, su.initial_version + ne.idx, su.stream_id, su.initial_version + ne.idx+    FROM new_events ne+    CROSS JOIN stream_upsert su+  ),+  all_update AS (+    UPDATE streams+    SET stream_version = stream_version + (SELECT count(*) FROM new_events)+    WHERE stream_id = 0+      AND EXISTS (SELECT 1 FROM stream_upsert)+    RETURNING stream_version - (SELECT count(*) FROM new_events) AS initial_global_version+  ),+  all_links AS (+    INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version, category)+    SELECT ne.event_id, 0, au.initial_global_version + ne.idx, su.stream_id, su.initial_version + ne.idx, su.category+    FROM new_events ne+    CROSS JOIN all_update au+    CROSS JOIN stream_upsert su+  )+SELECT su.stream_id,+       su.initial_version + (SELECT count(*) FROM new_events),+       au.initial_global_version + (SELECT count(*) FROM new_events)+FROM stream_upsert su+CROSS JOIN all_update au
+ test/fixtures/stream-head-isolation/appendExpectedVersion.sql view
@@ -0,0 +1,46 @@+WITH+  new_events AS (+    SELECT *+    FROM unnest($1::uuid[], $2::text[], $3::uuid[], $4::uuid[], $5::jsonb[], $6::jsonb[], $7::timestamptz[])+    WITH ORDINALITY AS t(event_id, event_type, causation_id, correlation_id, data, metadata, created_at, idx)+  ),+  stream_update AS (+    UPDATE streams+    SET stream_version = stream_version + (SELECT count(*) FROM new_events)+    WHERE stream_name = $8+      AND stream_version = $9+      AND deleted_at IS NULL+    RETURNING stream_id, category, stream_version - (SELECT count(*) FROM new_events) AS initial_version+  ),+  inserted_events AS (+    INSERT INTO events (event_id, event_type, causation_id, correlation_id, data, metadata, created_at)+    SELECT event_id, event_type, causation_id, correlation_id, data, metadata, created_at+    FROM new_events+    WHERE EXISTS (SELECT 1 FROM stream_update)+    ORDER BY idx+  ),+  source_links AS (+    INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version)+    SELECT ne.event_id, su.stream_id, su.initial_version + ne.idx, su.stream_id, su.initial_version + ne.idx+    FROM new_events ne+    CROSS JOIN stream_update su+  ),+  all_update AS (+    UPDATE streams+    SET stream_version = stream_version + (SELECT count(*) FROM new_events)+    WHERE stream_id = 0+      AND EXISTS (SELECT 1 FROM stream_update)+    RETURNING stream_version - (SELECT count(*) FROM new_events) AS initial_global_version+  ),+  all_links AS (+    INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version, category)+    SELECT ne.event_id, 0, au.initial_global_version + ne.idx, su.stream_id, su.initial_version + ne.idx, su.category+    FROM new_events ne+    CROSS JOIN all_update au+    CROSS JOIN stream_update su+  )+SELECT su.stream_id,+       su.initial_version + (SELECT count(*) FROM new_events),+       au.initial_global_version + (SELECT count(*) FROM new_events)+FROM stream_update su+CROSS JOIN all_update au
+ test/fixtures/stream-head-isolation/appendNoStream.sql view
@@ -0,0 +1,44 @@+WITH+  new_events AS (+    SELECT *+    FROM unnest($1::uuid[], $2::text[], $3::uuid[], $4::uuid[], $5::jsonb[], $6::jsonb[], $7::timestamptz[])+    WITH ORDINALITY AS t(event_id, event_type, causation_id, correlation_id, data, metadata, created_at, idx)+  ),+  stream_insert AS (+    INSERT INTO streams (stream_name, stream_version)+    VALUES ($8, (SELECT count(*) FROM new_events))+    ON CONFLICT (stream_name) DO NOTHING+    RETURNING stream_id, category, 0::bigint AS initial_version+  ),+  inserted_events AS (+    INSERT INTO events (event_id, event_type, causation_id, correlation_id, data, metadata, created_at)+    SELECT event_id, event_type, causation_id, correlation_id, data, metadata, created_at+    FROM new_events+    WHERE EXISTS (SELECT 1 FROM stream_insert)+    ORDER BY idx+  ),+  source_links AS (+    INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version)+    SELECT ne.event_id, si.stream_id, si.initial_version + ne.idx, si.stream_id, si.initial_version + ne.idx+    FROM new_events ne+    CROSS JOIN stream_insert si+  ),+  all_update AS (+    UPDATE streams+    SET stream_version = stream_version + (SELECT count(*) FROM new_events)+    WHERE stream_id = 0+      AND EXISTS (SELECT 1 FROM stream_insert)+    RETURNING stream_version - (SELECT count(*) FROM new_events) AS initial_global_version+  ),+  all_links AS (+    INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version, category)+    SELECT ne.event_id, 0, au.initial_global_version + ne.idx, si.stream_id, si.initial_version + ne.idx, si.category+    FROM new_events ne+    CROSS JOIN all_update au+    CROSS JOIN stream_insert si+  )+SELECT si.stream_id,+       si.initial_version + (SELECT count(*) FROM new_events),+       au.initial_global_version + (SELECT count(*) FROM new_events)+FROM stream_insert si+CROSS JOIN all_update au
+ test/fixtures/stream-head-isolation/appendStreamExists.sql view
@@ -0,0 +1,45 @@+WITH+  new_events AS (+    SELECT *+    FROM unnest($1::uuid[], $2::text[], $3::uuid[], $4::uuid[], $5::jsonb[], $6::jsonb[], $7::timestamptz[])+    WITH ORDINALITY AS t(event_id, event_type, causation_id, correlation_id, data, metadata, created_at, idx)+  ),+  stream_update AS (+    UPDATE streams+    SET stream_version = stream_version + (SELECT count(*) FROM new_events)+    WHERE stream_name = $8+      AND deleted_at IS NULL+    RETURNING stream_id, category, stream_version - (SELECT count(*) FROM new_events) AS initial_version+  ),+  inserted_events AS (+    INSERT INTO events (event_id, event_type, causation_id, correlation_id, data, metadata, created_at)+    SELECT event_id, event_type, causation_id, correlation_id, data, metadata, created_at+    FROM new_events+    WHERE EXISTS (SELECT 1 FROM stream_update)+    ORDER BY idx+  ),+  source_links AS (+    INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version)+    SELECT ne.event_id, su.stream_id, su.initial_version + ne.idx, su.stream_id, su.initial_version + ne.idx+    FROM new_events ne+    CROSS JOIN stream_update su+  ),+  all_update AS (+    UPDATE streams+    SET stream_version = stream_version + (SELECT count(*) FROM new_events)+    WHERE stream_id = 0+      AND EXISTS (SELECT 1 FROM stream_update)+    RETURNING stream_version - (SELECT count(*) FROM new_events) AS initial_global_version+  ),+  all_links AS (+    INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version, category)+    SELECT ne.event_id, 0, au.initial_global_version + ne.idx, su.stream_id, su.initial_version + ne.idx, su.category+    FROM new_events ne+    CROSS JOIN all_update au+    CROSS JOIN stream_update su+  )+SELECT su.stream_id,+       su.initial_version + (SELECT count(*) FROM new_events),+       au.initial_global_version + (SELECT count(*) FROM new_events)+FROM stream_update su+CROSS JOIN all_update au
+ test/fixtures/stream-head-isolation/getStreamStmt.sql view
@@ -0,0 +1,3 @@+SELECT stream_id, stream_name, stream_version, created_at, deleted_at, truncate_before+FROM streams+WHERE stream_name = $1
+ test/fixtures/stream-head-isolation/linkToStreamStmt.sql view
@@ -0,0 +1,33 @@+WITH+  event_list AS (+    SELECT event_id, idx+    FROM unnest($1::uuid[]) WITH ORDINALITY AS t(event_id, idx)+  ),+  stream_upsert AS (+    INSERT INTO streams (stream_name, stream_version)+    VALUES ($2, (SELECT count(*) FROM event_list))+    ON CONFLICT (stream_name)+    DO UPDATE SET stream_version = streams.stream_version + (SELECT count(*) FROM event_list)+      WHERE streams.deleted_at IS NULL+    RETURNING stream_id, stream_version - (SELECT count(*) FROM event_list) AS initial_version+  ),+  link_inserts AS (+    -- LEFT JOIN LATERAL surfaces missing-event rows as NULLs for original_*; the+    -- NOT NULL constraint on stream_events.original_stream_id then aborts the+    -- entire CTE, rolling back the stream_upsert's version bump. Before the F3+    -- fix this was a plain JOIN LATERAL, which silently dropped missing-event+    -- rows while still bumping stream_version → silent gap in the link target.+    INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version)+    SELECT el.event_id, su.stream_id, su.initial_version + el.idx,+           orig.original_stream_id, orig.original_stream_version+    FROM event_list el+    CROSS JOIN stream_upsert su+    LEFT JOIN LATERAL (+      SELECT se.original_stream_id, se.original_stream_version+      FROM stream_events se+      WHERE se.event_id = el.event_id AND se.stream_id <> 0+      LIMIT 1+    ) orig ON true+  )+SELECT su.stream_id, su.initial_version + (SELECT count(*) FROM event_list)+FROM stream_upsert su
+ test/fixtures/stream-head-isolation/readAllBackwardStmt.sql view
@@ -0,0 +1,11 @@+SELECT e.event_id, e.event_type,+       se.stream_version, se.stream_version AS global_position,+       se.original_stream_id, se.original_stream_version,+       e.data, e.metadata, e.causation_id, e.correlation_id,+       e.created_at+FROM stream_events se+JOIN events e ON e.event_id = se.event_id+WHERE se.stream_id = 0+  AND se.stream_version < $1+ORDER BY se.stream_version DESC+LIMIT $2
+ test/fixtures/stream-head-isolation/readAllForwardStmt.sql view
@@ -0,0 +1,11 @@+SELECT e.event_id, e.event_type,+       se.stream_version, se.stream_version AS global_position,+       se.original_stream_id, se.original_stream_version,+       e.data, e.metadata, e.causation_id, e.correlation_id,+       e.created_at+FROM stream_events se+JOIN events e ON e.event_id = se.event_id+WHERE se.stream_id = 0+  AND se.stream_version > $1+ORDER BY se.stream_version ASC+LIMIT $2
+ test/fixtures/stream-head-isolation/readCategoryForwardConsumerGroupStmt.sql view
@@ -0,0 +1,13 @@+SELECT e.event_id, e.event_type,+       se.stream_version, se.stream_version AS global_position,+       se.original_stream_id, se.original_stream_version,+       e.data, e.metadata, e.causation_id, e.correlation_id,+       e.created_at+FROM stream_events se+JOIN events e ON e.event_id = se.event_id+WHERE se.stream_id = 0+  AND se.category = $2+  AND se.stream_version > $1+  AND (((hashtextextended(se.original_stream_id::text, 0) % $4) + $4) % $4) = $3+ORDER BY se.stream_version ASC+LIMIT $5
+ test/fixtures/stream-head-isolation/readCategoryForwardStmt.sql view
@@ -0,0 +1,12 @@+SELECT e.event_id, e.event_type,+       se.stream_version, se.stream_version AS global_position,+       se.original_stream_id, se.original_stream_version,+       e.data, e.metadata, e.causation_id, e.correlation_id,+       e.created_at+FROM stream_events se+JOIN events e ON e.event_id = se.event_id+WHERE se.stream_id = 0+  AND se.category = $2+  AND se.stream_version > $1+ORDER BY se.stream_version ASC+LIMIT $3
+ test/fixtures/stream-head-isolation/readStreamBackwardStmt.sql view
@@ -0,0 +1,14 @@+SELECT e.event_id, e.event_type,+       se.stream_version, 0::bigint AS global_position,+       se.original_stream_id, se.original_stream_version,+       e.data, e.metadata, e.causation_id, e.correlation_id,+       e.created_at+FROM stream_events se+JOIN events e  ON e.event_id  = se.event_id+JOIN streams s ON s.stream_id = se.stream_id+WHERE s.stream_name = $1+  AND s.deleted_at IS NULL+  AND se.stream_version < $2+  AND se.stream_version >= s.truncate_before+ORDER BY se.stream_version DESC+LIMIT $3
+ test/fixtures/stream-head-isolation/readStreamForwardStmt.sql view
@@ -0,0 +1,14 @@+SELECT e.event_id, e.event_type,+       se.stream_version, 0::bigint AS global_position,+       se.original_stream_id, se.original_stream_version,+       e.data, e.metadata, e.causation_id, e.correlation_id,+       e.created_at+FROM stream_events se+JOIN events e  ON e.event_id  = se.event_id+JOIN streams s ON s.stream_id = se.stream_id+WHERE s.stream_name = $1+  AND s.deleted_at IS NULL+  AND se.stream_version > $2+  AND se.stream_version >= s.truncate_before+ORDER BY se.stream_version ASC+LIMIT $3