packages feed

shibuya-kafka-adapter-0.9.1.0: test/Shibuya/Adapter/Kafka/AckHandleTest.hs

module Shibuya.Adapter.Kafka.AckHandleTest (tests) where

import Control.Concurrent.Async qualified as Async
import Control.Concurrent.MVar (MVar, newEmptyMVar, putMVar, takeMVar)
import Control.Exception (try)
import Data.ByteString (ByteString)
import Data.IORef (IORef, atomicModifyIORef', atomicWriteIORef, newIORef, readIORef)
import Data.Int (Int64)
import Data.Map.Strict qualified as Map
import Effectful (Eff, IOE, Limit (..), Persistence (..), UnliftStrategy (..), liftIO, runEff, withEffToIO, (:>))
import Effectful.Dispatch.Dynamic (interpret)
import Effectful.Error.Static (Error, runErrorNoCallStack, throwError)
import Kafka.Consumer (RdKafkaRespErrT (..))
import Kafka.Consumer.Types (ConsumerRecord (..), Offset (..), PartitionOffset (..), RebalanceEvent (..), Timestamp (..), TopicPartition (..))
import Kafka.Effectful.Consumer.Effect (KafkaConsumer (..))
import Kafka.Types (BatchSize (..), KafkaError (..), PartitionId (..), Timeout (..), TopicName (..))
import Shibuya.Adapter (Adapter (..))
import Shibuya.Adapter.Kafka (kafkaRebalanceHandler)
import Shibuya.Adapter.Kafka.Config (KafkaAdapterConfig (..))
import Shibuya.Adapter.Kafka.Internal (KafkaAcknowledgementException (..), KafkaAdapterState (..), ingestedStream, kafkaSource, mkAckHandle, mkIngested, newKafkaAdapterState)
import Shibuya.App (ProcessorId (..), defaultAppConfig, getAppMaster, mkProcessor, runApp, waitApp)
import Shibuya.Core.Ack (AckDecision (..), HaltReason (..), RetryDelay (..))
import Shibuya.Core.AckHandle (AckHandle (..))
import Shibuya.Core.Ingested (Ingested)
import Shibuya.Internal.Runner.Master (ProcessorLifecycle (..), getLifecycleSnapshot)
import Shibuya.Telemetry.Effect (runTracingNoop)
import Streamly.Data.Fold qualified as Fold
import Streamly.Data.Stream qualified as Stream
import System.Random (mkStdGen, randomR)
import System.Timeout (timeout)
import Test.Tasty (TestTree, testGroup)
import Test.Tasty.HUnit (assertEqual, assertFailure, testCase)

data MockState = MockState
  { storeAttempts :: !Int,
    storedOffsets :: ![Offset],
    pauseAttempts :: !Int,
    seekCalls :: ![TopicPartition],
    seekTimeouts :: ![Timeout],
    storeBlock :: !(Maybe (MVar (), MVar ())),
    storeFailuresRemaining :: !Int,
    pauseFailuresRemaining :: !Int,
    seekFailuresRemaining :: !Int,
    storeError :: !KafkaError,
    pauseError :: !KafkaError,
    seekError :: !KafkaError
  }

tests :: TestTree
tests =
  testGroup
    "AckHandle"
    [ testCase "transient store failures retry and then succeed" testTransientStoreRetry,
      testCase "persistent transient store failure records fatal slot and throws" testPersistentStoreFailure,
      testCase "fatal store failure records fatal slot and throws after one attempt" testFatalStoreFailure,
      testCase "AckHalt pause failure records fatal slot and throws" testAckHaltPauseFailure,
      testCase "AckRetry seeks exact failed offset and does not store" testAckRetrySeeks,
      testCase "AckRetry caps the consumer-lock seek timeout" testAckRetryBoundsSeek,
      testCase "seek barrier prevents stale successor store" testBarrierSkipsSuccessorStore,
      testCase "earliest retry survives later retry and acknowledgement" testEarliestRetrySurvives,
      testCase "one delivery cannot resolve its own retry" testRetryRequiresRedelivery,
      testCase "repeated retry on one delivery is idempotent" testRepeatedRetry,
      testCase "fixed-seed sequences match the earliest-unresolved reference model" testReferenceModelSeeds,
      testCase "exhausted acknowledgement throws immediately" testPersistentStoreFailureThrows,
      testCase "revocation fences an old delivery callback" testRevocationFencesCallback,
      testCase "cancellation releases finalizer and consumer ownership" testCancellationReleasesOwnership,
      testCase "terminal acknowledgement failure reaches the core lifecycle" testTerminalFailureReachesCore,
      testCase "source observes fatal slot before polling" testSourceObservesFatalSlot
    ]

testTransientStoreRetry :: IO ()
testTransientStoreRetry = do
  mock <- newIORef defaultMockState {storeFailuresRemaining = 2}
  state <- newKafkaAdapterState
  result <- runFinalizer mock $ finalizeRecord state (recordAt 42) AckOk
  assertRight result
  final <- readIORef mock
  fatal <- readIORef state.fatalError
  assertEqual "store attempts" 3 final.storeAttempts
  assertEqual "fatal slot" Nothing fatal

testPersistentStoreFailure :: IO ()
testPersistentStoreFailure = do
  let err = KafkaResponseError RdKafkaRespErrTransport
  mock <- newIORef defaultMockState {storeFailuresRemaining = 99, storeError = err}
  state <- newKafkaAdapterState
  assertAckFailure err $ runFinalizer mock $ finalizeRecord state (recordAt 42) AckOk
  final <- readIORef mock
  fatal <- readIORef state.fatalError
  assertEqual "store attempts" 3 final.storeAttempts
  assertEqual "fatal slot" (Just err) fatal

testFatalStoreFailure :: IO ()
testFatalStoreFailure = do
  let err = KafkaBadConfiguration
  mock <- newIORef defaultMockState {storeFailuresRemaining = 99, storeError = err}
  state <- newKafkaAdapterState
  assertAckFailure err $ runFinalizer mock $ finalizeRecord state (recordAt 42) AckOk
  final <- readIORef mock
  fatal <- readIORef state.fatalError
  assertEqual "store attempts" 1 final.storeAttempts
  assertEqual "fatal slot" (Just err) fatal

testAckHaltPauseFailure :: IO ()
testAckHaltPauseFailure = do
  let err = KafkaResponseError RdKafkaRespErrTransport
  mock <- newIORef defaultMockState {pauseFailuresRemaining = 99, pauseError = err}
  state <- newKafkaAdapterState
  assertAckFailure err $ runFinalizer mock $ finalizeRecord state (recordAt 42) (AckHalt (HaltFatal "stop"))
  final <- readIORef mock
  fatal <- readIORef state.fatalError
  assertEqual "pause attempts" 3 final.pauseAttempts
  assertEqual "fatal slot" (Just err) fatal

testAckRetrySeeks :: IO ()
testAckRetrySeeks = do
  mock <- newIORef defaultMockState
  state <- newKafkaAdapterState
  result <- runFinalizer mock $ finalizeRecord state (recordAt 42) (AckRetry (RetryDelay 0))
  assertRight result
  final <- readIORef mock
  assertEqual "store attempts" 0 final.storeAttempts
  assertEqual
    "seek call"
    [TopicPartition (TopicName "orders") (PartitionId 0) (PartitionOffset 42)]
    final.seekCalls

testAckRetryBoundsSeek :: IO ()
testAckRetryBoundsSeek = do
  mock <- newIORef defaultMockState
  state <- newKafkaAdapterState
  let slowConfig = testConfig {pollTimeout = Timeout 5000}
  result <- runFinalizer mock $ do
    AckHandle finalize <- mkAckHandle state slowConfig (recordAt 42)
    finalize (AckRetry (RetryDelay 0))
  assertRight result
  final <- readIORef mock
  assertEqual "seek timeout" [Timeout 100] final.seekTimeouts

testBarrierSkipsSuccessorStore :: IO ()
testBarrierSkipsSuccessorStore = do
  mock <- newIORef defaultMockState
  state <- newKafkaAdapterState
  result <- runFinalizer mock $ do
    finalizeRecord state (recordAt 42) (AckRetry (RetryDelay 0))
    finalizeRecord state (recordAt 43) AckOk
    finalizeRecord state (recordAt 42) AckOk
  assertRight result
  final <- readIORef mock
  assertEqual "only retried message stored" 1 final.storeAttempts

testEarliestRetrySurvives :: IO ()
testEarliestRetrySurvives = do
  mock <- newIORef defaultMockState
  state <- newKafkaAdapterState
  result <- runFinalizer mock $ do
    finalizeRecord state (recordAt 42) (AckRetry (RetryDelay 0))
    finalizeRecord state (recordAt 43) (AckRetry (RetryDelay 0))
    finalizeRecord state (recordAt 43) AckOk
  assertRight result
  final <- readIORef mock
  assertEqual
    "both retries seek the earliest unresolved offset"
    [ TopicPartition (TopicName "orders") (PartitionId 0) (PartitionOffset 42),
      TopicPartition (TopicName "orders") (PartitionId 0) (PartitionOffset 42)
    ]
    final.seekCalls
  assertEqual "later acknowledgement remains fenced" 0 final.storeAttempts

testRetryRequiresRedelivery :: IO ()
testRetryRequiresRedelivery = do
  mock <- newIORef defaultMockState
  state <- newKafkaAdapterState
  result <- runFinalizer mock $ do
    AckHandle finalize <- mkAckHandle state testConfig (recordAt 42)
    finalize (AckRetry (RetryDelay 0))
    finalize AckOk
  assertRight result
  final <- readIORef mock
  assertEqual "same delivery cannot store after requesting retry" 0 final.storeAttempts

testRepeatedRetry :: IO ()
testRepeatedRetry = do
  mock <- newIORef defaultMockState
  state <- newKafkaAdapterState
  result <- runFinalizer mock $ do
    AckHandle finalize <- mkAckHandle state testConfig (recordAt 42)
    finalize (AckRetry (RetryDelay 0))
    finalize (AckRetry (RetryDelay 0))
  assertRight result
  final <- readIORef mock
  assertEqual "one successful retry performs one seek" 1 (length final.seekCalls)

testReferenceModelSeeds :: IO ()
testReferenceModelSeeds =
  -- EP-44's integrated release gate requires at least 1,000 recorded model
  -- cases. Keep the range deterministic so a failure names its replayable
  -- seed and the ordinary package suite exercises the full gate.
  mapM_ runSeed [400040 .. 401039]
  where
    runSeed seed = do
      let generator = mkStdGen seed
          (baseDelta, generator') = randomR (0, 5 :: Int64) generator
          (gap, generator'') = randomR (1, 5 :: Int64) generator'
          (laterFirst, _) = randomR (False, True) generator''
          base = 40 + baseDelta
          later = base + gap
          (firstOffset, secondOffset) = if laterFirst then (later, base) else (base, later)
          retryOrder = [firstOffset, secondOffset]
          expectedSeeks = map toTopicPartition (runningMinimum retryOrder)
      mock <- newIORef defaultMockState
      state <- newKafkaAdapterState
      result <- runFinalizer mock $ do
        first <- mkAckHandle state testConfig (recordAt firstOffset)
        second <- mkAckHandle state testConfig (recordAt secondOffset)
        finalizeHandle first (AckRetry (RetryDelay 0))
        finalizeHandle second (AckRetry (RetryDelay 0))
        prematureLater <- mkAckHandle state testConfig (recordAt later)
        finalizeHandle prematureLater AckOk
        replayBase <- mkAckHandle state testConfig (recordAt base)
        finalizeHandle replayBase AckOk
        replayLater <- mkAckHandle state testConfig (recordAt later)
        finalizeHandle replayLater AckOk
      assertRight result
      final <- readIORef mock
      assertEqual ("seed " <> show seed <> " seek boundary") expectedSeeks final.seekCalls
      assertEqual ("seed " <> show seed <> " stored offsets") [Offset base, Offset later] final.storedOffsets

    runningMinimum = \case
      [] -> []
      first : rest -> scanl min first rest

    toTopicPartition offset =
      TopicPartition (TopicName "orders") (PartitionId 0) (PartitionOffset offset)

    finalizeHandle (AckHandle finalize) = finalize

testPersistentStoreFailureThrows :: IO ()
testPersistentStoreFailureThrows = do
  let err = KafkaResponseError RdKafkaRespErrTransport
  mock <- newIORef defaultMockState {storeFailuresRemaining = 99, storeError = err}
  state <- newKafkaAdapterState
  assertAckFailure err $ runFinalizer mock $ finalizeRecord state (recordAt 42) AckOk

testRevocationFencesCallback :: IO ()
testRevocationFencesCallback = do
  mock <- newIORef defaultMockState
  state <- newKafkaAdapterState
  result <- runFinalizer mock $ do
    AckHandle finalize <- mkAckHandle state testConfig (recordAt 42)
    liftIO $
      kafkaRebalanceHandler
        state
        (error "consumer handle is not inspected")
        (RebalanceRevoke [(TopicName "orders", PartitionId 0)])
    finalize AckOk
  assertRight result
  final <- readIORef mock
  assertEqual "revoked callback cannot store" 0 final.storeAttempts

testCancellationReleasesOwnership :: IO ()
testCancellationReleasesOwnership = do
  started <- newEmptyMVar
  release <- newEmptyMVar
  mock <- newIORef defaultMockState {storeBlock = Just (started, release)}
  state <- newKafkaAdapterState
  result <- runFinalizer mock $ do
    AckHandle finalize <- mkAckHandle state testConfig (recordAt 42)
    withEffToIO (ConcUnlift Persistent Unlimited) $ \runInIO -> do
      worker <- Async.async (runInIO (finalize AckOk))
      takeMVar started
      Async.cancel worker
      atomicModifyIORef' mock (\mockState -> (mockState {storeBlock = Nothing}, ()))
      mbCompleted <- timeout 1000000 (runInIO (finalize AckOk))
      case mbCompleted of
        Nothing -> assertFailure "finalizer lock or consumer lock remained held after cancellation"
        Just () -> pure ()
  assertRight result
  final <- readIORef mock
  assertEqual "cancelled attempt plus successful retry" 2 final.storeAttempts

testTerminalFailureReachesCore :: IO ()
testTerminalFailureReachesCore = do
  let err = KafkaBadConfiguration
      processorId = ProcessorId "kafka-terminal-ack-failure"
  mock <- newIORef defaultMockState {storeFailuresRemaining = 99, storeError = err}
  state <- newKafkaAdapterState
  result <-
    timeout 5000000 $
      runEff . runErrorNoCallStack @KafkaError . runMockConsumer mock . runTracingNoop $ do
        ingested <- mkIngested state testConfig (recordAt 42)
        let adapter =
              Adapter
                { adapterName = "kafka:test-terminal-ack-failure",
                  source = Stream.fromList [ingested],
                  shutdown = pure ()
                }
            processor = mkProcessor adapter (\_ -> pure AckOk)
        appResult <- runApp defaultAppConfig [(processorId, processor)]
        case appResult of
          Left appError -> error $ "runApp failed: " <> show appError
          Right appHandle -> do
            waitApp appHandle
            lifecycle <- getLifecycleSnapshot (getAppMaster appHandle)
            pure (Map.lookup processorId lifecycle)
  case result of
    Just (Right (Just (LifecycleFailed _ _))) -> pure ()
    other -> assertFailure $ "expected retained terminal acknowledgement failure, got: " <> show other

testSourceObservesFatalSlot :: IO ()
testSourceObservesFatalSlot = do
  let err = KafkaBadConfiguration
  mock <- newIORef defaultMockState
  state <- newKafkaAdapterState
  atomicWriteIORef state.fatalError (Just err)
  result <-
    runEff . runErrorNoCallStack @KafkaError . runMockConsumer mock $
      Stream.fold Fold.drain $
        ingestedStream unreachableBuilder (kafkaSource state testConfig)
  assertEqual "source error" (Left err) result

runFinalizer ::
  IORef MockState ->
  Eff '[KafkaConsumer, Error KafkaError, IOE] a ->
  IO (Either KafkaError a)
runFinalizer mock action =
  runEff . runErrorNoCallStack @KafkaError . runMockConsumer mock $
    action

runMockConsumer ::
  (IOE :> es, Error KafkaError :> es) =>
  IORef MockState ->
  Eff (KafkaConsumer : es) a ->
  Eff es a
runMockConsumer mock =
  interpret $ \_env -> \case
    StoreOffsetMessage cr -> attemptStore mock cr.crOffset
    PausePartitions _ -> attemptPause mock
    SeekPartitions tps seekTimeout -> recordSeek mock tps seekTimeout
    PollMessage _ -> error "AckHandleTest: PollMessage not exercised"
    PollMessageBatch _ _ -> error "AckHandleTest: PollMessageBatch not exercised"
    PollMessageEither _ -> error "AckHandleTest: PollMessageEither not exercised"
    CommitOffsetMessage _ _ -> error "AckHandleTest: CommitOffsetMessage not exercised"
    CommitAllOffsets _ -> error "AckHandleTest: CommitAllOffsets not exercised"
    CommitPartitionsOffsets _ _ -> error "AckHandleTest: CommitPartitionsOffsets not exercised"
    StoreOffsets _ -> error "AckHandleTest: StoreOffsets not exercised"
    Assign _ -> error "AckHandleTest: Assign not exercised"
    ResumePartitions _ -> error "AckHandleTest: ResumePartitions not exercised"
    Committed _ _ -> error "AckHandleTest: Committed not exercised"
    Position _ -> error "AckHandleTest: Position not exercised"
    Assignment -> error "AckHandleTest: Assignment not exercised"
    Subscription -> error "AckHandleTest: Subscription not exercised"
    AskConsumerHandle -> error "AckHandleTest: AskConsumerHandle not exercised"

unreachableBuilder ::
  ConsumerRecord (Maybe ByteString) (Maybe ByteString) ->
  Eff es (Ingested es (Maybe ByteString))
unreachableBuilder _ = error "AckHandleTest: source should not yield records"

attemptStore ::
  (IOE :> es, Error KafkaError :> es) =>
  IORef MockState ->
  Offset ->
  Eff es ()
attemptStore mock offset = do
  (mbBlock, mbErr) <-
    liftIO $
      atomicModifyIORef' mock $ \s ->
        let remaining = s.storeFailuresRemaining
            s' =
              s
                { storeAttempts = s.storeAttempts + 1,
                  storedOffsets = s.storedOffsets <> [offset],
                  storeFailuresRemaining = max 0 (remaining - 1)
                }
         in (s', (s.storeBlock, if remaining > 0 then Just s.storeError else Nothing))
  liftIO $ case mbBlock of
    Nothing -> pure ()
    Just (started, release) -> putMVar started () >> takeMVar release
  maybe (pure ()) throwError mbErr

attemptPause :: (IOE :> es, Error KafkaError :> es) => IORef MockState -> Eff es ()
attemptPause mock = do
  mbErr <-
    liftIO $
      atomicModifyIORef' mock $ \s ->
        let remaining = s.pauseFailuresRemaining
            s' = s {pauseAttempts = s.pauseAttempts + 1, pauseFailuresRemaining = max 0 (remaining - 1)}
         in (s', if remaining > 0 then Just s.pauseError else Nothing)
  maybe (pure ()) throwError mbErr

recordSeek :: (IOE :> es, Error KafkaError :> es) => IORef MockState -> [TopicPartition] -> Timeout -> Eff es ()
recordSeek mock tps seekTimeout = do
  mbErr <-
    liftIO $
      atomicModifyIORef' mock $ \s ->
        let remaining = s.seekFailuresRemaining
            s' =
              s
                { seekCalls = s.seekCalls <> tps,
                  seekTimeouts = s.seekTimeouts <> [seekTimeout],
                  seekFailuresRemaining = max 0 (remaining - 1)
                }
         in (s', if remaining > 0 then Just s.seekError else Nothing)
  maybe (pure ()) throwError mbErr

finalizeRecord ::
  (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) =>
  KafkaAdapterState ->
  ConsumerRecord (Maybe ByteString) (Maybe ByteString) ->
  AckDecision ->
  Eff es ()
finalizeRecord state cr decision =
  do
    AckHandle finalize <- mkAckHandle state testConfig cr
    finalize decision

recordAt :: Int64 -> ConsumerRecord (Maybe ByteString) (Maybe ByteString)
recordAt offset =
  ConsumerRecord
    { crTopic = TopicName "orders",
      crPartition = PartitionId 0,
      crOffset = Offset offset,
      crTimestamp = NoTimestamp,
      crHeaders = mempty,
      crKey = Nothing,
      crValue = Just "payload"
    }

testConfig :: KafkaAdapterConfig
testConfig =
  KafkaAdapterConfig
    { topics = [TopicName "orders"],
      pollTimeout = Timeout 100,
      batchSize = BatchSize 100
    }

defaultMockState :: MockState
defaultMockState =
  MockState
    { storeAttempts = 0,
      storedOffsets = [],
      pauseAttempts = 0,
      seekCalls = [],
      seekTimeouts = [],
      storeBlock = Nothing,
      storeFailuresRemaining = 0,
      pauseFailuresRemaining = 0,
      seekFailuresRemaining = 0,
      storeError = KafkaResponseError RdKafkaRespErrTransport,
      pauseError = KafkaResponseError RdKafkaRespErrTransport,
      seekError = KafkaResponseError RdKafkaRespErrTransport
    }

assertRight :: (Show e) => Either e a -> IO ()
assertRight = \case
  Left err -> assertFailure $ "expected Right, got Left: " <> show err
  Right _ -> pure ()

assertAckFailure :: KafkaError -> IO (Either KafkaError ()) -> IO ()
assertAckFailure expected action = do
  result <- try @KafkaAcknowledgementException action
  case result of
    Left (KafkaAcknowledgementException actual) -> assertEqual "acknowledgement error" expected actual
    Right value -> assertFailure $ "expected KafkaAcknowledgementException, got: " <> show value