packages feed

otel-effectful-1.0.0: src/Effectful/OpenTelemetry/Protocol/Effect.hs

module Effectful.OpenTelemetry.Protocol.Effect {-# WARNING in "x-unstable-interface" "This is an unstable interface." #-} where

import Effectful
import Effectful.Concurrent (Concurrent)
import Effectful.Concurrent.STM (atomically, flushTQueue, newTQueueIO)
import Effectful.Dispatch.Static
import Effectful.OpenTelemetry.Exporter.STM qualified as Exporter
import Effectful.OpenTelemetry.Exporter.Type (Exporter (..), withExporter)
import Prelude

data OTLP a :: Effect

type instance DispatchOf (OTLP _) = 'Static 'WithSideEffects

newtype instance StaticRep (OTLP a) = OTLP {sink :: a -> IO ()}

-- | Run the 'OTLP' effect with a custom 'Exporter'.
-- Passes items to the exporter synchronously as they are emitted, without batching or retrying.
runOTLPWith :: forall a es r. (IOE :> es) => Exporter es a -> Eff (OTLP a ': es) r -> Eff es r
runOTLPWith Exporter{..} eff = withExporter \sink -> evalStaticRep OTLP{sink} eff

-- | Run the 'OTLP' effect, collecting telemetry in-memory rather than sending to a collector.
runInMemoryOTLP
    :: forall a es r
     . (IOE :> es, Concurrent :> es)
    => Eff (OTLP a ': es) r
    -> Eff es (r, [a])
runInMemoryOTLP eff = do
    queue <- newTQueueIO
    r <- runOTLPWith (Exporter.tqueue queue) eff
    items <- atomically $ flushTQueue queue
    pure (r, items)

-- | Run the 'OTLP' effect as a no-op action.
runNoOTLP :: forall a es r. (IOE :> es) => Eff (OTLP a ': es) r -> Eff es r
runNoOTLP = runOTLPWith mempty

export :: (OTLP a :> es) => a -> Eff es ()
export item = do
    sink <- exportIO
    unsafeEff_ $ sink item

exportIO :: forall a es. (OTLP a :> es) => Eff es (a -> IO ())
exportIO = do
    OTLP{..} <- getStaticRep
    pure sink