{-# LANGUAGE OverloadedStrings #-}
{-# LANGUAGE RankNTypes #-}
{-# LANGUAGE ScopedTypeVariables #-}
{-# LANGUAGE TypeApplications #-}
module Arbiter.Otel
(
Telemetry (..)
, withTelemetry
, withTelemetryIf
, withTelemetryFromEnv
, withExternalTelemetry
, telemetryLogConfig
, instrumentPool
, instrumentPools
, instrumentConfig
, runWorkerPools
, runSelectedWorkerPools
, runWorkerPoolsWith
, runSelectedWorkerPoolsWith
, withGauges
, startGauges
, 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
)
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)
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
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))
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
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
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
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))
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))
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)
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)