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 +29/−0
- bench/ShibuyaOverhead.hs +2/−1
- bench/StreamHeadCost.hs +119/−0
- bench/StreamHeadPaired.hs +145/−0
- kiroku-store.cabal +70/−6
- src/Kiroku/Store/Effect.hs +34/−0
- src/Kiroku/Store/Read.hs +50/−0
- src/Kiroku/Store/SQL.hs +156/−0
- src/Kiroku/Store/Subscription.hs +17/−1
- src/Kiroku/Store/Subscription/DeadLetter/SQL.hs +142/−0
- src/Kiroku/Store/Subscription/EventPublisher.hs +26/−4
- src/Kiroku/Store/Subscription/Types.hs +79/−2
- src/Kiroku/Store/Subscription/Worker.hs +1/−1
- src/Kiroku/Store/Types.hs +22/−1
- test/Main.hs +27/−0
- test/Test/BrowseQueryPlans.hs +109/−0
- test/Test/BrowseReads.hs +94/−0
- test/Test/BrowseReadsMock.hs +36/−0
- test/Test/DeadLetterQueryPlans.hs +98/−0
- test/Test/Helpers.hs +10/−6
- test/Test/PerformanceStructure.hs +90/−0
- test/Test/PublisherDropCounter.hs +59/−0
- test/Test/StreamHead.hs +108/−0
- test/Test/StreamHeadIsolation.hs +36/−0
- test/Test/StreamHeadMock.hs +28/−0
- test/Test/SubscriptionDeadLetters.hs +126/−0
- test/Test/SubscriptionDeadLettersMock.hs +23/−0
- test/Test/SubscriptionPauseResume.hs +18/−2
- test/fixtures/stream-head-isolation/appendAnyVersion.sql +46/−0
- test/fixtures/stream-head-isolation/appendExpectedVersion.sql +46/−0
- test/fixtures/stream-head-isolation/appendNoStream.sql +44/−0
- test/fixtures/stream-head-isolation/appendStreamExists.sql +45/−0
- test/fixtures/stream-head-isolation/getStreamStmt.sql +3/−0
- test/fixtures/stream-head-isolation/linkToStreamStmt.sql +33/−0
- test/fixtures/stream-head-isolation/readAllBackwardStmt.sql +11/−0
- test/fixtures/stream-head-isolation/readAllForwardStmt.sql +11/−0
- test/fixtures/stream-head-isolation/readCategoryForwardConsumerGroupStmt.sql +13/−0
- test/fixtures/stream-head-isolation/readCategoryForwardStmt.sql +12/−0
- test/fixtures/stream-head-isolation/readStreamBackwardStmt.sql +14/−0
- test/fixtures/stream-head-isolation/readStreamForwardStmt.sql +14/−0
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