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

-- | Worker-lifecycle OpenTelemetry instruments, bound through 'ObservabilityHooks'.
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

-- | The job-lifecycle instruments.
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
  }

-- | Arbiter's instrumentation scope, shared by every meter the project registers.
arbiterMeter :: MeterProvider -> IO Meter
arbiterMeter :: MeterProvider -> IO Meter
arbiterMeter MeterProvider
meterProvider = MeterProvider -> InstrumentationLibrary -> IO Meter
getMeter MeterProvider
meterProvider InstrumentationLibrary
"arbiter"

-- | Create the instruments on a meter provider.
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

-- | Second-scale histogram bounds.
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]
    }

-- | Observability hooks recording one queue's jobs to the instruments.
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
    -- A clock stepped backwards records zero.
    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)
    -- The declared labels.
    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))
    -- One attribute set per label, built once per hook set.
    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"
    -- The policy prefix.
    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)])

-- | Record rows affected by a reaper operation. Schema-wide reaper work has no
-- queue attribute.
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)]))

-- | Build an OTel attribute set from text key/value pairs.
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]

-- | The @policy_kind@ attribute for an admission metric, shared by counters and gauges.
rateLimitKind, concurrencyKind :: Text
rateLimitKind :: Text
rateLimitKind = Text
"rate_limit"
concurrencyKind :: Text
concurrencyKind = Text
"concurrency"