packages feed

kiroku-store-0.11.0.0: test/Test/PublisherDropCounter.hs

{-# LANGUAGE OverloadedRecordDot #-}

module Test.PublisherDropCounter (spec) where

import Control.Concurrent.STM
import Control.Exception (bracket)
import Data.Aeson (object)
import Data.IntMap.Strict qualified as IntMap
import Data.Text (Text)
import Data.Vector qualified as V
import Kiroku.Store
import Kiroku.Store.Subscription.EventPublisher qualified as Pub
import System.Timeout (timeout)
import Test.Helpers (makeEvent, waitForPublisher, withTestStore)
import Test.Hspec

spec :: Spec
spec = describe "publisher drop counter" $ do
    it "counts drops, keeps the newest batch and preserves the idempotent wrapper" $ withTestStore $ \store -> do
        baseline <- IntMap.size <$> readTVarIO (Pub.subscribers store.publisher)
        bracket (atomically (Pub.subscribePublisherWith store.publisher 1 DropOldest)) Pub.unsubscribe $ \sub -> do
            appendAndWait store "drop-a" 1
            readTVarIO sub.subscriptionDropped `shouldReturn` 0
            appendAndWait store "drop-b" 2
            readTVarIO sub.subscriptionDropped `shouldReturn` 1
            readTVarIO sub.subscriptionStatus `shouldReturn` Pub.Active
            newest sub.subscriptionQueue `shouldReturn` [GlobalPosition 2]
            bracket (atomically (Pub.subscribePublisher store.publisher 1 DropOldest)) (\(_, _, stop) -> stop) $ \(queue, _, stop) -> do
                appendAndWait store "drop-c" 3
                appendAndWait store "drop-d" 4
                newest queue `shouldReturn` [GlobalPosition 4]
                stop >> stop
            Pub.unsubscribe sub >> Pub.unsubscribe sub
        (IntMap.size <$> readTVarIO (Pub.subscribers store.publisher)) `shouldReturn` baseline
    it "does not count PauseAndResume or DropSubscription overflow" $ withTestStore $ \store ->
        mapM_
            ( \policy -> bracket (atomically (Pub.subscribePublisherWith store.publisher 1 policy)) Pub.unsubscribe $ \sub -> do
                Right pos <- runStoreIO store (appendToStream (StreamName "policies-a") AnyVersion [makeEvent "E" (object [])])
                bounded (waitForPublisher store pos.globalPosition)
                Right pos2 <- runStoreIO store (appendToStream (StreamName "policies-a") AnyVersion [makeEvent "E" (object [])])
                bounded (waitForPublisher store pos2.globalPosition)
                readTVarIO sub.subscriptionDropped `shouldReturn` 0
                readTVarIO sub.subscriptionStatus `shouldReturn` (if policy == PauseAndResume then Pub.Paused else Pub.Overflowed)
            )
            [PauseAndResume, DropSubscription]

appendAndWait :: KirokuStore -> Text -> Int -> IO ()
appendAndWait store name n = do
    Right _ <- runStoreIO store (appendToStream (StreamName name) NoStream [makeEvent "E" (object [])])
    bounded (waitForPublisher store (GlobalPosition (fromIntegral n)))

bounded :: IO a -> IO a
bounded action = timeout 5_000_000 action >>= maybe (fail "publisher timeout") pure

newest :: TBQueue DecodedBatch -> IO [GlobalPosition]
newest queue = do
    Just (UnchangedBatch events) <- atomically (tryReadTBQueue queue)
    atomically (isEmptyTBQueue queue) `shouldReturn` True
    pure (map (.globalPosition) (V.toList events))