packages feed

otel-effectful-1.0.0: src/Effectful/OpenTelemetry/Exporter/OTLP.hs

{-# OPTIONS_GHC -Wno-name-shadowing #-}

module Effectful.OpenTelemetry.Exporter.OTLP
    ( otlp
    , http
    , grpc
    , defaultSendTimeoutSeconds
    , defaultGrpcSendTimeout
    )
where

import Control.Concurrent.STM qualified as STM
import Control.Monad (forever, unless)
import Control.Monad.Extra (fromMaybeM, ifM)
import Data.Aeson qualified as Aeson
import Data.ByteString.Lazy (LazyByteString)
import Data.Functor (void)
import Effectful
import Effectful.Concurrent (Concurrent, threadDelay)
import Effectful.Concurrent.Async (Async, async)
import Effectful.Concurrent.Async qualified as Async
import Effectful.Concurrent.STM
    ( STM
    , TBQueue
    , atomically
    , flushTBQueue
    , isFullTBQueue
    , newTBQueueIO
    , tryReadTBQueue
    , writeTBQueue
    )
import Effectful.Environment (Environment)
import Effectful.Error.Static (runErrorNoCallStackWith)
import Effectful.Exception (bracket, handleSync, throwIO, trySync)
import Effectful.GrpcClient qualified as Grpc
import Effectful.Http2Client (HostName, PortNumber)
import Effectful.HttpClient
    ( ResponseTimeout
    , responseTimeoutMicro
    , responseTimeoutNone
    , runHttpClientTls
    )
import Effectful.OpenTelemetry.Exporter.Type (Exporter (..))
import Effectful.OpenTelemetry.Protocol.Environment
    ( ConfigError
    , detectResource
    , isSdkDisabled
    , lookupExportConfig
    , lookupTransport
    )
import Effectful.OpenTelemetry.Protocol.Exception
    ( ConnectionTimeout
    , OTLPClientError (..)
    , OTLPGrpcError (..)
    , OTLPHttpException (..)
    , SomeOTLPException (..)
    )
import Effectful.OpenTelemetry.Protocol.Export qualified as Export
import Effectful.OpenTelemetry.Protocol.GRPC qualified as GRPC
import Effectful.OpenTelemetry.Protocol.HTTP qualified as HTTP
import Effectful.OpenTelemetry.Protocol.Resource (Resource)
import Effectful.OpenTelemetry.Protocol.Resource qualified as Resource
import Effectful.OpenTelemetry.Protocol.Scope (Scope)
import Effectful.OpenTelemetry.Protocol.Transport
    ( Compression
    , Encoding (..)
    , Protocol (..)
    , Transport (..)
    )
import Effectful.Retry
    ( Retry
    , RetryPolicyM
    , capDelay
    , fullJitterBackoff
    , limitRetriesByCumulativeDelay
    , recoverAll
    )
import Effectful.Timeout (Timeout, timeout)
import Network.GRPC.HTTP2.Proto3Wire (RPC)
import Network.URI (URI)
import Numeric.Natural (Natural)
import Proto3.Wire.Encode qualified as Proto3.Encode
import Prelude

-- | A collector configured from
-- <https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#general-sdk-configuration the standard environment variables>.
--
-- 'mempty' if @OTEL_SDK_DISABLED = true@.
otlp
    :: forall a es
     . ( Export.Request a
       , IOE :> es
       , Concurrent :> es
       , Environment :> es
       , Retry :> es
       , Timeout :> es
       )
    => Resource
    -> Scope
    -> Exporter es a
otlp resource scope = Exporter \liftExporter ->
    ifM isSdkDisabled (liftExporter mempty) $
        runErrorNoCallStackWith (throwIO . SomeOTLPException @ConfigError) do
            resource <- (resource <>) <$> detectResource
            Transport{..} <- lookupTransport @a
            exportConfig <- lookupExportConfig @a $ Export.defaultConfig @a
            withExporter
                ( case protocol of
                    HTTP encoding endpoint ->
                        http @a
                            resource
                            scope
                            encoding
                            endpoint
                            exportConfig
                            compression
                            ( if timeoutMs == 0
                                then responseTimeoutNone
                                else responseTimeoutMicro $ fromIntegral timeoutMs * 1_000
                            )
                    GRPC host port rpc ->
                        grpc @a
                            resource
                            scope
                            host
                            port
                            rpc
                            exportConfig
                            compression
                            (Grpc.Timeout $ fromIntegral timeoutMs `div` 1_000)
                )
                (inject . liftExporter)

http
    :: forall a es
     . ( Export.Request a
       , IOE :> es
       , Concurrent :> es
       , Retry :> es
       , Timeout :> es
       )
    => Resource
    -> Scope
    -> Encoding
    -> URI
    -> Export.Config
    -> Compression
    -> ResponseTimeout
    -> Exporter es a
http resource scope encoding endpoint exportConfig compression responseTimeout =
    Exporter $ export exportConfig sender
  where
    encodeBatch :: Encoding -> [Resource.Items a] -> LazyByteString
    encodeBatch Json items = Aeson.encode (Export.exportJson @a items)
    encodeBatch Proto items = Proto3.Encode.toLazyByteString (Export.exportProto @a items)
    signalEndpoint = Export.appendExportPath @a endpoint
    sender :: Exporter es [a]
    sender = Exporter \send ->
        runHttpClientTls $ withEffToIO (ConcUnlift Persistent Unlimited) \unlift ->
            unlift . inject . send $
                unlift
                    . runErrorNoCallStackWith (throwIO . SomeOTLPException . OTLPHttpException)
                    . HTTP.sendPayload
                        signalEndpoint
                        encoding
                        compression
                        responseTimeout
                    . encodeBatch encoding
                    . Resource.wrapItems resource scope

grpc
    :: forall a es
     . ( Export.Request a
       , IOE :> es
       , Concurrent :> es
       , Retry :> es
       , Timeout :> es
       )
    => Resource
    -> Scope
    -> HostName
    -> PortNumber
    -> RPC
    -> Export.Config
    -> Compression
    -> Grpc.Timeout
    -> Exporter es a
grpc resource scope host port rpc exportConfig compression timeout =
    Exporter $ export exportConfig sender
  where
    sender :: Exporter es [a]
    sender = Exporter \send ->
        runErrorNoCallStackWith (throwIO . SomeOTLPException @ConnectionTimeout)
            . runErrorNoCallStackWith (throwIO . SomeOTLPException . OTLPClientError)
            . GRPC.runGrpc host port
            $ withEffToIO (ConcUnlift Persistent Unlimited) \unlift ->
                unlift . inject . send $
                    unlift
                        . runErrorNoCallStackWith (throwIO . SomeOTLPException . OTLPGrpcError)
                        . GRPC.sendRpc host port rpc compression timeout
                        . Export.exportProto @a
                        . Resource.wrapItems resource scope

defaultSendTimeoutSeconds :: Natural
defaultSendTimeoutSeconds = 30

defaultGrpcSendTimeout :: Grpc.Timeout
defaultGrpcSendTimeout = Grpc.Timeout $ fromIntegral defaultSendTimeoutSeconds

withTimeout :: (Timeout :> es) => Natural -> ([a] -> Eff es ()) -> [a] -> Eff es ()
withTimeout timeoutMs send = fromMaybeM (throwIO $ userError "timeout") . timeout (fromIntegral timeoutMs * 1_000) . send

-- | Retry an action with capped exponential backoff, giving up without rethrowing
-- once retries are exhausted.
retrying :: forall es. (IOE :> es, Retry :> es) => Eff es () -> Eff es ()
retrying =
    handleSync (const $ pure ())
        . recoverAll retryPolicy
        . const
  where
    retryPolicy :: RetryPolicyM (Eff es)
    retryPolicy =
        limitRetriesByCumulativeDelay maxElapsedRetryTimeµs
            . capDelay maxRetryIntervalµs
            $ fullJitterBackoff initialRetryIntervalµs

    -- NOTE: There is no specification for configuring retry parameters,
    -- so we hard code the defaults.
    initialRetryIntervalµs :: Int
    initialRetryIntervalµs = 5_000_000 -- 5 s
    maxRetryIntervalµs :: Int
    maxRetryIntervalµs = 30_000_000 -- 30 s
    maxElapsedRetryTimeµs :: Int
    maxElapsedRetryTimeµs = 300_000_000 -- 300 s

-- | Send items to the batch exporter, to be exported either individually or in batches,
-- depending on the given 'Export.Config'.
-- Items may be dropped on a full queue, exhausted retries, or timeout.
export
    :: forall es a b
     . ( Retry :> es
       , Concurrent :> es
       , IOE :> es
       , Timeout :> es
       )
    => Export.Config
    -> Exporter es [a]
    -> ((a -> IO ()) -> Eff es b)
    -> Eff es b
export Export.Config{batch = Nothing, ..} exporter action =
    withEffToIO (ConcUnlift Persistent Unlimited) \unlift ->
        unlift . action $ \item ->
            unlift . retrying . withExporter exporter $ \export ->
                withTimeout exportTimeoutMs (liftIO . export) [item]
export Export.Config{batch = Just Export.BatchConfig{..}, ..} exporter action = do
    queue <- newTBQueueIO maxQueueSize
    let worker :: Eff es (Async a)
        worker =
            async . forever . retrying . withExporter exporter $
                forever . \send -> do
                    threadDelay . fromIntegral $ 1_000 * scheduledDelayMs
                    items <- atomically $ drainN maxBatchSize queue
                    unless (null items) . withTimeout exportTimeoutMs (liftIO . send) $ items

        enqueue :: a -> IO ()
        enqueue item = STM.atomically do
            full <- isFullTBQueue queue
            unless full $ writeTBQueue queue item

        flush :: Async a -> Eff es ()
        flush task = do
            Async.cancel task
            remaining <- atomically $ flushTBQueue queue
            -- NOTE: We don't retry when flushing because we don't want to block
            -- program termination. 'withTimeout' still caps how long this
            -- can take, so a stuck connection can't hang shutdown either.
            -- To properly protect against data loss if the collector itself crashes,
            -- we may want to consider implementing persistent storage:
            -- https://opentelemetry.io/docs/collector/resiliency/#persistent-storage-write-ahead-log---wal
            unless (null remaining)
                . void
                . trySync
                $ withExporter exporter \export -> withTimeout exportTimeoutMs (liftIO . export) remaining

    bracket worker flush . const $ action enqueue
  where
    drainN :: Natural -> TBQueue a -> STM [a]
    drainN = (fmap reverse .) . drainNReverse

    drainNReverse :: Natural -> TBQueue a -> STM [a]
    drainNReverse 0 _ = pure []
    drainNReverse n queue =
        tryReadTBQueue queue >>= \case
            Nothing -> pure []
            Just x -> (x :) <$> drainNReverse (n - 1) queue