{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE ScopedTypeVariables #-}
module Arbiter.Core.Trace
(
TraceContext (..)
, stampTraceContext
, Tracer
, resolveTracer
, currentTraceContext
, withPublishSpan
, ConsumeSpan
, ConsumeShape (..)
, consumeSpanFor
, withConsumeSpan
, withJobParent
, capturingContext
, 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
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
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)))))
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]}
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)
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
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
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))
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"
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
}
data ConsumeSpan = ConsumeSpan
{ ConsumeSpan -> Text
consumeName :: Text
, ConsumeSpan -> AttributeMap
consumeAttrs :: AttributeMap
, ConsumeSpan -> ConsumeShape
consumeShape :: ConsumeShape
}
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)
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
}
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
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
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)
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
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)
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))]
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)
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)]
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)
]
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