packages feed

otel-effectful-1.0.0: src/Effectful/OpenTelemetry/Metrics/Histogram.hs

{-# OPTIONS_GHC -Wno-redundant-constraints #-}

module Effectful.OpenTelemetry.Metrics.Histogram
    ( Histogram
    , new
    , newWithBounds
    , defaultExplicitBounds
    , record
    , recordIO
    , measurement
    , Payload (..)
    )
where

import Control.Concurrent.STM (TVar)
import Control.Concurrent.STM qualified as STM
import Data.Aeson.Types (ToJSON (..), object, (.=))
import Data.List qualified as List
import Data.Scientific (Scientific)
import Data.Text (Text)
import Data.Word (Word64)
import Effectful
import Effectful.Dispatch.Static (unsafeEff_)
import Effectful.OpenTelemetry.Metrics.Effect (Metrics)
import Effectful.OpenTelemetry.Metrics.Effect qualified as Metrics
import Effectful.OpenTelemetry.Metrics.Instrument (Instrument)
import Effectful.OpenTelemetry.Metrics.Instrument qualified as Instrument
import Effectful.OpenTelemetry.Metrics.Measurement
    ( AggregationTemporality (..)
    , HistogramDataPoint (..)
    , Measurement
    , Metric (..)
    )
import Effectful.OpenTelemetry.Metrics.Measurement qualified as Measurement
import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))
import Effectful.OpenTelemetry.Timestamp (Timestamp)
import Effectful.OpenTelemetry.Timestamp qualified as Timestamp
import GHC.Generics (Generic)
import Proto3.Wire.Encode.Class qualified as Proto
import Prelude

-- | Recommended default histogram bucket bounds for general latency-style data.
defaultExplicitBounds :: [Scientific]
defaultExplicitBounds = [0, 5, 10, 25, 50, 75, 100, 250, 500, 750, 1000, 2500, 5000, 7500, 10000]

-- | A synchronous 'Instrument' which reports arbitrary values that are likely to be statistically
-- meaningful.
-- It is intended for statistics such as histograms, summaries, and percentile.
--
-- Example uses for Histogram:
--
-- - the request duration
-- - the size of the response payload
--
-- See <https://opentelemetry.io/docs/specs/otel/metrics/data-model/#histogram the OpenTelemetry spec>.
data Histogram = Histogram
    { name :: Text
    , startTime :: Timestamp
    , explicitBounds :: [Scientific]
    , state :: TVar State
    , metadata :: Metadata
    }

data State = State
    { time :: Timestamp
    , count :: !Word64
    , sumValues :: !Scientific
    , bucketCounts :: ![Word64]
    , minValue :: !(Maybe Scientific)
    , maxValue :: !(Maybe Scientific)
    }

-- | Create a new 'Histogram' with 'defaultExplicitBounds'
-- and register it to be sampled and exported.
new :: (Metrics :> es) => Text -> Metadata -> Eff es Histogram
new = flip newWithBounds defaultExplicitBounds

-- | Create a new 'Histogram' with the specified bucket bounds
-- and register it to be sampled and exported.
newWithBounds :: (Metrics :> es) => Text -> [Scientific] -> Metadata -> Eff es Histogram
newWithBounds name (List.sort -> explicitBounds) metadata =
    Metrics.register
        =<< unsafeEff_ do
            startTime <- Timestamp.now
            state <-
                STM.newTVarIO
                    State
                        { time = startTime
                        , count = 0
                        , sumValues = 0
                        , bucketCounts = replicate (length explicitBounds + 1) 0
                        , minValue = Nothing
                        , maxValue = Nothing
                        }
            pure Histogram{..}

-- | Update the statistics with the specified amount.
record :: (Metrics :> es) => Histogram -> Scientific -> Eff es ()
record = (unsafeEff_ .) . recordIO

recordIO :: Histogram -> Scientific -> IO ()
recordIO histogram value = do
    time <- Timestamp.now
    STM.atomically . STM.modifyTVar' histogram.state $ \s ->
        State
            { time
            , count = s.count + 1
            , sumValues = s.sumValues + value
            , bucketCounts = bumpBucket histogram.explicitBounds s.bucketCounts
            , minValue = Just $ maybe value (min value) s.minValue
            , maxValue = Just $ maybe value (max value) s.maxValue
            }
  where
    bumpBucket :: [Scientific] -> [Word64] -> [Word64]
    bumpBucket [] (bucket : bs) = (bucket + 1) : bs
    bumpBucket (bound : moreBounds) (bucket : bs)
        | value <= bound = (bucket + 1) : bs
        | otherwise = bucket : bumpBucket moreBounds bs
    bumpBucket _ [] = []

-- | A histogram 'Measurement' of a single observed data point.
measurement :: Metadata -> Text -> HistogramDataPoint -> Measurement
measurement metadata name dataPoint =
    Measurement.create
        name
        metadata
        Payload
            { dataPoints = [dataPoint]
            , aggregationTemporality = TemporalityCumulative
            }

instance Instrument Histogram where
    name = name
    sample Histogram{metadata = metadata@Metadata{..}, ..} = do
        State{..} <- STM.readTVarIO state
        pure . pure $
            ( time
            , measurement
                metadata
                name
                HistogramDataPoint
                    { attributes
                    , startTime
                    , time
                    , count
                    , sum = if count == 0 then Nothing else Just sumValues
                    , bucketCounts
                    , explicitBounds
                    , minValue
                    , maxValue
                    }
            )

data Payload = Payload
    { dataPoints :: [HistogramDataPoint]
    , aggregationTemporality :: AggregationTemporality
    }
    deriving stock (Generic, Show, Eq)

instance ToJSON Payload where
    toJSON Payload{..} =
        object
            [ "dataPoints" .= dataPoints
            , "aggregationTemporality" .= fromEnum aggregationTemporality
            ]

instance Proto.Encode Payload where
    encode Payload{..} =
        mconcat
            [ foldMap (Proto.encodeField 1) dataPoints
            , Proto.encodeField 2 aggregationTemporality
            ]

instance Metric Payload where
    metricJsonKey _ = "histogram"
    metricProtoFieldNumber _ = 9