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