{-# LANGUAGE FlexibleContexts #-}
{-# LANGUAGE OverloadedStrings #-}
module Arbiter.Otel.Metrics
( ArbiterMeters
, arbiterMeter
, newArbiterMeters
, otelHooks
, otelMaintenance
, attrs
, rateLimitKind
, concurrencyKind
) where
import Arbiter.Core.Concurrency.Spec (ConcurrencyKey (..))
import Arbiter.Core.Job.Types
( HasKind (..)
, ObservabilityHooks (..)
, PayloadKeys (..)
, defaultObservabilityHooks
, payloadKeys
)
import Arbiter.Core.RateLimit.Spec (RateLimitKey (..))
import Arbiter.Worker.Config (MaintenanceOp, maintenanceOpName)
import Control.Monad (mfilter)
import Control.Monad.IO.Class (MonadIO, liftIO)
import Data.Fixed (Fixed (MkFixed))
import Data.Foldable (traverse_)
import Data.Int (Int64)
import Data.Map.Strict qualified as Map
import Data.Set qualified as Set
import Data.Text (Text)
import Data.Time.Clock (diffUTCTime, nominalDiffTimeToSeconds)
import OpenTelemetry.Attributes (Attributes, toAttribute, unsafeAttributesFromListIgnoringLimits)
import OpenTelemetry.Metric.Core
( AdvisoryParameters (..)
, Counter
, Histogram
, Meter
, MeterProvider
, counterAdd
, defaultAdvisoryParameters
, getMeter
, histogramRecord
, meterCreateCounterInt64
, meterCreateHistogram
)
import Arbiter.Otel.MetricNames qualified as Name
data ArbiterMeters = ArbiterMeters
{ ArbiterMeters -> Counter Int64
claimed :: Counter Int64
, ArbiterMeters -> Counter Int64
processed :: Counter Int64
, ArbiterMeters -> Counter Int64
retried :: Counter Int64
, ArbiterMeters -> Counter Int64
admitted :: Counter Int64
, ArbiterMeters -> Counter Int64
maintained :: Counter Int64
, ArbiterMeters -> Histogram
duration :: Histogram
}
arbiterMeter :: MeterProvider -> IO Meter
arbiterMeter :: MeterProvider -> IO Meter
arbiterMeter MeterProvider
meterProvider = MeterProvider -> InstrumentationLibrary -> IO Meter
getMeter MeterProvider
meterProvider InstrumentationLibrary
"arbiter"
newArbiterMeters :: MeterProvider -> IO ArbiterMeters
newArbiterMeters :: MeterProvider -> IO ArbiterMeters
newArbiterMeters MeterProvider
meterProvider = do
meter <- MeterProvider -> IO Meter
arbiterMeter MeterProvider
meterProvider
let counter MetricName
name Text
desc =
Meter
-> Text
-> Maybe Text
-> Maybe Text
-> AdvisoryParameters
-> IO (Counter Int64)
meterCreateCounterInt64 Meter
meter (MetricName -> Text
Name.metricName MetricName
name) Maybe Text
forall a. Maybe a
Nothing (Text -> Maybe Text
forall a. a -> Maybe a
Just Text
desc) AdvisoryParameters
defaultAdvisoryParameters
ArbiterMeters
<$> counter Name.JobsClaimed "Jobs claimed by workers"
<*> counter Name.JobsProcessed "Jobs processed, by terminal outcome"
<*> counter Name.JobsRetries "Failed jobs scheduled for another attempt"
<*> counter Name.AdmissionAdmitted "Claimed jobs that passed an admission policy, by kind and policy"
<*> counter Name.MaintenanceRows "Rows a reaper op touched, by op"
<*> meterCreateHistogram
meter
(Name.metricName Name.HandlerDuration)
(Just "s")
(Just "Seconds a job spent in the handler (a batched pool records its batch's span for each job)")
durationBuckets
durationBuckets :: AdvisoryParameters
durationBuckets :: AdvisoryParameters
durationBuckets =
AdvisoryParameters
defaultAdvisoryParameters
{ advisoryExplicitBucketBoundaries = Just [0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1, 2.5, 5, 10, 30, 60, 300]
}
otelHooks :: forall m payload. (HasKind payload, MonadIO m) => ArbiterMeters -> Text -> ObservabilityHooks m payload
otelHooks :: forall (m :: * -> *) payload.
(HasKind payload, MonadIO m) =>
ArbiterMeters -> Text -> ObservabilityHooks m payload
otelHooks ArbiterMeters
meterSet Text
queue =
ObservabilityHooks m payload
forall (m :: * -> *) payload.
Applicative m =>
ObservabilityHooks m payload
defaultObservabilityHooks
{ onJobClaimed = \JobRead payload
job ClaimTime
_ -> IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> m ()) -> IO () -> m ()
forall a b. (a -> b) -> a -> b
$ do
Counter Int64 -> Int64 -> Attributes -> IO ()
forall a. Counter a -> a -> Attributes -> IO ()
counterAdd (ArbiterMeters -> Counter Int64
claimed ArbiterMeters
meterSet) Int64
1 (Maybe Text -> Attributes
queueAttr (JobRead payload -> Maybe Text
labelOf JobRead payload
job))
(RateLimitKey -> IO ()) -> Maybe RateLimitKey -> IO ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (Text -> Text -> IO ()
admissionThrough Text
rateLimitKind (Text -> IO ()) -> (RateLimitKey -> Text) -> RateLimitKey -> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. RateLimitKey -> Text
rlkPrefix) (PayloadKeys -> Maybe RateLimitKey
jobRateLimitKey (JobRead payload -> PayloadKeys
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> adm
payloadKeys JobRead payload
job))
(ConcurrencyKey -> IO ()) -> Maybe ConcurrencyKey -> IO ()
forall (t :: * -> *) (f :: * -> *) a b.
(Foldable t, Applicative f) =>
(a -> f b) -> t a -> f ()
traverse_ (Text -> Text -> IO ()
admissionThrough Text
concurrencyKind (Text -> IO ())
-> (ConcurrencyKey -> Text) -> ConcurrencyKey -> IO ()
forall b c a. (b -> c) -> (a -> b) -> a -> c
. ConcurrencyKey -> Text
ckPrefix) (PayloadKeys -> Maybe ConcurrencyKey
jobConcurrencyKey (JobRead payload -> PayloadKeys
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> adm
payloadKeys JobRead payload
job))
, onJobSuccess = \JobRead payload
job ClaimTime
start ClaimTime
end -> IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (IO () -> m ()) -> IO () -> m ()
forall a b. (a -> b) -> a -> b
$ do
let attributes :: Attributes
attributes = Maybe Text -> Attributes
successAttr (JobRead payload -> Maybe Text
labelOf JobRead payload
job)
Counter Int64 -> Int64 -> Attributes -> IO ()
forall a. Counter a -> a -> Attributes -> IO ()
counterAdd (ArbiterMeters -> Counter Int64
processed ArbiterMeters
meterSet) Int64
1 Attributes
attributes
Histogram -> Double -> Attributes -> IO ()
histogramRecord (ArbiterMeters -> Histogram
duration ArbiterMeters
meterSet) (ClaimTime -> ClaimTime -> Double
forall {a}. (Ord a, Fractional a) => ClaimTime -> ClaimTime -> a
secs ClaimTime
start ClaimTime
end) Attributes
attributes
, onJobFailure = \JobRead payload
job Text
_ ClaimTime
start ClaimTime
end -> IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (Histogram -> Double -> Attributes -> IO ()
histogramRecord (ArbiterMeters -> Histogram
duration ArbiterMeters
meterSet) (ClaimTime -> ClaimTime -> Double
forall {a}. (Ord a, Fractional a) => ClaimTime -> ClaimTime -> a
secs ClaimTime
start ClaimTime
end) (Maybe Text -> Attributes
failureAttr (JobRead payload -> Maybe Text
labelOf JobRead payload
job)))
, onJobRetry = \JobRead payload
job BackoffDelay
_ -> IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (Counter Int64 -> Int64 -> Attributes -> IO ()
forall a. Counter a -> a -> Attributes -> IO ()
counterAdd (ArbiterMeters -> Counter Int64
retried ArbiterMeters
meterSet) Int64
1 (Maybe Text -> Attributes
queueAttr (JobRead payload -> Maybe Text
labelOf JobRead payload
job)))
, onJobFailedAndMovedToDLQ = \Text
_ JobRead payload
job -> IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (Counter Int64 -> Int64 -> Attributes -> IO ()
forall a. Counter a -> a -> Attributes -> IO ()
counterAdd (ArbiterMeters -> Counter Int64
processed ArbiterMeters
meterSet) Int64
1 (Maybe Text -> Attributes
dlqAttr (JobRead payload -> Maybe Text
labelOf JobRead payload
job)))
, onJobCancelled = \JobRead payload
job Text
_ -> IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (Counter Int64 -> Int64 -> Attributes -> IO ()
forall a. Counter a -> a -> Attributes -> IO ()
counterAdd (ArbiterMeters -> Counter Int64
processed ArbiterMeters
meterSet) Int64
1 (Maybe Text -> Attributes
cancelledAttr (JobRead payload -> Maybe Text
labelOf JobRead payload
job)))
, onJobUnavailable = \JobRead payload
job Text
_ -> IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (Counter Int64 -> Int64 -> Attributes -> IO ()
forall a. Counter a -> a -> Attributes -> IO ()
counterAdd (ArbiterMeters -> Counter Int64
processed ArbiterMeters
meterSet) Int64
1 (Maybe Text -> Attributes
unavailableAttr (JobRead payload -> Maybe Text
labelOf JobRead payload
job)))
}
where
secs :: ClaimTime -> ClaimTime -> a
secs ClaimTime
start ClaimTime
end = case BackoffDelay -> Pico
nominalDiffTimeToSeconds (ClaimTime -> ClaimTime -> BackoffDelay
diffUTCTime ClaimTime
end ClaimTime
start) of
MkFixed Integer
picos -> a -> a -> a
forall a. Ord a => a -> a -> a
max a
0 (Integer -> a
forall a. Num a => Integer -> a
fromInteger Integer
picos a -> a -> a
forall a. Fractional a => a -> a -> a
/ a
1e12)
declared :: Set Text
declared = [Text] -> Set Text
forall a. Ord a => [a] -> Set a
Set.fromList (forall payload. HasKind payload => [Text]
kindsFor @payload)
labelOf :: JobRead payload -> Maybe Text
labelOf JobRead payload
job = (Text -> Bool) -> Maybe Text -> Maybe Text
forall (m :: * -> *) a. MonadPlus m => (a -> Bool) -> m a -> m a
mfilter (Text -> Set Text -> Bool
forall a. Ord a => a -> Set a -> Bool
`Set.member` Set Text
declared) (PayloadKeys -> Maybe Text
jobKind (JobRead payload -> PayloadKeys
forall payload key q insertedAt adm.
JobRecord payload key q insertedAt adm -> adm
payloadKeys JobRead payload
job))
byLabel :: (Maybe Text -> Attributes) -> Maybe Text -> Attributes
byLabel :: (Maybe Text -> Attributes) -> Maybe Text -> Attributes
byLabel Maybe Text -> Attributes
build = \Maybe Text
label -> Attributes
-> Maybe Text -> Map (Maybe Text) Attributes -> Attributes
forall k a. Ord k => a -> k -> Map k a -> a
Map.findWithDefault (Maybe Text -> Attributes
build Maybe Text
forall a. Maybe a
Nothing) Maybe Text
label Map (Maybe Text) Attributes
table
where
table :: Map (Maybe Text) Attributes
table = [(Maybe Text, Attributes)] -> Map (Maybe Text) Attributes
forall k a. Ord k => [(k, a)] -> Map k a
Map.fromList [(Maybe Text
label, Maybe Text -> Attributes
build Maybe Text
label) | Maybe Text
label <- Maybe Text
forall a. Maybe a
Nothing Maybe Text -> [Maybe Text] -> [Maybe Text]
forall a. a -> [a] -> [a]
: (Text -> Maybe Text) -> [Text] -> [Maybe Text]
forall a b. (a -> b) -> [a] -> [b]
map Text -> Maybe Text
forall a. a -> Maybe a
Just (Set Text -> [Text]
forall a. Set a -> [a]
Set.toList Set Text
declared)]
kindPair :: Maybe b -> [(Text, b)]
kindPair = (b -> [(Text, b)]) -> Maybe b -> [(Text, b)]
forall m a. Monoid m => (a -> m) -> Maybe a -> m
forall (t :: * -> *) m a.
(Foldable t, Monoid m) =>
(a -> m) -> t a -> m
foldMap (\b
kind -> [(Text
"kind", b
kind)])
queueAttr :: Maybe Text -> Attributes
queueAttr = (Maybe Text -> Attributes) -> Maybe Text -> Attributes
byLabel (\Maybe Text
label -> [(Text, Text)] -> Attributes
attrs ((Text
"queue", Text
queue) (Text, Text) -> [(Text, Text)] -> [(Text, Text)]
forall a. a -> [a] -> [a]
: Maybe Text -> [(Text, Text)]
forall {b}. Maybe b -> [(Text, b)]
kindPair Maybe Text
label))
outcomeAttr :: Text -> Maybe Text -> Attributes
outcomeAttr Text
outcome = (Maybe Text -> Attributes) -> Maybe Text -> Attributes
byLabel (\Maybe Text
label -> [(Text, Text)] -> Attributes
attrs ([(Text
"queue", Text
queue), (Text
"outcome", Text
outcome)] [(Text, Text)] -> [(Text, Text)] -> [(Text, Text)]
forall a. Semigroup a => a -> a -> a
<> Maybe Text -> [(Text, Text)]
forall {b}. Maybe b -> [(Text, b)]
kindPair Maybe Text
label))
successAttr :: Maybe Text -> Attributes
successAttr = Text -> Maybe Text -> Attributes
outcomeAttr Text
"success"
failureAttr :: Maybe Text -> Attributes
failureAttr = Text -> Maybe Text -> Attributes
outcomeAttr Text
"failure"
dlqAttr :: Maybe Text -> Attributes
dlqAttr = Text -> Maybe Text -> Attributes
outcomeAttr Text
"dlq"
cancelledAttr :: Maybe Text -> Attributes
cancelledAttr = Text -> Maybe Text -> Attributes
outcomeAttr Text
"cancelled"
unavailableAttr :: Maybe Text -> Attributes
unavailableAttr = Text -> Maybe Text -> Attributes
outcomeAttr Text
"unavailable"
admissionThrough :: Text -> Text -> IO ()
admissionThrough Text
policyKind Text
prefix =
Counter Int64 -> Int64 -> Attributes -> IO ()
forall a. Counter a -> a -> Attributes -> IO ()
counterAdd (ArbiterMeters -> Counter Int64
admitted ArbiterMeters
meterSet) Int64
1 ([(Text, Text)] -> Attributes
attrs [(Text
"queue", Text
queue), (Text
"policy_kind", Text
policyKind), (Text
"policy", Text
prefix)])
otelMaintenance :: (MonadIO m) => ArbiterMeters -> MaintenanceOp -> Int64 -> m ()
otelMaintenance :: forall (m :: * -> *).
MonadIO m =>
ArbiterMeters -> MaintenanceOp -> Int64 -> m ()
otelMaintenance ArbiterMeters
meterSet MaintenanceOp
operation Int64
rowCount = IO () -> m ()
forall a. IO a -> m a
forall (m :: * -> *) a. MonadIO m => IO a -> m a
liftIO (Counter Int64 -> Int64 -> Attributes -> IO ()
forall a. Counter a -> a -> Attributes -> IO ()
counterAdd (ArbiterMeters -> Counter Int64
maintained ArbiterMeters
meterSet) Int64
rowCount ([(Text, Text)] -> Attributes
attrs [(Text
"op", MaintenanceOp -> Text
maintenanceOpName MaintenanceOp
operation)]))
attrs :: [(Text, Text)] -> Attributes
attrs :: [(Text, Text)] -> Attributes
attrs [(Text, Text)]
kvs = [(Text, Attribute)] -> Attributes
unsafeAttributesFromListIgnoringLimits [(Text
key, Text -> Attribute
forall a. ToAttribute a => a -> Attribute
toAttribute Text
value) | (Text
key, Text
value) <- [(Text, Text)]
kvs]
rateLimitKind, concurrencyKind :: Text
rateLimitKind :: Text
rateLimitKind = Text
"rate_limit"
concurrencyKind :: Text
concurrencyKind = Text
"concurrency"