packages feed

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

module Test.StreamHead (spec) where

import Control.Concurrent.Async qualified as Async
import Control.Concurrent.MVar (newEmptyMVar, putMVar, takeMVar)
import Control.Lens ((&), (.~), (^.))
import Control.Monad (forM_, replicateM_, void)
import Data.Aeson (Value (Null))
import Data.Generics.Labels ()
import Data.IORef (modifyIORef', newIORef, readIORef)
import Data.Text (Text)
import Data.Vector qualified as V
import Effectful (runEff)
import Effectful.Error.Static (runErrorNoCallStack)
import Hasql.Pool qualified as Pool
import Hasql.Session qualified as Session
import Kiroku.Store
import System.Timeout (timeout)
import Test.Helpers (makeEvent, withTestStore, withTestStoreSettings)
import Test.Hspec

spec :: Spec
spec = describe "stream head" $ do
    forM_ [False, True] $ \resource -> describe (if resource then "resource runner" else "direct runner") $ do
        let readHead store name =
                if resource
                    then runEff . runErrorNoCallStack @StoreError . runKirokuStoreWith store . runStoreResource $ getStreamWithHead name
                    else runStoreIO store (getStreamWithHead name)
            check store name version headPosition = do
                Right (Just metadata) <- runStoreIO store (getStream name)
                metadata ^. #version `shouldBe` StreamVersion version
                readHead store name `shouldReturn` Right (Just (metadata, headPosition))
        it "distinguishes absent, empty and reserved $all in empty and populated stores" $ withTestStore $ \store -> do
            readHead store (StreamName "missing") `shouldReturn` Right Nothing
            raw store "INSERT INTO streams(stream_name) VALUES ('empty')"
            check store (StreamName "empty") 0 Nothing
            check store (StreamName "$all") 0 Nothing
            appendOne store (StreamName "origin")
            check store (StreamName "$all") 1 Nothing
        it "captures interleaved heads, logical lifecycle and a fresh identity after hard deletion" $ withTestStore $ \store -> do
            let a = StreamName "a"; b = StreamName "b"
            mapM_ (appendOne store) [a, b, a, b]
            originated <- appendPosition store a
            originated `shouldBe` GlobalPosition 5
            check store a 3 (Just originated)
            appendOne store b
            check store a 3 (Just originated)
            Right (Just _) <- runStoreIO store (setStreamTruncateBefore a (StreamVersion 3))
            check store a 3 (Just originated)
            Right (Just _) <- runStoreIO store (softDeleteStream a)
            check store a 3 (Just originated)
            Right (Just (old, _)) <- readHead store a
            Right (Just _) <- runStoreIO store (hardDeleteStream a)
            readHead store a `shouldReturn` Right Nothing
            raw store "INSERT INTO streams(stream_name) VALUES ('a')"
            check store a 0 Nothing
            Right (Just (fresh, _)) <- readHead store a
            fresh ^. #id `shouldNotBe` (old ^. #id)
            recreated <- appendPosition store a
            recreated `shouldBe` GlobalPosition 7
            check store a 1 (Just recreated)
        it "ignores links while tracking subsequent originated appends" $ withTestStore $ \store -> do
            let mixed = StreamName "mixed"; source = StreamName "source"; linked = StreamName "linked"
            appendOne store mixed
            appendOne store source
            Right events <- runStoreIO store (readAllForward (GlobalPosition 0) 10)
            let eid = (events V.! 1) ^. #eventId
            Right _ <- runStoreIO store (linkToStream linked [eid])
            check store linked 1 Nothing
            Right _ <- runStoreIO store (linkToStream mixed [eid])
            check store mixed 2 (Just ((events V.! 0) ^. #globalPosition))
            next <- appendPosition store mixed
            next `shouldBe` GlobalPosition 3
            check store mixed 3 (Just next)
    it "never invokes the event decode hook" $ do
        calls <- newIORef (0 :: Int)
        let hook _ = modifyIORef' calls (+ 1) >> error "unexpected event decode"
        withTestStoreSettings (\settings -> settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just hook}) $ \store -> do
            expected <- appendPosition store (StreamName "hook")
            Right (Just (_, actual)) <- runStoreIO store (getStreamWithHead (StreamName "hook"))
            actual `shouldBe` Just expected
            readIORef calls `shouldReturn` 0
    it "observes version and head from one snapshot while appends commit" $ withTestStore $ \store -> do
        let name = StreamName "concurrent"
        appendOne store name
        start <- newEmptyMVar
        let writer = takeMVar start >> replicateM_ 100 (appendOne store name)
            reader = replicateM_ 200 $ do
                Right (Just (info, Just (GlobalPosition position))) <- runStoreIO store (getStreamWithHead name)
                info ^. #version `shouldBe` StreamVersion position
        outcome <- timeout 20_000_000 $ Async.withAsync writer $ \worker -> do
            putMVar start ()
            reader
            Async.wait worker
        outcome `shouldBe` Just ()
        Right (Just (info, headPosition)) <- runStoreIO store (getStreamWithHead name)
        info ^. #version `shouldBe` StreamVersion 101
        headPosition `shouldBe` Just (GlobalPosition 101)

appendOne :: KirokuStore -> StreamName -> IO ()
appendOne store name = void (appendPosition store name)

appendPosition :: KirokuStore -> StreamName -> IO GlobalPosition
appendPosition store name = do
    result <- runStoreIO store (appendToStream name AnyVersion [makeEvent "StreamHeadEvent" Null])
    either (fail . show) (pure . (^. #globalPosition)) result

raw :: KirokuStore -> Text -> IO ()
raw store command = Pool.use (store ^. #pool) (Session.script command) >>= either (expectationFailure . show) pure