kiroku-store-0.10.0.0: test/Test/PublisherCallbackResilience.hs
{-# LANGUAGE DeriveAnyClass #-}
{-# LANGUAGE OverloadedStrings #-}
module Test.PublisherCallbackResilience (spec) where
import Control.Concurrent (threadDelay)
import Control.Concurrent.Async qualified as Async
import Control.Concurrent.MVar (newEmptyMVar, tryPutMVar)
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, 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 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
deriving stock (Show)
deriving anyclass (Exception)
timeoutMicros :: Int
timeoutMicros = 10_000_000
within :: String -> IO a -> IO a
within label action = do
result <- Async.race (threadDelay timeoutMicros) action
case result of
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
caughtUp <- newEmptyMVar
delivered <- newIORef ([] :: [EventType])
deliveredCount <- newTVarIO (0 :: Int)
let subName = SubscriptionName "publisher-decode-hook-resilience"
decodeOnce event
| event ^. #eventType == EventType "Boom" = do
alreadyFailed <- atomicModifyIORef' failedOnce (\old -> (True, old))
if alreadyFailed
then pure (Right event)
else throwIO CallbackBoom
| otherwise = pure (Right event)
observe evt = do
case evt of
KirokuEventPublisherLoopError{} -> () <$ tryPutMVar loopErrorSeen ()
_ -> pure ()
caughtUpEventHandler subName caughtUp Nothing evt
tweak settings =
settings
& #storeSettings .~ defaultStoreSettings{decodeHook = Just decodeOnce}
& #eventHandler .~ Just observe
handler event = do
modifyIORef' delivered ((event ^. #eventType) :)
atomically $ do
n <- readTVar deliveredCount
let n' = n + 1
writeTVar deliveredCount n'
pure $ if n' >= 2 then Stop else Continue
withTestStoreSettings tweak $ \store -> do
handle <- subscribe store (defaultSubscriptionConfig subName AllStreams handler)
waitForSubscriptionLive caughtUp
Right _ <- runStoreIO store $ appendToStream (StreamName "publisher-decode-boom") NoStream [makeEvent "Boom" (Aeson.object [])]
within "publisher loop error event" (waitForSubscriptionLive loopErrorSeen)
Right _ <- runStoreIO store $ appendToStream (StreamName "publisher-decode-ok") NoStream [makeEvent "Ok" (Aeson.object [])]
result <- waitWithTimeout timeoutMicros handle
case result of
Right (Right ()) -> pure ()
Left timeout -> expectationFailure timeout
Right (Left e) -> expectationFailure ("expected clean stop, got: " <> show e)
seen <- reverse <$> readIORef delivered
seen `shouldBe` [EventType "Boom", EventType "Ok"]
it "drops throwing eventHandler exceptions without killing publisher or worker" $ do
deliveredCount <- newTVarIO (0 :: Int)
let tweak settings = settings & #eventHandler .~ Just (\_ -> throwIO CallbackBoom)
handler _ = do
atomically $ do
n <- readTVar deliveredCount
let n' = n + 1
writeTVar deliveredCount n'
pure $ if n' >= 3 then Stop else Continue
withTestStoreSettings tweak $ \store -> do
handle <- subscribe store (defaultSubscriptionConfig (SubscriptionName "throwing-event-handler-resilience") AllStreams handler)
Right _ <- runStoreIO store $ appendToStream (StreamName "throwing-handler-1") NoStream [makeEvent "One" (Aeson.object [])]
Right _ <- runStoreIO store $ appendToStream (StreamName "throwing-handler-2") NoStream [makeEvent "Two" (Aeson.object [])]
Right _ <- runStoreIO store $ appendToStream (StreamName "throwing-handler-3") NoStream [makeEvent "Three" (Aeson.object [])]
result <- waitWithTimeout timeoutMicros handle
case result of
Right (Right ()) -> pure ()
Left timeout -> expectationFailure timeout
Right (Left e) -> expectationFailure ("expected clean stop, got: " <> show e)
finalCount <- atomically (readTVar deliveredCount)
finalCount `shouldBe` 3