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