otel-effectful-1.0.0: src/Effectful/OpenTelemetry/Protocol/GRPC.hs
{-# LANGUAGE TupleSections #-}
{-# OPTIONS_GHC -Wno-orphans #-}
module Effectful.OpenTelemetry.Protocol.GRPC {-# WARNING in "x-unstable-interface" "This is an unstable interface." #-} where
import Data.ByteString.Char8 qualified as ByteString
import Effectful
import Effectful.Error.Static (Error, throwError)
import Effectful.GrpcClient
( Decoding (..)
, Encoding (..)
, GrpcError (..)
, GrpcReply
, gzip
, open
, singleRequest
, uncompressed
)
import Effectful.GrpcClient qualified as Grpc
import Effectful.Http2Client
( ClientError
, HostName
, Http2Client
, PortNumber
, defaultGoAwayHandler
, ignoreFallbackHandler
, newHttp2FrameConnection
, runHttp2Client
)
import Effectful.OpenTelemetry.Protocol.Exception (ConnectionTimeout (..))
import Effectful.OpenTelemetry.Protocol.Transport (Compression (..))
import Effectful.Timeout (Timeout, timeout)
import Network.GRPC.HTTP2.Proto3Wire (Proto3WireEncoder (..), RPC)
import Network.URI (URI)
import Network.URI.Static (uri)
import Proto3.Wire.Encode (MessageBuilder)
import Prelude
type GRPC = Http2Client
-- | Default OTLP gRPC endpoint URI: @http:\/\/localhost:4317@.
defaultEndpoint :: URI
defaultEndpoint = [uri|http://localhost:4317|]
runGrpc
:: ( IOE :> es
, Error ConnectionTimeout :> es
, Error ClientError :> es
, Timeout :> es
)
=> HostName
-> PortNumber
-> Eff (GRPC ': es) a
-> Eff es a
runGrpc host port action = do
frame <-
maybe (throwError ConnectionTimeout) pure
=<< timeout 10_000_000 (newHttp2FrameConnection host port Nothing)
runHttp2Client frame 4096 4096 mempty defaultGoAwayHandler ignoreFallbackHandler action
instance Proto3WireEncoder () where
proto3WireEncode () = mempty
proto3WireDecode = pure ()
instance Proto3WireEncoder MessageBuilder where
proto3WireEncode = id
proto3WireDecode = error "The impossible happened: otel-effectful does not decode messages"
sendRpc
:: ( Http2Client :> es
, Error ClientError :> es
, Error GrpcError :> es
)
=> HostName
-> PortNumber
-> RPC
-> Compression
-> Grpc.Timeout
-> MessageBuilder
-> Eff es ()
sendRpc host port rpc compression grpcTimeout msg =
either (const $ throwError GrpcTooMuchConcurrency) (const @_ @(GrpcReply ()) $ pure ())
=<< open
authority
[]
grpcTimeout
(Encoding codec)
(Decoding codec)
(singleRequest rpc msg)
where
authority = ByteString.pack $ host <> ":" <> show port
codec = case compression of
NoCompression -> uncompressed
GZip -> gzip