shibuya-kafka-adapter (empty) → 0.1.0.0
raw patch · 13 files changed
+1159/−0 lines, 13 filesdep +basedep +bytestringdep +containers
Dependencies added: base, bytestring, containers, effectful-core, hw-kafka-client, hw-kafka-streamly, kafka-effectful, process, random, shibuya-core, shibuya-kafka-adapter, stm, streamly, streamly-core, tasty, tasty-hunit, text, time
Files
- CHANGELOG.md +35/−0
- LICENSE +21/−0
- README.md +42/−0
- shibuya-kafka-adapter.cabal +113/−0
- src/Shibuya/Adapter/Kafka.hs +135/−0
- src/Shibuya/Adapter/Kafka/Config.hs +48/−0
- src/Shibuya/Adapter/Kafka/Convert.hs +86/−0
- src/Shibuya/Adapter/Kafka/Internal.hs +115/−0
- test/Kafka/TestEnv.hs +191/−0
- test/Main.hs +16/−0
- test/Shibuya/Adapter/Kafka/AdapterTest.hs +60/−0
- test/Shibuya/Adapter/Kafka/ConvertTest.hs +120/−0
- test/Shibuya/Adapter/Kafka/IntegrationTest.hs +177/−0
+ CHANGELOG.md view
@@ -0,0 +1,35 @@+# Changelog++## 0.1.0.0 — 2026-04-18++Initial release.++`shibuya-kafka-adapter` bridges Apache Kafka to the+[Shibuya](https://github.com/shinzui/shibuya) queue-processing framework. It+builds on [`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).++### Features++- Poll-driven consumer that produces Shibuya `Envelope` values from Kafka+ `ConsumerRecord`s.+- Offset-commit semantics combining `noAutoOffsetStore`, explicit+ `storeOffsetMessage` on successful acknowledgement, and librdkafka+ auto-commit of the stored offsets.+- Partition-aware dispatch: `Envelope`s carry topic/partition/offset, and+ `AckHalt` pauses the originating partition via the consumer effect.+- W3C `traceparent` / `tracestate` header extraction from Kafka message+ headers, surfaced on the Shibuya `Envelope` for OpenTelemetry propagation.+- Kafka message timestamp conversion to Shibuya's `UTCTime` representation.+- Graceful shutdown that calls `commitAllOffsets` so stored offsets are+ flushed before the consumer handle closes.++### Known Limitations++- No automatic partition resume after `AckHalt` within the consumer session;+ resumption is left to the operator or the next rebalance.+- No dead-letter queue production. `AckDead` stores the offset so the stream+ advances past the poison message but does not publish it anywhere.
+ LICENSE view
@@ -0,0 +1,21 @@+MIT License++Copyright (c) 2026 Nadeem Bitar++Permission is hereby granted, free of charge, to any person obtaining a copy+of this software and associated documentation files (the "Software"), to deal+in the Software without restriction, including without limitation the rights+to use, copy, modify, merge, publish, distribute, sublicense, and/or sell+copies of the Software, and to permit persons to whom the Software is+furnished to do so, subject to the following conditions:++The above copyright notice and this permission notice shall be included in all+copies or substantial portions of the Software.++THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR+IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,+FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE+AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER+LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,+OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE+SOFTWARE.
+ README.md view
@@ -0,0 +1,42 @@+# shibuya-kafka-adapter++Kafka adapter for the [Shibuya](https://github.com/shinzui/shibuya) queue-processing framework.++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.++## Packages++- `shibuya-kafka-adapter` — the adapter library (`Shibuya.Adapter.Kafka`, `.Config`, `.Convert`).+- `shibuya-kafka-adapter-bench` — micro-benchmarks for the conversion hot path (`ConsumerRecord` → `Envelope`, W3C header extraction, timestamps).+- `shibuya-kafka-adapter-jitsurei` — runnable examples: `BasicConsumer`, `MultiTopic`, `MultiPartition`, `OffsetManagement`.++## Building++The repo ships a Nix flake and `direnv` config for a reproducible toolchain.++```sh+direnv allow # or: nix develop+cabal build all+cabal test shibuya-kafka-adapter+```++Benchmarks and examples:++```sh+cabal bench shibuya-kafka-adapter-bench+cabal run BasicConsumer+```++## Layout++```+shibuya-kafka-adapter/ library sources and tests+shibuya-kafka-adapter-bench/ tasty-bench micro-benchmarks+shibuya-kafka-adapter-jitsurei/ runnable usage examples+docs/plans/ execution plans+mori.dhall project manifest (mori registry)+```++## License++MIT. See package cabal files for details.
+ shibuya-kafka-adapter.cabal view
@@ -0,0 +1,113 @@+cabal-version: 3.12+name: shibuya-kafka-adapter+version: 0.1.0.0+synopsis: Kafka adapter for the Shibuya queue processing framework+description:+ A Shibuya adapter that integrates with Apache Kafka via kafka-effectful+ and hw-kafka-client. Provides polling, offset commit semantics, partition+ awareness, and graceful shutdown.++author: Nadeem Bitar+copyright: 2026 Nadeem Bitar+maintainer: nadeem@gmail.com+homepage: https://github.com/shinzui/shibuya-kafka-adapter+bug-reports: https://github.com/shinzui/shibuya-kafka-adapter/issues+license: MIT+license-file: LICENSE+build-type: Simple+tested-with: GHC ==9.12.2+category: Concurrency, Streaming+extra-doc-files:+ CHANGELOG.md+ README.md++source-repository head+ type: git+ location: https://github.com/shinzui/shibuya-kafka-adapter.git++common warnings+ ghc-options:+ -Wall -Wcompat -Widentities -Wincomplete-uni-patterns+ -Wincomplete-record-updates -Wredundant-constraints+ -fhide-source-paths -Wmissing-export-lists -Wpartial-fields+ -Wmissing-deriving-strategies++library+ import: warnings+ exposed-modules:+ Shibuya.Adapter.Kafka+ Shibuya.Adapter.Kafka.Config+ Shibuya.Adapter.Kafka.Convert+ Shibuya.Adapter.Kafka.Internal++ default-extensions:+ DeriveAnyClass+ DerivingStrategies+ DuplicateRecordFields+ LambdaCase+ NoFieldSelectors+ OverloadedLabels+ OverloadedRecordDot+ OverloadedStrings+ QuasiQuotes++ build-depends:+ , base ^>=4.21.0.0+ , bytestring ^>=0.12+ , containers ^>=0.7+ , effectful-core ^>=2.6.1.0+ , hw-kafka-client >=5.3 && <6+ , hw-kafka-streamly ^>=0.1+ , kafka-effectful ^>=0.1+ , shibuya-core ^>=0.1.0.0+ , stm ^>=2.5+ , streamly ^>=0.11+ , streamly-core ^>=0.3+ , text ^>=2.1+ , time ^>=1.14++ hs-source-dirs: src+ default-language: GHC2024++test-suite shibuya-kafka-adapter-test+ import: warnings+ default-language: GHC2024+ type: exitcode-stdio-1.0+ ghc-options: -threaded -rtsopts -with-rtsopts=-N+ hs-source-dirs: test+ main-is: Main.hs+ default-extensions:+ DeriveAnyClass+ DerivingStrategies+ DuplicateRecordFields+ LambdaCase+ NoFieldSelectors+ OverloadedLabels+ OverloadedRecordDot+ OverloadedStrings+ QuasiQuotes++ other-modules:+ Kafka.TestEnv+ Shibuya.Adapter.Kafka.AdapterTest+ Shibuya.Adapter.Kafka.ConvertTest+ Shibuya.Adapter.Kafka.IntegrationTest++ build-depends:+ , base ^>=4.21.0.0+ , bytestring+ , containers+ , effectful-core+ , hw-kafka-client+ , kafka-effectful+ , process+ , random+ , shibuya-core+ , shibuya-kafka-adapter+ , stm+ , streamly+ , streamly-core+ , tasty ^>=1.5+ , tasty-hunit ^>=0.10+ , text+ , time
+ src/Shibuya/Adapter/Kafka.hs view
@@ -0,0 +1,135 @@+{- | Kafka adapter for the Shibuya queue processing framework.++This adapter integrates with Apache Kafka via+[kafka-effectful](https://github.com/shinzui/kafka-effectful) and+[hw-kafka-client](https://github.com/haskell-works/hw-kafka-client).++== Example Usage++@+import Shibuya.App (runApp, mkProcessor)+import Shibuya.Adapter.Kafka (kafkaAdapter, defaultConfig)+import Kafka.Effectful.Consumer (runKafkaConsumer)+import Kafka.Consumer (brokersList, groupId, noAutoOffsetStore)++main :: IO ()+main = runEff+ . runError \@KafkaError+ . runKafkaConsumer props sub+ $ do+ adapter <- kafkaAdapter (defaultConfig [TopicName \"orders\"])+ result <- runApp IgnoreFailures 100+ [ (ProcessorId \"orders\", mkProcessor adapter myHandler)+ ]+ ...+@++== 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++== Fatal Error Propagation++Non-fatal Kafka errors (poll timeouts, partition EOFs, and the rest of the+non-fatal set defined by @hw-kafka-streamly@'s 'Kafka.Streamly.Source.isFatal')+are filtered out of the poll stream. Any error that survives that filter is+fatal by construction (for example, an SSL handshake failure, an authentication+failure, or an invalid broker configuration) and terminates the stream by+throwing through the 'Effectful.Error.Static.Error' @KafkaError@ effect. The+caller observes the failure by receiving a @Left err@ from the+@runError \@KafkaError@ scope around 'Shibuya.App.runApp'.++== 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.+-}+module Shibuya.Adapter.Kafka (+ -- * Adapter+ kafkaAdapter,++ -- * Configuration+ KafkaAdapterConfig (..),++ -- * Defaults+ defaultConfig,++ -- * Re-exports+ TopicName (..),+ BrokerAddress (..),+ ConsumerGroupId (..),+ OffsetReset (..),+ OffsetCommit (..),+ Timeout (..),+ BatchSize (..),+ KafkaError,+)+where++import Control.Concurrent.STM (TVar, atomically, newTVarIO, readTVarIO, writeTVar)+import Control.Monad.IO.Class (liftIO)+import Data.ByteString (ByteString)+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 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++{- | Create a Kafka adapter with the given configuration.++The adapter operates within an existing 'KafkaConsumer' effect scope.+Consumer lifecycle (connection, group membership) is managed by+@runKafkaConsumer@ from kafka-effectful.++The adapter uses @noAutoOffsetStore@ with manual @storeOffsetMessage@ ++auto-commit for offset management. On shutdown, @commitAllOffsets@ flushes+any stored offsets to the broker.++The returned 'Shibuya.Adapter.Adapter.shutdown' action must be invoked while+the 'KafkaConsumer' effect is still in scope. Invoking it after+@runKafkaConsumer@ has returned will throw a 'KafkaError' from+@commitAllOffsets@ against a consumer that is no longer valid. This is a+caller-side invariant; the adapter does not catch the error.+-}+kafkaAdapter ::+ (KafkaConsumer :> es, Error KafkaError :> es, IOE :> es) =>+ KafkaAdapterConfig ->+ Eff es (Adapter es (Maybe ByteString))+kafkaAdapter config = do+ shutdownVar <- liftIO $ newTVarIO False+ let messageSource = ingestedStream mkIngested (kafkaSource config)+ pure+ Adapter+ { adapterName = "kafka:" <> Text.intercalate "," (map unTopicName config.topics)+ , source = takeUntilShutdown shutdownVar messageSource+ , shutdown = do+ liftIO $ atomically $ writeTVar shutdownVar True+ commitAllOffsets OffsetCommit+ }++-- | 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)
+ src/Shibuya/Adapter/Kafka/Config.hs view
@@ -0,0 +1,48 @@+-- | Configuration types for the Kafka adapter.+module Shibuya.Adapter.Kafka.Config (+ -- * Main Configuration+ KafkaAdapterConfig (..),++ -- * Defaults+ defaultConfig,+)+where++import GHC.Generics (Generic)+import Kafka.Consumer.Types (OffsetReset (..))+import Kafka.Types (BatchSize (..), Timeout (..), TopicName)++{- | Configuration for the Kafka adapter.++Consumer properties (brokers, group ID, etc.) are provided when running+@runKafkaConsumer@ — the adapter operates /within/ the @KafkaConsumer@+effect scope, not outside it.+-}+data KafkaAdapterConfig = KafkaAdapterConfig+ { topics :: ![TopicName]+ -- ^ Topics to consume from+ , 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)++{- | Default adapter configuration for the given topics.++Defaults:++* @pollTimeout@: 1000ms+* @batchSize@: 100+* @offsetReset@: Earliest+-}+defaultConfig :: [TopicName] -> KafkaAdapterConfig+defaultConfig ts =+ KafkaAdapterConfig+ { topics = ts+ , pollTimeout = Timeout 1000+ , batchSize = BatchSize 100+ , offsetReset = Earliest+ }
+ src/Shibuya/Adapter/Kafka/Convert.hs view
@@ -0,0 +1,86 @@+-- | Type conversions between Kafka and Shibuya types.+module Shibuya.Adapter.Kafka.Convert (+ -- * Message Conversion+ consumerRecordToEnvelope,++ -- * Trace Context+ extractTraceHeaders,++ -- * Timestamp Conversion+ timestampToUTCTime,+)+where++import Data.ByteString (ByteString)+import Data.Int (Int64)+import Data.Text qualified as Text+import Data.Time (UTCTime)+import Data.Time.Clock.POSIX (posixSecondsToUTCTime)+import Kafka.Consumer.Types (ConsumerRecord (..), Offset (..), Timestamp (..))+import Kafka.Types (+ Headers,+ Millis (..),+ PartitionId (..),+ TopicName (..),+ headersToList,+ )+import Shibuya.Core.Types (Cursor (..), Envelope (..), MessageId (..), TraceHeaders)++{- | Convert a Kafka 'ConsumerRecord' to a Shibuya 'Envelope'.++Field mapping:++* @messageId@: @\"{topic}-{partition}-{offset}\"@ (globally unique within a cluster)+* @cursor@: @CursorInt offset@+* @partition@: @Just (show partitionId)@+* @enqueuedAt@: converted from Kafka timestamp if available+* @traceContext@: extracted from @traceparent@/@tracestate@ headers+* @payload@: the @crValue@ field (@Maybe ByteString@)+-}+consumerRecordToEnvelope ::+ 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+ , payload = cr.crValue+ }++-- | Build a globally unique message ID from topic, partition, and offset.+mkMessageId :: TopicName -> PartitionId -> Offset -> MessageId+mkMessageId (TopicName topic) (PartitionId pid) (Offset off) =+ MessageId (topic <> "-" <> Text.pack (show pid) <> "-" <> Text.pack (show off))++{- | Extract W3C trace headers from Kafka message headers.++Looks for @traceparent@ and @tracestate@ header keys.+Returns 'Nothing' if @traceparent@ is not present (it's required for valid context).+-}+extractTraceHeaders :: Headers -> Maybe TraceHeaders+extractTraceHeaders headers =+ 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'.++Returns 'Nothing' for 'NoTimestamp'.+Both 'CreateTime' and 'LogAppendTime' carry milliseconds since Unix epoch.+-}+timestampToUTCTime :: Timestamp -> Maybe UTCTime+timestampToUTCTime = \case+ CreateTime (Millis ms) -> Just (millisToUTCTime ms)+ LogAppendTime (Millis ms) -> Just (millisToUTCTime ms)+ NoTimestamp -> Nothing++-- | Convert milliseconds since Unix epoch to UTCTime.+millisToUTCTime :: Int64 -> UTCTime+millisToUTCTime ms = posixSecondsToUTCTime (fromIntegral ms / 1000)
+ src/Shibuya/Adapter/Kafka/Internal.hs view
@@ -0,0 +1,115 @@+{- | Internal implementation details for the Kafka adapter.+This module is not part of the public API and may change without notice.+-}+module Shibuya.Adapter.Kafka.Internal (+ -- * Stream Construction+ kafkaSource,+ ingestedStream,++ -- * Ingested Construction+ mkIngested,++ -- * AckHandle Construction+ mkAckHandle,+)+where++import Data.ByteString (ByteString)+import Data.Function ((&))+import Effectful (Eff, (:>))+import Effectful.Error.Static (Error, throwError)+import Kafka.Consumer.Types (ConsumerRecord (..))+import Kafka.Effectful.Consumer.Effect (+ KafkaConsumer,+ pausePartitions,+ pollMessageBatch,+ storeOffsetMessage,+ )+import Kafka.Streamly.Source (skipNonFatal)+import Kafka.Types (KafkaError)+import Shibuya.Adapter.Kafka.Config (KafkaAdapterConfig (..))+import Shibuya.Adapter.Kafka.Convert (consumerRecordToEnvelope)+import Shibuya.Core.Ack (AckDecision (..))+import Shibuya.Core.AckHandle (AckHandle (..))+import Shibuya.Core.Ingested (Ingested (..))+import Streamly.Data.Stream (Stream)+import Streamly.Data.Stream qualified as Stream++{- | 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.+-}+kafkaSource ::+ (KafkaConsumer :> es) =>+ KafkaAdapterConfig ->+ Stream (Eff es) (Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString)))+kafkaSource config =+ skipNonFatal $+ Stream.repeatM pollBatch+ & Stream.concatMap Stream.fromList+ where+ pollBatch =+ pollMessageBatch config.pollTimeout config.batchSize++{- | 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)+* 'AckHalt' -> 'pausePartitions' (do NOT store offset; message will be re-consumed)+-}+mkAckHandle ::+ (KafkaConsumer :> es) =>+ ConsumerRecord (Maybe ByteString) (Maybe ByteString) ->+ AckHandle es+mkAckHandle cr = AckHandle $ \case+ AckOk ->+ storeOffsetMessage cr+ AckRetry _ ->+ storeOffsetMessage cr+ AckDeadLetter _ ->+ storeOffsetMessage cr+ AckHalt _ ->+ pausePartitions [(cr.crTopic, cr.crPartition)]++{- | Combine conversion and ack handle to produce an 'Ingested'.++Lease is always 'Nothing' for Kafka (no visibility timeout mechanism).+-}+mkIngested ::+ (KafkaConsumer :> es) =>+ ConsumerRecord (Maybe ByteString) (Maybe ByteString) ->+ Ingested es (Maybe ByteString)+mkIngested cr =+ Ingested+ { envelope = consumerRecordToEnvelope cr+ , ack = mkAckHandle cr+ , lease = Nothing+ }++{- | Transform a poll stream of @Either KafkaError ConsumerRecord@ into a+stream of 'Ingested'.++A @Right cr@ is wrapped via the supplied builder (in production,+'mkIngested'). A @Left err@ that reaches this stage is fatal by construction+— 'Kafka.Streamly.Source.skipNonFatal' has already dropped non-fatal errors+— and is thrown via the 'Error' @KafkaError@ effect, terminating the stream.++Parameterizing over the builder function keeps this helper free of the+'KafkaConsumer' constraint, so it can be exercised in a unit test that+injects a synthetic @Left@ without standing up a real consumer.+-}+ingestedStream ::+ (Error KafkaError :> es) =>+ (ConsumerRecord (Maybe ByteString) (Maybe ByteString) -> Ingested es (Maybe ByteString)) ->+ Stream (Eff es) (Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString))) ->+ Stream (Eff es) (Ingested es (Maybe ByteString))+ingestedStream mkI =+ Stream.mapMaybeM $ \case+ Right cr -> pure (Just (mkI cr))+ Left err -> throwError err
+ test/Kafka/TestEnv.hs view
@@ -0,0 +1,191 @@+-- | Test environment helpers for Kafka integration tests.+module Kafka.TestEnv (+ -- * Test Environment+ TestEnv (..),+ withTestEnv,++ -- * Producing+ produceMessages,+ produceKeyedMessages,++ -- * Consuming via Adapter+ consumeN,++ -- * Topic Management+ createTopic,+ createTopicWithPartitions,+)+where++import Control.Monad (forM_)+import Control.Monad.IO.Class (liftIO)+import Data.ByteString (ByteString)+import Data.IORef (modifyIORef', newIORef, readIORef)+import Data.Text qualified as Text+import Effectful (runEff)+import Effectful.Error.Static (runError)+import Kafka.Consumer.ConsumerProperties (ConsumerProperties)+import Kafka.Consumer.Subscription (Subscription)+import Kafka.Consumer.Types (ConsumerGroupId (..), OffsetCommit (..), OffsetReset (..))+import Kafka.Effectful.Consumer (+ brokersList,+ groupId,+ noAutoOffsetStore,+ offsetReset,+ runKafkaConsumer,+ topics,+ )+import Kafka.Effectful.Consumer.Effect (commitAllOffsets)+import Kafka.Effectful.Producer (+ flushProducer,+ produceMessage,+ runKafkaProducer,+ )+import Kafka.Effectful.Producer qualified as P+import Kafka.Producer.ProducerProperties (ProducerProperties)+import Kafka.Producer.Types (ProducePartition (..), ProducerRecord (..))+import Kafka.Types (+ BatchSize (..),+ BrokerAddress (..),+ KafkaError,+ Timeout (..),+ TopicName (..),+ )+import Shibuya.Adapter (Adapter (..))+import Shibuya.Adapter.Kafka (KafkaAdapterConfig (..), kafkaAdapter)+import Shibuya.Core.Ack (AckDecision (..))+import Shibuya.Core.AckHandle (AckHandle (..))+import Shibuya.Core.Ingested (Ingested (..))+import Shibuya.Core.Types (Envelope (..))+import Streamly.Data.Fold qualified as Fold+import Streamly.Data.Stream qualified as Stream+import System.Process (callCommand)+import System.Random (randomRIO)++-- | Test environment with isolated topic and group.+data TestEnv = TestEnv+ { testBroker :: !BrokerAddress+ , testTopic :: !TopicName+ , testGroupId :: !ConsumerGroupId+ , testPrefix :: !String+ }++-- | Create a test environment with a random prefix for isolation.+withTestEnv :: (TestEnv -> IO a) -> IO a+withTestEnv f = do+ prefix <- randomPrefix+ let env =+ TestEnv+ { testBroker = BrokerAddress "localhost:9092"+ , testTopic = TopicName (Text.pack (prefix <> "-topic"))+ , testGroupId = ConsumerGroupId (Text.pack (prefix <> "-group"))+ , testPrefix = prefix+ }+ f env++-- | Generate a random 10-character alphanumeric prefix.+randomPrefix :: IO String+randomPrefix = do+ let chars = ['a' .. 'z'] ++ ['0' .. '9']+ mapM (\_ -> do i <- randomRIO (0, length chars - 1); pure (chars !! i)) [1 .. 10 :: Int]++-- | Create a topic with 1 partition via rpk.+createTopic :: TestEnv -> IO ()+createTopic env = createTopicWithPartitions env 1++-- | Create a topic with the specified number of partitions via rpk.+createTopicWithPartitions :: TestEnv -> Int -> IO ()+createTopicWithPartitions env n =+ callCommand $+ "rpk topic create "+ <> Text.unpack (unTopicName env.testTopic)+ <> " -p "+ <> show n+ <> " 2>/dev/null || true"++-- | Build consumer properties for testing.+mkConsumerProps :: TestEnv -> ConsumerProperties+mkConsumerProps env =+ brokersList [env.testBroker]+ <> groupId env.testGroupId+ <> noAutoOffsetStore++-- | Build producer properties for testing.+mkProducerProps :: TestEnv -> ProducerProperties+mkProducerProps env =+ P.brokersList [env.testBroker]++-- | Build subscription for testing.+mkSubscription :: TestEnv -> Subscription+mkSubscription env =+ topics [env.testTopic]+ <> offsetReset Earliest++-- | Produce messages with the given payloads to the test topic.+produceMessages :: TestEnv -> [ByteString] -> IO ()+produceMessages env payloads = do+ result <- runEff . runError @KafkaError $ do+ runKafkaProducer (mkProducerProps env) $ do+ forM_ payloads $ \payload ->+ produceMessage+ ProducerRecord+ { prTopic = env.testTopic+ , prPartition = UnassignedPartition+ , prKey = Nothing+ , prValue = Just payload+ , prHeaders = mempty+ }+ flushProducer+ case result of+ Left err -> error $ "Failed to produce: " <> show err+ Right () -> pure ()++-- | Produce keyed messages to the test topic.+produceKeyedMessages :: TestEnv -> [(ByteString, ByteString)] -> IO ()+produceKeyedMessages env pairs = do+ result <- runEff . runError @KafkaError $ do+ runKafkaProducer (mkProducerProps env) $ do+ forM_ pairs $ \(key, payload) ->+ produceMessage+ ProducerRecord+ { prTopic = env.testTopic+ , prPartition = UnassignedPartition+ , prKey = Just key+ , prValue = Just payload+ , prHeaders = mempty+ }+ flushProducer+ case result of+ Left err -> error $ "Failed to produce: " <> show err+ Right () -> pure ()++-- | Consume N messages from the test topic via the adapter, applying the given ack decision.+consumeN ::+ TestEnv ->+ Int ->+ AckDecision ->+ IO [Envelope (Maybe ByteString)]+consumeN env n ackDecision = do+ ref <- newIORef ([] :: [Envelope (Maybe ByteString)])+ result <- runEff . runError @KafkaError $ do+ runKafkaConsumer (mkConsumerProps env) (mkSubscription env) $ do+ let config =+ KafkaAdapterConfig+ { topics = [env.testTopic]+ , pollTimeout = Timeout 5000+ , batchSize = BatchSize 100+ , offsetReset = Earliest+ }+ Adapter{source} <- kafkaAdapter config+ Stream.fold Fold.drain+ $ Stream.mapM+ ( \(Ingested{envelope, ack = AckHandle finalize}) -> do+ liftIO $ modifyIORef' ref (envelope :)+ finalize ackDecision+ )+ $ Stream.take n source+ commitAllOffsets OffsetCommit+ case result of+ Left err -> error $ "Failed to consume: " <> show err+ Right () -> pure ()+ reverse <$> readIORef ref
+ test/Main.hs view
@@ -0,0 +1,16 @@+module Main (main) where++import Shibuya.Adapter.Kafka.AdapterTest qualified as AdapterTest+import Shibuya.Adapter.Kafka.ConvertTest qualified as ConvertTest+import Shibuya.Adapter.Kafka.IntegrationTest qualified as IntegrationTest+import Test.Tasty (defaultMain, testGroup)++main :: IO ()+main =+ defaultMain $+ testGroup+ "shibuya-kafka-adapter"+ [ AdapterTest.tests+ , ConvertTest.tests+ , IntegrationTest.tests+ ]
+ test/Shibuya/Adapter/Kafka/AdapterTest.hs view
@@ -0,0 +1,60 @@+{- | White-box tests for the Kafka adapter that do not require a running+broker. Exercises 'ingestedStream' with synthetic inputs to assert that+fatal 'KafkaError' values propagate through the 'Error' effect.+-}+module Shibuya.Adapter.Kafka.AdapterTest (tests) where++import Data.ByteString (ByteString)+import Effectful (runEff)+import Effectful.Error.Static (runError)+import Kafka.Consumer.Types (ConsumerRecord)+import Kafka.Types (KafkaError (..))+import Shibuya.Adapter.Kafka.Internal (ingestedStream)+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)++tests :: TestTree+tests =+ testGroup+ "Adapter"+ [ testCase "fatal KafkaError surfaces via Error effect" testFatalPropagation+ , testCase "non-empty prefix of Rights then Left still surfaces fatal" testFatalAfterRights+ ]++{- | Builder that must never be evaluated in these tests. The input streams+contain only @Left@ values, so the @Right@ branch of 'ingestedStream'+never fires.+-}+unreachableBuilder ::+ ConsumerRecord (Maybe ByteString) (Maybe ByteString) ->+ Ingested es (Maybe ByteString)+unreachableBuilder _ = error "AdapterTest: Right branch should not be reached"++testFatalPropagation :: IO ()+testFatalPropagation = do+ let fatalErr = KafkaBadConfiguration+ input :: [Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString))]+ input = [Left fatalErr]+ result <- runEff . runError @KafkaError $ do+ Stream.fold Fold.drain $ ingestedStream unreachableBuilder (Stream.fromList input)+ case result of+ Left (_cs, err) -> assertEqual "propagated error" fatalErr err+ Right () -> assertFailure "expected fatal error to propagate via Error effect"++testFatalAfterRights :: IO ()+testFatalAfterRights = do+ -- `skipNonFatal` has already filtered non-fatal errors upstream by the time+ -- a stream reaches `ingestedStream`, so any Left here is fatal by+ -- construction. This case confirms the stream aborts on the first fatal+ -- Left — the second element must not be drawn.+ let fatalErr = KafkaBadConfiguration+ input :: [Either KafkaError (ConsumerRecord (Maybe ByteString) (Maybe ByteString))]+ input = [Left fatalErr, Left (error "AdapterTest: second element must not be forced")]+ result <- runEff . runError @KafkaError $ do+ Stream.fold Fold.drain $ ingestedStream unreachableBuilder (Stream.fromList input)+ case result of+ Left (_cs, err) -> assertEqual "propagated error" fatalErr err+ Right () -> assertFailure "expected fatal error to propagate via Error effect"
+ test/Shibuya/Adapter/Kafka/ConvertTest.hs view
@@ -0,0 +1,120 @@+module Shibuya.Adapter.Kafka.ConvertTest (tests) where++import Data.ByteString (ByteString)+import Data.Time.Clock.POSIX (posixSecondsToUTCTime)+import Kafka.Consumer.Types (ConsumerRecord (..), Offset (..), Timestamp (..))+import Kafka.Types (+ Headers,+ Millis (..),+ PartitionId (..),+ TopicName (..),+ headersFromList,+ )+import Shibuya.Adapter.Kafka.Convert (+ consumerRecordToEnvelope,+ extractTraceHeaders,+ timestampToUTCTime,+ )+import Shibuya.Core.Types (Cursor (..), Envelope (..), MessageId (..))+import Test.Tasty (TestTree, testGroup)+import Test.Tasty.HUnit (assertEqual, testCase)++tests :: TestTree+tests =+ testGroup+ "Convert"+ [ testGroup "consumerRecordToEnvelope" envelopeTests+ , testGroup "extractTraceHeaders" traceHeaderTests+ , testGroup "timestampToUTCTime" timestampTests+ ]++-- | A minimal ConsumerRecord for testing.+mkRecord ::+ TopicName ->+ PartitionId ->+ Offset ->+ Timestamp ->+ Headers ->+ Maybe ByteString ->+ Maybe ByteString ->+ ConsumerRecord (Maybe ByteString) (Maybe ByteString)+mkRecord topic pid offset ts hdrs key value =+ ConsumerRecord+ { crTopic = topic+ , crPartition = pid+ , crOffset = offset+ , crTimestamp = ts+ , crHeaders = hdrs+ , crKey = key+ , crValue = value+ }++envelopeTests :: [TestTree]+envelopeTests =+ [ testCase "messageId is topic-partition-offset" $ do+ let cr = mkRecord (TopicName "orders") (PartitionId 2) (Offset 42) NoTimestamp mempty Nothing (Just "hello")+ env = consumerRecordToEnvelope cr+ assertEqual "messageId" (MessageId "orders-2-42") env.messageId+ , testCase "cursor is CursorInt of offset" $ do+ let cr = mkRecord (TopicName "t") (PartitionId 0) (Offset 99) NoTimestamp mempty Nothing Nothing+ env = consumerRecordToEnvelope cr+ assertEqual "cursor" (Just (CursorInt 99)) env.cursor+ , testCase "partition is show of partitionId" $ do+ let cr = mkRecord (TopicName "t") (PartitionId 5) (Offset 0) NoTimestamp mempty Nothing Nothing+ env = consumerRecordToEnvelope cr+ assertEqual "partition" (Just "5") env.partition+ , testCase "enqueuedAt from CreateTime" $ do+ let cr = mkRecord (TopicName "t") (PartitionId 0) (Offset 0) (CreateTime (Millis 1700000000000)) mempty Nothing Nothing+ env = consumerRecordToEnvelope cr+ assertEqual "enqueuedAt" (Just (posixSecondsToUTCTime 1700000000)) env.enqueuedAt+ , testCase "enqueuedAt Nothing for NoTimestamp" $ do+ let cr = mkRecord (TopicName "t") (PartitionId 0) (Offset 0) NoTimestamp mempty Nothing Nothing+ env = consumerRecordToEnvelope cr+ assertEqual "enqueuedAt" Nothing env.enqueuedAt+ , testCase "payload is crValue" $ do+ let cr = mkRecord (TopicName "t") (PartitionId 0) (Offset 0) NoTimestamp mempty Nothing (Just "payload-data")+ env = consumerRecordToEnvelope cr+ assertEqual "payload" (Just "payload-data") env.payload+ , testCase "payload Nothing when crValue is Nothing" $ do+ let cr = mkRecord (TopicName "t") (PartitionId 0) (Offset 0) NoTimestamp mempty Nothing Nothing+ env = consumerRecordToEnvelope cr+ assertEqual "payload" Nothing env.payload+ , testCase "traceContext extracted from headers" $ do+ let hdrs = headersFromList [("traceparent", "00-abc-def-01")]+ cr = mkRecord (TopicName "t") (PartitionId 0) (Offset 0) NoTimestamp hdrs Nothing Nothing+ env = consumerRecordToEnvelope cr+ assertEqual "traceContext" (Just [("traceparent", "00-abc-def-01")]) env.traceContext+ ]++traceHeaderTests :: [TestTree]+traceHeaderTests =+ [ testCase "extracts traceparent only" $ do+ let hdrs = headersFromList [("traceparent", "00-abc-def-01")]+ assertEqual "trace" (Just [("traceparent", "00-abc-def-01")]) (extractTraceHeaders hdrs)+ , testCase "extracts traceparent and tracestate" $ do+ let hdrs = headersFromList [("traceparent", "00-abc-def-01"), ("tracestate", "vendor=opaque")]+ assertEqual+ "trace"+ (Just [("traceparent", "00-abc-def-01"), ("tracestate", "vendor=opaque")])+ (extractTraceHeaders hdrs)+ , testCase "returns Nothing when no traceparent" $ do+ let hdrs = headersFromList [("other-header", "value")]+ assertEqual "trace" Nothing (extractTraceHeaders hdrs)+ , testCase "returns Nothing for empty headers" $ do+ assertEqual "trace" Nothing (extractTraceHeaders mempty)+ ]++timestampTests :: [TestTree]+timestampTests =+ [ testCase "CreateTime converts to UTCTime" $ do+ let result = timestampToUTCTime (CreateTime (Millis 1700000000000))+ assertEqual "time" (Just (posixSecondsToUTCTime 1700000000)) result+ , testCase "LogAppendTime converts to UTCTime" $ do+ let result = timestampToUTCTime (LogAppendTime (Millis 1700000000000))+ assertEqual "time" (Just (posixSecondsToUTCTime 1700000000)) result+ , testCase "NoTimestamp returns Nothing" $ do+ assertEqual "time" Nothing (timestampToUTCTime NoTimestamp)+ , testCase "zero millis converts to epoch" $ do+ let result = timestampToUTCTime (CreateTime (Millis 0))+ assertEqual "time" (Just (posixSecondsToUTCTime 0)) result+ ]
+ test/Shibuya/Adapter/Kafka/IntegrationTest.hs view
@@ -0,0 +1,177 @@+module Shibuya.Adapter.Kafka.IntegrationTest (tests) where++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.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.Effectful.Consumer (+ brokersList,+ groupId,+ noAutoOffsetStore,+ offsetReset,+ runKafkaConsumer,+ topics,+ )+import Kafka.Effectful.Consumer.Effect (pollMessageBatch)+import Kafka.TestEnv (+ TestEnv (..),+ consumeN,+ createTopic,+ createTopicWithPartitions,+ produceKeyedMessages,+ produceMessages,+ withTestEnv,+ )+import Kafka.Types (+ BatchSize (..),+ KafkaError,+ Timeout (..),+ TopicName (..),+ )+import Shibuya.Adapter (Adapter (..))+import Shibuya.Adapter.Kafka (KafkaAdapterConfig (..), kafkaAdapter)+import Shibuya.Core.Ack (AckDecision (..))+import Shibuya.Core.AckHandle (AckHandle (..))+import Shibuya.Core.Ingested (Ingested (..))+import Shibuya.Core.Types (Cursor (..), Envelope (..), MessageId (..))+import Streamly.Data.Fold qualified as Fold+import Streamly.Data.Stream qualified as Stream+import Test.Tasty (TestTree, testGroup)+import Test.Tasty.HUnit (assertBool, assertEqual, testCase)++tests :: TestTree+tests =+ testGroup+ "Integration"+ [ testCase "Basic produce-consume" testBasicProduceConsume+ , testCase "Offset commit verification" testOffsetCommit+ , testCase "Multi-partition distribution" testMultiPartition+ , testCase "Batch polling" testBatchPolling+ , testCase "Graceful shutdown" testGracefulShutdown+ ]++testBasicProduceConsume :: IO ()+testBasicProduceConsume = withTestEnv $ \env -> do+ createTopic env+ let payloads = ["msg-1", "msg-2", "msg-3", "msg-4", "msg-5"]+ produceMessages env payloads+ envelopes <- consumeN env 5 AckOk++ -- Verify all 5 messages received+ assertEqual "message count" 5 (length envelopes)++ -- Verify payloads+ let receivedPayloads = map (\(Envelope{payload}) -> payload) envelopes+ assertEqual "payloads" (map Just payloads) receivedPayloads++ -- Verify messageId format: topic-partition-offset+ case envelopes of+ (Envelope{messageId = MessageId firstIdText} : _) ->+ assertBool "messageId contains topic" (Text.isPrefixOf (unTopicName env.testTopic) firstIdText)+ [] -> error "unreachable: already verified 5 envelopes"++ -- Verify cursor is populated+ assertBool "cursor is Just" (all (\(Envelope{cursor}) -> case cursor of Just (CursorInt _) -> True; _ -> False) envelopes)++ -- Verify partition is populated+ assertBool "partition is Just" (all (\(Envelope{partition}) -> case partition of Just _ -> True; _ -> False) envelopes)++testOffsetCommit :: IO ()+testOffsetCommit = withTestEnv $ \env -> do+ createTopic env+ let payloads = ["oc-1", "oc-2", "oc-3"]+ produceMessages env payloads++ -- Consume all 3, AckOk each (stores offsets), then commit+ _ <- consumeN env 3 AckOk++ -- Create new consumer in same group - should get no messages+ result <- runEff . runError @KafkaError $ do+ let props = brokersList [env.testBroker] <> groupId env.testGroupId <> noAutoOffsetStore+ sub = topics [env.testTopic] <> offsetReset Earliest+ runKafkaConsumer props sub $ do+ -- Poll a few times to allow group join + rebalance, verify no re-delivery+ allResults <- forM [1 .. 3 :: Int] $ \_ -> do+ results <- pollMessageBatch (Timeout 3000) (BatchSize 100)+ pure [cr | Right cr <- results]+ let totalMessages = concat allResults+ liftIO $ assertEqual "no re-delivery" 0 (length totalMessages)+ case result of+ Left err -> error $ "Failed: " <> show err+ Right () -> pure ()++testMultiPartition :: IO ()+testMultiPartition = withTestEnv $ \env -> do+ createTopicWithPartitions env 3+ let pairs =+ [ ("key-a", "msg-a")+ , ("key-b", "msg-b")+ , ("key-c", "msg-c")+ , ("key-d", "msg-d")+ , ("key-e", "msg-e")+ , ("key-f", "msg-f")+ ]+ produceKeyedMessages env pairs+ envelopes <- consumeN env 6 AckOk++ assertEqual "message count" 6 (length envelopes)++ let partitions = mapMaybe (\(Envelope{partition}) -> partition) envelopes+ assertEqual "all have partition" 6 (length partitions)++ let uniquePartitions = nub partitions+ assertBool+ ("expected multiple partitions, got: " <> show uniquePartitions)+ (length uniquePartitions >= 2)++testBatchPolling :: IO ()+testBatchPolling = withTestEnv $ \env -> do+ createTopic env+ let payloads = map (\i -> BS8.pack ("batch-" <> show i)) [1 .. 20 :: Int]+ produceMessages env payloads+ envelopes <- consumeN env 20 AckOk++ assertEqual "message count" 20 (length envelopes)++ let receivedPayloads = sort $ mapMaybe (\(Envelope{payload}) -> payload) envelopes+ assertEqual "payloads" (sort payloads) receivedPayloads++testGracefulShutdown :: IO ()+testGracefulShutdown = withTestEnv $ \env -> do+ createTopic env+ let payloads = ["sd-1", "sd-2", "sd-3"]+ produceMessages env payloads++ ref <- newIORef ([] :: [Envelope (Maybe 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+ let config =+ KafkaAdapterConfig+ { topics = [env.testTopic]+ , pollTimeout = Timeout 5000+ , batchSize = BatchSize 100+ , offsetReset = Earliest+ }+ Adapter{source, shutdown} <- kafkaAdapter config+ Stream.fold Fold.drain+ $ Stream.mapM+ ( \(Ingested{envelope, ack = AckHandle finalize}) -> do+ liftIO $ modifyIORef' ref (envelope :)+ finalize AckOk+ )+ $ Stream.take 3 source+ shutdown+ case result of+ Left err -> error $ "Failed: " <> show err+ Right () -> pure ()+ envelopes <- reverse <$> readIORef ref+ assertEqual "consumed 3 before shutdown" 3 (length envelopes)