packages feed

kafka-effectful-0.3.1.0: src/Kafka/Effectful/OpenTelemetry/Producer/Interpreter.hs

-- The 'EffectHandler' type synonym in effectful-core expands to a
-- constraint that GHC's redundant-constraint check flags on the
-- handler's signature, even though the constraint is required for
-- 'interpret' to type-check. Suppress the warning at the file level.
{-# OPTIONS_GHC -Wno-redundant-constraints #-}

-- | OpenTelemetry-traced interpreter for the 'KafkaProducer' effect.
--
-- Drop-in alternative to 'Kafka.Effectful.Producer.Interpreter.runKafkaProducer'
-- that opens a Producer-kind span around every record-sending operation
-- ('produceMessage', 'produceMessage'', 'produceMessageSync',
-- 'produceMessageBatch'), populates the span with the spec-aligned
-- @messaging.*@ attribute set, and injects the current OTel context as
-- W3C @traceparent@\/@tracestate@ headers on the record before handing
-- it off to the underlying @hw-kafka-client@ produce call.
--
-- Non-sending operations (flush, transactional begin\/commit\/abort, etc.)
-- are passed through unchanged — they do not represent message sends and
-- therefore do not get a span.
--
-- The design parallels the upstream
-- @hs-opentelemetry-instrumentation-hw-kafka-client@\'s
-- @OpenTelemetry.Instrumentation.Kafka.produceMessage@.
--
-- @since 0.2.0.0
module Kafka.Effectful.OpenTelemetry.Producer.Interpreter
  ( -- * Interpreter
    runKafkaProducerTraced,
  )
where

import Control.Concurrent.MVar qualified as Concurrent
import Data.Foldable (for_)
import Data.Text qualified as Text
import Effectful (Eff, IOE, (:>))
import Effectful qualified
import Effectful.Dispatch.Dynamic (EffectHandler, interpret)
import Effectful.Error.Static (Error, throwError)
import Effectful.Exception qualified as Exception
import Kafka.Effectful.OpenTelemetry.Propagation
  ( injectTraceContextIntoRecord,
  )
import Kafka.Effectful.OpenTelemetry.Semantic
  ( producerRecordAttributesWith,
    producerSpanName,
  )
import Kafka.Effectful.Producer.Effect (KafkaProducer (..))
import Kafka.Producer (ProducerRecord (prTopic))
import Kafka.Producer qualified as K
import Kafka.Producer.ProducerProperties (ProducerProperties)
import Kafka.Transaction qualified as K
import Kafka.Types (KafkaError)
import OpenTelemetry.Attributes.Key (unkey)
import OpenTelemetry.Context qualified as Context
import OpenTelemetry.Context.ThreadLocal (getContext)
import OpenTelemetry.SemanticConventions (error_type)
import OpenTelemetry.SemanticsConfig (getSemanticsOptions, lookupStability)
import OpenTelemetry.Trace.Core
  ( Span,
    SpanArguments (kind),
    SpanKind (Producer),
    SpanStatus (Error),
    Tracer,
    addAttribute,
    addAttributesToSpanArguments,
    defaultSpanArguments,
    inSpan'',
    setStatus,
  )

-- | Run the 'KafkaProducer' effect with OpenTelemetry tracing.
--
-- Identical in shape to 'Kafka.Effectful.Producer.Interpreter.runKafkaProducer',
-- plus an additional 'Tracer' argument used to open a Producer-kind
-- span around every record-sending operation. The current
-- 'OpenTelemetry.Context.Context' is injected as W3C trace-context
-- headers on the outgoing record before the underlying
-- @hw-kafka-client@ produce call runs, so that downstream consumers
-- can extract the context and continue the trace.
--
-- The producer handle is acquired and released via 'Exception.bracket',
-- exactly as 'runKafkaProducer' does. Errors are thrown via the
-- 'Error' effect.
--
-- @since 0.2.0.0
runKafkaProducerTraced ::
  (IOE :> es, Error KafkaError :> es) =>
  Tracer ->
  ProducerProperties ->
  Eff (KafkaProducer : es) a ->
  Eff es a
runKafkaProducerTraced tracer props action =
  Exception.bracket
    acquire
    (Effectful.liftIO . K.closeProducer)
    (\producer -> interpret (handleTracedProducer tracer producer) action)
  where
    acquire = do
      result <- Effectful.liftIO $ K.newProducer props
      case result of
        Left err -> throwError err
        Right producer -> pure producer

handleTracedProducer ::
  (IOE :> es, Error KafkaError :> es) =>
  Tracer ->
  K.KafkaProducer ->
  EffectHandler KafkaProducer es
handleTracedProducer tracer producer _env = \case
  ProduceMessage record ->
    withProducerSpan tracer record $ \span_ instrumentedRecord -> do
      mbErr <- Effectful.liftIO $ K.produceMessage producer instrumentedRecord
      for_ mbErr $ \err -> do
        recordKafkaError span_ err
        throwError err
  ProduceMessage' record cb ->
    withProducerSpan tracer record $ \span_ instrumentedRecord -> do
      res <-
        Effectful.liftIO $
          K.produceMessage' producer instrumentedRecord cb
      case res of
        Left (K.ImmediateError err) -> do
          recordKafkaError span_ err
          throwError err
        Right () -> pure ()
  ProduceMessageSync record ->
    withProducerSpan tracer record $ \span_ instrumentedRecord -> do
      var <- Effectful.liftIO Concurrent.newEmptyMVar
      res <-
        Effectful.liftIO $
          K.produceMessage' producer instrumentedRecord (Concurrent.putMVar var)
      case res of
        Left (K.ImmediateError err) -> do
          recordKafkaError span_ err
          throwError err
        Right () -> do
          Effectful.liftIO $ K.flushProducer producer
          report <- Effectful.liftIO $ Concurrent.takeMVar var
          case report of
            K.DeliverySuccess _ offset -> pure offset
            K.DeliveryFailure _ err -> do
              recordKafkaError span_ err
              throwError err
            K.NoMessageError err -> do
              recordKafkaError span_ err
              throwError err
  ProduceMessageBatch records -> do
    -- One span per record so the Producer-kind attributes
    -- (partition, key) are per-record, matching the upstream
    -- reference. Spans are opened sequentially as the list is
    -- traversed.
    results <-
      traverse
        ( \r ->
            withProducerSpan tracer r $ \span_ instrumentedRecord -> do
              mbErr <-
                Effectful.liftIO $
                  K.produceMessage producer instrumentedRecord
              for_ mbErr (recordKafkaError span_)
              pure (r, mbErr)
        )
        records
    pure [(r, err) | (r, Just err) <- results]
  FlushProducer ->
    Effectful.liftIO $ K.flushProducer producer
  InitTransactions timeout ->
    throwOnJust $ K.initTransactions producer timeout
  BeginTransaction ->
    throwOnJust $ K.beginTransaction producer
  CommitTransaction timeout ->
    Effectful.liftIO $ K.commitTransaction producer timeout
  AbortTransaction timeout ->
    throwOnJust $ K.abortTransaction producer timeout
  SendOffsetsToTransaction consumer record timeout ->
    Effectful.liftIO $
      K.commitOffsetMessageTransaction producer consumer record timeout
  AskProducerHandle ->
    pure producer
  where
    throwOnJust action' = do
      mbErr <- Effectful.liftIO action'
      for_ mbErr throwError

-- | Open a Producer-kind span around a record-sending action.
--
-- Builds the @messaging.*@ attribute set from the record, opens a
-- span named @\"send \<topic\>\"@, injects the current OTel context
-- as W3C trace-context headers onto a clone of the record, and runs
-- the supplied action with that instrumented record.
withProducerSpan ::
  (IOE :> es) =>
  Tracer ->
  ProducerRecord ->
  (Span -> ProducerRecord -> Eff es a) ->
  Eff es a
withProducerSpan tracer record action = do
  semOpts <- Effectful.liftIO $ lookupStability "messaging" <$> getSemanticsOptions
  inSpan'' tracer (producerSpanName (prTopic record)) (spanArgs semOpts) $ \newSpan -> do
    ctx <- getContext
    instrumentedRecord <-
      Effectful.liftIO $
        injectTraceContextIntoRecord (Context.insertSpan newSpan ctx) record
    action newSpan instrumentedRecord
  where
    spanArgs semOpts =
      addAttributesToSpanArguments
        (producerRecordAttributesWith semOpts record)
        defaultSpanArguments {kind = Producer}

recordKafkaError ::
  (IOE :> es) =>
  Span ->
  KafkaError ->
  Eff es ()
recordKafkaError span_ err = do
  let errText = Text.pack (show err)
  Effectful.liftIO $ do
    addAttribute span_ (unkey error_type) errText
    setStatus span_ (Error errText)