kiroku-store-0.11.0.0: test/Test/SubscriptionDeadLetters.hs
{-# LANGUAGE OverloadedLabels #-}
module Test.SubscriptionDeadLetters (spec) where
import Control.Lens ((&), (.~), (^.))
import Control.Monad (forM_, void)
import Data.Aeson (object, (.=))
import Data.Generics.Labels ()
import Data.Int (Int32, Int64)
import Data.Text qualified as T
import Data.Time.Clock (getCurrentTime)
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 hiding (HistoryRetentionInventoryQuery (..), member)
import Kiroku.Store.SQL qualified as SQL
import Test.Helpers
import Test.Hspec
spec :: Spec
spec = describe "SubscriptionDeadLetters" $ do
it "validates page sizes including both boundaries" $ do
forM_ [0, -1, 1001] $ \n -> mkSubscriptionDeadLetterLimit n `shouldBe` Left (SubscriptionDeadLetterLimitOutOfRange n)
forM_ [1, 1000] $ \n -> fmap subscriptionDeadLetterLimitValue (mkSubscriptionDeadLetterLimit n) `shouldBe` Right n
subscriptionDeadLetterLimitValue defaultSubscriptionDeadLetterLimit `shouldBe` 100
around withTestStore $ do
it "returns empty pages through both Store interpreters" $ \store -> do
let query = defaultSubscriptionDeadLetterQuery (SubscriptionName "absent")
runStoreIO store (subscriptionDeadLetters query) `shouldReturn` Right (SubscriptionDeadLetterPage V.empty Nothing)
(runEff . runErrorNoCallStack @StoreError . runKirokuStoreWith store . runStoreResource $ subscriptionDeadLetters query)
`shouldReturn` Right (SubscriptionDeadLetterPage V.empty Nothing)
it "pages without omissions, isolates names and matches the internal member read" $ \store -> do
events <- seedEvents store 7
mapM_ (insertDeadLetterForEvent store "page") (take 5 events)
mapM_ (insertDeadLetterForEvent store "other") (drop 5 events)
let query = (defaultSubscriptionDeadLetterQuery (SubscriptionName "page")){limit = validLimit 2}
first <- page store query
second <- page store query{after = first ^. #nextCursor}
third <- page store query{after = second ^. #nextCursor}
map positions [first, second, third] `shouldBe` [[5, 4], [3, 2], [1]]
third ^. #nextCursor `shouldBe` Nothing
Right internal <- Pool.use (store ^. #pool) (Session.statement ("page", 0) SQL.readDeadLettersStmt)
map SQL.deadLetterGlobalPosition (V.toList internal) `shouldBe` [5, 4, 3, 2, 1]
forM_ [5, 1000] $ \n -> do
final <- page store query{limit = validLimit n}
positions final `shouldBe` [5, 4, 3, 2, 1]
final ^. #nextCursor `shouldBe` Nothing
it "merges historical members and resolves same-position ties with an exclusive cursor" $ \store -> do
[event] <- seedEvents store 1
forM_ [0, 1, 9] $ \member -> insertDeadLetterWith store "ties" member (object []) "tie" 2 event
let query = (defaultSubscriptionDeadLetterQuery (SubscriptionName "ties")){limit = validLimit 1}
first <- page store query
second <- page store query{after = first ^. #nextCursor}
third <- page store query{after = second ^. #nextCursor}
map (\p -> map (\row -> row ^. #consumerGroupMember) (V.toList (p ^. #deadLetters))) [first, second, third] `shouldBe` [[9], [1], [0]]
third ^. #nextCursor `shouldBe` Nothing
forM_ [0, 1, 9, 7] $ \member -> do
selected <- page store (query & #consumerGroupMember .~ Just member)
positions selected `shouldBe` if member == 7 then [] else [1]
it "preserves arbitrary and decode-failure JSON unchanged" $ \store -> do
events <- seedEvents store 2
let reasons = [object ["kind" .= ("other" :: T.Text), "detail" .= object ["code" .= (42 :: Int)]], object ["kind" .= ("decode_failure" :: T.Text), "detail" .= ("bad version" :: T.Text)]]
forM_ (zip events reasons) $ \(event, reason) -> insertDeadLetterWith store "json" 0 reason "summary" 3 event
result <- page store (defaultSubscriptionDeadLetterQuery (SubscriptionName "json"))
map (\row -> row ^. #reason) (V.toList (result ^. #deadLetters)) `shouldBe` reverse reasons
now <- getCurrentTime
forM_ (V.toList (result ^. #deadLetters)) $ \row -> do
row ^. #attemptCount `shouldBe` 3
row ^. #createdAt `shouldSatisfy` (<= now)
it "never invokes a configured event decode hook for dead-letter reads" $ \_ ->
withTestStoreSettings (\settings -> (settings & #storeSettings .~ defaultStoreSettings{decodeHook = Just (\_ -> fail "dead-letter read invoked decode hook")})) $ \store -> do
Right _ <- runStoreIO store $ appendToStream (StreamName "hook-1") NoStream [makeEvent "E" (object [])]
Right () <- Pool.use (store ^. #pool) (Session.script "INSERT INTO dead_letters(subscription_name,global_position,event_id,reason,reason_summary,attempt_count) SELECT 'hook',1,event_id,'{}','fixture',1 FROM events")
result <- page store (defaultSubscriptionDeadLetterQuery (SubscriptionName "hook"))
positions result `shouldBe` [1]
it "keeps a cursor valid after its event and row are hard-deleted" $ \store -> do
events <- seedEvents store 3
mapM_ (insertDeadLetterForEvent store "delete") events
let query = (defaultSubscriptionDeadLetterQuery (SubscriptionName "delete")){limit = validLimit 1}
first <- page store query
Right _ <- runStoreIO store $ hardDeleteStream (StreamName "dead-3")
older <- page store query{after = first ^. #nextCursor, limit = validLimit 10}
positions older `shouldBe` [2, 1]
it "includes legal maxBound pairs on first pages" $ \store -> do
[event] <- seedEvents store 1
insertDeadLetterForEvent store "boundary" event
Right () <- Pool.use (store ^. #pool) (Session.script "UPDATE dead_letters SET global_position=9223372036854775807, dead_letter_id=9223372036854775807 WHERE subscription_name='boundary'")
let query = defaultSubscriptionDeadLetterQuery (SubscriptionName "boundary")
forM_ [Nothing, Just 0] $ \member -> do
result <- page store (query & #consumerGroupMember .~ member)
positions result `shouldBe` [maxBound]
older <- page store (query & #consumerGroupMember .~ member & #after .~ Just (SubscriptionDeadLetterCursor (GlobalPosition maxBound) maxBound))
positions older `shouldBe` []
it "reads a real worker-produced poison reason without changing its checkpoint" $ \store -> do
events <- seedEvents store 3
waitForPublisher store (GlobalPosition 3)
handle <- subscribe store $ defaultSubscriptionConfig (SubscriptionName "worker") AllStreams $ \event -> pure $ case event ^. #globalPosition of
GlobalPosition 2 -> DeadLetter (DeadLetterPoison "boom")
GlobalPosition 3 -> Stop
_ -> Continue
waitWithTimeout 10_000_000 handle >>= \case
Right (Right ()) -> pure ()
other -> expectationFailure (show other)
result <- page store (defaultSubscriptionDeadLetterQuery (SubscriptionName "worker"))
case V.toList (result ^. #deadLetters) of
[row] -> do
row ^. #eventId `shouldBe` (events !! 1) ^. #eventId
row ^. #reason `shouldBe` object ["kind" .= ("poison" :: T.Text), "detail" .= ("boom" :: T.Text)]
row ^. #reasonSummary `shouldBe` "poison: boom"
_ -> expectationFailure (show result)
Pool.use (store ^. #pool) (Session.statement ("worker", 0) SQL.getCheckpointMemberStmt) `shouldReturn` Right (Just 3)
validLimit :: Int32 -> SubscriptionDeadLetterLimit
validLimit = either (error . show) (\value -> value) . mkSubscriptionDeadLetterLimit
page :: KirokuStore -> SubscriptionDeadLetterQuery -> IO SubscriptionDeadLetterPage
page store query = runStoreIO store (subscriptionDeadLetters query) >>= either (fail . show) pure
positions :: SubscriptionDeadLetterPage -> [Int64]
positions result = map (\row -> let GlobalPosition p = row ^. #globalPosition in p) (V.toList (result ^. #deadLetters))
seedEvents :: KirokuStore -> Int -> IO [RecordedEvent]
seedEvents store count = do
forM_ [1 .. count] $ \n -> void $ runStoreIO store $ appendToStream (StreamName ("dead-" <> T.pack (show n))) NoStream [makeEvent "E" (object [])]
Right events <- runStoreIO store $ readAllForward (GlobalPosition 0) (fromIntegral count)
pure (V.toList events)