packages feed

otel-effectful-1.0.0: test/Effectful/OpenTelemetry/Exporter/GrafanaSpec.hs

{-# LANGUAGE ConstraintKinds #-}
{-# OPTIONS_GHC -Wno-name-shadowing #-}
{-# OPTIONS_GHC -Wno-orphans #-}

module Effectful.OpenTelemetry.Exporter.GrafanaSpec where

import Control.Concurrent (forkIO, killThread)
import Control.Exception qualified as Exception
import Control.Monad (forever)
import Data.Aeson (Value (..))
import Data.Bifunctor (second)
import Data.ByteString qualified as ByteString
import Data.ByteString.Builder.Extra (defaultChunkSize)
import Data.Foldable (for_)
import Data.Functor (void, (<&>))
import Data.IORef (IORef, modifyIORef', newIORef, readIORef)
import Data.Maybe (fromMaybe)
import Data.Text qualified as Text
import Effectful
import Effectful.Concurrent (Concurrent, runConcurrent)
import Effectful.Concurrent.STM (atomically, flushTQueue, newTQueueIO)
import Effectful.Environment (Environment, lookupEnv, runEnvironment)
import Effectful.HUnit (HUnit)
import Effectful.Hspec
import Effectful.Http2Client (HostName, PortNumber)
import Effectful.HttpClient (responseTimeoutDefault)
import Effectful.OpenTelemetry.Exporter qualified as Exporter
import Effectful.OpenTelemetry.Exporter.Grafana.Loki qualified as Loki
import Effectful.OpenTelemetry.Exporter.Grafana.Mimir qualified as Mimir
import Effectful.OpenTelemetry.Exporter.Grafana.Polling (withListening)
import Effectful.OpenTelemetry.Exporter.Grafana.Tempo qualified as Tempo
import Effectful.OpenTelemetry.Exporter.STM qualified as Exporter.STM
import Effectful.OpenTelemetry.Logging
import Effectful.OpenTelemetry.Metrics
import Effectful.OpenTelemetry.Protocol
    ( Compression (..)
    , Encoding (..)
    , Resource (..)
    , Scope (..)
    )
import Effectful.OpenTelemetry.Protocol.Attributes qualified as Attributes
import Effectful.OpenTelemetry.Protocol.Export qualified as Export
import Effectful.OpenTelemetry.Protocol.Resource qualified as Resource
import Effectful.OpenTelemetry.Protocol.Scope qualified as Scope
import Effectful.OpenTelemetry.Protocol.Transport (Protocol (..))
import Effectful.OpenTelemetry.Tracing
import Effectful.OpenTelemetry.Tracing.Span.Kind qualified as Kind
import Effectful.QuickCheck (arbitrary, generate)
import Effectful.Retry (Retry, runRetry)
import Effectful.Timeout (Timeout, runTimeout, timeout)
import Network.Simple.TCP qualified as TCP
import Network.Socket (socketPort)
import Network.URI (URI (..), URIAuth (..), parseAbsoluteURI)
import Network.URI.Static (uri)
import Text.Read (readMaybe)
import Util
import Prelude hiding (log, span)

spec :: (IOE :> es, HUnit :> es, Hspec :> es) => Eff es ()
spec = do
    compressionSpec
    exportSpec

compressionSpec :: (IOE :> es, Hspec :> es) => Eff es ()
compressionSpec = runEnvironment . runConcurrent . runRetry . runTimeout . describe "Compression" . parallel $ do
    httpEndpoint <- alloyHttpEndpoint
    (grpcHost, grpcPort) <- alloyGrpcAuthority
    protocols <- alloyProtocols
    let httpHost = uriRegNameOrDefault httpEndpoint
        httpPort = uriPortOrDefault httpEndpoint
    for_ protocols \protocol -> do
        let (readyEndpoint, upstreamHost, upstreamPort) = case protocol of
                HTTP{} -> (httpEndpoint, httpHost, httpPort)
                GRPC{} -> (viaLocalPort httpEndpoint grpcPort, grpcHost, grpcPort)
            mkRun :: PortNumber -> Compression -> RunSignals_
            mkRun port = case protocol of
                HTTP encoding _ -> runFor . HTTP encoding $ viaLocalPort httpEndpoint port
                GRPC{} -> runFor $ GRPC "127.0.0.1" port (Export.exportGrpcRPC @Span)
        it (protocolLabel protocol) . withListening "Alloy" readyEndpoint $ do
            plain <- liftIO . bytesSentFor upstreamHost upstreamPort $ \port -> mkRun port NoCompression
            gzipped <- liftIO . bytesSentFor upstreamHost upstreamPort $ \port -> mkRun port GZip
            gzipped `shouldSatisfy` (> 0)
            (4 * gzipped) `shouldSatisfy` (< plain)
  where
    withByteCountingProxy
        :: HostName
        -> PortNumber
        -> (PortNumber -> IO Int -> IO a)
        -> IO a
    withByteCountingProxy upstreamHost upstreamPort body =
        TCP.listen (TCP.Host "127.0.0.1") "0" \(listening, _) -> do
            -- 'TCP.listen' reports the address it was asked to bind, whose port is
            -- still 0 for an ephemeral bind; the socket knows the real one.
            localPort <- socketPort listening
            sent <- newIORef 0
            let serve = forever $ TCP.acceptFork listening \(client, _) ->
                    TCP.connect upstreamHost (show upstreamPort) \(upstream, _) -> do
                        back <- forkIO $ relay upstream client Nothing
                        relay client upstream $ Just sent
                        killThread back
            Exception.bracket (forkIO serve) killThread . const $
                body localPort (readIORef sent)
      where
        relay :: TCP.Socket -> TCP.Socket -> Maybe (IORef Int) -> IO ()
        relay from to counter = Exception.handle @Exception.IOException (const $ pure ()) loop
          where
            loop =
                TCP.recv from defaultChunkSize >>= \case
                    Nothing -> pure ()
                    Just chunk -> do
                        for_ counter \c -> modifyIORef' c (+ ByteString.length chunk)
                        TCP.send to chunk
                        loop

    bytesSentFor :: HostName -> PortNumber -> (PortNumber -> RunSignals_) -> IO Int
    bytesSentFor host port mkRun =
        withByteCountingProxy host port \localPort readBytes -> do
            void
                . mkRun localPort testResource testScope
                . inSpan
                    "compression-probe"
                    Kind.Internal
                    (Attributes.fromList [("filler", String . Text.replicate 2_000 $ "filler")])
                $ pure ()
            readBytes

    viaLocalPort :: URI -> PortNumber -> URI
    viaLocalPort endpoint port =
        endpoint
            { uriAuthority =
                Just
                    URIAuth
                        { uriUserInfo = ""
                        , uriRegName = "127.0.0.1"
                        , uriPort = ':' : show port
                        }
            }

exportSpec :: (IOE :> es, HUnit :> es, Hspec :> es) => Eff es ()
exportSpec = runEnvironment . runConcurrent . runRetry . runTimeout . describe "Export" . parallel $ do
    protocols <- alloyProtocols
    endpoints <- readEndpoints
    for_ protocols (`grafanaSpec` endpoints)

data Endpoints = Endpoints
    { loki :: URI
    , tempo :: URI
    , mimir :: URI
    }

readEndpoints :: (Environment :> es) => Eff es Endpoints
readEndpoints = do
    loki <- envURI "LOKI_URI" [uri|http://localhost:3100|]
    tempo <- envURI "TEMPO_URI" [uri|http://localhost:3200|]
    mimir <- envURI "MIMIR_URI" [uri|http://localhost:3300|]
    pure Endpoints{..}

envURI :: (Environment :> es) => String -> URI -> Eff es URI
envURI var def =
    lookupEnv var <&> \case
        Nothing -> def
        Just raw -> fromMaybe (error $ "invalid " <> var <> ": " <> raw) $ parseAbsoluteURI raw

alloyHttpEndpoint :: (Environment :> es) => Eff es URI
alloyHttpEndpoint = envURI "ALLOY_OTLP_HTTP_ENDPOINT" [uri|http://localhost:4318|]

alloyGrpcAuthority :: (Environment :> es) => Eff es (HostName, PortNumber)
alloyGrpcAuthority =
    lookupEnv "ALLOY_OTLP_GRPC_AUTHORITY" <&> \case
        Nothing -> ("127.0.0.1", 4317)
        Just authority -> case break (== ':') authority of
            (h, ':' : p) -> (h, read p)
            _ -> error $ "invalid ALLOY_OTLP_GRPC_AUTHORITY: " <> authority

alloyProtocols :: (Environment :> es) => Eff es [Protocol]
alloyProtocols = do
    httpEndpoint <- alloyHttpEndpoint
    (grpcHost, grpcPort) <- alloyGrpcAuthority
    pure
        [ HTTP Json httpEndpoint
        , HTTP Proto httpEndpoint
        , GRPC grpcHost grpcPort (Export.exportGrpcRPC @Span)
        ]

testResource :: Resource
testResource =
    Resource
        { attributes = Attributes.fromList [("service.name", "otel-effectful-test")]
        }

testScope :: Scope
testScope = Scope{name = "grafana-spec", version = "0.1.0", attributes = mempty}

data Telemetry = Telemetry
    { spans :: [Span]
    , logs :: [LogRecord]
    , measurements :: [Measurement]
    }
    deriving stock (Show)

type SignalStack = '[Metrics, Logging, Tracing, Environment, Timeout, Retry, Concurrent, IOE]

type RunSignals = forall a. Resource -> Scope -> Eff SignalStack a -> IO (a, Telemetry)

type RunSignals_ = Resource -> Scope -> Eff SignalStack () -> IO ((), Telemetry)

type ExporterCtx sig es =
    (Export.Request sig, IOE :> es, Concurrent :> es, Retry :> es, Timeout :> es)

type MkExporter =
    forall sig es
     . (ExporterCtx sig es)
    => Resource
    -> Scope
    -> Exporter.Exporter es sig

httpExporter
    :: forall sig es
     . (ExporterCtx sig es)
    => URI
    -> Encoding
    -> Compression
    -> Resource
    -> Scope
    -> Exporter.Exporter es sig
httpExporter endpoint encoding compression res scope =
    Exporter.http
        res
        scope
        encoding
        endpoint
        (Export.defaultConfig @sig)
        compression
        responseTimeoutDefault

grpcExporter
    :: forall sig es
     . (ExporterCtx sig es)
    => HostName
    -> PortNumber
    -> Compression
    -> Resource
    -> Scope
    -> Exporter.Exporter es sig
grpcExporter host port compression res scope =
    Exporter.grpc
        res
        scope
        host
        port
        (Export.exportGrpcRPC @sig)
        (Export.defaultConfig @sig)
        compression
        Exporter.defaultGrpcSendTimeout

runSignals :: MkExporter -> RunSignals
runSignals mkExporter res scope act =
    runEff . runConcurrent . runRetry . runTimeout . runEnvironment $ do
        spanQueue <- newTQueueIO
        logQueue <- newTQueueIO
        measurementQueue <- newTQueueIO
        result <-
            runTracingWith (Exporter.STM.tqueue spanQueue <> mkExporter res scope)
                . runLoggingWith (Exporter.STM.tqueue logQueue <> mkExporter res scope)
                . runMetricsWith (Just 60_000) (Exporter.STM.tqueue measurementQueue <> mkExporter res scope)
                $ act
        spans <- atomically $ flushTQueue spanQueue
        logs <- atomically $ flushTQueue logQueue
        measurements <- atomically $ flushTQueue measurementQueue
        pure (result, Telemetry{spans, logs, measurements})

runFor :: Protocol -> Compression -> RunSignals
runFor (HTTP encoding endpoint) compression = runSignals $ httpExporter endpoint encoding compression
runFor (GRPC host port _) compression = runSignals $ grpcExporter host port compression

grafanaSpec
    :: forall es
     . (IOE :> es, HUnit :> es, Hspec :> es, Retry :> es, Timeout :> es)
    => Protocol
    -> Endpoints
    -> Eff es ()
grafanaSpec protocol endpoints =
    describe (protocolLabel protocol)
        . modifyMaxSuccess (const 10)
        . parallel
        $ do
            (resourceAttrs, scopeAttrs) <- liftIO $ generate arbitrary
            let resource = Resource{attributes = testResource.attributes <> resourceAttrs}
                scope =
                    Scope
                        { name = testScope.name
                        , version = testScope.version
                        , attributes = scopeAttrs
                        }
                run :: RunSignals
                run = runFor protocol GZip
                runTest
                    :: forall a
                     . Eff '[Metrics, Logging, Tracing, Environment, Timeout, Retry, Concurrent, IOE] a
                    -> Eff es (a, Telemetry)
                runTest act =
                    timeout 10_000_000 (liftIO $ run resource scope act) >>= \case
                        Just result -> pure result
                        Nothing -> expectationFailure "signal push timed out" >> error "unreachable"

            Tempo.spec protocol endpoints.tempo resource scope $ fmap (second spans) . runTest
            Loki.spec protocol endpoints.loki $ fmap (second logs) . runTest
            Mimir.spec protocol endpoints.mimir $ fmap (second measurements) . runTest

uriRegNameOrDefault :: URI -> HostName
uriRegNameOrDefault endpoint = case endpoint.uriAuthority of
    Just URIAuth{uriRegName} | not (null uriRegName) -> uriRegName
    _ -> "127.0.0.1"

uriPortOrDefault :: URI -> PortNumber
uriPortOrDefault endpoint = case endpoint.uriAuthority of
    Just URIAuth{uriPort = ':' : (readMaybe -> Just port)} -> port
    _ -> if endpoint.uriScheme == "https:" then 443 else 80