kiroku-store 0.9.0.1 → 0.10.0.0
raw patch · 47 files changed
+2640/−311 lines, 47 filesdep ~basedep ~contravariant-extrasdep ~ephemeral-pgPVP ok
version bump matches the API change (PVP)
Dependency ranges changed: base, contravariant-extras, ephemeral-pg, generic-lens, hasql, hasql-pool, lens, text
API changes (from Hackage documentation)
- Kiroku.Store.Subscription.Stream: instance GHC.Internal.Exception.Type.Exception Kiroku.Store.Subscription.Stream.InvalidStreamBufferSize
- Kiroku.Store.Subscription.Types: ConsumerGroup :: !Int32 -> !Int32 -> ConsumerGroup
- Kiroku.Store.Subscription.Types: [member] :: ConsumerGroup -> !Int32
- Kiroku.Store.Subscription.Types: [size] :: ConsumerGroup -> !Int32
- Kiroku.Store.Subscription.Types: instance GHC.Internal.Exception.Type.Exception Kiroku.Store.Subscription.Types.InvalidConsumerGroup
+ Kiroku.Store: KirokuEventPublisherDecodeFailed :: !GlobalPosition -> !EventId -> !DecodeFailure -> KirokuEvent
+ Kiroku.Store: KirokuEventSubscriptionGroupSizeMismatch :: !ConsumerGroupSizeMismatch -> !SubscriptionGroupContext -> KirokuEvent
+ Kiroku.Store: KirokuEventSubscriptionHandlerStalled :: !SubscriptionName -> !GlobalPosition -> !EventId -> !NominalDiffTime -> !SubscriptionGroupContext -> KirokuEvent
+ Kiroku.Store: KirokuEventSubscriptionTargetBound :: !SubscriptionName -> !SubscriptionTarget -> !SubscriptionGroupContext -> KirokuEvent
+ Kiroku.Store: StopUndecodable :: !DecodeFailure -> SubscriptionStopReason
+ Kiroku.Store.Error: EventDecodeFailed :: !DecodeFailure -> StoreError
+ Kiroku.Store.Observability: DeadLetterDecodeFailure :: !DecodeFailure -> DeadLetterReason
+ Kiroku.Store.Observability: KirokuEventPublisherDecodeFailed :: !GlobalPosition -> !EventId -> !DecodeFailure -> KirokuEvent
+ Kiroku.Store.Observability: KirokuEventSubscriptionGroupSizeMismatch :: !ConsumerGroupSizeMismatch -> !SubscriptionGroupContext -> KirokuEvent
+ Kiroku.Store.Observability: KirokuEventSubscriptionHandlerStalled :: !SubscriptionName -> !GlobalPosition -> !EventId -> !NominalDiffTime -> !SubscriptionGroupContext -> KirokuEvent
+ Kiroku.Store.Observability: KirokuEventSubscriptionTargetBound :: !SubscriptionName -> !SubscriptionTarget -> !SubscriptionGroupContext -> KirokuEvent
+ Kiroku.Store.Observability: StopUndecodable :: !DecodeFailure -> SubscriptionStopReason
+ Kiroku.Store.SQL: [dlGroupSize] :: DeadLetterParams -> !Int32
+ Kiroku.Store.SQL: [dlTargetCategory] :: DeadLetterParams -> !Maybe Text
+ Kiroku.Store.SQL: [dlTargetKind] :: DeadLetterParams -> !Text
+ Kiroku.Store.SQL: saveAllCheckpointMemberStmt :: Statement (Text, Int32, Int64, Int32) ()
+ Kiroku.Store.SQL: saveCategoryCheckpointMemberStmt :: Statement (Text, Int32, Int64, Int32, Text) ()
+ Kiroku.Store.Settings: DecodeFailure :: !EventId -> !Text -> DecodeFailure
+ Kiroku.Store.Settings: Decoded :: !RecordedEvent -> DecodedEvent
+ Kiroku.Store.Settings: TransformedBatch :: !Vector DecodedEvent -> DecodedBatch
+ Kiroku.Store.Settings: UnchangedBatch :: !Vector RecordedEvent -> DecodedBatch
+ Kiroku.Store.Settings: Undecodable :: !RecordedEvent -> !DecodeFailure -> DecodedEvent
+ Kiroku.Store.Settings: [decodeFailureEventId] :: DecodeFailure -> !EventId
+ Kiroku.Store.Settings: [decodeFailureReason] :: DecodeFailure -> !Text
+ Kiroku.Store.Settings: data DecodeFailure
+ Kiroku.Store.Settings: data DecodedBatch
+ Kiroku.Store.Settings: data DecodedEvent
+ Kiroku.Store.Settings: decodeEvent :: StoreSettings -> RecordedEvent -> IO DecodedEvent
+ Kiroku.Store.Settings: decodedBatchLastEvent :: DecodedBatch -> RecordedEvent
+ Kiroku.Store.Settings: decodedBatchLength :: DecodedBatch -> Int
+ Kiroku.Store.Settings: decodedEventRecorded :: DecodedEvent -> RecordedEvent
+ Kiroku.Store.Settings: filterDecodedBatch :: (RecordedEvent -> Bool) -> DecodedBatch -> DecodedBatch
+ Kiroku.Store.Settings: instance GHC.Classes.Eq Kiroku.Store.Settings.DecodeFailure
+ Kiroku.Store.Settings: instance GHC.Classes.Eq Kiroku.Store.Settings.DecodedBatch
+ Kiroku.Store.Settings: instance GHC.Classes.Eq Kiroku.Store.Settings.DecodedEvent
+ Kiroku.Store.Settings: instance GHC.Internal.Generics.Generic Kiroku.Store.Settings.DecodeFailure
+ Kiroku.Store.Settings: instance GHC.Internal.Show.Show Kiroku.Store.Settings.DecodeFailure
+ Kiroku.Store.Settings: instance GHC.Internal.Show.Show Kiroku.Store.Settings.DecodedBatch
+ Kiroku.Store.Settings: instance GHC.Internal.Show.Show Kiroku.Store.Settings.DecodedEvent
+ Kiroku.Store.Subscription.Checkpoint: ConsumerGroupResizeReport :: !Vector Int32 -> !Int -> !ConsumerGroupSize -> !GlobalPosition -> ConsumerGroupResizeReport
+ Kiroku.Store.Subscription.Checkpoint: SubscriptionTargetRebindReport :: !Int -> !Vector (Maybe SubscriptionTarget) -> !SubscriptionTarget -> !GlobalPosition -> SubscriptionTargetRebindReport
+ Kiroku.Store.Subscription.Checkpoint: [newSize] :: ConsumerGroupResizeReport -> !ConsumerGroupSize
+ Kiroku.Store.Subscription.Checkpoint: [previousBindings] :: SubscriptionTargetRebindReport -> !Vector (Maybe SubscriptionTarget)
+ Kiroku.Store.Subscription.Checkpoint: [previousMemberCount] :: ConsumerGroupResizeReport -> !Int
+ Kiroku.Store.Subscription.Checkpoint: [previousSizes] :: ConsumerGroupResizeReport -> !Vector Int32
+ Kiroku.Store.Subscription.Checkpoint: [reboundMemberCount] :: SubscriptionTargetRebindReport -> !Int
+ Kiroku.Store.Subscription.Checkpoint: [reboundPosition] :: SubscriptionTargetRebindReport -> !GlobalPosition
+ Kiroku.Store.Subscription.Checkpoint: [reboundTarget] :: SubscriptionTargetRebindReport -> !SubscriptionTarget
+ Kiroku.Store.Subscription.Checkpoint: [resumePosition] :: ConsumerGroupResizeReport -> !GlobalPosition
+ Kiroku.Store.Subscription.Checkpoint: data ConsumerGroupResizeReport
+ Kiroku.Store.Subscription.Checkpoint: data SubscriptionTargetRebindReport
+ Kiroku.Store.Subscription.Checkpoint: instance GHC.Classes.Eq Kiroku.Store.Subscription.Checkpoint.ConsumerGroupResizeReport
+ Kiroku.Store.Subscription.Checkpoint: instance GHC.Classes.Eq Kiroku.Store.Subscription.Checkpoint.SubscriptionTargetRebindReport
+ Kiroku.Store.Subscription.Checkpoint: instance GHC.Internal.Generics.Generic Kiroku.Store.Subscription.Checkpoint.ConsumerGroupResizeReport
+ Kiroku.Store.Subscription.Checkpoint: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Checkpoint.ConsumerGroupResizeReport
+ Kiroku.Store.Subscription.Checkpoint: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Checkpoint.SubscriptionTargetRebindReport
+ Kiroku.Store.Subscription.Checkpoint: rebindSubscriptionTargetTx :: SubscriptionName -> SubscriptionTarget -> GlobalPosition -> Transaction SubscriptionTargetRebindReport
+ Kiroku.Store.Subscription.Checkpoint: resizeConsumerGroupTx :: SubscriptionName -> ConsumerGroupSize -> Transaction ConsumerGroupResizeReport
+ Kiroku.Store.Subscription.Fsm: DeadLetterDecodeFailure :: !DecodeFailure -> DeadLetterReason
+ Kiroku.Store.Subscription.Fsm: StopUndecodable :: !DecodeFailure -> SubscriptionStopReason
+ Kiroku.Store.Subscription.Stream: data StreamBufferSize
+ Kiroku.Store.Subscription.Stream: defaultStreamBufferSize :: StreamBufferSize
+ Kiroku.Store.Subscription.Stream: instance GHC.Classes.Eq Kiroku.Store.Subscription.Stream.StreamBufferSize
+ Kiroku.Store.Subscription.Stream: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Stream.StreamBufferSize
+ Kiroku.Store.Subscription.Stream: mkStreamBufferSize :: Natural -> Either InvalidStreamBufferSize StreamBufferSize
+ Kiroku.Store.Subscription.Stream: streamBufferSizeValue :: StreamBufferSize -> Natural
+ Kiroku.Store.Subscription.Types: AdoptUnbound :: TargetBindingPolicy
+ Kiroku.Store.Subscription.Types: ConsumerGroupSizeMismatch :: !SubscriptionName -> !Int32 -> !Vector Int32 -> ConsumerGroupSizeMismatch
+ Kiroku.Store.Subscription.Types: DeadLetterDecodeFailure :: !DecodeFailure -> DeadLetterReason
+ Kiroku.Store.Subscription.Types: InvalidBatchSize :: Int32 -> InvalidBatchSize
+ Kiroku.Store.Subscription.Types: InvalidHandlerStallWarnAfter :: NominalDiffTime -> InvalidHandlerStallWarnAfter
+ Kiroku.Store.Subscription.Types: RequireBound :: TargetBindingPolicy
+ Kiroku.Store.Subscription.Types: SomeSubscriptionStartupFailure :: e -> SomeSubscriptionStartupFailure
+ Kiroku.Store.Subscription.Types: SubscriptionTargetMismatch :: !SubscriptionName -> !SubscriptionTarget -> !Vector (Maybe SubscriptionTarget) -> SubscriptionTargetMismatch
+ Kiroku.Store.Subscription.Types: SubscriptionUndecodable :: DecodeFailure -> SubscriptionUndecodable
+ Kiroku.Store.Subscription.Types: [configuredSize] :: ConsumerGroupSizeMismatch -> !Int32
+ Kiroku.Store.Subscription.Types: [configuredTarget] :: SubscriptionTargetMismatch -> !SubscriptionTarget
+ Kiroku.Store.Subscription.Types: [handlerStallWarnAfter] :: SubscriptionConfigM (m :: Type -> Type) -> !Maybe NominalDiffTime
+ Kiroku.Store.Subscription.Types: [mismatchName] :: ConsumerGroupSizeMismatch -> !SubscriptionName
+ Kiroku.Store.Subscription.Types: [observedSizes] :: ConsumerGroupSizeMismatch -> !Vector Int32
+ Kiroku.Store.Subscription.Types: [observedTargets] :: SubscriptionTargetMismatch -> !Vector (Maybe SubscriptionTarget)
+ Kiroku.Store.Subscription.Types: [subscriptionDecodeFailure] :: SubscriptionUndecodable -> DecodeFailure
+ Kiroku.Store.Subscription.Types: [targetBindingPolicy] :: SubscriptionConfigM (m :: Type -> Type) -> !TargetBindingPolicy
+ Kiroku.Store.Subscription.Types: [targetMismatchName] :: SubscriptionTargetMismatch -> !SubscriptionName
+ Kiroku.Store.Subscription.Types: [undecodableHandler] :: SubscriptionConfigM (m :: Type -> Type) -> !Maybe (RecordedEvent -> DecodeFailure -> m SubscriptionResult)
+ Kiroku.Store.Subscription.Types: batchSizeValue :: BatchSize -> Int32
+ Kiroku.Store.Subscription.Types: consumerGroupSizeValue :: ConsumerGroupSize -> Int32
+ Kiroku.Store.Subscription.Types: data BatchSize
+ Kiroku.Store.Subscription.Types: data ConsumerGroupSize
+ Kiroku.Store.Subscription.Types: data ConsumerGroupSizeMismatch
+ Kiroku.Store.Subscription.Types: data SomeSubscriptionStartupFailure
+ Kiroku.Store.Subscription.Types: data SubscriptionTargetMismatch
+ Kiroku.Store.Subscription.Types: data TargetBindingPolicy
+ Kiroku.Store.Subscription.Types: defaultBatchSize :: BatchSize
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.BatchSize
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.ConsumerGroupGuardConflict
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.ConsumerGroupSize
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.ConsumerGroupSizeMismatch
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.InvalidBatchSize
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.InvalidConsumerGroup
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.InvalidHandlerStallWarnAfter
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.SubscriptionTargetMismatch
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.SubscriptionUndecodable
+ Kiroku.Store.Subscription.Types: instance GHC.Classes.Eq Kiroku.Store.Subscription.Types.TargetBindingPolicy
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Exception.Type.Exception Kiroku.Store.Subscription.Types.ConsumerGroupSizeMismatch
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Exception.Type.Exception Kiroku.Store.Subscription.Types.InvalidHandlerStallWarnAfter
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Exception.Type.Exception Kiroku.Store.Subscription.Types.SomeSubscriptionStartupFailure
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Exception.Type.Exception Kiroku.Store.Subscription.Types.SubscriptionTargetMismatch
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Exception.Type.Exception Kiroku.Store.Subscription.Types.SubscriptionUndecodable
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.BatchSize
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.ConsumerGroupSize
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.ConsumerGroupSizeMismatch
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.InvalidBatchSize
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.InvalidHandlerStallWarnAfter
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.SomeSubscriptionStartupFailure
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.SubscriptionTargetMismatch
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.SubscriptionUndecodable
+ Kiroku.Store.Subscription.Types: instance GHC.Internal.Show.Show Kiroku.Store.Subscription.Types.TargetBindingPolicy
+ Kiroku.Store.Subscription.Types: member :: ConsumerGroup -> Int32
+ Kiroku.Store.Subscription.Types: mkBatchSize :: Int32 -> Either InvalidBatchSize BatchSize
+ Kiroku.Store.Subscription.Types: mkConsumerGroup :: Int32 -> ConsumerGroupSize -> Either InvalidConsumerGroup ConsumerGroup
+ Kiroku.Store.Subscription.Types: mkConsumerGroupSize :: Int32 -> Either InvalidConsumerGroup ConsumerGroupSize
+ Kiroku.Store.Subscription.Types: newtype InvalidBatchSize
+ Kiroku.Store.Subscription.Types: newtype InvalidHandlerStallWarnAfter
+ Kiroku.Store.Subscription.Types: newtype SubscriptionUndecodable
+ Kiroku.Store.Subscription.Types: size :: ConsumerGroup -> Int32
+ Kiroku.Store.Subscription.Worker: instance GHC.Classes.Eq Kiroku.Store.Subscription.Worker.StalledDelivery
- Kiroku.Store.SQL: DeadLetterParams :: !Text -> !Int32 -> !Int64 -> !UUID -> !Value -> !Text -> !Int32 -> DeadLetterParams
+ Kiroku.Store.SQL: DeadLetterParams :: !Text -> !Int32 -> !Text -> !Maybe Text -> !Int32 -> !Int64 -> !UUID -> !Value -> !Text -> !Int32 -> DeadLetterParams
- Kiroku.Store.SQL: saveCheckpointMemberStmt :: Statement (Text, Int32, Int64) ()
+ Kiroku.Store.SQL: saveCheckpointMemberStmt :: Statement (Text, Int32, Int64, Int32, Text, Maybe Text) ()
- Kiroku.Store.Settings: StoreSettings :: !Maybe (EventData -> IO EventData) -> !Maybe (RecordedEvent -> IO RecordedEvent) -> StoreSettings
+ Kiroku.Store.Settings: StoreSettings :: !Maybe (EventData -> IO EventData) -> !Maybe (RecordedEvent -> IO (Either DecodeFailure RecordedEvent)) -> StoreSettings
- Kiroku.Store.Settings: [decodeHook] :: StoreSettings -> !Maybe (RecordedEvent -> IO RecordedEvent)
+ Kiroku.Store.Settings: [decodeHook] :: StoreSettings -> !Maybe (RecordedEvent -> IO (Either DecodeFailure RecordedEvent))
- Kiroku.Store.Settings: decodeEvents :: StoreSettings -> Vector RecordedEvent -> IO (Vector RecordedEvent)
+ Kiroku.Store.Settings: decodeEvents :: StoreSettings -> Vector RecordedEvent -> IO DecodedBatch
- Kiroku.Store.Subscription.EventPublisher: Subscriber :: !TBQueue (Vector RecordedEvent) -> !TVar SubscriberStatus -> !OverflowPolicy -> Subscriber
+ Kiroku.Store.Subscription.EventPublisher: Subscriber :: !TBQueue DecodedBatch -> !TVar SubscriberStatus -> !OverflowPolicy -> Subscriber
- Kiroku.Store.Subscription.EventPublisher: [subQueue] :: Subscriber -> !TBQueue (Vector RecordedEvent)
+ Kiroku.Store.Subscription.EventPublisher: [subQueue] :: Subscriber -> !TBQueue DecodedBatch
- Kiroku.Store.Subscription.EventPublisher: subscribePublisher :: EventPublisher -> Natural -> OverflowPolicy -> STM (TBQueue (Vector RecordedEvent), TVar SubscriberStatus, IO ())
+ Kiroku.Store.Subscription.EventPublisher: subscribePublisher :: EventPublisher -> Natural -> OverflowPolicy -> STM (TBQueue DecodedBatch, TVar SubscriberStatus, IO ())
- Kiroku.Store.Subscription.Fsm: BatchFetched :: !Vector RecordedEvent -> Input
+ Kiroku.Store.Subscription.Fsm: BatchFetched :: !DecodedBatch -> Input
- Kiroku.Store.Subscription.Fsm: ConnectionLost :: !UsageError -> Input
+ Kiroku.Store.Subscription.Fsm: ConnectionLost :: !GlobalPosition -> !UsageError -> Input
- Kiroku.Store.Subscription.Fsm: DeliverBatch :: !Vector RecordedEvent -> Effect
+ Kiroku.Store.Subscription.Fsm: DeliverBatch :: !DecodedBatch -> Effect
- Kiroku.Store.Subscription.Stream: subscriptionAckStream :: KirokuStore -> SubscriptionConfig -> Natural -> IO (Stream IO AckItem, IO ())
+ Kiroku.Store.Subscription.Stream: subscriptionAckStream :: KirokuStore -> SubscriptionConfig -> StreamBufferSize -> IO (Stream IO AckItem, IO ())
- Kiroku.Store.Subscription.Stream: subscriptionStream :: KirokuStore -> SubscriptionConfig -> Natural -> IO (Stream IO RecordedEvent, IO ())
+ Kiroku.Store.Subscription.Stream: subscriptionStream :: KirokuStore -> SubscriptionConfig -> StreamBufferSize -> IO (Stream IO RecordedEvent, IO ())
- Kiroku.Store.Subscription.Types: SubscriptionConfig :: !SubscriptionName -> !SubscriptionTarget -> !EventHandlerM m -> !Int32 -> !Natural -> !OverflowPolicy -> !Maybe ConsumerGroup -> !Bool -> !MissingCheckpointPolicy -> !RetryPolicy -> !EventTypeFilter -> !Maybe (RecordedEvent -> Bool) -> SubscriptionConfigM (m :: Type -> Type)
+ Kiroku.Store.Subscription.Types: SubscriptionConfig :: !SubscriptionName -> !SubscriptionTarget -> !EventHandlerM m -> !Maybe (RecordedEvent -> DecodeFailure -> m SubscriptionResult) -> !Maybe NominalDiffTime -> !BatchSize -> !Natural -> !OverflowPolicy -> !Maybe ConsumerGroup -> !Bool -> !MissingCheckpointPolicy -> !TargetBindingPolicy -> !RetryPolicy -> !EventTypeFilter -> !Maybe (RecordedEvent -> Bool) -> SubscriptionConfigM (m :: Type -> Type)
- Kiroku.Store.Subscription.Types: [batchSize] :: SubscriptionConfigM (m :: Type -> Type) -> !Int32
+ Kiroku.Store.Subscription.Types: [batchSize] :: SubscriptionConfigM (m :: Type -> Type) -> !BatchSize
- Kiroku.Store.Subscription.Worker: LiveFromPublisherQueue :: !TBQueue (Vector RecordedEvent) -> !TVar SubscriberStatus -> LiveSource
+ Kiroku.Store.Subscription.Worker: LiveFromPublisherQueue :: !TBQueue DecodedBatch -> !TVar SubscriberStatus -> LiveSource
Files
- CHANGELOG.md +43/−0
- LICENSE +11/−0
- bench/CheckpointTargetCost.hs +101/−0
- bench/Main.hs +4/−1
- bench/ShibuyaOverhead.hs +15/−5
- kiroku-store.cabal +24/−1
- src/Kiroku/Store/Effect.hs +14/−7
- src/Kiroku/Store/Error.hs +56/−30
- src/Kiroku/Store/Observability.hs +12/−2
- src/Kiroku/Store/SQL.hs +43/−13
- src/Kiroku/Store/Settings.hs +80/−12
- src/Kiroku/Store/Subscription.hs +6/−9
- src/Kiroku/Store/Subscription/Checkpoint.hs +74/−1
- src/Kiroku/Store/Subscription/Checkpoint/SQL.hs +171/−9
- src/Kiroku/Store/Subscription/Effect.hs +6/−2
- src/Kiroku/Store/Subscription/EventPublisher.hs +10/−6
- src/Kiroku/Store/Subscription/Fsm.hs +17/−12
- src/Kiroku/Store/Subscription/Stream.hs +24/−10
- src/Kiroku/Store/Subscription/Types.hs +163/−22
- src/Kiroku/Store/Subscription/Worker.hs +190/−83
- test/Main.hs +108/−21
- test/Test/CatchupDbErrorNoPrematureSwitch.hs +4/−1
- test/Test/CategoryIdleNoSpin.hs +3/−3
- test/Test/ConsumerGroup.hs +8/−6
- test/Test/ConsumerGroupEffect.hs +3/−3
- test/Test/ConsumerGroupResize.hs +257/−0
- test/Test/ConsumerGroupSql.hs +4/−4
- test/Test/EventTypeFilter.hs +2/−2
- test/Test/FailureInjection.hs +4/−1
- test/Test/HandlerStall.hs +254/−0
- test/Test/Helpers.hs +10/−0
- test/Test/InterpreterHooks.hs +12/−3
- test/Test/PerformanceStructure.hs +95/−2
- test/Test/PublisherCallbackResilience.hs +418/−6
- test/Test/PublisherIdleAdvance.hs +1/−1
- test/Test/PublisherRestartNoRebroadcast.hs +12/−3
- test/Test/StreamBridgeTermination.hs +7/−13
- test/Test/SubscriptionCheckpointInventory.hs +1/−1
- test/Test/SubscriptionCheckpointReset.hs +1/−1
- test/Test/SubscriptionCheckpointWorker.hs +5/−4
- test/Test/SubscriptionPauseResume.hs +8/−2
- test/Test/SubscriptionReconnect.hs +57/−1
- test/Test/SubscriptionRegistry.hs +3/−3
- test/Test/SubscriptionRetryDeadLetter.hs +3/−3
- test/Test/SubscriptionState.hs +8/−2
- test/Test/SubscriptionTarget.hs +195/−0
- test/Test/UniqueViolationMapping.hs +93/−0
CHANGELOG.md view
@@ -1,5 +1,48 @@ # Changelog +## 0.10.0.0 — 2026-10-10++### Breaking Changes++* `ConsumerGroup` and `ConsumerGroupSize` are opaque validated types. Construct+ them with `mkConsumerGroupSize` and `mkConsumerGroup`; use `mkBatchSize` for+ positive subscription batches and `mkStreamBufferSize` for bridge capacities.+* `decodeHook` returns `Either DecodeFailure RecordedEvent` in IO;+ `decodeEvents` returns `DecodedBatch`. Return `Right` from successful hooks.+* `SubscriptionConfigM` adds `undecodableHandler` (default `Nothing`) and+ `handlerStallWarnAfter` (default `Nothing`). Observers must handle typed+ publisher decode failures, group-size mismatch, advisory handler stalls and+ `StopUndecodable`.+* Apply migrations 0013 and 0014 from kiroku-store-migrations 0.7.0.0 with+ subscription workers stopped. Startup verifies persisted group size and+ target binding. Declare adoption of legacy targets explicitly; incompatible+ restart refuses before delivery through `SomeSubscriptionStartupFailure`.++### New Features++* Public `resizeConsumerGroupTx` equalizes every new member at the old minimum+ checkpoint and returns `ConsumerGroupResizeReport`, composing with caller SQL.+ Stop every member before resizing, including same-size hash-assignment changes.+* Public `rebindSubscriptionTargetTx` atomically changes a stopped subscription's+ target. Resize preserves target bindings.+* Typed undecodable events use bounded retries and preserve the checkpoint before+ a failed event. Explicit callbacks can skip, stop, retry or dead-letter with+ `DeadLetterDecodeFailure`; reads return `EventDecodeFailed` without partial data.+* A positive `handlerStallWarnAfter` enables a scoped advisory watchdog. Invalid+ intervals refuse on `wait` before checkpoint initialization. Warnings do not+ acknowledge, retry or checkpoint an event; the default path has no tracking work.++### Bug Fixes++* Live reconnect retains processed progress. Typed hook failures no longer stall+ the shared publisher; absent hooks retain the unchanged-vector fast path.+ Hook exceptions remain programming failures.+* Match unique constraints by exact name across append, transaction, link and+ multi-stream error attribution. `events_pkey` and `stream_events_pkey` return+ `DuplicateEvent` with a parseable caller ID. `ux_stream_events_stream_version`+ returns `UnexpectedServerError "23505"` with the original message for an+ invariant failure. Unknown append constraints retain the expected-version fallback.+ ## 0.9.0.1 — 2026-09-25 ### Bug Fixes
+ LICENSE view
@@ -0,0 +1,11 @@+Copyright (c) 2026 Nadeem Bitar.++Redistribution and use in source and binary forms, with or without modification, are permitted provided that the following conditions are met:++1. Redistributions of source code must retain the above copyright notice, this list of conditions and the following disclaimer.++2. Redistributions in binary form must reproduce the above copyright notice, this list of conditions and the following disclaimer in the documentation and/or other materials provided with the distribution.++3. Neither the name of the copyright holder nor the names of its contributors may be used to endorse or promote products derived from this software without specific prior written permission.++THIS SOFTWARE IS PROVIDED BY THE COPYRIGHT HOLDERS AND CONTRIBUTORS "AS IS" AND ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE ARE DISCLAIMED. IN NO EVENT SHALL THE COPYRIGHT HOLDER OR CONTRIBUTORS BE LIABLE FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION) HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF SUCH DAMAGE.
+ bench/CheckpointTargetCost.hs view
@@ -0,0 +1,101 @@+{-# LANGUAGE MultilineStrings #-}+{-# LANGUAGE OverloadedRecordDot #-}++{- | Small supplementary EP-2 comparison of the changed checkpoint-save path.+No acceptance threshold: local ratios retain their uncertainty. The control+statement and row layout are from 23a03a1b8b56773a140d0683c4c8d4d34b9d7e36.+-}+module Main where++import Contravariant.Extras (contrazip4)+import Control.Lens ((^.))+import Control.Monad (forM_, replicateM_)+import Data.Generics.Labels ()+import Data.IORef (atomicModifyIORef', newIORef)+import Data.Int (Int32, Int64)+import Data.Text (Text)+import EphemeralPg qualified as Pg+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.Store.SQL qualified as SQL+import Kiroku.Test.Postgres (ephemeralConfig, migrateTestDatabase)++main :: IO ()+main = do+ original <- ephemeralConfig+ let keys = ["fsync", "synchronous_commit", "full_page_writes", "shared_buffers", "wal_level"]+ durable = original{Pg.postgresSettings = filter (\(key, _) -> key `notElem` keys) original.postgresSettings <> [("fsync", "on"), ("synchronous_commit", "on"), ("full_page_writes", "on"), ("shared_buffers", "128MB"), ("wal_level", "replica")]}+ result <- Pg.withCachedConfig durable Pg.defaultCacheConfig $ \database -> do+ let connection = Pg.connectionString database+ migrateTestDatabase connection+ withStore (defaultConnectionSettings connection) measure+ either (fail . show) pure result++measure :: KirokuStore -> IO ()+measure store = do+ settings <- use store $ Session.statement () (preparable "SELECT current_setting('server_version_num')::int4 >= 180000 AND current_setting('server_version_num')::int4 < 190000 AND current_setting('fsync') = 'on' AND current_setting('synchronous_commit') = 'on' AND current_setting('full_page_writes') = 'on'" E.noParams (D.singleRow (D.column (D.nonNullable D.bool))))+ if settings then pure () else fail "requires durable PostgreSQL 18"+ use store $+ Session.script+ """+ CREATE SCHEMA ep2_control;+ CREATE TABLE ep2_control.subscriptions (+ subscription_id BIGSERIAL PRIMARY KEY,+ subscription_name TEXT NOT NULL,+ stream_name TEXT NOT NULL DEFAULT '$all',+ last_seen BIGINT NOT NULL DEFAULT 0,+ consumer_group_member INT NOT NULL DEFAULT 0,+ consumer_group_size INT NOT NULL DEFAULT 1,+ created_at TIMESTAMPTZ NOT NULL DEFAULT now(),+ updated_at TIMESTAMPTZ NOT NULL DEFAULT now()+ );+ CREATE UNIQUE INDEX ON ep2_control.subscriptions (subscription_name, consumer_group_member);+ """+ controlCounter <- newIORef 0+ candidateCounter <- newIORef 0+ let save candidate = do+ p <- atomicModifyIORef' (if candidate then candidateCounter else controlCounter) (\n -> (n + 1, n + 1))+ let m = fromIntegral (p `mod` 4)+ if candidate+ then use store $ Session.statement ("target-cost", m, p, 4, "performance") SQL.saveCategoryCheckpointMemberStmt+ else use store $ Session.statement ("target-cost", m, p, 4) controlSave+ trial pair candidate = do+ walBefore <- use store (Session.statement () walPosition)+ start <- getMonotonicTimeNSec+ replicateM_ 2000 (save candidate)+ end <- getMonotonicTimeNSec+ walAfter <- use store (Session.statement () walPosition)+ putStrLn $ show pair <> "," <> (if candidate then "candidate" else "control") <> ",2000," <> show (fromIntegral (end - start) / 2000 :: Double) <> "," <> show (fromIntegral (walAfter - walBefore) / 2000 :: Double)+ replicateM_ 200 (save False >> save True)+ putStrLn "pair,arm,saves,ns_per_save,wal_bytes_per_save"+ forM_ [1 .. 3 :: Int] $ \pair ->+ if odd pair+ then trial pair False >> trial pair True+ else trial pair True >> trial pair False+ -- Every completed save advances one durable member; both arms do equal work.+ expected <- use store $ Session.statement () (preparable "SELECT count(*) = 4 AND sum(last_seen) = 24794 FROM ep2_control.subscriptions" E.noParams (D.singleRow (D.column (D.nonNullable D.bool))))+ actual <- use store $ Session.statement () (preparable "SELECT count(*) = 4 AND sum(last_seen) = 24794 AND bool_and(target_kind = 'category' AND target_category = 'performance') FROM kiroku.subscriptions WHERE subscription_name = 'target-cost'" E.noParams (D.singleRow (D.column (D.nonNullable D.bool))))+ if expected && actual then putStrLn "durable equal-work checks: passed" else fail "durable equal-work checks failed"++use :: KirokuStore -> Session.Session a -> IO a+use store session = Pool.use (store ^. #pool) session >>= either (fail . show) pure++walPosition :: Statement () Int64+walPosition = preparable "SELECT (pg_current_wal_insert_lsn() - '0/0'::pg_lsn)::bigint" E.noParams (D.singleRow (D.column (D.nonNullable D.int8)))++controlSave :: Statement (Text, Int32, Int64, Int32) ()+controlSave =+ preparable+ """+ INSERT INTO ep2_control.subscriptions (subscription_name, consumer_group_member, last_seen, updated_at, consumer_group_size)+ VALUES ($1, $2, $3, now(), $4)+ ON CONFLICT (subscription_name, consumer_group_member)+ DO UPDATE SET last_seen = GREATEST(subscriptions.last_seen, EXCLUDED.last_seen), updated_at = now(), consumer_group_size = EXCLUDED.consumer_group_size+ """+ (contrazip4 (E.param (E.nonNullable E.text)) (E.param (E.nonNullable E.int4)) (E.param (E.nonNullable E.int8)) (E.param (E.nonNullable E.int4)))+ D.noResult
bench/Main.hs view
@@ -619,12 +619,15 @@ { name = subName , target = Category (CategoryName cat) , handler = handler- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing
bench/ShibuyaOverhead.hs view
@@ -29,6 +29,7 @@ import EphemeralPg qualified as Pg import Kiroku.Store import Kiroku.Store.Subscription.Stream (subscriptionStream)+import Kiroku.Store.Subscription.Stream qualified as Buffer import Kiroku.Test.Postgres (ephemeralConfig) import Shibuya.Adapter (Adapter (..)) import Shibuya.App (ProcessorId (..), defaultAppConfig, mkProcessor, runApp, stopApp)@@ -120,12 +121,15 @@ { name = subName , target = AllStreams , handler = handler- , batchSize = 500+ , batchSize = either (error . show) Prelude.id (mkBatchSize 500) , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -147,19 +151,22 @@ { name = subName , target = AllStreams , handler = \_ -> pure Continue- , batchSize = 500+ , batchSize = either (error . show) Prelude.id (mkBatchSize 500) , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing } t0 <- getCurrentTime- (stream, cancelStream) <- subscriptionStream store cfg 256+ (stream, cancelStream) <- subscriptionStream store cfg (either (error . show) Prelude.id (Buffer.mkStreamBufferSize 256)) Stream.fold Fold.drain (Stream.take n stream) cancelStream t1 <- getCurrentTime@@ -179,18 +186,21 @@ { name = subName , target = AllStreams , handler = \_ -> pure Continue- , batchSize = 500+ , batchSize = either (error . show) Prelude.id (mkBatchSize 500) , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing } - (ioStream, cancelAction) <- liftIO $ subscriptionStream store cfg 256+ (ioStream, cancelAction) <- liftIO $ subscriptionStream store cfg (either (error . show) Prelude.id (Buffer.mkStreamBufferSize 256)) let effStream = Stream.morphInner liftIO ioStream ingestedStream = fmap (mkIngested cancelAction) effStream
kiroku-store.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.0 name: kiroku-store-version: 0.9.0.1+version: 0.10.0.0 synopsis: High-performance PostgreSQL event store description: Kiroku is a PostgreSQL-backed event store for Haskell applications. It@@ -12,6 +12,7 @@ author: Nadeem Bitar maintainer: nadeem@gmail.com license: BSD-3-Clause+license-file: LICENSE build-type: Simple category: Database, Eventing extra-doc-files: CHANGELOG.md@@ -105,9 +106,11 @@ Test.Concurrency Test.ConsumerGroup Test.ConsumerGroupEffect+ Test.ConsumerGroupResize Test.ConsumerGroupSql Test.EventTypeFilter Test.FailureInjection+ Test.HandlerStall Test.Helpers Test.HistoryRetention Test.HistoryRetentionMock@@ -134,8 +137,10 @@ Test.SubscriptionRegistry Test.SubscriptionRetryDeadLetter Test.SubscriptionState+ Test.SubscriptionTarget Test.Transaction Test.TruncateBefore+ Test.UniqueViolationMapping Test.VisibleGlobalHeadPosition Test.VisibleGlobalHeadPositionMock @@ -271,3 +276,21 @@ , time >=1.12 && <1.15 , uuid >=1.3 && <1.4 , vector >=0.13 && <0.14++benchmark kiroku-checkpoint-target-cost+ import: common+ type: exitcode-stdio-1.0+ main-is: CheckpointTargetCost.hs+ hs-source-dirs: bench+ ghc-options: -threaded -rtsopts "-with-rtsopts=-N4 -A32m"+ build-depends:+ , base+ , contravariant-extras+ , ephemeral-pg+ , generic-lens+ , hasql+ , hasql-pool+ , kiroku-store+ , kiroku-test-support+ , lens+ , text
src/Kiroku/Store/Effect.hs view
@@ -56,7 +56,7 @@ import Kiroku.Store.HistoryRetention.Types import Kiroku.Store.Observability (KirokuEvent (..), emitOrDrop) import Kiroku.Store.SQL qualified as SQL-import Kiroku.Store.Settings (decodeEvents, enrichEvents)+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.Types (@@ -237,24 +237,24 @@ evs <- usePool (store ^. #pool) $ Session.statement (name, startVer, limit) SQL.readStreamForwardStmt- liftIO $ decodeEvents (store ^. #storeSettings) evs+ decodeReadEvents (store ^. #storeSettings) evs ReadStreamBackward (StreamName name) (StreamVersion startVer) limit -> do let cursor = if startVer == 0 then maxBound else startVer evs <- usePool (store ^. #pool) $ Session.statement (name, cursor, limit) SQL.readStreamBackwardStmt- liftIO $ decodeEvents (store ^. #storeSettings) evs+ decodeReadEvents (store ^. #storeSettings) evs ReadAllForward (GlobalPosition startPos) limit -> do evs <- usePool (store ^. #pool) $ Session.statement (startPos, limit) SQL.readAllForwardStmt- liftIO $ decodeEvents (store ^. #storeSettings) evs+ decodeReadEvents (store ^. #storeSettings) evs ReadAllBackward (GlobalPosition startPos) limit -> do let cursor = if startPos == 0 then maxBound else startPos evs <- usePool (store ^. #pool) $ Session.statement (cursor, limit) SQL.readAllBackwardStmt- liftIO $ decodeEvents (store ^. #storeSettings) evs+ decodeReadEvents (store ^. #storeSettings) evs GetVisibleGlobalHeadPosition -> usePool (store ^. #pool) $ Session.statement () SQL.visibleGlobalHeadPositionStmt@@ -291,7 +291,7 @@ evs <- usePool (store ^. #pool) $ Session.statement (startPos, cat, limit) SQL.readCategoryForwardStmt- liftIO $ decodeEvents (store ^. #storeSettings) evs+ decodeReadEvents (store ^. #storeSettings) evs AppendMultiStream [] -> pure [] AppendMultiStream ops -> do@@ -350,7 +350,7 @@ FilterCausationAncestors (EventId eid) -> usePool (store ^. #pool) $ Session.statement eid SQL.findCausationAncestorsStmt- liftIO $ decodeEvents (store ^. #storeSettings) evs+ decodeReadEvents (store ^. #storeSettings) evs SoftDeleteStream (StreamName name) -> do rejectInvalidApplicationStream name usePool (store ^. #pool) $@@ -667,3 +667,10 @@ packing, or per-version dispatch. They are not part of the supported public surface and may change without notice. -}++-- No hook preserves the direct read passthrough. Successful hooks map directly+-- into a result vector, avoiding the subscription-specific wrapper/unwrapper.+decodeReadEvents :: (IOE :> es, Error StoreError :> es) => StoreSettings -> Vector RecordedEvent -> Eff es (Vector RecordedEvent)+decodeReadEvents settings events = case decodeHook settings of+ Nothing -> pure events+ Just hook -> V.mapM (\event -> liftIO (hook event) >>= either (throwError . EventDecodeFailed) pure) events
src/Kiroku/Store/Error.hs view
@@ -19,8 +19,11 @@ extractStreamNameFromDetail, ) where +import Control.Applicative ((<|>)) import Control.Exception (Exception) import Data.ByteString qualified as BS+import Data.Char (isAlphaNum)+import Data.List (find) import Data.Text (Text) import Data.Text qualified as T import Data.Text.Encoding qualified as TE@@ -30,6 +33,7 @@ import Hasql.Errors qualified as Errors import Hasql.Pool (UsageError (..)) import Kiroku.Store.HistoryRetention.Types (HistoryRetentionConflict)+import Kiroku.Store.Settings (DecodeFailure) import Kiroku.Store.Types {- | Errors that can occur during store operations.@@ -97,6 +101,9 @@ StreamAlreadyExists !StreamName | {- | A caller-supplied @event_id@ collides with an existing event. + Append maps both @events_pkey@ and the composite+ @stream_events_pkey@ to this constructor.+ The constructor carries 'Just' the id when the PostgreSQL detail string could be parsed, 'Nothing' otherwise. A 'Nothing' payload is rare in practice; it occurs when the server's locale changes@@ -147,9 +154,9 @@ -} TransientTransactionFailure !Text !Text | {- | PostgreSQL raised a server error whose @SQLSTATE@ code is- outside the set this store recognises (currently @23505@- unique violation, @23503@ foreign key violation, and the- class-40 codes carried by 'TransientTransactionFailure'). The+ outside the set this store recognises, or an append violated+ @ux_stream_events_stream_version@ (@23505@), indicating an+ internal stream-version invariant failure. The first 'Text' is the @SQLSTATE@ code, the second is the human-readable message. This is *not* generally retryable — investigate.@@ -161,6 +168,8 @@ prefer the specific constructors. -} ConnectionError !Text+ | -- | A decode hook refused an event; the read returns no partial result.+ EventDecodeFailed !DecodeFailure deriving stock (Eq, Show, Generic) deriving anyclass (Exception) @@ -184,13 +193,17 @@ PostgreSQL error code mapping: 23505 (unique_violation) + events_pkey -> 'DuplicateEvent'+ 23505 (unique_violation) + stream_events_pkey -> 'DuplicateEvent'+ 23505 + ux_stream_events_stream_version -> 'UnexpectedServerError' 23505 (unique_violation) + ix_streams_stream_name -> 'StreamAlreadyExists' 23505 (unique_violation) + other -> 'WrongExpectedVersion' 23503 (foreign_key_violation) -> 'StreamNotFound'+ 40001 / 40P01 -> 'TransientTransactionFailure' any other server code -> 'UnexpectedServerError' -The constraint-name matching depends on the literal strings @events_pkey@-and @ix_streams_stream_name@. If a future schema migration renames a+Constraint names are compared exactly, with the quoted message taking+precedence over a delimiter-aware detail fallback. Matching depends on the+four literal names above. If a future schema migration renames a constraint, the @23505@ branch falls through to the generic 'WrongExpectedVersion' mapping; keep the names stable in @kiroku-store-migrations/migrations@.@@ -214,15 +227,12 @@ mapTransactionUsageError usageErr = case extractServerError usageErr of Just (Errors.ServerError "23505" message detail _ _)- | containsConstraint "events_pkey" message detail ->+ | uniqueConstraintName message detail == Just "events_pkey" -> DuplicateEvent (EventId <$> (detail >>= extractUuidFromDetail))- | containsConstraint "stream_events_pkey" message detail ->+ | uniqueConstraintName message detail == Just "stream_events_pkey" -> DuplicateEvent (EventId <$> (detail >>= extractFirstUuidFromCompositeDetail)) _ -> mapGenericUsageError usageErr- where- containsConstraint name message detail =- name `T.isInfixOf` message || maybe False (T.isInfixOf name) detail -- | Generic, non-append-shaped mapping for hasql pool usage errors. mapGenericUsageError :: UsageError -> StoreError@@ -246,16 +256,13 @@ mapLinkUsageError target usageErr = case extractServerError usageErr of Just (Errors.ServerError "23505" message detail _ _)- | containsConstraint "stream_events_pkey" message detail ->+ | uniqueConstraintName message detail == Just "stream_events_pkey" -> EventAlreadyLinked target (extractCompositeEventId detail) Just (Errors.ServerError "23502" _ _ _ _) -> LinkSourceEventMissing target _ -> mapGenericUsageError usageErr where- containsConstraint name message detail =- name `T.isInfixOf` message || maybe False (T.isInfixOf name) detail- extractCompositeEventId (Just d) = EventId <$> extractFirstUuidFromCompositeDetail d extractCompositeEventId Nothing = Nothing @@ -286,27 +293,47 @@ - message: "duplicate key value violates unique constraint \"events_pkey\"" - detail: "Key (event_id)=(uuid-value) already exists." -We check both message and detail for the constraint name.+We extract one constraint name, preferring the quoted message over detail,+then compare it exactly. Composite link keys use their first UUID. A+stream-version index violation is an invariant failure, not a precondition+conflict, and preserves the SQLSTATE and original message. When the events_pkey case fires but the detail string cannot be parsed (e.g., the server's locale produced an unexpected format), the 'DuplicateEvent' constructor carries 'Nothing' rather than a fabricated-all-zeroes UUID — see 'extractEventId'.+all-zeroes UUID. -} mapUniqueViolation :: Text -> ExpectedVersion -> Text -> Maybe Text -> StoreError-mapUniqueViolation streamName expected message detail- | containsConstraint "events_pkey" = DuplicateEvent (extractEventId detail)- | containsConstraint "ix_streams_stream_name" = StreamAlreadyExists (StreamName streamName)- | otherwise =- -- Generic unique violation — treat as version conflict- WrongExpectedVersion (StreamName streamName) expected (StreamVersion 0)- where- containsConstraint name =- name `T.isInfixOf` message || maybe False (T.isInfixOf name) detail+mapUniqueViolation streamName expected message detail =+ case uniqueConstraintName message detail of+ Just "events_pkey" -> DuplicateEvent (EventId <$> (detail >>= extractUuidFromDetail))+ Just "stream_events_pkey" -> DuplicateEvent (EventId <$> (detail >>= extractFirstUuidFromCompositeDetail))+ Just "ix_streams_stream_name" -> StreamAlreadyExists (StreamName streamName)+ Just "ux_stream_events_stream_version" -> UnexpectedServerError "23505" message+ _ ->+ -- Generic unique violation — treat as version conflict+ WrongExpectedVersion (StreamName streamName) expected (StreamVersion 0) - -- Try to extract event_id from detail like "Key (event_id)=(uuid) already exists."- extractEventId (Just d) = EventId <$> extractUuidFromDetail d- extractEventId Nothing = Nothing+{- | PostgreSQL's quoted constraint wins even when it is unknown. Legacy+detail-only errors may name an owned constraint as a complete identifier;+underscores and dollar signs are identifier characters, not delimiters.+This helper is used only after a failed statement has returned SQLSTATE 23505.+-}+uniqueConstraintName :: Text -> Maybe Text -> Maybe Text+uniqueConstraintName message detail =+ quotedConstraint message <|> (detail >>= fromDetail)+ where+ quotedConstraint input =+ case T.breakOn "unique constraint \"" input of+ (_, rest)+ | not (T.null rest) ->+ let (name, closing) = T.breakOn "\"" (T.drop (T.length "unique constraint \"") rest)+ in if T.null name || T.null closing then Nothing else Just name+ _ -> Nothing+ fromDetail input =+ quotedConstraint input <|> find (`elem` ownedConstraints) (T.split (not . identifierChar) input)+ identifierChar c = isAlphaNum c || c == '_' || c == '$'+ ownedConstraints = ["events_pkey", "stream_events_pkey", "ix_streams_stream_name", "ux_stream_events_stream_version"] {- | Append-precondition failures observable inside a 'Hasql.Transaction.Transaction' body.@@ -476,8 +503,7 @@ attributeMultiStreamError ops@((StreamName firstName, firstExpected) : _) usageErr = case extractServerError usageErr of Just (Errors.ServerError "23505" message (Just detail) _ _)- | "ix_streams_stream_name" `T.isInfixOf` message- || "ix_streams_stream_name" `T.isInfixOf` detail+ | uniqueConstraintName message (Just detail) == Just "ix_streams_stream_name" , Just sn <- extractStreamNameFromDetail detail , Just (StreamName name, expected) <- lookupStream sn ops -> mapUsageError name expected usageErr
src/Kiroku/Store/Observability.hs view
@@ -54,16 +54,20 @@ import Control.Exception (SomeAsyncException, SomeException, asyncExceptionFromException, catch, throwIO) import Data.Foldable (for_) import Data.Int (Int32)+import Data.Time (NominalDiffTime) import Data.Time.Clock (UTCTime) import Hasql.Pool (UsageError) import Kiroku.Store.HistoryRetention.Types (HistoryRetentionConflict, HistoryRetentionLeaseId, HistoryRetentionLeaseOwner, HistoryRetentionPruneResult)+import Kiroku.Store.Settings (DecodeFailure) import Kiroku.Store.Subscription.Fsm (DeadLetterReason (..), SubscriptionStopReason (..)) import Kiroku.Store.Subscription.Types ( CheckpointInitialization,+ ConsumerGroupSizeMismatch, SubscriptionCheckpointMissing, SubscriptionName,+ SubscriptionTarget, )-import Kiroku.Store.Types (GlobalPosition, StreamId, StreamName)+import Kiroku.Store.Types (EventId, GlobalPosition, StreamId, StreamName) {- | A structured operational event emitted by 'Kiroku.Store' itself. @@ -72,7 +76,9 @@ module Haddock for context. -} data KirokuEvent- = {- | The dedicated @LISTEN@ connection encountered a non-async+ = KirokuEventSubscriptionTargetBound !SubscriptionName !SubscriptionTarget !SubscriptionGroupContext+ | KirokuEventSubscriptionGroupSizeMismatch !ConsumerGroupSizeMismatch !SubscriptionGroupContext+ | {- | The dedicated @LISTEN@ connection encountered a non-async exception and the listener loop is about to attempt reconnection. The 'Int' is the consecutive failure count starting at @1@; it drives the exponential-backoff delay (capped at 30 seconds) and@@ -100,6 +106,10 @@ failing callback is stalling live broadcast until it is fixed. -} KirokuEventPublisherLoopError !SomeException+ | -- | Typed hook failure at broadcast time; publisher still advances.+ KirokuEventPublisherDecodeFailed !GlobalPosition !EventId !DecodeFailure+ | -- | One handler invocation remains pending after the configured interval.+ KirokuEventSubscriptionHandlerStalled !SubscriptionName !GlobalPosition !EventId !NominalDiffTime !SubscriptionGroupContext | {- | A subscription's worker thread encountered a 'UsageError' in the database phase identified by 'SubscriptionDbPhase'. Checkpoint-load errors fail startup loudly, fetch-batch errors are retried at the same
src/Kiroku/Store/SQL.hs view
@@ -61,6 +61,8 @@ saveCheckpointStmt, getCheckpointMemberStmt, saveCheckpointMemberStmt,+ saveAllCheckpointMemberStmt,+ saveCategoryCheckpointMemberStmt, -- * Dead-letter statements DeadLetterParams (..),@@ -69,7 +71,7 @@ readDeadLettersStmt, ) where -import Contravariant.Extras (contrazip2, contrazip3, contrazip4, contrazip5)+import Contravariant.Extras (contrazip2, contrazip3, contrazip4, contrazip5, contrazip6) import Control.Lens ((^.)) import Data.Aeson (Value) import Data.Functor.Contravariant ((>$<))@@ -1255,17 +1257,35 @@ same @GREATEST(...)@ monotonicity as 'saveCheckpointStmt' so a save never moves a member's checkpoint backward. -}-saveCheckpointMemberStmt :: Statement (Text, Int32, Int64) ()+saveCheckpointMemberStmt :: Statement (Text, Int32, Int64, Int32, Text, Maybe Text) () saveCheckpointMemberStmt = preparable saveCheckpointMemberSQL- ( contrazip3+ ( contrazip6 (E.param (E.nonNullable E.text)) (E.param (E.nonNullable E.int4)) (E.param (E.nonNullable E.int8))+ (E.param (E.nonNullable E.int4))+ (E.param (E.nonNullable E.text))+ (E.param (E.nullable E.text)) ) D.noResult +-- | Bound targets use a fixed kind rather than encoding a constant per save.+saveAllCheckpointMemberStmt :: Statement (Text, Int32, Int64, Int32) ()+saveAllCheckpointMemberStmt =+ preparable+ (saveCheckpointMemberSQLWith "'all', NULL")+ (contrazip4 (E.param (E.nonNullable E.text)) (E.param (E.nonNullable E.int4)) (E.param (E.nonNullable E.int8)) (E.param (E.nonNullable E.int4)))+ D.noResult++saveCategoryCheckpointMemberStmt :: Statement (Text, Int32, Int64, Int32, Text) ()+saveCategoryCheckpointMemberStmt =+ preparable+ (saveCheckpointMemberSQLWith "'category', $5")+ (contrazip5 (E.param (E.nonNullable E.text)) (E.param (E.nonNullable E.int4)) (E.param (E.nonNullable E.int8)) (E.param (E.nonNullable E.int4)) (E.param (E.nonNullable E.text)))+ D.noResult+ getCheckpointMemberSQL :: Text getCheckpointMemberSQL = """@@ -1276,14 +1296,18 @@ """ saveCheckpointMemberSQL :: Text-saveCheckpointMemberSQL =- """- INSERT INTO subscriptions (subscription_name, consumer_group_member, last_seen, updated_at)- VALUES ($1, $2, $3, now())- ON CONFLICT (subscription_name, consumer_group_member)- DO UPDATE SET last_seen = GREATEST(subscriptions.last_seen, EXCLUDED.last_seen), updated_at = now()- """+saveCheckpointMemberSQL = saveCheckpointMemberSQLWith "$5, $6" +saveCheckpointMemberSQLWith :: Text -> Text+saveCheckpointMemberSQLWith binding =+ "INSERT INTO subscriptions (subscription_name, consumer_group_member, last_seen, updated_at, consumer_group_size, target_kind, target_category) VALUES ($1, $2, $3, now(), $4, "+ <> binding+ <> ") "+ <> """+ ON CONFLICT (subscription_name, consumer_group_member)+ DO UPDATE SET last_seen = GREATEST(subscriptions.last_seen, EXCLUDED.last_seen), updated_at = now(), consumer_group_size = EXCLUDED.consumer_group_size, target_kind = EXCLUDED.target_kind, target_category = EXCLUDED.target_category+ """+ -- --------------------------------------------------------------------------- -- Dead-letter Statements -- ---------------------------------------------------------------------------@@ -1292,6 +1316,9 @@ data DeadLetterParams = DeadLetterParams { dlSubscriptionName :: !Text , dlMember :: !Int32+ , dlTargetKind :: !Text+ , dlTargetCategory :: !(Maybe Text)+ , dlGroupSize :: !Int32 , dlGlobalPosition :: !Int64 , dlEventId :: !UUID , dlReason :: !Value@@ -1322,6 +1349,9 @@ <> ((^. #dlReason) >$< E.param (E.nonNullable E.jsonb)) <> ((^. #dlReasonSummary) >$< E.param (E.nonNullable E.text)) <> ((^. #dlAttemptCount) >$< E.param (E.nonNullable E.int4))+ <> ((^. #dlGroupSize) >$< E.param (E.nonNullable E.int4))+ <> ((^. #dlTargetKind) >$< E.param (E.nonNullable E.text))+ <> ((^. #dlTargetCategory) >$< E.param (E.nullable E.text)) {- | Atomically record an event in @kiroku.dead_letters@ and advance the subscription's checkpoint past it, in a single statement.@@ -1350,10 +1380,10 @@ VALUES ($1, $2, $3, $4, $5, $6, $7) ON CONFLICT (subscription_name, consumer_group_member, global_position, event_id) DO NOTHING )- INSERT INTO subscriptions (subscription_name, consumer_group_member, last_seen, updated_at)- VALUES ($1, $2, $3, now())+ INSERT INTO subscriptions (subscription_name, consumer_group_member, last_seen, updated_at, consumer_group_size, target_kind, target_category)+ VALUES ($1, $2, $3, now(), $8, $9, $10) ON CONFLICT (subscription_name, consumer_group_member)- DO UPDATE SET last_seen = GREATEST(subscriptions.last_seen, EXCLUDED.last_seen), updated_at = now()+ DO UPDATE SET last_seen = GREATEST(subscriptions.last_seen, EXCLUDED.last_seen), updated_at = now(), consumer_group_size = EXCLUDED.consumer_group_size, target_kind = EXCLUDED.target_kind, target_category = EXCLUDED.target_category """ -- | Read the dead letters recorded for one subscription member, newest first.
src/Kiroku/Store/Settings.hs view
@@ -10,8 +10,8 @@ hooks to operate on opaque bytes. Both fields default to 'Nothing'. With the defaults, the helpers below-take a 'pure' fast path that allocates nothing extra; no traversal of-the events list or vector occurs.+retain the original event list/vector without traversal. Subscription+decoding adds one batch constructor, with no per-event wrappers. A typical use case is enriching every appended event with an OpenTelemetry trace context drawn from the calling thread:@@ -22,7 +22,7 @@ ctx <- captureCurrentSpan -- OpenTelemetry, OTLP, whatever pure (ed & #metadata %~ injectTraceContext ctx) , 'decodeHook' = Just $ \\re ->- pure (re & #metadata %~ Just . redactPII)+ pure (Right (re & #metadata %~ Just . redactPII)) } @ @@ -41,14 +41,67 @@ StoreSettings (..), defaultStoreSettings, enrichEvents,+ DecodeFailure (..),+ DecodedEvent (..),+ DecodedBatch (..), decodeEvents,+ decodeEvent,+ decodedEventRecorded,+ decodedBatchLength,+ decodedBatchLastEvent,+ filterDecodedBatch, ) where +import Data.Text (Text) import Data.Vector (Vector) import Data.Vector qualified as V import GHC.Generics (Generic)-import Kiroku.Store.Types (EventData, RecordedEvent)+import Kiroku.Store.Types (EventData, EventId, RecordedEvent) +{- | A hook could not decode this event. Return it through 'Left'; throwing+from the hook remains a programming error. Default retry exhaustion is surfaced+by SubscriptionUndecodable, carrying this failure.+-}+data DecodeFailure = DecodeFailure+ { decodeFailureEventId :: !EventId+ , decodeFailureReason :: !Text+ }+ deriving stock (Eq, Show, Generic)++-- | A transformed event, or its raw value retained for retry and disposition.+data DecodedEvent+ = Decoded !RecordedEvent+ | Undecodable !RecordedEvent !DecodeFailure+ deriving stock (Eq, Show)++{- | No hook means no traversal or per-event wrappers. Hook results are shared+by all live publisher subscribers, whose retry/disposition is independent.+-}+data DecodedBatch+ = UnchangedBatch !(Vector RecordedEvent)+ | TransformedBatch !(Vector DecodedEvent)+ deriving stock (Eq, Show)++decodedEventRecorded :: DecodedEvent -> RecordedEvent+decodedEventRecorded = \case+ Decoded event -> event+ Undecodable event _ -> event++decodedBatchLength :: DecodedBatch -> Int+decodedBatchLength = \case+ UnchangedBatch events -> V.length events+ TransformedBatch events -> V.length events++decodedBatchLastEvent :: DecodedBatch -> RecordedEvent+decodedBatchLastEvent = \case+ UnchangedBatch events -> V.last events+ TransformedBatch events -> decodedEventRecorded (V.last events)++filterDecodedBatch :: (RecordedEvent -> Bool) -> DecodedBatch -> DecodedBatch+filterDecodedBatch predicate = \case+ UnchangedBatch events -> UnchangedBatch (V.filter predicate events)+ TransformedBatch events -> TransformedBatch (V.filter (predicate . decodedEventRecorded) events)+ {- | Interpreter-level hooks for cross-cutting concerns at the event-data boundary. All fields default to 'Nothing' (no-op). @@ -61,14 +114,20 @@ the caller. Used to decrypt payloads, redact PII, or attach derived metadata. -When a field is 'Nothing', the interpreter takes a @pure@ fast path-that does not allocate or traverse.+When a field is 'Nothing', reads and enrichment return their input directly.+Subscription decoding retains the vector in one batch constructor without+traversal or per-event wrappers. -} data StoreSettings = StoreSettings { enrichEvent :: !(Maybe (EventData -> IO EventData)) -- ^ Append-path hook. Runs once per appended event before encoding.- , decodeHook :: !(Maybe (RecordedEvent -> IO RecordedEvent))- -- ^ Read- and subscription-path hook. Runs once per surfaced event.+ , decodeHook :: !(Maybe (RecordedEvent -> IO (Either DecodeFailure RecordedEvent)))+ {- ^ Read- and subscription-path hook. Return 'Left' for an undecodable+ event: reads fail with a typed store error, subscriptions use their+ optional undecodable handler or retry and stop by default. Exceptions+ remain programming failures. Runs once per surfaced event; an undecodable+ event's retry re-applies the hook to its original value.+ -} } deriving stock (Generic) @@ -89,9 +148,18 @@ Just f -> traverse f xs {- | Apply 'decodeHook' to a vector of events. When the hook is-'Nothing', returns the vector unchanged with no traversal.+'Nothing', retains the vector unchanged without traversal or per-event wrappers. -}-decodeEvents :: StoreSettings -> Vector RecordedEvent -> IO (Vector RecordedEvent)+decodeEvents :: StoreSettings -> Vector RecordedEvent -> IO DecodedBatch decodeEvents ss xs = case decodeHook ss of- Nothing -> pure xs- Just f -> V.mapM f xs+ Nothing -> pure (UnchangedBatch xs)+ Just f -> TransformedBatch <$> V.mapM (applyDecode f) xs++-- | Re-apply the hook to one raw undecodable event on a subscriber retry.+decodeEvent :: StoreSettings -> RecordedEvent -> IO DecodedEvent+decodeEvent ss event = case decodeHook ss of+ Nothing -> pure (Decoded event)+ Just f -> applyDecode f event++applyDecode :: (RecordedEvent -> IO (Either DecodeFailure RecordedEvent)) -> RecordedEvent -> IO DecodedEvent+applyDecode f event = either (Undecodable event) Decoded <$> f event
src/Kiroku/Store/Subscription.hs view
@@ -15,12 +15,10 @@ import Control.Concurrent.Async qualified as Async import Control.Concurrent.STM (atomically, modifyTVar', newTVarIO, readTVarIO)-import Control.Exception (bracket, bracketOnError, finally, mask, throwIO)+import Control.Exception (bracket, bracketOnError, finally, mask) import Control.Lens ((^.))-import Control.Monad (when) import Control.Monad.IO.Class (MonadIO, liftIO) import Control.Monad.IO.Unlift (MonadUnliftIO, withRunInIO)-import Data.Foldable (for_) import Data.Generics.Labels () import Data.Int (Int32) import Data.Map.Strict (Map)@@ -109,6 +107,11 @@ Investigate the slow handler and either fix the slowness, raise 'queueCapacity', or switch to 'Kiroku.Store.Subscription.Types.DropOldest' if the consumer can tolerate event loss.+* @Left e@ where @e@ is+ 'Kiroku.Store.Subscription.Types.SubscriptionUndecodable' — the default+ decode retry policy exhausted. The exception carries 'DecodeFailure';+ the checkpoint remains before the failed event. Fix the hook and restart+ the same subscription to replay it. * @Left e@ where @e@ is a 'Hasql.Pool.UsageError' from checkpoint load — startup could not read the saved checkpoint. The worker stops rather than silently replaying from global position 0.@@ -123,12 +126,6 @@ -} subscribe :: (MonadIO m) => KirokuStore -> SubscriptionConfig -> m SubscriptionHandle subscribe store config = liftIO $ do- -- Fail fast on a misconfigured group, before any thread is spawned. A bad- -- (member, size) is a programmer error; throwing here keeps subscribe's- -- non-Either signature intact for every existing caller (see EP-2 Decision Log).- for_ (consumerGroup config) $ \(ConsumerGroup m n) ->- when (n < 1 || m < 0 || m >= n) $- throwIO (InvalidConsumerGroup m n) mask $ \_restore -> bracketOnError ( case (consumerGroup config, target config) of
src/Kiroku/Store/Subscription/Checkpoint.hs view
@@ -1,14 +1,20 @@ {- | Explicit mutation operations for durable subscription checkpoints. Ordinary subscription checkpoint saves are monotonic. This module owns the-separate, deliberately named reset operation for callers that need to move+separate reset, resize and rebind operations for callers that need to move persisted progress backward or forward as part of a larger transaction. -} module Kiroku.Store.Subscription.Checkpoint ( SubscriptionCheckpointResetReport (..), resetSubscriptionCheckpointsTx,+ ConsumerGroupResizeReport (..),+ resizeConsumerGroupTx,+ SubscriptionTargetRebindReport (..),+ rebindSubscriptionTargetTx, ) where +import Data.Int (Int32)+import Data.List (nub, sort) import Data.List.NonEmpty (NonEmpty) import Data.List.NonEmpty qualified as NonEmpty import Data.Vector (Vector)@@ -17,8 +23,11 @@ import Hasql.Transaction qualified as Tx import Kiroku.Store.Subscription.Checkpoint.SQL qualified as SQL import Kiroku.Store.Subscription.Types (+ ConsumerGroupSize, SubscriptionCheckpointKey (..), SubscriptionName (..),+ SubscriptionTarget,+ consumerGroupSizeValue, ) import Kiroku.Store.Types (GlobalPosition (..)) @@ -68,3 +77,67 @@ resetKey (_, Nothing) = Nothing missingName (name, Nothing) = Just (SubscriptionName name) missingName (_, Just _) = Nothing++{- | Evidence of explicit topology equalization. Previous sizes are sorted and+distinct; all new members resume from 'resumePosition'. A missing group starts+at zero. A repeat reports the now-equalized topology without moving progress.+-}+data ConsumerGroupResizeReport = ConsumerGroupResizeReport+ { previousSizes :: !(Vector Int32)+ , previousMemberCount :: !Int+ , newSize :: !ConsumerGroupSize+ , resumePosition :: !GlobalPosition+ }+ deriving stock (Eq, Show, Generic)++{- | Stop every worker for the name before calling this operation. Lock its+checkpoint set, rewind every new member to the old minimum, and remove obsolete+members atomically. Existing member identities are retained. A hash-assignment+change also requires this equalization even when the group size is unchanged.+The caller can compose or roll back the resize with application-owned SQL.+-}+resizeConsumerGroupTx ::+ SubscriptionName -> ConsumerGroupSize -> Tx.Transaction ConsumerGroupResizeReport+resizeConsumerGroupTx (SubscriptionName name) newSize = do+ Tx.statement name SQL.lockCheckpointNameStmt+ rows <- Tx.statement name SQL.lockCheckpointRowsStmt+ let position = if Vector.null rows then 0 else Vector.minimum (Vector.map (\(_, _, p, _) -> p) rows)+ previousSizes = Vector.fromList . sort . nub $ [n | (_, n, _, _) <- Vector.toList rows]+ let binding = if Vector.null rows then Nothing else let (_, _, _, b) = Vector.head rows in b+ (kind, category) = SQL.targetColumns binding+ Tx.statement (name, consumerGroupSizeValue newSize, position, kind, category) SQL.resizeCheckpointMembersStmt+ pure+ ConsumerGroupResizeReport+ { previousSizes = previousSizes+ , previousMemberCount = Vector.length rows+ , newSize = newSize+ , resumePosition = GlobalPosition position+ }++-- | Rebind evidence includes every distinct prior binding; Nothing denotes legacy rows.+data SubscriptionTargetRebindReport = SubscriptionTargetRebindReport+ { reboundMemberCount :: !Int+ , previousBindings :: !(Vector (Maybe SubscriptionTarget))+ , reboundTarget :: !SubscriptionTarget+ , reboundPosition :: !GlobalPosition+ }+ deriving stock (Eq, Show)++{- | Stop all workers first. Bind every existing member and explicitly reset its+position in the caller's transaction. A missing name fails the transaction;+repeating the operation leaves target and progress unchanged.+-}+rebindSubscriptionTargetTx ::+ SubscriptionName -> SubscriptionTarget -> GlobalPosition -> Tx.Transaction SubscriptionTargetRebindReport+rebindSubscriptionTargetTx (SubscriptionName name) target position@(GlobalPosition pos) = do+ Tx.statement name SQL.lockCheckpointNameStmt+ rows <- Tx.statement name SQL.lockCheckpointRowsStmt+ let (kind, category) = SQL.targetColumns (Just target)+ count <- Tx.statement (name, kind, category, pos) SQL.rebindCheckpointTargetStmt+ pure+ SubscriptionTargetRebindReport+ { reboundMemberCount = fromIntegral count+ , previousBindings = Vector.fromList . nub $ [binding | (_, _, _, binding) <- Vector.toList rows]+ , reboundTarget = target+ , reboundPosition = position+ }
src/Kiroku/Store/Subscription/Checkpoint/SQL.hs view
@@ -1,27 +1,45 @@ {-# LANGUAGE MultilineStrings #-}+{-# LANGUAGE RankNTypes #-} -- | Package-internal SQL for subscription checkpoint lifecycle operations. module Kiroku.Store.Subscription.Checkpoint.SQL ( initializeSubscriptionCheckpointSession,+ initializeWorkerCheckpointSession,+ lockCheckpointNameStmt,+ lockCheckpointRowsStmt,+ resizeCheckpointMembersStmt, resetSubscriptionCheckpointsStmt,+ targetColumns,+ saveBoundCheckpointSession,+ rebindCheckpointTargetStmt, ) where -import Contravariant.Extras (contrazip2, contrazip3)+import Contravariant.Extras (contrazip2, contrazip4, contrazip5, contrazip6) import Data.Int (Int32, Int64)+import Data.List (nub, sort) import Data.Text (Text) import Data.Vector (Vector)+import Data.Vector qualified as Vector import Hasql.Decoders qualified as D import Hasql.Encoders qualified as E import Hasql.Session qualified as Session import Hasql.Statement (Statement, preparable)+import Hasql.Transaction qualified as Tx+import Hasql.Transaction.Sessions qualified as TxSessions+import Kiroku.Store.SQL qualified as SQL import Kiroku.Store.Subscription.Types ( CheckpointInitialization (..),+ ConsumerGroupSizeMismatch (..), MissingCheckpointPolicy (..),+ SomeSubscriptionStartupFailure (..), SubscriptionCheckpointKey (..), SubscriptionCheckpointMissing (..), SubscriptionName (..),+ SubscriptionTarget (..),+ SubscriptionTargetMismatch (..),+ TargetBindingPolicy (..), )-import Kiroku.Store.Types (GlobalPosition (..))+import Kiroku.Store.Types (CategoryName (..), GlobalPosition (..)) {- | Resolve one checkpoint key in a Hasql session. @@ -37,8 +55,55 @@ Int32 -> MissingCheckpointPolicy -> Session.Session (Either SubscriptionCheckpointMissing CheckpointInitialization)-initializeSubscriptionCheckpointSession subscriptionName@(SubscriptionName name) member policy = do- first <- Session.statement (name, member, policyCode policy) initializeSubscriptionCheckpointStmt+initializeSubscriptionCheckpointSession name member policy =+ initializeCheckpointWith Session.statement name member 1 Nothing policy++{- | Startup-only validation and insertion share one checkout and one transaction.+The name lock serializes competing topologies, including an initially absent+row set. It is never acquired by ordinary checkpoint saves or event appends.+-}+initializeWorkerCheckpointSession ::+ SubscriptionName ->+ Int32 ->+ Int32 ->+ SubscriptionTarget ->+ TargetBindingPolicy ->+ MissingCheckpointPolicy ->+ Session.Session (Either SomeSubscriptionStartupFailure (CheckpointInitialization, Bool))+initializeWorkerCheckpointSession subscriptionName@(SubscriptionName name) member configured target bindingPolicy policy =+ TxSessions.transaction TxSessions.ReadCommitted TxSessions.Write $ do+ Tx.statement name lockCheckpointNameStmt+ rows <- Tx.statement name readCheckpointRowsStmt+ let sizes = Vector.fromList . sort . nub . fmap (\(_, n, _, _) -> n) $ Vector.toList rows+ bindings = Vector.fromList . nub $ [binding | (_, _, _, binding) <- Vector.toList rows]+ unbound = not (Vector.null rows) && bindings == Vector.singleton Nothing+ matches = Vector.all (== Just target) bindings+ if Vector.any (/= configured) sizes+ then pure (Left (SomeSubscriptionStartupFailure (ConsumerGroupSizeMismatch subscriptionName configured sizes)))+ else+ if not matches && not (unbound && bindingPolicy == AdoptUnbound)+ then pure (Left (SomeSubscriptionStartupFailure (SubscriptionTargetMismatch subscriptionName target bindings)))+ else do+ -- Check exact-key absence before adoption so a refused startup never+ -- mutates sibling bindings. Both operations share the name lock.+ resolution <- initializeCheckpointWith Tx.statement subscriptionName member configured (Just target) policy+ case resolution of+ Left missing -> pure (Left (SomeSubscriptionStartupFailure missing))+ Right initialized -> do+ if unbound then Tx.statement (name, targetColumns (Just target)) adoptCheckpointTargetStmt else pure ()+ pure (Right (initialized, unbound))++initializeCheckpointWith ::+ (Monad m) =>+ (forall a b. a -> Statement a b -> m b) ->+ SubscriptionName ->+ Int32 ->+ Int32 ->+ Maybe SubscriptionTarget ->+ MissingCheckpointPolicy ->+ m (Either SubscriptionCheckpointMissing CheckpointInitialization)+initializeCheckpointWith statement subscriptionName@(SubscriptionName name) member groupSize binding policy = do+ first <- statement (name, member, policyCode policy, groupSize, fst (targetColumns binding), snd (targetColumns binding)) initializeSubscriptionCheckpointStmt case first of Just result -> pure (Right (decodeResult result)) Nothing -> case policy of@@ -48,7 +113,7 @@ -- statement's snapshot. A fresh statement snapshot observes -- the committed winner; singleRow turns a violated invariant -- into a structured Hasql session error.- position <- Session.statement (name, member) readInitializedCheckpointStmt+ position <- statement (name, member) readInitializedCheckpointStmt pure (Right (ExistingCheckpoint key (GlobalPosition position))) where key = SubscriptionCheckpointKey subscriptionName member@@ -63,7 +128,7 @@ FromCurrentHead -> "from_current_head" FailIfMissing -> "fail_if_missing" -initializeSubscriptionCheckpointStmt :: Statement (Text, Int32, Text) (Maybe (Int64, Bool))+initializeSubscriptionCheckpointStmt :: Statement (Text, Int32, Text, Int32, Text, Maybe Text) (Maybe (Int64, Bool)) initializeSubscriptionCheckpointStmt = preparable """@@ -78,8 +143,8 @@ ), inserted AS ( INSERT INTO subscriptions- (subscription_name, consumer_group_member, last_seen, updated_at)- SELECT $1, $2, desired.last_seen, now()+ (subscription_name, consumer_group_member, last_seen, updated_at, consumer_group_size, target_kind, target_category)+ SELECT $1, $2, desired.last_seen, now(), $4, $5, $6 FROM desired WHERE desired.last_seen IS NOT NULL ON CONFLICT (subscription_name, consumer_group_member) DO NOTHING@@ -94,10 +159,13 @@ AND consumer_group_member = $2 LIMIT 1 """- ( contrazip3+ ( contrazip6 (E.param (E.nonNullable E.text)) (E.param (E.nonNullable E.int4)) (E.param (E.nonNullable E.text))+ (E.param (E.nonNullable E.int4))+ (E.param (E.nonNullable E.text))+ (E.param (E.nullable E.text)) ) ( D.rowMaybe $ (,)@@ -158,3 +226,97 @@ <$> D.column (D.nonNullable D.text) <*> D.column (D.nullable D.int4) )++-- | A separate key domain from the optional member guard.+lockCheckpointNameStmt :: Statement Text ()+lockCheckpointNameStmt =+ preparable+ "SELECT pg_advisory_xact_lock(hashtextextended('kiroku:checkpoint-topology:' || $1, 0))"+ (E.param (E.nonNullable E.text))+ D.noResult++readCheckpointRowsStmt :: Statement Text (Vector (Int32, Int32, Int64, Maybe SubscriptionTarget))+readCheckpointRowsStmt = checkpointRowsStmt ""++lockCheckpointRowsStmt :: Statement Text (Vector (Int32, Int32, Int64, Maybe SubscriptionTarget))+lockCheckpointRowsStmt = checkpointRowsStmt " FOR UPDATE"++checkpointRowsStmt :: Text -> Statement Text (Vector (Int32, Int32, Int64, Maybe SubscriptionTarget))+checkpointRowsStmt suffix =+ preparable+ ("SELECT consumer_group_member, consumer_group_size, last_seen, target_kind, target_category FROM subscriptions WHERE subscription_name = $1 ORDER BY consumer_group_member" <> suffix)+ (E.param (E.nonNullable E.text))+ (D.rowVector ((,,,) <$> D.column (D.nonNullable D.int4) <*> D.column (D.nonNullable D.int4) <*> D.column (D.nonNullable D.int8) <*> targetBindingRow))++{- | Keep existing row identities, equalize every new member, remove obsolete+members. Caller already holds the name and row locks; workers must be stopped.+-}+resizeCheckpointMembersStmt :: Statement (Text, Int32, Int64, Text, Maybe Text) ()+resizeCheckpointMembersStmt =+ preparable+ """+ WITH removed AS (+ DELETE FROM subscriptions+ WHERE subscription_name = $1 AND consumer_group_member >= $2+ )+ INSERT INTO subscriptions+ (subscription_name, consumer_group_member, consumer_group_size, last_seen, updated_at, target_kind, target_category)+ SELECT $1, member, $2, $3, now(), $4, $5 FROM generate_series(0, $2 - 1) AS member+ ON CONFLICT (subscription_name, consumer_group_member) DO UPDATE+ SET consumer_group_size = EXCLUDED.consumer_group_size,+ last_seen = EXCLUDED.last_seen, updated_at = now()+ """+ ( contrazip5+ (E.param (E.nonNullable E.text))+ (E.param (E.nonNullable E.int4))+ (E.param (E.nonNullable E.int8))+ (E.param (E.nonNullable E.text))+ (E.param (E.nullable E.text))+ )+ D.noResult++{- | Select a fixed-kind statement once for the configured target. Both forms+write the full binding and retain the same unconditional monotonic upsert.+-}+saveBoundCheckpointSession :: SubscriptionTarget -> Text -> Int32 -> Int64 -> Int32 -> Session.Session ()+saveBoundCheckpointSession AllStreams name member position groupSize =+ Session.statement (name, member, position, groupSize) SQL.saveAllCheckpointMemberStmt+saveBoundCheckpointSession (Category (CategoryName category)) name member position groupSize =+ Session.statement (name, member, position, groupSize, category) SQL.saveCategoryCheckpointMemberStmt++-- | One encoder for every checkpoint target write, including legacy provisioning.+targetColumns :: Maybe SubscriptionTarget -> (Text, Maybe Text)+targetColumns Nothing = ("unbound", Nothing)+targetColumns (Just AllStreams) = ("all", Nothing)+targetColumns (Just (Category (CategoryName category))) = ("category", Just category)++-- The CHECK constraints make this decoder total over valid database rows.+targetBindingRow :: D.Row (Maybe SubscriptionTarget)+targetBindingRow = decode <$> D.column (D.nonNullable D.text) <*> D.column (D.nullable D.text)+ where+ decode "all" Nothing = Just AllStreams+ decode "category" (Just category) = Just (Category (CategoryName category))+ decode _ _ = Nothing++adoptCheckpointTargetStmt :: Statement (Text, (Text, Maybe Text)) ()+adoptCheckpointTargetStmt =+ preparable+ "UPDATE subscriptions SET target_kind = $2, target_category = $3 WHERE subscription_name = $1"+ (contrazip2 (E.param (E.nonNullable E.text)) (contrazip2 (E.param (E.nonNullable E.text)) (E.param (E.nullable E.text))))+ D.noResult++-- A singleRow decoder refuses an absent name and aborts the surrounding transaction.+rebindCheckpointTargetStmt :: Statement (Text, Text, Maybe Text, Int64) Int64+rebindCheckpointTargetStmt =+ preparable+ """+ WITH rebound AS (+ UPDATE subscriptions+ SET target_kind = $2, target_category = $3, last_seen = $4, updated_at = now()+ WHERE subscription_name = $1+ RETURNING subscription_id+ )+ SELECT count(*) FROM rebound HAVING count(*) > 0+ """+ (contrazip4 (E.param (E.nonNullable E.text)) (E.param (E.nonNullable E.text)) (E.param (E.nullable E.text)) (E.param (E.nonNullable E.int8)))+ (D.singleRow (D.column (D.nonNullable D.int8)))
src/Kiroku/Store/Subscription/Effect.hs view
@@ -121,11 +121,15 @@ runSubscription store = interpret $ \env -> \case Subscribe config -> localUnliftIO env (ConcUnlift Persistent (Limited 1)) $ \unlift -> do- -- Record update intentionally: every field except 'handler' is preserved,+ -- Record update intentionally: policy fields are preserved; both callbacks are unlifted, -- including 'consumerGroup' and 'consumerGroupGuard' (added in EP-2). Do not -- switch to a full record literal here — that would silently reset new fields -- and drop a caller's consumer-group membership.- let ioConfig = config{handler = \evt -> unlift (handler config evt)}+ let ioConfig =+ config+ { handler = \evt -> unlift (handler config evt)+ , undecodableHandler = fmap (\callback event failure -> unlift (callback event failure)) (undecodableHandler config)+ } Sub.subscribe store ioConfig -- | Interpret Subscription by reading the store handle from 'KirokuStoreResource'.
src/Kiroku/Store/Subscription/EventPublisher.hs view
@@ -58,14 +58,13 @@ import Data.Int (Int32) import Data.IntMap.Strict (IntMap) import Data.IntMap.Strict qualified as IntMap-import Data.Vector (Vector) import Data.Vector qualified as V import Hasql.Pool (Pool) import Hasql.Pool qualified as Pool import Hasql.Session qualified as Session import Kiroku.Store.Observability (KirokuEvent (..), emitOrDrop) import Kiroku.Store.SQL qualified as SQL-import Kiroku.Store.Settings (StoreSettings, decodeEvents)+import Kiroku.Store.Settings (DecodedBatch (..), DecodedEvent (..), StoreSettings, decodeEvents) import Kiroku.Store.Subscription.Types (OverflowPolicy (..)) import Kiroku.Store.Types (GlobalPosition (..), RecordedEvent (..)) import Numeric.Natural (Natural)@@ -91,7 +90,7 @@ terminates the subscription with 'SubscriptionOverflowed'. -} data Subscriber = Subscriber- { subQueue :: !(TBQueue (Vector RecordedEvent))+ { subQueue :: !(TBQueue DecodedBatch) , subStatus :: !(TVar SubscriberStatus) , subPolicy :: !OverflowPolicy }@@ -183,7 +182,7 @@ -- | Queue capacity (number of batches) Natural -> OverflowPolicy ->- STM (TBQueue (Vector RecordedEvent), TVar SubscriberStatus, IO ())+ STM (TBQueue DecodedBatch, TVar SubscriberStatus, IO ()) subscribePublisher pub cap policy = do queue <- newTBQueue cap status <- newTVar Active@@ -271,7 +270,12 @@ -- observes the same transformed view, and the cost is -- paid in one place instead of per-subscriber. events <- decodeEvents stSettings rawEvents- let lastEvent = V.last events+ case events of+ UnchangedBatch _ -> pure ()+ TransformedBatch decoded -> for_ decoded $ \case+ Decoded _ -> pure ()+ Undecodable event failure -> emitOrDrop mHandler (KirokuEventPublisherDecodeFailed (globalPosition event) (eventId event) failure)+ let lastEvent = V.last rawEvents newPos = globalPosition lastEvent -- Snapshot the current subscriber set, then deliver outside -- the snapshot's STM transaction. Each delivery is its own@@ -291,7 +295,7 @@ for_ (IntMap.elems (subs' `IntMap.difference` subs)) (deliverBatchSTM events) writeTVar posVar newPos -- If we got a full batch, there may be more — loop immediately- if V.length events >= fromIntegral publisherBatchSize+ if V.length rawEvents >= fromIntegral publisherBatchSize then fullFetch else pure ()
src/Kiroku/Store/Subscription/Fsm.hs view
@@ -60,10 +60,9 @@ import Data.Text (Text) import Data.Text qualified as T import Data.Time (NominalDiffTime)-import Data.Vector (Vector)-import Data.Vector qualified as V import Hasql.Pool qualified as Pool-import Kiroku.Store.Types (GlobalPosition (..), RecordedEvent (..))+import Kiroku.Store.Settings (DecodeFailure (..), DecodedBatch, decodedBatchLastEvent, decodedBatchLength)+import Kiroku.Store.Types (EventId (..), GlobalPosition (..), RecordedEvent (..)) {- | Why a subscription's worker thread stopped. @@ -92,6 +91,8 @@ exception). The 'SomeException' carries the cause. -} StopWorkerCrashed !SomeException+ | -- | Default undecodable-event retries exhausted; checkpoint remains before it.+ StopUndecodable !DecodeFailure deriving stock (Show) {- | How long to wait before redelivering a retried event.@@ -126,6 +127,8 @@ DeadLetterMaxAttempts !Int | -- | A custom reason: a summary plus structured JSON detail. DeadLetterOther !Text !Value+ | -- | A consumer explicitly chose to skip an undecodable event.+ DeadLetterDecodeFailure !DecodeFailure deriving stock (Eq, Show) -- | The short operator-facing summary stored in @dead_letters.reason_summary@.@@ -135,6 +138,7 @@ DeadLetterInvalid detail -> "invalid payload: " <> detail DeadLetterMaxAttempts n -> "max retry attempts exceeded (" <> T.pack (show n) <> ")" DeadLetterOther summary _ -> summary+ DeadLetterDecodeFailure failure -> "decode failure: " <> decodeFailureReason failure -- | The structured JSON detail stored in @dead_letters.reason@ (JSONB). deadLetterReasonJson :: DeadLetterReason -> Value@@ -143,6 +147,7 @@ DeadLetterInvalid detail -> object ["kind" .= ("invalid_payload" :: Text), "detail" .= detail] DeadLetterMaxAttempts n -> object ["kind" .= ("max_attempts_exceeded" :: Text), "attempts" .= n] DeadLetterOther summary detail -> object ["kind" .= ("other" :: Text), "summary" .= summary, "detail" .= detail]+ DeadLetterDecodeFailure failure -> object ["kind" .= ("decode_failure" :: Text), "event_id" .= (case decodeFailureEventId failure of EventId eid -> show eid), "detail" .= decodeFailureReason failure] {- | What unblocks a 'Paused' worker. @@ -225,7 +230,7 @@ -} data Input = -- | A non-empty history\/live batch arrived.- BatchFetched !(Vector RecordedEvent)+ BatchFetched !DecodedBatch | -- | A fetch returned no rows (catch-up is complete). FetchEmpty | -- | A fetch hit a database error.@@ -245,7 +250,7 @@ | -- | The worker drained the stale queue and is ready to recover (re-catch-up). QueueDrained | -- | The worker lost its database pool while live.- ConnectionLost !Pool.UsageError+ ConnectionLost !GlobalPosition !Pool.UsageError | -- | The caller cancelled the worker. Cancelled deriving stock (Show)@@ -262,7 +267,7 @@ | -- | Obtain the next live batch via the active live strategy. RunLive | -- | Call the handler per event, checkpointing at the batch tail.- DeliverBatch !(Vector RecordedEvent)+ DeliverBatch !DecodedBatch | {- | Persist the checkpoint at this position. (Named 'Checkpoint' rather than @SaveCheckpoint@ to avoid clashing with the 'Kiroku.Store.Observability.SubscriptionDbPhase' constructor of that name.)@@ -287,10 +292,10 @@ -- The last event's position in a batch, falling back to the given cursor when -- the batch is empty (the driver never feeds 'step' an empty 'BatchFetched', -- so the fallback is defensive).-lastPos :: GlobalPosition -> Vector RecordedEvent -> GlobalPosition+lastPos :: GlobalPosition -> DecodedBatch -> GlobalPosition lastPos fallback evs- | V.null evs = fallback- | otherwise = globalPosition (V.last evs)+ | decodedBatchLength evs == 0 = fallback+ | otherwise = globalPosition (decodedBatchLastEvent evs) {- | The single transition function. Given the current state and an input, return the next state and the effects to perform.@@ -312,7 +317,7 @@ QueueOverflowed -> (Stopped StopOverflowed, [Halt StopOverflowed]) QueueBackpressured -> (CatchingUp c n, []) -- defensive: catch-up reads the DB, not the queue QueueDrained -> (CatchingUp c n, [])- ConnectionLost _ -> (Reconnecting c 1, [EmitReconnecting 1, Backoff 1])+ ConnectionLost observed _ -> (Reconnecting (max c observed) 1, [EmitReconnecting 1, Backoff 1]) Cancelled -> (Stopped StopCancelled, [Halt StopCancelled]) Live c -> case input of BatchFetched evs -> (Live (lastPos c evs), [DeliverBatch evs])@@ -323,7 +328,7 @@ QueueOverflowed -> (Stopped StopOverflowed, [Halt StopOverflowed]) QueueBackpressured -> (Paused c ResumeOnDrain, [EmitPaused]) QueueDrained -> (Live c, [RunLive])- ConnectionLost _ -> (Reconnecting c 1, [EmitReconnecting 1, Backoff 1])+ ConnectionLost observed _ -> (Reconnecting (max c observed) 1, [EmitReconnecting 1, Backoff 1]) Cancelled -> (Stopped StopCancelled, [Halt StopCancelled]) Paused c rc -> case input of -- The worker drained the stale queue and cleared the pause flag; recover@@ -338,7 +343,7 @@ BatchFetched evs -> (CatchingUp (lastPos c evs) 0, [DeliverBatch evs]) FetchEmpty -> (Live c, []) FetchFailed _ -> (Reconnecting c (n + 1), [EmitReconnecting (n + 1), Backoff (n + 1)])- ConnectionLost _ -> (Reconnecting c (n + 1), [EmitReconnecting (n + 1), Backoff (n + 1)])+ ConnectionLost observed _ -> (Reconnecting (max c observed) (n + 1), [EmitReconnecting (n + 1), Backoff (n + 1)]) CaughtUp -> (CatchingUp c 0, []) HandlerStopped _ -> (Stopped StopHandlerRequested, [Halt StopHandlerRequested]) Cancelled -> (Stopped StopCancelled, [Halt StopCancelled])
src/Kiroku/Store/Subscription/Stream.hs view
@@ -32,6 +32,10 @@ -- * Ack-coupled pull stream AckItem (..), InvalidStreamBufferSize (..),+ StreamBufferSize,+ mkStreamBufferSize,+ streamBufferSizeValue,+ defaultStreamBufferSize, subscriptionAckStream, ) where@@ -54,8 +58,7 @@ writeTBQueue, writeTVar, )-import Control.Exception (Exception, SomeException, fromException, mask, onException, throwIO)-import Control.Monad (when)+import Control.Exception (SomeException, fromException, mask, onException, throwIO) import Data.IORef (atomicModifyIORef', newIORef) import Kiroku.Store.Connection (KirokuStore) import Kiroku.Store.Subscription (subscribe)@@ -91,16 +94,29 @@ -- ^ one-shot reply the consumer must fill exactly once } -{- | Thrown by 'subscriptionAckStream' when the requested bridge queue capacity-is zero.+{- | Construction error from 'mkStreamBufferSize' for a zero capacity. A zero-capacity 'TBQueue' would make the bridge handler block forever on its first delivery, before a stream consumer can ever see the event or reply to it. -} newtype InvalidStreamBufferSize = InvalidStreamBufferSize Natural deriving stock (Eq, Show)- deriving anyclass (Exception) +-- | Positive capacity, validated before allocating a bridge or starting a worker.+newtype StreamBufferSize = StreamBufferSize Natural+ deriving stock (Eq, Show)++mkStreamBufferSize :: Natural -> Either InvalidStreamBufferSize StreamBufferSize+mkStreamBufferSize n+ | n >= 1 = Right (StreamBufferSize n)+ | otherwise = Left (InvalidStreamBufferSize n)++streamBufferSizeValue :: StreamBufferSize -> Natural+streamBufferSizeValue (StreamBufferSize n) = n++defaultStreamBufferSize :: StreamBufferSize+defaultStreamBufferSize = StreamBufferSize 256+ data BridgeTermination = BridgeClosedCleanly | BridgeCrashed !SomeException@@ -141,7 +157,7 @@ KirokuStore -> SubscriptionConfig -> -- | TBQueue capacity for the bridge; must be at least 1.- Natural ->+ StreamBufferSize -> IO (Stream IO RecordedEvent, IO ()) subscriptionStream store config bufferSize = do (ackStream, cancelAction) <- subscriptionAckStream store config bufferSize@@ -174,12 +190,10 @@ KirokuStore -> SubscriptionConfig -> -- | TBQueue capacity for the bridge; must be at least 1.- Natural ->+ StreamBufferSize -> IO (Stream IO AckItem, IO ()) subscriptionAckStream store config bufferSize = mask $ \restore -> do- when (bufferSize < 1) $- throwIO (InvalidStreamBufferSize bufferSize)- queue <- newTBQueueIO bufferSize+ queue <- newTBQueueIO (streamBufferSizeValue bufferSize) closedVar <- newTVarIO Nothing -- Tracks the previous (eventId, attempt) so a consecutive redelivery of the -- same event (the worker's bounded retry) is reported with an incremented
src/Kiroku/Store/Subscription/Types.hs view
@@ -1,3 +1,5 @@+{-# LANGUAGE ExistentialQuantification #-}+ {- | Core configuration and result types for subscriptions. This module defines the user-facing vocabulary for starting a subscription and@@ -26,9 +28,19 @@ SubscriptionCheckpoint (..), SubscriptionCheckpointInventory (..), SubscriptionTarget (..),+ TargetBindingPolicy (..),+ SubscriptionTargetMismatch (..),+ SomeSubscriptionStartupFailure (..),+ BatchSize,+ mkBatchSize,+ batchSizeValue,+ defaultBatchSize,+ InvalidBatchSize (..), SubscriptionResult (..), OverflowPolicy (..), SubscriptionOverflowed (..),+ SubscriptionUndecodable (..),+ InvalidHandlerStallWarnAfter (..), EventHandlerM, EventHandler, SubscriptionConfigM (..),@@ -52,19 +64,29 @@ defaultRetryPolicy, -- * Consumer groups- ConsumerGroup (..),+ ConsumerGroup,+ member,+ size,+ ConsumerGroupSize,+ consumerGroupSizeValue,+ mkConsumerGroupSize,+ mkConsumerGroup,+ ConsumerGroupSizeMismatch (..), InvalidConsumerGroup (..), ConsumerGroupGuardConflict (..), ) where -import Control.Exception (Exception, SomeException)+import Control.Exception (Exception (..), SomeException) import Data.Int (Int32) import Data.Set (Set) import Data.Set qualified as Set import Data.Text (Text)+import Data.Time (NominalDiffTime) import Data.Time.Clock (UTCTime)+import Data.Typeable (cast) import Data.Vector (Vector) import GHC.Generics (Generic)+import Kiroku.Store.Settings (DecodeFailure) import Kiroku.Store.Subscription.Fsm ( DeadLetterReason (..), RetryDelay (..),@@ -190,8 +212,13 @@ { checkpointKey :: SubscriptionCheckpointKey } deriving stock (Eq, Show, Generic)- deriving anyclass (Exception) +instance Exception SubscriptionCheckpointMissing where+ toException = toException . SomeSubscriptionStartupFailure+ fromException exception = do+ SomeSubscriptionStartupFailure refusal <- fromException exception+ cast refusal+ -- | One checkpoint row that has been durably persisted by a subscription. data SubscriptionCheckpoint = SubscriptionCheckpoint { subscriptionName :: !SubscriptionName@@ -307,6 +334,16 @@ deriving stock (Show) deriving anyclass (Exception) +{- | Default decode retries exhausted. Wait returns this exception and the+stopped event carries StopUndecodable; the durable checkpoint is still before+the failed event. Fix the hook and restart to replay it.+-}+newtype SubscriptionUndecodable = SubscriptionUndecodable+ { subscriptionDecodeFailure :: DecodeFailure+ }+ deriving stock (Eq, Show)+ deriving anyclass (Exception)+ -- | Handler callback invoked for each event, parameterized by monad. type EventHandlerM m = RecordedEvent -> m SubscriptionResult @@ -318,7 +355,18 @@ { name :: !SubscriptionName , target :: !SubscriptionTarget , handler :: !(EventHandlerM m)- , batchSize :: !Int32+ , undecodableHandler :: !(Maybe (RecordedEvent -> DecodeFailure -> m SubscriptionResult))+ {- ^ Optional disposition of a raw undecodable event. Default Nothing:+ retry at one-second spacing and then stop without skipping the event.+ A callback uses the ordinary Continue/Stop/Retry/DeadLetter vocabulary.+ -}+ , handlerStallWarnAfter :: !(Maybe NominalDiffTime)+ {- ^ Optional positive warning interval for a handler holding one event.+ Default Nothing allocates no diagnostic cell, watchdog or timer and does+ no diagnostic work per delivery. Warnings are advisory and never finalize+ an acknowledgement or advance a checkpoint. Nonpositive values fail startup.+ -}+ , batchSize :: !BatchSize -- ^ Number of events to fetch per batch during catch-up (default: 100) , queueCapacity :: !Natural {- ^ Maximum number of /batches/ the publisher may enqueue for this@@ -342,8 +390,7 @@ {- ^ 'Nothing' (the default) = ordinary single-consumer subscription. 'Just cg' = this worker is member 'member cg' of a group of size 'size cg'. The invariant @size >= 1@ and @0 <= member < size@ is- enforced once at 'Kiroku.Store.Subscription.subscribe' time, which- throws 'InvalidConsumerGroup' on violation.+ enforced at construction by 'mkConsumerGroupSize' and 'mkConsumerGroup'. -} , consumerGroupGuard :: !Bool {- ^ When 'True' (default 'False'), the worker performs a one-shot@@ -361,6 +408,8 @@ always take precedence, so changing this field never rewinds or advances durable progress. -}+ , targetBindingPolicy :: !TargetBindingPolicy+ -- ^ Adopt legacy unbound checkpoints (default) or require an existing binding. , retryPolicy :: !RetryPolicy {- ^ Bounds redelivery of an event for which the handler returned 'Retry' before the worker dead-letters it. Default: 'defaultRetryPolicy'@@ -427,12 +476,15 @@ { name = name' , target = target' , handler = handler'- , batchSize = 100+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = PauseAndResume , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -470,26 +522,56 @@ -- | Handle defaulting to 'IO'. type SubscriptionHandle = SubscriptionHandleM IO --- | Static consumer-group membership for a subscription.-data ConsumerGroup = ConsumerGroup- { member :: !Int32- -- ^ 0-based member index; must satisfy @0 <= member < size@.- , size :: !Int32- -- ^ total members in the group; must be @>= 1@.- }+-- | A positive group size, validated once before starting any workers.+newtype ConsumerGroupSize = ConsumerGroupSize Int32 deriving stock (Eq, Show) -{- | Thrown by 'Kiroku.Store.Subscription.subscribe' when a 'ConsumerGroup'-violates @size >= 1@ or @0 <= member < size@. Carries the offending values for-diagnostics.--}+consumerGroupSizeValue :: ConsumerGroupSize -> Int32+consumerGroupSizeValue (ConsumerGroupSize n) = n++mkConsumerGroupSize :: Int32 -> Either InvalidConsumerGroup ConsumerGroupSize+mkConsumerGroupSize n+ | n >= 1 = Right (ConsumerGroupSize n)+ | otherwise = Left (InvalidConsumerGroup 0 n)++-- | Static membership. Use 'mkConsumerGroup'; accessors cannot update the pair.+data ConsumerGroup = ConsumerGroup !Int32 !ConsumerGroupSize+ deriving stock (Eq, Show)++member :: ConsumerGroup -> Int32+member (ConsumerGroup m _) = m++size :: ConsumerGroup -> Int32+size (ConsumerGroup _ n) = consumerGroupSizeValue n++mkConsumerGroup :: Int32 -> ConsumerGroupSize -> Either InvalidConsumerGroup ConsumerGroup+mkConsumerGroup m n+ | m >= 0 && m < consumerGroupSizeValue n = Right (ConsumerGroup m n)+ | otherwise = Left (InvalidConsumerGroup m (consumerGroupSizeValue n))++-- | Construction error; invalid configuration never reaches a worker. data InvalidConsumerGroup = InvalidConsumerGroup { invalidMember :: !Int32 , invalidSize :: !Int32 }- deriving stock (Show)- deriving anyclass (Exception)+ deriving stock (Eq, Show) +{- | Stored topology disagrees with the configured size. Stop all members and+use 'Kiroku.Store.Subscription.Checkpoint.resizeConsumerGroupTx' before restart.+-}+data ConsumerGroupSizeMismatch = ConsumerGroupSizeMismatch+ { mismatchName :: !SubscriptionName+ , configuredSize :: !Int32+ , observedSizes :: !(Vector Int32)+ }+ deriving stock (Eq, Show)++instance Exception ConsumerGroupSizeMismatch where+ toException = toException . SomeSubscriptionStartupFailure+ fromException exception = do+ SomeSubscriptionStartupFailure refusal <- fromException exception+ cast refusal+ {- | Thrown at subscription startup when 'consumerGroupGuard' is 'True' and another holder currently holds the advisory lock for this @(name, member)@. Indicates two processes are configured as the same group member. This is a@@ -500,5 +582,64 @@ { conflictName :: !SubscriptionName , conflictMember :: !Int32 }- deriving stock (Show)- deriving anyclass (Exception)+ deriving stock (Eq, Show)++instance Exception ConsumerGroupGuardConflict where+ toException = toException . SomeSubscriptionStartupFailure+ fromException exception = do+ SomeSubscriptionStartupFailure refusal <- fromException exception+ cast refusal++-- | A nonpositive handler-stall warning interval was configured.+newtype InvalidHandlerStallWarnAfter = InvalidHandlerStallWarnAfter NominalDiffTime+ deriving stock (Eq, Show)++instance Exception InvalidHandlerStallWarnAfter where+ toException = toException . SomeSubscriptionStartupFailure+ fromException exception = do+ SomeSubscriptionStartupFailure refusal <- fromException exception+ cast refusal++-- | A catchable family of subscription startup refusals. Concrete catches work too.+data SomeSubscriptionStartupFailure = forall e. (Exception e) => SomeSubscriptionStartupFailure e++instance Show SomeSubscriptionStartupFailure where+ show (SomeSubscriptionStartupFailure e) = show e++instance Exception SomeSubscriptionStartupFailure++-- | Policy for checkpoint rows created before target identity was persisted.+data TargetBindingPolicy = AdoptUnbound | RequireBound+ deriving stock (Eq, Show)++-- | Stored bindings differ from the requested target, or strict startup found unbound rows.+data SubscriptionTargetMismatch = SubscriptionTargetMismatch+ { targetMismatchName :: !SubscriptionName+ , configuredTarget :: !SubscriptionTarget+ , observedTargets :: !(Vector (Maybe SubscriptionTarget))+ }+ deriving stock (Eq, Show)++instance Exception SubscriptionTargetMismatch where+ toException = toException . SomeSubscriptionStartupFailure+ fromException exception = do+ SomeSubscriptionStartupFailure refusal <- fromException exception+ cast refusal++-- | A positive catch-up fetch size. The constructor is deliberately hidden.+newtype BatchSize = BatchSize Int32+ deriving stock (Eq, Show)++newtype InvalidBatchSize = InvalidBatchSize Int32+ deriving stock (Eq, Show)++mkBatchSize :: Int32 -> Either InvalidBatchSize BatchSize+mkBatchSize n+ | n >= 1 = Right (BatchSize n)+ | otherwise = Left (InvalidBatchSize n)++batchSizeValue :: BatchSize -> Int32+batchSizeValue (BatchSize n) = n++defaultBatchSize :: BatchSize+defaultBatchSize = BatchSize 100
src/Kiroku/Store/Subscription/Worker.hs view
@@ -34,17 +34,21 @@ import Contravariant.Extras (contrazip2) import Control.Concurrent (threadDelay) import Control.Concurrent.Async qualified as Async-import Control.Concurrent.STM (TBQueue, TVar, atomically, check, orElse, readTBQueue, readTVar, registerDelay, tryReadTBQueue, writeTVar)-import Control.Exception (SomeException, bracket, fromException, throwIO, try)+import Control.Concurrent.STM (TBQueue, TVar, atomically, check, newTVarIO, orElse, readTBQueue, readTVar, registerDelay, tryReadTBQueue, writeTVar)+import Control.Concurrent.STM qualified as STM+import Control.Exception (SomeException, bracket, finally, fromException, mask, throwIO, toException, try)+import Control.Monad (when) import Control.Monad.IO.Class (MonadIO, liftIO) import Data.IORef (IORef, newIORef, readIORef, writeIORef) import Data.Int (Int32) import Data.Map.Strict (Map) import Data.Map.Strict qualified as Map import Data.Text (Text)+import Data.Time (NominalDiffTime) import Data.Vector (Vector) import Data.Vector qualified as V import Data.Word (Word64)+import GHC.Clock (getMonotonicTimeNSec) import Hasql.Decoders qualified as D import Hasql.Encoders qualified as E import Hasql.Pool (Pool)@@ -60,7 +64,7 @@ emitOrDrop, ) import Kiroku.Store.SQL qualified as SQL-import Kiroku.Store.Settings (StoreSettings, decodeEvents)+import Kiroku.Store.Settings (DecodedBatch (..), DecodedEvent (..), StoreSettings, decodeEvent, decodeEvents, decodedBatchLastEvent, decodedBatchLength, decodedEventRecorded, filterDecodedBatch) import Kiroku.Store.Subscription.Checkpoint.SQL qualified as CheckpointSQL import Kiroku.Store.Subscription.EventPublisher (SubscriberStatus) import Kiroku.Store.Subscription.EventPublisher qualified as Pub@@ -173,7 +177,7 @@ = {- | Non-group AllStreams: read the publisher's bounded queue; the status TVar carries Paused/Overflowed backpressure signals. -}- LiveFromPublisherQueue !(TBQueue (Vector RecordedEvent)) !(TVar SubscriberStatus)+ LiveFromPublisherQueue !(TBQueue DecodedBatch) !(TVar SubscriberStatus) | {- | 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).@@ -236,7 +240,63 @@ -} StoreSettings -> m ()-runWorker pool liveSource stateVar pubPosVar catGenVar config mHandler stSettings = liftIO $ do+runWorker pool liveSource stateVar pubPosVar catGenVar config mHandler stSettings = liftIO $+ withHandlerStallDiagnostics config (emitOrDrop mHandler) $ \deliveryConfig ->+ runWorkerBody pool liveSource stateVar pubPosVar catGenVar deliveryConfig mHandler stSettings++-- Select the delivery handler once per worker. The disabled arm returns the+-- original config directly, with no per-delivery diagnostic branch or work.+withHandlerStallDiagnostics :: SubscriptionConfig -> (KirokuEvent -> IO ()) -> (SubscriptionConfig -> IO a) -> IO a+withHandlerStallDiagnostics config emit action = case handlerStallWarnAfter config of+ Nothing -> action config+ Just threshold+ | threshold <= 0 -> throwIO (InvalidHandlerStallWarnAfter threshold)+ | otherwise -> do+ pending <- newTVarIO Nothing+ let tracked event = mask $ \restore -> do+ started <- getMonotonicTimeNSec+ let delivery = StalledDelivery (globalPosition event) (eventId event) started+ atomically (writeTVar pending (Just delivery))+ restore (handler config event) `finally` atomically (writeTVar pending Nothing)+ Async.withAsync (watchHandler pending threshold) $ \watchdog -> do+ Async.link watchdog+ action config{handler = tracked}+ where+ -- Park while idle. Reuse a pending interval across completed/replaced+ -- invocations, rather than abandoning a registerDelay timer per event.+ watchHandler pending threshold = awaitDelivery+ where+ awaitDelivery = do+ _ <- atomically $ readTVar pending >>= maybe STM.retry pure+ waitInterval threshold+ waitInterval delay = do+ timer <- registerDelay (durationMicros delay)+ atomically (readTVar timer >>= check)+ current <- atomically (readTVar pending)+ case current of+ Nothing -> awaitDelivery+ Just delivery@(StalledDelivery pos eid started) -> do+ now <- getMonotonicTimeNSec+ let elapsed = fromRational (fromIntegral (now - started) / 1_000_000_000)+ if elapsed < threshold+ then waitInterval (threshold - elapsed)+ else do+ stillPending <- atomically ((== Just delivery) <$> readTVar pending)+ when stillPending $+ emit (KirokuEventSubscriptionHandlerStalled (name config) pos eid elapsed (groupCtxOf config))+ waitInterval threshold++-- One cell per enabled worker, not one timer/thread per handler invocation.+data StalledDelivery = StalledDelivery !GlobalPosition !EventId !Word64+ deriving stock (Eq)++-- registerDelay accepts an Int. Cap huge intervals safely and round positive+-- sub-microsecond intervals up instead of creating a zero-delay polling loop.+durationMicros :: NominalDiffTime -> Int+durationMicros duration = fromInteger (max 1 (min (toInteger (maxBound :: Int)) (ceiling (duration * 1_000_000))))++runWorkerBody :: Pool -> LiveSource -> TVar SubscriptionState -> TVar GlobalPosition -> TVar (Map Text Word64) -> SubscriptionConfig -> Maybe (KirokuEvent -> IO ()) -> StoreSettings -> IO ()+runWorkerBody pool liveSource stateVar pubPosVar catGenVar config mHandler stSettings = do let emit = emitOrDrop mHandler subName = name config groupCtx = groupCtxOf config@@ -246,7 +306,7 @@ -- Optional startup guardrail: when consumerGroupGuard is on, fail fast -- if another holder currently holds this (name, member)'s advisory lock. case (consumerGroupGuard config, consumerGroup config) of- (True, Just (ConsumerGroup m _)) -> guardMember pool subName m+ (True, Just cg) -> guardMember pool subName (member cg) _ -> pure () resolution <- loadCheckpoint pool config emit checkpoint <- case resolution of@@ -314,7 +374,7 @@ case fetchResult of Left err -> pure (FetchFailed err) Right events- | V.null events -> pure CaughtUp+ | decodedBatchLength events == 0 -> pure CaughtUp | otherwise -> pure (BatchFetched events) Live c -> case liveSource of LiveFromPublisherQueue liveQueue statusVar -> do@@ -331,8 +391,8 @@ -- in the queue. Drop those stale entries so live -- mode cannot replay them or rewind the checkpoint. events <- readTBQueue liveQueue- let fresh = V.filter ((> c) . globalPosition) events- pure (if V.null fresh then FetchEmpty else BatchFetched fresh)+ let fresh = filterDecodedBatch ((> c) . globalPosition) events+ pure (if decodedBatchLength fresh == 0 then FetchEmpty else BatchFetched fresh) LiveFromCategoryNotify cat -> liveExitToInput =<< liveLoopCategoryNotify pool config stateVar catGenVar cat emit posRef c stSettings LiveFromGroupPolling ->@@ -366,7 +426,7 @@ case fetchResult of Left err -> pure (FetchFailed err) Right events- | V.null events -> pure FetchEmpty+ | decodedBatchLength events == 0 -> pure FetchEmpty | otherwise -> pure (BatchFetched events) -- Defensive totality: 'Retrying' is a surfaced observability state -- that the delivery primitive writes into the state TVar and then@@ -405,7 +465,7 @@ FetchHistory _ -> go es RunLive -> go es DeliverBatch events -> do- result <- processEvents pool config stateVar events emit posRef+ result <- processEvents pool config stateVar events emit posRef stSettings case result of Nothing -> pure (Just (HandlerStopped (lastPosOf events))) Just _ -> go es@@ -414,15 +474,18 @@ StopOverflowed -> throwIO (SubscriptionOverflowed subName) StopCancelled -> throwIO Async.AsyncCancelled StopWorkerCrashed ex -> throwIO ex+ StopUndecodable failure -> throwIO (SubscriptionUndecodable failure) - lastPosOf events = globalPosition (V.last events)+ lastPosOf events = globalPosition (decodedBatchLastEvent events) -- Map a DB-driven live loop's exit onto the next FSM input: a clean stop -- becomes 'HandlerStopped' (at the last processed position); a fetch error -- becomes 'ConnectionLost', driving the FSM into 'Reconnecting'. liveExitToInput = \case LiveHandlerStopped -> HandlerStopped <$> readIORef posRef- LiveFetchError err -> pure (ConnectionLost err)+ LiveFetchError err -> do+ position <- readIORef posRef+ pure (ConnectionLost position err) -- Read and discard every batch currently in the live queue (non-blocking). -- Used when resuming from 'Paused': the discarded events are re-read from@@ -444,7 +507,7 @@ -- The consumer-group context for this config's lifecycle events: 'NonGroup' for -- an ordinary subscription, @GroupMember member size@ for a group member. groupCtxOf :: SubscriptionConfig -> SubscriptionGroupContext-groupCtxOf config = maybe NonGroup (\(ConsumerGroup m n) -> GroupMember m n) (consumerGroup config)+groupCtxOf config = maybe NonGroup (\cg -> GroupMember (member cg) (size cg)) (consumerGroup config) {- Startup-only conflict probe for the consumer-group guardrail. Uses a transaction-scoped advisory lock ('pg_try_advisory_xact_lock') which auto-releases@@ -477,6 +540,7 @@ classifyStopReason e | Just (_ :: SubscriptionOverflowed) <- fromException e = StopOverflowed | Just (_ :: Async.AsyncCancelled) <- fromException e = StopCancelled+ | Just (SubscriptionUndecodable failure) <- fromException e = StopUndecodable failure | otherwise = StopWorkerCrashed e -- The consumer-group member index for this config, or 0 for a non-group@@ -487,6 +551,9 @@ configMember :: SubscriptionConfig -> Int32 configMember config = maybe 0 member (consumerGroup config) +configSize :: SubscriptionConfig -> Int32+configSize config = maybe 1 size (consumerGroup config)+ -- Resolve the exact checkpoint key through the shared initializer. A database -- error is emitted and rethrown so startup fails loudly. A semantic -- 'FailIfMissing' result remains typed so the caller can emit the distinct@@ -502,18 +569,31 @@ mHook <- readIORef loadCheckpointHookRef injected <- maybe (pure Nothing) (\hook -> hook config) mHook result <- case injected of- Just hooked -> pure hooked+ Just hooked -> pure (fmap (either (Left . SomeSubscriptionStartupFailure) (Right . (,False))) hooked) Nothing -> Pool.use pool $- CheckpointSQL.initializeSubscriptionCheckpointSession+ CheckpointSQL.initializeWorkerCheckpointSession subName mem+ (configSize config)+ (target config)+ (targetBindingPolicy config) (missingCheckpointPolicy config) case result of Left err -> do emit (KirokuEventSubscriptionDbError subName LoadCheckpoint err (groupCtxOf config)) throwIO err- Right resolution -> pure resolution+ Right (Left refusal@(SomeSubscriptionStartupFailure concrete)) ->+ case fromException (toException concrete) of+ Just missing -> pure (Left missing)+ Nothing -> do+ case fromException (toException concrete) of+ Just mismatch -> emit (KirokuEventSubscriptionGroupSizeMismatch mismatch (groupCtxOf config))+ Nothing -> pure ()+ throwIO refusal+ Right (Right (resolution, adopted)) -> do+ when adopted $ emit (KirokuEventSubscriptionTargetBound subName (target config) (groupCtxOf config))+ pure (Right resolution) -- How a DB-driven live loop ('liveLoopCategoryNotify' / 'liveLoopDbDriven') -- exited. The driver maps these onto FSM inputs: a clean handler stop becomes@@ -584,11 +664,11 @@ case fetchResult of Left err -> pure (Left err) Right events -> do- emit (KirokuEventSubscriptionFetched (name config) (V.length events) (groupCtxOf config))- if V.null events+ emit (KirokuEventSubscriptionFetched (name config) (decodedBatchLength events) (groupCtxOf config))+ if decodedBatchLength events == 0 then pure (Right (Just c)) else do- result <- processEvents pool config stateVar events emit posRef+ result <- processEvents pool config stateVar events emit posRef stSettings case result of Nothing -> pure (Right Nothing) -- handler said Stop Just newPos -> drainTo newPos@@ -636,11 +716,11 @@ case fetchResult of Left err -> pure (Left err) Right events -> do- emit (KirokuEventSubscriptionFetched (name config) (V.length events) (groupCtxOf config))- if V.null events+ emit (KirokuEventSubscriptionFetched (name config) (decodedBatchLength events) (groupCtxOf config))+ if decodedBatchLength events == 0 then pure (Right (Just c)) else do- result <- processEvents pool config stateVar events emit posRef+ result <- processEvents pool config stateVar events emit posRef stSettings case result of Nothing -> pure (Right Nothing) -- handler said Stop Just newPos -> drainTo newPos@@ -659,7 +739,7 @@ GlobalPosition -> (KirokuEvent -> IO ()) -> StoreSettings ->- IO (Either Pool.UsageError (Vector RecordedEvent))+ IO (Either Pool.UsageError DecodedBatch) fetchBatch pool config cursor@(GlobalPosition pos) emit stSettings = do mHook <- readIORef fetchBatchHookRef injected <- maybe (pure Nothing) (\hook -> hook config cursor) mHook@@ -668,16 +748,18 @@ Nothing -> case (consumerGroup config, target config) of (Nothing, AllStreams) -> do- result <- Pool.use pool (Session.statement (pos, batchSize config) SQL.readAllForwardStmt)+ result <- Pool.use pool (Session.statement (pos, batchSizeValue (batchSize config)) SQL.readAllForwardStmt) handle result (Nothing, Category (CategoryName cat)) -> do- result <- Pool.use pool (Session.statement (pos, cat, batchSize config) SQL.readCategoryForwardStmt)+ result <- Pool.use pool (Session.statement (pos, cat, batchSizeValue (batchSize config)) SQL.readCategoryForwardStmt) handle result- (Just (ConsumerGroup m n), AllStreams) -> do- result <- Pool.use pool (Session.statement (pos, m, n, batchSize config) SQL.readAllForwardConsumerGroupStmt)+ (Just cg, AllStreams) -> do+ let m = member cg; n = size cg+ result <- Pool.use pool (Session.statement (pos, m, n, batchSizeValue (batchSize config)) SQL.readAllForwardConsumerGroupStmt) handle result- (Just (ConsumerGroup m n), Category (CategoryName cat)) -> do- result <- Pool.use pool (Session.statement (pos, cat, m, n, batchSize config) SQL.readCategoryForwardConsumerGroupStmt)+ (Just cg, Category (CategoryName cat)) -> do+ let m = member cg; n = size cg+ result <- Pool.use pool (Session.statement (pos, cat, m, n, batchSizeValue (batchSize config)) SQL.readCategoryForwardConsumerGroupStmt) handle result where handle = \case@@ -711,74 +793,95 @@ Pool -> SubscriptionConfig -> TVar SubscriptionState ->- Vector RecordedEvent ->+ DecodedBatch -> (KirokuEvent -> IO ()) -> IORef GlobalPosition ->+ StoreSettings -> IO (Maybe GlobalPosition)-processEvents pool config stateVar events emit posRef = do- -- The state the driver wrote for this batch (CatchingUp / Live); restored- -- after each retry so the observable state does not stick on 'Retrying'.+processEvents pool config stateVar batch emit posRef stSettings = do driving <- atomically (readTVar stateVar)- -- Emit one centralized per-batch delivery event for *every* target and both- -- phases. This is the single delivery primitive, so this one emit uniformly- -- covers catch-up for every target, AllStreams live, and the DB-driven live- -- loops (which still also emit KirokuEventSubscriptionFetched per fetch). let phase = case driving of CatchingUp{} -> DeliveredCatchUp _ -> DeliveredLive- emit (KirokuEventSubscriptionDelivered subName (V.length events) phase groupCtx)- go driving 0+ emit (KirokuEventSubscriptionDelivered subName (decodedBatchLength batch) phase groupCtx)+ -- Choose the vector representation once per batch. The unchanged arm walks+ -- RecordedEvent directly, with no per-event Decoded/Undecodable allocation.+ case batch of+ UnchangedBatch events -> walk events id (\event pos -> deliver driving event pos 1)+ TransformedBatch events -> walk events decodedEventRecorded (\event pos -> dispatch driving event pos 1) where subName = name config groupCtx = groupCtxOf config maxAttempts = retryMaxAttempts (retryPolicy config) - go driving i- | i >= V.length events = do- let lastEvent = V.last events- newPos = globalPosition lastEvent- writeIORef posRef newPos- saveCheckpoint pool config newPos emit- pure (Just newPos)- | otherwise = do- let event = events V.! i- evtPos = globalPosition event- writeIORef posRef evtPos- if shouldDeliver (eventTypeFilter config) (selector config) event- then deliver driving i event evtPos 1- else -- Filtered out by the type filter or the selector: skip the- -- handler entirely (so a non-matching event never reaches the- -- bridge and is never retried or dead-lettered), but keep walking- -- the batch so the batch-tail checkpoint advances the cursor past- -- it. The subscription never stalls on a long run of filtered-out- -- events.- go driving (i + 1)+ walk :: Vector a -> (a -> RecordedEvent) -> (a -> GlobalPosition -> IO Bool) -> IO (Maybe GlobalPosition)+ walk events rawOf consume = go 0+ where+ go i+ | i >= V.length events = do+ let newPos = globalPosition (rawOf (V.last events))+ writeIORef posRef newPos+ saveCheckpoint pool config newPos emit+ pure (Just newPos)+ | otherwise = do+ let item = events V.! i+ event = rawOf item+ evtPos = globalPosition event+ keepGoing <-+ if shouldDeliver (eventTypeFilter config) (selector config) event+ then consume item evtPos+ else pure True+ if keepGoing+ then writeIORef posRef evtPos >> go (i + 1)+ else pure Nothing - -- Deliver one event; @attempt@ is the 1-based delivery attempt (1 = first).- deliver driving i event evtPos attempt = do+ deliver driving event evtPos attempt = do+ writeIORef posRef evtPos result <- handler config event- case result of- Continue -> go driving (i + 1)- Stop -> do- -- Save checkpoint up to the event we just processed- saveCheckpoint pool config evtPos emit- pure Nothing- DeadLetter reason -> do- writeDeadLetter pool config evtPos event reason attempt emit- go driving (i + 1)- Retry delay- -- Exhausted the retry budget: dead-letter and advance past it.- | attempt >= maxAttempts -> do- writeDeadLetter pool config evtPos event (DeadLetterMaxAttempts attempt) attempt emit- go driving (i + 1)- -- Redeliver the same event after the requested delay.+ resolve driving event evtPos attempt result (deliver driving event evtPos (attempt + 1))++ dispatch driving outcome evtPos attempt = case outcome of+ Decoded event -> deliver driving event evtPos attempt+ Undecodable raw failure -> case undecodableHandler config of+ Nothing+ | attempt >= maxAttempts -> throwIO (SubscriptionUndecodable failure) | otherwise -> do- atomically (writeTVar stateVar (Retrying evtPos attempt))- emit (KirokuEventSubscriptionRetrying subName evtPos attempt groupCtx)- threadDelay (retryDelayMicros delay)- atomically (writeTVar stateVar driving)- deliver driving i event evtPos (attempt + 1)+ pause driving evtPos attempt (RetryDelay 1)+ retryDecode driving raw evtPos (attempt + 1)+ Just callback -> do+ result <- callback raw failure+ resolve driving raw evtPos attempt result (retryDecode driving raw evtPos (attempt + 1)) + retryDecode driving raw evtPos attempt = do+ outcome <- decodeEvent stSettings raw+ case outcome of+ Decoded event | not (shouldDeliver (eventTypeFilter config) (selector config) event) -> pure True+ _ -> dispatch driving outcome evtPos attempt++ -- Ordinary and explicitly chosen undecodable dispositions share exactly+ -- one checkpoint/dead-letter/retry resolver. The absent callback never+ -- enters its exhausted-Retry dead-letter branch.+ resolve driving event evtPos attempt result retry = case result of+ Continue -> pure True+ Stop -> do+ writeIORef posRef evtPos+ saveCheckpoint pool config evtPos emit+ pure False+ DeadLetter reason -> do+ writeDeadLetter pool config evtPos event reason attempt emit+ pure True+ Retry delay+ | attempt >= maxAttempts -> do+ writeDeadLetter pool config evtPos event (DeadLetterMaxAttempts attempt) attempt emit+ pure True+ | otherwise -> pause driving evtPos attempt delay >> retry++ pause driving evtPos attempt delay = do+ atomically (writeTVar stateVar (Retrying evtPos attempt))+ emit (KirokuEventSubscriptionRetrying subName evtPos attempt groupCtx)+ threadDelay (retryDelayMicros delay)+ atomically (writeTVar stateVar driving)+ -- Atomically record an event in @kiroku.dead_letters@ and advance the -- subscription's checkpoint past it (one statement; the checkpoint does not -- advance if the insert fails). On a database error the worker surfaces a@@ -797,11 +900,15 @@ writeDeadLetter pool config gp@(GlobalPosition pos) event reason attempt emit = do let subName@(SubscriptionName name') = name config mem = configMember config+ (kind, category) = CheckpointSQL.targetColumns (Just (target config)) EventId uuid = eventId event params = SQL.DeadLetterParams { SQL.dlSubscriptionName = name' , SQL.dlMember = mem+ , SQL.dlTargetKind = kind+ , SQL.dlTargetCategory = category+ , SQL.dlGroupSize = configSize config , SQL.dlGlobalPosition = pos , SQL.dlEventId = uuid , SQL.dlReason = deadLetterReasonJson reason@@ -831,7 +938,7 @@ mem = configMember config mHook <- readIORef saveCheckpointHookRef mapM_ (\hook -> hook config position) mHook- result <- Pool.use pool (Session.statement (name', mem, pos) SQL.saveCheckpointMemberStmt)+ result <- Pool.use pool (CheckpointSQL.saveBoundCheckpointSession (target config) name' mem pos (configSize config)) case result of Left err -> emit (KirokuEventSubscriptionDbError subName SaveCheckpoint err (groupCtxOf config)) Right () -> pure ()
test/Main.hs view
@@ -31,9 +31,11 @@ import Test.Concurrency qualified as Concurrency import Test.ConsumerGroup qualified as ConsumerGroup import Test.ConsumerGroupEffect qualified as ConsumerGroupEffect+import Test.ConsumerGroupResize qualified as ConsumerGroupResize import Test.ConsumerGroupSql qualified as ConsumerGroupSql import Test.EventTypeFilter qualified as EventTypeFilter import Test.FailureInjection qualified as FailureInjection+import Test.HandlerStall qualified as HandlerStall import Test.Helpers import Test.HistoryRetention qualified as HistoryRetention import Test.HistoryRetentionMock qualified as HistoryRetentionMock@@ -61,13 +63,17 @@ import Test.SubscriptionRegistry qualified as SubscriptionRegistry import Test.SubscriptionRetryDeadLetter qualified as SubscriptionRetryDeadLetter import Test.SubscriptionState qualified as SubscriptionState+import Test.SubscriptionTarget qualified as SubscriptionTarget import Test.Transaction qualified as Transaction import Test.TruncateBefore qualified as TruncateBefore+import Test.UniqueViolationMapping qualified as UniqueViolationMapping import Test.VisibleGlobalHeadPosition qualified as VisibleGlobalHeadPosition import Test.VisibleGlobalHeadPositionMock qualified as VisibleGlobalHeadPositionMock main :: IO () main = withSharedMigratedPostgres $ hspec $ do+ UniqueViolationMapping.spec+ SubscriptionTarget.spec Category.spec Properties.spec Concurrency.spec@@ -83,6 +89,7 @@ InterpreterHooks.spec Causation.spec ConsumerGroupSql.spec+ ConsumerGroupResize.spec ConsumerGroup.spec ConsumerGroupEffect.spec describe "performance structure" $ do@@ -90,6 +97,7 @@ NotifyGuard.spec StreamNameLookup.noOpSpec CategoryIdleNoSpin.spec+ HandlerStall.spec PublisherCallbackResilience.spec PublisherIdleAdvance.spec PublisherRestartNoRebroadcast.spec@@ -261,6 +269,22 @@ (r3 ^. #globalPosition) `shouldBe` GlobalPosition 3 describe "duplicate event ID" $ do+ it "returns the caller event id on same-stream retry and leaves one stored event" $ \store -> do+ let eid = EventId (UUID.fromWords 0x01234567 0x89ab7def 0x80123456 0x7890abcd)+ event = makeEvent "Created" (Aeson.object []) & #eventId .~ Just eid+ stream = StreamName "same-stream-duplicate"+ Right _ <- runStoreIO store $ appendToStream stream NoStream [event]+ forM_ [AnyVersion, StreamExists, ExactVersion (StreamVersion 1)] $ \expected -> do+ result <- runStoreIO store $ appendToStream stream expected [event]+ result `shouldBe` Left (DuplicateEvent (Just eid))+ countEvents store `shouldReturn` 1+ Right stored <- runStoreIO store $ readStreamForward stream (StreamVersion 0) 100+ Right global <- runStoreIO store $ readAllForward (GlobalPosition 0) 100+ map (\e -> e ^. #eventId) (V.toList stored) `shouldBe` [eid]+ map (\e -> e ^. #eventId) (V.toList global) `shouldBe` [eid]+ Right info <- runStoreIO store $ getStream stream+ fmap (\s -> s ^. #version) info `shouldBe` Just (StreamVersion 1)+ it "rejects duplicate event IDs" $ \store -> do let eid = EventId (case UUID.fromString "01234567-89ab-7def-8012-34567890abcd" of Just u -> u; Nothing -> error "bad uuid") let event1 =@@ -1084,12 +1108,15 @@ { name = SubscriptionName "catchup-test" , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1121,12 +1148,15 @@ { name = SubscriptionName "live-test" , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1167,12 +1197,15 @@ { name = SubscriptionName "ckpt-test" , target = AllStreams , handler = handler1- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1197,12 +1230,15 @@ { name = SubscriptionName "ckpt-test" , target = AllStreams , handler = handler2- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1243,12 +1279,15 @@ { name = SubscriptionName "transition-no-duplicates" , target = AllStreams , handler = handler'- , batchSize = 2+ , batchSize = either (error . show) Prelude.id (mkBatchSize 2) , queueCapacity = 32 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1297,12 +1336,15 @@ { name = SubscriptionName "cancel-replay-test" , target = AllStreams , handler = handler1- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1321,12 +1363,15 @@ { name = SubscriptionName "cancel-replay-test" , target = AllStreams , handler = handler2- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1358,12 +1403,15 @@ { name = SubscriptionName "stop-boundary-test" , target = AllStreams , handler = handler1- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1386,12 +1434,15 @@ { name = SubscriptionName "stop-boundary-test" , target = AllStreams , handler = handler2- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1440,12 +1491,15 @@ { name = SubscriptionName "f18-live-test" , target = Category (CategoryName "order") , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1495,12 +1549,15 @@ { name = SubscriptionName "cat-sub-test" , target = Category (CategoryName "order") , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1542,12 +1599,15 @@ { name = SubscriptionName "invoice-category-ordering" , target = Category (CategoryName "invoice") , handler = handler'- , batchSize = 2+ , batchSize = either (error . show) Prelude.id (mkBatchSize 2) , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1574,12 +1634,15 @@ { name = SubscriptionName "cancel-test" , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1605,12 +1668,15 @@ { name = SubscriptionName "empty-store-test" , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1646,7 +1712,7 @@ { name = SubscriptionName "debounce-test" , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , -- Sized to absorb 50 publisher batches without -- backpressure overflow. The default 16 was -- intermittently overrun under load (the@@ -1661,6 +1727,9 @@ , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1711,12 +1780,15 @@ { name = SubscriptionName "f6-overflow-test" , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 1 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1772,12 +1844,15 @@ { name = SubscriptionName "f6-overflow-test" , target = AllStreams , handler = handler2- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1816,12 +1891,15 @@ { name = SubscriptionName "eff-catchup-test" , target = AllStreams , handler = effHandler- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1851,12 +1929,15 @@ { name = SubscriptionName "withsub-normal" , target = AllStreams , handler = \_ -> pure Continue- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -1879,12 +1960,15 @@ { name = SubscriptionName "withsub-throw" , target = AllStreams , handler = \_ -> pure Continue- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -2061,12 +2145,15 @@ { name = SubscriptionName "lifecycle-test" , target = AllStreams , handler = h- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing
test/Test/CatchupDbErrorNoPrematureSwitch.hs view
@@ -44,12 +44,15 @@ { name = subName , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing
test/Test/CategoryIdleNoSpin.hs view
@@ -47,7 +47,7 @@ import Data.Generics.Labels () import Data.Text qualified as T import Kiroku.Store-import Test.Helpers (caughtUpEventHandler, makeEvent, waitForPublisher, waitForSubscriptionLive, withTestStoreSettings)+import Test.Helpers (caughtUpEventHandler, makeEvent, validConsumerGroup, waitForPublisher, waitForSubscriptionLive, withTestStoreSettings) import Test.Hspec -- | Poll an @IO Bool@ predicate until it holds or the microsecond budget runs out.@@ -142,7 +142,7 @@ pure Continue memberConfig subName member = (defaultSubscriptionConfig subName (Category (CategoryName "quiet")) deliver)- { consumerGroup = Just ConsumerGroup{member = member, size = 2}+ { consumerGroup = Just (validConsumerGroup member 2) } withTestStoreSettings (\s -> s & #eventHandler .~ Just obsHandler) $ \store -> bracket (subscribe store (memberConfig (names !! 0) 0)) cancel $ \_ ->@@ -192,7 +192,7 @@ -- append advances the global position it gates on. cfg = (defaultSubscriptionConfig subName AllStreams deliver)- { consumerGroup = Just ConsumerGroup{member = 0, size = 3}+ { consumerGroup = Just (validConsumerGroup 0 3) } bracket (subscribe store cfg) cancel $ \_handle -> do waitForSubscriptionLive liveBarrier
test/Test/ConsumerGroup.hs view
@@ -33,9 +33,10 @@ import Kiroku.Store import Kiroku.Store.SQL qualified as SQL import Kiroku.Store.Subscription.Stream (subscriptionStream)+import Kiroku.Store.Subscription.Stream qualified as Buffer import Kiroku.Test.Postgres (withMigratedTestDatabase) import Streamly.Data.Stream qualified as Stream-import Test.Helpers (makeEvent, waitForPublisher, waitWithTimeout, withTestStore, withTestStoreSettings)+import Test.Helpers (makeEvent, validConsumerGroup, waitForPublisher, waitWithTimeout, withTestStore, withTestStoreSettings) import Test.Hspec -- | Extract @(originalStreamId, globalPosition)@ as raw 'Int64's from an event.@@ -83,7 +84,7 @@ Text -> Text -> Int32 -> Int32 -> EventHandler -> SubscriptionConfig memberConfig nm cat m n h = (defaultSubscriptionConfig (SubscriptionName nm) (Category (CategoryName cat)) h)- { consumerGroup = Just (ConsumerGroup{member = m, size = n})+ { consumerGroup = Just (validConsumerGroup m n) } -- | A size-@n@ @$all@-group config for member @m@ with the given handler.@@ -91,7 +92,7 @@ Text -> Int32 -> Int32 -> EventHandler -> SubscriptionConfig memberConfigAll nm m n h = (defaultSubscriptionConfig (SubscriptionName nm) AllStreams h)- { consumerGroup = Just (ConsumerGroup{member = m, size = n})+ { consumerGroup = Just (validConsumerGroup m n) } {- | Run a subscription built from the given config-completer, collecting the@@ -252,7 +253,7 @@ -- must ignore it. runStmtP store $ Session.statement- ("rz-sub" :: Text, 0 :: Int32, 10_000_000 :: Int64)+ ("rz-sub" :: Text, 0 :: Int32, 10_000_000 :: Int64, 4 :: Int32, "category", Just "rz") SQL.saveCheckpointMemberStmt -- Run 2: member 2 restarts and must resume from its OWN checkpoint@@ -292,9 +293,10 @@ let cfg = (defaultSubscriptionConfig (SubscriptionName "guard-sub") (Category (CategoryName "guardcat")) (\_ -> pure Continue))- { consumerGroup = Just (ConsumerGroup{member = 3, size = 4})+ { consumerGroup = Just (validConsumerGroup 3 4) , consumerGroupGuard = True , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -329,7 +331,7 @@ k `shouldSatisfy` (\x -> x > 0 && x < total) (stream, cancelStream) <-- subscriptionStream store (memberConfig "bridge-sub" "bridge" 0 2 (\_ -> pure Continue)) 64+ subscriptionStream store (memberConfig "bridge-sub" "bridge" 0 2 (\_ -> pure Continue)) (either (error . show) Prelude.id (Buffer.mkStreamBufferSize 64)) pulled <- Stream.toList (Stream.take k stream) cancelStream sort (map (snd . pairOf) pulled) `shouldBe` slicePos
test/Test/ConsumerGroupEffect.hs view
@@ -52,7 +52,7 @@ import Kiroku.Store import Kiroku.Store.SQL qualified as SQL import Kiroku.Store.Subscription.Effect qualified as SubEff-import Test.Helpers (makeEvent, waitForPublisher, waitWithTimeout, withTestStore)+import Test.Helpers (makeEvent, validConsumerGroup, waitForPublisher, waitWithTimeout, withTestStore) import Test.Hspec -- | Extract @(originalStreamId, globalPosition)@ as raw 'Int64's from an event.@@ -109,7 +109,7 @@ pure (if c >= k then Stop else Continue) cfg = (defaultSubscriptionConfig (SubscriptionName nm) (Category (CategoryName cat)) effHandler)- { consumerGroup = Just (ConsumerGroup{member = m, size = n})+ { consumerGroup = Just (validConsumerGroup m n) } runEff $ SubEff.runSubscription store $ do handle <- SubEff.subscribe cfg@@ -193,7 +193,7 @@ pure (if c >= 5 then Stop else Continue) ) )- { consumerGroup = Just (ConsumerGroup{member = 0, size = 1})+ { consumerGroup = Just (validConsumerGroup 0 1) } handle <- SubEff.subscribe cfg liftIO $ do
+ test/Test/ConsumerGroupResize.hs view
@@ -0,0 +1,257 @@+{-# LANGUAGE MultilineStrings #-}+{-# LANGUAGE NumericUnderscores #-}++module Test.ConsumerGroupResize (spec) where++import Control.Concurrent (threadDelay)+import Control.Concurrent.Async qualified as Async+import Control.Concurrent.STM+import Control.Exception (SomeException, fromException)+import Control.Lens ((^.))+import Control.Monad (forM_, void)+import Data.Aeson qualified as Aeson+import Data.Generics.Labels ()+import Data.IORef+import Data.Int (Int32, Int64)+import Data.Set qualified as Set+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Vector qualified as Vector+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 (preparable)+import Hasql.Transaction qualified as Tx+import Hasql.Transaction.Sessions qualified as TxSessions+import Kiroku.Store+import Kiroku.Store.SQL qualified as SQL+import System.Timeout (timeout)+import Test.Helpers (makeEvent, validConsumerGroup, waitWithTimeout, withTestStore, withTestStoreSettings)+import Test.Hspec++spec :: Spec+spec = describe "consumer-group resize" $ do+ it "rejects invalid sizes and member indices at construction" $ do+ mkConsumerGroupSize 0 `shouldBe` Left (InvalidConsumerGroup 0 0)+ mkConsumerGroupSize (-1) `shouldBe` Left (InvalidConsumerGroup 0 (-1))+ let size2 = groupSize 2+ mkConsumerGroup (-1) size2 `shouldBe` Left (InvalidConsumerGroup (-1) 2)+ mkConsumerGroup 2 size2 `shouldBe` Left (InvalidConsumerGroup 2 2)+ mkConsumerGroup 1 size2 `shouldBe` Right (validConsumerGroup 1 2)++ it "refuses a configured size that disagrees with stored topology before creating a member" $ do+ refusals <- newIORef []+ let observe = \case+ KirokuEventSubscriptionGroupSizeMismatch refusal group -> modifyIORef' refusals ((refusal, group) :)+ _ -> pure ()+ withTestStoreSettings (\settings -> settings{eventHandler = Just observe}) $ \store -> do+ seedCheckpoint store "mismatch" 0 2 5+ called <- newIORef False+ handle <- subscribe store $ config "mismatch" 2 3 $ \_ -> writeIORef called True >> pure Continue+ expectMismatch handle "mismatch" 3 [2]+ readIORef called `shouldReturn` False+ rows store "mismatch" `shouldReturn` [(0, 2, 5)]+ readIORef refusals `shouldReturn` [(ConsumerGroupSizeMismatch (SubscriptionName "mismatch") 3 (Vector.singleton 2), GroupMember 2 3)]++ it "also refuses an ordinary subscription over an existing larger group" $+ withTestStore $ \store -> do+ seedCheckpoint store "ordinary-mismatch" 0 2 0+ handle <- subscribe store (defaultSubscriptionConfig (SubscriptionName "ordinary-mismatch") AllStreams (\_ -> pure Continue))+ expectMismatch handle "ordinary-mismatch" 1 [2]++ it "refuses an underestimated derived topology until resized" $+ withTestStore $ \store -> do+ -- The upgrade test proves a lone legacy member derives size 1.+ seedCheckpoint store "underestimated" 0 1 8+ refused <- subscribe store (config "underestimated" 0 2 (\_ -> pure Continue))+ expectMismatch refused "underestimated" 2 [1]+ void $ resize store "underestimated" 2+ accepted <- subscribe store (config "underestimated" 0 2 (\_ -> pure Continue))+ awaitLive accepted+ cancel accepted+ rows store "underestimated" `shouldReturn` [(0, 2, 8), (1, 2, 8)]++ it "refuses mixed stored sizes even when its own row agrees" $+ withTestStore $ \store -> do+ seedCheckpoint store "mixed" 0 2 0+ seedCheckpoint store "mixed" 1 3 0+ handle <- subscribe store (config "mixed" 0 2 (\_ -> pure Continue))+ expectMismatch handle "mixed" 2 [2, 3]++ it "serializes concurrent initializers with competing group sizes" $+ withTestStore $ \store -> do+ [firstHandle, secondHandle] <- Async.mapConcurrently (\n -> subscribe store (config "concurrent" 0 n (\_ -> pure Continue))) [2, 3]+ -- The winner reaches Live; the loser exits with a typed refusal.+ result <- timeout 5_000_000 $ Async.race (wait firstHandle) (wait secondHandle)+ case result of+ Just (Left outcome) -> assertMismatch outcome+ Just (Right outcome) -> assertMismatch outcome+ Nothing -> expectationFailure "competing topology startup timed out"+ mapM_ cancel [firstHandle, secondHandle]+ stored <- rows store "concurrent"+ length stored `shouldBe` 1+ stored `shouldSatisfy` \case [(0, n, 0)] -> n == 2 || n == 3; _ -> False++ it "persists configured topology on initialization, ordinary saves and dead-letter saves" $+ withTestStore $ \store -> do+ void $ runStoreIO store $ appendToStream (StreamName "persist-1") NoStream [makeEvent "E" (Aeson.object [])]+ -- A size-1 group exercises all rows deterministically. Size >1 is+ -- checked below via explicit ordinary and dead-letter SQL paths.+ handle <- subscribe store (config "saved" 0 1 (\_ -> pure Stop))+ expectClean handle+ rows store "saved" `shouldReturn` [(0, 1, 1)]+ live <- subscribe store ((config "initialized" 2 3 (\_ -> pure Continue)){missingCheckpointPolicy = FromCurrentHead})+ awaitLive live+ cancel live+ rows store "initialized" `shouldReturn` [(2, 3, 1)]+ seedCheckpoint store "saved-size" 1 4 20+ seedCheckpoint store "saved-size" 1 4 10+ rows store "saved-size" `shouldReturn` [(1, 4, 20)]+ Right events <- runStoreIO store (readAllForward (GlobalPosition 0) 10)+ let event = Vector.head events+ EventId eid = event ^. #eventId+ params = SQL.DeadLetterParams "dead-letter-size" 1 "unbound" Nothing 4 1 eid (Aeson.object []) "test" 1+ runSession store (Session.statement params SQL.insertDeadLetterAndCheckpointStmt)+ rows store "dead-letter-size" `shouldReturn` [(1, 4, 1)]+ let StreamId sid = event ^. #originalStreamId+ ownerStmt =+ preparable+ "SELECT (((hashtextextended($1::bigint::text, 0) % 4) + 4) % 4)::int4"+ (E.param (E.nonNullable E.int8))+ (D.singleRow (D.column (D.nonNullable D.int4)))+ owner <- runSession store (Session.statement sid ownerStmt)+ saved <- subscribe store (config "worker-save-size" owner 4 (\_ -> pure Stop))+ expectClean saved+ rows store "worker-save-size" `shouldReturn` [(owner, 4, 1)]+ deadLettered <- subscribe store (config "worker-dead-letter-size" owner 4 (\_ -> pure (DeadLetter (DeadLetterPoison "explicit consumer decision"))))+ awaitLive deadLettered+ cancel deadLettered+ rows store "worker-dead-letter-size" `shouldReturn` [(owner, 4, 1)]++ it "delivers every seeded event after equalizing size 2 to size 3" $+ withTestStore $ \store -> do+ forM_ [1 .. 40 :: Int] $ \i -> do+ result <- runStoreIO store $ appendToStream (StreamName ("resize-" <> Text.pack (show i))) NoStream [makeEvent "E" (Aeson.object [])]+ result `shouldSatisfy` either (const False) (const True)+ Right events <- runStoreIO store (readAllForward (GlobalPosition 0) 100)+ let expected = Set.fromList [event ^. #eventId | event <- Vector.toList events]+ -- Member 0 has seen nothing, member 1 has crossed the full log.+ -- Remember its actual delivered set; a naive size change loses+ -- streams that move from old member 0 to new member 1.+ oldFast <- runSession store (Session.statement (0, 1, 2, 100) SQL.readAllForwardConsumerGroupStmt)+ naive <- mapM (\(m, p) -> runSession store (Session.statement (p, m, 3, 100) SQL.readAllForwardConsumerGroupStmt)) [(0, 0), (1, 40), (2, 0)]+ let alreadySeen = Set.fromList [e ^. #eventId | e <- Vector.toList oldFast]+ naiveSeen = Set.fromList [e ^. #eventId | batch <- naive, e <- Vector.toList batch]+ Set.union alreadySeen naiveSeen `shouldNotBe` expected+ seedCheckpoint store "skewed" 0 2 0+ seedCheckpoint store "skewed" 1 2 40+ refused <- subscribe store (config "skewed" 1 3 (\_ -> pure Continue))+ expectMismatch refused "skewed" 3 [2]+ report <- resize store "skewed" 3+ report `shouldBe` ConsumerGroupResizeReport (Vector.singleton 2) 2 (groupSize 3) (GlobalPosition 0)+ delivered <- newTVarIO alreadySeen+ handles <- mapM (\m -> subscribe store (config "skewed" m 3 (\e -> atomically (modifyTVar' delivered (Set.insert (e ^. #eventId))) >> pure Continue))) [0, 1, 2]+ complete <- timeout 15_000_000 (atomically (readTVar delivered >>= check . (== expected)))+ mapM_ cancel handles+ complete `shouldBe` Just ()+ readTVarIO delivered `shouldReturn` expected++ it "is idempotent when repeated at the same size and retains surviving member identities" $+ withTestStore $ \store -> do+ seedCheckpoint store "repeat" 0 2 5+ seedCheckpoint store "repeat" 1 2 20+ oldIds <- memberIds store "repeat"+ void $ resize store "repeat" 3+ originalRows <- rows store "repeat"+ resizedIds <- memberIds store "repeat"+ take 2 resizedIds `shouldBe` oldIds+ second <- resize store "repeat" 3+ rows store "repeat" `shouldReturn` originalRows+ memberIds store "repeat" `shouldReturn` resizedIds+ originalRows `shouldBe` [(0, 3, 5), (1, 3, 5), (2, 3, 5)]+ resumePosition second `shouldBe` GlobalPosition 5+ previousMemberCount second `shouldBe` 3++ it "starts a missing group at zero and removes obsolete members on shrink" $+ withTestStore $ \store -> do+ report <- resize store "missing" 3+ previousMemberCount report `shouldBe` 0+ resumePosition report `shouldBe` GlobalPosition 0+ rows store "missing" `shouldReturn` [(0, 3, 0), (1, 3, 0), (2, 3, 0)]+ void $ resize store "missing" 1+ rows store "missing" `shouldReturn` [(0, 1, 0)]++ it "rolls back the complete resize with the caller transaction" $+ withTestStore $ \store -> do+ seedCheckpoint store "rollback" 0 2 5+ seedCheckpoint store "rollback" 1 2 20+ void $ runSession store $ TxSessions.transaction TxSessions.ReadCommitted TxSessions.Write $ do+ report <- resizeConsumerGroupTx (SubscriptionName "rollback") (groupSize 3)+ Tx.condemn+ pure report+ rows store "rollback" `shouldReturn` [(0, 2, 5), (1, 2, 20)]++groupSize :: Int32 -> ConsumerGroupSize+groupSize = either (error . show) Prelude.id . mkConsumerGroupSize++config :: Text -> Int32 -> Int32 -> EventHandler -> SubscriptionConfig+config name m n handler =+ (defaultSubscriptionConfig (SubscriptionName name) AllStreams handler)+ { consumerGroup = Just (validConsumerGroup m n)+ }++seedCheckpoint :: KirokuStore -> Text -> Int32 -> Int32 -> Int64 -> IO ()+seedCheckpoint store name m n p = runSession store (Session.statement (name, m, p, n, "unbound", Nothing) SQL.saveCheckpointMemberStmt)++resize :: KirokuStore -> Text -> Int32 -> IO ConsumerGroupResizeReport+resize store name n = runSession store $ TxSessions.transaction TxSessions.ReadCommitted TxSessions.Write (resizeConsumerGroupTx (SubscriptionName name) (groupSize n))++rows :: KirokuStore -> Text -> IO [(Int32, Int32, Int64)]+rows store name = Vector.toList <$> runSession store (Session.statement name stmt)+ where+ stmt =+ preparable+ "SELECT consumer_group_member, consumer_group_size, last_seen FROM subscriptions WHERE subscription_name = $1 ORDER BY consumer_group_member"+ (E.param (E.nonNullable E.text))+ (D.rowVector ((,,) <$> D.column (D.nonNullable D.int4) <*> D.column (D.nonNullable D.int4) <*> D.column (D.nonNullable D.int8)))++runSession :: KirokuStore -> Session.Session a -> IO a+runSession store session = Pool.use (store ^. #pool) session >>= either (error . show) pure++expectMismatch :: SubscriptionHandle -> Text -> Int32 -> [Int32] -> IO ()+expectMismatch handle name n observed = do+ result <- waitWithTimeout 5_000_000 handle+ case result of+ Right (Left exception) -> fromException exception `shouldBe` Just (ConsumerGroupSizeMismatch (SubscriptionName name) n (Vector.fromList observed))+ other -> expectationFailure ("expected topology mismatch, got " <> show other)++assertMismatch :: Either SomeException () -> IO ()+assertMismatch = \case+ Left exception -> (fromException exception :: Maybe ConsumerGroupSizeMismatch) `shouldSatisfy` maybe False (const True)+ Right () -> expectationFailure "expected competing-size startup refusal"++expectClean :: SubscriptionHandle -> IO ()+expectClean handle =+ waitWithTimeout 5_000_000 handle >>= \case+ Right (Right ()) -> pure ()+ other -> expectationFailure (show other)++awaitLive :: SubscriptionHandle -> IO ()+awaitLive handle = do+ result <- timeout 5_000_000 loop+ result `shouldBe` Just ()+ where+ loop =+ currentState handle >>= \case+ Just state | stateName state == "live" -> pure ()+ _ -> threadDelay 1_000 >> loop++memberIds :: KirokuStore -> Text -> IO [Int64]+memberIds store name =+ runSession store $+ Session.statement name $+ preparable+ "SELECT subscription_id FROM subscriptions WHERE subscription_name = $1 ORDER BY consumer_group_member"+ (E.param (E.nonNullable E.text))+ (D.rowList (D.column (D.nonNullable D.int8)))
test/Test/ConsumerGroupSql.hs view
@@ -183,16 +183,16 @@ let subName = "proj-acct" :: Text it "stores and reads independent checkpoints per member" $ \store -> do- runStmt store $ Session.statement (subName, 0 :: Int32, 7 :: Int64) SQL.saveCheckpointMemberStmt- runStmt store $ Session.statement (subName, 1 :: Int32, 13 :: Int64) SQL.saveCheckpointMemberStmt+ runStmt store $ Session.statement (subName, 0 :: Int32, 7 :: Int64, 2 :: Int32, "unbound", Nothing) SQL.saveCheckpointMemberStmt+ runStmt store $ Session.statement (subName, 1 :: Int32, 13 :: Int64, 2 :: Int32, "unbound", Nothing) SQL.saveCheckpointMemberStmt m0 <- runStmt store $ Session.statement (subName, 0 :: Int32) SQL.getCheckpointMemberStmt m1 <- runStmt store $ Session.statement (subName, 1 :: Int32) SQL.getCheckpointMemberStmt m0 `shouldBe` Just 7 m1 `shouldBe` Just 13 it "never moves a member checkpoint backward (GREATEST monotonicity)" $ \store -> do- runStmt store $ Session.statement (subName, 0 :: Int32, 20 :: Int64) SQL.saveCheckpointMemberStmt- runStmt store $ Session.statement (subName, 0 :: Int32, 5 :: Int64) SQL.saveCheckpointMemberStmt+ runStmt store $ Session.statement (subName, 0 :: Int32, 20 :: Int64, 2 :: Int32, "unbound", Nothing) SQL.saveCheckpointMemberStmt+ runStmt store $ Session.statement (subName, 0 :: Int32, 5 :: Int64, 2 :: Int32, "unbound", Nothing) SQL.saveCheckpointMemberStmt m0 <- runStmt store $ Session.statement (subName, 0 :: Int32) SQL.getCheckpointMemberStmt m0 `shouldBe` Just 20
test/Test/EventTypeFilter.hs view
@@ -181,7 +181,7 @@ cfg = (defaultSubscriptionConfig subName AllStreams handler') { eventTypeFilter = OnlyEventTypes (Set.fromList [EventType "A"])- , batchSize = 2000+ , batchSize = either (error . show) Prelude.id (mkBatchSize 2000) } handle <- subscribe store cfg reached <- waitForCheckpoint store subT 1002@@ -254,7 +254,7 @@ cfg = (defaultSubscriptionConfig subName AllStreams handler') { selector = Just keepOnly- , batchSize = 2000+ , batchSize = either (error . show) Prelude.id (mkBatchSize 2000) } handle <- subscribe store cfg reached <- waitForCheckpoint store subT 1002
test/Test/FailureInjection.hs view
@@ -60,12 +60,15 @@ { name = SubscriptionName "f14-down-window" , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing
+ test/Test/HandlerStall.hs view
@@ -0,0 +1,254 @@+module Test.HandlerStall (spec) where++import Control.Concurrent (threadDelay)+import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)+import Control.Concurrent.STM qualified as STM+import Control.Exception (fromException)+import Control.Lens ((&), (.~), (^.))+import Control.Monad (forM_)+import Data.Aeson qualified as Aeson+import Data.Maybe (isNothing)+import Data.Time (NominalDiffTime)+import Data.Vector qualified as V+import GHC.Clock (getMonotonicTimeNSec)+import Kiroku.Store+import System.Timeout (timeout)+import Test.Helpers (caughtUpEventHandler, makeEvent, validConsumerGroup, waitForPublisher, waitForSubscriptionLive, withTestStoreSettings)+import Test.Hspec++within :: IO a -> IO a+within action = timeout 5_000_000 action >>= maybe (fail "handler stall test timed out") pure++appendEvents :: KirokuStore -> IO ()+appendEvents store = do+ Right _ <-+ runStoreIO store $+ appendToStream+ (StreamName "stall-events")+ NoStream+ [makeEvent "First" (Aeson.object []), makeEvent "Second" (Aeson.object [])]+ pure ()++checkpoint :: KirokuStore -> IO GlobalPosition+checkpoint store = do+ Right (SubscriptionCheckpointInventory _ rows) <- runStoreIO store subscriptionCheckpointInventory+ case [savedPosition | SubscriptionCheckpoint name _ savedPosition _ <- V.toList rows, name == SubscriptionName "stall-test"] of+ [savedPosition] -> pure savedPosition+ other -> fail ("unexpected checkpoints: " <> show other)++spec :: Spec+spec = describe "handler stall" $ do+ forM_+ [ (AllStreams, Nothing, False)+ , (AllStreams, Nothing, True)+ , (Category (CategoryName "stall"), Nothing, True)+ , (AllStreams, Just (validConsumerGroup 0 1), True)+ , (Category (CategoryName "stall"), Just (validConsumerGroup 0 1), False)+ ]+ $ \(target', group, live) ->+ it ("warns without acknowledging or checkpointing a blocked handler for " <> show (target', group, live)) $ do+ release <- newEmptyMVar+ caughtUp <- newEmptyMVar+ warning <- STM.newEmptyTMVarIO+ count <- STM.newTVarIO (0 :: Int)+ let subName = SubscriptionName "stall-test"+ observe event = do+ caughtUpEventHandler subName caughtUp Nothing event+ case event of+ KirokuEventSubscriptionHandlerStalled name pos eid elapsed groupContext | name == subName -> STM.atomically $ do+ STM.modifyTVar' count (+ 1)+ _ <- STM.tryPutTMVar warning (pos, eid, elapsed, groupContext)+ pure ()+ _ -> pure ()+ tweak settings = settings & #eventHandler .~ Just observe+ handler event =+ if event ^. #globalPosition == GlobalPosition 1+ then takeMVar release >> pure Continue+ else pure Stop+ withTestStoreSettings tweak $ \store -> do+ if live then pure () else appendEvents store >> waitForPublisher store (GlobalPosition 2)+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig subName target' handler)+ { consumerGroup = group+ , handlerStallWarnAfter = Just 0.05+ }+ )+ if live then waitForSubscriptionLive caughtUp >> appendEvents store else pure ()+ (pos, eid, elapsed, groupContext) <- within (STM.atomically (STM.readTMVar warning))+ pos `shouldBe` GlobalPosition 1+ Right events <- runStoreIO store (readAllForward (GlobalPosition 0) 1)+ eid `shouldBe` (V.head events ^. #eventId)+ elapsed `shouldSatisfy` (>= 0.05)+ groupContext `shouldBe` maybe NonGroup (\_ -> GroupMember 0 1) group+ checkpoint store `shouldReturn` GlobalPosition 0+ putMVar release ()+ within (wait handle) >>= (`shouldSatisfy` either (const False) (const True))+ checkpoint store `shouldReturn` GlobalPosition 2+ stoppedCount <- STM.readTVarIO count+ threadDelay 120_000+ STM.readTVarIO count `shouldReturn` stoppedCount+ currentState handle >>= (`shouldSatisfy` isNothing)++ it "warns periodically while pending and joins the watchdog on cancellation" $ do+ blocked <- newEmptyMVar+ warnings <- STM.newTQueueIO+ let observe event = case event of+ KirokuEventSubscriptionHandlerStalled _ pos _ elapsed _ -> do+ now <- getMonotonicTimeNSec+ STM.atomically (STM.writeTQueue warnings (pos, elapsed, now))+ _ -> pure ()+ withTestStoreSettings (\s -> s & #eventHandler .~ Just observe) $ \store -> do+ appendEvents store+ waitForPublisher store (GlobalPosition 2)+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig (SubscriptionName "stall-test") AllStreams (\_ -> takeMVar blocked >> pure Stop))+ { handlerStallWarnAfter = Just 0.05+ }+ )+ (_, firstElapsed, firstAt) <- within (STM.atomically (STM.readTQueue warnings))+ (_, secondElapsed, secondAt) <- within (STM.atomically (STM.readTQueue warnings))+ secondElapsed `shouldSatisfy` (> firstElapsed)+ secondAt - firstAt `shouldSatisfy` (>= 45_000_000)+ checkpoint store `shouldReturn` GlobalPosition 0+ within (cancel handle)+ _ <- STM.atomically (STM.flushTQueue warnings)+ threadDelay 120_000+ STM.atomically (STM.isEmptyTQueue warnings) `shouldReturn` True+ currentState handle >>= (`shouldSatisfy` isNothing)++ it "warns for the current invocation after replacing a quick handler" $ do+ release <- newEmptyMVar+ warning <- STM.newEmptyTMVarIO+ let observe event = case event of+ KirokuEventSubscriptionHandlerStalled _ pos _ elapsed _ -> STM.atomically $ do+ _ <- STM.tryPutTMVar warning (pos, elapsed)+ pure ()+ _ -> pure ()+ handler event =+ if event ^. #globalPosition == GlobalPosition 1+ then pure Continue+ else takeMVar release >> pure Stop+ withTestStoreSettings (\s -> s & #eventHandler .~ Just observe) $ \store -> do+ appendEvents store+ waitForPublisher store (GlobalPosition 2)+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig (SubscriptionName "stall-test") AllStreams handler)+ { handlerStallWarnAfter = Just 0.05+ }+ )+ (pos, elapsed) <- within (STM.atomically (STM.readTMVar warning))+ pos `shouldBe` GlobalPosition 2+ elapsed `shouldSatisfy` (>= 0.05)+ -- Normal progress is persisted at the batch boundary.+ checkpoint store `shouldReturn` GlobalPosition 0+ putMVar release ()+ _ <- within (wait handle)+ pure ()++ it "contains a throwing diagnostic callback and still completes the handler" $ do+ release <- newEmptyMVar+ observed <- STM.newEmptyTMVarIO+ let observe KirokuEventSubscriptionHandlerStalled{} = do+ STM.atomically $ do+ _ <- STM.tryPutTMVar observed ()+ pure ()+ ioError (userError "diagnostic callback failed")+ observe _ = pure ()+ withTestStoreSettings (\s -> s & #eventHandler .~ Just observe) $ \store -> do+ appendEvents store+ waitForPublisher store (GlobalPosition 2)+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig (SubscriptionName "stall-test") AllStreams (\_ -> takeMVar release >> pure Stop))+ { handlerStallWarnAfter = Just 0.05+ }+ )+ within (STM.atomically (STM.readTMVar observed))+ putMVar release ()+ within (wait handle) >>= (`shouldSatisfy` either (const False) (const True))+ checkpoint store `shouldReturn` GlobalPosition 1++ it "clears tracking and joins the watchdog when a handler throws" $ do+ release <- newEmptyMVar+ count <- STM.newTVarIO (0 :: Int)+ let observe KirokuEventSubscriptionHandlerStalled{} = STM.atomically (STM.modifyTVar' count (+ 1))+ observe _ = pure ()+ withTestStoreSettings (\s -> s & #eventHandler .~ Just observe) $ \store -> do+ appendEvents store+ waitForPublisher store (GlobalPosition 2)+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig (SubscriptionName "stall-test") AllStreams (\_ -> takeMVar release >> ioError (userError "handler failed")))+ { handlerStallWarnAfter = Just 0.05+ }+ )+ within (STM.atomically (STM.readTVar count >>= STM.check . (> 0)))+ putMVar release ()+ within (wait handle) >>= (`shouldSatisfy` either (const True) (const False))+ stoppedCount <- STM.readTVarIO count+ threadDelay 120_000+ STM.readTVarIO count `shouldReturn` stoppedCount+ checkpoint store `shouldReturn` GlobalPosition 0++ it "emits nothing when the handler completes before the threshold" $ do+ count <- STM.newTVarIO (0 :: Int)+ let observe KirokuEventSubscriptionHandlerStalled{} = STM.atomically (STM.modifyTVar' count (+ 1))+ observe _ = pure ()+ withTestStoreSettings (\s -> s & #eventHandler .~ Just observe) $ \store -> do+ appendEvents store+ waitForPublisher store (GlobalPosition 2)+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig (SubscriptionName "stall-test") AllStreams (\_ -> pure Stop))+ { handlerStallWarnAfter = Just 0.1+ }+ )+ within (wait handle) >>= (`shouldSatisfy` either (const False) (const True))+ threadDelay 150_000+ STM.readTVarIO count `shouldReturn` 0++ it "keeps warnings disabled by default even when a handler is blocked" $ do+ entered <- newEmptyMVar+ release <- newEmptyMVar+ count <- STM.newTVarIO (0 :: Int)+ let observe KirokuEventSubscriptionHandlerStalled{} = STM.atomically (STM.modifyTVar' count (+ 1))+ observe _ = pure ()+ withTestStoreSettings (\s -> s & #eventHandler .~ Just observe) $ \store -> do+ appendEvents store+ waitForPublisher store (GlobalPosition 2)+ let config = defaultSubscriptionConfig (SubscriptionName "stall-test") AllStreams (\_ -> putMVar entered () >> takeMVar release >> pure Stop)+ handlerStallWarnAfter config `shouldBe` Nothing+ handle <- subscribe store config+ within (takeMVar entered)+ threadDelay 150_000+ STM.readTVarIO count `shouldReturn` 0+ putMVar release ()+ _ <- within (wait handle)+ pure ()++ forM_ [0, -1 :: NominalDiffTime] $ \threshold ->+ it ("refuses a nonpositive warning interval before checkpoint initialization: " <> show threshold) $+ withTestStoreSettings Prelude.id $ \store -> do+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig (SubscriptionName "stall-test") AllStreams (\_ -> pure Stop))+ { handlerStallWarnAfter = Just threshold+ }+ )+ within (wait handle) >>= \case+ Left err -> do+ (fromException err :: Maybe InvalidHandlerStallWarnAfter) `shouldBe` Just (InvalidHandlerStallWarnAfter threshold)+ (fromException err :: Maybe SomeSubscriptionStartupFailure) `shouldSatisfy` (not . isNothing)+ Right () -> expectationFailure "invalid interval was accepted"+ Right (SubscriptionCheckpointInventory _ rows) <- runStoreIO store subscriptionCheckpointInventory+ V.null rows `shouldBe` True
test/Test/Helpers.hs view
@@ -22,6 +22,7 @@ -- * Event construction makeEvent,+ validConsumerGroup, -- * Subscription wait waitWithTimeout,@@ -137,6 +138,9 @@ SQL.DeadLetterParams { SQL.dlSubscriptionName = subscriptionName , SQL.dlMember = 0+ , SQL.dlTargetKind = "unbound"+ , SQL.dlTargetCategory = Nothing+ , SQL.dlGroupSize = 1 , SQL.dlGlobalPosition = globalPosition , SQL.dlEventId = eid , SQL.dlReason = Aeson.object [("source", Aeson.String "test")]@@ -339,3 +343,9 @@ case passthrough of Nothing -> pure () Just f -> f evt++{- | Build known-valid membership for test fixtures; invalid cases test the+smart constructors directly rather than passing a malformed worker config.+-}+validConsumerGroup :: Int32 -> Int32 -> ConsumerGroup+validConsumerGroup m n = either (error . show) Prelude.id (mkConsumerGroupSize n >>= mkConsumerGroup m)
test/Test/InterpreterHooks.hs view
@@ -22,7 +22,7 @@ import Data.IORef (modifyIORef', newIORef, readIORef) import Data.Vector qualified as V import Kiroku.Store-import Test.Helpers (makeEvent, waitForPublisher, waitWithTimeout, withTestStore, withTestStoreSettings)+import Test.Helpers (caughtUpEventHandler, makeEvent, waitForPublisher, waitForSubscriptionLive, waitWithTimeout, withTestStore, withTestStoreSettings) import Test.Hspec spec :: Spec@@ -88,9 +88,15 @@ readHookFiresSpec :: Spec readHookFiresSpec = do+ it "retains the no-hook batch without evaluating its events" $ do+ batch <- decodeEvents defaultStoreSettings (V.replicate 100 (error "no-hook decoding traversed the vector"))+ case batch of+ UnchangedBatch events -> V.length events `shouldBe` 100+ TransformedBatch _ -> expectationFailure "default decoding allocated per-event outcomes"+ it "applies decodeHook to readAllForward results" $ do let marker = Aeson.object [("decoded", Aeson.String "yes")]- inject re = pure $ re & #metadata .~ Just marker+ inject re = pure $ Right (re & #metadata .~ Just marker) tweak cs = cs & #storeSettings@@ -113,13 +119,15 @@ (re ^. #metadata) `shouldBe` Just marker it "applies decodeHook to subscription handlers across catch-up and live phases" $ do+ caughtUp <- newEmptyMVar let marker = Aeson.object [("sub", Aeson.String "tagged")]- inject re = pure $ re & #metadata .~ Just marker+ inject re = pure $ Right (re & #metadata .~ Just marker) subName = SubscriptionName "hook-sub-1" tweak cs = cs & #storeSettings .~ defaultStoreSettings{decodeHook = Just inject}+ & #eventHandler .~ Just (caughtUpEventHandler subName caughtUp Nothing) withTestStoreSettings tweak $ \store -> do -- Pre-append one event so the worker's catch-up path runs -- before live mode kicks in.@@ -141,6 +149,7 @@ else pure Continue config = defaultSubscriptionConfig subName AllStreams handlerFn sub <- subscribe store config+ waitForSubscriptionLive caughtUp -- Append a second event to land in live mode. Right _ <- runStoreIO store $
test/Test/PerformanceStructure.hs view
@@ -2,6 +2,8 @@ module Test.PerformanceStructure (spec) where +import Control.Concurrent (myThreadId)+import Control.Exception (fromException) import Control.Lens ((^.)) import Control.Monad (forM_, unless) import Data.Aeson (Value (..))@@ -9,7 +11,7 @@ import Data.Aeson.KeyMap qualified as KeyMap import Data.ByteString (ByteString) import Data.Generics.Labels ()-import Data.IORef (IORef, modifyIORef', newIORef, readIORef)+import Data.IORef (IORef, modifyIORef', newIORef, readIORef, writeIORef) import Data.Int (Int64) import Data.Text (Text) import Data.Text qualified as T@@ -21,8 +23,10 @@ import Hasql.Statement qualified as Statement import Kiroku.Store import Kiroku.Store.SQL qualified as SQL+import Kiroku.Store.Subscription.Stream qualified as Buffer+import Kiroku.Store.Subscription.Worker (withLoadCheckpointHookForTest) import Kiroku.Test.Fixtures.CategoryScaling (categoryScalingFixtureSql, categoryScalingHead)-import Test.Helpers (makeEvent, withTestStore, withTestStoreSettings)+import Test.Helpers (makeEvent, validConsumerGroup, waitWithTimeout, withTestStore, withTestStoreSettings) import Test.Hspec spec :: Spec@@ -34,6 +38,95 @@ noOpAppendSpec :: Spec noOpAppendSpec = describe "no-op paths use no pooled connection" $ do+ it "rejects invalid batch and buffer sizes before pool checkout" $ do+ checkouts <- newIORef (0 :: Int)+ withObservedStore checkouts $ \_ -> do+ before <- readIORef checkouts+ mkBatchSize 0 `shouldBe` Left (InvalidBatchSize 0)+ mkBatchSize (-1) `shouldBe` Left (InvalidBatchSize (-1))+ Buffer.mkStreamBufferSize 0 `shouldBe` Left (Buffer.InvalidStreamBufferSize 0)+ after <- readIORef checkouts+ after - before `shouldBe` 0++ it "rejects invalid consumer-group values before pool checkout" $ do+ checkouts <- newIORef (0 :: Int)+ withObservedStore checkouts $ \_ -> do+ before <- readIORef checkouts+ mkConsumerGroupSize 0 `shouldBe` Left (InvalidConsumerGroup 0 0)+ let Right n = mkConsumerGroupSize 2+ mkConsumerGroup 2 n `shouldBe` Left (InvalidConsumerGroup 2 2)+ after <- readIORef checkouts+ after - before `shouldBe` 0++ it "refuses mismatched topology in exactly one initialization checkout with no handler" $ do+ workerThread <- newIORef Nothing+ checkouts <- newIORef (0 :: Int)+ delivered <- newIORef False+ let observe (ConnectionObservation _ InUseConnectionStatus) = do+ thread <- myThreadId+ selected <- readIORef workerThread+ if selected == Just thread then modifyIORef' checkouts (+ 1) else pure ()+ observe _ = pure ()+ withTestStoreSettings (\settings -> settings{observationHandler = Just observe}) $ \store -> do+ seeded <- Pool.use (store ^. #pool) (Session.statement ("topology-checkout", 0, 0, 2, "unbound", Nothing) SQL.saveCheckpointMemberStmt)+ seeded `shouldBe` Right ()+ withLoadCheckpointHookForTest+ ( \_ -> do+ thread <- myThreadId+ writeIORef workerThread (Just thread)+ pure Nothing+ )+ $ do+ let cfg =+ (defaultSubscriptionConfig (SubscriptionName "topology-checkout") AllStreams (\_ -> writeIORef delivered True >> pure Continue))+ { consumerGroup = Just (validConsumerGroup 0 3)+ }+ handle <- subscribe store cfg+ waitWithTimeout 5_000_000 handle >>= \case+ Right (Left exception) ->+ (fromException exception :: Maybe ConsumerGroupSizeMismatch) `shouldSatisfy` maybe False (const True)+ other -> expectationFailure ("expected mismatch, got " <> show other)+ readIORef checkouts `shouldReturn` 1+ readIORef delivered `shouldReturn` False++ it "refuses mismatched target in exactly one initialization checkout with no handler" $ do+ workerThread <- newIORef Nothing+ checkouts <- newIORef (0 :: Int)+ delivered <- newIORef False+ let observe (ConnectionObservation _ InUseConnectionStatus) = do+ thread <- myThreadId+ selected <- readIORef workerThread+ if selected == Just thread then modifyIORef' checkouts (+ 1) else pure ()+ observe _ = pure ()+ withTestStoreSettings (\settings -> settings{observationHandler = Just observe}) $ \store -> do+ seeded <- Pool.use (store ^. #pool) (Session.statement ("target-checkout", 0, 0, 1, "all", Nothing) SQL.saveCheckpointMemberStmt)+ seeded `shouldBe` Right ()+ withLoadCheckpointHookForTest+ ( \_ -> do+ thread <- myThreadId+ writeIORef workerThread (Just thread)+ pure Nothing+ )+ $ do+ let cfg =+ (defaultSubscriptionConfig (SubscriptionName "target-checkout") (Category (CategoryName "other")) (\_ -> writeIORef delivered True >> pure Continue))++ handle <- subscribe store cfg+ waitWithTimeout 5_000_000 handle >>= \case+ Right (Left exception) ->+ (fromException exception :: Maybe SubscriptionTargetMismatch) `shouldSatisfy` maybe False (const True)+ other -> expectationFailure ("expected mismatch, got " <> show other)+ readIORef checkouts `shouldReturn` 1+ readIORef delivered `shouldReturn` False++ it "keeps ordinary checkpoint saves as one unconditional monotonic upsert" $ do+ forM_ [Statement.toSql SQL.saveCheckpointMemberStmt, Statement.toSql SQL.saveAllCheckpointMemberStmt, Statement.toSql SQL.saveCategoryCheckpointMemberStmt] $ \statement -> do+ let sql = T.toLower statement+ T.count "insert into" sql `shouldBe` 1+ sql `shouldSatisfy` T.isInfixOf "greatest(subscriptions.last_seen, excluded.last_seen)"+ sql `shouldNotSatisfy` T.isInfixOf "where"+ sql `shouldNotSatisfy` T.isInfixOf "returning"+ it "rejects an empty appendToStream batch before pool checkout" $ do checkouts <- newIORef (0 :: Int) withObservedStore checkouts $ \store -> do
test/Test/PublisherCallbackResilience.hs view
@@ -6,14 +6,30 @@ import Control.Concurrent (threadDelay) import Control.Concurrent.Async qualified as Async import Control.Concurrent.MVar (newEmptyMVar, tryPutMVar)-import Control.Concurrent.STM (atomically, newTVarIO, readTVar, writeTVar)-import Control.Exception (Exception, throwIO)+import Control.Concurrent.STM (atomically, modifyTVar', newTVarIO, readTVar, writeTVar)+import Control.Exception (Exception, fromException, throwIO) import Control.Lens ((&), (.~), (^.))+import Control.Monad (forM_, when)+import Control.Monad.IO.Class (liftIO)+import Data.Aeson ((.=)) import Data.Aeson qualified as Aeson import Data.Generics.Labels ()-import Data.IORef (atomicModifyIORef', modifyIORef', newIORef, readIORef)+import Data.IORef (atomicModifyIORef', modifyIORef', newIORef, readIORef, writeIORef)+import Data.Int (Int64)+import Data.Maybe (isNothing)+import Data.Set qualified as Set+import Data.Text (Text)+import Data.Vector qualified as V+import Effectful (runEff)+import Effectful.State.Static.Local qualified as State+import GHC.Clock (getMonotonicTimeNSec)+import Hasql.Pool qualified as Pool+import Hasql.Session qualified as Session import Kiroku.Store-import Test.Helpers (caughtUpEventHandler, makeEvent, waitForSubscriptionLive, waitWithTimeout, withTestStoreSettings)+import Kiroku.Store.SQL qualified as SQL+import Kiroku.Store.Subscription.Effect qualified as SubEff+import Kiroku.Store.Subscription.Fsm (SubscriptionState (..))+import Test.Helpers (caughtUpEventHandler, makeEvent, validConsumerGroup, waitForPublisher, waitForSubscriptionLive, waitWithTimeout, withTestStoreSettings) import Test.Hspec data CallbackBoom = CallbackBoom@@ -30,8 +46,404 @@ Left () -> fail ("timed out waiting for " <> label) Right a -> pure a +failureOf :: RecordedEvent -> DecodeFailure+failureOf event = DecodeFailure (event ^. #eventId) "cannot decrypt"++poisonHook :: RecordedEvent -> IO (Either DecodeFailure RecordedEvent)+poisonHook event =+ pure $+ if event ^. #eventType == EventType "Boom" then Left (failureOf event) else Right event++appendTypes :: KirokuStore -> Text -> [Text] -> IO ()+appendTypes store stream types = do+ Right _ <- runStoreIO store $ appendToStream (StreamName stream) NoStream (map (\typ -> makeEvent typ (Aeson.object [])) types)+ pure ()++expectCleanStop :: SubscriptionHandle -> IO ()+expectCleanStop handle =+ within "clean subscription stop" (wait handle) >>= \case+ Right () -> pure ()+ Left err -> expectationFailure ("expected clean stop, got " <> show err)++expectDecodeStop :: SubscriptionHandle -> IO DecodeFailure+expectDecodeStop handle =+ within "typed decode stop" (wait handle) >>= \case+ Left err | Just (SubscriptionUndecodable failure) <- fromException err -> pure failure+ other -> fail ("expected SubscriptionUndecodable, got " <> show other)++readCheckpoint :: KirokuStore -> Text -> IO (Maybe Int64)+readCheckpoint store name' = do+ Right value <- Pool.use (store ^. #pool) (Session.statement (name', 0) SQL.getCheckpointMemberStmt)+ pure value++readDeadLetters :: KirokuStore -> Text -> IO [SQL.DeadLetterRecord]+readDeadLetters store name' = do+ Right values <- Pool.use (store ^. #pool) (Session.statement (name', 0) SQL.readDeadLettersStmt)+ pure (V.toList values)++awaitCheckpoint :: KirokuStore -> Text -> Int64 -> IO ()+awaitCheckpoint store name' expected = within "durable checkpoint" loop+ where+ loop = do+ value <- readCheckpoint store name'+ if value == Just expected then pure () else threadDelay 1_000 >> loop+ spec :: Spec spec = describe "publisher callback resilience" $ do+ it "stops only the default subscriber while a sibling dead-letters and continues live" $ do+ defaultLive <- newEmptyMVar+ siblingLive <- newEmptyMVar+ siblingDone <- newEmptyMVar+ attempts <- newIORef []+ defaultSeen <- newIORef ([] :: [EventType])+ siblingSeen <- newIORef ([] :: [EventType])+ publisherFailures <- newTVarIO (0 :: Int)+ retryEvents <- newTVarIO (0 :: Int)+ stopped <- newIORef Nothing+ let defaultName = SubscriptionName "decode-default-live"+ siblingName = SubscriptionName "decode-sibling-live"+ hook event+ | event ^. #eventType == EventType "Boom" = do+ now <- getMonotonicTimeNSec+ modifyIORef' attempts (now :)+ pure (Left (failureOf event))+ | otherwise = pure (Right event)+ observe event = do+ caughtUpEventHandler defaultName defaultLive Nothing event+ caughtUpEventHandler siblingName siblingLive Nothing event+ case event of+ KirokuEventPublisherDecodeFailed{} -> atomically $ modifyTVar' publisherFailures (+ 1)+ KirokuEventSubscriptionRetrying name' _ _ _ | name' == defaultName -> atomically $ modifyTVar' retryEvents (+ 1)+ KirokuEventSubscriptionStopped name' _ (StopUndecodable failure) _ | name' == defaultName -> modifyIORef' stopped (const (Just failure))+ _ -> pure ()+ tweak settings =+ settings+ & #storeSettings .~ defaultStoreSettings{decodeHook = Just hook}+ & #eventHandler .~ Just observe+ record ref event = modifyIORef' ref ((event ^. #eventType) :) >> pure Continue+ siblingHandler event = do+ result <- record siblingSeen event+ when (event ^. #eventType == EventType "After") $ () <$ tryPutMVar siblingDone ()+ pure result+ withTestStoreSettings tweak $ \store -> do+ defaultHandle <- subscribe store ((defaultSubscriptionConfig defaultName AllStreams (record defaultSeen)){retryPolicy = RetryPolicy 3})+ sibling <-+ subscribe+ store+ ( (defaultSubscriptionConfig siblingName AllStreams siblingHandler)+ { undecodableHandler = Just (\_ failure -> pure (DeadLetter (DeadLetterDecodeFailure failure)))+ }+ )+ waitForSubscriptionLive defaultLive+ waitForSubscriptionLive siblingLive+ appendTypes store "decode-live" ["Before", "Boom", "After"]+ failure <- expectDecodeStop defaultHandle+ decodeFailureReason failure `shouldBe` "cannot decrypt"+ readIORef stopped `shouldReturn` Just failure+ atomically (readTVar retryEvents) `shouldReturn` 2+ within "healthy sibling delivery" (waitForSubscriptionLive siblingDone)+ awaitCheckpoint store "decode-sibling-live" 3+ waitForPublisher store (GlobalPosition 3)+ atomically (readTVar publisherFailures) `shouldReturn` 1+ reverse <$> readIORef defaultSeen `shouldReturn` [EventType "Before"]+ reverse <$> readIORef siblingSeen `shouldReturn` [EventType "Before", EventType "After"]+ readCheckpoint store "decode-default-live" >>= (`shouldSatisfy` maybe False (< 2))+ (map SQL.deadLetterGlobalPosition <$> readDeadLetters store "decode-default-live") `shouldReturn` []+ letters <- readDeadLetters store "decode-sibling-live"+ map SQL.deadLetterGlobalPosition letters `shouldBe` [2]+ map SQL.deadLetterReasonSummary letters `shouldBe` ["decode failure: cannot decrypt"]+ map SQL.deadLetterReason letters `shouldBe` [Aeson.object ["kind" .= ("decode_failure" :: Text), "event_id" .= (case decodeFailureEventId failure of EventId eid -> show eid), "detail" .= ("cannot decrypt" :: Text)]]+ times <- reverse <$> readIORef attempts+ length times `shouldBe` 3+ zipWith (-) (drop 1 times) times `shouldSatisfy` all (>= 900_000_000)+ currentState sibling >>= \case+ Just Live{} -> pure ()+ other -> expectationFailure ("expected healthy sibling, got " <> show other)+ currentState defaultHandle >>= (`shouldSatisfy` isNothing)+ cancel sibling++ it "shares successful live decoding across subscribers" $ do+ firstLive <- newEmptyMVar+ secondLive <- newEmptyMVar+ calls <- newIORef (0 :: Int)+ let firstName = SubscriptionName "decode-fanout-first"+ secondName = SubscriptionName "decode-fanout-second"+ hook event = modifyIORef' calls (+ 1) >> pure (Right event)+ observe event = do+ caughtUpEventHandler firstName firstLive Nothing event+ caughtUpEventHandler secondName secondLive Nothing event+ tweak settings =+ settings+ & #storeSettings .~ defaultStoreSettings{decodeHook = Just hook}+ & #eventHandler .~ Just observe+ withTestStoreSettings tweak $ \store -> do+ first <- subscribe store (defaultSubscriptionConfig firstName AllStreams (\_ -> pure Stop))+ second <- subscribe store (defaultSubscriptionConfig secondName AllStreams (\_ -> pure Stop))+ waitForSubscriptionLive firstLive+ waitForSubscriptionLive secondLive+ appendTypes store "decode-fanout" ["Good"]+ expectCleanStop first+ expectCleanStop second+ readIORef calls `shouldReturn` 1++ forM_ [(Category (CategoryName "decode"), Nothing), (AllStreams, Just (validConsumerGroup 0 1))] $ \(target', group) ->+ it ("handles an undecodable event after confirmed DB-driven live transition for " <> show (target', group)) $ do+ caughtUp <- newEmptyMVar+ seen <- newIORef ([] :: [EventType])+ let subName = SubscriptionName "decode-db-live"+ tweak settings =+ settings+ & #storeSettings .~ defaultStoreSettings{decodeHook = Just poisonHook}+ & #eventHandler .~ Just (caughtUpEventHandler subName caughtUp Nothing)+ handler event = modifyIORef' seen ((event ^. #eventType) :) >> pure Stop+ withTestStoreSettings tweak $ \store -> do+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig subName target' handler)+ { consumerGroup = group+ , undecodableHandler = Just (\_ failure -> pure (DeadLetter (DeadLetterDecodeFailure failure)))+ }+ )+ waitForSubscriptionLive caughtUp+ appendTypes store "decode-db-live" ["Boom", "After"]+ expectCleanStop handle+ readIORef seen `shouldReturn` [EventType "After"]+ readCheckpoint store "decode-db-live" `shouldReturn` Just 2+ map SQL.deadLetterGlobalPosition <$> readDeadLetters store "decode-db-live" `shouldReturn` [1]++ it "does not apply an undecodable disposition to an explicitly filtered-out event" $ do+ let tweak settings = settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just poisonHook}+ withTestStoreSettings tweak $ \store -> do+ appendTypes store "decode-filter" ["Boom", "After"]+ waitForPublisher store (GlobalPosition 2)+ handle <-+ subscribe+ store+ ( ( defaultSubscriptionConfig+ (SubscriptionName "decode-filter")+ AllStreams+ ( \event -> do+ event ^. #eventType `shouldBe` EventType "After"+ pure Stop+ )+ )+ { eventTypeFilter = OnlyEventTypes (Set.singleton (EventType "After"))+ , retryPolicy = RetryPolicy 1+ }+ )+ expectCleanStop handle+ readCheckpoint store "decode-filter" `shouldReturn` Just 2+ map SQL.deadLetterGlobalPosition <$> readDeadLetters store "decode-filter" `shouldReturn` []++ it "recovers a typed one-shot failure on the default retry without a callback" $ do+ failedOnce <- newIORef False+ caughtUp <- newEmptyMVar+ seen <- newIORef ([] :: [EventType])+ let subName = SubscriptionName "decode-one-shot"+ hook event+ | event ^. #eventType == EventType "Boom" = do+ already <- atomicModifyIORef' failedOnce (\old -> (True, old))+ pure (if already then Right event else Left (failureOf event))+ | otherwise = pure (Right event)+ tweak settings =+ settings+ & #storeSettings .~ defaultStoreSettings{decodeHook = Just hook}+ & #eventHandler .~ Just (caughtUpEventHandler subName caughtUp Nothing)+ handler event = do+ modifyIORef' seen ((event ^. #eventType) :)+ pure (if event ^. #eventType == EventType "After" then Stop else Continue)+ withTestStoreSettings tweak $ \store -> do+ handle <- subscribe store (defaultSubscriptionConfig subName AllStreams handler)+ waitForSubscriptionLive caughtUp+ appendTypes store "decode-recovery" ["Boom", "After"]+ expectCleanStop handle+ reverse <$> readIORef seen `shouldReturn` [EventType "Boom", EventType "After"]+ readCheckpoint store "decode-one-shot" `shouldReturn` Just 2+ (map SQL.deadLetterGlobalPosition <$> readDeadLetters store "decode-one-shot") `shouldReturn` []++ forM_ [AllStreams, Category (CategoryName "decode")] $ \target' ->+ forM_ [Nothing, Just (validConsumerGroup 0 1)] $ \group ->+ it ("stops a persistent undecodable event during catch-up for " <> show (target', group)) $ do+ let tweak settings = settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just poisonHook}+ handler _ = expectationFailure "undecodable event reached ordinary handler" >> pure Stop+ withTestStoreSettings tweak $ \store -> do+ appendTypes store "decode-catchup" ["Boom"]+ waitForPublisher store (GlobalPosition 1)+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig (SubscriptionName "decode-catchup") target' handler)+ { consumerGroup = group+ , retryPolicy = RetryPolicy 1+ }+ )+ _ <- expectDecodeStop handle+ readCheckpoint store "decode-catchup" `shouldReturn` Just 0+ (map SQL.deadLetterGlobalPosition <$> readDeadLetters store "decode-catchup") `shouldReturn` []++ forM_ [Continue, Stop, Retry (RetryDelay 0)] $ \disposition ->+ it ("honors an explicit undecodable disposition " <> show disposition) $ do+ callbacks <- newIORef (0 :: Int)+ seen <- newIORef ([] :: [EventType])+ let tweak settings = settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just poisonHook}+ handler event = modifyIORef' seen ((event ^. #eventType) :) >> pure Stop+ callback _ _ = modifyIORef' callbacks (+ 1) >> pure disposition+ withTestStoreSettings tweak $ \store -> do+ appendTypes store "decode-disposition" ["Boom", "After"]+ waitForPublisher store (GlobalPosition 2)+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig (SubscriptionName "decode-disposition") AllStreams handler)+ { undecodableHandler = Just callback+ , retryPolicy = RetryPolicy 3+ }+ )+ expectCleanStop handle+ case disposition of+ Stop -> do+ readIORef seen `shouldReturn` []+ readCheckpoint store "decode-disposition" `shouldReturn` Just 1+ (map SQL.deadLetterGlobalPosition <$> readDeadLetters store "decode-disposition") `shouldReturn` []+ Retry _ -> do+ readIORef callbacks `shouldReturn` 3+ letters <- readDeadLetters store "decode-disposition"+ map SQL.deadLetterReasonSummary letters `shouldBe` ["max retry attempts exceeded (3)"]+ readCheckpoint store "decode-disposition" `shouldReturn` Just 2+ _ -> do+ readIORef callbacks `shouldReturn` 1+ readIORef seen `shouldReturn` [EventType "After"]+ (map SQL.deadLetterGlobalPosition <$> readDeadLetters store "decode-disposition") `shouldReturn` []++ it "re-applies the hook on a callback Retry and then calls the ordinary handler" $ do+ failedOnce <- newIORef False+ callbacks <- newIORef (0 :: Int)+ seen <- newIORef ([] :: [EventType])+ let hook event = do+ already <- atomicModifyIORef' failedOnce (\old -> (True, old))+ pure (if already then Right event else Left (failureOf event))+ tweak settings = settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just hook}+ handler event = modifyIORef' seen ((event ^. #eventType) :) >> pure Stop+ withTestStoreSettings tweak $ \store -> do+ appendTypes store "decode-callback-retry" ["Boom"]+ waitForPublisher store (GlobalPosition 1)+ handle <-+ subscribe+ store+ ( (defaultSubscriptionConfig (SubscriptionName "decode-callback-retry") AllStreams handler)+ { undecodableHandler = Just (\_ _ -> modifyIORef' callbacks (+ 1) >> pure (Retry (RetryDelay 0)))+ }+ )+ expectCleanStop handle+ readIORef callbacks `shouldReturn` 1+ readIORef seen `shouldReturn` [EventType "Boom"]++ it "fails reads with EventDecodeFailed instead of returning a partial vector" $ do+ let tweak settings = settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just poisonHook}+ withTestStoreSettings tweak $ \store -> do+ appendTypes store "decode-read" ["Before", "Boom", "After"]+ forM_+ [ readAllForward (GlobalPosition 0) 10+ , readAllBackward (GlobalPosition 10) 10+ , readStreamForward (StreamName "decode-read") (StreamVersion 0) 10+ , readStreamBackward (StreamName "decode-read") (StreamVersion 10) 10+ , readCategory (CategoryName "decode") (GlobalPosition 0) 10+ ]+ $ \readEvents ->+ runStoreIO store readEvents >>= \case+ Left (EventDecodeFailed failure) -> decodeFailureReason failure `shouldBe` "cannot decrypt"+ other -> expectationFailure ("expected typed read failure, got " <> show other)++ it "cancels promptly while waiting for the default decode retry" $ do+ retrying <- newEmptyMVar+ let subName = SubscriptionName "decode-cancel"+ observe event = case event of+ KirokuEventSubscriptionRetrying{} -> () <$ tryPutMVar retrying ()+ _ -> pure ()+ tweak settings =+ settings+ & #storeSettings .~ defaultStoreSettings{decodeHook = Just poisonHook}+ & #eventHandler .~ Just observe+ withTestStoreSettings tweak $ \store -> do+ appendTypes store "decode-cancel" ["Boom"]+ waitForPublisher store (GlobalPosition 1)+ handle <- subscribe store (defaultSubscriptionConfig subName AllStreams (\_ -> pure Continue))+ waitForSubscriptionLive retrying+ startedAt <- getMonotonicTimeNSec+ cancel handle+ finishedAt <- getMonotonicTimeNSec+ finishedAt - startedAt `shouldSatisfy` (< 500_000_000)+ currentState handle >>= (`shouldSatisfy` isNothing)+ readCheckpoint store "decode-cancel" `shouldReturn` Just 0++ it "replays the failed event from its durable checkpoint after the hook is fixed" $ do+ repaired <- newIORef False+ seen <- newIORef ([] :: [EventType])+ let subName = SubscriptionName "decode-restart"+ hook event = do+ healthy <- readIORef repaired+ pure (if healthy then Right event else Left (failureOf event))+ tweak settings = settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just hook}+ handler event = modifyIORef' seen ((event ^. #eventType) :) >> pure Stop+ withTestStoreSettings tweak $ \store -> do+ appendTypes store "decode-restart" ["Boom"]+ waitForPublisher store (GlobalPosition 1)+ failed <- subscribe store ((defaultSubscriptionConfig subName AllStreams handler){retryPolicy = RetryPolicy 1})+ _ <- expectDecodeStop failed+ readCheckpoint store "decode-restart" `shouldReturn` Just 0+ writeIORef repaired True+ recovered <- subscribe store (defaultSubscriptionConfig subName AllStreams handler)+ expectCleanStop recovered+ readIORef seen `shouldReturn` [EventType "Boom"]+ readCheckpoint store "decode-restart" `shouldReturn` Just 1++ it "releases a bracketed subscriber during a pending default decode retry" $ do+ retrying <- newEmptyMVar+ let observe event = case event of+ KirokuEventSubscriptionRetrying{} -> () <$ tryPutMVar retrying ()+ _ -> pure ()+ tweak settings =+ settings+ & #storeSettings .~ defaultStoreSettings{decodeHook = Just poisonHook}+ & #eventHandler .~ Just observe+ within "bracketed store/subscription shutdown" $ withTestStoreSettings tweak $ \store -> do+ appendTypes store "decode-shutdown" ["Boom"]+ waitForPublisher store (GlobalPosition 1)+ handle <-+ withSubscription+ store+ ((defaultSubscriptionConfig (SubscriptionName "decode-shutdown") AllStreams (\_ -> pure Continue)){retryPolicy = RetryPolicy 100})+ (\h -> waitForSubscriptionLive retrying >> pure h)+ currentState handle >>= (`shouldSatisfy` isNothing)+ within "cancelled worker completion" (wait handle) >>= \case+ Left err | Just Async.AsyncCancelled <- fromException err -> pure ()+ other -> expectationFailure ("expected bracket cancellation, got " <> show other)++ it "unlifts the undecodable callback in the same persistent effect environment as the handler" $ do+ seen <- newIORef ([] :: [Int])+ let tweak settings = settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just poisonHook}+ withTestStoreSettings tweak $ \store -> do+ appendTypes store "decode-effectful" ["Boom", "After"]+ waitForPublisher store (GlobalPosition 2)+ runEff $ SubEff.runSubscription store $ State.evalState (0 :: Int) $ do+ let config =+ ( defaultSubscriptionConfig+ (SubscriptionName "decode-effectful")+ AllStreams+ ( \_ -> do+ State.modify @Int (+ 1)+ count <- State.get @Int+ liftIO (modifyIORef' seen (count :))+ pure Stop+ )+ )+ { undecodableHandler = Just (\_ _ -> State.modify @Int (+ 1) >> pure Continue)+ }+ handle <- SubEff.subscribe config+ liftIO (expectCleanStop handle)+ readIORef seen `shouldReturn` [2]+ it "keeps the publisher alive when decodeHook throws once" $ do failedOnce <- newIORef False loopErrorSeen <- newEmptyMVar@@ -44,9 +456,9 @@ | event ^. #eventType == EventType "Boom" = do alreadyFailed <- atomicModifyIORef' failedOnce (\old -> (True, old)) if alreadyFailed- then pure event+ then pure (Right event) else throwIO CallbackBoom- | otherwise = pure event+ | otherwise = pure (Right event) observe evt = do case evt of KirokuEventPublisherLoopError{} -> () <$ tryPutMVar loopErrorSeen ()
test/Test/PublisherIdleAdvance.hs view
@@ -40,7 +40,7 @@ { decodeHook = Just $ \event -> do atomicModifyIORef' counter (\n -> (n + 1, ()))- pure event+ pure (Right event) } appendEvents :: KirokuStore -> Int -> String -> IO GlobalPosition
test/Test/PublisherRestartNoRebroadcast.hs view
@@ -49,12 +49,15 @@ { name = subName , target = AllStreams , handler = handler'- , batchSize = 10+ , batchSize = either (error . show) Prelude.id (mkBatchSize 10) , queueCapacity = 1 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -112,12 +115,15 @@ if evt ^. #globalPosition == GlobalPosition 4 then pure Stop else pure Continue- , batchSize = 10+ , batchSize = either (error . show) Prelude.id (mkBatchSize 10) , queueCapacity = 1 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -192,12 +198,15 @@ { name = subName , target = AllStreams , handler = handler'- , batchSize = 10+ , batchSize = either (error . show) Prelude.id (mkBatchSize 10) , queueCapacity = 1 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing
test/Test/StreamBridgeTermination.hs view
@@ -13,6 +13,7 @@ import Data.Map.Strict qualified as Map import Kiroku.Store import Kiroku.Store.Subscription.Stream (AckItem (..), InvalidStreamBufferSize (..), subscriptionAckStream)+import Kiroku.Store.Subscription.Stream qualified as Buffer import Kiroku.Store.Subscription.Worker (withFetchBatchHookForTest) import Streamly.Data.Stream qualified as Stream import Test.Helpers (makeEvent, waitForPublisher, withTestStore)@@ -39,15 +40,8 @@ spec :: Spec spec = describe "stream bridge termination" $ do- it "rejects a zero-sized bridge buffer" $- withTestStore $ \store -> do- let cfg = defaultSubscriptionConfig (SubscriptionName "bridge-zero-buffer-sub") AllStreams (\_ -> pure Continue)- subscriptionAckStream store cfg 0- `shouldThrow` ( \e ->- case e of- InvalidStreamBufferSize 0 -> True- _ -> False- )+ it "rejects a zero-sized bridge buffer at construction" $+ Buffer.mkStreamBufferSize 0 `shouldBe` Left (InvalidStreamBufferSize 0) it "rethrows the worker exception to the consumer when the worker dies" $ withTestStore $ \store -> do@@ -56,7 +50,7 @@ let cfg = defaultSubscriptionConfig (SubscriptionName "bridge-crash-sub") AllStreams (\_ -> pure Continue) injectBoom _ _ = throwIO TestBoom withFetchBatchHookForTest injectBoom $ do- (stream, cancelStream) <- subscriptionAckStream store cfg 16+ (stream, cancelStream) <- subscriptionAckStream store cfg (either (error . show) Prelude.id (Buffer.mkStreamBufferSize 16)) pulled <- within "stream pull to throw TestBoom" (try (Stream.uncons stream)) cancelStream case pulled of@@ -68,7 +62,7 @@ it "ends the stream after a clean worker stop" $ withTestStore $ \store -> do let cfg = defaultSubscriptionConfig (SubscriptionName "bridge-clean-stop-sub") AllStreams (\_ -> pure Continue)- (stream0, cancelStream) <- subscriptionAckStream store cfg 16+ (stream0, cancelStream) <- subscriptionAckStream store cfg (either (error . show) Prelude.id (Buffer.mkStreamBufferSize 16)) finally ( do pos <- appendOne store (StreamName "bridge-clean-stop") (EventType "BridgeCleanStop")@@ -88,7 +82,7 @@ it "cancelAction returns promptly even when the bridge queue is full" $ withTestStore $ \store -> do let cfg = defaultSubscriptionConfig (SubscriptionName "bridge-full-cancel-sub") AllStreams (\_ -> pure Continue)- (stream0, cancelStream) <- subscriptionAckStream store cfg 1+ (stream0, cancelStream) <- subscriptionAckStream store cfg (either (error . show) Prelude.id (Buffer.mkStreamBufferSize 1)) _pos1 <- appendOne store (StreamName "bridge-full-cancel-1") (EventType "BridgeFullCancel1") pos2 <- appendOne store (StreamName "bridge-full-cancel-2") (EventType "BridgeFullCancel2") waitForPublisher store pos2@@ -110,7 +104,7 @@ let subscriptionName = SubscriptionName "bridge-idempotent-cancel-sub" cfg = defaultSubscriptionConfig subscriptionName AllStreams (\_ -> pure Continue) key = (subscriptionName, 0)- (_stream, cancelStream) <- subscriptionAckStream store cfg 1+ (_stream, cancelStream) <- subscriptionAckStream store cfg (either (error . show) Prelude.id (Buffer.mkStreamBufferSize 1)) statesBefore <- subscriptionStates store Map.member key statesBefore `shouldBe` True
test/Test/SubscriptionCheckpointInventory.hs view
@@ -161,7 +161,7 @@ saveCheckpoint :: KirokuStore -> Text -> Int32 -> Int64 -> IO () saveCheckpoint store name member position = do- result <- Pool.use (store ^. #pool) $ Session.statement (name, member, position) SQL.saveCheckpointMemberStmt+ result <- Pool.use (store ^. #pool) $ Session.statement (name, member, position, max 1 (member + 1), "unbound", Nothing) SQL.saveCheckpointMemberStmt case result of Left err -> error ("saveCheckpoint failed: " <> show err) Right () -> pure ()
test/Test/SubscriptionCheckpointReset.hs view
@@ -168,7 +168,7 @@ saveCheckpoint store name member position = do result <- Pool.use (store ^. #pool) $- Session.statement (name, member, position) SQL.saveCheckpointMemberStmt+ Session.statement (name, member, position, max 1 (member + 1), "unbound", Nothing) SQL.saveCheckpointMemberStmt case result of Left err -> expectationFailure ("ordinary checkpoint save failed: " <> show err) Right () -> pure ()
test/Test/SubscriptionCheckpointWorker.hs view
@@ -21,8 +21,9 @@ import Kiroku.Store import Kiroku.Store.Subscription.Effect qualified as SubEff import Kiroku.Store.Subscription.Stream (subscriptionAckStream)+import Kiroku.Store.Subscription.Stream qualified as Buffer import Streamly.Data.Stream qualified as Stream-import Test.Helpers (makeEvent, waitForPublisher, waitWithTimeout, withTestStore, withTestStoreSettings)+import Test.Helpers (makeEvent, validConsumerGroup, waitForPublisher, waitWithTimeout, withTestStore, withTestStoreSettings) import Test.Hspec spec :: Spec@@ -79,7 +80,7 @@ config = (defaultSubscriptionConfig name AllStreams handler) { missingCheckpointPolicy = FromCurrentHead- , batchSize = 3+ , batchSize = either (error . show) Prelude.id (mkBatchSize 3) } subscribeThread <- Async.async (takeMVar gate >> subscribe store config) appendThread <- Async.async $ do@@ -121,7 +122,7 @@ atomically (writeTVar handlerCalled True) pure Continue )- { consumerGroup = Just (ConsumerGroup member 2)+ { consumerGroup = Just (validConsumerGroup member 2) , missingCheckpointPolicy = FromCurrentHead } handles <- mapM (subscribe store . config) [0, 1]@@ -195,7 +196,7 @@ (defaultSubscriptionConfig name AllStreams (\_ -> pure Continue)) { missingCheckpointPolicy = FailIfMissing }- (stream, cancelStream) <- subscriptionAckStream store config 1+ (stream, cancelStream) <- subscriptionAckStream store config (either (error . show) Prelude.id (Buffer.mkStreamBufferSize 1)) pulled <- finally (try (Stream.uncons stream)) cancelStream case pulled of Left err
test/Test/SubscriptionPauseResume.hs view
@@ -81,12 +81,15 @@ { name = SubscriptionName "pause-resume-test" , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 1 , overflowPolicy = PauseAndResume , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -146,12 +149,15 @@ { name = SubscriptionName "dropsub-overflow-test" , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 1 , overflowPolicy = DropSubscription , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing
test/Test/SubscriptionReconnect.hs view
@@ -93,12 +93,15 @@ { name = subName , target = Category (CategoryName "rc") , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = PauseAndResume , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -134,3 +137,56 @@ let isReconnecting (KirokuEventSubscriptionReconnecting n _ _) = n == subName isReconnecting _ = False any isReconnecting evts `shouldBe` True++ it "resumes a database-driven live worker from processed progress" $ do+ delivered <- newIORef []+ caughtUp <- newTVarIO False+ processed <- newTVarIO False+ failed <- newTVarIO False+ recoveryCursor <- newIORef Nothing+ let subName = SubscriptionName "mid-live-progress"+ observe = \case+ KirokuEventSubscriptionCaughtUp n _ _ | n == subName -> atomically (writeTVar caughtUp True)+ _ -> pure ()+ inject config cursor+ | SubTypes.name config /= subName = pure Nothing+ | otherwise = do+ live <- readTVarIO caughtUp+ done <- readTVarIO processed+ alreadyFailed <- readTVarIO failed+ if live && done && not alreadyFailed+ then do+ cursor `shouldBe` GlobalPosition 1+ atomically (writeTVar failed True)+ pure (Just (Left Pool.AcquisitionTimeoutUsageError))+ else do+ if alreadyFailed+ then do+ previous <- readIORef recoveryCursor+ case previous of+ Nothing -> modifyIORef' recoveryCursor (const (Just cursor))+ Just _ -> pure ()+ else pure ()+ pure Nothing+ handler event = do+ modifyIORef' delivered ((event ^. #globalPosition) :)+ if event ^. #globalPosition == GlobalPosition 1+ then atomically (writeTVar processed True) >> pure Continue+ else pure Stop+ cfg = defaultSubscriptionConfig subName (Category (CategoryName "midlive")) handler+ withTestStoreSettings (& #eventHandler .~ Just observe) $ \store ->+ withFetchBatchHookForTest inject $ do+ handle <- subscribe store cfg+ atomically (readTVar caughtUp >>= check)+ Right _ <- runStoreIO store $ appendToStream (StreamName "midlive-1") NoStream [makeEvent "A" (Aeson.object [])]+ atomically (readTVar processed >>= check)+ -- Wake the live fetch once more after its first successful batch.+ Right _ <- runStoreIO store $ appendToStream (StreamName "midlive-2") NoStream [makeEvent "B" (Aeson.object [])]+ result <- waitWithTimeout 10_000_000 handle+ case result of+ Right (Right ()) -> pure ()+ other -> expectationFailure (show other)+ readTVarIO failed `shouldReturn` True+ readIORef recoveryCursor `shouldReturn` Just (GlobalPosition 1)+ reverse <$> readIORef delivered `shouldReturn` [GlobalPosition 1, GlobalPosition 2]+ readCheckpoint store "mid-live-progress" `shouldReturn` Just 2
test/Test/SubscriptionRegistry.hs view
@@ -34,7 +34,7 @@ import Data.Text (Text) import Kiroku.Store import Kiroku.Store.Subscription.EventPublisher qualified as Pub-import Test.Helpers (makeEvent, waitForPublisher, withTestStore)+import Test.Helpers (makeEvent, validConsumerGroup, waitForPublisher, withTestStore) import Test.Hspec -- | A plain @$all@ subscription config whose handler never stops.@@ -45,14 +45,14 @@ groupCont :: Text -> Text -> Int32 -> Int32 -> SubscriptionConfig groupCont nm cat m n = (defaultSubscriptionConfig (SubscriptionName nm) (Category (CategoryName cat)) (\_ -> pure Continue))- { consumerGroup = Just (ConsumerGroup{member = m, size = n})+ { consumerGroup = Just (validConsumerGroup m n) } -- | A size-@n@ @$all@ group config for member @m@ whose handler never stops. groupAllCont :: Text -> Int32 -> Int32 -> SubscriptionConfig groupAllCont nm m n = (defaultSubscriptionConfig (SubscriptionName nm) AllStreams (\_ -> pure Continue))- { consumerGroup = Just (ConsumerGroup{member = m, size = n})+ { consumerGroup = Just (validConsumerGroup m n) } publisherSubscriberCount :: KirokuStore -> IO Int
test/Test/SubscriptionRetryDeadLetter.hs view
@@ -80,7 +80,7 @@ 2 -> pure (DeadLetter (DeadLetterPoison "boom")) 3 -> pure Stop _ -> pure Continue- cfg = (defaultSubscriptionConfig subName AllStreams handler'){batchSize = 100}+ cfg = (defaultSubscriptionConfig subName AllStreams handler'){batchSize = defaultBatchSize} handle <- subscribe store cfg result <- waitWithTimeout 20_000_000 handle case result of@@ -113,7 +113,7 @@ if c <= 2 then pure (Retry (RetryDelay 0)) else pure Continue 3 -> pure Stop _ -> pure Continue- cfg = (defaultSubscriptionConfig subName AllStreams handler'){batchSize = 100}+ cfg = (defaultSubscriptionConfig subName AllStreams handler'){batchSize = defaultBatchSize} handle <- subscribe store cfg result <- waitWithTimeout 20_000_000 handle case result of@@ -145,7 +145,7 @@ _ -> pure Continue cfg = (defaultSubscriptionConfig subName AllStreams handler')- { batchSize = 100+ { batchSize = defaultBatchSize , retryPolicy = RetryPolicy{retryMaxAttempts = 3} } handle <- subscribe store cfg
test/Test/SubscriptionState.hs view
@@ -84,12 +84,15 @@ { name = subName , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = PauseAndResume , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing@@ -127,12 +130,15 @@ { name = subName , target = AllStreams , handler = handler'- , batchSize = 100+ , batchSize = defaultBatchSize , queueCapacity = 16 , overflowPolicy = PauseAndResume , consumerGroup = Nothing , consumerGroupGuard = False , missingCheckpointPolicy = FromBeginning+ , targetBindingPolicy = AdoptUnbound+ , undecodableHandler = Nothing+ , handlerStallWarnAfter = Nothing , retryPolicy = defaultRetryPolicy , eventTypeFilter = AllEventTypes , selector = Nothing
+ test/Test/SubscriptionTarget.hs view
@@ -0,0 +1,195 @@+{-# LANGUAGE MultilineStrings #-}++module Test.SubscriptionTarget (spec) where++import Control.Concurrent (threadDelay)+import Control.Concurrent.Async qualified as Async+import Control.Exception (Exception, SomeException, fromException, throwIO, toException, try)+import Control.Lens ((^.))+import Control.Monad (void)+import Data.Aeson qualified as Aeson+import Data.Generics.Labels ()+import Data.IORef (modifyIORef', newIORef, readIORef, writeIORef)+import Data.Int (Int32, Int64)+import Data.Text (Text)+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 (preparable)+import Hasql.Transaction qualified as Tx+import Hasql.Transaction.Sessions qualified as TxSessions+import Kiroku.Store+import Kiroku.Store.Subscription.Stream qualified as Buffer+import System.Timeout (timeout)+import Test.Helpers (makeEvent, validConsumerGroup, waitWithTimeout, withTestStore, withTestStoreSettings)+import Test.Hspec++spec :: Spec+spec = do+ describe "subscription configuration" $ do+ it "validates batch and bridge buffer boundaries as values" $ do+ mkBatchSize 0 `shouldBe` Left (InvalidBatchSize 0)+ mkBatchSize (-1) `shouldBe` Left (InvalidBatchSize (-1))+ fmap batchSizeValue (mkBatchSize 1) `shouldBe` Right 1+ batchSizeValue defaultBatchSize `shouldBe` 100+ Buffer.mkStreamBufferSize 0 `shouldBe` Left (Buffer.InvalidStreamBufferSize 0)+ fmap Buffer.streamBufferSizeValue (Buffer.mkStreamBufferSize 1) `shouldBe` Right 1+ it "exposes every runtime startup refusal through the parent and concrete catches" $ do+ let name = SubscriptionName "hierarchy"+ key = SubscriptionCheckpointKey name 0+ assertHierarchy (SubscriptionCheckpointMissing key)+ assertHierarchy (ConsumerGroupGuardConflict name 0)+ assertHierarchy (ConsumerGroupSizeMismatch name 2 (V.singleton 1))+ assertHierarchy (SubscriptionTargetMismatch name AllStreams (V.singleton Nothing))+ describe "checkpoint target" $ do+ it "refuses accidental target reuse before delivery, then resumes after explicit rebind" $+ withTestStore $ \store -> do+ Right _ <- runStoreIO store $ appendToStream (StreamName "target-1") NoStream [makeEvent "E" (Aeson.object [])]+ first <- subscribe store (cfg "retarget" AllStreams (\_ -> pure Stop))+ clean first+ called <- newIORef False+ let newTarget = Category (CategoryName "target")+ refused <- subscribe store (cfg "retarget" newTarget (\_ -> writeIORef called True >> pure Stop))+ mismatch refused+ readIORef called `shouldReturn` False+ report <- rebind store "retarget" newTarget 0+ reboundMemberCount report `shouldBe` 1+ previousBindings report `shouldBe` V.singleton (Just AllStreams)+ accepted <- subscribe store (cfg "retarget" newTarget (\_ -> writeIORef called True >> pure Stop))+ clean accepted+ readIORef called `shouldReturn` True+ rows store "retarget" `shouldReturn` [(0, "category", Just "target", 1)]+ it "adopts every legacy member once and emits one bound event for the name" $ do+ events <- newIORef []+ let observe e = modifyIORef' events (e :)+ withTestStoreSettings (\s -> s{eventHandler = Just observe}) $ \store -> do+ void $ resize store "adopt" 2+ startLive store ((cfg "adopt" AllStreams (\_ -> pure Continue)){consumerGroup = Just (validConsumerGroup 0 2)})+ startLive store ((cfg "adopt" AllStreams (\_ -> pure Continue)){consumerGroup = Just (validConsumerGroup 1 2), targetBindingPolicy = RequireBound})+ rows store "adopt" `shouldReturn` [(0, "all", Nothing, 0), (1, "all", Nothing, 0)]+ seen <- readIORef events+ length [() | KirokuEventSubscriptionTargetBound (SubscriptionName "adopt") AllStreams _ <- seen] `shouldBe` 1+ it "refuses unbound rows under RequireBound without changing rows or invoking the handler" $+ withTestStore $ \store -> do+ void $ resize store "strict" 1+ called <- newIORef False+ handle <- subscribe store ((cfg "strict" AllStreams (\_ -> writeIORef called True >> pure Continue)){targetBindingPolicy = RequireBound})+ mismatch handle+ readIORef called `shouldReturn` False+ rows store "strict" `shouldReturn` [(0, "unbound", Nothing, 0)]+ it "refuses mixed bound/unbound siblings before inserting a missing member" $+ withTestStore $ \store -> do+ void $ resize store "mixed-target" 3+ runSession store (Session.script "UPDATE subscriptions SET target_kind = 'all' WHERE subscription_name = 'mixed-target' AND consumer_group_member = 0; DELETE FROM subscriptions WHERE subscription_name = 'mixed-target' AND consumer_group_member = 2")+ handle <- subscribe store ((cfg "mixed-target" AllStreams (\_ -> expectationFailure "handler ran" >> pure Stop)){consumerGroup = Just (validConsumerGroup 2 3)})+ mismatch handle+ rows store "mixed-target" `shouldReturn` [(0, "all", Nothing, 0), (1, "unbound", Nothing, 0)]+ it "does not adopt siblings when FailIfMissing refuses an absent member" $+ withTestStore $ \store -> do+ void $ resize store "missing-member" 2+ runSession store (Session.script "DELETE FROM subscriptions WHERE subscription_name = 'missing-member' AND consumer_group_member = 1")+ handle <- subscribe store ((cfg "missing-member" AllStreams (\_ -> pure Stop)){consumerGroup = Just (validConsumerGroup 1 2), missingCheckpointPolicy = FailIfMissing})+ waitWithTimeout 5_000_000 handle >>= \case+ Right (Left e) -> (fromException e :: Maybe SubscriptionCheckpointMissing) `shouldSatisfy` maybe False (const True)+ other -> expectationFailure (show other)+ rows store "missing-member" `shouldReturn` [(0, "unbound", Nothing, 0)]+ it "serializes conflicting targets on an initially absent name" $+ withTestStore $ \store -> do+ [a, b] <- Async.mapConcurrently (\target -> subscribe store (cfg "competing-targets" target (\_ -> pure Continue))) [AllStreams, Category (CategoryName "other")]+ outcome <- timeout 5_000_000 (Async.race (wait a) (wait b))+ case outcome of+ Just (Left (Left e)) -> assertMismatch e+ Just (Right (Left e)) -> assertMismatch e+ other -> expectationFailure (show other)+ mapM_ cancel [a, b]+ length <$> rows store "competing-targets" `shouldReturn` 1+ it "preserves category binding on resize and initializes fresh names under RequireBound" $+ withTestStore $ \store -> do+ let target = Category (CategoryName "preserved")+ startLive store ((cfg "bound-resize" target (\_ -> pure Continue)){targetBindingPolicy = RequireBound})+ void $ resize store "bound-resize" 3+ rows store "bound-resize" `shouldReturn` [(m, "category", Just "preserved", 0) | m <- [0 .. 2]]+ startLive store ((cfg "bound-resize" target (\_ -> pure Continue)){consumerGroup = Just (validConsumerGroup 2 3), targetBindingPolicy = RequireBound})+ it "rebinds all members atomically, repeats idempotently, and rolls back with the caller" $+ withTestStore $ \store -> do+ void $ resize store "rebind-group" 3+ report <- rebind store "rebind-group" AllStreams 7+ reboundMemberCount report `shouldBe` 3+ previousBindings report `shouldBe` V.singleton Nothing+ rows store "rebind-group" `shouldReturn` [(m, "all", Nothing, 7) | m <- [0 .. 2]]+ again <- rebind store "rebind-group" AllStreams 7+ previousBindings again `shouldBe` V.singleton (Just AllStreams)+ void $ runSession store $ TxSessions.transaction TxSessions.ReadCommitted TxSessions.Write $ do+ result <- rebindSubscriptionTargetTx (SubscriptionName "rebind-group") (Category (CategoryName "rollback")) (GlobalPosition 0)+ Tx.condemn+ pure result+ rows store "rebind-group" `shouldReturn` [(m, "all", Nothing, 7) | m <- [0 .. 2]]+ absent <- Pool.use (store ^. #pool) $ TxSessions.transaction TxSessions.ReadCommitted TxSessions.Write (rebindSubscriptionTargetTx (SubscriptionName "absent") AllStreams (GlobalPosition 0))+ absent `shouldSatisfy` either (const True) (const False)+ it "persists target binding through an atomic dead-letter save" $+ withTestStore $ \store -> do+ Right _ <- runStoreIO store $ appendToStream (StreamName "poison-target-1") NoStream [makeEvent "E" (Aeson.object [])]+ withSubscription store (cfg "target-dead-letter" (Category (CategoryName "poison")) (\_ -> pure (DeadLetter (DeadLetterPoison "chosen by consumer")))) $ \_ -> do+ let awaitDurable = do+ persisted <- rows store "target-dead-letter"+ if persisted == [(0, "category", Just "poison", 1)]+ then pure ()+ else threadDelay 1_000 >> awaitDurable+ timeout 5_000_000 awaitDurable `shouldReturn` Just ()+ rows store "target-dead-letter" `shouldReturn` [(0, "category", Just "poison", 1)]++assertHierarchy :: (Exception e, Eq e) => e -> IO ()+assertHierarchy refusal = do+ parent <- try (throwIO refusal) :: IO (Either SomeSubscriptionStartupFailure ())+ case parent of+ Left caught -> fromException (toException caught) `shouldBe` Just refusal+ Right () -> expectationFailure "parent catch did not match"++cfg :: Text -> SubscriptionTarget -> EventHandler -> SubscriptionConfig+cfg name = defaultSubscriptionConfig (SubscriptionName name)++clean :: SubscriptionHandle -> IO ()+clean handle =+ waitWithTimeout 5_000_000 handle >>= \case+ Right (Right ()) -> pure ()+ other -> expectationFailure (show other)++assertMismatch :: SomeException -> IO ()+assertMismatch e = do+ (fromException e :: Maybe SubscriptionTargetMismatch) `shouldSatisfy` maybe False (const True)+ (fromException e :: Maybe SomeSubscriptionStartupFailure) `shouldSatisfy` maybe False (const True)++mismatch :: SubscriptionHandle -> IO ()+mismatch handle =+ waitWithTimeout 5_000_000 handle >>= \case+ Right (Left e) -> assertMismatch e+ other -> expectationFailure (show other)++startLive :: KirokuStore -> SubscriptionConfig -> IO ()+startLive store config = withSubscription store config $ \handle -> do+ result <- timeout 5_000_000 (awaitLive handle)+ result `shouldBe` Just ()+ where+ awaitLive handle =+ currentState handle >>= \case+ Just state | stateName state == "live" -> pure ()+ _ -> threadDelay 1_000 >> awaitLive handle++runSession :: KirokuStore -> Session.Session a -> IO a+runSession store session = Pool.use (store ^. #pool) session >>= either (fail . show) pure++resize :: KirokuStore -> Text -> Int32 -> IO ConsumerGroupResizeReport+resize store name n = let Right size = mkConsumerGroupSize n in runSession store $ TxSessions.transaction TxSessions.ReadCommitted TxSessions.Write (resizeConsumerGroupTx (SubscriptionName name) size)+rebind :: KirokuStore -> Text -> SubscriptionTarget -> Int64 -> IO SubscriptionTargetRebindReport+rebind store name target position = runSession store $ TxSessions.transaction TxSessions.ReadCommitted TxSessions.Write (rebindSubscriptionTargetTx (SubscriptionName name) target (GlobalPosition position))++rows :: KirokuStore -> Text -> IO [(Int32, Text, Maybe Text, Int64)]+rows store name = V.toList <$> runSession store (Session.statement name stmt)+ where+ stmt =+ preparable+ "SELECT consumer_group_member, target_kind, target_category, last_seen FROM subscriptions WHERE subscription_name = $1 ORDER BY consumer_group_member"+ (E.param (E.nonNullable E.text))+ (D.rowVector ((,,,) <$> D.column (D.nonNullable D.int4) <*> D.column (D.nonNullable D.text) <*> D.column (D.nullable D.text) <*> D.column (D.nonNullable D.int8)))
+ test/Test/UniqueViolationMapping.hs view
@@ -0,0 +1,93 @@+module Test.UniqueViolationMapping (spec) where++import Control.Monad (forM_)+import Data.Text (Text)+import Data.UUID qualified as UUID+import Hasql.Errors qualified as Errors+import Hasql.Pool (UsageError (..))+import Kiroku.Store.Error+import Kiroku.Store.Types+import Test.Hspec++spec :: Spec+spec = describe "unique violation mapping" $ do+ let eid = EventId (UUID.fromWords 0x01234567 0x89ab7def 0x80123456 0x7890abcd)+ scalar = "Key (event_id)=(01234567-89ab-7def-8012-34567890abcd) already exists."+ composite = "Key (event_id, stream_id)=(01234567-89ab-7def-8012-34567890abcd, 42) already exists."+ fallback = WrongExpectedVersion (StreamName "orders-1") AnyVersion (StreamVersion 0)+ mapAppend = mapUsageError "orders-1" AnyVersion+ cases =+ [ ("events_pkey", scalar, DuplicateEvent (Just eid))+ , ("stream_events_pkey", composite, DuplicateEvent (Just eid))+ , ("ix_streams_stream_name", "Key (stream_name)=(orders-1) already exists.", StreamAlreadyExists (StreamName "orders-1"))+ , ("ux_stream_events_stream_version", "Key (stream_id, stream_version)=(42, 1) already exists.", UnexpectedServerError "23505" "detail-only error")+ , ("unknown_constraint", scalar, fallback)+ ]+ forM_ cases $ \(name, detail, expected) -> do+ it ("classifies the exact message constraint " <> show name) $ do+ let message = constraintMessage name+ outcome = case expected of+ UnexpectedServerError code _ -> UnexpectedServerError code message+ other -> other+ mapAppend (serverUsage message (Just detail)) `shouldBe` outcome+ it ("classifies a detail-only constraint " <> show name) $+ mapAppend (serverUsage "detail-only error" (Just ("constraint: " <> name <> "; " <> detail)))+ `shouldBe` expected++ it "does not match stream_events_pkey as events_pkey" $+ mapAppend (serverUsage (constraintMessage "stream_events_pkey") (Just composite))+ `shouldBe` DuplicateEvent (Just eid)+ forM_ ["events_pkey", "stream_events_pkey", "ix_streams_stream_name", "ux_stream_events_stream_version"] $ \name ->+ forM_ ["prefix_" <> name, name <> "_suffix", name <> "$suffix"] $ \other -> do+ it ("rejects the message lookalike " <> show other) $+ mapAppend (serverUsage (constraintMessage other) (Just scalar)) `shouldBe` fallback+ it ("rejects the detail lookalike " <> show other) $+ mapAppend (serverUsage "detail-only error" (Just (other <> "; " <> scalar))) `shouldBe` fallback++ it "uses the message constraint ahead of misleading detail tokens" $+ mapAppend (serverUsage (constraintMessage "ux_stream_events_stream_version") (Just ("events_pkey " <> scalar)))+ `shouldBe` UnexpectedServerError "23505" (constraintMessage "ux_stream_events_stream_version")+ it "retains an unknown message constraint even when detail names a known constraint" $+ mapAppend (serverUsage (constraintMessage "unknown_constraint") (Just ("events_pkey " <> scalar)))+ `shouldBe` fallback+ forM_ ["events_pkey", "stream_events_pkey"] $ \name ->+ forM_ [Nothing, Just "unparseable localized detail"] $ \detail ->+ it ("does not fabricate an event id for " <> show (name, detail)) $+ mapAppend (serverUsage (constraintMessage name) detail) `shouldBe` DuplicateEvent Nothing+ it "maps a composite duplicate inside a transaction with its event id" $+ mapTransactionUsageError (serverUsage (constraintMessage "stream_events_pkey") (Just composite))+ `shouldBe` DuplicateEvent (Just eid)+ it "preserves the link-specific duplicate classification" $+ mapLinkUsageError (StreamName "orders-1") (serverUsage (constraintMessage "stream_events_pkey") (Just composite))+ `shouldBe` EventAlreadyLinked (StreamName "orders-1") (Just eid)+ it "accepts a quoted constraint in legacy detail" $+ mapAppend (serverUsage "detail-only error" (Just (constraintMessage "stream_events_pkey" <> "; " <> composite)))+ `shouldBe` DuplicateEvent (Just eid)+ it "keeps an unnamed unique violation on the conservative fallback" $+ mapAppend (serverUsage "localized message" Nothing) `shouldBe` fallback+ it "does not classify a transaction constraint lookalike as a duplicate" $+ mapTransactionUsageError (serverUsage (constraintMessage "other_events_pkey") (Just scalar))+ `shouldBe` UnexpectedServerError "23505" (constraintMessage "other_events_pkey")+ it "does not classify a link constraint lookalike as already linked" $+ mapLinkUsageError (StreamName "orders-1") (serverUsage (constraintMessage "stream_events_pkey_backup") (Just composite))+ `shouldBe` UnexpectedServerError "23505" (constraintMessage "stream_events_pkey_backup")+ it "attributes a multi-stream name violation to the named operation" $+ attributeMultiStreamError+ [(StreamName "orders-1", AnyVersion), (StreamName "orders-2", NoStream)]+ (serverUsage (constraintMessage "ix_streams_stream_name") (Just "Key (stream_name)=(orders-2) already exists."))+ `shouldBe` StreamAlreadyExists (StreamName "orders-2")+ it "does not attribute a multi-stream constraint lookalike to a later operation" $+ attributeMultiStreamError+ [(StreamName "orders-1", AnyVersion), (StreamName "orders-2", NoStream)]+ (serverUsage (constraintMessage "ix_streams_stream_name_backup") (Just "Key (stream_name)=(orders-2) already exists."))+ `shouldBe` fallback++constraintMessage :: Text -> Text+constraintMessage name = "duplicate key value violates unique constraint \"" <> name <> "\""++serverUsage :: Text -> Maybe Text -> UsageError+serverUsage message detail =+ SessionUsageError $+ Errors.StatementSessionError 1 0 "" [] True $+ Errors.ServerStatementError $+ Errors.ServerError "23505" message detail Nothing Nothing