{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE ScopedTypeVariables #-}

-- | W3C trace context carried on enqueued jobs, and the producer and consumer spans
-- around them. The OpenTelemetry API is inert until an SDK installs a provider.
module Arbiter.Core.Trace
  ( -- * Job trace context
    TraceContext (..)
  , stampTraceContext

    -- * Tracer
  , Tracer
  , resolveTracer

    -- * Producer
  , currentTraceContext
  , withPublishSpan

    -- * Consumer
  , ConsumeSpan
  , ConsumeShape (..)
  , consumeSpanFor
  , withConsumeSpan
  , withJobParent
  , capturingContext

    -- * Span enrichment
    -- $enrichment
  , markSpanError
  , recordJobFailure
  , recordJobCancelled
  ) where

import Control.Applicative ((<|>))
import Control.Exception (fromException)
import Control.Monad (guard)
import Control.Monad.IO.Class (MonadIO)
import Data.ByteString qualified as BS
import Data.HashMap.Strict qualified as HM
import Data.List.NonEmpty (NonEmpty ((:|)))
import Data.List.NonEmpty qualified as NE
import Data.Maybe (mapMaybe, maybeToList)
import Data.Text (Text)
import Data.Text.Encoding (decodeUtf8, encodeUtf8)
import Data.Text.Lazy (toStrict)
import Data.Text.Lazy.Builder (toLazyText)
import Data.Text.Lazy.Builder.Int (decimal)
import OpenTelemetry.Attributes (Attribute, toAttribute)
import OpenTelemetry.Attributes.Map (AttributeMap)
import OpenTelemetry.Context (Context, insertSpan, lookupSpan)
import OpenTelemetry.Context.ThreadLocal (attachContext, detachContext, getContext)
import OpenTelemetry.Propagator.W3CTraceContext (decodeSpanContext, encodeSpanContext)
import OpenTelemetry.Trace.Core
  ( ExceptionClassification (..)
  , ExceptionHandler
  , ExceptionResponse (..)
  , NewEvent (..)
  , NewLink (..)
  , SpanArguments (..)
  , SpanContext
  , SpanKind (..)
  , SpanStatus (..)
  , Tracer
  , TracerOptions (..)
  , addEvent
  , defaultSpanArguments
  , getActiveSpan
  , getGlobalTracerProvider
  , getSpanContext
  , inSpan''
  , isValid
  , makeTracer
  , setStatus
  , tracerIsEnabled
  , tracerOptions
  , withActiveSpan
  , wrapSpanContext
  )
import UnliftIO (MonadUnliftIO, bracket)
import UnliftIO.Async (AsyncCancelled (..))

import Arbiter.Core.Exceptions
  ( JobForceCancelled (..)
  , JobGoneException (..)
  , JobNackException (..)
  )
import Arbiter.Core.Job.Schema (TableName)
import Arbiter.Core.Job.Types (HasKind (..), Job, JobRead, JobWrite, TraceContext (..))
import Arbiter.Core.Job.Types qualified as JT

-- $enrichment
-- Reach for @hs-opentelemetry-api@ directly for custom spans and attributes.

-- | Fill a job's trace context. A job that carries one of its own keeps it.
stampTraceContext :: Maybe TraceContext -> JobWrite payload -> JobWrite payload
stampTraceContext :: forall payload.
Maybe TraceContext -> JobWrite payload -> JobWrite payload
stampTraceContext Maybe TraceContext
Nothing JobWrite payload
job = JobWrite payload
job
stampTraceContext Maybe TraceContext
ctx JobWrite payload
job = Maybe TraceContext -> JobWrite payload -> JobWrite payload
forall payload.
Maybe TraceContext -> JobWrite payload -> JobWrite payload
JT.setTraceContext (JobWrite payload -> Maybe TraceContext
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> Maybe TraceContext
JT.traceContext JobWrite payload
job Maybe TraceContext -> Maybe TraceContext -> Maybe TraceContext
forall a. Maybe a -> Maybe a -> Maybe a
forall (f :: * -> *) a. Alternative f => f a -> f a -> f a
<|> Maybe TraceContext
ctx) JobWrite payload
job

-- | The ambient span's trace context, or 'Nothing' when no span is active. Read from
-- the thread-local context of the enqueuing thread.
currentTraceContext :: IO (Maybe TraceContext)
currentTraceContext :: IO (Maybe TraceContext)
currentTraceContext = IO (Maybe TraceContext)
-> (Span -> IO (Maybe TraceContext))
-> Maybe Span
-> IO (Maybe TraceContext)
forall b a. b -> (a -> b) -> Maybe a -> b
maybe (Maybe TraceContext -> IO (Maybe TraceContext)
forall a. a -> IO a
forall (f :: * -> *) a. Applicative f => a -> f a
pure Maybe TraceContext
forall a. Maybe a
Nothing) Span -> IO (Maybe TraceContext)
encoded (Maybe Span -> IO (Maybe TraceContext))
-> IO (Maybe Span) -> IO (Maybe TraceContext)
forall (m :: * -> *) a b. Monad m => (a -> m b) -> m a -> m b
=<< IO (Maybe Span)
forall (m :: * -> *). MonadIO m => m (Maybe Span)
getActiveSpan
  where
    encoded :: Span -> IO (Maybe TraceContext)
encoded Span
activeSpan = do
      valid <- SpanContext -> Bool
isValid (SpanContext -> Bool) -> IO SpanContext -> IO Bool
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> Span -> IO SpanContext
forall (m :: * -> *). MonadIO m => Span -> m SpanContext
getSpanContext Span
activeSpan
      if not valid
        then pure Nothing
        else do
          (parent, state) <- encodeSpanContext activeSpan
          pure (Just (TraceContext (decodeUtf8 parent) (decodeUtf8 state <$ guard (not (BS.null state)))))

-- | Resolve the arbiter tracer once, or 'Nothing' when nothing is collecting.
resolveTracer :: (MonadIO m) => m (Maybe Tracer)
resolveTracer :: forall (m :: * -> *). MonadIO m => m (Maybe Tracer)
resolveTracer = do
  provider <- m TracerProvider
forall (m :: * -> *). MonadIO m => m TracerProvider
getGlobalTracerProvider
  let tracer = TracerProvider -> InstrumentationLibrary -> TracerOptions -> Tracer
makeTracer TracerProvider
provider InstrumentationLibrary
"arbiter" TracerOptions
arbiterTracerOptions
  pure (if tracerIsEnabled tracer then Just tracer else Nothing)

arbiterTracerOptions :: TracerOptions
arbiterTracerOptions :: TracerOptions
arbiterTracerOptions = TracerOptions
tracerOptions {tracerExceptionHandlerOptions = [routineControlFlow]}

-- | Record control-flow exceptions for nacks, reclaims, and job cancellations
-- without an error status. Ignore worker cancellation exceptions.
routineControlFlow :: ExceptionHandler
routineControlFlow :: ExceptionHandler
routineControlFlow SomeException
exception
  | Just JobNackException
JobNackException <- SomeException -> Maybe JobNackException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = Maybe ExceptionResponse
recorded
  | Just (JobGoneException Text
_ [Int64]
_) <- SomeException -> Maybe JobGoneException
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = Maybe ExceptionResponse
recorded
  | Just (JobForceCancelled [Int64]
_ [Int64]
_) <- SomeException -> Maybe JobForceCancelled
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = Maybe ExceptionResponse
recorded
  | Just AsyncCancelled
AsyncCancelled <- SomeException -> Maybe AsyncCancelled
forall e. Exception e => SomeException -> Maybe e
fromException SomeException
exception = ExceptionResponse -> Maybe ExceptionResponse
forall a. a -> Maybe a
Just (ExceptionClassification -> AttributeMap -> ExceptionResponse
ExceptionResponse ExceptionClassification
IgnoredException AttributeMap
forall a. Monoid a => a
mempty)
  | Bool
otherwise = Maybe ExceptionResponse
forall a. Maybe a
Nothing
  where
    recorded :: Maybe ExceptionResponse
recorded = ExceptionResponse -> Maybe ExceptionResponse
forall a. a -> Maybe a
Just (ExceptionClassification -> AttributeMap -> ExceptionResponse
ExceptionResponse ExceptionClassification
RecordedException AttributeMap
forall a. Monoid a => a
mempty)

-- | Run an action inside a named span, unwrapped when no tracer was resolved.
spanning :: (MonadUnliftIO m) => Maybe Tracer -> Text -> SpanArguments -> m a -> m a
spanning :: forall (m :: * -> *) a.
MonadUnliftIO m =>
Maybe Tracer -> Text -> SpanArguments -> m a -> m a
spanning Maybe Tracer
mTracer Text
name SpanArguments
args m a
action =
  m a -> (Tracer -> m a) -> Maybe Tracer -> m a
forall b a. b -> (a -> b) -> Maybe a -> b
maybe m a
action (\Tracer
tracer -> Tracer -> Text -> SpanArguments -> (Span -> m a) -> m a
forall (m :: * -> *) a.
(MonadUnliftIO m, HasCallStack) =>
Tracer -> Text -> SpanArguments -> (Span -> m a) -> m a
inSpan'' Tracer
tracer Text
name SpanArguments
args (m a -> Span -> m a
forall a b. a -> b -> a
const m a
action)) Maybe Tracer
mTracer

-- | Run an action inside a @publish \<queue\>@ producer span over @n@ jobs.
withPublishSpan :: (HasKind payload, MonadUnliftIO m) => TableName -> [JobWrite payload] -> m a -> m a
withPublishSpan :: forall payload (m :: * -> *) a.
(HasKind payload, MonadUnliftIO m) =>
Text -> [JobWrite payload] -> m a -> m a
withPublishSpan Text
queue [JobWrite payload]
jobs m a
action =
  m (Maybe Tracer)
forall (m :: * -> *). MonadIO m => m (Maybe Tracer)
resolveTracer m (Maybe Tracer) -> (Maybe Tracer -> m a) -> m a
forall a b. m a -> (a -> m b) -> m b
forall (m :: * -> *) a b. Monad m => m a -> (a -> m b) -> m b
>>= \Maybe Tracer
tracer -> Maybe Tracer -> Text -> SpanArguments -> m a -> m a
forall (m :: * -> *) a.
MonadUnliftIO m =>
Maybe Tracer -> Text -> SpanArguments -> m a -> m a
spanning Maybe Tracer
tracer (Text
"publish " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
queue) (Text -> [JobWrite payload] -> SpanArguments
forall payload.
HasKind payload =>
Text -> [JobWrite payload] -> SpanArguments
producerArgs Text
queue [JobWrite payload]
jobs) m a
action

-- | Mark the currently active span failed.
markSpanError :: (MonadIO m) => Text -> m ()
markSpanError :: forall (m :: * -> *). MonadIO m => Text -> m ()
markSpanError Text
msg = (Span -> m ()) -> m ()
forall (m :: * -> *). MonadIO m => (Span -> m ()) -> m ()
withActiveSpan (\Span
activeSpan -> Span -> SpanStatus -> m ()
forall (m :: * -> *). MonadIO m => Span -> SpanStatus -> m ()
setStatus Span
activeSpan (Text -> SpanStatus
Error Text
msg))

-- | Record one job's failure as an event on the active span. A no-op when none is active.
-- The span status is left alone.
recordJobFailure :: (MonadIO m) => JobRead payload -> Text -> m ()
recordJobFailure :: forall (m :: * -> *) payload.
MonadIO m =>
JobRead payload -> Text -> m ()
recordJobFailure = Text -> Text -> JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
MonadIO m =>
Text -> Text -> JobRead payload -> Text -> m ()
jobEvent Text
"job.failed" Text
"arbiter.job.failure_reason"

-- | Record one job's cancellation as an event on the active span.
recordJobCancelled :: (MonadIO m) => JobRead payload -> Text -> m ()
recordJobCancelled :: forall (m :: * -> *) payload.
MonadIO m =>
JobRead payload -> Text -> m ()
recordJobCancelled = Text -> Text -> JobRead payload -> Text -> m ()
forall (m :: * -> *) payload.
MonadIO m =>
Text -> Text -> JobRead payload -> Text -> m ()
jobEvent Text
"job.cancelled" Text
"arbiter.job.cancel_reason"

jobEvent :: (MonadIO m) => Text -> Text -> JobRead payload -> Text -> m ()
jobEvent :: forall (m :: * -> *) payload.
MonadIO m =>
Text -> Text -> JobRead payload -> Text -> m ()
jobEvent Text
name Text
reasonKey JobRead payload
job Text
msg = (Span -> m ()) -> m ()
forall (m :: * -> *). MonadIO m => (Span -> m ()) -> m ()
withActiveSpan (Span -> NewEvent -> m ()
forall (m :: * -> *). MonadIO m => Span -> NewEvent -> m ()
`addEvent` NewEvent
event)
  where
    event :: NewEvent
event =
      NewEvent
        { newEventName :: Text
newEventName = Text
name
        , newEventAttributes :: AttributeMap
newEventAttributes =
            [(Text, Attribute)] -> AttributeMap
forall k v. Hashable k => [(k, v)] -> HashMap k v
HM.fromList
              [ (Text
"messaging.message.id", JobRead payload -> Attribute
forall payload. JobRead payload -> Attribute
messageId JobRead payload
job)
              , (Text
reasonKey, Text -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute Text
msg)
              ]
        , newEventTimestamp :: Maybe Timestamp
newEventTimestamp = Maybe Timestamp
forall a. Maybe a
Nothing
        }

-- | Queue attributes for consumer spans. Resolve them one time for each pool
-- and reuse them for claimed jobs.
data ConsumeSpan = ConsumeSpan
  { ConsumeSpan -> Text
consumeName :: Text
  , ConsumeSpan -> AttributeMap
consumeAttrs :: AttributeMap
  , ConsumeSpan -> ConsumeShape
consumeShape :: ConsumeShape
  }

-- | What one consumer span covers.
data ConsumeShape = PerJob | PerBatch
  deriving stock (ConsumeShape -> ConsumeShape -> Bool
(ConsumeShape -> ConsumeShape -> Bool)
-> (ConsumeShape -> ConsumeShape -> Bool) -> Eq ConsumeShape
forall a. (a -> a -> Bool) -> (a -> a -> Bool) -> Eq a
$c== :: ConsumeShape -> ConsumeShape -> Bool
== :: ConsumeShape -> ConsumeShape -> Bool
$c/= :: ConsumeShape -> ConsumeShape -> Bool
/= :: ConsumeShape -> ConsumeShape -> Bool
Eq, Int -> ConsumeShape -> ShowS
[ConsumeShape] -> ShowS
ConsumeShape -> String
(Int -> ConsumeShape -> ShowS)
-> (ConsumeShape -> String)
-> ([ConsumeShape] -> ShowS)
-> Show ConsumeShape
forall a.
(Int -> a -> ShowS) -> (a -> String) -> ([a] -> ShowS) -> Show a
$cshowsPrec :: Int -> ConsumeShape -> ShowS
showsPrec :: Int -> ConsumeShape -> ShowS
$cshow :: ConsumeShape -> String
show :: ConsumeShape -> String
$cshowList :: [ConsumeShape] -> ShowS
showList :: [ConsumeShape] -> ShowS
Show)

-- | The consumer-span shape for a queue.
consumeSpanFor :: TableName -> ConsumeShape -> ConsumeSpan
consumeSpanFor :: Text -> ConsumeShape -> ConsumeSpan
consumeSpanFor Text
queue ConsumeShape
shape =
  ConsumeSpan
    { consumeName :: Text
consumeName = Text
"process " Text -> Text -> Text
forall a. Semigroup a => a -> a -> a
<> Text
queue
    , consumeAttrs :: AttributeMap
consumeAttrs = Text -> Text -> AttributeMap
messagingAttrs Text
queue Text
"process"
    , consumeShape :: ConsumeShape
consumeShape = ConsumeShape
shape
    }

-- | Run a job handler inside a @process \<queue\>@ consumer span linked to its producer.
withConsumeSpan
  :: (MonadUnliftIO m) => Maybe Tracer -> ConsumeSpan -> NonEmpty (JobRead payload) -> m a -> m a
withConsumeSpan :: forall (m :: * -> *) payload a.
MonadUnliftIO m =>
Maybe Tracer
-> ConsumeSpan -> NonEmpty (JobRead payload) -> m a -> m a
withConsumeSpan Maybe Tracer
mTracer ConsumeSpan
consumeSpan NonEmpty (JobRead payload)
jobs =
  Maybe Tracer -> Text -> SpanArguments -> m a -> m a
forall (m :: * -> *) a.
MonadUnliftIO m =>
Maybe Tracer -> Text -> SpanArguments -> m a -> m a
spanning Maybe Tracer
mTracer (ConsumeSpan -> Text
consumeName ConsumeSpan
consumeSpan) SpanArguments
args
  where
    (JobRead payload
firstJob :| [JobRead payload]
_) = NonEmpty (JobRead payload)
jobs
    args :: SpanArguments
args = case ConsumeSpan -> ConsumeShape
consumeShape ConsumeSpan
consumeSpan of
      ConsumeShape
PerBatch -> ConsumeSpan -> NonEmpty (JobRead payload) -> SpanArguments
forall payload.
ConsumeSpan -> NonEmpty (JobRead payload) -> SpanArguments
batchConsumerArgs ConsumeSpan
consumeSpan NonEmpty (JobRead payload)
jobs
      ConsumeShape
PerJob -> ConsumeSpan -> JobRead payload -> SpanArguments
forall payload. ConsumeSpan -> JobRead payload -> SpanArguments
consumerArgs ConsumeSpan
consumeSpan JobRead payload
firstJob

-- | Capture the caller context for an action on a child thread. Return the
-- unchanged action if the context has no span.
capturingContext :: (MonadUnliftIO m) => m (m a -> m a)
capturingContext :: forall (m :: * -> *) a. MonadUnliftIO m => m (m a -> m a)
capturingContext = (\Context
ctx -> (m a -> m a) -> (Span -> m a -> m a) -> Maybe Span -> m a -> m a
forall b a. b -> (a -> b) -> Maybe a -> b
maybe m a -> m a
forall a. a -> a
id ((m a -> m a) -> Span -> m a -> m a
forall a b. a -> b -> a
const (Context -> m a -> m a
forall (m :: * -> *) a. MonadUnliftIO m => Context -> m a -> m a
withContext Context
ctx)) (Context -> Maybe Span
lookupSpan Context
ctx)) (Context -> m a -> m a) -> m Context -> m (m a -> m a)
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> m Context
forall (m :: * -> *). MonadIO m => m Context
getContext

-- | Run an action under @ctx@, restoring the caller's own afterwards.
withContext :: (MonadUnliftIO m) => Context -> m a -> m a
withContext :: forall (m :: * -> *) a. MonadUnliftIO m => Context -> m a -> m a
withContext Context
ctx m a
action = m Token -> (Token -> m ()) -> (Token -> m a) -> m a
forall (m :: * -> *) a b c.
MonadUnliftIO m =>
m a -> (a -> m b) -> (a -> m c) -> m c
bracket (Context -> m Token
forall (m :: * -> *). MonadIO m => Context -> m Token
attachContext Context
ctx) Token -> m ()
forall (m :: * -> *). MonadIO m => Token -> m ()
detachContext (m a -> Token -> m a
forall a b. a -> b -> a
const m a
action)

-- | Run an action with the job's stored trace context as the ambient parent.
withJobParent :: (MonadUnliftIO m) => JobRead payload -> m a -> m a
withJobParent :: forall (m :: * -> *) payload a.
MonadUnliftIO m =>
JobRead payload -> m a -> m a
withJobParent JobRead payload
job m a
action = m a -> (SpanContext -> m a) -> Maybe SpanContext -> m a
forall b a. b -> (a -> b) -> Maybe a -> b
maybe m a
action SpanContext -> m a
attached (JobRead payload -> Maybe SpanContext
forall payload. JobRead payload -> Maybe SpanContext
spanContextForJob JobRead payload
job)
  where
    attached :: SpanContext -> m a
attached SpanContext
spanContext = (Context -> m a -> m a) -> m a -> Context -> m a
forall a b c. (a -> b -> c) -> b -> a -> c
flip Context -> m a -> m a
forall (m :: * -> *) a. MonadUnliftIO m => Context -> m a -> m a
withContext m a
action (Context -> m a) -> (Context -> Context) -> Context -> m a
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Span -> Context -> Context
insertSpan (SpanContext -> Span
wrapSpanContext SpanContext
spanContext) (Context -> m a) -> m Context -> m a
forall (m :: * -> *) a b. Monad m => (a -> m b) -> m a -> m b
=<< m Context
forall (m :: * -> *). MonadIO m => m Context
getContext

-- | A span link reconstructed from a job's stored W3C trace context.
spanLinkForJob :: JobRead payload -> Maybe NewLink
spanLinkForJob :: forall payload. JobRead payload -> Maybe NewLink
spanLinkForJob JobRead payload
job =
  (\SpanContext
spanContext -> NewLink {linkContext :: SpanContext
linkContext = SpanContext
spanContext, linkAttributes :: AttributeMap
linkAttributes = AttributeMap
forall a. Monoid a => a
mempty}) (SpanContext -> NewLink) -> Maybe SpanContext -> Maybe NewLink
forall (f :: * -> *) a b. Functor f => (a -> b) -> f a -> f b
<$> JobRead payload -> Maybe SpanContext
forall payload. JobRead payload -> Maybe SpanContext
spanContextForJob JobRead payload
job

spanContextForJob :: JobRead payload -> Maybe SpanContext
spanContextForJob :: forall payload. JobRead payload -> Maybe SpanContext
spanContextForJob JobRead payload
job = do
  ctx <- JobRead payload -> Maybe TraceContext
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> Maybe TraceContext
JT.traceContext JobRead payload
job
  decodeSpanContext (Just (encodeUtf8 (traceparent ctx))) (encodeUtf8 <$> tracestate ctx)

-- | A lone job contributes its own attributes. A batch reports its size, per the
-- messaging conventions.
producerArgs :: (HasKind payload) => TableName -> [JobWrite payload] -> SpanArguments
producerArgs :: forall payload.
HasKind payload =>
Text -> [JobWrite payload] -> SpanArguments
producerArgs Text
queue [JobWrite payload]
jobs =
  SpanArguments
defaultSpanArguments {kind = Producer, attributes = messagingAttrs queue "publish" <> published}
  where
    published :: AttributeMap
published = case [JobWrite payload]
jobs of
      [JobWrite payload
job] -> JobWrite payload -> AttributeMap
forall payload. HasKind payload => JobWrite payload -> AttributeMap
writeAttrs JobWrite payload
job
      [JobWrite payload]
_ -> [(Text, Attribute)] -> AttributeMap
forall k v. Hashable k => [(k, v)] -> HashMap k v
HM.fromList [(Text
"messaging.batch.message_count", Int -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute ([JobWrite payload] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [JobWrite payload]
jobs))]

-- | What a job carries whichever end of the queue reads it.
jobShapeAttrs :: Job payload key q ins adm -> [(Text, Attribute)]
jobShapeAttrs :: forall payload key q ins adm.
Job payload key q ins adm -> [(Text, Attribute)]
jobShapeAttrs Job payload key q ins adm
job =
  (Text
"arbiter.priority", Int -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute (Int32 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral (Job payload key q ins adm -> Int32
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> Int32
JT.priority Job payload key q ins adm
job) :: Int))
    (Text, Attribute) -> [(Text, Attribute)] -> [(Text, Attribute)]
forall a. a -> [a] -> [a]
: (Text -> [(Text, Attribute)]) -> Maybe Text -> [(Text, Attribute)]
forall m a. Monoid m => (a -> m) -> Maybe a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap (\Text
group -> [(Text
"arbiter.group_key", Text -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute Text
group)]) (Job payload key q ins adm -> Maybe Text
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> Maybe Text
JT.groupKey Job payload key q ins adm
job)

writeAttrs :: (HasKind payload) => JobWrite payload -> AttributeMap
writeAttrs :: forall payload. HasKind payload => JobWrite payload -> AttributeMap
writeAttrs JobWrite payload
job = [(Text, Attribute)] -> AttributeMap
forall k v. Hashable k => [(k, v)] -> HashMap k v
HM.fromList (Maybe Text -> [(Text, Attribute)]
kindAttr (payload -> Maybe Text
forall payload. HasKind payload => payload -> Maybe Text
kindOf (JobWrite payload -> payload
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> payload
JT.payload JobWrite payload
job)) [(Text, Attribute)] -> [(Text, Attribute)] -> [(Text, Attribute)]
forall a. Semigroup a => a -> a -> a
<> JobWrite payload -> [(Text, Attribute)]
forall payload key q ins adm.
Job payload key q ins adm -> [(Text, Attribute)]
jobShapeAttrs JobWrite payload
job)

-- | The payload's variant label, absent for a payload that declares none.
kindAttr :: Maybe Text -> [(Text, Attribute)]
kindAttr :: Maybe Text -> [(Text, Attribute)]
kindAttr Maybe Text
label = [(Text
"arbiter.kind", Text -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute Text
value) | Text
value <- Maybe Text -> [Text]
forall a. Maybe a -> [a]
maybeToList Maybe Text
label]

consumerArgs :: ConsumeSpan -> JobRead payload -> SpanArguments
consumerArgs :: forall payload. ConsumeSpan -> JobRead payload -> SpanArguments
consumerArgs ConsumeSpan
consumeSpan JobRead payload
job =
  SpanArguments
defaultSpanArguments
    { kind = Consumer
    , links = maybeToList (spanLinkForJob job)
    , attributes = consumeAttrs consumeSpan <> jobAttrs job
    }

batchConsumerArgs :: ConsumeSpan -> NonEmpty (JobRead payload) -> SpanArguments
batchConsumerArgs :: forall payload.
ConsumeSpan -> NonEmpty (JobRead payload) -> SpanArguments
batchConsumerArgs ConsumeSpan
consumeSpan NonEmpty (JobRead payload)
jobs =
  SpanArguments
defaultSpanArguments
    { kind = Consumer
    , links = kept
    , attributes = consumeAttrs consumeSpan <> HM.fromList (count : droppedAttr)
    }
  where
    ([NewLink]
kept, [NewLink]
over) = Int -> [NewLink] -> ([NewLink], [NewLink])
forall a. Int -> [a] -> ([a], [a])
splitAt Int
maxSpanLinks ((JobRead payload -> Maybe NewLink)
-> [JobRead payload] -> [NewLink]
forall a b. (a -> Maybe b) -> [a] -> [b]
mapMaybe JobRead payload -> Maybe NewLink
forall payload. JobRead payload -> Maybe NewLink
spanLinkForJob (NonEmpty (JobRead payload) -> [JobRead payload]
forall a. NonEmpty a -> [a]
NE.toList NonEmpty (JobRead payload)
jobs))
    count :: (Text, Attribute)
count = (Text
"messaging.batch.message_count", Int -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute (NonEmpty (JobRead payload) -> Int
forall a. NonEmpty a -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length NonEmpty (JobRead payload)
jobs :: Int))
    droppedAttr :: [(Text, Attribute)]
droppedAttr = [(Text
"arbiter.batch.links_dropped", Int -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute ([NewLink] -> Int
forall a. [a] -> Int
forall (t :: * -> *) a. Foldable t => t a -> Int
length [NewLink]
over :: Int)) | Bool -> Bool
not ([NewLink] -> Bool
forall a. [a] -> Bool
forall (t :: * -> *) a. Foldable t => t a -> Bool
null [NewLink]
over)]

-- | The link count a span keeps under the SDK's default limits.
maxSpanLinks :: Int
maxSpanLinks :: Int
maxSpanLinks = Int
128

messagingAttrs :: Text -> Text -> AttributeMap
messagingAttrs :: Text -> Text -> AttributeMap
messagingAttrs Text
queue Text
operation =
  [(Text, Attribute)] -> AttributeMap
forall k v. Hashable k => [(k, v)] -> HashMap k v
HM.fromList
    [ (Text
"messaging.system", Text -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute (Text
"arbiter" :: Text))
    , (Text
"messaging.destination.name", Text -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute Text
queue)
    , (Text
"messaging.operation.type", Text -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute Text
operation)
    , (Text
"messaging.operation.name", Text -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute Text
operation)
    ]

-- | Textual, per the messaging semantic conventions.
messageId :: JobRead payload -> Attribute
messageId :: forall payload. JobRead payload -> Attribute
messageId = Text -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute (Text -> Attribute)
-> (JobRead payload -> Text) -> JobRead payload -> Attribute
forall b c a. (b -> c) -> (a -> b) -> a -> c
. LazyText -> Text
toStrict (LazyText -> Text)
-> (JobRead payload -> LazyText) -> JobRead payload -> Text
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Builder -> LazyText
toLazyText (Builder -> LazyText)
-> (JobRead payload -> Builder) -> JobRead payload -> LazyText
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Int64 -> Builder
forall a. Integral a => a -> Builder
decimal (Int64 -> Builder)
-> (JobRead payload -> Int64) -> JobRead payload -> Builder
forall b c a. (b -> c) -> (a -> b) -> a -> c
. JobRead payload -> Int64
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> key
JT.primaryKey

jobAttrs :: JobRead payload -> AttributeMap
jobAttrs :: forall payload. JobRead payload -> AttributeMap
jobAttrs JobRead payload
job =
  [(Text, Attribute)] -> AttributeMap
forall k v. Hashable k => [(k, v)] -> HashMap k v
HM.fromList ([(Text, Attribute)] -> AttributeMap)
-> [(Text, Attribute)] -> AttributeMap
forall a b. (a -> b) -> a -> b
$
    [ (Text
"messaging.message.id", JobRead payload -> Attribute
forall payload. JobRead payload -> Attribute
messageId JobRead payload
job)
    , (Text
"messaging.message.retry.count", Int -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute (Int -> Int -> Int
forall a. Ord a => a -> a -> a
max Int
0 (Int32 -> Int
forall a b. (Integral a, Num b) => a -> b
fromIntegral (JobRead payload -> Int32
forall payload q insertedAt adm.
JobRecord payload Int64 q insertedAt adm -> Int32
JT.attempts JobRead payload
job) Int -> Int -> Int
forall a. Num a => a -> a -> a
- Int
1) :: Int))
    ]
      [(Text, Attribute)] -> [(Text, Attribute)] -> [(Text, Attribute)]
forall a. Semigroup a => a -> a -> a
<> Maybe Text -> [(Text, Attribute)]
kindAttr (PayloadKeys -> Maybe Text
JT.jobKind (JobRead payload -> PayloadKeys
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> adm
JT.payloadKeys JobRead payload
job))
      [(Text, Attribute)] -> [(Text, Attribute)] -> [(Text, Attribute)]
forall a. Semigroup a => a -> a -> a
<> JobRead payload -> [(Text, Attribute)]
forall payload key q ins adm.
Job payload key q ins adm -> [(Text, Attribute)]
jobShapeAttrs JobRead payload
job