{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE RankNTypes #-}
{-# LANGUAGE ScopedTypeVariables #-}
{-# LANGUAGE TypeApplications #-}

-- | OpenTelemetry for arbiter: traces, metrics and logs.
--
-- Spans and trace-context propagation are built in. This module is the SDK side.
-- 'runWorkerPools' installs the exporters and instruments the pools. The @With@ variants
-- take a handle the caller installed itself.
--
-- @
-- env <- createHasqlEnv ...
-- runHasqlDb env $ Otel.runWorkerPools [namedWorkerPool emailCfg]
-- @
module Arbiter.Otel
  ( -- * Setup
    Telemetry (..)
  , withTelemetry
  , withTelemetryIf
  , withTelemetryFromEnv
  , withExternalTelemetry
  , telemetryLogConfig

    -- * Instrumenting a pool
  , instrumentPool
  , instrumentPools
  , instrumentConfig

    -- * Running pools
  , runWorkerPools
  , runSelectedWorkerPools
  , runWorkerPoolsWith
  , runSelectedWorkerPoolsWith

    -- * Gauges
  , withGauges
  , startGauges

    -- * Metric names
  , arbiterMetricNames
  ) where

import Arbiter.Core.HighLevel qualified as Arb
import Arbiter.Core.Job.Types (HasKind)
import Arbiter.Core.MonadArbiter (MonadArbiter, RegistryOf, getSchema)
import Arbiter.Core.QueueRegistry (RegistryTables, TableForPayload, registryQueueKinds)
import Arbiter.Core.Threads (labelArbiterThread)
import Arbiter.Worker (NamedWorkerPool (..))
import Arbiter.Worker qualified as Worker
import Arbiter.Worker.Config (WorkerConfig (..), withHooks, withMaintenance)
import Arbiter.Worker.Logger (LogConfig, LogLevel (..), defaultLogConfig, tryLog)
import Data.Proxy (Proxy (..))
import Data.Text (Text)
import GHC.TypeLits (KnownSymbol)
import UnliftIO (MonadUnliftIO, withRunInIO)
import UnliftIO.Async (withAsync)

import Arbiter.Otel.Gauges (startGauges, withGaugeLoop)
import Arbiter.Otel.MetricNames (arbiterMetricNames)
import Arbiter.Otel.Metrics (otelHooks, otelMaintenance)
import Arbiter.Otel.Telemetry
  ( Telemetry (..)
  , telemetryLogConfig
  , withExternalTelemetry
  , withTelemetry
  , withTelemetryFromEnv
  , withTelemetryIf
  )

-- | Give a pool consumer spans, lifecycle metrics, and the telemetry log destination.
-- The pool's registry queue name labels its metrics and gauges. Apply it once per pool.
instrumentPool :: (MonadUnliftIO m) => Telemetry -> NamedWorkerPool m -> NamedWorkerPool m
instrumentPool :: forall (m :: * -> *).
MonadUnliftIO m =>
Telemetry -> NamedWorkerPool m -> NamedWorkerPool m
instrumentPool Telemetry
tel (NamedWorkerPool Text
queue WorkerConfig m payload
cfg) =
  Text -> WorkerConfig m payload -> NamedWorkerPool m
forall (m :: * -> *) payload.
(EncodeJobResult (ResultOf m payload), QueueOperation m payload,
 RegistryAdmissionPolicies (RegistryOf m),
 RegistryTables (RegistryOf m)) =>
Text -> WorkerConfig m payload -> NamedWorkerPool m
NamedWorkerPool Text
queue (Telemetry
-> Text -> WorkerConfig m payload -> WorkerConfig m payload
forall payload (m :: * -> *).
(HasKind payload, MonadUnliftIO m) =>
Telemetry
-> Text -> WorkerConfig m payload -> WorkerConfig m payload
labelledConfig Telemetry
tel Text
queue WorkerConfig m payload
cfg)

-- | 'instrumentPool' over a pool list.
instrumentPools :: (MonadUnliftIO m) => Telemetry -> [NamedWorkerPool m] -> [NamedWorkerPool m]
instrumentPools :: forall (m :: * -> *).
MonadUnliftIO m =>
Telemetry -> [NamedWorkerPool m] -> [NamedWorkerPool m]
instrumentPools = (NamedWorkerPool m -> NamedWorkerPool m)
-> [NamedWorkerPool m] -> [NamedWorkerPool m]
forall a b. (a -> b) -> [a] -> [b]
map ((NamedWorkerPool m -> NamedWorkerPool m)
 -> [NamedWorkerPool m] -> [NamedWorkerPool m])
-> (Telemetry -> NamedWorkerPool m -> NamedWorkerPool m)
-> Telemetry
-> [NamedWorkerPool m]
-> [NamedWorkerPool m]
forall b c a. (b -> c) -> (a -> b) -> a -> c
. Telemetry -> NamedWorkerPool m -> NamedWorkerPool m
forall (m :: * -> *).
MonadUnliftIO m =>
Telemetry -> NamedWorkerPool m -> NamedWorkerPool m
instrumentPool

-- | Run the registry's depth and health gauges alongside @action@. Nothing is scanned
-- when the handle has metrics off.
withGauges
  :: forall m b
   . (MonadArbiter m, RegistryTables (RegistryOf m))
  => Telemetry
  -> LogConfig
  -> m b
  -> m b
withGauges :: forall (m :: * -> *) b.
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Telemetry -> LogConfig -> m b -> m b
withGauges Telemetry
tel LogConfig
baseLog m b
action = do
  schema <- m Text
forall (m :: * -> *). MonadArbiter m => m Text
getSchema
  withRunInIO $ \forall a. m a -> IO a
runDb ->
    Telemetry
-> LogConfig
-> (forall a. m a -> IO a)
-> Text
-> [(Text, [Text])]
-> NominalDiffTime
-> (IO () -> IO b)
-> IO b
forall (m :: * -> *) b.
MonadArbiter m =>
Telemetry
-> LogConfig
-> (forall a. m a -> IO a)
-> Text
-> [(Text, [Text])]
-> NominalDiffTime
-> (IO () -> IO b)
-> IO b
withGaugeLoop Telemetry
tel LogConfig
baseLog m a -> IO a
forall a. m a -> IO a
runDb Text
schema (Proxy (RegistryOf m) -> [(Text, [Text])]
forall (registry :: JobPayloadRegistry).
RegistryTables registry =>
Proxy registry -> [(Text, [Text])]
registryQueueKinds (forall (t :: JobPayloadRegistry). Proxy t
forall {k} (t :: k). Proxy t
Proxy @(RegistryOf m))) (Telemetry -> NominalDiffTime
gaugeRefresh Telemetry
tel) ((IO () -> IO b) -> IO b) -> (IO () -> IO b) -> IO b
forall a b. (a -> b) -> a -> b
$
      \IO ()
loop -> IO () -> (Async () -> IO b) -> IO b
forall (m :: * -> *) a b.
MonadUnliftIO m =>
m a -> (Async a -> m b) -> m b
withAsync (Text -> Maybe Text -> IO ()
forall (m :: * -> *). MonadIO m => Text -> Maybe Text -> m ()
labelArbiterThread Text
"gauges" Maybe Text
forall a. Maybe a
Nothing IO () -> IO () -> IO ()
forall a b. IO a -> IO b -> IO b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> IO ()
loop) (IO b -> Async () -> IO b
forall a b. a -> b -> a
const (m b -> IO b
forall a. m a -> IO a
runDb m b
action))

-- | 'Arbiter.Worker.runWorkerPools' with the SDK installed from the environment, the
-- pools instrumented and the gauges running. Take the handle yourself with
-- 'runWorkerPoolsWith' when something else needs it.
runWorkerPools
  :: forall m
   . (MonadArbiter m, RegistryTables (RegistryOf m))
  => [NamedWorkerPool m]
  -> m ()
runWorkerPools :: forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
[NamedWorkerPool m] -> m ()
runWorkerPools [NamedWorkerPool m]
pools =
  let baseLog :: LogConfig
baseLog = [NamedWorkerPool m] -> LogConfig
forall (m :: * -> *). [NamedWorkerPool m] -> LogConfig
poolsLogConfig [NamedWorkerPool m]
pools
   in LogConfig -> (Telemetry -> m ()) -> m ()
forall (m :: * -> *) a.
MonadUnliftIO m =>
LogConfig -> (Telemetry -> m a) -> m a
withTelemetryHere LogConfig
baseLog ((Telemetry -> m ()) -> m ()) -> (Telemetry -> m ()) -> m ()
forall a b. (a -> b) -> a -> b
$ \Telemetry
tel -> Telemetry -> LogConfig -> [NamedWorkerPool m] -> m ()
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Telemetry -> LogConfig -> [NamedWorkerPool m] -> m ()
runWorkerPoolsWith Telemetry
tel LogConfig
baseLog [NamedWorkerPool m]
pools

-- | 'runWorkerPools' over an explicit queue list.
runSelectedWorkerPools
  :: forall m
   . (MonadArbiter m, RegistryTables (RegistryOf m))
  => [Text]
  -> [NamedWorkerPool m]
  -> m ()
runSelectedWorkerPools :: forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
[Text] -> [NamedWorkerPool m] -> m ()
runSelectedWorkerPools [Text]
enabled [NamedWorkerPool m]
pools =
  let baseLog :: LogConfig
baseLog = [NamedWorkerPool m] -> LogConfig
forall (m :: * -> *). [NamedWorkerPool m] -> LogConfig
poolsLogConfig [NamedWorkerPool m]
pools
   in LogConfig -> (Telemetry -> m ()) -> m ()
forall (m :: * -> *) a.
MonadUnliftIO m =>
LogConfig -> (Telemetry -> m a) -> m a
withTelemetryHere LogConfig
baseLog ((Telemetry -> m ()) -> m ()) -> (Telemetry -> m ()) -> m ()
forall a b. (a -> b) -> a -> b
$ \Telemetry
tel -> Telemetry -> LogConfig -> [Text] -> [NamedWorkerPool m] -> m ()
forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Telemetry -> LogConfig -> [Text] -> [NamedWorkerPool m] -> m ()
runSelectedWorkerPoolsWith Telemetry
tel LogConfig
baseLog [Text]
enabled [NamedWorkerPool m]
pools

-- | The first pool's log config, used as the gauge loop's base.
poolsLogConfig :: [NamedWorkerPool m] -> LogConfig
poolsLogConfig :: forall (m :: * -> *). [NamedWorkerPool m] -> LogConfig
poolsLogConfig (NamedWorkerPool Text
_ WorkerConfig m payload
cfg : [NamedWorkerPool m]
_) = WorkerConfig m payload -> LogConfig
forall (m :: * -> *) payload. WorkerConfig m payload -> LogConfig
logConfig WorkerConfig m payload
cfg
poolsLogConfig [] = LogConfig
defaultLogConfig

-- | 'runWorkerPools' over a handle the caller installed itself, with the gauge loop's
-- base log config.
runWorkerPoolsWith
  :: forall m
   . (MonadArbiter m, RegistryTables (RegistryOf m))
  => Telemetry
  -> LogConfig
  -> [NamedWorkerPool m]
  -> m ()
runWorkerPoolsWith :: forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Telemetry -> LogConfig -> [NamedWorkerPool m] -> m ()
runWorkerPoolsWith Telemetry
tel LogConfig
baseLog [NamedWorkerPool m]
pools =
  Telemetry -> LogConfig -> m () -> m ()
forall (m :: * -> *) b.
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Telemetry -> LogConfig -> m b -> m b
withGauges Telemetry
tel LogConfig
baseLog ([NamedWorkerPool m] -> m ()
forall (m :: * -> *).
(MonadUnliftIO m, RegistryTables (RegistryOf m)) =>
[NamedWorkerPool m] -> m ()
Worker.runWorkerPools (Telemetry -> [NamedWorkerPool m] -> [NamedWorkerPool m]
forall (m :: * -> *).
MonadUnliftIO m =>
Telemetry -> [NamedWorkerPool m] -> [NamedWorkerPool m]
instrumentPools Telemetry
tel [NamedWorkerPool m]
pools))

-- | 'runSelectedWorkerPools' over a handle the caller installed itself.
runSelectedWorkerPoolsWith
  :: forall m
   . (MonadArbiter m, RegistryTables (RegistryOf m))
  => Telemetry
  -> LogConfig
  -> [Text]
  -> [NamedWorkerPool m]
  -> m ()
runSelectedWorkerPoolsWith :: forall (m :: * -> *).
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Telemetry -> LogConfig -> [Text] -> [NamedWorkerPool m] -> m ()
runSelectedWorkerPoolsWith Telemetry
tel LogConfig
baseLog [Text]
enabled [NamedWorkerPool m]
pools =
  Telemetry -> LogConfig -> m () -> m ()
forall (m :: * -> *) b.
(MonadArbiter m, RegistryTables (RegistryOf m)) =>
Telemetry -> LogConfig -> m b -> m b
withGauges Telemetry
tel LogConfig
baseLog ([Text] -> [NamedWorkerPool m] -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
[Text] -> [NamedWorkerPool m] -> m ()
Worker.runSelectedWorkerPools [Text]
enabled (Telemetry -> [NamedWorkerPool m] -> [NamedWorkerPool m]
forall (m :: * -> *).
MonadUnliftIO m =>
Telemetry -> [NamedWorkerPool m] -> [NamedWorkerPool m]
instrumentPools Telemetry
tel [NamedWorkerPool m]
pools))

-- | Install the SDK around an action in the database monad, logging how it started.
withTelemetryHere :: (MonadUnliftIO m) => LogConfig -> (Telemetry -> m a) -> m a
withTelemetryHere :: forall (m :: * -> *) a.
MonadUnliftIO m =>
LogConfig -> (Telemetry -> m a) -> m a
withTelemetryHere LogConfig
baseLog Telemetry -> m a
use = ((forall a. m a -> IO a) -> IO a) -> m a
forall b. ((forall a. m a -> IO a) -> IO b) -> m b
forall (m :: * -> *) b.
MonadUnliftIO m =>
((forall a. m a -> IO a) -> IO b) -> m b
withRunInIO (((forall a. m a -> IO a) -> IO a) -> m a)
-> ((forall a. m a -> IO a) -> IO a) -> m a
forall a b. (a -> b) -> a -> b
$ \forall a. m a -> IO a
runDb ->
  (Telemetry -> IO a) -> IO a
forall a. (Telemetry -> IO a) -> IO a
withTelemetryFromEnv ((Telemetry -> IO a) -> IO a) -> (Telemetry -> IO a) -> IO a
forall a b. (a -> b) -> a -> b
$ \Telemetry
tel ->
    m a -> IO a
forall a. m a -> IO a
runDb (LogConfig -> LogLevel -> Text -> m ()
forall (m :: * -> *).
MonadUnliftIO m =>
LogConfig -> LogLevel -> Text -> m ()
tryLog (Telemetry -> LogConfig -> LogConfig
telemetryLogConfig Telemetry
tel LogConfig
baseLog) LogLevel
Info (Telemetry -> Text
telemetrySummary Telemetry
tel) m () -> m a -> m a
forall a b. m a -> m b -> m b
forall (m :: * -> *) a b. Monad m => m a -> m b -> m b
>> Telemetry -> m a
use Telemetry
tel)

-- | 'instrumentPool' over a bare config, labelled by the payload's registry queue.
instrumentConfig
  :: forall m payload
   . (HasKind payload, KnownSymbol (TableForPayload payload (RegistryOf m)), MonadUnliftIO m)
  => Telemetry
  -> WorkerConfig m payload
  -> WorkerConfig m payload
instrumentConfig :: forall (m :: * -> *) payload.
(HasKind payload,
 KnownSymbol (TableForPayload payload (RegistryOf m)),
 MonadUnliftIO m) =>
Telemetry -> WorkerConfig m payload -> WorkerConfig m payload
instrumentConfig Telemetry
tel = Telemetry
-> Text -> WorkerConfig m payload -> WorkerConfig m payload
forall payload (m :: * -> *).
(HasKind payload, MonadUnliftIO m) =>
Telemetry
-> Text -> WorkerConfig m payload -> WorkerConfig m payload
labelledConfig Telemetry
tel (forall payload (m :: * -> *).
KnownSymbol (TableForPayload payload (RegistryOf m)) =>
Text
Arb.queueTable @payload @m)

labelledConfig
  :: (HasKind payload, MonadUnliftIO m)
  => Telemetry
  -> Text
  -> WorkerConfig m payload
  -> WorkerConfig m payload
labelledConfig :: forall payload (m :: * -> *).
(HasKind payload, MonadUnliftIO m) =>
Telemetry
-> Text -> WorkerConfig m payload -> WorkerConfig m payload
labelledConfig Telemetry
tel Text
queue WorkerConfig m payload
cfg =
  (WorkerConfig m payload -> WorkerConfig m payload)
-> (ArbiterMeters
    -> WorkerConfig m payload -> WorkerConfig m payload)
-> Maybe ArbiterMeters
-> WorkerConfig m payload
-> WorkerConfig m payload
forall b a. b -> (a -> b) -> Maybe a -> b
maybe WorkerConfig m payload -> WorkerConfig m payload
forall a. a -> a
id ArbiterMeters -> WorkerConfig m payload -> WorkerConfig m payload
instrument (Telemetry -> Maybe ArbiterMeters
meters Telemetry
tel) (WorkerConfig m payload -> WorkerConfig m payload)
-> WorkerConfig m payload -> WorkerConfig m payload
forall a b. (a -> b) -> a -> b
$ WorkerConfig m payload
cfg {logConfig = telemetryLogConfig tel (logConfig cfg)}
  where
    instrument :: ArbiterMeters -> WorkerConfig m payload -> WorkerConfig m payload
instrument ArbiterMeters
meterSet = (ObservabilityHooks m payload -> ObservabilityHooks m payload)
-> WorkerConfig m payload -> WorkerConfig m payload
forall (m :: * -> *) payload.
(ObservabilityHooks m payload -> ObservabilityHooks m payload)
-> WorkerConfig m payload -> WorkerConfig m payload
withHooks (ArbiterMeters -> Text -> ObservabilityHooks m payload
forall (m :: * -> *) payload.
(HasKind payload, MonadIO m) =>
ArbiterMeters -> Text -> ObservabilityHooks m payload
otelHooks ArbiterMeters
meterSet Text
queue ObservabilityHooks m payload
-> ObservabilityHooks m payload -> ObservabilityHooks m payload
forall a. Semigroup a => a -> a -> a
<>) (WorkerConfig m payload -> WorkerConfig m payload)
-> (WorkerConfig m payload -> WorkerConfig m payload)
-> WorkerConfig m payload
-> WorkerConfig m payload
forall b c a. (b -> c) -> (a -> b) -> a -> c
. (MaintenanceOp -> Int64 -> m ())
-> WorkerConfig m payload -> WorkerConfig m payload
forall (m :: * -> *) payload.
MonadUnliftIO m =>
(MaintenanceOp -> Int64 -> m ())
-> WorkerConfig m payload -> WorkerConfig m payload
withMaintenance (ArbiterMeters -> MaintenanceOp -> Int64 -> m ()
forall (m :: * -> *).
MonadIO m =>
ArbiterMeters -> MaintenanceOp -> Int64 -> m ()
otelMaintenance ArbiterMeters
meterSet)