shibuya-kafka-adapter 0.1.0.0 → 0.2.0.0
raw patch · 6 files changed
+368/−15 lines, 6 filesdep +asyncdep +hs-opentelemetry-apiPVP ok
version bump matches the API change (PVP)
Dependencies added: async, hs-opentelemetry-api
API changes (from Hackage documentation)
+ Shibuya.Adapter.Kafka.Tracing: traced :: forall (es :: [Effect]) v. (Tracing :> es, IOE :> es) => TopicName -> Stream (Eff es) (Ingested es v) -> Stream (Eff es) (Ingested es v)
Files
- CHANGELOG.md +25/−0
- README.md +19/−1
- shibuya-kafka-adapter.cabal +19/−14
- src/Shibuya/Adapter/Kafka/Tracing.hs +101/−0
- test/Main.hs +2/−0
- test/Shibuya/Adapter/Kafka/TracingTest.hs +202/−0
CHANGELOG.md view
@@ -1,5 +1,30 @@ # Changelog +## 0.2.0.0 — 2026-04-18++Additive release. Adds one new exposed module, no changes to existing+public types or to `kafkaAdapter`'s signature.++### Added++- `Shibuya.Adapter.Kafka.Tracing` exposing a single stream transformer+ `traced :: (Tracing :> es, IOE :> es) => TopicName -> Stream (Eff es)+ (Ingested es v) -> Stream (Eff es) (Ingested es v)`. For each emitted+ `Ingested`, `traced` rewrites the envelope's `AckHandle` so that+ when the downstream handler calls `finalize` the call is enclosed+ in a Consumer-kind `shibuya.process.message` span parented on the+ envelope's carried W3C `traceparent` (or a root span when absent).+ The span carries the v1.27 messaging-conventions attribute set+ (`messaging.system`, `messaging.destination.name`,+ `messaging.message.id`, and — when partition is known —+ `messaging.destination.partition.id`) from+ `Shibuya.Telemetry.Semantic`.++- Explicit `hs-opentelemetry-api ^>=0.3` build-depends edge on the+ library. The package was previously in the closure transitively via+ `shibuya-core`, but cabal does not let a library import from a+ transitively-present dependency.+ ## 0.1.0.0 — 2026-04-18 Initial release.
README.md view
@@ -6,9 +6,27 @@ ## Packages -- `shibuya-kafka-adapter` — the adapter library (`Shibuya.Adapter.Kafka`, `.Config`, `.Convert`).+- `shibuya-kafka-adapter` — the adapter library (`Shibuya.Adapter.Kafka`, `.Config`, `.Convert`, `.Tracing`). - `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`.++## Tracing (opt-in)++`Shibuya.Adapter.Kafka.Tracing.traced` is an opt-in stream transformer that wraps each emitted `Ingested` so that the downstream handler's eventual `finalize` call runs inside a Consumer-kind `shibuya.process.message` OpenTelemetry span. The span inherits the envelope's W3C `traceparent` as parent (from `Envelope.traceContext`) or opens a fresh root span when no parent is present, and is populated with the v1.27 messaging-conventions attributes (`messaging.system=kafka`, `messaging.destination.name`, `messaging.message.id`, and `messaging.destination.partition.id` when the partition is known). A caller that does not import this module pays nothing — no spans are opened and the adapter's public surface is unchanged.++Typical wiring:++```haskell+import Shibuya.Adapter.Kafka (kafkaAdapter, defaultConfig)+import Shibuya.Adapter.Kafka.Tracing (traced)+import Shibuya.Telemetry.Effect (runTracing)++runTracing tracer $ do+ Adapter{source} <- kafkaAdapter (defaultConfig [TopicName "orders"])+ Stream.fold Fold.drain+ $ Stream.mapM userHandler+ $ traced (TopicName "orders") source+``` ## Building
shibuya-kafka-adapter.cabal view
@@ -1,6 +1,6 @@ cabal-version: 3.12 name: shibuya-kafka-adapter-version: 0.1.0.0+version: 0.2.0.0 synopsis: Kafka adapter for the Shibuya queue processing framework description: A Shibuya adapter that integrates with Apache Kafka via kafka-effectful@@ -39,6 +39,7 @@ Shibuya.Adapter.Kafka.Config Shibuya.Adapter.Kafka.Convert Shibuya.Adapter.Kafka.Internal+ Shibuya.Adapter.Kafka.Tracing default-extensions: DeriveAnyClass@@ -52,19 +53,20 @@ 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+ , base ^>=4.21.0.0+ , bytestring ^>=0.12+ , containers ^>=0.7+ , effectful-core ^>=2.6.1.0+ , hs-opentelemetry-api ^>=0.3+ , 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@@ -92,12 +94,15 @@ Shibuya.Adapter.Kafka.AdapterTest Shibuya.Adapter.Kafka.ConvertTest Shibuya.Adapter.Kafka.IntegrationTest+ Shibuya.Adapter.Kafka.TracingTest build-depends:+ , async ^>=2.2 , base ^>=4.21.0.0 , bytestring , containers , effectful-core+ , hs-opentelemetry-api ^>=0.3 , hw-kafka-client , kafka-effectful , process
+ src/Shibuya/Adapter/Kafka/Tracing.hs view
@@ -0,0 +1,101 @@+{- | Opt-in per-message tracing for the Shibuya Kafka adapter.++Exposes a single stream transformer, 'traced', that wraps each 'Ingested'+emitted by 'Shibuya.Adapter.Kafka.kafkaAdapter' so that when a downstream+handler eventually calls the envelope's 'AckHandle.finalize', the call is+enclosed in an OpenTelemetry Consumer-kind span named+@shibuya.process.message@ (see 'Shibuya.Telemetry.Semantic.processMessageSpanName').++The span:++* Is a child of the W3C trace context carried on the envelope, when present.+ The parent is recovered via 'Shibuya.Telemetry.Propagation.extractTraceContext'+ from @envelope.traceContext@. When absent, the span is a fresh root span.+* Carries the v1.27 messaging-conventions attribute set:+ @messaging.system@, @messaging.destination.name@, @messaging.message.id@,+ and — when a partition is known — @messaging.destination.partition.id@.++Typical usage:++@+Adapter{source} <- kafkaAdapter (defaultConfig [TopicName \"orders\"])+Stream.fold Fold.drain+ $ Stream.mapM userHandler+ $ traced (TopicName \"orders\") source+@++A caller that does not import this module pays nothing: the adapter's public+surface is unchanged and no span is opened.+-}+module Shibuya.Adapter.Kafka.Tracing (+ traced,+)+where++import Data.Text (Text)+import Effectful (Eff, IOE, (:>))+import Kafka.Types (TopicName (..))+import OpenTelemetry.Trace.Core (Span)+import Shibuya.Core.AckHandle (AckHandle (..))+import Shibuya.Core.Ingested (Ingested (..))+import Shibuya.Core.Types (Envelope (..), MessageId (..))+import Shibuya.Telemetry.Effect (+ Tracing,+ addAttribute,+ withExtractedContext,+ withSpan',+ )+import Shibuya.Telemetry.Propagation (extractTraceContext)+import Shibuya.Telemetry.Semantic (+ attrMessagingDestinationName,+ attrMessagingDestinationPartitionId,+ attrMessagingMessageId,+ attrMessagingSystem,+ consumerSpanArgs,+ processMessageSpanName,+ )+import Streamly.Data.Stream (Stream)+import Streamly.Data.Stream qualified as Stream++{- | Rewrite each 'Ingested' so that its 'AckHandle.finalize' opens a+Consumer-kind span parented on the envelope's carried trace context+(when any) and populated with the v1.27 messaging attributes.++The topic name is supplied by the caller rather than parsed from the+envelope's 'MessageId' — see the Decision Log of+@docs/plans/9-add-shibuya-kafka-tracing-module.md@ for the rationale.+-}+traced ::+ forall es v.+ (Tracing :> es, IOE :> es) =>+ TopicName ->+ Stream (Eff es) (Ingested es v) ->+ Stream (Eff es) (Ingested es v)+traced (TopicName topicName) =+ Stream.mapM $ \ing -> do+ let envelope = ing.envelope+ AckHandle finalize = ing.ack+ parentCtx = envelope.traceContext >>= extractTraceContext+ wrappedFinalize decision =+ withExtractedContext parentCtx $+ withSpan' processMessageSpanName consumerSpanArgs $+ \sp -> do+ populateAttrs sp topicName envelope+ finalize decision+ pure ing{ack = AckHandle wrappedFinalize}++-- | Set the messaging-conventions v1.27 attribute set on a span.+populateAttrs ::+ (Tracing :> es, IOE :> es) =>+ Span ->+ Text ->+ Envelope v ->+ Eff es ()+populateAttrs sp topicName envelope = do+ addAttribute sp attrMessagingSystem ("kafka" :: Text)+ addAttribute sp attrMessagingDestinationName topicName+ let MessageId msgIdText = envelope.messageId+ addAttribute sp attrMessagingMessageId msgIdText+ case envelope.partition of+ Just p -> addAttribute sp attrMessagingDestinationPartitionId p+ Nothing -> pure ()
test/Main.hs view
@@ -3,6 +3,7 @@ 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 Shibuya.Adapter.Kafka.TracingTest qualified as TracingTest import Test.Tasty (defaultMain, testGroup) main :: IO ()@@ -13,4 +14,5 @@ [ AdapterTest.tests , ConvertTest.tests , IntegrationTest.tests+ , TracingTest.tests ]
+ test/Shibuya/Adapter/Kafka/TracingTest.hs view
@@ -0,0 +1,202 @@+{- | Unit tests for 'Shibuya.Adapter.Kafka.Tracing.traced'.++Drives the 'traced' stream transformer over a single synthetic 'Ingested',+captures the resulting OTel span with an in-memory 'SpanProcessor', and+asserts span shape: trace parenting, attribute set, and that the original+'AckDecision' still flows through the wrapped 'AckHandle.finalize'.++No Kafka broker is required. The processor implementation mirrors+@OpenTelemetry.Exporter.InMemory.Span.inMemoryListExporter@ from+hs-opentelemetry-exporter-in-memory; the latter is not used directly+because its current Hackage release pins @hs-opentelemetry-api <0.3@,+while this repository resolves to 0.3.1.0.+-}+module Shibuya.Adapter.Kafka.TracingTest (tests) where++import Control.Concurrent.Async (async)+import Control.Exception (bracket)+import Data.ByteString (ByteString)+import Data.IORef (IORef, atomicModifyIORef, newIORef, readIORef, writeIORef)+import Data.Text (Text)+import Data.Text.Encoding qualified as TE+import Effectful (Eff, IOE, liftIO, runEff, (:>))+import Kafka.Types (TopicName (..))+import OpenTelemetry.Attributes (Attribute (..), PrimitiveAttribute (..), lookupAttribute)+import OpenTelemetry.Processor.Span (ShutdownResult (..), SpanProcessor (..))+import OpenTelemetry.Trace.Core (+ ImmutableSpan (..),+ SpanContext (..),+ Tracer,+ createTracerProvider,+ emptyTracerProviderOptions,+ getSpanContext,+ makeTracer,+ shutdownTracerProvider,+ tracerOptions,+ )+import OpenTelemetry.Trace.Id (Base (..), spanIdBaseEncodedText, traceIdBaseEncodedText)+import Shibuya.Adapter.Kafka.Tracing (traced)+import Shibuya.Core.Ack (AckDecision (..))+import Shibuya.Core.AckHandle (AckHandle (..))+import Shibuya.Core.Ingested (Ingested (..))+import Shibuya.Core.Types (Cursor (..), Envelope (..), MessageId (..), TraceHeaders)+import Shibuya.Telemetry.Effect (runTracing)+import Streamly.Data.Fold qualified as Fold+import Streamly.Data.Stream qualified as Stream+import Test.Tasty (TestTree, testGroup)+import Test.Tasty.HUnit (Assertion, assertEqual, assertFailure, testCase)++tests :: TestTree+tests =+ testGroup+ "Tracing"+ [ testCase "envelope traceparent becomes span parent" testParentedSpan+ , testCase "missing traceContext yields root span" testRootSpan+ , testCase "messaging attributes populated from envelope" testAttributes+ , testCase "AckDecision threads through wrapped finalize" testAckPassthrough+ ]++-- | SpanProcessor that records ended spans to an IORef.+recordingProcessor :: IO (SpanProcessor, IORef [ImmutableSpan])+recordingProcessor = do+ ref <- newIORef []+ let processor =+ SpanProcessor+ { spanProcessorOnStart = \_ _ -> pure ()+ , spanProcessorOnEnd = \spanRef -> do+ s <- readIORef spanRef+ atomicModifyIORef ref (\l -> (s : l, ()))+ , spanProcessorShutdown = async (pure ShutdownSuccess)+ , spanProcessorForceFlush = pure ()+ }+ pure (processor, ref)++withRecordingTracer :: (IORef [ImmutableSpan] -> Tracer -> IO a) -> IO a+withRecordingTracer k = do+ (processor, spansRef) <- recordingProcessor+ bracket+ (createTracerProvider [processor] emptyTracerProviderOptions)+ shutdownTracerProvider+ $ \tp ->+ k spansRef (makeTracer tp "shibuya-kafka-adapter-test" tracerOptions)++-- Known-good W3C traceparent with recoverable trace and span ids.+expectedTraceIdHex :: Text+expectedTraceIdHex = "0af7651916cd43dd8448eb211c80319c"++expectedSpanIdHex :: Text+expectedSpanIdHex = "b7ad6b7169203331"++sampleTraceparent :: ByteString+sampleTraceparent =+ TE.encodeUtf8 ("00-" <> expectedTraceIdHex <> "-" <> expectedSpanIdHex <> "-01")++mkEnvelope :: Maybe TraceHeaders -> Envelope ()+mkEnvelope tc =+ Envelope+ { messageId = MessageId "orders-2-42"+ , cursor = Just (CursorInt 42)+ , partition = Just "2"+ , enqueuedAt = Nothing+ , traceContext = tc+ , payload = ()+ }++mkIngestedFor ::+ (IOE :> es) =>+ IORef (Maybe AckDecision) ->+ Envelope () ->+ Ingested es ()+mkIngestedFor ackRef env =+ Ingested+ { envelope = env+ , ack = AckHandle $ \decision -> liftIO $ writeIORef ackRef (Just decision)+ , lease = Nothing+ }++-- | Drive 'traced' over a single envelope, calling finalize with 'AckOk'.+runOneThroughTraced ::+ Tracer ->+ Envelope () ->+ IORef (Maybe AckDecision) ->+ IO ()+runOneThroughTraced tracer env ackRef = runEff . runTracing tracer $ do+ let ing = mkIngestedFor ackRef env+ Stream.fold Fold.drain $+ Stream.mapM callFinalize $+ traced (TopicName "orders") $+ Stream.fromList [ing]+ where+ callFinalize :: Ingested es () -> Eff es ()+ callFinalize i =+ let AckHandle fin = i.ack+ in fin AckOk++testParentedSpan :: Assertion+testParentedSpan = withRecordingTracer $ \spansRef tracer -> do+ ackRef <- newIORef Nothing+ runOneThroughTraced tracer (mkEnvelope (Just [("traceparent", sampleTraceparent)])) ackRef+ spans <- readIORef spansRef+ case spans of+ [s] -> do+ assertEqual+ "span traceId matches parent"+ expectedTraceIdHex+ (traceIdBaseEncodedText Base16 s.spanContext.traceId)+ case s.spanParent of+ Just parentSpan -> do+ parentCtx <- getSpanContext parentSpan+ assertEqual+ "parent traceId"+ expectedTraceIdHex+ (traceIdBaseEncodedText Base16 parentCtx.traceId)+ assertEqual+ "parent spanId"+ expectedSpanIdHex+ (spanIdBaseEncodedText Base16 parentCtx.spanId)+ Nothing -> assertFailure "expected spanParent to be set"+ other -> assertFailure ("expected exactly one span, got " <> show (length other))++testRootSpan :: Assertion+testRootSpan = withRecordingTracer $ \spansRef tracer -> do+ ackRef <- newIORef Nothing+ runOneThroughTraced tracer (mkEnvelope Nothing) ackRef+ spans <- readIORef spansRef+ case spans of+ [s] -> case s.spanParent of+ Nothing -> pure ()+ Just _ -> assertFailure "expected root span to have no parent"+ other -> assertFailure ("expected exactly one span, got " <> show (length other))++testAttributes :: Assertion+testAttributes = withRecordingTracer $ \spansRef tracer -> do+ ackRef <- newIORef Nothing+ runOneThroughTraced tracer (mkEnvelope (Just [("traceparent", sampleTraceparent)])) ackRef+ spans <- readIORef spansRef+ case spans of+ [s] -> do+ let attrs = s.spanAttributes+ assertEqual+ "messaging.system"+ (Just (AttributeValue (TextAttribute "kafka")))+ (lookupAttribute attrs "messaging.system")+ assertEqual+ "messaging.destination.name"+ (Just (AttributeValue (TextAttribute "orders")))+ (lookupAttribute attrs "messaging.destination.name")+ assertEqual+ "messaging.message.id"+ (Just (AttributeValue (TextAttribute "orders-2-42")))+ (lookupAttribute attrs "messaging.message.id")+ assertEqual+ "messaging.destination.partition.id"+ (Just (AttributeValue (TextAttribute "2")))+ (lookupAttribute attrs "messaging.destination.partition.id")+ other -> assertFailure ("expected exactly one span, got " <> show (length other))++testAckPassthrough :: Assertion+testAckPassthrough = withRecordingTracer $ \_ tracer -> do+ ackRef <- newIORef Nothing+ runOneThroughTraced tracer (mkEnvelope Nothing) ackRef+ decision <- readIORef ackRef+ assertEqual "underlying finalize received the decision" (Just AckOk) decision