packages feed

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