packages feed

kiroku-store-0.10.0.0: test/Test/HandlerStall.hs

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