shibuya-kafka-adapter 0.7.0.0 → 0.8.0.1
raw patch · 12 files changed
+880/−108 lines, 12 filesdep ~shibuya-corePVP ok
version bump matches the API change (PVP)
Dependency ranges changed: shibuya-core
API changes (from Hackage documentation)
- Shibuya.Adapter.Kafka: Earliest :: OffsetReset
- Shibuya.Adapter.Kafka: Latest :: OffsetReset
- Shibuya.Adapter.Kafka: [offsetReset] :: KafkaAdapterConfig -> !OffsetReset
- Shibuya.Adapter.Kafka: data OffsetReset
- Shibuya.Adapter.Kafka.Config: [offsetReset] :: KafkaAdapterConfig -> !OffsetReset
+ Shibuya.Adapter.Kafka: data KafkaAdapterState
+ Shibuya.Adapter.Kafka: kafkaAdapterWith :: forall (es :: [Effect]). (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) => KafkaAdapterState -> KafkaAdapterConfig -> Eff es (Adapter es (Maybe ByteString))
+ Shibuya.Adapter.Kafka: kafkaRebalanceHandler :: KafkaAdapterState -> KafkaConsumer -> RebalanceEvent -> IO ()
+ Shibuya.Adapter.Kafka: newKafkaAdapterState :: IO KafkaAdapterState
+ Shibuya.Adapter.Kafka.Convert: extractTraceHeadersFromList :: [(ByteString, ByteString)] -> Maybe TraceHeaders
+ Shibuya.Adapter.Kafka.Internal: KafkaAdapterState :: !TVar Bool -> !IORef (Map PartitionKey Offset) -> !IORef (Maybe KafkaError) -> !MVar () -> KafkaAdapterState
+ Shibuya.Adapter.Kafka.Internal: [consumerLock] :: KafkaAdapterState -> !MVar ()
+ Shibuya.Adapter.Kafka.Internal: [fatalError] :: KafkaAdapterState -> !IORef (Maybe KafkaError)
+ Shibuya.Adapter.Kafka.Internal: [seekBarrier] :: KafkaAdapterState -> !IORef (Map PartitionKey Offset)
+ Shibuya.Adapter.Kafka.Internal: [shutdownVar] :: KafkaAdapterState -> !TVar Bool
+ Shibuya.Adapter.Kafka.Internal: data KafkaAdapterState
+ Shibuya.Adapter.Kafka.Internal: dropStaleRecords :: forall (es :: [Effect]). IOE :> es => KafkaAdapterState -> Stream (Eff es) (Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString))) -> Stream (Eff es) (Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString)))
+ Shibuya.Adapter.Kafka.Internal: newKafkaAdapterState :: IO KafkaAdapterState
+ Shibuya.Adapter.Kafka.Internal: withConsumerLock :: forall (es :: [Effect]) a. IOE :> es => KafkaAdapterState -> Eff es a -> Eff es a
- Shibuya.Adapter.Kafka: KafkaAdapterConfig :: ![TopicName] -> !Timeout -> !BatchSize -> !OffsetReset -> KafkaAdapterConfig
+ Shibuya.Adapter.Kafka: KafkaAdapterConfig :: ![TopicName] -> !Timeout -> !BatchSize -> KafkaAdapterConfig
- Shibuya.Adapter.Kafka.Config: KafkaAdapterConfig :: ![TopicName] -> !Timeout -> !BatchSize -> !OffsetReset -> KafkaAdapterConfig
+ Shibuya.Adapter.Kafka.Config: KafkaAdapterConfig :: ![TopicName] -> !Timeout -> !BatchSize -> KafkaAdapterConfig
- Shibuya.Adapter.Kafka.Internal: kafkaSource :: forall (es :: [Effect]). KafkaConsumer :> es => KafkaAdapterConfig -> Stream (Eff es) (Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString)))
+ Shibuya.Adapter.Kafka.Internal: kafkaSource :: forall (es :: [Effect]). (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) => KafkaAdapterState -> KafkaAdapterConfig -> Stream (Eff es) (Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString)))
- Shibuya.Adapter.Kafka.Internal: mkAckHandle :: forall (es :: [Effect]). KafkaConsumer :> es => ConsumerRecord (Maybe ByteString) (Maybe ByteString) -> AckHandle es
+ Shibuya.Adapter.Kafka.Internal: mkAckHandle :: forall (es :: [Effect]). (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) => KafkaAdapterState -> KafkaAdapterConfig -> ConsumerRecord (Maybe ByteString) (Maybe ByteString) -> AckHandle es
- Shibuya.Adapter.Kafka.Internal: mkIngested :: forall (es :: [Effect]). KafkaConsumer :> es => ConsumerRecord (Maybe ByteString) (Maybe ByteString) -> Ingested es (Maybe ByteString)
+ Shibuya.Adapter.Kafka.Internal: mkIngested :: forall (es :: [Effect]). (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) => KafkaAdapterState -> KafkaAdapterConfig -> ConsumerRecord (Maybe ByteString) (Maybe ByteString) -> Ingested es (Maybe ByteString)
Files
- CHANGELOG.md +54/−0
- README.md +16/−0
- shibuya-kafka-adapter.cabal +4/−3
- src/Shibuya/Adapter/Kafka.hs +138/−40
- src/Shibuya/Adapter/Kafka/Config.hs +5/−6
- src/Shibuya/Adapter/Kafka/Convert.hs +17/−16
- src/Shibuya/Adapter/Kafka/Internal.hs +216/−33
- test/Kafka/TestEnv.hs +1/−2
- test/Main.hs +3/−1
- test/Shibuya/Adapter/Kafka/AckHandleTest.hs +246/−0
- test/Shibuya/Adapter/Kafka/ConvertTest.hs +7/−0
- test/Shibuya/Adapter/Kafka/IntegrationTest.hs +173/−7
CHANGELOG.md view
@@ -1,5 +1,59 @@ # Changelog +## 0.8.0.1 — 2026-07-05++### Other Changes++- Require `shibuya-core ^>=0.8.0.1`, picking up the upstream patch release+ with hot-path allocation fixes for `Async` and `Ahead` processing.+- Keep the example and benchmark packages on the shared `0.8.0.1` repo+ version line.++## 0.8.0.0 — 2026-07-04++### Breaking Changes++- Require `shibuya-core ^>=0.8.0.0`.+- Remove the dead `KafkaAdapterConfig.offsetReset` field. Offset reset policy+ belongs to the `Subscription` passed to `runKafkaConsumer`; `topics` remains+ adapter metadata and is checked against the live subscription with a stderr+ warning on mismatch.++### New Features++- Add `kafkaAdapterWith`, `newKafkaAdapterState`, and+ `kafkaRebalanceHandler` for callers that want rebalance logging and eager+ cleanup of retry barriers for revoked partitions.+- Compute Kafka headers once in `consumerRecordToEnvelope` and expose+ `extractTraceHeadersFromList` for callers that already have a materialized+ header list.++### Bug Fixes++- `AckRetry` now seeks the partition back to the failed message instead of+ storing the offset, preserving at-least-once redelivery.+- Handler-exception retries from core no longer allow later buffered messages+ to commit past the failed offset.+- Ack-path Kafka errors are classified inside `finalize`; transient failures+ are retried briefly, and persistent failures terminate the source stream as+ adapter errors instead of handler errors.+- `AckHalt` pause failures no longer cancel the halt decision.+- Idle shutdown exits promptly and ignores Kafka's no-offset response when+ there is nothing to commit.+- `AckDeadLetter` still stores the offset, but now emits a prominent stderr+ warning because this adapter does not include a DLQ producer.++### Other Changes++- Use the `shibuya-core 0.8.0.0` adapter-facing `mkEnvelope` and+ `mkIngested` smart constructors, and update runnable examples to the+ `runApp defaultAppConfig` / `Message` handler API.+- Document the Serial-only processing contract, the dead-letter limitation,+ Kafka's lack of delivery-attempt counts, halt/eviction behavior, and shutdown+ ordering.+- Keep the example and benchmark packages on the shared `0.8.0.0` repo version+ line.+ ## 0.7.0.0 — 2026-06-05 ### Changed
README.md view
@@ -4,6 +4,22 @@ Integrates with Apache Kafka via [`kafka-effectful`](https://github.com/shinzui/kafka-effectful) for the consumer effect (polling, offset store, partition pause) and [`hw-kafka-streamly`](https://hackage.haskell.org/package/hw-kafka-streamly) for error classification (`skipNonFatal`), on top of [`hw-kafka-client`](https://github.com/haskell-works/hw-kafka-client). Provides polling, offset commit semantics, partition awareness, and graceful shutdown. +## Runtime Contract++Run this adapter with serial message processing only. librdkafka stores the highest offset per partition and does not track gaps, so concurrent finalization can commit past an earlier message that failed, halted, or requested retry. The adapter cannot inspect the processor concurrency policy at construction time; avoiding `Async` and `Ahead` processing is a caller contract until a gap-tracking commit layer exists.++Offset reset policy belongs to the `Subscription` passed to `runKafkaConsumer`, for example `topics [TopicName "orders"] <> offsetReset Earliest`. `KafkaAdapterConfig.topics` is metadata used for adapter naming and a construction-time warning if it differs from the live subscription.++`AckRetry` seeks the partition back to the failed message and does not store that message's offset. `AckDeadLetter` stores the offset so the consumer group moves on, prints `[shibuya-kafka-adapter] WARNING: dead-lettered message DROPPED` to stderr, and makes the message unrecoverable from the group's committed position. This adapter does not include a DLQ producer.++Kafka does not expose a per-message delivery-attempt counter through this consumer API, so `Envelope.attempt` is always `Nothing`. Handlers that need bounded retries must use their own store or return `AckHalt` to stop the stream.++`AckHalt` pauses the originating partition and stops the processor. Polling then stops; after `max.poll.interval.ms` (librdkafka default: 300000 ms, or 5 minutes) the broker may evict this consumer and rebalance the partition to another group member. A single-member group stalls until restart.++On shutdown, the adapter commits offsets stored so far and signals the source stream to stop. Let the surrounding `runKafkaConsumer` scope end normally after stopping the app so the consumer close path can flush offsets stored during the drain window.++For rebalance visibility, create state with `newKafkaAdapterState`, install `Kafka.Consumer.setCallback (Kafka.Consumer.rebalanceCallback (kafkaRebalanceHandler state))` before consumer creation, and pass the same state to `kafkaAdapterWith`. The helper logs rebalance events and clears retry barriers for revoked partitions; in-flight fencing for cooperative rebalances is out of scope.+ ## Packages - `shibuya-kafka-adapter` — the adapter library (`Shibuya.Adapter.Kafka`, `.Config`, `.Convert`).
shibuya-kafka-adapter.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.12 name: shibuya-kafka-adapter-version: 0.7.0.0+version: 0.8.0.1 synopsis: Kafka adapter for the Shibuya queue processing framework description: A Shibuya adapter that integrates with Apache Kafka via kafka-effectful@@ -61,7 +61,7 @@ , hw-kafka-client >=5.3 && <6 , hw-kafka-streamly ^>=0.2 , kafka-effectful ^>=0.3.0.0- , shibuya-core ^>=0.7.0.0+ , shibuya-core ^>=0.8.0.1 , stm ^>=2.5 , streamly ^>=0.11 , streamly-core ^>=0.3@@ -92,6 +92,7 @@ other-modules: Kafka.TestEnv+ Shibuya.Adapter.Kafka.AckHandleTest Shibuya.Adapter.Kafka.AdapterTest Shibuya.Adapter.Kafka.ConvertTest Shibuya.Adapter.Kafka.IntegrationTest@@ -108,7 +109,7 @@ , kafka-effectful , process , random- , shibuya-core ^>=0.7.0.0+ , shibuya-core ^>=0.8.0.1 , shibuya-kafka-adapter , stm , streamly
src/Shibuya/Adapter/Kafka.hs view
@@ -7,7 +7,7 @@ == Example Usage @-import Shibuya.App (runApp, mkProcessor)+import Shibuya.App (defaultAppConfig, runApp, mkProcessor) import Shibuya.Adapter.Kafka (kafkaAdapter, defaultConfig) import Kafka.Effectful.Consumer (runKafkaConsumer) import Kafka.Consumer (brokersList, groupId, noAutoOffsetStore)@@ -18,7 +18,7 @@ . runKafkaConsumer props sub $ do adapter <- kafkaAdapter (defaultConfig [TopicName \"orders\"])- result <- runApp IgnoreFailures 100+ result <- runApp defaultAppConfig [ (ProcessorId \"orders\", mkProcessor adapter myHandler) ] ...@@ -26,13 +26,38 @@ == Message Lifecycle -1. Messages are polled from Kafka in batches-2. Each message is wrapped as an 'Ingested' with an 'AckHandle'-3. On 'AckOk', the offset is stored (auto-commit flushes to broker)-4. On 'AckRetry', the offset is stored (Kafka cannot un-read messages)-5. On 'AckDeadLetter', the offset is stored (DLQ production in future milestone)-6. On 'AckHalt', the partition is paused and offset is NOT stored+1. Messages are polled from Kafka in batches.+2. Each message is wrapped as an @Ingested@ value with an @AckHandle@ via+ Shibuya's adapter-facing smart constructors.+3. On @AckOk@, the offset is stored locally; auto-commit or consumer close+ later flushes stored offsets to the broker.+4. On @AckRetry@, the offset is not stored. The adapter seeks the partition+ back to the failed message so Kafka can redeliver it.+5. On @AckDeadLetter@, the offset is stored after a loud stderr warning.+6. On @AckHalt@, the partition is paused and offset is not stored. +== Serial Operation Required++This adapter must be run with serial message processing. librdkafka stores the+highest offset per partition without gap tracking, so concurrent finalization+can commit past an earlier message that failed, halted, or requested retry.+The 'Adapter' value does not contain the processor concurrency policy, so this+is a caller contract rather than a runtime guard: do not use @Async@ or @Ahead@+processing with this adapter until a gap-tracking commit layer exists.++== Dead Letters Are Dropped++This adapter does not include a DLQ producer. @AckDeadLetter@ stores the+message offset so the consumer group moves on, emits a warning to stderr, and+makes the message unrecoverable from that consumer group's committed position.+Core tracing still records the dead-letter decision and reason on the+per-message span.++Kafka does not expose a per-message delivery counter through this consumer+API, so 'Shibuya.Core.Types.Envelope.attempt' is always 'Nothing'. Handlers+cannot safely cap retries by counting attempts from the envelope; use an+external store, or return @AckHalt@ to stop the stream.+ == Fatal Error Propagation Non-fatal Kafka errors (poll timeouts, partition EOFs, and the rest of the@@ -46,17 +71,31 @@ == AckHalt Partition Pause Semantics -'AckHalt' pauses the originating partition by calling @pausePartitions@ from-@kafka-effectful@. The partition is not automatically resumed within the-current consumer session — the adapter has no side channel for a handler to-request resumption. A new consumer session (a new call to @runKafkaConsumer@)-starts with no paused partitions, so resumption happens implicitly on restart-but not mid-session. If mid-session resume is required, the caller must-manage that explicitly outside the 'Adapter' surface.+@AckHalt@ pauses the originating partition by calling @pausePartitions@ from+@kafka-effectful@ and the processor stops. Polling therefore stops. After+@max.poll.interval.ms@ (librdkafka default: 300000 ms, or 5 minutes) the broker+may evict this consumer from its group and rebalance the partition to another+member, which resumes from the last committed offset. A single-member group+simply stalls until restart. Paused state is session-local and does not outlive+the current consumer.++== Rebalance Callback Helper++'kafkaRebalanceHandler' is optional. Install it with+@Kafka.Consumer.setCallback (Kafka.Consumer.rebalanceCallback (kafkaRebalanceHandler state))@+before creating the consumer when you want stderr visibility into assignment+changes and eager cleanup of retry barriers for revoked partitions. Without it,+the seek barrier still self-heals when messages are finalized at or below the+barrier offset. Cooperative rebalance fencing of in-flight work is outside this+adapter's scope. -} module Shibuya.Adapter.Kafka ( -- * Adapter kafkaAdapter,+ kafkaAdapterWith,+ KafkaAdapterState,+ newKafkaAdapterState,+ kafkaRebalanceHandler, -- * Configuration KafkaAdapterConfig (..),@@ -68,7 +107,6 @@ TopicName (..), BrokerAddress (..), ConsumerGroupId (..),- OffsetReset (..), OffsetCommit (..), Timeout (..), BatchSize (..),@@ -76,20 +114,24 @@ ) where -import Control.Concurrent.STM (TVar, atomically, newTVarIO, readTVarIO, writeTVar)+import Control.Concurrent.STM (atomically, writeTVar) import Control.Monad.IO.Class (liftIO) import Data.ByteString (ByteString)+import Data.IORef (atomicModifyIORef')+import Data.Map.Strict qualified as Map+import Data.Set qualified as Set import Data.Text qualified as Text import Effectful (Eff, IOE, (:>))-import Effectful.Error.Static (Error)-import Kafka.Consumer.Types (ConsumerGroupId (..), OffsetCommit (..), OffsetReset (..))-import Kafka.Effectful.Consumer.Effect (KafkaConsumer, commitAllOffsets)-import Kafka.Types (BatchSize (..), BrokerAddress (..), KafkaError, Timeout (..), TopicName (..))+import Effectful.Error.Static (Error, catchError, throwError)+import Kafka.Consumer (RdKafkaRespErrT (..))+import Kafka.Consumer.Types (ConsumerGroupId (..), OffsetCommit (..), RebalanceEvent (..))+import Kafka.Consumer.Types qualified as KC+import Kafka.Effectful.Consumer.Effect (KafkaConsumer, commitAllOffsets, subscription)+import Kafka.Types (BatchSize (..), BrokerAddress (..), KafkaError (..), PartitionId, Timeout (..), TopicName (..)) import Shibuya.Adapter (Adapter (..)) import Shibuya.Adapter.Kafka.Config (KafkaAdapterConfig (..), defaultConfig)-import Shibuya.Adapter.Kafka.Internal (ingestedStream, kafkaSource, mkIngested)-import Streamly.Data.Stream (Stream)-import Streamly.Data.Stream qualified as Stream+import Shibuya.Adapter.Kafka.Internal (KafkaAdapterState (..), dropStaleRecords, ingestedStream, kafkaSource, mkIngested, newKafkaAdapterState, withConsumerLock)+import System.IO (hPutStrLn, stderr) {- | Create a Kafka adapter with the given configuration. @@ -99,7 +141,10 @@ The adapter uses @noAutoOffsetStore@ with manual @storeOffsetMessage@ + auto-commit for offset management. On shutdown, @commitAllOffsets@ flushes-any stored offsets to the broker.+offsets stored so far. Messages finalized during the drain window store offsets+after that explicit commit; let the surrounding @runKafkaConsumer@ scope end+normally after 'Shibuya.App.stopApp' returns so the consumer close path can+flush the final stored offsets under the same auto-commit mode. The returned 'Shibuya.Adapter.Adapter.shutdown' action must be invoked while the 'KafkaConsumer' effect is still in scope. Invoking it after@@ -112,24 +157,77 @@ KafkaAdapterConfig -> Eff es (Adapter es (Maybe ByteString)) kafkaAdapter config = do- shutdownVar <- liftIO $ newTVarIO False- let messageSource = ingestedStream mkIngested (kafkaSource config)+ state <- liftIO newKafkaAdapterState+ kafkaAdapterWith state config++{- | Create a Kafka adapter using caller-owned adapter state.++Use this when the same @KafkaAdapterState@ must also be referenced by+'kafkaRebalanceHandler', which is installed in consumer properties before the+consumer is created.+-}+kafkaAdapterWith ::+ (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) =>+ KafkaAdapterState ->+ KafkaAdapterConfig ->+ Eff es (Adapter es (Maybe ByteString))+kafkaAdapterWith state config = do+ warnOnSubscriptionMismatch config+ let messageSource =+ ingestedStream (mkIngested state config) $+ dropStaleRecords state $+ kafkaSource state config pure Adapter { adapterName = "kafka:" <> Text.intercalate "," (map unTopicName config.topics)- , source = takeUntilShutdown shutdownVar messageSource+ , source = messageSource , shutdown = do- liftIO $ atomically $ writeTVar shutdownVar True- commitAllOffsets OffsetCommit+ liftIO $ atomically $ writeTVar state.shutdownVar True+ withConsumerLock state (commitAllOffsets OffsetCommit)+ `catchError` \_ err -> case err of+ KafkaResponseError RdKafkaRespErrNoOffset -> pure ()+ _ -> throwError err } --- | Take from stream until shutdown signal is set.-takeUntilShutdown ::- (IOE :> es) =>- TVar Bool ->- Stream (Eff es) a ->- Stream (Eff es) a-takeUntilShutdown shutdownVar =- Stream.takeWhileM $ \_ -> do- isShutdown <- liftIO $ readTVarIO shutdownVar- pure (not isShutdown)+warnOnSubscriptionMismatch ::+ (KafkaConsumer :> es, IOE :> es) =>+ KafkaAdapterConfig ->+ Eff es ()+warnOnSubscriptionMismatch config = do+ liveSubscription <- subscription+ let configured = Set.fromList config.topics+ subscribed = Set.fromList (map fst liveSubscription)+ if configured == subscribed+ then pure ()+ else+ liftIO $+ hPutStrLn stderr $+ "[shibuya-kafka-adapter] WARNING: config topics differ from live Kafka subscription; configured="+ <> show (Set.toList configured)+ <> " subscribed="+ <> show (Set.toList subscribed)++{- | Rebalance callback helper for caller-installed Kafka callbacks.++Install with 'Kafka.Consumer.setCallback' and+'Kafka.Consumer.rebalanceCallback' before creating the consumer. The callback+logs every rebalance event to stderr and clears pending retry barriers for+revoked partitions. It does not fence in-flight work.+-}+kafkaRebalanceHandler ::+ KafkaAdapterState ->+ KC.KafkaConsumer ->+ RebalanceEvent ->+ IO ()+kafkaRebalanceHandler state _consumer event = do+ hPutStrLn stderr $ "[shibuya-kafka-adapter] rebalance: " <> show event+ case event of+ RebalanceRevoke revoked ->+ clearRevokedBarriers revoked+ _ ->+ pure ()+ where+ clearRevokedBarriers :: [(TopicName, PartitionId)] -> IO ()+ clearRevokedBarriers revoked =+ atomicModifyIORef' state.seekBarrier $ \barriers ->+ (foldr Map.delete barriers revoked, ())
src/Shibuya/Adapter/Kafka/Config.hs view
@@ -9,7 +9,6 @@ where import GHC.Generics (Generic)-import Kafka.Consumer.Types (OffsetReset (..)) import Kafka.Types (BatchSize (..), Timeout (..), TopicName) {- | Configuration for the Kafka adapter.@@ -20,13 +19,15 @@ -} data KafkaAdapterConfig = KafkaAdapterConfig { topics :: ![TopicName]- -- ^ Topics to consume from+ {- ^ Topics expected by the adapter, used for observability metadata and+ checked against the live consumer subscription at construction. The+ actual subscription, including offset-reset policy, is supplied to+ @runKafkaConsumer@ by the caller.+ -} , pollTimeout :: !Timeout -- ^ Timeout for each poll call (default: 1000ms) , batchSize :: !BatchSize -- ^ Maximum messages per poll batch (default: 100)- , offsetReset :: !OffsetReset- -- ^ Where to start when no committed offset exists (default: Earliest) } deriving stock (Show, Eq, Generic) @@ -36,7 +37,6 @@ * @pollTimeout@: 1000ms * @batchSize@: 100-* @offsetReset@: Earliest -} defaultConfig :: [TopicName] -> KafkaAdapterConfig defaultConfig ts =@@ -44,5 +44,4 @@ { topics = ts , pollTimeout = Timeout 1000 , batchSize = BatchSize 100- , offsetReset = Earliest }
src/Shibuya/Adapter/Kafka/Convert.hs view
@@ -5,6 +5,7 @@ -- * Trace Context extractTraceHeaders,+ extractTraceHeadersFromList, -- * Timestamp Conversion timestampToUTCTime,@@ -29,7 +30,7 @@ ) import OpenTelemetry.Attributes (Attribute, toAttribute, unkey) import OpenTelemetry.SemanticConventions qualified as Sem-import Shibuya.Core.Types (Cursor (..), Envelope (..), MessageId (..), TraceHeaders)+import Shibuya.Core.Types (Cursor (..), Envelope (..), MessageId (..), TraceHeaders, mkEnvelope) {- | Convert a Kafka 'ConsumerRecord' to a Shibuya 'Envelope'. @@ -58,17 +59,16 @@ ConsumerRecord (Maybe ByteString) (Maybe ByteString) -> Envelope (Maybe ByteString) consumerRecordToEnvelope cr =- Envelope- { messageId = mkMessageId cr.crTopic cr.crPartition cr.crOffset- , cursor = Just (CursorInt (fromIntegral (unOffset cr.crOffset)))- , partition = Just (Text.pack (show (unPartitionId cr.crPartition)))- , enqueuedAt = timestampToUTCTime cr.crTimestamp- , traceContext = extractTraceHeaders cr.crHeaders- , headers = Just (headersToList cr.crHeaders)- , attempt = Nothing- , attributes = kafkaSpanAttributes cr.crPartition cr.crOffset- , payload = cr.crValue- }+ let headerList = headersToList cr.crHeaders+ in (mkEnvelope (mkMessageId cr.crTopic cr.crPartition cr.crOffset) cr.crValue)+ { cursor = Just (CursorInt (fromIntegral (unOffset cr.crOffset)))+ , partition = Just (Text.pack (show (unPartitionId cr.crPartition)))+ , enqueuedAt = timestampToUTCTime cr.crTimestamp+ , traceContext = extractTraceHeadersFromList headerList+ , headers = Just headerList+ , attempt = Nothing+ , attributes = kafkaSpanAttributes cr.crPartition cr.crOffset+ } {- | OpenTelemetry attributes that the framework's per-message span should carry for a Kafka-sourced envelope.@@ -104,14 +104,15 @@ Returns 'Nothing' if @traceparent@ is not present (it's required for valid context). -} extractTraceHeaders :: Headers -> Maybe TraceHeaders-extractTraceHeaders headers =+extractTraceHeaders = extractTraceHeadersFromList . headersToList++-- | Extract W3C trace headers from an already materialized Kafka header list.+extractTraceHeadersFromList :: [(ByteString, ByteString)] -> Maybe TraceHeaders+extractTraceHeadersFromList headerList = case (lookup "traceparent" headerList, lookup "tracestate" headerList) of (Nothing, _) -> Nothing (Just tp, Nothing) -> Just [("traceparent", tp)] (Just tp, Just ts) -> Just [("traceparent", tp), ("tracestate", ts)]- where- headerList :: [(ByteString, ByteString)]- headerList = headersToList headers {- | Convert a Kafka 'Timestamp' to 'UTCTime'.
src/Shibuya/Adapter/Kafka/Internal.hs view
@@ -2,8 +2,14 @@ This module is not part of the public API and may change without notice. -} module Shibuya.Adapter.Kafka.Internal (+ -- * Adapter State+ KafkaAdapterState (..),+ newKafkaAdapterState,+ withConsumerLock,+ -- * Stream Construction kafkaSource,+ dropStaleRecords, ingestedStream, -- * Ingested Construction@@ -14,83 +20,260 @@ ) where +import Control.Concurrent (threadDelay)+import Control.Concurrent.MVar (MVar, newMVar, putMVar, takeMVar)+import Control.Concurrent.STM (TVar, newTVarIO, readTVarIO) import Data.ByteString (ByteString) import Data.Function ((&))-import Effectful (Eff, (:>))-import Effectful.Error.Static (Error, throwError)-import Kafka.Consumer.Types (ConsumerRecord (..))+import Data.IORef (IORef, atomicModifyIORef', newIORef, readIORef)+import Data.Map.Strict (Map)+import Data.Map.Strict qualified as Map+import Data.Time.Clock (NominalDiffTime)+import Effectful (Eff, IOE, (:>))+import Effectful qualified+import Effectful.Error.Static (Error, catchError, throwError)+import Effectful.Exception qualified as Exception+import Kafka.Consumer.Types (ConsumerRecord (..), Offset (..), PartitionOffset (..), TopicPartition (..)) import Kafka.Effectful.Consumer.Effect ( KafkaConsumer, pausePartitions, pollMessageBatch,+ seekPartitions, storeOffsetMessage, )-import Kafka.Streamly.Stream (skipNonFatal)-import Kafka.Types (KafkaError)+import Kafka.Streamly.Stream (isFatal, skipNonFatal)+import Kafka.Types (KafkaError, PartitionId, Timeout (..), TopicName) import Shibuya.Adapter.Kafka.Config (KafkaAdapterConfig (..)) import Shibuya.Adapter.Kafka.Convert (consumerRecordToEnvelope)-import Shibuya.Core.Ack (AckDecision (..))+import Shibuya.Core.Ack (AckDecision (..), RetryDelay (..)) import Shibuya.Core.AckHandle (AckHandle (..))-import Shibuya.Core.Ingested (Ingested (..))+import Shibuya.Core.Ingested (Ingested)+import Shibuya.Core.Ingested qualified as Core import Streamly.Data.Stream (Stream) import Streamly.Data.Stream qualified as Stream+import System.IO (hPutStrLn, stderr) +type PartitionKey = (TopicName, PartitionId)++-- | Mutable state shared by the source stream and ack handles.+data KafkaAdapterState = KafkaAdapterState+ { shutdownVar :: !(TVar Bool)+ , seekBarrier :: !(IORef (Map PartitionKey Offset))+ , fatalError :: !(IORef (Maybe KafkaError))+ , consumerLock :: !(MVar ())+ {- ^ Serializes every librdkafka consumer operation. Under+ 'Shibuya.App.runApp' the consumer handle is shared between the ingester+ thread (which polls) and the processor thread (which seeks, stores,+ pauses, and commits during finalize). Running @rd_kafka_consume_batch_queue@+ concurrently with a seek/store on the same handle corrupts librdkafka's+ internal fetch queue and crashes with a native SIGSEGV, so every consumer+ call is wrapped in this mutex. Each poll is additionally bounded (see+ @maxPollHoldMillis@) so the lock is released frequently and finalize is+ never starved.+ -}+ }++{- | Allocate mutable state shared by the Kafka source, ack handles, and+optional rebalance callback.+-}+newKafkaAdapterState :: IO KafkaAdapterState+newKafkaAdapterState =+ KafkaAdapterState+ <$> newTVarIO False+ <*> newIORef Map.empty+ <*> newIORef Nothing+ <*> newMVar ()++{- | Run a librdkafka consumer operation while holding the shared consumer+lock, guaranteeing no other consumer call runs concurrently on the same+handle. See 'consumerLock' for why this is mandatory. The lock is always+released, including when the action throws through the 'Error' @KafkaError@+effect.+-}+withConsumerLock :: (IOE :> es) => KafkaAdapterState -> Eff es a -> Eff es a+withConsumerLock state =+ Exception.bracket_+ (Effectful.liftIO (takeMVar state.consumerLock))+ (Effectful.liftIO (putMVar state.consumerLock ()))++{- | Upper bound (milliseconds) on how long any single blocking consumer call+may hold the consumer lock. Every consumer call runs under 'consumerLock' so+the poll never races a concurrent seek/store on the shared handle (see+'consumerLock'). A call that blocks for its whole timeout — an empty poll, or a+seek that cannot complete — would then hold the lock for that long and starve+(and, under GHC's deadlock detector, wedge) the finalize path that must seek to+redeliver a retried message. Capping the poll and seek timeouts at this value+keeps the lock available roughly every @maxPollHoldMillis@ so finalize can+interleave and the ingester promptly re-polls to observe redelivered records.+It also bounds shutdown latency. The cap leaves ample headroom over a healthy+seek (which acknowledges in single-digit milliseconds), so it only bites when a+call is genuinely stuck — in which case surfacing it fast (via the ack path's+bounded retry and fatal slot) is the desired behavior.+-}+maxPollHoldMillis :: Int+maxPollHoldMillis = 100++{- | Cap a caller-supplied timeout at 'maxPollHoldMillis' so a single blocking+consumer call cannot hold the consumer lock longer than that. Applied to both+the poll and the retry seek.+-}+boundedLockTimeout :: Timeout -> Timeout+boundedLockTimeout t = Timeout (min (unTimeout t) maxPollHoldMillis)+ {- | Create a stream of 'ConsumerRecord's by repeatedly polling the broker. -Calls 'pollMessageBatch' in a loop, preserving errors as @Left@ values.-Non-fatal errors (timeouts, partition EOF, etc.) are filtered out via-'skipNonFatal' from hw-kafka-streamly. Fatal errors are preserved for-upstream handling.+Calls 'pollMessageBatch' in a loop under 'consumerLock', preserving errors as+@Left@ values. Each poll uses a timeout capped at @maxPollHoldMillis@ so the+lock stays available to the finalize path (which seeks, stores, pauses, and+commits). Non-fatal errors (timeouts, partition EOF, etc.) are filtered out+via 'skipNonFatal' from hw-kafka-streamly. Fatal errors are preserved for+upstream handling; a fatal error recorded by the ack path (see 'fatalError')+terminates the stream. -} kafkaSource ::- (KafkaConsumer :> es) =>+ (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) =>+ KafkaAdapterState -> KafkaAdapterConfig -> Stream (Eff es) (Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString)))-kafkaSource config =+kafkaSource state config = skipNonFatal $- Stream.repeatM pollBatch+ Stream.unfoldrM step () & Stream.concatMap Stream.fromList where- pollBatch =- pollMessageBatch config.pollTimeout config.batchSize+ pollT = boundedLockTimeout config.pollTimeout+ step () = do+ mbFatal <- Effectful.liftIO $ readIORef state.fatalError+ case mbFatal of+ Just err -> throwError err+ Nothing -> pure ()+ isShutdown <- Effectful.liftIO $ readTVarIO state.shutdownVar+ if isShutdown+ then pure Nothing+ else do+ batch <- withConsumerLock state (pollMessageBatch pollT config.batchSize)+ pure (Just (batch, ())) +-- | Drop records already buffered above a pending retry barrier.+dropStaleRecords ::+ (IOE :> es) =>+ KafkaAdapterState ->+ Stream (Eff es) (Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString))) ->+ Stream (Eff es) (Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString)))+dropStaleRecords state =+ Stream.filterM $ \case+ Left _ -> pure True+ Right cr -> do+ barriers <- Effectful.liftIO $ readIORef state.seekBarrier+ pure $ case Map.lookup (partitionKey cr) barriers of+ Nothing -> True+ Just barrierOff -> cr.crOffset <= barrierOff+ {- | Create an 'AckHandle' for a single 'ConsumerRecord'. Maps 'AckDecision' to Kafka operations: * 'AckOk' -> 'storeOffsetMessage' (mark offset ready for commit)-* 'AckRetry' -> 'storeOffsetMessage' (Kafka cannot un-read; see Decision Log)-* 'AckDeadLetter' -> 'storeOffsetMessage' (DLQ deferred to future milestone)+* 'AckRetry' -> record seek barrier and seek partition back to the failed offset+* 'AckDeadLetter' -> warn to stderr, then 'storeOffsetMessage' (DLQ deferred) * 'AckHalt' -> 'pausePartitions' (do NOT store offset; message will be re-consumed) -} mkAckHandle ::- (KafkaConsumer :> es) =>+ (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) =>+ KafkaAdapterState ->+ KafkaAdapterConfig -> ConsumerRecord (Maybe ByteString) (Maybe ByteString) -> AckHandle es-mkAckHandle cr = AckHandle $ \case+mkAckHandle state config cr = AckHandle $ \case AckOk ->- storeOffsetMessage cr- AckRetry _ ->- storeOffsetMessage cr- AckDeadLetter _ ->- storeOffsetMessage cr+ ackAttempt state (storeGuarded state cr)+ AckRetry (RetryDelay delay) -> do+ Effectful.liftIO $ delayRetry delay+ Effectful.liftIO $+ atomicModifyIORef' state.seekBarrier $ \barriers ->+ (Map.insert (partitionKey cr) cr.crOffset barriers, ())+ ackAttempt state $+ withConsumerLock state $+ seekPartitions+ [ TopicPartition+ { tpTopicName = cr.crTopic+ , tpPartition = cr.crPartition+ , tpOffset = PartitionOffset (unOffset cr.crOffset)+ }+ ]+ (boundedLockTimeout config.pollTimeout)+ AckDeadLetter reason -> do+ Effectful.liftIO $+ hPutStrLn stderr $+ "[shibuya-kafka-adapter] WARNING: dead-lettered message DROPPED (no DLQ producer): "+ <> show (cr.crTopic, cr.crPartition, cr.crOffset)+ <> " reason="+ <> show reason+ ackAttempt state (storeGuarded state cr) AckHalt _ ->- pausePartitions [(cr.crTopic, cr.crPartition)]+ ackAttempt state (withConsumerLock state (pausePartitions [(cr.crTopic, cr.crPartition)])) +ackAttempt ::+ (Error KafkaError :> es, IOE :> es) =>+ KafkaAdapterState ->+ Eff es () ->+ Eff es ()+ackAttempt state action = go (1 :: Int)+ where+ maxAttempts = 3+ retryDelayMicros = 50000++ go attempt =+ action `catchError` \_ err ->+ if isFatal err || attempt >= maxAttempts+ then Effectful.liftIO $ recordFatalError state err+ else do+ Effectful.liftIO $ threadDelay retryDelayMicros+ go (attempt + 1)++recordFatalError :: KafkaAdapterState -> KafkaError -> IO ()+recordFatalError state err =+ atomicModifyIORef' state.fatalError $ \case+ Just existing -> (Just existing, ())+ Nothing -> (Just err, ())++storeGuarded ::+ (KafkaConsumer :> es, IOE :> es) =>+ KafkaAdapterState ->+ ConsumerRecord (Maybe ByteString) (Maybe ByteString) ->+ Eff es ()+storeGuarded state cr = do+ shouldStore <-+ Effectful.liftIO $+ atomicModifyIORef' state.seekBarrier $ \barriers ->+ case Map.lookup (partitionKey cr) barriers of+ Nothing -> (barriers, True)+ Just barrierOff+ | cr.crOffset <= barrierOff -> (Map.delete (partitionKey cr) barriers, True)+ | otherwise -> (barriers, False)+ if shouldStore then withConsumerLock state (storeOffsetMessage cr) else pure ()++partitionKey :: ConsumerRecord k v -> PartitionKey+partitionKey cr = (cr.crTopic, cr.crPartition)++delayRetry :: NominalDiffTime -> IO ()+delayRetry delay+ | delay <= 0 = pure ()+ | otherwise = threadDelay (floor (realToFrac delay * (1000000 :: Double)))+ {- | Combine conversion and ack handle to produce an 'Ingested'. Lease is always 'Nothing' for Kafka (no visibility timeout mechanism). -} mkIngested ::- (KafkaConsumer :> es) =>+ (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) =>+ KafkaAdapterState ->+ KafkaAdapterConfig -> ConsumerRecord (Maybe ByteString) (Maybe ByteString) -> Ingested es (Maybe ByteString)-mkIngested cr =- Ingested- { envelope = consumerRecordToEnvelope cr- , ack = mkAckHandle cr- , lease = Nothing- }+mkIngested state config cr =+ Core.mkIngested+ (consumerRecordToEnvelope cr)+ (mkAckHandle state config cr) {- | Transform a poll stream of @Either KafkaError ConsumerRecord@ into a stream of 'Ingested'.
test/Kafka/TestEnv.hs view
@@ -76,7 +76,7 @@ prefix <- randomPrefix let env = TestEnv- { testBroker = BrokerAddress "localhost:9092"+ { testBroker = BrokerAddress "127.0.0.1:9092" , testTopic = TopicName (Text.pack (prefix <> "-topic")) , testGroupId = ConsumerGroupId (Text.pack (prefix <> "-group")) , testPrefix = prefix@@ -174,7 +174,6 @@ { topics = [env.testTopic] , pollTimeout = Timeout 5000 , batchSize = BatchSize 100- , offsetReset = Earliest } Adapter{source} <- kafkaAdapter config Stream.fold Fold.drain
test/Main.hs view
@@ -1,5 +1,6 @@ module Main (main) where +import Shibuya.Adapter.Kafka.AckHandleTest qualified as AckHandleTest import Shibuya.Adapter.Kafka.AdapterTest qualified as AdapterTest import Shibuya.Adapter.Kafka.ConvertTest qualified as ConvertTest import Shibuya.Adapter.Kafka.IntegrationTest qualified as IntegrationTest@@ -10,7 +11,8 @@ defaultMain $ testGroup "shibuya-kafka-adapter"- [ AdapterTest.tests+ [ AckHandleTest.tests+ , AdapterTest.tests , ConvertTest.tests , IntegrationTest.tests ]
+ test/Shibuya/Adapter/Kafka/AckHandleTest.hs view
@@ -0,0 +1,246 @@+module Shibuya.Adapter.Kafka.AckHandleTest (tests) where++import Data.ByteString (ByteString)+import Data.IORef (IORef, atomicModifyIORef', atomicWriteIORef, newIORef, readIORef)+import Data.Int (Int64)+import Effectful (Eff, IOE, liftIO, runEff, (:>))+import Effectful.Dispatch.Dynamic (interpret)+import Effectful.Error.Static (Error, runErrorNoCallStack, throwError)+import Kafka.Consumer (RdKafkaRespErrT (..))+import Kafka.Consumer.Types (ConsumerRecord (..), Offset (..), PartitionOffset (..), Timestamp (..), TopicPartition (..))+import Kafka.Effectful.Consumer.Effect (KafkaConsumer (..))+import Kafka.Types (BatchSize (..), KafkaError (..), PartitionId (..), Timeout (..), TopicName (..))+import Shibuya.Adapter.Kafka.Config (KafkaAdapterConfig (..))+import Shibuya.Adapter.Kafka.Internal (KafkaAdapterState (..), ingestedStream, kafkaSource, mkAckHandle, newKafkaAdapterState)+import Shibuya.Core.Ack (AckDecision (..), HaltReason (..), RetryDelay (..))+import Shibuya.Core.AckHandle (AckHandle (..))+import Shibuya.Core.Ingested (Ingested)+import Streamly.Data.Fold qualified as Fold+import Streamly.Data.Stream qualified as Stream+import Test.Tasty (TestTree, testGroup)+import Test.Tasty.HUnit (assertEqual, assertFailure, testCase)++data MockState = MockState+ { storeAttempts :: !Int+ , pauseAttempts :: !Int+ , seekCalls :: ![TopicPartition]+ , 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 without throwing" testPersistentStoreFailure+ , testCase "fatal store failure records fatal slot after one attempt" testFatalStoreFailure+ , testCase "AckHalt pause failure does not throw and records fatal slot" testAckHaltPauseFailure+ , testCase "AckRetry seeks exact failed offset and does not store" testAckRetrySeeks+ , testCase "seek barrier prevents stale successor store" testBarrierSkipsSuccessorStore+ , 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+ 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" (Just err) fatal++testFatalStoreFailure :: IO ()+testFatalStoreFailure = do+ let err = KafkaBadConfiguration+ mock <- newIORef defaultMockState{storeFailuresRemaining = 99, storeError = err}+ state <- newKafkaAdapterState+ result <- runFinalizer mock $ finalizeRecord state (recordAt 42) AckOk+ assertRight result+ 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+ result <- runFinalizer mock $ finalizeRecord state (recordAt 42) (AckHalt (HaltFatal "stop"))+ assertRight result+ 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++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++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 _ -> attemptStore mock+ PausePartitions _ -> attemptPause mock+ SeekPartitions tps _ -> recordSeek mock tps+ PollMessage _ -> error "AckHandleTest: PollMessage not exercised"+ PollMessageBatch _ _ -> error "AckHandleTest: PollMessageBatch 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) ->+ Ingested es (Maybe ByteString)+unreachableBuilder _ = error "AckHandleTest: source should not yield records"++attemptStore :: (IOE :> es, Error KafkaError :> es) => IORef MockState -> Eff es ()+attemptStore mock = do+ mbErr <-+ liftIO $+ atomicModifyIORef' mock $ \s ->+ let remaining = s.storeFailuresRemaining+ s' = s{storeAttempts = s.storeAttempts + 1, storeFailuresRemaining = max 0 (remaining - 1)}+ in (s', if remaining > 0 then Just s.storeError else Nothing)+ 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] -> Eff es ()+recordSeek mock tps = do+ mbErr <-+ liftIO $+ atomicModifyIORef' mock $ \s ->+ let remaining = s.seekFailuresRemaining+ s' = s{seekCalls = s.seekCalls <> tps, 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 =+ let AckHandle finalize = mkAckHandle state testConfig cr+ in 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+ , pauseAttempts = 0+ , seekCalls = []+ , 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 ()
test/Shibuya/Adapter/Kafka/ConvertTest.hs view
@@ -16,6 +16,7 @@ import Shibuya.Adapter.Kafka.Convert ( consumerRecordToEnvelope, extractTraceHeaders,+ extractTraceHeadersFromList, timestampToUTCTime, ) import Shibuya.Core.Types (Cursor (..), Envelope (..), MessageId (..))@@ -135,6 +136,12 @@ assertEqual "trace" Nothing (extractTraceHeaders hdrs) , testCase "returns Nothing for empty headers" $ do assertEqual "trace" Nothing (extractTraceHeaders mempty)+ , testCase "extracts from materialized header list" $ do+ let raw = [("traceparent", "00-list-parent-01"), ("tracestate", "vendor=list")]+ assertEqual+ "trace"+ (Just [("traceparent", "00-list-parent-01"), ("tracestate", "vendor=list")])+ (extractTraceHeadersFromList raw) ] timestampTests :: [TestTree]
test/Shibuya/Adapter/Kafka/IntegrationTest.hs view
@@ -1,16 +1,17 @@ module Shibuya.Adapter.Kafka.IntegrationTest (tests) where +import Control.Exception (throwIO) import Control.Monad (forM) import Control.Monad.IO.Class (liftIO) import Data.ByteString (ByteString) import Data.ByteString.Char8 qualified as BS8-import Data.IORef (modifyIORef', newIORef, readIORef)+import Data.IORef (modifyIORef', newIORef, readIORef, writeIORef) import Data.List (nub, sort) import Data.Maybe (mapMaybe) import Data.Text qualified as Text import Effectful (runEff) import Effectful.Error.Static (runError)-import Kafka.Consumer.Types (OffsetReset (..))+import Kafka.Consumer.Types (OffsetCommit (..), OffsetReset (..)) import Kafka.Effectful.Consumer ( brokersList, groupId,@@ -19,7 +20,7 @@ runKafkaConsumer, topics, )-import Kafka.Effectful.Consumer.Effect (pollMessageBatch)+import Kafka.Effectful.Consumer.Effect (commitAllOffsets, pollMessageBatch) import Kafka.TestEnv ( TestEnv (..), consumeN,@@ -37,14 +38,17 @@ ) import Shibuya.Adapter (Adapter (..)) import Shibuya.Adapter.Kafka (KafkaAdapterConfig (..), kafkaAdapter)-import Shibuya.Core.Ack (AckDecision (..))+import Shibuya.App (ProcessorId (..), defaultAppConfig, mkProcessor, runApp, waitApp)+import Shibuya.Core.Ack (AckDecision (..), RetryDelay (..)) import Shibuya.Core.AckHandle (AckHandle (..))-import Shibuya.Core.Ingested (Ingested (..))+import Shibuya.Core.Ingested (Ingested (..), Message (..)) import Shibuya.Core.Types (Cursor (..), Envelope (..), MessageId (..))+import Shibuya.Telemetry.Effect (runTracingNoop) import Streamly.Data.Fold qualified as Fold import Streamly.Data.Stream qualified as Stream+import System.Timeout (timeout) import Test.Tasty (TestTree, testGroup)-import Test.Tasty.HUnit (assertBool, assertEqual, testCase)+import Test.Tasty.HUnit (assertBool, assertEqual, assertFailure, testCase) tests :: TestTree tests =@@ -55,6 +59,10 @@ , testCase "Multi-partition distribution" testMultiPartition , testCase "Batch polling" testBatchPolling , testCase "Graceful shutdown" testGracefulShutdown+ , testCase "Idle graceful shutdown completes promptly" testIdleGracefulShutdown+ , testCase "AckRetry redelivers within the same session" testAckRetryRedelivery+ , testCase "AckRetry is not committed past when session exits" testAckRetryAbandonedSession+ , testCase "Handler exception redelivers instead of skipping" testHandlerExceptionRedelivery ] testBasicProduceConsume :: IO ()@@ -159,7 +167,6 @@ { topics = [env.testTopic] , pollTimeout = Timeout 5000 , batchSize = BatchSize 100- , offsetReset = Earliest } Adapter{source, shutdown} <- kafkaAdapter config Stream.fold Fold.drain@@ -175,3 +182,162 @@ Right () -> pure () envelopes <- reverse <$> readIORef ref assertEqual "consumed 3 before shutdown" 3 (length envelopes)++testIdleGracefulShutdown :: IO ()+testIdleGracefulShutdown = withTestEnv $ \env -> do+ createTopic env++ timedResult <- timeout 3000000 $ runEff . runError @KafkaError $ do+ let props = brokersList [env.testBroker] <> groupId env.testGroupId <> noAutoOffsetStore+ sub = topics [env.testTopic] <> offsetReset Earliest+ runKafkaConsumer props sub $ do+ let config =+ KafkaAdapterConfig+ { topics = [env.testTopic]+ , pollTimeout = Timeout 250+ , batchSize = BatchSize 100+ }+ Adapter{source, shutdown} <- kafkaAdapter config+ shutdown+ Stream.fold Fold.drain source+ case timedResult of+ Nothing -> assertFailure "idle shutdown did not terminate promptly"+ Just (Left (_cs, err)) -> assertFailure $ "idle shutdown failed: " <> show err+ Just (Right ()) -> pure ()++testAckRetryRedelivery :: IO ()+testAckRetryRedelivery = withTestEnv $ \env -> do+ createTopic env+ let payloads = ["r-1", "r-2", "r-3"]+ produceMessages env payloads++ retried <- newIORef False+ seen <- newIORef ([] :: [ByteString])+ result <- runEff . runError @KafkaError $ do+ let props = brokersList [env.testBroker] <> groupId env.testGroupId <> noAutoOffsetStore+ sub = topics [env.testTopic] <> offsetReset Earliest+ runKafkaConsumer props sub $ do+ Adapter{source} <- kafkaAdapter (testConfig env)+ Stream.fold Fold.drain+ $ Stream.mapM+ ( \(Ingested{envelope, ack = AckHandle finalize}) -> do+ let payload = maybe "" id envelope.payload+ liftIO $ modifyIORef' seen (payload :)+ hasRetried <- liftIO $ readIORef retried+ if payload == "r-2" && not hasRetried+ then do+ liftIO $ writeIORef retried True+ finalize (AckRetry (RetryDelay 0))+ else finalize AckOk+ )+ $ Stream.take 4 source+ commitAllOffsets OffsetCommit+ case result of+ Left (_cs, err) -> assertFailure $ "AckRetry redelivery failed: " <> show err+ Right () -> pure ()++ delivered <- reverse <$> readIORef seen+ assertBool ("expected r-2 at least twice, saw " <> show delivered) (countPayload "r-2" delivered >= 2)+ assertBool ("expected r-3 after retry, saw " <> show delivered) ("r-3" `elem` delivered)++ noRedelivery <- runEff . runError @KafkaError $ do+ let props = brokersList [env.testBroker] <> groupId env.testGroupId <> noAutoOffsetStore+ sub = topics [env.testTopic] <> offsetReset Earliest+ runKafkaConsumer props sub $ do+ batches <- forM [1 .. 3 :: Int] $ \_ ->+ pollMessageBatch (Timeout 500) (BatchSize 100)+ liftIO $ assertEqual "no redelivery after final AckOk commit" 0 (length [cr | Right cr <- concat batches])+ case noRedelivery of+ Left (_cs, err) -> assertFailure $ "post-commit verification failed: " <> show err+ Right () -> pure ()++testAckRetryAbandonedSession :: IO ()+testAckRetryAbandonedSession = withTestEnv $ \env -> do+ createTopic env+ let payloads = ["ab-1", "ab-2", "ab-3"]+ produceMessages env payloads++ firstSession <- runEff . runError @KafkaError $ do+ let props = brokersList [env.testBroker] <> groupId env.testGroupId <> noAutoOffsetStore+ sub = topics [env.testTopic] <> offsetReset Earliest+ runKafkaConsumer props sub $ do+ Adapter{source} <- kafkaAdapter (testConfig env)+ Stream.fold Fold.drain+ $ Stream.mapM+ ( \(Ingested{envelope, ack = AckHandle finalize}) -> do+ case envelope.payload of+ Just "ab-2" -> finalize (AckRetry (RetryDelay 0))+ _ -> finalize AckOk+ )+ $ Stream.take 2 source+ case firstSession of+ Left (_cs, err) -> assertFailure $ "first session failed: " <> show err+ Right () -> pure ()++ redelivered <- newIORef ([] :: [ByteString])+ secondSession <- runEff . runError @KafkaError $ do+ let props = brokersList [env.testBroker] <> groupId env.testGroupId <> noAutoOffsetStore+ sub = topics [env.testTopic] <> offsetReset Earliest+ runKafkaConsumer props sub $ do+ Adapter{source} <- kafkaAdapter (testConfig env)+ Stream.fold Fold.drain+ $ Stream.mapM+ ( \(Ingested{envelope, ack = AckHandle finalize}) -> do+ maybe (pure ()) (liftIO . modifyIORef' redelivered . (:)) envelope.payload+ finalize AckOk+ )+ $ Stream.take 2 source+ commitAllOffsets OffsetCommit+ case secondSession of+ Left (_cs, err) -> assertFailure $ "second session failed: " <> show err+ Right () -> pure ()++ delivered <- readIORef redelivered+ assertBool ("expected ab-2 redelivery, saw " <> show delivered) ("ab-2" `elem` delivered)++testHandlerExceptionRedelivery :: IO ()+testHandlerExceptionRedelivery = withTestEnv $ \env -> do+ createTopic env+ let payloads = ["n-1", "n-2", "n-3"]+ produceMessages env payloads++ thrown <- newIORef False+ seen <- newIORef ([] :: [ByteString])+ result <- runEff . runError @KafkaError . runTracingNoop $ do+ let props = brokersList [env.testBroker] <> groupId env.testGroupId <> noAutoOffsetStore+ sub = topics [env.testTopic] <> offsetReset Earliest+ runKafkaConsumer props sub $ do+ upstream <- kafkaAdapter (testConfig env)+ let finiteAdapter = upstream{source = Stream.take 4 upstream.source}+ handler Message{envelope} = do+ let payload = maybe "" id envelope.payload+ liftIO $ modifyIORef' seen (payload :)+ hasThrown <- liftIO $ readIORef thrown+ if payload == "n-2" && not hasThrown+ then liftIO $ do+ writeIORef thrown True+ throwIO (userError "planned handler exception")+ else pure AckOk+ appResult <- runApp defaultAppConfig [(ProcessorId "handler-exception", mkProcessor finiteAdapter handler)]+ case appResult of+ Left appErr -> liftIO $ assertFailure $ "runApp failed: " <> show appErr+ Right appHandle -> waitApp appHandle+ commitAllOffsets OffsetCommit+ case result of+ Left (_cs, err) -> assertFailure $ "handler exception scenario failed: " <> show err+ Right () -> pure ()++ delivered <- reverse <$> readIORef seen+ assertBool ("expected n-2 at least twice, saw " <> show delivered) (countPayload "n-2" delivered >= 2)+ assertBool ("expected n-3 after handler exception retry, saw " <> show delivered) ("n-3" `elem` delivered)++testConfig :: TestEnv -> KafkaAdapterConfig+testConfig env =+ KafkaAdapterConfig+ { topics = [env.testTopic]+ , pollTimeout = Timeout 500+ , batchSize = BatchSize 100+ }++countPayload :: ByteString -> [ByteString] -> Int+countPayload target = length . filter (== target)