otel-effectful-1.0.0: src/Effectful/OpenTelemetry/Metrics/Effect.hs
{-# OPTIONS_GHC -Wno-redundant-constraints #-}
module Effectful.OpenTelemetry.Metrics.Effect
( -- * Effect
Metrics
, runMetrics
, runHttpMetrics
, runGrpcMetrics
, runConsoleMetrics
, runInMemoryMetrics
, runMetricsWith
, runNoMetrics
-- * Instruments
, register
)
where
import Control.Concurrent qualified as IO
import Control.Concurrent.STM qualified as STM
import Control.Monad (forever, unless)
import Control.Monad.Extra (concatMapM, ifM)
import Data.HashMap.Strict (HashMap)
import Data.HashMap.Strict qualified as HashMap
import Data.List qualified as List
import Data.Maybe (fromMaybe)
import Data.Text (Text)
import Effectful
import Effectful.Concurrent (Concurrent)
import Effectful.Concurrent.STM (TVar)
import Effectful.Dispatch.Static
import Effectful.Environment (Environment, lookupEnv)
import Effectful.Exception (bracket, finally)
import Effectful.Http2Client (HostName, PortNumber)
import Effectful.HttpClient (responseTimeoutDefault)
import Effectful.OpenTelemetry.Exporter (Exporter)
import Effectful.OpenTelemetry.Exporter qualified as Exporter
import Effectful.OpenTelemetry.Exporter.Console qualified as Console
import Effectful.OpenTelemetry.Exporter.Environment qualified as Exporter.Environment
import Effectful.OpenTelemetry.Exporter.OTLP (defaultGrpcSendTimeout)
import Effectful.OpenTelemetry.Metrics.Instrument (Instrument, SomeInstrument (..))
import Effectful.OpenTelemetry.Metrics.Instrument qualified as Instrument
import Effectful.OpenTelemetry.Metrics.Measurement (Measurement (..))
import Effectful.OpenTelemetry.Protocol
( Compression
, Encoding
, OTLP
, Resource
, Scope
, exportIO
, runInMemoryOTLP
, runNoOTLP
, runOTLPWith
)
import Effectful.OpenTelemetry.Protocol.Environment qualified as Environment
import Effectful.OpenTelemetry.Protocol.Export (exportSignalEnvName)
import Effectful.OpenTelemetry.Protocol.Export qualified as Export
import Effectful.OpenTelemetry.Timestamp (Timestamp)
import Effectful.OpenTelemetry.Timestamp qualified as Timestamp
import Effectful.Retry (Retry)
import Effectful.Timeout (Timeout)
import Network.URI (URI)
import Numeric.Natural (Natural)
import Text.Read (readMaybe)
import Prelude
data Metrics :: Effect
type instance DispatchOf Metrics = 'Static 'WithSideEffects
data instance StaticRep Metrics = Metrics
{ instruments :: TVar (HashMap Text SomeInstrument)
, lastExport :: TVar Timestamp
, export :: Measurement -> IO ()
}
-- | Adds an 'Instrument' to the pool of instruments that are periodically exported.
register :: (Metrics :> es, Instrument i) => i -> Eff es i
register i = do
Metrics{..} <- getStaticRep
unsafeEff_ . STM.atomically . STM.modifyTVar instruments $
HashMap.insert (Instrument.name i) (SomeInstrument i)
pure i
-- | Samples all instruments, then exports the samples that are newer than the last export time.
-- Updates the last export time to the latest exported sample time.
exportInstruments :: StaticRep Metrics -> IO ()
exportInstruments Metrics{..} = do
lastExport' <- STM.readTVarIO lastExport
samples <-
fmap (filter $ (lastExport' <) . fst)
. concatMapM Instrument.sample
. HashMap.elems
=<< STM.readTVarIO instruments
unless (null samples) do
let (timestamps, metrics) = unzip samples
mapM_ export metrics
STM.atomically . STM.writeTVar lastExport . maximum $ timestamps
newMetrics :: (IOE :> es, OTLP Measurement :> es) => Eff es (StaticRep Metrics)
newMetrics = do
export <- exportIO
instruments <- liftIO . STM.newTVarIO $ mempty
lastExport <- liftIO . STM.newTVarIO $ Timestamp.epoch
pure Metrics{..}
-- | Periodically call 'exportInstruments' in the background.
-- The background thread runs for the duration of the given action.
-- Exports all remaining instruments when done.
runMetricsThread
:: (OTLP Measurement :> es, IOE :> es)
=> Natural
-> Eff (Metrics ': es) a
-> Eff es a
runMetricsThread intervalMs eff = do
metrics <- newMetrics
bracket
( liftIO . IO.forkIO . forever $ do
IO.threadDelay . fromIntegral $ intervalMs * 1_000
exportInstruments metrics
)
( \threadId -> liftIO do
IO.killThread threadId
exportInstruments metrics
)
. const
. evalStaticRep metrics
$ eff
-- | Runs the given action and exports all instruments when done,
-- even when the action raises an exception.
runMetricsState
:: (OTLP Measurement :> es, IOE :> es)
=> Eff (Metrics ': es) a
-> Eff es a
runMetricsState eff = do
metrics <- newMetrics
evalStaticRep metrics $ inject eff `finally` liftIO (exportInstruments metrics)
defaultExportIntervalMs :: Natural
defaultExportIntervalMs = 60_000
-- | Run the 'Metrics' effect, sending telemetry to an exporter.
-- Reads the configuration from <https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#general-sdk-configuration the standard environment variables>.
-- Delegates to 'runNoMetrics' if @OTEL_SDK_DISABLED = true@.
runMetrics
:: ( IOE :> es
, Concurrent :> es
, Environment :> es
, Retry :> es
, Timeout :> es
)
=> Resource
-> Scope
-> Eff (Metrics ': es) a
-> Eff es a
runMetrics resource scope eff =
ifM Environment.isSdkDisabled (runNoMetrics eff) do
intervalMs <-
fromMaybe defaultExportIntervalMs
. (readMaybe =<<)
<$> lookupEnv (List.intercalate "_" ["OTEL", exportSignalEnvName @Measurement, "EXPORT_INTERVAL"])
exporter <- Environment.runConfigError $ Exporter.Environment.lookup @Measurement resource scope
runOTLPWith exporter
. runMetricsThread intervalMs
. inject
$ eff
-- | Run the 'Metrics' effect, sending telemetry to a collector at the given 'URI'
-- over HTTP with the given 'Encoding'.
runHttpMetrics
:: (IOE :> es, Concurrent :> es, Retry :> es, Timeout :> es)
=> Resource
-> Scope
-> Encoding
-> Compression
-> URI
-> Eff (Metrics ': es) a
-> Eff es a
runHttpMetrics resource scope encoding compression endpoint = do
runOTLPWith
( Exporter.http
@Measurement
resource
scope
encoding
endpoint
(Export.defaultConfig @Measurement)
compression
responseTimeoutDefault
)
. runMetricsThread defaultExportIntervalMs
. inject
-- | Run the 'Metrics' effect, sending telemetry to a gRPC collector.
runGrpcMetrics
:: (IOE :> es, Concurrent :> es, Retry :> es, Timeout :> es)
=> Resource
-> Scope
-> HostName
-> PortNumber
-> Compression
-> Eff (Metrics ': es) a
-> Eff es a
runGrpcMetrics resource scope host port compression = do
runOTLPWith
( Exporter.grpc @Measurement
resource
scope
host
port
(Export.exportGrpcRPC @Measurement)
(Export.defaultConfig @Measurement)
compression
defaultGrpcSendTimeout
)
. runMetricsThread defaultExportIntervalMs
. inject
-- | Run the 'Metrics' effect, printing telemetry to the console rather than sending to a collector.
runConsoleMetrics
:: (IOE :> es)
=> Maybe Natural
-- ^ Background sampling interval in ms.
-- If set, a background thread periodically samples instruments.
-- Otherwise, they are only sampled once at the end.
-> Eff (Metrics ': es) a
-> Eff es a
runConsoleMetrics intervalMs = runMetricsWith intervalMs Console.stdout
-- | Run the 'Metrics' effect, collecting telemetry in-memory rather than sending to a collector.
runInMemoryMetrics
:: (IOE :> es, Concurrent :> es)
=> Eff (Metrics ': es) a
-> Eff es (a, [Measurement])
runInMemoryMetrics = runInMemoryOTLP @Measurement . runMetricsState . inject
-- | Run the 'Metrics' effect with a given 'Exporter'.
runMetricsWith
:: (IOE :> es)
=> Maybe Natural
-- ^ Background sampling interval in ms.
-- If set, a background thread periodically samples instruments.
-- Otherwise, they are only sampled once at the end.
-> Exporter es Measurement
-> Eff (Metrics ': es) a
-> Eff es a
runMetricsWith intervalMs exporter =
runOTLPWith exporter . maybe runMetricsState runMetricsThread intervalMs . inject
-- | Run the 'Metrics' effect as a no-op action.
runNoMetrics :: (IOE :> es) => Eff (Metrics ': es) a -> Eff es a
runNoMetrics = runNoOTLP @Measurement . runMetricsState . inject