kiroku-store 0.8.0.2 → 0.9.0.0
raw patch · 10 files changed
+815/−232 lines, 10 filesPVP ok
version bump matches the API change (PVP)
API changes (from Hackage documentation)
+ Kiroku.Store.SQL: appendParamsEncoder :: Params AppendParams
+ Kiroku.Store.SQL: appendResultDecoder :: Result (Maybe AppendResult)
+ Kiroku.Store.SQL: readCategoryConsumerGroupEncoder :: Params (Int64, Text, Int32, Int32, Int32)
+ Kiroku.Store.SQL: readCategoryEncoder :: Params (Int64, Text, Int32)
+ Kiroku.Store.SQL: recordedEventRow :: Row RecordedEvent
Files
- CHANGELOG.md +46/−0
- bench/Explain.hs +3/−3
- bench/Main.hs +184/−122
- bench/RegressionGate.hs +307/−24
- kiroku-store.cabal +2/−1
- src/Kiroku/Store/SQL.hs +42/−39
- src/Kiroku/Store/Subscription.hs +5/−2
- src/Kiroku/Store/Subscription/Worker.hs +14/−9
- test/Test/CategoryIdleNoSpin.hs +77/−16
- test/Test/PerformanceStructure.hs +135/−16
CHANGELOG.md view
@@ -1,5 +1,51 @@ # Changelog +## 0.9.0.0 — 2026-09-25++### Breaking Changes++* Requires schema migration `0012` from kiroku-store-migrations 0.6.0.0. The+ append statements now write the source stream's category onto each `$all`+ junction row, so against an older schema every append fails with SQLSTATE+ `42703` (undefined column). Conversely, once `0012` is applied, any process+ still on kiroku-store 0.8 or older fails every append with SQLSTATE `23514`,+ because its `$all` rows lack the category the new check constraint requires.+ There is no rolling-deploy path: stop every appending process, apply `0012`+ (in a maintenance window on a large store), then start the new code. Reads by+ old code keep working on the new schema. The `kiroku-upgrade` Seihou+ blueprint's `0.8.0.2 -> 0.9.0.0` edge walks a consuming project through it.++### Bug Fixes++* Category reads no longer cost work proportional to the number of streams in+ the category (BUG-2). `readCategory`, category subscriptions, and+ consumer-group category subscriptions used a LATERAL join that probed every+ stream of the category on each call: a caught-up poll of a 20,000-stream+ category read 60,387 shared buffers (about 30,000 for a member of a size-2+ group), and the figure grew with every stream ever created. Both statements+ are now one range scan of `ix_stream_events_all_by_category` from+ `(category, checkpoint)` that stops at the limit; the same poll reads 6+ buffers (3 for the group member). Results and ordering are unchanged.++### Other Changes++* Consumer-group members of a `Category` subscription now wake in live mode+ only when their category receives an append, as ordinary category+ subscriptions already did, instead of on every append anywhere in the+ store. Members of an idle category do no live database work; a member whose+ sibling owns the new event does one empty fetch. `AllStreams` group members+ are unchanged.+* `Kiroku.Store.SQL` additionally exports `appendParamsEncoder`,+ `appendResultDecoder`, `readCategoryEncoder`,+ `readCategoryConsumerGroupEncoder`, and `recordedEventRow`, used by the+ benchmark controls.+* New structural tests pin both category statements to+ `ix_stream_events_all_by_category` without a Sort and hold a caught-up poll+ on a 20,000-stream category to 32 buffers. The workload gate gains+ `category-read` (the new statements against the LATERAL ones) and+ `append-category-column` (the append cost of the new column and index), and+ the benchmark suite gains a `category-scaling` group.+ ## 0.8.0.2 — 2026-09-21 ### Bug Fixes
bench/Explain.hs view
@@ -98,7 +98,7 @@ 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, stream_version - (SELECT count(*) FROM new_events) AS initial_version+ 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)@@ -121,8 +121,8 @@ 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)- SELECT ne.event_id, 0, au.initial_global_version + ne.idx, su.stream_id, su.initial_version + ne.idx+ 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
bench/Main.hs view
@@ -26,6 +26,8 @@ import Hasql.Transaction qualified as Tx import Hasql.Transaction.Sessions qualified as TxSessions import Kiroku.Store+import Kiroku.Store.SQL qualified as SQL+import Kiroku.Test.Fixtures.CategoryScaling (categoryScalingFixtureSql, categoryScalingHead) import Kiroku.Test.Postgres (ephemeralConfig, migrateTestDatabase, withMigratedTestDatabase, withSharedMigratedPostgres) import Test.Tasty.Bench @@ -118,7 +120,7 @@ ON CONFLICT (stream_name) DO UPDATE SET stream_version = streams.stream_version + 1 WHERE streams.deleted_at IS NULL- RETURNING stream_id, stream_version - 1 AS initial_version+ RETURNING stream_id, category, stream_version - 1 AS initial_version ), inserted_events AS ( INSERT INTO events (event_id, event_type, causation_id, correlation_id, data, metadata, created_at)@@ -140,8 +142,8 @@ RETURNING stream_version - 1 AS initial_global_version ), all_links AS (- INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version)- SELECT ne.event_id, 0, au.initial_global_version + 1, su.stream_id, su.initial_version + 1+ 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 + 1, su.stream_id, su.initial_version + 1, su.category FROM new_event ne CROSS JOIN all_update au CROSS JOIN stream_upsert su@@ -168,7 +170,7 @@ 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, stream_version - (SELECT count(*) FROM new_events) AS initial_version+ 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)@@ -191,8 +193,8 @@ 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)- SELECT ne.event_id, 0, au.initial_global_version + ne.idx, su.stream_id, su.initial_version + ne.idx+ 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@@ -280,7 +282,7 @@ UPDATE streams SET stream_version = stream_version + 1 WHERE stream_id = $1::bigint- RETURNING stream_id, stream_version - 1 AS initial_version+ RETURNING stream_id, category, stream_version - 1 AS initial_version ), inserted_event AS ( INSERT INTO events (event_id, event_type, causation_id, correlation_id, data, metadata, created_at)@@ -298,8 +300,8 @@ RETURNING stream_version - 1 AS initial_global_version ), all_link AS (- INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version)- SELECT $2::uuid, 0, au.initial_global_version + 1, su.stream_id, su.initial_version + 1+ INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version, category)+ SELECT $2::uuid, 0, au.initial_global_version + 1, su.stream_id, su.initial_version + 1, su.category FROM all_update au CROSS JOIN stream_update su )@@ -352,7 +354,7 @@ stream_insert AS ( INSERT INTO streams (stream_name, stream_version) VALUES ($1::text, 1)- RETURNING stream_id, 0::bigint AS initial_version+ RETURNING stream_id, category, 0::bigint AS initial_version ), inserted_event AS ( INSERT INTO events (event_id, event_type, causation_id, correlation_id, data, metadata, created_at)@@ -370,8 +372,8 @@ RETURNING stream_version - 1 AS initial_global_version ), all_link AS (- INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version)- SELECT $2::uuid, 0, au.initial_global_version + 1, si.stream_id, si.initial_version + 1+ INSERT INTO stream_events (event_id, stream_id, stream_version, original_stream_id, original_stream_version, category)+ SELECT $2::uuid, 0, au.initial_global_version + 1, si.stream_id, si.initial_version + 1, si.category FROM all_update au CROSS JOIN stream_insert si )@@ -709,6 +711,43 @@ seedCheckpointInventory inventory10000Store 10_000 action inventory100Store inventory10000Store +{- | A store seeded with 'categoryScalingFixtureSql': categories @performance@+(200 streams), @idle@ (20,000 single-event streams), and @noise@, head at+'categoryScalingHead'. Must run inside 'withSharedMigratedPostgres'.+-}+withCategoryScalingStore :: (KirokuStore -> IO a) -> IO a+withCategoryScalingStore action =+ withMigratedTestDatabase $ \connection ->+ withStore (defaultConnectionSettings connection) $ \store -> do+ seeded <- Pool.use (store ^. #pool) (Session.script categoryScalingFixtureSql)+ case seeded of+ Left err -> error ("Category-scaling benchmark setup failed: " <> show err)+ Right () -> action store++-- | Run the plain category read 10 times at one cursor and limit.+runCategoryPolls :: KirokuStore -> Text -> Int64 -> Int32 -> IO ()+runCategoryPolls store category cursor limit =+ mapM_+ ( \_ ->+ Pool.use (store ^. #pool) (Session.statement (cursor, category, limit) SQL.readCategoryForwardStmt)+ >>= forceCategoryPoll+ )+ [1 .. 10 :: Int]++-- | Run the consumer-group category read 10 times for one member.+runGroupCategoryPolls :: KirokuStore -> Text -> Int64 -> Int32 -> Int32 -> Int32 -> IO ()+runGroupCategoryPolls store category cursor member size limit =+ mapM_+ ( \_ ->+ Pool.use (store ^. #pool) (Session.statement (cursor, category, member, size, limit) SQL.readCategoryForwardConsumerGroupStmt)+ >>= forceCategoryPoll+ )+ [1 .. 10 :: Int]++forceCategoryPoll :: Either Pool.UsageError (V.Vector RecordedEvent) -> IO ()+forceCategoryPoll (Right events) = V.length events `seq` pure ()+forceCategoryPoll (Left err) = error ("Category-scaling benchmark read failed: " <> show err)+ main :: IO () main = do -- Start ephemeral PostgreSQL once for all benchmarks@@ -817,131 +856,154 @@ subCounter <- newIORef (0 :: Int) withInventoryBenchmarkStores $ \inventory100Store inventory10000Store ->- defaultMain- [ bgroup- "append"+ withCategoryScalingStore $ \scalingStore ->+ defaultMain [ bgroup- "single-event"- [ bench "NoStream (new stream)" $ whnfIO $ do- sn <- nextStream "bench-single"- r' <- runStoreIO store $ appendToStream sn NoStream [makeEvent "BenchEvent"]- forceAppend r'- , bench "AnyVersion (new stream)" $ whnfIO $ do- sn <- nextStream "bench-any"- r' <- runStoreIO store $ appendToStream sn AnyVersion [makeEvent "BenchEvent"]- forceAppend r'+ "append"+ [ bgroup+ "single-event"+ [ bench "NoStream (new stream)" $ whnfIO $ do+ sn <- nextStream "bench-single"+ r' <- runStoreIO store $ appendToStream sn NoStream [makeEvent "BenchEvent"]+ forceAppend r'+ , bench "AnyVersion (new stream)" $ whnfIO $ do+ sn <- nextStream "bench-any"+ r' <- runStoreIO store $ appendToStream sn AnyVersion [makeEvent "BenchEvent"]+ forceAppend r'+ ]+ , bgroup+ "batch-10"+ [ bench "NoStream" $ whnfIO $ do+ sn <- nextStream "bench-b10"+ let events = map (\i -> makeEvent ("E" <> T.pack (show i))) [1 .. 10 :: Int]+ r' <- runStoreIO store $ appendToStream sn NoStream events+ forceAppend r'+ ]+ , bgroup+ "batch-100"+ [ bench "NoStream" $ whnfIO $ do+ sn <- nextStream "bench-b100"+ let events = map (\i -> makeEvent ("E" <> T.pack (show i))) [1 .. 100 :: Int]+ r' <- runStoreIO store $ appendToStream sn NoStream events+ forceAppend r'+ ]+ , bgroup+ "sequential"+ [ bench "10 appends to same stream" $ whnfIO $ do+ sn <- nextStream "bench-seq"+ r0 <- runStoreIO store $ appendToStream sn NoStream [makeEvent "Init"]+ forceAppend r0+ let res0 = case r0 of+ Right ok -> ok+ Left e -> error ("Sequential append failed: " <> show e)+ let go _ 0 = pure ()+ go v n = do+ r' <- runStoreIO store $ appendToStream sn (ExactVersion v) [makeEvent "Seq"]+ case r' of+ Right res -> go (res ^. #streamVersion) (n - 1 :: Int)+ Left e -> error ("Sequential append failed: " <> show e)+ go (res0 ^. #streamVersion) 9+ ] ] , bgroup- "batch-10"- [ bench "NoStream" $ whnfIO $ do- sn <- nextStream "bench-b10"- let events = map (\i -> makeEvent ("E" <> T.pack (show i))) [1 .. 10 :: Int]- r' <- runStoreIO store $ appendToStream sn NoStream events- forceAppend r'+ "raw-append-shape"+ [ bgroup+ "AnyVersion"+ [ bench "scalar singleton (new stream)" $+ whnfIO $+ runRawScalarAppendAnyVersionNewStream store rawCounter+ , bench "production arrays/unnest (new stream)" $+ whnfIO $+ runRawProductionAppendAnyVersionNewStream store rawCounter+ , bench "two-roundtrip (new stream)" $+ whnfIO $+ runRawTwoRoundtripAppendNewStream store rawCounter+ , bench "two-roundtrip + BEGIN/COMMIT (new stream)" $+ whnfIO $+ runRawTwoRoundtripAppendNewStreamTx store rawCounter+ , bench "scalar singleton (hot stream)" $+ whnfIO $+ runRawScalarAppendAnyVersionHotStream store+ , bench "production arrays/unnest (hot stream)" $+ whnfIO $+ runRawProductionAppendAnyVersionHotStream store+ , bench "two-roundtrip (hot stream)" $+ whnfIO $+ runRawTwoRoundtripAppendExistingHotStream store+ , bench "two-roundtrip + BEGIN/COMMIT (hot stream)" $+ whnfIO $+ runRawTwoRoundtripAppendExistingHotStreamTx store+ ] ] , bgroup- "batch-100"- [ bench "NoStream" $ whnfIO $ do- sn <- nextStream "bench-b100"- let events = map (\i -> makeEvent ("E" <> T.pack (show i))) [1 .. 100 :: Int]- r' <- runStoreIO store $ appendToStream sn NoStream events- forceAppend r'+ "read"+ [ bench "stream forward (100-event page)" $ whnfIO $ do+ r' <- runStoreIO store $ readStreamForward readStreamName (StreamVersion 0) 100+ forceRead r'+ , bench "$all forward (100-event page)" $ whnfIO $ do+ r' <- runStoreIO store $ readAllForward (GlobalPosition 0) 100+ forceRead r' ] , bgroup- "sequential"- [ bench "10 appends to same stream" $ whnfIO $ do- sn <- nextStream "bench-seq"- r0 <- runStoreIO store $ appendToStream sn NoStream [makeEvent "Init"]- forceAppend r0- let res0 = case r0 of- Right ok -> ok- Left e -> error ("Sequential append failed: " <> show e)- let go _ 0 = pure ()- go v n = do- r' <- runStoreIO store $ appendToStream sn (ExactVersion v) [makeEvent "Seq"]- case r' of- Right res -> go (res ^. #streamVersion) (n - 1 :: Int)- Left e -> error ("Sequential append failed: " <> show e)- go (res0 ^. #streamVersion) 9+ "category"+ [ bench "category forward (100-event page)" $ whnfIO $ do+ -- Read from cat1 category (has 10 streams × 100 events = 1000 events)+ r' <- runStoreIO store $ readCategory (CategoryName "cat1") (GlobalPosition 0) 100+ forceRead r'+ , bench "exhausted-category" $ whnfIO $ do+ -- cat1 events are inserted early in setup; a high cursor proves+ -- category reads do not scan the rest of $all looking for matches.+ r' <- runStoreIO store $ readCategory (CategoryName "cat1") (GlobalPosition 90_000) 100+ forceRead r'+ , bench "$all forward (100-event page baseline)" $ whnfIO $ do+ r' <- runStoreIO store $ readAllForward (GlobalPosition 0) 100+ forceRead r' ]- ]- , bgroup- "raw-append-shape"- [ bgroup- "AnyVersion"- [ bench "scalar singleton (new stream)" $- whnfIO $- runRawScalarAppendAnyVersionNewStream store rawCounter- , bench "production arrays/unnest (new stream)" $- whnfIO $- runRawProductionAppendAnyVersionNewStream store rawCounter- , bench "two-roundtrip (new stream)" $- whnfIO $- runRawTwoRoundtripAppendNewStream store rawCounter- , bench "two-roundtrip + BEGIN/COMMIT (new stream)" $+ , -- F19 — Concurrent-writer stress as structured benchmarks.+ -- The legacy ad-hoc B9 measurement (still present above+ -- for historical comparability) prints throughput and+ -- latency once; these bgroup entries surface the same+ -- workload through tasty-bench so it participates in the+ -- baseline-regression workflow (Justfile bench-regression).+ bgroup+ "concurrent"+ [ bench "8 writers x 10 appends" $ whnfIO $ runConcurrentWriters store concCounter 8 10+ , bench "32 writers x 10 appends" $ whnfIO $ runConcurrentWriters store concCounter 32 10+ ]+ , bgroup+ "reliability-audit"+ [ bench "hot invoice-payment 10 AnyVersion appends" $ whnfIO $ runHotInvoicePayment store 10+ , bench "appendMultiStream 3 existing streams" $ whnfIO $ runAppendMultiStream store+ , bench "subscription category catch-up 100 events" $ whnfIO $ runSubscriptionCatchup store subCounter+ ]+ , -- BUG-2: a caught-up category poll must not cost work+ -- proportional to the number of streams in the category.+ -- Each cell runs the statement 10 times.+ bgroup+ "category-scaling"+ [ bench "plain caught-up poll (200 streams)" $ whnfIO $- runRawTwoRoundtripAppendNewStreamTx store rawCounter- , bench "scalar singleton (hot stream)" $+ runCategoryPolls scalingStore "performance" categoryScalingHead 100+ , bench "plain caught-up poll (20000 streams)" $ whnfIO $- runRawScalarAppendAnyVersionHotStream store- , bench "production arrays/unnest (hot stream)" $+ runCategoryPolls scalingStore "idle" categoryScalingHead 100+ , bench "group caught-up poll (20000 streams)" $ whnfIO $- runRawProductionAppendAnyVersionHotStream store- , bench "two-roundtrip (hot stream)" $+ runGroupCategoryPolls scalingStore "idle" categoryScalingHead 1 2 100+ , bench "plain page from 0 (20000 streams)" $ whnfIO $- runRawTwoRoundtripAppendExistingHotStream store- , bench "two-roundtrip + BEGIN/COMMIT (hot stream)" $+ runCategoryPolls scalingStore "idle" 0 100+ , bench "exhausted category (200 streams)" $+ -- performance ends at 20,000; 60,000 other rows follow. whnfIO $- runRawTwoRoundtripAppendExistingHotStreamTx store+ runCategoryPolls scalingStore "performance" 20_000 100 ]- ]- , bgroup- "read"- [ bench "stream forward (100-event page)" $ whnfIO $ do- r' <- runStoreIO store $ readStreamForward readStreamName (StreamVersion 0) 100- forceRead r'- , bench "$all forward (100-event page)" $ whnfIO $ do- r' <- runStoreIO store $ readAllForward (GlobalPosition 0) 100- forceRead r'- ]- , bgroup- "category"- [ bench "category forward (100-event page)" $ whnfIO $ do- -- Read from cat1 category (has 10 streams × 100 events = 1000 events)- r' <- runStoreIO store $ readCategory (CategoryName "cat1") (GlobalPosition 0) 100- forceRead r'- , bench "exhausted-category" $ whnfIO $ do- -- cat1 events are inserted early in setup; a high cursor proves- -- category reads do not scan the rest of $all looking for matches.- r' <- runStoreIO store $ readCategory (CategoryName "cat1") (GlobalPosition 90_000) 100- forceRead r'- , bench "$all forward (100-event page baseline)" $ whnfIO $ do- r' <- runStoreIO store $ readAllForward (GlobalPosition 0) 100- forceRead r'- ]- , -- F19 — Concurrent-writer stress as structured benchmarks.- -- The legacy ad-hoc B9 measurement (still present above- -- for historical comparability) prints throughput and- -- latency once; these bgroup entries surface the same- -- workload through tasty-bench so it participates in the- -- baseline-regression workflow (Justfile bench-regression).- bgroup- "concurrent"- [ bench "8 writers x 10 appends" $ whnfIO $ runConcurrentWriters store concCounter 8 10- , bench "32 writers x 10 appends" $ whnfIO $ runConcurrentWriters store concCounter 32 10- ]- , bgroup- "reliability-audit"- [ bench "hot invoice-payment 10 AnyVersion appends" $ whnfIO $ runHotInvoicePayment store 10- , bench "appendMultiStream 3 existing streams" $ whnfIO $ runAppendMultiStream store- , bench "subscription category catch-up 100 events" $ whnfIO $ runSubscriptionCatchup store subCounter- ]- , bgroup- "subscription-checkpoint-inventory"- [ bench "100 rows" $ whnfIO $ runCheckpointInventoryBenchmark inventory100Store- , bench "10000 rows" $ whnfIO $ runCheckpointInventoryBenchmark inventory10000Store+ , bgroup+ "subscription-checkpoint-inventory"+ [ bench "100 rows" $ whnfIO $ runCheckpointInventoryBenchmark inventory100Store+ , bench "10000 rows" $ whnfIO $ runCheckpointInventoryBenchmark inventory10000Store+ ] ]- ] case result of Left err -> error ("Failed to start ephemeral PostgreSQL: " <> show err) Right () -> pure ()
bench/RegressionGate.hs view
@@ -1,19 +1,28 @@+{-# LANGUAGE MultilineStrings #-}+ module Main where import Control.Lens ((^.)) import Control.Monad (forM, unless) import Data.Aeson qualified as Aeson import Data.Generics.Labels ()+import Data.IORef (IORef, atomicModifyIORef', newIORef)+import Data.Int (Int32, Int64) import Data.Maybe (isNothing) import Data.Text (Text) import Data.Text qualified as T import Data.Time.Clock (getCurrentTime) import Data.Vector qualified as V+import Hasql.Decoders qualified as D import Hasql.Pool qualified as Pool+import Hasql.Session qualified as Session+import Hasql.Statement (Statement)+import Hasql.Statement qualified as Statement import Hasql.Transaction qualified as Tx import Hasql.Transaction.Sessions qualified as TxSessions import Kiroku.Store import Kiroku.Store.SQL qualified as SQL+import Kiroku.Test.Fixtures.CategoryScaling (categoryScalingFixtureSql, categoryScalingHead) import Kiroku.Test.Postgres (withMigratedTestDatabase, withSharedMigratedPostgres) import Test.Tasty (localOption) import Test.Tasty.Bench@@ -24,33 +33,307 @@ withMigratedTestDatabase $ \controlConnectionString -> withMigratedTestDatabase $ \candidateConnectionString -> withStore (defaultConnectionSettings controlConnectionString) $ \controlStore ->- withStore (defaultConnectionSettings candidateConnectionString) $ \candidateStore -> do- let fourStreams = namedStreams "workload-gate-4" 4- eightStreams = namedStreams "workload-gate-8" 8- seedStreams controlStore (fourStreams <> eightStreams)- seedStreams candidateStore (fourStreams <> eightStreams)+ withStore (defaultConnectionSettings candidateConnectionString) $ \candidateStore ->+ withAppendCategoryStores $ \appendControlStore appendCandidateStore ->+ withCategoryScalingStore $ \scalingStore -> do+ let fourStreams = namedStreams "workload-gate-4" 4+ eightStreams = namedStreams "workload-gate-8" 8+ seedStreams controlStore (fourStreams <> eightStreams)+ seedStreams candidateStore (fourStreams <> eightStreams) - runSequentialMultiAppend controlStore fourStreams- runProductionMultiAppend candidateStore fourStreams- runSequentialMultiAppend controlStore eightStreams- runProductionMultiAppend candidateStore eightStreams+ runSequentialMultiAppend controlStore fourStreams+ runProductionMultiAppend candidateStore fourStreams+ runSequentialMultiAppend controlStore eightStreams+ runProductionMultiAppend candidateStore eightStreams - defaultMain- [ localOption WallTime $- bgroup- "append-multi-stream"- [ bench "sequential-control-4" $- whnfIO (runSequentialMultiAppend controlStore fourStreams)- , bcompareWithin 0 0.90 "sequential-control-4" $- bench "production-pipeline-4" $- whnfIO (runProductionMultiAppend candidateStore fourStreams)- , bench "sequential-control-8" $- whnfIO (runSequentialMultiAppend controlStore eightStreams)- , bcompareWithin 0 0.90 "sequential-control-8" $- bench "production-pipeline-8" $- whnfIO (runProductionMultiAppend candidateStore eightStreams)+ appendControlCounter <- newIORef 0+ appendCandidateCounter <- newIORef 0+ runAppendWorkload preCategoryAppendAnyVersion appendControlStore appendControlCounter+ runAppendWorkload SQL.appendAnyVersion appendCandidateStore appendCandidateCounter++ defaultMain+ [ localOption WallTime $+ bgroup+ "append-multi-stream"+ [ bench "sequential-control-4" $+ whnfIO (runSequentialMultiAppend controlStore fourStreams)+ , bcompareWithin 0 0.90 "sequential-control-4" $+ bench "production-pipeline-4" $+ whnfIO (runProductionMultiAppend candidateStore fourStreams)+ , bench "sequential-control-8" $+ whnfIO (runSequentialMultiAppend controlStore eightStreams)+ , bcompareWithin 0 0.90 "sequential-control-8" $+ bench "production-pipeline-8" $+ whnfIO (runProductionMultiAppend candidateStore eightStreams)+ ]+ , -- BUG-2 / plan 91 G4: carrying the category onto each $all+ -- row (migration 0012) adds a column and one partial-index+ -- insert per event. The control runs the pre-0012 append on+ -- a database without the index or CHECK.+ localOption WallTime $+ bgroup+ "append-category-column"+ [ bench "control-append-40" $+ whnfIO (runAppendWorkload preCategoryAppendAnyVersion appendControlStore appendControlCounter)+ , bcompareWithin 0 1.05 "control-append-40" $+ bench "candidate-append-40" $+ whnfIO (runAppendWorkload SQL.appendAnyVersion appendCandidateStore appendCandidateCounter)+ ]+ , -- BUG-2 / plan 91 G3: the index-range category reads against+ -- the LATERAL statements they replaced, on one database seeded+ -- with the category-scaling fixture. The unpartitioned read+ -- must not be slower where LATERAL was already cheap, and both+ -- caught-up polls on 20,000 streams must be 5x faster.+ localOption WallTime $+ bgroup+ "category-read"+ [ bench "control-exhausted-category" $+ whnfIO (runPlainReads lateralCategoryRead scalingStore "performance" categoryScalingHead)+ , bcompareWithin 0 1.05 "control-exhausted-category" $+ bench "candidate-exhausted-category" $+ whnfIO (runPlainReads SQL.readCategoryForwardStmt scalingStore "performance" categoryScalingHead)+ , bench "control-page-200-streams-from-0" $+ whnfIO (runPlainReads lateralCategoryRead scalingStore "performance" 0)+ , bcompareWithin 0 1.05 "control-page-200-streams-from-0" $+ bench "candidate-page-200-streams-from-0" $+ whnfIO (runPlainReads SQL.readCategoryForwardStmt scalingStore "performance" 0)+ , bench "control-page-20000-streams-from-0" $+ whnfIO (runPlainReads lateralCategoryRead scalingStore "idle" 0)+ , bcompareWithin 0 1.05 "control-page-20000-streams-from-0" $+ bench "candidate-page-20000-streams-from-0" $+ whnfIO (runPlainReads SQL.readCategoryForwardStmt scalingStore "idle" 0)+ , bench "control-plain-caught-up-20000-streams" $+ whnfIO (runPlainReads lateralCategoryRead scalingStore "idle" categoryScalingHead)+ , bcompareWithin 0 0.20 "control-plain-caught-up-20000-streams" $+ bench "candidate-plain-caught-up-20000-streams" $+ whnfIO (runPlainReads SQL.readCategoryForwardStmt scalingStore "idle" categoryScalingHead)+ , bench "control-group-caught-up-20000-streams" $+ whnfIO (runGroupReads lateralCategoryGroupRead scalingStore "idle" categoryScalingHead)+ , bcompareWithin 0 0.20 "control-group-caught-up-20000-streams" $+ bench "candidate-group-caught-up-20000-streams" $+ whnfIO (runGroupReads SQL.readCategoryForwardConsumerGroupStmt scalingStore "idle" categoryScalingHead)+ ] ]- ]++-- | A migrated store seeded with 'categoryScalingFixtureSql'.+withCategoryScalingStore :: (KirokuStore -> IO a) -> IO a+withCategoryScalingStore action =+ withMigratedTestDatabase $ \connectionString ->+ withStore (defaultConnectionSettings connectionString) $ \store -> do+ seeded <- Pool.use (store ^. #pool) (Session.script categoryScalingFixtureSql)+ case seeded of+ Left err -> error ("category-read gate setup failed: " <> show err)+ Right () -> action store++{- | Ten executions of an unpartitioned category read (limit 100). Ten, not a+hundred: a LATERAL control poll on 20,000 streams costs about 15 ms, and a+hundred per iteration exceeds tasty-bench's timeout.+-}+runPlainReads ::+ Statement (Int64, Text, Int32) (V.Vector RecordedEvent) ->+ KirokuStore ->+ Text ->+ Int64 ->+ IO ()+runPlainReads statement store category cursor =+ mapM_+ (\_ -> Pool.use (store ^. #pool) (Session.statement (cursor, category, 100) statement) >>= forceReads)+ [1 .. 10 :: Int]++-- | Ten executions of a consumer-group category read, member 1 of 2 (limit 100).+runGroupReads ::+ Statement (Int64, Text, Int32, Int32, Int32) (V.Vector RecordedEvent) ->+ KirokuStore ->+ Text ->+ Int64 ->+ IO ()+runGroupReads statement store category cursor =+ mapM_+ (\_ -> Pool.use (store ^. #pool) (Session.statement (cursor, category, 1, 2, 100) statement) >>= forceReads)+ [1 .. 10 :: Int]++forceReads :: Either Pool.UsageError (V.Vector RecordedEvent) -> IO ()+forceReads (Right events) = V.length events `seq` pure ()+forceReads (Left err) = error ("category-read gate read failed: " <> show err)++-- | 'SQL.readCategoryForwardStmt' as it was before plan 91 (git 12d50d5).+lateralCategoryRead :: Statement (Int64, Text, Int32) (V.Vector RecordedEvent)+lateralCategoryRead =+ Statement.preparable+ lateralCategoryReadSQL+ SQL.readCategoryEncoder+ (D.rowVector SQL.recordedEventRow)++-- | 'SQL.readCategoryForwardConsumerGroupStmt' as it was before plan 91 (git 12d50d5).+lateralCategoryGroupRead :: Statement (Int64, Text, Int32, Int32, Int32) (V.Vector RecordedEvent)+lateralCategoryGroupRead =+ Statement.preparable+ lateralCategoryGroupReadSQL+ SQL.readCategoryConsumerGroupEncoder+ (D.rowVector SQL.recordedEventRow)++lateralCategoryReadSQL :: Text+lateralCategoryReadSQL =+ """+ 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 streams s+ JOIN LATERAL (+ SELECT se.*+ FROM stream_events se+ WHERE se.stream_id = 0+ AND se.original_stream_id = s.stream_id+ AND se.stream_version > $1+ ORDER BY se.stream_version ASC+ LIMIT $3+ ) se ON true+ JOIN events e ON e.event_id = se.event_id+ WHERE s.category = $2+ ORDER BY se.stream_version ASC+ LIMIT $3+ """++lateralCategoryGroupReadSQL :: Text+lateralCategoryGroupReadSQL =+ """+ 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 streams s+ JOIN LATERAL (+ SELECT se.*+ FROM stream_events se+ WHERE se.stream_id = 0+ AND se.original_stream_id = s.stream_id+ AND se.stream_version > $1+ ORDER BY se.stream_version ASC+ LIMIT $5+ ) se ON true+ JOIN events e ON e.event_id = se.event_id+ WHERE s.category = $2+ AND (((hashtextextended(s.stream_id::text, 0) % $4) + $4) % $4) = $3+ ORDER BY se.stream_version ASC+ LIMIT $5+ """++{- | Two freshly migrated stores for the append-category-column gate. The+control database has the category index and CHECK from migration 0012+dropped, so it pays exactly the pre-0012 write cost when driven by+'preCategoryAppendAnyVersion'; the candidate is left as migrated.+-}+withAppendCategoryStores :: (KirokuStore -> KirokuStore -> IO a) -> IO a+withAppendCategoryStores action =+ withMigratedTestDatabase $ \controlConnectionString ->+ withMigratedTestDatabase $ \candidateConnectionString ->+ withStore (defaultConnectionSettings controlConnectionString) $ \controlStore ->+ withStore (defaultConnectionSettings candidateConnectionString) $ \candidateStore -> do+ dropped <-+ Pool.use+ (controlStore ^. #pool)+ ( Session.script+ """+ DROP INDEX kiroku.ix_stream_events_all_by_category;+ ALTER TABLE kiroku.stream_events DROP CONSTRAINT ck_stream_events_all_category;+ """+ )+ case dropped of+ Left err -> error ("append-category control setup failed: " <> show err)+ Right () -> action controlStore candidateStore++{- | One gate iteration: 20 single-event appends to fresh streams, then 20 to+one hot stream. Each store gets its own counter, so control and candidate+create the same stream names and do the same work.+-}+runAppendWorkload ::+ Statement SQL.AppendParams (Maybe AppendResult) ->+ KirokuStore ->+ IORef Int ->+ IO ()+runAppendWorkload statement store counter = do+ iteration <- atomicModifyIORef' counter (\n -> (n + 1, n))+ now <- getCurrentTime+ let freshNames =+ [ "append-gate-" <> T.pack (show iteration) <> "-" <> T.pack (show index)+ | index <- [1 .. 20 :: Int]+ ]+ names = freshNames <> replicate 20 "append-gate-hot"+ mapM_+ ( \name -> do+ enriched <- enrichEvents (store ^. #storeSettings) [makeEvent "AppendCategoryGate"]+ prepared <- prepareEvents enriched+ result <-+ Pool.use (store ^. #pool) $+ Session.statement (buildAppendParams name now prepared) statement+ case result of+ Right (Just appendResult) -> forceAppendResults [appendResult]+ Right Nothing -> error "append-category gate append returned no row"+ Left err -> error ("append-category gate append failed: " <> show err)+ )+ names++-- | 'SQL.appendAnyVersion' as it was before migration 0012 (git 12d50d5).+preCategoryAppendAnyVersion :: Statement SQL.AppendParams (Maybe AppendResult)+preCategoryAppendAnyVersion =+ Statement.preparable+ preCategoryAppendAnyVersionSQL+ SQL.appendParamsEncoder+ SQL.appendResultDecoder++preCategoryAppendAnyVersionSQL :: Text+preCategoryAppendAnyVersionSQL =+ """+ 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, 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)+ SELECT ne.event_id, 0, au.initial_global_version + ne.idx, su.stream_id, su.initial_version + ne.idx+ 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+ """ namedStreams :: Text -> Int -> [(StreamName, Text)] namedStreams prefix count =
kiroku-store.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.0 name: kiroku-store-version: 0.8.0.2+version: 0.9.0.0 synopsis: High-performance PostgreSQL event store description: Kiroku is a PostgreSQL-backed event store for Haskell applications. It@@ -209,6 +209,7 @@ , aeson >=2.1 && <2.3 , base >=4.18 && <5 , 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
src/Kiroku/Store/SQL.hs view
@@ -3,6 +3,8 @@ module Kiroku.Store.SQL ( -- * Append statements AppendParams (..),+ appendParamsEncoder,+ appendResultDecoder, appendExpectedVersion, appendStreamExists, appendNoStream,@@ -17,6 +19,7 @@ readAllForwardStmt, readAllBackwardStmt, readCategoryForwardStmt,+ readCategoryEncoder, getStreamStmt, eventExistsInStreamStmt, lookupStreamNamesStmt,@@ -25,8 +28,12 @@ -- * Consumer-group read statements readCategoryForwardConsumerGroupStmt,+ readCategoryConsumerGroupEncoder, readAllForwardConsumerGroupStmt, + -- * Row decoders+ recordedEventRow,+ -- * Causation / correlation statements findByCorrelationStmt, findCausationDescendantsStmt,@@ -179,7 +186,7 @@ WHERE stream_name = $8 AND stream_version = $9 AND deleted_at IS NULL- RETURNING stream_id, stream_version - (SELECT count(*) FROM new_events) AS initial_version+ 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)@@ -202,8 +209,8 @@ 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)- SELECT ne.event_id, 0, au.initial_global_version + ne.idx, su.stream_id, su.initial_version + ne.idx+ 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@@ -230,7 +237,7 @@ SET stream_version = stream_version + (SELECT count(*) FROM new_events) WHERE stream_name = $8 AND deleted_at IS NULL- RETURNING stream_id, stream_version - (SELECT count(*) FROM new_events) AS initial_version+ 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)@@ -253,8 +260,8 @@ 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)- SELECT ne.event_id, 0, au.initial_global_version + ne.idx, su.stream_id, su.initial_version + ne.idx+ 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@@ -280,7 +287,7 @@ INSERT INTO streams (stream_name, stream_version) VALUES ($8, (SELECT count(*) FROM new_events)) ON CONFLICT (stream_name) DO NOTHING- RETURNING stream_id, 0::bigint AS initial_version+ 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)@@ -303,8 +310,8 @@ 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)- SELECT ne.event_id, 0, au.initial_global_version + ne.idx, si.stream_id, si.initial_version + ne.idx+ 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@@ -336,7 +343,7 @@ 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, stream_version - (SELECT count(*) FROM new_events) AS initial_version+ 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)@@ -359,8 +366,8 @@ 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)- SELECT ne.event_id, 0, au.initial_global_version + ne.idx, su.stream_id, su.initial_version + ne.idx+ 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@@ -793,7 +800,14 @@ -- Category Read Statements -- --------------------------------------------------------------------------- --- | Read events from streams matching a category, in global position order.+{- | Read events from streams matching a category, in global position order.++Scans @ix_stream_events_all_by_category@ from @(category, startPosition)@ and+stops at the limit, so a poll's cost follows the rows it returns, not the+number of streams in the category (BUG-2). The @category@ column on @$all@+junction rows is written by the append statements (migration @0012@). Params:+@(startPosition, category, limit)@.+-} readCategoryForwardStmt :: Statement (Int64, Text, Int32) (Vector RecordedEvent) readCategoryForwardStmt = preparable@@ -816,18 +830,11 @@ se.original_stream_id, se.original_stream_version, e.data, e.metadata, e.causation_id, e.correlation_id, e.created_at- FROM streams s- JOIN LATERAL (- SELECT se.*- FROM stream_events se- WHERE se.stream_id = 0- AND se.original_stream_id = s.stream_id- AND se.stream_version > $1- ORDER BY se.stream_version ASC- LIMIT $3- ) se ON true+ FROM stream_events se JOIN events e ON e.event_id = se.event_id- WHERE s.category = $2+ WHERE se.stream_id = 0+ AND se.category = $2+ AND se.stream_version > $1 ORDER BY se.stream_version ASC LIMIT $3 """@@ -845,9 +852,12 @@ @member_of(stream_id) = (((hashtextextended(stream_id::text, 0) % size) + size) % size)@ -The predicate is applied to @s.stream_id@ in the outer @WHERE@ so whole-unassigned streams are pruned before the lateral join. Params:-@(startPosition, category, member, size, limit)@.+The read scans @ix_stream_events_all_by_category@ from+@(category, startPosition)@ and applies the predicate to+@se.original_stream_id@, which the index carries as an @INCLUDE@ column, so+other members' rows are skipped on index tuples. A poll touches about+@limit * size@ index entries at most, never one probe per stream in the+category (BUG-2). Params: @(startPosition, category, member, size, limit)@. -} readCategoryForwardConsumerGroupStmt :: Statement (Int64, Text, Int32, Int32, Int32) (Vector RecordedEvent)@@ -874,19 +884,12 @@ se.original_stream_id, se.original_stream_version, e.data, e.metadata, e.causation_id, e.correlation_id, e.created_at- FROM streams s- JOIN LATERAL (- SELECT se.*- FROM stream_events se- WHERE se.stream_id = 0- AND se.original_stream_id = s.stream_id- AND se.stream_version > $1- ORDER BY se.stream_version ASC- LIMIT $5- ) se ON true+ FROM stream_events se JOIN events e ON e.event_id = se.event_id- WHERE s.category = $2- AND (((hashtextextended(s.stream_id::text, 0) % $4) + $4) % $4) = $3+ 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 """
src/Kiroku/Store/Subscription.hs view
@@ -140,9 +140,12 @@ (queueCapacity config) (overflowPolicy config) pure (LiveFromPublisherQueue queue statusVar, unsubscribe)- (Nothing, Category (CategoryName cat)) ->+ -- Plain categories and consumer-group category members both wake+ -- on the category's NOTIFY generation; a member's fetch applies+ -- its partition predicate in SQL.+ (_, Category (CategoryName cat)) -> pure (LiveFromCategoryNotify cat, pure ())- (Just _, _) ->+ (Just _, AllStreams) -> pure (LiveFromGroupPolling, pure ()) ) (\(_, unsubscribe) -> unsubscribe)
src/Kiroku/Store/Subscription/Worker.hs view
@@ -174,12 +174,13 @@ TVar carries Paused/Overflowed backpressure signals. -} LiveFromPublisherQueue !(TBQueue (Vector RecordedEvent)) !(TVar SubscriberStatus)- | {- | Non-group Category: wake on the named category's NOTIFY generation- counter and re-query the database.+ | {- | Category, plain or consumer-group member: wake on the named+ category's NOTIFY generation counter and re-query the database (with the+ partition predicate, for a member). -} LiveFromCategoryNotify !Text- | {- | Consumer-group member, for either target: wake when the global- position advances and re-query with the partition predicate.+ | {- | Consumer-group member of AllStreams: wake when the global position+ advances and re-query with the partition predicate. -} LiveFromGroupPolling @@ -535,10 +536,14 @@ -- 30s safety poll) reconciles notifications lost while the listener connection is -- reconnecting, preserving at-least-once delivery with bounded latency. ----- This loop serves only non-group `Category` subscriptions. Consumer-group members--- cannot use the per-category signal: their interest is--- `hashtextextended(stream_id) % size = member`, a Postgres hash the worker cannot--- cheaply replicate from the payload, so they stay on `liveLoopDbDriven`.+-- This loop serves every `Category` subscription, plain or consumer-group. A+-- member cannot tell from a NOTIFY payload whether the stream is in its slice+-- (that is `hashtextextended(stream_id) % size = member`, a Postgres hash), but+-- it can gate on the category: the category generation advances on every append+-- to a stream of the category, a superset of the member's own streams, and+-- `fetchBatch` applies the partition predicate in SQL. So a member of an idle+-- category does no live database work while other categories are busy, and a+-- member whose sibling received the append does one empty fetch. liveLoopCategoryNotify :: Pool -> SubscriptionConfig ->@@ -588,7 +593,7 @@ Nothing -> pure (Right Nothing) -- handler said Stop Just newPos -> drainTo newPos --- Phase 2: live (DB-driven, consumer-group members only). Bypasses the broadcast+-- Phase 2: live (DB-driven, consumer-group members of AllStreams). Bypasses the broadcast -- and re-queries the database when the publisher's GLOBAL position advances, -- letting `fetchBatch` apply the partition predicate baked into the consumer-group -- SQL. A partitioned member cannot read the broadcast `liveQueue` because it
test/Test/CategoryIdleNoSpin.hs view
@@ -12,8 +12,12 @@ blocks on its own per-category NOTIFY generation and does __zero__ database fetches while a different category receives sustained traffic. A real append to the subscribed category still wakes it (liveness).- * Consumer-group member- ('Kiroku.Store.Subscription.Worker.liveLoopDbDriven', corrected gate): an idle+ * Consumer-group members of a @Category@ (plan 91 M5): they share the plain+ category loop and its per-category generation, so both members of an idle+ category do __zero__ fetches while another category is busy, and an append+ to the category wakes them and the owning member delivers it.+ * Consumer-group member of @AllStreams@+ ('Kiroku.Store.Subscription.Worker.liveLoopDbDriven', corrected gate): the member gates on the /last observed global position/ rather than its per-partition cursor, so it wakes at most once per global advance — a bounded number of fetches — instead of the unbounded spin the original cursor-gate@@ -27,8 +31,10 @@ see the plan's Decision Log.) To confirm these specs actually pin the regression: temporarily restore the old-cursor-gated 'liveLoopDbDriven' body and route @(Nothing, Category{})@ back through-it — both idle fetch counts then explode and the assertions fail.+cursor-gated 'liveLoopDbDriven' body and route @Category@ subscriptions back+through it — the idle fetch counts then explode and the assertions fail. Routing+only @(Just _, Category{})@ back through the corrected 'liveLoopDbDriven' fails+the zero-fetch assertion for group category members. -} module Test.CategoryIdleNoSpin (spec) where @@ -112,9 +118,64 @@ deliveredFinal <- readTVarIO deliveredVar deliveredFinal `shouldBe` 1 - it "an idle consumer-group member does not spin while a different category advances the global position" $ do- let subName = SubscriptionName "grp-sub"+ it "idle consumer-group category members do zero fetches while another category is active, then the owner wakes on its own event" $ do+ -- Two members (0 and 1 of size 2) over category "quiet". Each runs under+ -- its own subscription name so each has its own live barrier and fetch+ -- counter; the partition predicate depends only on (member, size).+ let names = [SubscriptionName "quiet-grp-0", SubscriptionName "quiet-grp-1"] deliveredVar <- newTVarIO (0 :: Int)+ fetchVars <- mapM (const (newTVarIO (0 :: Int))) names+ barriers <- mapM (const newEmptyMVar) names+ let countFetch evt = case evt of+ KirokuEventSubscriptionFetched n _ _ ->+ sequence_+ [ atomically (modifyTVar' v (+ 1))+ | (subName, v) <- zip names fetchVars+ , n == subName+ ]+ _ -> pure ()+ obsHandler evt = do+ sequence_ [caughtUpEventHandler subName barrier Nothing evt | (subName, barrier) <- zip names barriers]+ countFetch evt+ deliver _evt = do+ atomically (modifyTVar' deliveredVar (+ 1))+ pure Continue+ memberConfig subName member =+ (defaultSubscriptionConfig subName (Category (CategoryName "quiet")) deliver)+ { consumerGroup = Just ConsumerGroup{member = member, size = 2}+ }+ withTestStoreSettings (\s -> s & #eventHandler .~ Just obsHandler) $ \store ->+ bracket (subscribe store (memberConfig (names !! 0) 0)) cancel $ \_ ->+ bracket (subscribe store (memberConfig (names !! 1) 1)) cancel $ \_ -> do+ -- Settle: each member reaches live mode and finishes its initial+ -- post-catch-up drain (one empty fetch).+ mapM_ waitForSubscriptionLive barriers+ mapM_ (\v -> waitUntil 5_000_000 ((>= 1) <$> readTVarIO v)) fetchVars+ bases <- mapM readTVarIO fetchVars++ -- Drive a different category. Before plan 91 M5 each member woke on+ -- every global advance and ran one empty fetch per wake; now both+ -- wait on the "quiet" category's generation, which never moves.+ let busyCount = 20 :: Int+ busyStreams = ["busy-" <> T.pack (show i) | i <- [1 .. busyCount]]+ appendEach store busyStreams "Busy"+ waitForPublisher store (GlobalPosition (fromIntegral busyCount))+ threadDelay 500_000++ afterIdle <- mapM readTVarIO fetchVars+ zipWith (-) afterIdle bases `shouldBe` [0, 0]+ readTVarIO deliveredVar `shouldReturn` 0++ -- Liveness: an append to the quiet category wakes both members;+ -- exactly the owner of "quiet-1" delivers it.+ appendEach store ["quiet-1"] "Quiet"+ waitUntil 5_000_000 ((>= 1) <$> readTVarIO deliveredVar)+ threadDelay 200_000+ readTVarIO deliveredVar `shouldReturn` 1++ it "an idle AllStreams consumer-group member wakes a bounded number of times, not a busy spin" $ do+ let subName = SubscriptionName "grp-all-sub"+ deliveredVar <- newTVarIO (0 :: Int) fetchVar <- newTVarIO (0 :: Int) liveBarrier <- newEmptyMVar let countFetch evt = case evt of@@ -126,18 +187,16 @@ let deliver _evt = do atomically (modifyTVar' deliveredVar (+ 1)) pure Continue- -- Member 0 of 3 over category "grp" — which receives NO events, so- -- this member's partition fetch is always empty. The flood lands in- -- a different category, advancing only the global position.+ -- Member 0 of 3 over $all. Its fetches apply the partition+ -- predicate, so it receives only its slice of the flood, but every+ -- append advances the global position it gates on. cfg =- (defaultSubscriptionConfig subName (Category (CategoryName "grp")) deliver)+ (defaultSubscriptionConfig subName AllStreams deliver) { consumerGroup = Just ConsumerGroup{member = 0, size = 3} } bracket (subscribe store cfg) cancel $ \_handle -> do waitForSubscriptionLive liveBarrier - -- The corrected group loop gates BEFORE draining, so on an empty- -- store it blocks with zero fetches until the global position moves. let floodCount = 20 :: Int floodStreams = ["flood-" <> T.pack (show i) | i <- [1 .. floodCount]] appendEach store floodStreams "Flood"@@ -145,9 +204,11 @@ threadDelay 500_000 afterIdle <- readTVarIO fetchVar- deliveredIdle <- readTVarIO deliveredVar+ delivered <- readTVarIO deliveredVar -- The corrected gate wakes at most once per observed global position- -- (<= floodCount), bounded — NOT the unbounded busy-spin of the old- -- cursor gate, which racks up thousands of empty fetches in 500ms.+ -- and drains to empty, so fetches stay bounded by the flood size+ -- (plus the drain's empty fetch per wake) — NOT the unbounded+ -- busy-spin of the old cursor gate, which racks up thousands of+ -- empty fetches in 500ms. afterIdle `shouldSatisfy` (< 50)- deliveredIdle `shouldBe` 0+ delivered `shouldSatisfy` (<= floodCount)
test/Test/PerformanceStructure.hs view
@@ -21,13 +21,15 @@ import Hasql.Statement qualified as Statement import Kiroku.Store import Kiroku.Store.SQL qualified as SQL-import Test.Helpers (withTestStore, withTestStoreSettings)+import Kiroku.Test.Fixtures.CategoryScaling (categoryScalingFixtureSql, categoryScalingHead)+import Test.Helpers (makeEvent, withTestStore, withTestStoreSettings) import Test.Hspec spec :: Spec spec = do noOpAppendSpec queryPlanSpec+ categoryReadCostSpec noOpAppendSpec :: Spec noOpAppendSpec =@@ -88,17 +90,6 @@ expectIndex "ux_stream_events_stream_version" plan expectNoNodeType "Sort" plan - it "category high-cursor reads use ix_stream_events_all_by_origin" $ \store -> do- plan <-- explainProductionStatement- store- SQL.readCategoryForwardStmt- [ ("$3", "100::int4")- , ("$2", "'performance'::text")- , ("$1", "15000::bigint")- ]- expectIndex "ix_stream_events_all_by_origin" plan- it "dead-letter reads use ix_dead_letters_subscription_position without Sort" $ \store -> do plan <- explainProductionStatement@@ -128,6 +119,125 @@ Left err -> expectationFailure ("could not inspect retention triggers: " <> show err) Right shape -> shape `shouldBe` (6, 0) +{- | BUG-2. A category read's cost must follow the rows it returns, not the+number of streams in the category. The store holds the category-scaling+fixture: @performance@ with 200 streams, @idle@ with 20,000 one-event streams,+and @noise@, head at 'categoryScalingHead'.+-}+categoryReadCostSpec :: Spec+categoryReadCostSpec =+ describe "category read cost" $+ aroundAll withCategoryScalingStore $ do+ it "category high-cursor reads use ix_stream_events_all_by_category without Sort" $ \store -> do+ plan <-+ explainProductionStatement+ store+ SQL.readCategoryForwardStmt+ [ ("$3", "100::int4")+ , ("$2", "'performance'::text")+ , ("$1", "15000::bigint")+ ]+ expectIndex "ix_stream_events_all_by_category" plan+ expectNoNodeType "Sort" plan++ it "consumer-group category reads use ix_stream_events_all_by_category without Sort" $ \store -> do+ plan <-+ explainProductionStatement+ store+ SQL.readCategoryForwardConsumerGroupStmt+ [ ("$5", "100::int4")+ , ("$4", "2::int4")+ , ("$3", "1::int4")+ , ("$2", "'performance'::text")+ , ("$1", "15000::bigint")+ ]+ expectIndex "ix_stream_events_all_by_category" plan+ expectNoNodeType "Sort" plan++ it "category caught-up poll on 20000 idle streams reads at most 32 buffers" $ \store -> do+ let cursor = T.pack (show categoryScalingHead) <> "::bigint"+ plainAt =+ [ ("$3", "100::int4")+ , ("$2", "'idle'::text")+ , ("$1", cursor)+ ]+ groupAt member =+ [ ("$5", "100::int4")+ , ("$4", "2::int4")+ , ("$3", T.pack (show member) <> "::int4")+ , ("$2", "'idle'::text")+ , ("$1", cursor)+ ]+ groupMembers = [0, 1 :: Int]++ plainIdle <- explainAnalyzeBuffers store SQL.readCategoryForwardStmt plainAt+ expectBufferBudget "plain caught-up poll" plainIdle+ snd plainIdle `shouldBe` 0+ groupIdle <- mapM (explainAnalyzeBuffers store SQL.readCategoryForwardConsumerGroupStmt . groupAt) groupMembers+ mapM_ (expectBufferBudget "group caught-up poll") groupIdle+ map snd groupIdle `shouldBe` [0, 0]++ appended <- runStoreIO store $ appendToStream (StreamName "idle-1") AnyVersion [makeEvent "IdleWake" Null]+ appended `shouldSatisfy` either (const False) (const True)++ plainOne <- explainAnalyzeBuffers store SQL.readCategoryForwardStmt plainAt+ expectBufferBudget "plain poll after one append" plainOne+ snd plainOne `shouldBe` 1+ groupOne <- mapM (explainAnalyzeBuffers store SQL.readCategoryForwardConsumerGroupStmt . groupAt) groupMembers+ mapM_ (expectBufferBudget "group poll after one append") groupOne+ -- Exactly one member owns idle-1.+ sum (map snd groupOne) `shouldBe` 1++withCategoryScalingStore :: (KirokuStore -> IO ()) -> IO ()+withCategoryScalingStore action =+ withTestStore $ \store -> do+ result <- Pool.use (store ^. #pool) (Session.script categoryScalingFixtureSql)+ case result of+ Left err -> expectationFailure ("failed to seed category-scaling fixture: " <> show err)+ Right () -> action store++-- | The G1 budget: a category poll that returns at most one row.+expectBufferBudget :: String -> (Int64, Int64) -> Expectation+expectBufferBudget label (buffers, _) =+ unless (buffers <= 32) $+ expectationFailure $+ label+ <> ": expected at most 32 shared buffers, but the plan read "+ <> show buffers++{- | Execute a production statement under @EXPLAIN (ANALYZE, BUFFERS)@ and+return the top plan node's shared buffers (hit plus read) and actual rows.+Planning buffers are excluded; they are not paid by a prepared statement.+-}+explainAnalyzeBuffers ::+ KirokuStore ->+ Statement params result ->+ [(Text, Text)] ->+ IO (Int64, Int64)+explainAnalyzeBuffers store productionStatement replacements = do+ plan <-+ explainWith+ "EXPLAIN (ANALYZE, BUFFERS, COSTS OFF, TIMING OFF, FORMAT JSON)\n"+ store+ productionStatement+ replacements+ case plan of+ Array entries+ | Just (Object entry) <- headMay entries+ , Just (Object top) <- KeyMap.lookup "Plan" entry ->+ pure+ ( numberField "Shared Hit Blocks" top + numberField "Shared Read Blocks" top+ , numberField "Actual Rows" top+ )+ _ -> expectationFailure ("unexpected EXPLAIN shape: " <> show plan) >> fail "unreachable"+ where+ headMay values = case foldr (:) [] values of+ value : _ -> Just value+ [] -> Nothing+ numberField key object = case KeyMap.lookup key object of+ Just value | Aeson.Success (number :: Double) <- Aeson.fromJSON value -> round number+ _ -> 0+ withObservedStore :: IORef Int -> (KirokuStore -> IO ()) -> IO () withObservedStore checkouts = withTestStoreSettings $ \settings ->@@ -181,12 +291,13 @@ RETURNING event_id ), all_links AS ( INSERT INTO stream_events- (event_id, stream_id, stream_version, original_stream_id, original_stream_version)+ (event_id, stream_id, stream_version, original_stream_id, original_stream_version, category) SELECT fixture.event_id, 0, fixture.global_position, fixture.stream_id,- fixture.stream_version+ fixture.stream_version,+ 'performance' FROM fixture_events AS fixture JOIN inserted_events USING (event_id) RETURNING event_id@@ -241,10 +352,18 @@ Statement params result -> [(Text, Text)] -> IO Value-explainProductionStatement store productionStatement replacements = do+explainProductionStatement = explainWith "EXPLAIN (FORMAT JSON, COSTS OFF)\n"++explainWith ::+ Text ->+ KirokuStore ->+ Statement params result ->+ [(Text, Text)] ->+ IO Value+explainWith explainPrefix store productionStatement replacements = do let productionSql = Statement.toSql productionStatement explainedSql =- "EXPLAIN (FORMAT JSON, COSTS OFF)\n"+ explainPrefix <> foldl' (\sql (placeholder, literal) -> T.replace placeholder literal sql) productionSql replacements explainStatement :: Statement () ByteString explainStatement =