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