packages feed

otel-effectful (empty) → 1.0.0

raw patch · 85 files changed

+7937/−0 lines, 85 filesdep +aesondep +ansi-terminaldep +basesetup-changed

Dependencies added: aeson, ansi-terminal, base, base16-bytestring, base64-bytestring, bytestring, composition-extra, containers, effectful, extra, hspec-effectful, http-client, http-client-effectful, http-types, http2-client-effectful, http2-client-grpc-effectful, http2-grpc-proto3-wire, hunit-effectful, network, network-simple, network-uri, otel-effectful, prettyprinter, prettyprinter-ansi-terminal, proto3-wire, quickcheck-effectful, random, retry-effectful, safe, scientific, stm, text, time, unordered-containers, vector, zlib

Files

+ CHANGELOG.md view
@@ -0,0 +1,12 @@+# Changelog++All notable changes to this project will be documented in this file.++The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/),+and this project adheres to the [Haskell Package Versioning Policy](https://pvp.haskell.org/).++## [1.0.0] - 2026-10-01++### Added++- Initial release.
+ LICENCE view
@@ -0,0 +1,287 @@+                      EUROPEAN UNION PUBLIC LICENCE v. 1.2+                      EUPL © the European Union 2007, 2016++This European Union Public Licence (the ‘EUPL’) applies to the Work (as defined+below) which is provided under the terms of this Licence. Any use of the Work,+other than as authorised under this Licence is prohibited (to the extent such+use is covered by a right of the copyright holder of the Work).++The Work is provided under the terms of this Licence when the Licensor (as+defined below) has placed the following notice immediately following the+copyright notice for the Work:++        Licensed under the EUPL++or has expressed by any other means his willingness to license under the EUPL.++1. Definitions++In this Licence, the following terms have the following meaning:++- ‘The Licence’: this Licence.++- ‘The Original Work’: the work or software distributed or communicated by the+  Licensor under this Licence, available as Source Code and also as Executable+  Code as the case may be.++- ‘Derivative Works’: the works or software that could be created by the+  Licensee, based upon the Original Work or modifications thereof. This Licence+  does not define the extent of modification or dependence on the Original Work+  required in order to classify a work as a Derivative Work; this extent is+  determined by copyright law applicable in the country mentioned in Article 15.++- ‘The Work’: the Original Work or its Derivative Works.++- ‘The Source Code’: the human-readable form of the Work which is the most+  convenient for people to study and modify.++- ‘The Executable Code’: any code which has generally been compiled and which is+  meant to be interpreted by a computer as a program.++- ‘The Licensor’: the natural or legal person that distributes or communicates+  the Work under the Licence.++- ‘Contributor(s)’: any natural or legal person who modifies the Work under the+  Licence, or otherwise contributes to the creation of a Derivative Work.++- ‘The Licensee’ or ‘You’: any natural or legal person who makes any usage of+  the Work under the terms of the Licence.++- ‘Distribution’ or ‘Communication’: any act of selling, giving, lending,+  renting, distributing, communicating, transmitting, or otherwise making+  available, online or offline, copies of the Work or providing access to its+  essential functionalities at the disposal of any other natural or legal+  person.++2. Scope of the rights granted by the Licence++The Licensor hereby grants You a worldwide, royalty-free, non-exclusive,+sublicensable licence to do the following, for the duration of copyright vested+in the Original Work:++- use the Work in any circumstance and for all usage,+- reproduce the Work,+- modify the Work, and make Derivative Works based upon the Work,+- communicate to the public, including the right to make available or display+  the Work or copies thereof to the public and perform publicly, as the case may+  be, the Work,+- distribute the Work or copies thereof,+- lend and rent the Work or copies thereof,+- sublicense rights in the Work or copies thereof.++Those rights can be exercised on any media, supports and formats, whether now+known or later invented, as far as the applicable law permits so.++In the countries where moral rights apply, the Licensor waives his right to+exercise his moral right to the extent allowed by law in order to make effective+the licence of the economic rights here above listed.++The Licensor grants to the Licensee royalty-free, non-exclusive usage rights to+any patents held by the Licensor, to the extent necessary to make use of the+rights granted on the Work under this Licence.++3. Communication of the Source Code++The Licensor may provide the Work either in its Source Code form, or as+Executable Code. If the Work is provided as Executable Code, the Licensor+provides in addition a machine-readable copy of the Source Code of the Work+along with each copy of the Work that the Licensor distributes or indicates, in+a notice following the copyright notice attached to the Work, a repository where+the Source Code is easily and freely accessible for as long as the Licensor+continues to distribute or communicate the Work.++4. Limitations on copyright++Nothing in this Licence is intended to deprive the Licensee of the benefits from+any exception or limitation to the exclusive rights of the rights owners in the+Work, of the exhaustion of those rights or of other applicable limitations+thereto.++5. Obligations of the Licensee++The grant of the rights mentioned above is subject to some restrictions and+obligations imposed on the Licensee. Those obligations are the following:++Attribution right: The Licensee shall keep intact all copyright, patent or+trademarks notices and all notices that refer to the Licence and to the+disclaimer of warranties. The Licensee must include a copy of such notices and a+copy of the Licence with every copy of the Work he/she distributes or+communicates. The Licensee must cause any Derivative Work to carry prominent+notices stating that the Work has been modified and the date of modification.++Copyleft clause: If the Licensee distributes or communicates copies of the+Original Works or Derivative Works, this Distribution or Communication will be+done under the terms of this Licence or of a later version of this Licence+unless the Original Work is expressly distributed only under this version of the+Licence — for example by communicating ‘EUPL v. 1.2 only’. The Licensee+(becoming Licensor) cannot offer or impose any additional terms or conditions on+the Work or Derivative Work that alter or restrict the terms of the Licence.++Compatibility clause: If the Licensee Distributes or Communicates Derivative+Works or copies thereof based upon both the Work and another work licensed under+a Compatible Licence, this Distribution or Communication can be done under the+terms of this Compatible Licence. For the sake of this clause, ‘Compatible+Licence’ refers to the licences listed in the appendix attached to this Licence.+Should the Licensee's obligations under the Compatible Licence conflict with+his/her obligations under this Licence, the obligations of the Compatible+Licence shall prevail.++Provision of Source Code: When distributing or communicating copies of the Work,+the Licensee will provide a machine-readable copy of the Source Code or indicate+a repository where this Source will be easily and freely available for as long+as the Licensee continues to distribute or communicate the Work.++Legal Protection: This Licence does not grant permission to use the trade names,+trademarks, service marks, or names of the Licensor, except as required for+reasonable and customary use in describing the origin of the Work and+reproducing the content of the copyright notice.++6. Chain of Authorship++The original Licensor warrants that the copyright in the Original Work granted+hereunder is owned by him/her or licensed to him/her and that he/she has the+power and authority to grant the Licence.++Each Contributor warrants that the copyright in the modifications he/she brings+to the Work are owned by him/her or licensed to him/her and that he/she has the+power and authority to grant the Licence.++Each time You accept the Licence, the original Licensor and subsequent+Contributors grant You a licence to their contributions to the Work, under the+terms of this Licence.++7. Disclaimer of Warranty++The Work is a work in progress, which is continuously improved by numerous+Contributors. It is not a finished work and may therefore contain defects or+‘bugs’ inherent to this type of development.++For the above reason, the Work is provided under the Licence on an ‘as is’ basis+and without warranties of any kind concerning the Work, including without+limitation merchantability, fitness for a particular purpose, absence of defects+or errors, accuracy, non-infringement of intellectual property rights other than+copyright as stated in Article 6 of this Licence.++This disclaimer of warranty is an essential part of the Licence and a condition+for the grant of any rights to the Work.++8. Disclaimer of Liability++Except in the cases of wilful misconduct or damages directly caused to natural+persons, the Licensor will in no event be liable for any direct or indirect,+material or moral, damages of any kind, arising out of the Licence or of the use+of the Work, including without limitation, damages for loss of goodwill, work+stoppage, computer failure or malfunction, loss of data or any commercial+damage, even if the Licensor has been advised of the possibility of such damage.+However, the Licensor will be liable under statutory product liability laws as+far such laws apply to the Work.++9. Additional agreements++While distributing the Work, You may choose to conclude an additional agreement,+defining obligations or services consistent with this Licence. However, if+accepting obligations, You may act only on your own behalf and on your sole+responsibility, not on behalf of the original Licensor or any other Contributor,+and only if You agree to indemnify, defend, and hold each Contributor harmless+for any liability incurred by, or claims asserted against such Contributor by+the fact You have accepted any warranty or additional liability.++10. Acceptance of the Licence++The provisions of this Licence can be accepted by clicking on an icon ‘I agree’+placed under the bottom of a window displaying the text of this Licence or by+affirming consent in any other similar way, in accordance with the rules of+applicable law. Clicking on that icon indicates your clear and irrevocable+acceptance of this Licence and all of its terms and conditions.++Similarly, you irrevocably accept this Licence and all of its terms and+conditions by exercising any rights granted to You by Article 2 of this Licence,+such as the use of the Work, the creation by You of a Derivative Work or the+Distribution or Communication by You of the Work or copies thereof.++11. Information to the public++In case of any Distribution or Communication of the Work by means of electronic+communication by You (for example, by offering to download the Work from a+remote location) the distribution channel or media (for example, a website) must+at least provide to the public the information requested by the applicable law+regarding the Licensor, the Licence and the way it may be accessible, concluded,+stored and reproduced by the Licensee.++12. Termination of the Licence++The Licence and the rights granted hereunder will terminate automatically upon+any breach by the Licensee of the terms of the Licence.++Such a termination will not terminate the licences of any person who has+received the Work from the Licensee under the Licence, provided such persons+remain in full compliance with the Licence.++13. Miscellaneous++Without prejudice of Article 9 above, the Licence represents the complete+agreement between the Parties as to the Work.++If any provision of the Licence is invalid or unenforceable under applicable+law, this will not affect the validity or enforceability of the Licence as a+whole. Such provision will be construed or reformed so as necessary to make it+valid and enforceable.++The European Commission may publish other linguistic versions or new versions of+this Licence or updated versions of the Appendix, so far this is required and+reasonable, without reducing the scope of the rights granted by the Licence. New+versions of the Licence will be published with a unique version number.++All linguistic versions of this Licence, approved by the European Commission,+have identical value. Parties can take advantage of the linguistic version of+their choice.++14. Jurisdiction++Without prejudice to specific agreement between parties,++- any litigation resulting from the interpretation of this License, arising+  between the European Union institutions, bodies, offices or agencies, as a+  Licensor, and any Licensee, will be subject to the jurisdiction of the Court+  of Justice of the European Union, as laid down in article 272 of the Treaty on+  the Functioning of the European Union,++- any litigation arising between other parties and resulting from the+  interpretation of this License, will be subject to the exclusive jurisdiction+  of the competent court where the Licensor resides or conducts its primary+  business.++15. Applicable Law++Without prejudice to specific agreement between parties,++- this Licence shall be governed by the law of the European Union Member State+  where the Licensor has his seat, resides or has his registered office,++- this licence shall be governed by Belgian law if the Licensor has no seat,+  residence or registered office inside a European Union Member State.++Appendix++‘Compatible Licences’ according to Article 5 EUPL are:++- GNU General Public License (GPL) v. 2, v. 3+- GNU Affero General Public License (AGPL) v. 3+- Open Software License (OSL) v. 2.1, v. 3.0+- Eclipse Public License (EPL) v. 1.0+- CeCILL v. 2.0, v. 2.1+- Mozilla Public Licence (MPL) v. 2+- GNU Lesser General Public Licence (LGPL) v. 2.1, v. 3+- Creative Commons Attribution-ShareAlike v. 3.0 Unported (CC BY-SA 3.0) for+  works other than software+- European Union Public Licence (EUPL) v. 1.1, v. 1.2+- Québec Free and Open-Source Licence — Reciprocity (LiLiQ-R) or Strong+  Reciprocity (LiLiQ-R+).++The European Commission may update this Appendix to later versions of the above+licences without producing a new version of the EUPL, as long as they provide+the rights granted in Article 2 of this Licence and protect the covered Source+Code from exclusive appropriation.++All other changes or additions to this Appendix require the production of a new+EUPL version.
+ Setup.hs view
@@ -0,0 +1,2 @@+import Distribution.Simple+main = defaultMain
+ otel-effectful.cabal view
@@ -0,0 +1,210 @@+cabal-version: 3.0+name: otel-effectful+version: 1.0.0+synopsis: OpenTelemetry implementation for the Effectful ecosystem+description:+  Implementation of @<https://opentelemetry.io/ OpenTelemetry>@ for the @<https://hackage.haskell.org/package/effectful effectful>@ ecosystem.++homepage: https://digital-autonomy.institute+license: EUPL-1.2+license-file: LICENCE+author: IDA+maintainer: IDA+bug-reports: https://issues.digital-autonomy.institute+category: Test+build-type: Simple+extra-doc-files:+  CHANGELOG.md++common common+  default-language: Haskell2010+  ghc-options:+    -Weverything+    -Wno-unsafe+    -Wno-missing-safe-haskell-mode+    -Wno-missing-export-lists+    -Wno-missing-import-lists+    -Wno-missing-kind-signatures+    -Wno-all-missed-specialisations+    -Wno-missing-role-annotations+    -Wno-x-unstable-interface++  default-extensions:+    ApplicativeDo+    BangPatterns+    BlockArguments+    DataKinds+    DefaultSignatures+    DeriveAnyClass+    DeriveGeneric+    DerivingStrategies+    DerivingVia+    ExistentialQuantification+    ExplicitNamespaces+    FlexibleContexts+    FlexibleInstances+    GeneralizedNewtypeDeriving+    ImportQualifiedPost+    LambdaCase+    MultiParamTypeClasses+    NamedFieldPuns+    NoImplicitPrelude+    NumericUnderscores+    OverloadedLabels+    OverloadedRecordDot+    OverloadedStrings+    QuasiQuotes+    RankNTypes+    RecordWildCards+    RecursiveDo+    ScopedTypeVariables+    TupleSections+    TypeApplications+    TypeFamilies+    TypeOperators+    ViewPatterns++  build-depends:+    aeson >=2.2 && <2.3,+    base >=4.16 && <5,+    containers >=0.7 && <0.9,+    effectful >=2.6 && <2.8,+    extra >=1.0 && <1.9,+    http-client-effectful >=1.0 && <1.1,+    http-types >=0.12 && <0.13,+    http2-client-effectful >=1.0 && <1.1,+    network-uri >=2.6 && <2.7,+    random >=1.2 && <1.3,+    retry-effectful >=0.1 && <0.2,+    scientific >=0.3 && <0.4,+    text >=2.1 && <2.2,++library+  import: common+  hs-source-dirs: src+  build-depends:+    ansi-terminal >=1.0 && <1.2,+    base16-bytestring >=1.0 && <1.1,+    bytestring >=0.12 && <0.13,+    http2-client-grpc-effectful >=1.0 && <1.1,+    http2-grpc-proto3-wire >=0.1 && <0.2,+    prettyprinter >=1.7 && <1.8,+    prettyprinter-ansi-terminal >=1.1 && <1.2,+    proto3-wire >=1.4 && <1.5,+    safe >=0.3 && <0.4,+    stm >=2.5 && <2.6,+    time >=1.12 && <1.17,+    unordered-containers >=0.2 && <0.3,+    vector >=0.13 && <0.14,+    zlib >=0.7 && <0.8,++  exposed-modules:+    Effectful.OpenTelemetry+    Effectful.OpenTelemetry.Exporter+    Effectful.OpenTelemetry.Exporter.Console+    Effectful.OpenTelemetry.Exporter.Environment+    Effectful.OpenTelemetry.Exporter.OTLP+    Effectful.OpenTelemetry.Exporter.STM+    Effectful.OpenTelemetry.Exporter.Type+    Effectful.OpenTelemetry.Logging+    Effectful.OpenTelemetry.Logging.Effect+    Effectful.OpenTelemetry.Logging.LogRecord+    Effectful.OpenTelemetry.Logging.Severity+    Effectful.OpenTelemetry.Metrics+    Effectful.OpenTelemetry.Metrics.Counter+    Effectful.OpenTelemetry.Metrics.Effect+    Effectful.OpenTelemetry.Metrics.Gauge+    Effectful.OpenTelemetry.Metrics.Histogram+    Effectful.OpenTelemetry.Metrics.Instrument+    Effectful.OpenTelemetry.Metrics.Measurement+    Effectful.OpenTelemetry.Metrics.Metadata+    Effectful.OpenTelemetry.Metrics.ObservableCounter+    Effectful.OpenTelemetry.Metrics.ObservableGauge+    Effectful.OpenTelemetry.Metrics.ObservableUpDownCounter+    Effectful.OpenTelemetry.Metrics.RTS+    Effectful.OpenTelemetry.Metrics.Sum+    Effectful.OpenTelemetry.Metrics.UpDownCounter+    Effectful.OpenTelemetry.Protocol+    Effectful.OpenTelemetry.Protocol.AnyValue+    Effectful.OpenTelemetry.Protocol.Attributes+    Effectful.OpenTelemetry.Protocol.Effect+    Effectful.OpenTelemetry.Protocol.Environment+    Effectful.OpenTelemetry.Protocol.Exception+    Effectful.OpenTelemetry.Protocol.Export+    Effectful.OpenTelemetry.Protocol.GRPC+    Effectful.OpenTelemetry.Protocol.HTTP+    Effectful.OpenTelemetry.Protocol.Resource+    Effectful.OpenTelemetry.Protocol.Scope+    Effectful.OpenTelemetry.Protocol.Transport+    Effectful.OpenTelemetry.Timestamp+    Effectful.OpenTelemetry.Tracing+    Effectful.OpenTelemetry.Tracing.Effect+    Effectful.OpenTelemetry.Tracing.Propagator+    Effectful.OpenTelemetry.Tracing.Span+    Effectful.OpenTelemetry.Tracing.Span.Context+    Effectful.OpenTelemetry.Tracing.Span.Event+    Effectful.OpenTelemetry.Tracing.Span.ID+    Effectful.OpenTelemetry.Tracing.Span.Kind+    Effectful.OpenTelemetry.Tracing.Span.Status+    Effectful.OpenTelemetry.Tracing.Trace.Flags+    Effectful.OpenTelemetry.Tracing.Trace.ID+    Effectful.OpenTelemetry.Tracing.Trace.State+    Proto3.Wire.Encode.Class++  other-modules:+    Prettyprinter.Extra+    System.Random.Extra++test-suite spec+  import: common+  type: exitcode-stdio-1.0+  hs-source-dirs: test+  main-is: Main.hs+  ghc-options:+    -threaded+    -rtsopts+    "-with-rtsopts=-N -T"++  build-tool-depends: hspec-effectful-discover:hspec-effectful-discover+  build-depends:+    base16-bytestring >=1.0 && <1.1,+    base64-bytestring >=1.2 && <1.3,+    bytestring >=0.12 && <0.13,+    composition-extra >=1.0 && <2.2,+    containers >=0.7 && <0.9,+    hspec-effectful >=1.1 && <1.2,+    http-client >=0.7 && <0.8,+    hunit-effectful >=1.0 && <1.1,+    network >=3.2 && <3.3,+    network-simple >=0.4 && <0.5,+    otel-effectful,+    quickcheck-effectful >=1.0 && <1.1,++  other-modules:+    Arbitrary+    Effectful.OpenTelemetry.Exporter.DeadSpec+    Effectful.OpenTelemetry.Exporter.Grafana.Loki+    Effectful.OpenTelemetry.Exporter.Grafana.Mimir+    Effectful.OpenTelemetry.Exporter.Grafana.Polling+    Effectful.OpenTelemetry.Exporter.Grafana.Tempo+    Effectful.OpenTelemetry.Exporter.GrafanaSpec+    Effectful.OpenTelemetry.ExporterSpec+    Effectful.OpenTelemetry.Logging.SeveritySpec+    Effectful.OpenTelemetry.LoggingSpec+    Effectful.OpenTelemetry.Metrics.MeasurementSpec+    Effectful.OpenTelemetry.MetricsSpec+    Effectful.OpenTelemetry.Protocol.AnyValueSpec+    Effectful.OpenTelemetry.Protocol.AttributesSpec+    Effectful.OpenTelemetry.Protocol.EffectSpec+    Effectful.OpenTelemetry.Protocol.ExceptionSpec+    Effectful.OpenTelemetry.Protocol.TransportSpec+    Effectful.OpenTelemetry.TimestampSpec+    Effectful.OpenTelemetry.Tracing.PropagatorSpec+    Effectful.OpenTelemetry.Tracing.Span.ContextSpec+    Effectful.OpenTelemetry.Tracing.Span.IDSpec+    Effectful.OpenTelemetry.Tracing.Span.KindSpec+    Effectful.OpenTelemetry.Tracing.Trace.FlagsSpec+    Effectful.OpenTelemetry.Tracing.Trace.IDSpec+    Effectful.OpenTelemetry.Tracing.Trace.StateSpec+    Effectful.OpenTelemetry.TracingSpec+    Util
+ src/Effectful/OpenTelemetry.hs view
@@ -0,0 +1,173 @@+{-# LANGUAGE Trustworthy #-}++-- |+-- Module      : Effectful.OpenTelemetry+-- Copyright   : (c) 2026 Institute for Digital Autonomy+-- License     : EUPL-1.2+-- Maintainer  : IDA+--+-- An implementation of <https://opentelemetry.io/ OpenTelemetry> for the <https://hackage.haskell.org/package/effectful Effectful> ecosystem.+--+-- This library provides 'Tracing', 'Logging', and 'Metrics' effects, which record the+-- OpenTelemetry signals: spans, log records, and instrument measurements.+--+-- Run them in any combination with 'runTracing', 'runLogging' and 'runMetrics', or all together+-- with 'runOpenTelemetry'.+-- These runners read environment variables to determine where telemetry should be sent.+--+-- Lower-level runners allow you to wire up the configuration in some other way.+--+-- The modules in this library are designed for qualified import. See the examples below.+--+-- = Instrumenting an application+--+-- > import Effectful+-- > import Effectful.OpenTelemetry+-- > import Effectful.OpenTelemetry.Logging qualified as Logging+-- > import Effectful.OpenTelemetry.Logging.Severity qualified as Severity+-- > import Effectful.OpenTelemetry.Metrics.Counter qualified as Counter+-- > import Effectful.OpenTelemetry.Protocol.Attributes qualified as Attributes+-- > import Effectful.OpenTelemetry.Protocol.Resource qualified as Resource+-- > import Effectful.OpenTelemetry.Protocol.Scope qualified as Scope+-- > import Effectful.OpenTelemetry.Tracing qualified as Tracing+-- > import Effectful.OpenTelemetry.Tracing.Span.Kind qualified as Span.Kind+-- >+-- > resource :: Resource+-- > resource = Resource{attributes = Attributes.fromList [("service.name", "my-service")]}+-- >+-- > scope :: Scope+-- > scope = Scope{name = "my-service", version = "1.0", attributes = mempty}+-- >+-- > app :: (Tracing :> es, Logging :> es, Metrics :> es) => Eff es ()+-- > app = do+-- >     counter <- Counter.new "myCounter" mempty+-- >     Tracing.inSpan "increment-counter" Span.Kind.Internal mempty do+-- >         Counter.add counter 1+-- >         Logging.log Severity.Info "incremented counter" mempty Nothing+-- >+-- > main :: IO ()+-- > main = runEff $ runOpenTelemetry resource scope app+module Effectful.OpenTelemetry+    ( -- * Effects+      Logging+    , Metrics+    , Tracing+    , runOpenTelemetry+    , runOpenTelemetry'+    , runInMemoryOpenTelemetry+    , runOpenTelemetryWith+    , runNoOpenTelemetry+    , runTracing+    , runLogging+    , runMetrics+    , runTracingWith+    , runLoggingWith+    , runMetricsWith++      -- * Exporting signals+    , Exporter++      -- * Signals+    , LogRecord+    , Measurement+    , Span++      -- * Protocol+    , Resource (Resource)+    , Scope (Scope)++      -- * Exceptions+    , SomeOTLPException (..)+    )+where++import Control.Monad ((>=>))+import Control.Monad.Extra (ifM)+import Effectful+import Effectful.Concurrent (Concurrent, runConcurrent)+import Effectful.Environment (Environment, runEnvironment)+import Effectful.OpenTelemetry.Exporter (Exporter)+import Effectful.OpenTelemetry.Exporter qualified as Exporter+import Effectful.OpenTelemetry.Logging+import Effectful.OpenTelemetry.Metrics+import Effectful.OpenTelemetry.Protocol+import Effectful.OpenTelemetry.Protocol.Environment qualified as Environment+import Effectful.OpenTelemetry.Protocol.Exception+import Effectful.OpenTelemetry.Tracing hiding (inject)+import Effectful.Retry (Retry, runRetry)+import Effectful.Timeout (Timeout, runTimeout)+import Numeric.Natural (Natural)+import Prelude++-- | Run all OpenTelmetry effects, sending telemetry to an exporter.+-- Reads the configuration from <https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#general-sdk-configuration the standard environment variables>.+-- Delegates to 'runNoOpenTelemetry' if @OTEL_SDK_DISABLED = true@.+runOpenTelemetry+    :: (IOE :> es)+    => Resource+    -> Scope+    -> Eff (Logging ': Metrics ': Tracing ': es) a+    -> Eff es a+runOpenTelemetry resource scope =+    runConcurrent+        . runEnvironment+        . runRetry+        . runTimeout+        . runOpenTelemetry' resource scope+        . inject++runOpenTelemetry'+    :: ( IOE :> es+       , Concurrent :> es+       , Environment :> es+       , Retry :> es+       , Timeout :> es+       )+    => Resource+    -> Scope+    -> Eff (Logging ': Metrics ': Tracing ': es) a+    -> Eff es a+runOpenTelemetry' resource scope eff =+    ifM+        Environment.isSdkDisabled+        (runNoOpenTelemetry eff)+        ( runTracing resource scope+            . runMetrics resource scope+            . runLogging resource scope+            $ eff+        )++-- | Run all OpenTelmetry effects, collecting telemetry in-memory+-- rather than sending to a collector.+runInMemoryOpenTelemetry+    :: (IOE :> es, Concurrent :> es)+    => Eff (Logging ': Metrics ': Tracing ': es) a+    -> Eff es (a, [LogRecord], [Measurement], [Span])+runInMemoryOpenTelemetry =+    runInMemoryTracing+        . runInMemoryMetrics+        . runInMemoryLogging+        >=> \(((a, logs), metrics), spans) -> pure (a, logs, metrics, spans)++-- | Run all OpenTelemetry effects with the given 'Exproter's.+-- Passes records and spans to the exporters synchronously as they are produced,+-- without batching or retrying.+runOpenTelemetryWith+    :: (IOE :> es)+    => Maybe Natural+    -- ^ Background sampling interval in ms.+    -- If set, a background thread periodically samples instruments.+    -- Otherwise, they are only sampled once at the end.+    -> Exporter es LogRecord+    -> Exporter es Measurement+    -> Exporter es Span+    -> Eff (Logging ': Metrics ': Tracing ': es) a+    -> Eff es a+runOpenTelemetryWith intervalMs logs metrics spans =+    runTracingWith spans+        . runMetricsWith intervalMs (Exporter.inject metrics)+        . runLoggingWith (Exporter.inject logs)++-- | Run all OpenTelmetry effects as no-op actions.+runNoOpenTelemetry :: (IOE :> es) => Eff (Logging ': Metrics ': Tracing ': es) a -> Eff es a+runNoOpenTelemetry = runNoTracing . runNoMetrics . runNoLogging
+ src/Effectful/OpenTelemetry/Exporter.hs view
@@ -0,0 +1,12 @@+module Effectful.OpenTelemetry.Exporter+    ( module Effectful.OpenTelemetry.Exporter.Type+    , module Effectful.OpenTelemetry.Exporter.Console+    , module Effectful.OpenTelemetry.Exporter.OTLP+    , module Effectful.OpenTelemetry.Exporter.STM+    )+where++import Effectful.OpenTelemetry.Exporter.Console+import Effectful.OpenTelemetry.Exporter.OTLP+import Effectful.OpenTelemetry.Exporter.STM+import Effectful.OpenTelemetry.Exporter.Type
+ src/Effectful/OpenTelemetry/Exporter/Console.hs view
@@ -0,0 +1,49 @@+module Effectful.OpenTelemetry.Exporter.Console where++import Effectful.OpenTelemetry.Exporter.Type (Exporter)+import Effectful.OpenTelemetry.Exporter.Type qualified as Exporter+import Prettyprinter (unAnnotate)+import Prettyprinter.Extra (PrettyAnn (..))+import Prettyprinter.Render.Terminal (AnsiStyle, hPutDoc)+import System.Console.ANSI (hSupportsANSI)+import System.Environment (lookupEnv)+import System.IO (Handle)+import System.IO qualified as Handle+import Prelude++-- | Print telemetry to the given 'Handle', one payload per line.+--+-- The output is styled using ANSI escape sequences if supported (see 'hSupportsANSI').+-- Styling can be disabled by setting the @NO_COLOR@ environment variable.+console :: (PrettyAnn AnsiStyle a) => Handle -> Exporter es a+console handle = Exporter.fromIO \a -> do+    colour <-+        lookupEnv "NO_COLOR" >>= \case+            Just _ -> pure False+            Nothing -> hSupportsANSI handle+    if colour+        then ansiIO handle a+        else plainIO handle a++-- | Print telemetry to the given 'Handle', one payload per line.+-- The output is styled using ANSI escape sequences.+ansi :: (PrettyAnn AnsiStyle a) => Handle -> Exporter es a+ansi = Exporter.fromIO . ansiIO++ansiIO :: (PrettyAnn AnsiStyle a) => Handle -> a -> IO ()+ansiIO handle a = hPutDoc handle $ prettyAnn a <> "\n"++-- | Print telemetry to the given 'Handle', one payload per line.+plain :: (PrettyAnn AnsiStyle a) => Handle -> Exporter es a+plain = Exporter.fromIO . plainIO++plainIO :: (PrettyAnn AnsiStyle a) => Handle -> a -> IO ()+plainIO handle a = hPutDoc handle $ unAnnotate @AnsiStyle (prettyAnn a) <> "\n"++-- | Print telemetry to standard output via 'console'.+stdout :: (PrettyAnn AnsiStyle a) => Exporter es a+stdout = console Handle.stdout++-- | Print telemetry to standard error via 'console'.+stderr :: (PrettyAnn AnsiStyle a) => Exporter es a+stderr = console Handle.stderr
+ src/Effectful/OpenTelemetry/Exporter/Environment.hs view
@@ -0,0 +1,54 @@+module Effectful.OpenTelemetry.Exporter.Environment where++import Control.Monad.Extra (mconcatMapM)+import Data.Char (toLower)+import Data.List (intercalate)+import Data.List.Extra (nubOrd, splitOn, trim)+import Effectful+import Effectful.Concurrent (Concurrent)+import Effectful.Environment (Environment, lookupEnv)+import Effectful.Error.Static (Error, throwError)+import Effectful.OpenTelemetry.Exporter.Console+import Effectful.OpenTelemetry.Exporter.OTLP+import Effectful.OpenTelemetry.Exporter.Type+import Effectful.OpenTelemetry.Protocol.Environment qualified as Environment+import Effectful.OpenTelemetry.Protocol.Export (Request (..))+import Effectful.OpenTelemetry.Protocol.Export qualified as Export+import Effectful.OpenTelemetry.Protocol.Resource (Resource)+import Effectful.OpenTelemetry.Protocol.Scope (Scope)+import Effectful.Retry (Retry)+import Effectful.Timeout (Timeout)+import Prettyprinter.Extra (PrettyAnn)+import Prettyprinter.Render.Terminal (AnsiStyle)+import Prelude++lookup+    :: forall a es es'+     . ( Export.Request a+       , PrettyAnn AnsiStyle a+       , Environment :> es+       , Error Environment.ConfigError :> es+       , IOE :> es'+       , Environment :> es'+       , Retry :> es'+       , Timeout :> es'+       , Concurrent :> es'+       )+    => Resource -> Scope -> Eff es (Exporter es' a)+lookup resource scope =+    lookupEnv (intercalate "_" ["OTEL", transportSignalEnvName @a, "EXPORTER"]) >>= \case+        Nothing -> pure otlp'+        Just (trim -> "") -> pure otlp'+        Just raw ->+            mconcatMapM \case+                "none" -> mempty+                "otlp" -> pure otlp'+                "console" -> pure stdout+                other -> throwError $ Environment.UnknownExporter other+                . nubOrd+                . filter (not . null)+                . map (map toLower . trim)+                . splitOn ","+                $ raw+  where+    otlp' = otlp resource scope
+ src/Effectful/OpenTelemetry/Exporter/OTLP.hs view
@@ -0,0 +1,295 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}++module Effectful.OpenTelemetry.Exporter.OTLP+    ( otlp+    , http+    , grpc+    , defaultSendTimeoutSeconds+    , defaultGrpcSendTimeout+    )+where++import Control.Concurrent.STM qualified as STM+import Control.Monad (forever, unless)+import Control.Monad.Extra (fromMaybeM, ifM)+import Data.Aeson qualified as Aeson+import Data.ByteString.Lazy (LazyByteString)+import Data.Functor (void)+import Effectful+import Effectful.Concurrent (Concurrent, threadDelay)+import Effectful.Concurrent.Async (Async, async)+import Effectful.Concurrent.Async qualified as Async+import Effectful.Concurrent.STM+    ( STM+    , TBQueue+    , atomically+    , flushTBQueue+    , isFullTBQueue+    , newTBQueueIO+    , tryReadTBQueue+    , writeTBQueue+    )+import Effectful.Environment (Environment)+import Effectful.Error.Static (runErrorNoCallStackWith)+import Effectful.Exception (bracket, handleSync, throwIO, trySync)+import Effectful.GrpcClient qualified as Grpc+import Effectful.Http2Client (HostName, PortNumber)+import Effectful.HttpClient+    ( ResponseTimeout+    , responseTimeoutMicro+    , responseTimeoutNone+    , runHttpClientTls+    )+import Effectful.OpenTelemetry.Exporter.Type (Exporter (..))+import Effectful.OpenTelemetry.Protocol.Environment+    ( ConfigError+    , detectResource+    , isSdkDisabled+    , lookupExportConfig+    , lookupTransport+    )+import Effectful.OpenTelemetry.Protocol.Exception+    ( ConnectionTimeout+    , OTLPClientError (..)+    , OTLPGrpcError (..)+    , OTLPHttpException (..)+    , SomeOTLPException (..)+    )+import Effectful.OpenTelemetry.Protocol.Export qualified as Export+import Effectful.OpenTelemetry.Protocol.GRPC qualified as GRPC+import Effectful.OpenTelemetry.Protocol.HTTP qualified as HTTP+import Effectful.OpenTelemetry.Protocol.Resource (Resource)+import Effectful.OpenTelemetry.Protocol.Resource qualified as Resource+import Effectful.OpenTelemetry.Protocol.Scope (Scope)+import Effectful.OpenTelemetry.Protocol.Transport+    ( Compression+    , Encoding (..)+    , Protocol (..)+    , Transport (..)+    )+import Effectful.Retry+    ( Retry+    , RetryPolicyM+    , capDelay+    , fullJitterBackoff+    , limitRetriesByCumulativeDelay+    , recoverAll+    )+import Effectful.Timeout (Timeout, timeout)+import Network.GRPC.HTTP2.Proto3Wire (RPC)+import Network.URI (URI)+import Numeric.Natural (Natural)+import Proto3.Wire.Encode qualified as Proto3.Encode+import Prelude++-- | A collector configured from+-- <https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#general-sdk-configuration the standard environment variables>.+--+-- 'mempty' if @OTEL_SDK_DISABLED = true@.+otlp+    :: forall a es+     . ( Export.Request a+       , IOE :> es+       , Concurrent :> es+       , Environment :> es+       , Retry :> es+       , Timeout :> es+       )+    => Resource+    -> Scope+    -> Exporter es a+otlp resource scope = Exporter \liftExporter ->+    ifM isSdkDisabled (liftExporter mempty) $+        runErrorNoCallStackWith (throwIO . SomeOTLPException @ConfigError) do+            resource <- (resource <>) <$> detectResource+            Transport{..} <- lookupTransport @a+            exportConfig <- lookupExportConfig @a $ Export.defaultConfig @a+            withExporter+                ( case protocol of+                    HTTP encoding endpoint ->+                        http @a+                            resource+                            scope+                            encoding+                            endpoint+                            exportConfig+                            compression+                            ( if timeoutMs == 0+                                then responseTimeoutNone+                                else responseTimeoutMicro $ fromIntegral timeoutMs * 1_000+                            )+                    GRPC host port rpc ->+                        grpc @a+                            resource+                            scope+                            host+                            port+                            rpc+                            exportConfig+                            compression+                            (Grpc.Timeout $ fromIntegral timeoutMs `div` 1_000)+                )+                (inject . liftExporter)++http+    :: forall a es+     . ( Export.Request a+       , IOE :> es+       , Concurrent :> es+       , Retry :> es+       , Timeout :> es+       )+    => Resource+    -> Scope+    -> Encoding+    -> URI+    -> Export.Config+    -> Compression+    -> ResponseTimeout+    -> Exporter es a+http resource scope encoding endpoint exportConfig compression responseTimeout =+    Exporter $ export exportConfig sender+  where+    encodeBatch :: Encoding -> [Resource.Items a] -> LazyByteString+    encodeBatch Json items = Aeson.encode (Export.exportJson @a items)+    encodeBatch Proto items = Proto3.Encode.toLazyByteString (Export.exportProto @a items)+    signalEndpoint = Export.appendExportPath @a endpoint+    sender :: Exporter es [a]+    sender = Exporter \send ->+        runHttpClientTls $ withEffToIO (ConcUnlift Persistent Unlimited) \unlift ->+            unlift . inject . send $+                unlift+                    . runErrorNoCallStackWith (throwIO . SomeOTLPException . OTLPHttpException)+                    . HTTP.sendPayload+                        signalEndpoint+                        encoding+                        compression+                        responseTimeout+                    . encodeBatch encoding+                    . Resource.wrapItems resource scope++grpc+    :: forall a es+     . ( Export.Request a+       , IOE :> es+       , Concurrent :> es+       , Retry :> es+       , Timeout :> es+       )+    => Resource+    -> Scope+    -> HostName+    -> PortNumber+    -> RPC+    -> Export.Config+    -> Compression+    -> Grpc.Timeout+    -> Exporter es a+grpc resource scope host port rpc exportConfig compression timeout =+    Exporter $ export exportConfig sender+  where+    sender :: Exporter es [a]+    sender = Exporter \send ->+        runErrorNoCallStackWith (throwIO . SomeOTLPException @ConnectionTimeout)+            . runErrorNoCallStackWith (throwIO . SomeOTLPException . OTLPClientError)+            . GRPC.runGrpc host port+            $ withEffToIO (ConcUnlift Persistent Unlimited) \unlift ->+                unlift . inject . send $+                    unlift+                        . runErrorNoCallStackWith (throwIO . SomeOTLPException . OTLPGrpcError)+                        . GRPC.sendRpc host port rpc compression timeout+                        . Export.exportProto @a+                        . Resource.wrapItems resource scope++defaultSendTimeoutSeconds :: Natural+defaultSendTimeoutSeconds = 30++defaultGrpcSendTimeout :: Grpc.Timeout+defaultGrpcSendTimeout = Grpc.Timeout $ fromIntegral defaultSendTimeoutSeconds++withTimeout :: (Timeout :> es) => Natural -> ([a] -> Eff es ()) -> [a] -> Eff es ()+withTimeout timeoutMs send = fromMaybeM (throwIO $ userError "timeout") . timeout (fromIntegral timeoutMs * 1_000) . send++-- | Retry an action with capped exponential backoff, giving up without rethrowing+-- once retries are exhausted.+retrying :: forall es. (IOE :> es, Retry :> es) => Eff es () -> Eff es ()+retrying =+    handleSync (const $ pure ())+        . recoverAll retryPolicy+        . const+  where+    retryPolicy :: RetryPolicyM (Eff es)+    retryPolicy =+        limitRetriesByCumulativeDelay maxElapsedRetryTimeµs+            . capDelay maxRetryIntervalµs+            $ fullJitterBackoff initialRetryIntervalµs++    -- NOTE: There is no specification for configuring retry parameters,+    -- so we hard code the defaults.+    initialRetryIntervalµs :: Int+    initialRetryIntervalµs = 5_000_000 -- 5 s+    maxRetryIntervalµs :: Int+    maxRetryIntervalµs = 30_000_000 -- 30 s+    maxElapsedRetryTimeµs :: Int+    maxElapsedRetryTimeµs = 300_000_000 -- 300 s++-- | Send items to the batch exporter, to be exported either individually or in batches,+-- depending on the given 'Export.Config'.+-- Items may be dropped on a full queue, exhausted retries, or timeout.+export+    :: forall es a b+     . ( Retry :> es+       , Concurrent :> es+       , IOE :> es+       , Timeout :> es+       )+    => Export.Config+    -> Exporter es [a]+    -> ((a -> IO ()) -> Eff es b)+    -> Eff es b+export Export.Config{batch = Nothing, ..} exporter action =+    withEffToIO (ConcUnlift Persistent Unlimited) \unlift ->+        unlift . action $ \item ->+            unlift . retrying . withExporter exporter $ \export ->+                withTimeout exportTimeoutMs (liftIO . export) [item]+export Export.Config{batch = Just Export.BatchConfig{..}, ..} exporter action = do+    queue <- newTBQueueIO maxQueueSize+    let worker :: Eff es (Async a)+        worker =+            async . forever . retrying . withExporter exporter $+                forever . \send -> do+                    threadDelay . fromIntegral $ 1_000 * scheduledDelayMs+                    items <- atomically $ drainN maxBatchSize queue+                    unless (null items) . withTimeout exportTimeoutMs (liftIO . send) $ items++        enqueue :: a -> IO ()+        enqueue item = STM.atomically do+            full <- isFullTBQueue queue+            unless full $ writeTBQueue queue item++        flush :: Async a -> Eff es ()+        flush task = do+            Async.cancel task+            remaining <- atomically $ flushTBQueue queue+            -- NOTE: We don't retry when flushing because we don't want to block+            -- program termination. 'withTimeout' still caps how long this+            -- can take, so a stuck connection can't hang shutdown either.+            -- To properly protect against data loss if the collector itself crashes,+            -- we may want to consider implementing persistent storage:+            -- https://opentelemetry.io/docs/collector/resiliency/#persistent-storage-write-ahead-log---wal+            unless (null remaining)+                . void+                . trySync+                $ withExporter exporter \export -> withTimeout exportTimeoutMs (liftIO . export) remaining++    bracket worker flush . const $ action enqueue+  where+    drainN :: Natural -> TBQueue a -> STM [a]+    drainN = (fmap reverse .) . drainNReverse++    drainNReverse :: Natural -> TBQueue a -> STM [a]+    drainNReverse 0 _ = pure []+    drainNReverse n queue =+        tryReadTBQueue queue >>= \case+            Nothing -> pure []+            Just x -> (x :) <$> drainNReverse (n - 1) queue
+ src/Effectful/OpenTelemetry/Exporter/STM.hs view
@@ -0,0 +1,21 @@+module Effectful.OpenTelemetry.Exporter.STM where++import Control.Concurrent.STM hiding (atomically)+import Control.Concurrent.STM qualified as STM+import Effectful.OpenTelemetry.Exporter.Type+import Prelude++atomically :: (a -> STM ()) -> Exporter es a+atomically = fromIO . (STM.atomically .)++tvar :: TVar a -> Exporter es a+tvar = atomically . writeTVar++tmvar :: TMVar a -> Exporter es a+tmvar = atomically . putTMVar++tqueue :: TQueue a -> Exporter es a+tqueue = atomically . writeTQueue++tbqueue :: TBQueue a -> Exporter es a+tbqueue = atomically . writeTBQueue
+ src/Effectful/OpenTelemetry/Exporter/Type.hs view
@@ -0,0 +1,25 @@+module Effectful.OpenTelemetry.Exporter.Type where++import Effectful hiding (inject)+import Effectful qualified+import Effectful.Dispatch.Static (unsafeEff_)+import Prelude++newtype Exporter es a = Exporter {withExporter :: forall r. ((a -> IO ()) -> Eff es r) -> Eff es r}++-- | Send each event to both exporters.+instance Semigroup (Exporter es a) where+    Exporter f <> Exporter g = Exporter \k -> f \x -> g \y -> k (x <> y)++-- | The no-op exporter, which discards every event.+instance Monoid (Exporter es a) where+    mempty = Exporter ($ mempty)++fromIO :: (a -> IO ()) -> Exporter es a+fromIO f = Exporter ($ f)++-- | Embed an 'Exporter' into a broader effect stack @es'@. Generalised 'Effectful.inject'.+inject :: forall es es' a. (IOE :> es', Subset es es') => Exporter es a -> Exporter es' a+inject (Exporter withE) = Exporter \k ->+    withEffToIO (ConcUnlift Persistent Unlimited) \unlift ->+        unlift . Effectful.inject . withE $ unsafeEff_ . unlift . k
+ src/Effectful/OpenTelemetry/Logging.hs view
@@ -0,0 +1,10 @@+module Effectful.OpenTelemetry.Logging+    ( module Effectful.OpenTelemetry.Logging.Effect+    , module Effectful.OpenTelemetry.Logging.LogRecord+    , module Effectful.OpenTelemetry.Logging.Severity+    )+where++import Effectful.OpenTelemetry.Logging.Effect+import Effectful.OpenTelemetry.Logging.LogRecord+import Effectful.OpenTelemetry.Logging.Severity
+ src/Effectful/OpenTelemetry/Logging/Effect.hs view
@@ -0,0 +1,196 @@+module Effectful.OpenTelemetry.Logging.Effect+    ( -- * Effect+      Logging+    , runLogging+    , runHttpLogging+    , runGrpcLogging+    , runConsoleLogging+    , runInMemoryLogging+    , runLoggingWith+    , runNoLogging+    , withTracing++      -- * Emit logs+    , log+    )+where++import Control.Monad.Extra (ifM)+import Data.Aeson.Types (Value)+import Data.Text (Text)+import Effectful+import Effectful.Concurrent (Concurrent)+import Effectful.Dispatch.Static+import Effectful.Environment (Environment)+import Effectful.Http2Client (HostName, PortNumber)+import Effectful.HttpClient (responseTimeoutDefault)+import Effectful.OpenTelemetry.Exporter (Exporter)+import Effectful.OpenTelemetry.Exporter qualified as Exporter+import Effectful.OpenTelemetry.Exporter.Console qualified as Console+import Effectful.OpenTelemetry.Exporter.Environment qualified as Exporter.Environment+import Effectful.OpenTelemetry.Exporter.OTLP (defaultGrpcSendTimeout)+import Effectful.OpenTelemetry.Logging.LogRecord (LogRecord (..))+import Effectful.OpenTelemetry.Logging.Severity (Severity)+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+import Effectful.OpenTelemetry.Protocol.Effect+    ( OTLP+    , exportIO+    , runInMemoryOTLP+    , runNoOTLP+    , runOTLPWith+    )+import Effectful.OpenTelemetry.Protocol.Environment qualified as Environment+import Effectful.OpenTelemetry.Protocol.Export qualified as Export+import Effectful.OpenTelemetry.Protocol.Resource (Resource)+import Effectful.OpenTelemetry.Protocol.Scope (Scope)+import Effectful.OpenTelemetry.Protocol.Transport+    ( Compression+    , Encoding+    )+import Effectful.OpenTelemetry.Timestamp qualified as Timestamp+import Effectful.OpenTelemetry.Tracing.Effect (Tracing, currentContextIO)+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span (Context)+import Effectful.Retry (Retry)+import Effectful.Timeout (Timeout)+import Network.URI (URI)+import Prelude hiding (log)++data Logging :: Effect++type instance DispatchOf Logging = 'Static 'WithSideEffects++data instance StaticRep Logging = Logging+    { getContext :: IO (Maybe Span.Context)+    , export :: LogRecord -> IO ()+    }++withTracing :: (Logging :> es, Tracing :> es) => Eff es a -> Eff es a+withTracing eff = do+    getContext <- currentContextIO+    localStaticRep (\logging -> logging{getContext}) eff++-- | Emit a single log record with the current 'Timestamp.Timestamp' and 'Tracing' context, if any.+-- This function can be used without 'Tracing' as well.+log+    :: (Logging :> es)+    => Severity+    -> Value+    -> Attributes+    -> Maybe Text+    -> Eff es ()+log severity body attributes eventName = do+    Logging{..} <- getStaticRep+    unsafeEff_ do+        timestamp <- Timestamp.now+        context <- getContext+        export LogRecord{observedTimestamp = timestamp, ..}++runLoggingState+    :: (IOE :> es, OTLP LogRecord :> es, Tracing :> es)+    => Eff (Logging ': es) a+    -> Eff es a+runLoggingState eff = do+    getContext <- currentContextIO+    runLoggingState' getContext eff++runLoggingState'+    :: (IOE :> es, OTLP LogRecord :> es)+    => IO (Maybe Span.Context)+    -> Eff (Logging ': es) a+    -> Eff es a+runLoggingState' getContext eff = do+    export <- exportIO+    evalStaticRep (Logging{..}) eff++-- | Run the 'Logging' effect, sending telemetry to an exporter.+-- Uses the 'Tracing' effect to enrich logs with span context.+-- Reads the configuration from <https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#general-sdk-configuration the standard environment variables>.+-- Delegates to 'runNoLogging' if @OTEL_SDK_DISABLED = true@.+runLogging+    :: ( Tracing :> es+       , IOE :> es+       , Concurrent :> es+       , Environment :> es+       , Retry :> es+       , Timeout :> es+       )+    => Resource+    -> Scope+    -> Eff (Logging ': es) a+    -> Eff es a+runLogging resource scope eff =+    ifM Environment.isSdkDisabled (runNoLogging eff) do+        exporter <- Environment.runConfigError $ Exporter.Environment.lookup @LogRecord resource scope+        runLoggingWith exporter eff++-- | Run the 'Logging' effect, sending telemetry to a collector at the given 'URI'+-- over HTTP with the given 'Encoding'.+-- Uses the 'Tracing' effect to enrich logs with span context.+runHttpLogging+    :: (IOE :> es, Concurrent :> es, Retry :> es, Timeout :> es, Tracing :> es)+    => Resource+    -> Scope+    -> Encoding+    -> Compression+    -> URI+    -> Eff (Logging ': es) a+    -> Eff es a+runHttpLogging resource scope encoding compression endpoint =+    runLoggingWith $+        Exporter.http @LogRecord+            resource+            scope+            encoding+            endpoint+            (Export.defaultConfig @LogRecord)+            compression+            responseTimeoutDefault++-- | Run the 'Logging' effect, sending telemetry to a gRPC collector.+-- Uses the 'Tracing' effect to enrich logs with span context.+runGrpcLogging+    :: (IOE :> es, Concurrent :> es, Retry :> es, Timeout :> es, Tracing :> es)+    => Resource+    -> Scope+    -> HostName+    -> PortNumber+    -> Compression+    -> Eff (Logging ': es) a+    -> Eff es a+runGrpcLogging resource scope host port compression =+    runLoggingWith $+        Exporter.grpc @LogRecord+            resource+            scope+            host+            port+            (Export.exportGrpcRPC @LogRecord)+            (Export.defaultConfig @LogRecord)+            compression+            defaultGrpcSendTimeout++-- | Run the 'Logging' effect, printing telemetry to the console rather than sending to a collector.+runConsoleLogging :: (IOE :> es, Tracing :> es) => Eff (Logging ': es) a -> Eff es a+runConsoleLogging = runLoggingWith Console.stdout++-- | Run the 'Logging' effect, collecting telemetry in-memory rather than sending to a collector.+-- Uses the 'Tracing' effect to enrich logs with span context.+runInMemoryLogging+    :: (IOE :> es, Concurrent :> es, Tracing :> es)+    => Eff (Logging ': es) a+    -> Eff es (a, [LogRecord])+runInMemoryLogging = runInMemoryOTLP @LogRecord . runLoggingState . inject++-- | Run the 'Logging' effect with a given 'Exporter'.+-- Passes records to the exporter synchronously as they are emitted, without batching or retrying.+-- Uses the 'Tracing' effect to enrich logs with span context.+runLoggingWith+    :: (IOE :> es, Tracing :> es)+    => Exporter es LogRecord+    -> Eff (Logging ': es) a+    -> Eff es a+runLoggingWith exporter = runOTLPWith exporter . runLoggingState . inject++-- | Run the 'Logging' effect as a no-op action.+runNoLogging :: (IOE :> es) => Eff (Logging ': es) a -> Eff es a+runNoLogging = runNoOTLP @LogRecord . runLoggingState' (pure Nothing) . inject
+ src/Effectful/OpenTelemetry/Logging/LogRecord.hs view
@@ -0,0 +1,128 @@+module Effectful.OpenTelemetry.Logging.LogRecord where++import Data.Aeson.Types (ToJSON (..), Value (..), object, (.=))+import Data.Function ((&))+import Data.Functor ((<&>))+import Data.Maybe (catMaybes)+import Data.Text (Text)+import Effectful.OpenTelemetry.Logging.Severity (Severity)+import Effectful.OpenTelemetry.Protocol (AnyValue (..), 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.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span+import GHC.Generics (Generic)+import Network.GRPC.HTTP2.Proto3Wire (RPC (..))+import Prettyprinter (Pretty (..))+import Prettyprinter.Extra (PrettyAnn (..))+import Prettyprinter.Extra qualified as Pretty+import Prettyprinter.Render.Terminal (AnsiStyle)+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude hiding (String)++-- | Represents the recording of an event.+--+-- See <https://opentelemetry.io/docs/concepts/signals/logs/#log-record the OpenTelemetry spec>.+data LogRecord = LogRecord+    { timestamp :: Timestamp+    -- ^ Time when the event occured.+    , observedTimestamp :: Timestamp+    -- ^ Time when the event was observed by the collection system.+    , context :: Maybe Span.Context+    -- ^ The tracing context, present if recorded within a span.+    , severity :: Severity+    -- ^ Also known as the log level.+    , body :: Value+    -- ^ Describes the event as structured data.+    -- For a human-readable free-form message, use the 'String' constructor.+    , attributes :: Attributes+    -- ^ Additional information about the event.+    , eventName :: Maybe Text+    -- ^ Identifies the class or type of the event.+    }+    deriving stock (Generic, Eq, Show)++instance ToJSON LogRecord where+    toJSON LogRecord{..} =+        object $+            [ "timeUnixNano" .= timestamp+            , "observedTimeUnixNano" .= observedTimestamp+            , "severityNumber" .= (Number . fromIntegral . fromEnum) severity+            , "body" .= AnyValue body+            , "attributes" .= attributes+            ]+                <> catMaybes+                    [ ("flags" .=) <$> (context <&> (.traceFlags))+                    , ("traceId" .=) <$> (context <&> (.traceId))+                    , ("spanId" .=) <$> (context <&> (.spanId))+                    , ("eventName" .=) <$> eventName+                    ]++instance Proto.Encode LogRecord where+    encode LogRecord{..} =+        mconcat+            [ Proto.encodeField 1 timestamp+            , Proto.encodeField 11 observedTimestamp+            , Proto.encodeField 2 severity+            , -- 3: severity_text: not needed+              Proto.encodeField 5 body+            , Proto.encodeField 6 attributes+            , -- 7: dropped_attributes_count: Not supported+              flip foldMap context \Span.Context{..} ->+                mconcat+                    [ Proto.encodeField 8 traceFlags+                    , Proto.encodeField 9 traceId+                    , Proto.encodeField 10 spanId+                    ]+            , foldMap (Proto.encodeField 12) eventName+            ]++instance PrettyAnn AnsiStyle LogRecord where+    prettyAnn LogRecord{..} =+        Pretty.unwords . filter (not . Pretty.null) $+            [ prettyAnn timestamp+            , prettyAnn severity+            , eventName & maybe mempty \n -> pretty $ "[" <> n <> "]"+            , prettyAnn body+            , prettyAnn attributes+            ]++instance Export.Request LogRecord where+    exportHttpPathComponents = ["v1", "logs"]+    exportGrpcRPC =+        RPC+            { pkg = "opentelemetry.proto.collector.logs.v1"+            , srv = "LogsService"+            , meth = "Export"+            }+    exportJson = object . pure . ("resourceLogs" .=) . fmap resourceLogs+      where+        resourceLogs :: Resource.Items LogRecord -> Value+        resourceLogs Resource.Items{resource, scopeItems} =+            object+                [ "resource" .= resource+                , "scopeLogs" .= fmap scopeLogs scopeItems+                ]+        scopeLogs :: Scope.Items LogRecord -> Value+        scopeLogs Scope.Items{..} =+            object+                [ "scope" .= scope+                , "logRecords" .= items+                ]++    transportSignalEnvName = "LOGS"+    batchEnvPrefix = Just "BLRP"++    -- https://opentelemetry.io/docs/specs/otel/logs/sdk/#batching-processor+    defaultConfig =+        Export.Config+            { batch =+                Just+                    Export.BatchConfig+                        { maxQueueSize = 2_048+                        , scheduledDelayMs = 1_000+                        , maxBatchSize = 512+                        }+            , exportTimeoutMs = 30_000+            }
+ src/Effectful/OpenTelemetry/Logging/Severity.hs view
@@ -0,0 +1,123 @@+module Effectful.OpenTelemetry.Logging.Severity (Severity (..), toJSON) where++import Data.Aeson.Types (ToJSON (..), Value (String))+import Data.Char (toLower)+import Data.Text qualified as Text+import GHC.Generics (Generic)+import Prettyprinter (Pretty (..), annotate)+import Prettyprinter.Extra (PrettyAnn (..))+import Prettyprinter.Render.Terminal (AnsiStyle, Color (..), color, colorDull)+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | The severity of a log record.+-- Smaller values correspond to less severe events.+data Severity+    = -- | Fine-grained, highly verbose execution details.+      -- Typically used during active development or troubleshooting to track execution flow.+      -- Almost always disabled in production due to noise.+      Trace+    | Trace2+    | Trace3+    | Trace4+    | -- | Diagnostic event. Typically used for highlighting internal system events.+      Debug+    | Debug2+    | Debug3+    | Debug4+    | -- | Informational event. Typically used when something expected occurred.+      Info+    | Info2+    | Info3+    | Info4+    | -- | An unexpected event that may indicate potential future issues.+      Warn+    | Warn2+    | Warn3+    | Warn4+    | -- | An operation has failed, but the system can continue to run.+      Error+    | Error2+    | Error3+    | Error4+    | -- | An unrecoverable, terminal event that compromises the entire process.+      -- Indicates imminent crash or shutdown.+      Fatal+    | Fatal2+    | Fatal3+    | Fatal4+    deriving stock (Generic, Show, Eq, Ord, Bounded)++instance Read Severity where+    readsPrec _ s =+        [(it, "") | it <- [minBound .. maxBound], fmap toLower (show it) == fmap toLower s]++instance Enum Severity where+    toEnum 1 = Trace+    toEnum 2 = Trace2+    toEnum 3 = Trace3+    toEnum 4 = Trace4+    toEnum 5 = Debug+    toEnum 6 = Debug2+    toEnum 7 = Debug3+    toEnum 8 = Debug4+    toEnum 9 = Info+    toEnum 10 = Info2+    toEnum 11 = Info3+    toEnum 12 = Info4+    toEnum 13 = Warn+    toEnum 14 = Warn2+    toEnum 15 = Warn3+    toEnum 16 = Warn4+    toEnum 17 = Error+    toEnum 18 = Error2+    toEnum 19 = Error3+    toEnum 20 = Error4+    toEnum 21 = Fatal+    toEnum 22 = Fatal2+    toEnum 23 = Fatal3+    toEnum 24 = Fatal4+    toEnum _ = error "Enum.Severity.toEnum: bad argument"+    fromEnum Trace = 1+    fromEnum Trace2 = 2+    fromEnum Trace3 = 3+    fromEnum Trace4 = 4+    fromEnum Debug = 5+    fromEnum Debug2 = 6+    fromEnum Debug3 = 7+    fromEnum Debug4 = 8+    fromEnum Info = 9+    fromEnum Info2 = 10+    fromEnum Info3 = 11+    fromEnum Info4 = 12+    fromEnum Warn = 13+    fromEnum Warn2 = 14+    fromEnum Warn3 = 15+    fromEnum Warn4 = 16+    fromEnum Error = 17+    fromEnum Error2 = 18+    fromEnum Error3 = 19+    fromEnum Error4 = 20+    fromEnum Fatal = 21+    fromEnum Fatal2 = 22+    fromEnum Fatal3 = 23+    fromEnum Fatal4 = 24++instance ToJSON Severity where+    toJSON = String . Text.toUpper . Text.show++instance {-# OVERLAPPING #-} Proto.EncodeField Severity where+    encodeField n = Proto.int32 n . fromIntegral . fromEnum++instance PrettyAnn AnsiStyle Severity where+    prettyAnn severity =+        annotate colour . pretty . Text.justifyLeft 5 ' ' . Text.toUpper $+            Text.show severity+      where+        colour :: AnsiStyle+        colour+            | severity >= Error = color Red+            | severity >= Warn = color Yellow+            | severity >= Info = color Green+            | severity >= Debug = colorDull Blue+            | otherwise = color Magenta
+ src/Effectful/OpenTelemetry/Metrics.hs view
@@ -0,0 +1,24 @@+module Effectful.OpenTelemetry.Metrics+    ( module Effectful.OpenTelemetry.Metrics.Counter+    , module Effectful.OpenTelemetry.Metrics.Effect+    , module Effectful.OpenTelemetry.Metrics.Gauge+    , module Effectful.OpenTelemetry.Metrics.Histogram+    , module Effectful.OpenTelemetry.Metrics.Measurement+    , module Effectful.OpenTelemetry.Metrics.ObservableCounter+    , module Effectful.OpenTelemetry.Metrics.ObservableGauge+    , module Effectful.OpenTelemetry.Metrics.ObservableUpDownCounter+    , module Effectful.OpenTelemetry.Metrics.Sum+    , module Effectful.OpenTelemetry.Metrics.UpDownCounter+    )+where++import Effectful.OpenTelemetry.Metrics.Counter (Counter)+import Effectful.OpenTelemetry.Metrics.Effect+import Effectful.OpenTelemetry.Metrics.Gauge (Gauge)+import Effectful.OpenTelemetry.Metrics.Histogram (Histogram)+import Effectful.OpenTelemetry.Metrics.Measurement+import Effectful.OpenTelemetry.Metrics.ObservableCounter (ObservableCounter)+import Effectful.OpenTelemetry.Metrics.ObservableGauge (ObservableGauge)+import Effectful.OpenTelemetry.Metrics.ObservableUpDownCounter (ObservableUpDownCounter)+import Effectful.OpenTelemetry.Metrics.Sum (Monotonicity (..))+import Effectful.OpenTelemetry.Metrics.UpDownCounter (UpDownCounter)
+ src/Effectful/OpenTelemetry/Metrics/Counter.hs view
@@ -0,0 +1,77 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-redundant-constraints #-}++module Effectful.OpenTelemetry.Metrics.Counter+    ( Counter+    , new+    , add+    , addIO+    )+where++import Control.Concurrent.STM (TVar)+import Control.Concurrent.STM qualified as STM+import Data.Bifunctor (bimap)+import Data.Scientific (Scientific)+import Data.Text (Text)+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 (NumberDataPoint (..))+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import Effectful.OpenTelemetry.Metrics.Sum (Monotonicity (..))+import Effectful.OpenTelemetry.Metrics.Sum qualified as Sum+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Timestamp qualified as Timestamp+import Prelude++-- | A synchronous 'Instrument' which supports non-negative increments.+--+-- Example uses for 'Counter':+--+-- - count the number of bytes received+-- - count the number of requests completed+-- - count the number of accounts created+-- - count the number of checkpoints run+-- - count the number of HTTP 5xx errors+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/api/#counter the OpenTelemetry spec>.+data Counter = Counter+    { name :: Text+    , startTime :: Timestamp+    , value :: TVar (Timestamp, Scientific)+    , metadata :: Metadata+    }++-- | Create a new 'Counter' and register it to be sampled and exported.+new :: (Metrics :> es) => Text -> Metadata -> Eff es Counter+new name metadata =+    Metrics.register =<< unsafeEff_ do+        startTime <- Timestamp.now+        value <- STM.newTVarIO (startTime, 0)+        pure Counter{..}++-- | Increment a 'Counter'.+-- WARNING: This function is partial because it throws on negative increment.+add :: (Metrics :> es) => Counter -> Scientific -> Eff es ()+add = (unsafeEff_ .) . addIO++addIO :: Counter -> Scientific -> IO ()+addIO Counter{..} inc+    | inc < 0 =+        error $ "Counter.add: negative increment (" <> show inc <> ") added to counter " <> show name+    | otherwise = do+        time <- Timestamp.now+        STM.atomically . STM.modifyTVar' value $ bimap (const time) (+ inc)++instance Instrument Counter where+    name = name+    sample Counter{..} = do+        (time, value) <- STM.readTVarIO value+        pure . pure $+            ( time+            , Sum.measurement Monotonic metadata name NumberDataPoint{attributes = metadata.attributes, ..}+            )
+ src/Effectful/OpenTelemetry/Metrics/Effect.hs view
@@ -0,0 +1,253 @@+{-# OPTIONS_GHC -Wno-redundant-constraints #-}++module Effectful.OpenTelemetry.Metrics.Effect+    ( -- * Effect+      Metrics+    , runMetrics+    , runHttpMetrics+    , runGrpcMetrics+    , runConsoleMetrics+    , runInMemoryMetrics+    , runMetricsWith+    , runNoMetrics++      -- * Instruments+    , register+    )+where++import Control.Concurrent qualified as IO+import Control.Concurrent.STM qualified as STM+import Control.Monad (forever, unless)+import Control.Monad.Extra (concatMapM, ifM)+import Data.HashMap.Strict (HashMap)+import Data.HashMap.Strict qualified as HashMap+import Data.List qualified as List+import Data.Maybe (fromMaybe)+import Data.Text (Text)+import Effectful+import Effectful.Concurrent (Concurrent)+import Effectful.Concurrent.STM (TVar)+import Effectful.Dispatch.Static+import Effectful.Environment (Environment, lookupEnv)+import Effectful.Exception (bracket, finally)+import Effectful.Http2Client (HostName, PortNumber)+import Effectful.HttpClient (responseTimeoutDefault)+import Effectful.OpenTelemetry.Exporter (Exporter)+import Effectful.OpenTelemetry.Exporter qualified as Exporter+import Effectful.OpenTelemetry.Exporter.Console qualified as Console+import Effectful.OpenTelemetry.Exporter.Environment qualified as Exporter.Environment+import Effectful.OpenTelemetry.Exporter.OTLP (defaultGrpcSendTimeout)+import Effectful.OpenTelemetry.Metrics.Instrument (Instrument, SomeInstrument (..))+import Effectful.OpenTelemetry.Metrics.Instrument qualified as Instrument+import Effectful.OpenTelemetry.Metrics.Measurement (Measurement (..))+import Effectful.OpenTelemetry.Protocol+    ( Compression+    , Encoding+    , OTLP+    , Resource+    , Scope+    , exportIO+    , runInMemoryOTLP+    , runNoOTLP+    , runOTLPWith+    )+import Effectful.OpenTelemetry.Protocol.Environment qualified as Environment+import Effectful.OpenTelemetry.Protocol.Export (exportSignalEnvName)+import Effectful.OpenTelemetry.Protocol.Export qualified as Export+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Timestamp qualified as Timestamp+import Effectful.Retry (Retry)+import Effectful.Timeout (Timeout)+import Network.URI (URI)+import Numeric.Natural (Natural)+import Text.Read (readMaybe)+import Prelude++data Metrics :: Effect++type instance DispatchOf Metrics = 'Static 'WithSideEffects++data instance StaticRep Metrics = Metrics+    { instruments :: TVar (HashMap Text SomeInstrument)+    , lastExport :: TVar Timestamp+    , export :: Measurement -> IO ()+    }++-- | Adds an 'Instrument' to the pool of instruments that are periodically exported.+register :: (Metrics :> es, Instrument i) => i -> Eff es i+register i = do+    Metrics{..} <- getStaticRep+    unsafeEff_ . STM.atomically . STM.modifyTVar instruments $+        HashMap.insert (Instrument.name i) (SomeInstrument i)+    pure i++-- | Samples all instruments, then exports the samples that are newer than the last export time.+-- Updates the last export time to the latest exported sample time.+exportInstruments :: StaticRep Metrics -> IO ()+exportInstruments Metrics{..} = do+    lastExport' <- STM.readTVarIO lastExport+    samples <-+        fmap (filter $ (lastExport' <) . fst)+            . concatMapM Instrument.sample+            . HashMap.elems+            =<< STM.readTVarIO instruments+    unless (null samples) do+        let (timestamps, metrics) = unzip samples+        mapM_ export metrics+        STM.atomically . STM.writeTVar lastExport . maximum $ timestamps++newMetrics :: (IOE :> es, OTLP Measurement :> es) => Eff es (StaticRep Metrics)+newMetrics = do+    export <- exportIO+    instruments <- liftIO . STM.newTVarIO $ mempty+    lastExport <- liftIO . STM.newTVarIO $ Timestamp.epoch+    pure Metrics{..}++-- | Periodically call 'exportInstruments' in the background.+-- The background thread runs for the duration of the given action.+-- Exports all remaining instruments when done.+runMetricsThread+    :: (OTLP Measurement :> es, IOE :> es)+    => Natural+    -> Eff (Metrics ': es) a+    -> Eff es a+runMetricsThread intervalMs eff = do+    metrics <- newMetrics+    bracket+        ( liftIO . IO.forkIO . forever $ do+            IO.threadDelay . fromIntegral $ intervalMs * 1_000+            exportInstruments metrics+        )+        ( \threadId -> liftIO do+            IO.killThread threadId+            exportInstruments metrics+        )+        . const+        . evalStaticRep metrics+        $ eff++-- | Runs the given action and exports all instruments when done,+-- even when the action raises an exception.+runMetricsState+    :: (OTLP Measurement :> es, IOE :> es)+    => Eff (Metrics ': es) a+    -> Eff es a+runMetricsState eff = do+    metrics <- newMetrics+    evalStaticRep metrics $ inject eff `finally` liftIO (exportInstruments metrics)++defaultExportIntervalMs :: Natural+defaultExportIntervalMs = 60_000++-- | Run the 'Metrics' effect, sending telemetry to an exporter.+-- Reads the configuration from <https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#general-sdk-configuration the standard environment variables>.+-- Delegates to 'runNoMetrics' if @OTEL_SDK_DISABLED = true@.+runMetrics+    :: ( IOE :> es+       , Concurrent :> es+       , Environment :> es+       , Retry :> es+       , Timeout :> es+       )+    => Resource+    -> Scope+    -> Eff (Metrics ': es) a+    -> Eff es a+runMetrics resource scope eff =+    ifM Environment.isSdkDisabled (runNoMetrics eff) do+        intervalMs <-+            fromMaybe defaultExportIntervalMs+                . (readMaybe =<<)+                <$> lookupEnv (List.intercalate "_" ["OTEL", exportSignalEnvName @Measurement, "EXPORT_INTERVAL"])+        exporter <- Environment.runConfigError $ Exporter.Environment.lookup @Measurement resource scope+        runOTLPWith exporter+            . runMetricsThread intervalMs+            . inject+            $ eff++-- | Run the 'Metrics' effect, sending telemetry to a collector at the given 'URI'+-- over HTTP with the given 'Encoding'.+runHttpMetrics+    :: (IOE :> es, Concurrent :> es, Retry :> es, Timeout :> es)+    => Resource+    -> Scope+    -> Encoding+    -> Compression+    -> URI+    -> Eff (Metrics ': es) a+    -> Eff es a+runHttpMetrics resource scope encoding compression endpoint = do+    runOTLPWith+        ( Exporter.http+            @Measurement+            resource+            scope+            encoding+            endpoint+            (Export.defaultConfig @Measurement)+            compression+            responseTimeoutDefault+        )+        . runMetricsThread defaultExportIntervalMs+        . inject++-- | Run the 'Metrics' effect, sending telemetry to a gRPC collector.+runGrpcMetrics+    :: (IOE :> es, Concurrent :> es, Retry :> es, Timeout :> es)+    => Resource+    -> Scope+    -> HostName+    -> PortNumber+    -> Compression+    -> Eff (Metrics ': es) a+    -> Eff es a+runGrpcMetrics resource scope host port compression = do+    runOTLPWith+        ( Exporter.grpc @Measurement+            resource+            scope+            host+            port+            (Export.exportGrpcRPC @Measurement)+            (Export.defaultConfig @Measurement)+            compression+            defaultGrpcSendTimeout+        )+        . runMetricsThread defaultExportIntervalMs+        . inject++-- | Run the 'Metrics' effect, printing telemetry to the console rather than sending to a collector.+runConsoleMetrics+    :: (IOE :> es)+    => Maybe Natural+    -- ^ Background sampling interval in ms.+    -- If set, a background thread periodically samples instruments.+    -- Otherwise, they are only sampled once at the end.+    -> Eff (Metrics ': es) a+    -> Eff es a+runConsoleMetrics intervalMs = runMetricsWith intervalMs Console.stdout++-- | Run the 'Metrics' effect, collecting telemetry in-memory rather than sending to a collector.+runInMemoryMetrics+    :: (IOE :> es, Concurrent :> es)+    => Eff (Metrics ': es) a+    -> Eff es (a, [Measurement])+runInMemoryMetrics = runInMemoryOTLP @Measurement . runMetricsState . inject++-- | Run the 'Metrics' effect with a given 'Exporter'.+runMetricsWith+    :: (IOE :> es)+    => Maybe Natural+    -- ^ Background sampling interval in ms.+    -- If set, a background thread periodically samples instruments.+    -- Otherwise, they are only sampled once at the end.+    -> Exporter es Measurement+    -> Eff (Metrics ': es) a+    -> Eff es a+runMetricsWith intervalMs exporter =+    runOTLPWith exporter . maybe runMetricsState runMetricsThread intervalMs . inject++-- | Run the 'Metrics' effect as a no-op action.+runNoMetrics :: (IOE :> es) => Eff (Metrics ': es) a -> Eff es a+runNoMetrics = runNoOTLP @Measurement . runMetricsState . inject
+ src/Effectful/OpenTelemetry/Metrics/Gauge.hs view
@@ -0,0 +1,94 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-redundant-constraints #-}++module Effectful.OpenTelemetry.Metrics.Gauge+    ( Gauge+    , new+    , set+    , setIO+    , measurement+    , Payload (..)+    )+where++import Control.Concurrent.STM (TMVar)+import Control.Concurrent.STM qualified as STM+import Data.Aeson.Types (ToJSON (..))+import Data.Maybe (maybeToList)+import Data.Scientific (Scientific)+import Data.Text (Text)+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+    ( Measurement+    , Metric (..)+    , NumberDataPoint (..)+    )+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++-- | A synchronous 'Instrument' which records non-additive values when changes occur.+--+-- Example uses for 'Gauge':+--+-- - subscribe to change events for the background noise level+-- - subscribe to change events for the CPU fan speed+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/data-model/#gauge the OpenTelemetry spec>.+data Gauge = Gauge+    { name :: Text+    , startTime :: Timestamp+    , value :: TMVar (Timestamp, Scientific)+    , metadata :: Metadata+    }++-- | Create a new 'Gauge' and register it to be sampled and exported.+new :: (Metrics :> es) => Text -> Metadata -> Eff es Gauge+new name metadata =+    Metrics.register =<< unsafeEff_ do+        startTime <- Timestamp.now+        value <- STM.newEmptyTMVarIO+        pure Gauge{..}++-- | Set a 'Gauge' to the given value.+set :: (Metrics :> es) => Gauge -> Scientific -> Eff es ()+set = (unsafeEff_ .) . setIO++setIO :: Gauge -> Scientific -> IO ()+setIO gauge value = do+    time <- Timestamp.now+    STM.atomically . STM.writeTMVar gauge.value $ (time, value)++-- | A gauge 'Measurement' of a single observed data point.+measurement :: Metadata -> Text -> NumberDataPoint -> Measurement+measurement metadata name dataPoint =+    Measurement.create name metadata Payload{dataPoints = [dataPoint]}++instance Instrument Gauge where+    name = name+    sample Gauge{metadata = metadata@Metadata{..}, ..} = do+        measurements <- maybeToList <$> STM.atomically (STM.tryReadTMVar value)+        pure+            [ (time, measurement metadata name NumberDataPoint{..})+            | (time, value) <- measurements+            ]++newtype Payload = Payload {dataPoints :: [NumberDataPoint]}+    deriving stock (Generic, Show, Eq)+    deriving anyclass (ToJSON)++instance Proto.Encode Payload where+    encode Payload{..} = foldMap (Proto.encodeField 1) dataPoints++instance Metric Payload where+    metricJsonKey _ = "gauge"+    metricProtoFieldNumber _ = 5
+ src/Effectful/OpenTelemetry/Metrics/Histogram.hs view
@@ -0,0 +1,176 @@+{-# 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
+ src/Effectful/OpenTelemetry/Metrics/Instrument.hs view
@@ -0,0 +1,29 @@+{-# LANGUAGE GADTs #-}++module Effectful.OpenTelemetry.Metrics.Instrument where++import Data.Text (Text)+import Effectful.OpenTelemetry.Metrics.Measurement (Measurement)+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Prelude++-- | Used to report 'Measurement's.+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/api/#instrument the OpenTelemetry spec>.+class Instrument i where+    -- | The name of the instrument.+    name :: i -> Text++    -- | Take a reading of the instrument, along with the time it was recorded.+    --+    -- An instrument may report several 'Measurement's: an asynchronous+    -- instrument backed by a multiple-instrument callback reports one per+    -- metric it covers. Reporting none means there is nothing to export yet.+    sample :: i -> IO [(Timestamp, Measurement)]++-- | A wrapper type for any 'Instrument' @i@ for all @i@.+data SomeInstrument = forall i. (Instrument i) => SomeInstrument i++instance Instrument SomeInstrument where+    name (SomeInstrument i) = name i+    sample (SomeInstrument i) = sample i
+ src/Effectful/OpenTelemetry/Metrics/Measurement.hs view
@@ -0,0 +1,236 @@+{-# LANGUAGE DuplicateRecordFields #-}+{-# LANGUAGE GADTs #-}++module Effectful.OpenTelemetry.Metrics.Measurement where++import Data.Aeson.KeyMap qualified as KeyMap+import Data.Aeson.Types (Key, ToJSON (..), Value (..), object, (.=))+import Data.Foldable qualified as Foldable+import Data.Functor ((<&>))+import Data.Scientific (Scientific)+import Data.Scientific qualified as Scientific+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Typeable (Typeable, cast)+import Data.Word (Word64)+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import Effectful.OpenTelemetry.Protocol.Attributes (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.Timestamp (Timestamp)+import GHC.Generics (Generic)+import Network.GRPC.HTTP2.Proto3Wire (RPC (..))+import Prettyprinter (Pretty (..))+import Prettyprinter.Extra (PrettyAnn (..))+import Prettyprinter.Extra qualified as Pretty+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | Represents a data point reported via the OpenTelemetry metrics API.+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/api/#measurement the OpenTelemetry spec>.+data Measurement where+    Measurement+        :: (Metric d)+        => { name :: Text+           , description :: Maybe Text+           , unit :: Maybe Text+           , metric :: d+           }+        -> Measurement++instance Show Measurement where+    show Measurement{..} =+        "Measurement{name="+            <> show name+            <> ",description="+            <> show description+            <> ",unit="+            <> show unit+            <> ",metric="+            <> show (toJSON metric)+            <> "}"++instance ToJSON Measurement where+    toJSON Measurement{..} =+        object+            [ "name" .= name+            , "description" .= description+            , "unit" .= unit+            , metricJsonKey metric .= metric+            ]++-- https://github.com/open-telemetry/opentelemetry-proto/blob/main/opentelemetry/proto/metrics/v1/metrics.proto+instance Proto.Encode Measurement where+    encode Measurement{..} =+        mconcat+            [ Proto.encodeField 1 name+            , foldMap (Proto.encodeField 2) description+            , foldMap (Proto.encodeField 3) unit+            , Proto.encodeField (metricProtoFieldNumber metric) metric+            ]++instance PrettyAnn ann Measurement where+    prettyAnn Measurement{..} =+        Pretty.unwords . filter (not . Pretty.null) $+            [ pretty name <> ":"+            , case toJSON metric of+                Object (KeyMap.lookup "dataPoints" -> Just (Array dataPoints)) ->+                    Pretty.intercalate ", " $+                        Foldable.toList dataPoints <&> \case+                            Object dataPoint ->+                                prettyAnn . Object $+                                    KeyMap.filterWithKey+                                        (\key _ -> key `elem` ["asDouble", "asInt", "count", "max", "min", "sum"])+                                        dataPoint+                            val -> prettyAnn val+                val -> prettyAnn val+            , maybe "" pretty unit+            ]++instance Export.Request Measurement where+    exportHttpPathComponents = ["v1", "metrics"]+    exportGrpcRPC =+        RPC+            { pkg = "opentelemetry.proto.collector.metrics.v1"+            , srv = "MetricsService"+            , meth = "Export"+            }+    exportJson = object . pure . ("resourceMetrics" .=) . fmap resourceMetrics+      where+        resourceMetrics :: Resource.Items Measurement -> Value+        resourceMetrics Resource.Items{..} =+            object+                [ "resource" .= resource+                , "scopeMetrics" .= fmap scopeMetrics scopeItems+                ]+        scopeMetrics :: Scope.Items Measurement -> Value+        scopeMetrics Scope.Items{..} =+            object+                [ "scope" .= scope+                , "metrics" .= items+                ]++    transportSignalEnvName = "METRICS"+    exportSignalEnvName = "METRIC"+    batchEnvPrefix = Nothing++    defaultConfig =+        Export.Config+            { batch = Nothing+            , exportTimeoutMs = 30_000+            }++-- | Build a 'Measurement' for a metric of the given name, taking its+-- description and unit from the instrument's 'Metadata'.+create :: (Metric d) => Text -> Metadata -> d -> Measurement+create name Metadata{..} metric = Measurement{..}++-- | How to embed a metric in an OTLP payload.+class (Typeable d, ToJSON d, Proto.Encode d) => Metric d where+    -- | JSON field name under which this metric is nested.+    metricJsonKey :: d -> Key++    -- | Proto field number under which this metric is nested.+    metricProtoFieldNumber :: d -> Proto.FieldNumber++-- | Recover a 'Measurement''s underlying metric, if it is of the given type.+toMetric :: (Metric d) => Measurement -> Maybe d+toMetric Measurement{metric} = cast metric++-- | A single data point in a timeseries that describes the time-varying+-- scalar value of a metric.+data NumberDataPoint = NumberDataPoint+    { attributes :: Attributes+    , startTime :: Timestamp+    , time :: Timestamp+    , value :: Scientific+    }+    deriving stock (Generic, Show, Eq)++instance ToJSON NumberDataPoint where+    toJSON dp =+        object $+            [ "attributes" .= dp.attributes+            , "startTimeUnixNano" .= dp.startTime+            , "timeUnixNano" .= dp.time+            ]+                <> case Scientific.floatingOrInteger @Double @Integer dp.value of+                    Left _ -> ["asDouble" .= dp.value]+                    Right i -> ["asInt" .= Text.show i]++instance Proto.Encode NumberDataPoint where+    encode dp =+        mconcat+            [ Proto.encodeField 7 dp.attributes+            , Proto.encodeField 2 dp.startTime+            , Proto.encodeField 3 dp.time+            , either (Proto.double 4) (Proto.sfixed64 6) $+                Scientific.floatingOrInteger dp.value+            ]++-- | A single data point in a timeseries that describes the time-varying values of a histogram.+data HistogramDataPoint = HistogramDataPoint+    { attributes :: Attributes+    , startTime :: Timestamp+    , time :: Timestamp+    , count :: Word64+    , sum :: Maybe Scientific+    , bucketCounts :: [Word64]+    , explicitBounds :: [Scientific]+    , minValue :: Maybe Scientific+    , maxValue :: Maybe Scientific+    }+    deriving stock (Generic, Show, Eq)++instance ToJSON HistogramDataPoint where+    toJSON dp =+        object $+            [ "attributes" .= dp.attributes+            , "startTimeUnixNano" .= dp.startTime+            , "timeUnixNano" .= dp.time+            , "count" .= Text.show dp.count+            , "bucketCounts" .= (Text.show <$> dp.bucketCounts)+            , "explicitBounds" .= dp.explicitBounds+            ]+                <> ["sum" .= v | Just v <- [dp.sum]]+                <> ["min" .= v | Just v <- [dp.minValue]]+                <> ["max" .= v | Just v <- [dp.maxValue]]++instance Proto.Encode HistogramDataPoint where+    encode dp =+        mconcat+            [ Proto.encodeField 2 dp.startTime+            , Proto.encodeField 3 dp.time+            , Proto.fixed64 4 dp.count+            , foldMap (Proto.double 5 . Scientific.toRealFloat) dp.sum+            , foldMap (Proto.fixed64 6) dp.bucketCounts+            , foldMap (Proto.double 7 . Scientific.toRealFloat) dp.explicitBounds+            , Proto.encodeField 9 dp.attributes+            , foldMap (Proto.double 11 . Scientific.toRealFloat) dp.minValue+            , foldMap (Proto.double 12 . Scientific.toRealFloat) dp.maxValue+            ]++-- | Defines how a metric aggregator reports aggregated values.+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/data-model/#temporality the OpenTelemetry spec>.+data AggregationTemporality+    = -- | Aggregator reports changes since /last report time/.+      -- This means that successive data points /advance/ the starting timestamp.+      TemporalityDelta+    | -- | Aggregator reports changes since /a fixed start time/.+      -- This means that successive data points /repeat/ the starting timestamp.+      TemporalityCumulative+    deriving stock (Generic, Show, Eq, Bounded)++instance Enum AggregationTemporality where+    toEnum 1 = TemporalityDelta+    toEnum 2 = TemporalityCumulative+    toEnum _ = error "Enum.AggregationTemporality.toEnum: bad argument"++    fromEnum TemporalityDelta = 1+    fromEnum TemporalityCumulative = 2++instance {-# OVERLAPPING #-} Proto.EncodeField AggregationTemporality where+    encodeField n = Proto.int32 n . fromIntegral . fromEnum
+ src/Effectful/OpenTelemetry/Metrics/Metadata.hs view
@@ -0,0 +1,34 @@+module Effectful.OpenTelemetry.Metrics.Metadata where++import Control.Applicative ((<|>))+import Data.Text (Text)+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+import GHC.Generics (Generic)+import Prelude++-- | The optional descriptive parameters an instrument may carry.+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/api/#instrument the OpenTelemetry spec>.+data Metadata = Metadata+    { attributes :: Attributes+    , description :: Maybe Text+    , unit :: Maybe Text+    }+    deriving stock (Generic, Show, Eq)++-- | Merges 'attributes', right-biased for the 'Maybe' fields.+instance Semigroup Metadata where+    a <> b =+        Metadata+            { attributes = a.attributes <> b.attributes+            , description = b.description <|> a.description+            , unit = b.unit <|> a.unit+            }++instance Monoid Metadata where+    mempty =+        Metadata+            { attributes = mempty+            , description = Nothing+            , unit = Nothing+            }
+ src/Effectful/OpenTelemetry/Metrics/ObservableCounter.hs view
@@ -0,0 +1,57 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-redundant-constraints #-}++module Effectful.OpenTelemetry.Metrics.ObservableCounter (ObservableCounter, new) where++import Data.Scientific (Scientific)+import Data.Text (Text)+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 (NumberDataPoint (..))+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import Effectful.OpenTelemetry.Metrics.Sum (Monotonicity (..))+import Effectful.OpenTelemetry.Metrics.Sum qualified as Sum+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Timestamp qualified as Timestamp+import Prelude++-- | An asynchronous 'Instrument' which reports monotonically increasing values when it is observed.+-- Example uses for Asynchronous Counter:+--+-- - CPU time, which could be reported for each thread, each process or the entire system+-- - the number of page faults for each process+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/api/#asynchronous-counter the OpenTelemetry spec>.+data ObservableCounter = ObservableCounter+    { name :: Text+    , startTime :: Timestamp+    , metadata :: Metadata+    , observe :: IO Scientific+    }++-- | Create a new 'ObservableCounter'. The callback must return the cumulative total, not a delta.+new+    :: (Metrics :> es)+    => Text+    -> Metadata+    -> IO Scientific+    -- ^ The cumulative total+    -> Eff es ObservableCounter+new name metadata observe =+    Metrics.register =<< unsafeEff_ do+        startTime <- Timestamp.now+        pure ObservableCounter{..}++instance Instrument ObservableCounter where+    name = name+    sample ObservableCounter{..} = do+        time <- Timestamp.now+        value <- observe+        pure . pure $+            ( time+            , Sum.measurement Monotonic metadata name NumberDataPoint{attributes = metadata.attributes, ..}+            )
+ src/Effectful/OpenTelemetry/Metrics/ObservableGauge.hs view
@@ -0,0 +1,57 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-redundant-constraints #-}++module Effectful.OpenTelemetry.Metrics.ObservableGauge (ObservableGauge, new) where++import Data.Maybe (maybeToList)+import Data.Scientific (Scientific)+import Data.Text (Text)+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.Gauge qualified as Gauge+import Effectful.OpenTelemetry.Metrics.Instrument (Instrument)+import Effectful.OpenTelemetry.Metrics.Instrument qualified as Instrument+import Effectful.OpenTelemetry.Metrics.Measurement (NumberDataPoint (..))+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Timestamp qualified as Timestamp+import Prelude++-- | An asynchronous 'Instrument' which reports non-additive values when it is observed.+--+-- Example uses for 'ObservableGauge':+--+-- - the current room temperature+-- - the CPU fan speed+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/api/#asynchronous-gauge the OpenTelemetry spec>.+data ObservableGauge = ObservableGauge+    { name :: Text+    , startTime :: Timestamp+    , metadata :: Metadata+    , observe :: IO (Maybe Scientific)+    }++-- | Create a new 'ObservableGauge' with an observation callback.+new+    :: (Metrics :> es)+    => Text+    -> Metadata+    -> IO (Maybe Scientific)+    -> Eff es ObservableGauge+new name metadata observe =+    Metrics.register =<< unsafeEff_ do+        startTime <- Timestamp.now+        pure ObservableGauge{..}++instance Instrument ObservableGauge where+    name = name+    sample ObservableGauge{..} = do+        time <- Timestamp.now+        values <- maybeToList <$> observe+        pure+            [ (time, Gauge.measurement metadata name NumberDataPoint{attributes = metadata.attributes, ..})+            | value <- values+            ]
+ src/Effectful/OpenTelemetry/Metrics/ObservableUpDownCounter.hs view
@@ -0,0 +1,57 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-redundant-constraints #-}++module Effectful.OpenTelemetry.Metrics.ObservableUpDownCounter (ObservableUpDownCounter, new) where++import Data.Scientific (Scientific)+import Data.Text (Text)+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 (NumberDataPoint (..))+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import Effectful.OpenTelemetry.Metrics.Sum (Monotonicity (..))+import Effectful.OpenTelemetry.Metrics.Sum qualified as Sum+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Timestamp qualified as Timestamp+import Prelude++-- | An asynchronous 'Instrument' which reports additive values when it is observed.+--+-- Example uses for 'ObservableUpDownCounter':+--+-- - the process heap size+-- - the approximate number of items in a lock-free circular buffer+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/api/#asynchronous-updowncounter the OpenTelemetry spec>.+data ObservableUpDownCounter = ObservableUpDownCounter+    { name :: Text+    , startTime :: Timestamp+    , metadata :: Metadata+    , observe :: IO Scientific+    }++-- | Create a new 'ObservableUpDownCounter' with an observation callback.+new+    :: (Metrics :> es)+    => Text+    -> Metadata+    -> IO Scientific+    -> Eff es ObservableUpDownCounter+new name metadata observe =+    Metrics.register =<< unsafeEff_ do+        startTime <- Timestamp.now+        pure ObservableUpDownCounter{..}++instance Instrument ObservableUpDownCounter where+    name = name+    sample ObservableUpDownCounter{..} = do+        time <- Timestamp.now+        value <- observe+        pure . pure $+            ( time+            , Sum.measurement NonMonotonic metadata name NumberDataPoint{attributes = metadata.attributes, ..}+            )
+ src/Effectful/OpenTelemetry/Metrics/RTS.hs view
@@ -0,0 +1,242 @@+{-# LANGUAGE CPP #-}+{-# OPTIONS_GHC -Wno-redundant-constraints #-}++module Effectful.OpenTelemetry.Metrics.RTS where++import Control.Monad.Extra (whenM)+import Data.Functor (void)+import Data.Text (Text)+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.Gauge qualified as Gauge+import Effectful.OpenTelemetry.Metrics.Instrument (Instrument)+import Effectful.OpenTelemetry.Metrics.Instrument qualified as Instrument+import Effectful.OpenTelemetry.Metrics.Measurement (Measurement, NumberDataPoint (..))+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import Effectful.OpenTelemetry.Metrics.Sum (Monotonicity (..))+import Effectful.OpenTelemetry.Metrics.Sum qualified as Sum+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Timestamp qualified as Timestamp+import GHC.Stats (GCDetails (..), RTSStats (..))+import GHC.Stats qualified as Stats+import Prelude++-- | An asynchronous multi-'Instrument' which reports RTS metrics.+--+-- WARNING: Sampling throws an error if RTS stats are not enabled.+-- See 'Stats.getRTSStatsEnabled' for details.+--+-- Counters:+--+-- [@rts.gcs@] @gcs@: Total number of GCs.+-- [@rts.major_gcs@] @major_gcs@: Total number of major (oldest generation)+--   GCs.+-- [@rts.allocated_bytes@] @allocated_bytes@: Total bytes allocated.+-- [@rts.cumulative_live_bytes@] @cumulative_live_bytes@: Sum of live bytes+--   across all major GCs. Divided by @major_gcs@ gives the average live+--   data over the lifetime of the program.+-- [@rts.copied_bytes@] @copied_bytes@: Sum of @copied_bytes@ across all GCs.+-- [@rts.par_copied_bytes@] @par_copied_bytes@: Sum of @copied_bytes@ across+--   all parallel GCs.+-- [@rts.cumulative_par_balanced_copied_bytes@]+--   @cumulative_par_balanced_copied_bytes@: Sum of @par_balanced_copied@+--   bytes across all parallel GCs. (@base >= 4.11@)+-- [@rts.init_cpu_ns@] @init_cpu_ns@: Total CPU time used by the init phase.+--   (@base >= 4.12@)+-- [@rts.init_elapsed_ns@] @init_elapsed_ns@: Total elapsed time used by the+--   init phase. (@base >= 4.12@)+-- [@rts.mutator_cpu_ns@] @mutator_cpu_ns@: Total CPU time used by the+--   mutator.+-- [@rts.mutator_elapsed_ns@] @mutator_elapsed_ns@: Total elapsed time used+--   by the mutator.+-- [@rts.gc_cpu_ns@] @gc_cpu_ns@: Total CPU time used by the GC.+-- [@rts.gc_elapsed_ns@] @gc_elapsed_ns@: Total elapsed time used by the GC.+-- [@rts.cpu_ns@] @cpu_ns@: Total CPU time (at the previous GC).+-- [@rts.elapsed_ns@] @elapsed_ns@: Total elapsed time (at the previous GC).+-- [@rts.nonmoving_gc_sync_cpu_ns@] @nonmoving_gc_sync_cpu_ns@: The total CPU+--   time used during the post-mark pause phase of the concurrent nonmoving+--   GC. (@base >= 4.15@)+-- [@rts.nonmoving_gc_sync_elapsed_ns@] @nonmoving_gc_sync_elapsed_ns@: The+--   total time elapsed during the post-mark pause phase of the concurrent+--   nonmoving GC. (@base >= 4.15@)+-- [@rts.nonmoving_gc_cpu_ns@] @nonmoving_gc_cpu_ns@: The total CPU time used+--   by the nonmoving GC. (@base >= 4.15@)+-- [@rts.nonmoving_gc_elapsed_ns@] @nonmoving_gc_elapsed_ns@: The total time+--   elapsed during which there is a nonmoving GC active. (@base >= 4.15@)+--+-- Gauges:+--+-- [@rts.gc.allocated_bytes@] @gcdetails_allocated_bytes@: Number of bytes+--   allocated since the previous GC.+-- [@rts.gc.live_bytes@] @gcdetails_live_bytes@: Total amount of live data+--   in the heap (includes large + compact data). Updated after every GC.+--   Data in uncollected generations (in minor GCs) are considered live.+-- [@rts.gc.large_objects_bytes@] @gcdetails_large_objects_bytes@: Total+--   amount of live data in large objects.+-- [@rts.gc.compact_bytes@] @gcdetails_compact_bytes@: Total amount of live+--   data in compact regions.+-- [@rts.gc.slop_bytes@] @gcdetails_slop_bytes@: Total amount of slop+--   (wasted memory).+-- [@rts.gc.mem_in_use_bytes@] @gcdetails_mem_in_use_bytes@: Total amount of+--   memory in use by the RTS.+-- [@rts.gc.copied_bytes@] @gcdetails_copied_bytes@: Total amount of data+--   copied during this GC.+-- [@rts.gc.par_balanced_copied_bytes@] @gcdetails_par_balanced_copied_bytes@:+--   In parallel GC, the amount of balanced data copied by all threads.+--   (@base >= 4.11@)+-- [@rts.gc.block_fragmentation_bytes@] @gcdetails_block_fragmentation_bytes@:+--   The amount of memory lost due to block fragmentation in bytes. Block+--   fragmentation is the difference between the amount of blocks retained+--   by the RTS and the blocks that are in use. This occurs when megablocks+--   are only sparsely used, eg, when data that cannot be moved retains a+--   megablock. (@base >= 4.18@)+-- [@rts.gc.threads@] @gcdetails_threads@: Number of threads used in this GC.+-- [@rts.max_live_bytes@] @max_live_bytes@: Maximum live data (including+--   large objects + compact regions) in the heap. Updated after a major GC.+-- [@rts.max_large_objects_bytes@] @max_large_objects_bytes@: Maximum live+--   data in large objects.+-- [@rts.max_compact_bytes@] @max_compact_bytes@: Maximum live data in+--   compact regions.+-- [@rts.max_slop_bytes@] @max_slop_bytes@: Maximum slop.+-- [@rts.max_mem_in_use_bytes@] @max_mem_in_use_bytes@: Maximum memory in+--   use by the RTS.+-- [@rts.nonmoving_gc_sync_max_elapsed_ns@]+--   @nonmoving_gc_sync_max_elapsed_ns@: The maximum elapsed length of any+--   post-mark pause phase of the concurrent nonmoving GC. (@base >= 4.15@)+-- [@rts.nonmoving_gc_max_elapsed_ns@] @nonmoving_gc_max_elapsed_ns@: The+--   maximum time elapsed during any nonmoving GC cycle. (@base >= 4.15@)+-- [@rts.gc.gen@] @gcdetails_gen@: The generation number of this GC.+-- [@rts.gc.sync_elapsed_ns@] @gcdetails_sync_elapsed_ns@: The time elapsed+--   during synchronisation before GC.+-- [@rts.gc.cpu_ns@] @gcdetails_cpu_ns@: The CPU time used during GC itself.+-- [@rts.gc.elapsed_ns@] @gcdetails_elapsed_ns@: The time elapsed during GC+--   itself.+newtype RTS = RTS {startTime :: Timestamp}++-- | Create a new 'RTS' and register it to be sampled and exported.+new :: (Metrics :> es) => Eff es RTS+new =+    Metrics.register =<< unsafeEff_ do+        startTime <- Timestamp.now+        pure RTS{..}++-- | Register 'RTS' metrics to be sampled and exported.+-- WARNING: Errors if RTS stats are not enabled. Use 'run' to suppress this error.+register :: (Metrics :> es) => Eff es ()+register = void new++-- | If RTS stats are enabled, register 'RTS' metrics to be sampled and exported.+-- Otherwise no-op.+-- See 'Stats.getRTSStatsEnabled' for details.+run :: (IOE :> es, Metrics :> es) => Eff es a -> Eff es a+run eff = do+    whenM (liftIO Stats.getRTSStatsEnabled) register+    eff++instance Instrument RTS where+    name _ = "rts"+    sample rts = do+        time <- Timestamp.now+        stats <- Stats.getRTSStats+        pure [(time, m) | m <- measurements rts time stats]++{- FOURMOLU_DISABLE -}+measurements :: RTS -> Timestamp -> RTSStats -> [Measurement]+measurements RTS{startTime} time rts =+    [ Sum.measurement Monotonic (metadata "{gc}" "Total number of GCs") "rts.gcs" (point rts.gcs)+    , Sum.measurement Monotonic (metadata "{gc}" "Total number of major (oldest generation) GCs") "rts.major_gcs" (point rts.major_gcs)+    , Sum.measurement Monotonic (metadata "By" "Total bytes allocated") "rts.allocated_bytes" (point rts.allocated_bytes)+    , Sum.measurement Monotonic (metadata "By" "Sum of live bytes across all major GCs") "rts.cumulative_live_bytes" (point rts.cumulative_live_bytes)+    , Sum.measurement Monotonic (metadata "By" "Sum of copied bytes across all GCs") "rts.copied_bytes" (point rts.copied_bytes)+    , Sum.measurement Monotonic (metadata "By" "Sum of copied bytes across all parallel GCs") "rts.par_copied_bytes" (point rts.par_copied_bytes)+#if MIN_VERSION_base(4,11,0)+    , Sum.measurement+        Monotonic+        (metadata "By" "Sum of balanced copied bytes across all parallel GCs")+        "rts.cumulative_par_balanced_copied_bytes"+        (point rts.cumulative_par_balanced_copied_bytes)+#endif+#if MIN_VERSION_base(4,12,0)+    , Sum.measurement Monotonic (metadata "ns" "Total CPU time used by the init phase") "rts.init_cpu_ns" (point rts.init_cpu_ns)+    , Sum.measurement Monotonic (metadata "ns" "Total elapsed time used by the init phase") "rts.init_elapsed_ns" (point rts.init_elapsed_ns)+#endif+    , Sum.measurement Monotonic (metadata "ns" "Total CPU time used by the mutator") "rts.mutator_cpu_ns" (point rts.mutator_cpu_ns)+    , Sum.measurement Monotonic (metadata "ns" "Total elapsed time used by the mutator") "rts.mutator_elapsed_ns" (point rts.mutator_elapsed_ns)+    , Sum.measurement Monotonic (metadata "ns" "Total CPU time used by the GC") "rts.gc_cpu_ns" (point rts.gc_cpu_ns)+    , Sum.measurement Monotonic (metadata "ns" "Total elapsed time used by the GC") "rts.gc_elapsed_ns" (point rts.gc_elapsed_ns)+    , Sum.measurement Monotonic (metadata "ns" "Total CPU time (at the previous GC)") "rts.cpu_ns" (point rts.cpu_ns)+    , Sum.measurement Monotonic (metadata "ns" "Total elapsed time (at the previous GC)") "rts.elapsed_ns" (point rts.elapsed_ns)+#if MIN_VERSION_base(4,15,0)+    , Sum.measurement+        Monotonic+        (metadata "ns" "CPU time used during the post-mark pause phase of the nonmoving GC")+        "rts.nonmoving_gc_sync_cpu_ns"+        (point rts.nonmoving_gc_sync_cpu_ns)+    , Sum.measurement+        Monotonic+        (metadata "ns" "Elapsed time during the post-mark pause phase of the nonmoving GC")+        "rts.nonmoving_gc_sync_elapsed_ns"+        (point rts.nonmoving_gc_sync_elapsed_ns)+    , Sum.measurement+        Monotonic+        (metadata "ns" "Total CPU time used by the nonmoving GC")+        "rts.nonmoving_gc_cpu_ns"+        (point rts.nonmoving_gc_cpu_ns)+    , Sum.measurement+        Monotonic+        (metadata "ns" "Total time elapsed during which the nonmoving GC is active")+        "rts.nonmoving_gc_elapsed_ns"+        (point rts.nonmoving_gc_elapsed_ns)+#endif+    , Gauge.measurement (metadata "By" "Bytes allocated since the previous GC") "rts.gc.allocated_bytes" (point rts.gc.gcdetails_allocated_bytes)+    , Gauge.measurement (metadata "By" "Total amount of live data in the heap") "rts.gc.live_bytes" (point rts.gc.gcdetails_live_bytes)+    , Gauge.measurement (metadata "By" "Total amount of live data in large objects") "rts.gc.large_objects_bytes" (point rts.gc.gcdetails_large_objects_bytes)+    , Gauge.measurement (metadata "By" "Total amount of live data in compact regions") "rts.gc.compact_bytes" (point rts.gc.gcdetails_compact_bytes)+    , Gauge.measurement (metadata "By" "Total amount of slop (wasted memory)") "rts.gc.slop_bytes" (point rts.gc.gcdetails_slop_bytes)+    , Gauge.measurement (metadata "By" "Total amount of memory in use by the RTS") "rts.gc.mem_in_use_bytes" (point rts.gc.gcdetails_mem_in_use_bytes)+    , Gauge.measurement (metadata "By" "Total amount of data copied during the most recent GC") "rts.gc.copied_bytes" (point rts.gc.gcdetails_copied_bytes)+#if MIN_VERSION_base(4,11,0)+    , Gauge.measurement+        (metadata "By" "Amount of balanced data copied by all threads in the most recent parallel GC")+        "rts.gc.par_balanced_copied_bytes"+        (point rts.gc.gcdetails_par_balanced_copied_bytes)+#endif+#if MIN_VERSION_base(4,18,0)+    , Gauge.measurement+        (metadata "By" "Memory lost due to block fragmentation in the most recent GC")+        "rts.gc.block_fragmentation_bytes"+        (point rts.gc.gcdetails_block_fragmentation_bytes)+#endif+    , Gauge.measurement (metadata "{thread}" "Number of threads used in the most recent GC") "rts.gc.threads" (point rts.gc.gcdetails_threads)+    , Gauge.measurement (metadata "By" "Maximum live data in the heap") "rts.max_live_bytes" (point rts.max_live_bytes)+    , Gauge.measurement (metadata "By" "Maximum live data in large objects") "rts.max_large_objects_bytes" (point rts.max_large_objects_bytes)+    , Gauge.measurement (metadata "By" "Maximum live data in compact regions") "rts.max_compact_bytes" (point rts.max_compact_bytes)+    , Gauge.measurement (metadata "By" "Maximum slop") "rts.max_slop_bytes" (point rts.max_slop_bytes)+    , Gauge.measurement (metadata "By" "Maximum memory in use by the RTS") "rts.max_mem_in_use_bytes" (point rts.max_mem_in_use_bytes)+#if MIN_VERSION_base(4,15,0)+    , Gauge.measurement+        (metadata "ns" "Maximum elapsed time of any post-mark pause phase of the nonmoving GC")+        "rts.nonmoving_gc_sync_max_elapsed_ns"+        (point rts.nonmoving_gc_sync_max_elapsed_ns)+    , Gauge.measurement+        (metadata "ns" "Maximum time elapsed during any nonmoving GC cycle")+        "rts.nonmoving_gc_max_elapsed_ns"+        (point rts.nonmoving_gc_max_elapsed_ns)+#endif+    , Gauge.measurement (metadata "1" "Generation number of the most recent GC") "rts.gc.gen" (point rts.gc.gcdetails_gen)+    , Gauge.measurement+        (metadata "ns" "Time elapsed during synchronisation before the most recent GC")+        "rts.gc.sync_elapsed_ns"+        (point rts.gc.gcdetails_sync_elapsed_ns)+    , Gauge.measurement (metadata "ns" "CPU time used during the most recent GC") "rts.gc.cpu_ns" (point rts.gc.gcdetails_cpu_ns)+    , Gauge.measurement (metadata "ns" "Time elapsed during the most recent GC") "rts.gc.elapsed_ns" (point rts.gc.gcdetails_elapsed_ns)+    ]+  where+    metadata :: Text -> Text -> Metadata+    metadata unit description = Metadata{attributes = mempty, description = Just description, unit = Just unit}++    point :: (Integral a) => a -> NumberDataPoint+    point (fromIntegral -> value) = NumberDataPoint{attributes = mempty, ..}+{- FOURMOLU_ENABLE -}
+ src/Effectful/OpenTelemetry/Metrics/Sum.hs view
@@ -0,0 +1,60 @@+module Effectful.OpenTelemetry.Metrics.Sum where++import Data.Aeson.Types (ToJSON (..), object, (.=))+import Data.Text (Text)+import Effectful.OpenTelemetry.Metrics.Measurement+    ( AggregationTemporality (..)+    , Measurement+    , Metric (..)+    , NumberDataPoint (..)+    )+import Effectful.OpenTelemetry.Metrics.Measurement qualified as Measurement+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import GHC.Generics (Generic)+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | Whether a sum only ever increases.+data Monotonicity = Monotonic | NonMonotonic+    deriving stock (Eq)++-- | A cumulative sum 'Measurement' of a single data point.+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/data-model/#sums the OpenTelemetry spec>.+measurement :: Monotonicity -> Metadata -> Text -> NumberDataPoint -> Measurement+measurement monotonicity metadata name dataPoint =+    Measurement.create+        name+        metadata+        Payload+            { dataPoints = [dataPoint]+            , isMonotonic = monotonicity == Monotonic+            , aggregationTemporality = TemporalityCumulative+            }++data Payload = Payload+    { dataPoints :: [NumberDataPoint]+    , isMonotonic :: Bool+    , aggregationTemporality :: AggregationTemporality+    }+    deriving stock (Generic, Show, Eq)++instance ToJSON Payload where+    toJSON Payload{..} =+        object+            [ "dataPoints" .= dataPoints+            , "isMonotonic" .= isMonotonic+            , "aggregationTemporality" .= fromEnum aggregationTemporality+            ]++instance Proto.Encode Payload where+    encode Payload{..} =+        mconcat+            [ foldMap (Proto.encodeField 1) dataPoints+            , Proto.encodeField 2 aggregationTemporality+            , Proto.bool 3 isMonotonic+            ]++instance Metric Payload where+    metricJsonKey _ = "sum"+    metricProtoFieldNumber _ = 7
+ src/Effectful/OpenTelemetry/Metrics/UpDownCounter.hs view
@@ -0,0 +1,64 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-redundant-constraints #-}++module Effectful.OpenTelemetry.Metrics.UpDownCounter (UpDownCounter, new, add) where++import Control.Concurrent.STM (TVar)+import Control.Concurrent.STM qualified as STM+import Data.Bifunctor (bimap)+import Data.Scientific (Scientific)+import Data.Text (Text)+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 (NumberDataPoint (..))+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import Effectful.OpenTelemetry.Metrics.Sum (Monotonicity (..))+import Effectful.OpenTelemetry.Metrics.Sum qualified as Sum+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Timestamp qualified as Timestamp+import Prelude++-- | A synchronous 'Instrument' which supports increments and decrements.+--+-- Example uses for 'UpDownCounter':+--+-- - the number of active requests+-- - the number of items in a queue+--+-- See <https://opentelemetry.io/docs/specs/otel/metrics/api/#updowncounter the OpenTelemetry spec>.+data UpDownCounter = UpDownCounter+    { name :: Text+    , startTime :: Timestamp+    , value :: TVar (Timestamp, Scientific)+    , metadata :: Metadata+    }++-- | Create a new 'UpDownCounter' and register it to be sampled and exported.+new :: (Metrics :> es) => Text -> Metadata -> Eff es UpDownCounter+new name metadata =+    Metrics.register =<< unsafeEff_ do+        startTime <- Timestamp.now+        value <- STM.newTVarIO (startTime, 0)+        pure UpDownCounter{..}++-- | Add to an 'UpDownCounter'. The amount may be negative.+add :: (Metrics :> es) => UpDownCounter -> Scientific -> Eff es ()+add = (unsafeEff_ .) . addIO++addIO :: UpDownCounter -> Scientific -> IO ()+addIO UpDownCounter{..} inc = do+    time <- Timestamp.now+    STM.atomically . STM.modifyTVar' value $ bimap (const time) (+ inc)++instance Instrument UpDownCounter where+    name = name+    sample UpDownCounter{..} = do+        (time, value) <- STM.readTVarIO value+        pure . pure $+            ( time+            , Sum.measurement NonMonotonic metadata name NumberDataPoint{attributes = metadata.attributes, ..}+            )
+ src/Effectful/OpenTelemetry/Protocol.hs view
@@ -0,0 +1,18 @@+module Effectful.OpenTelemetry.Protocol+    ( module Effectful.OpenTelemetry.Protocol.AnyValue+    , module Effectful.OpenTelemetry.Protocol.Attributes+    , module Effectful.OpenTelemetry.Protocol.Effect+    , module Effectful.OpenTelemetry.Protocol.Resource+    , module Effectful.OpenTelemetry.Protocol.Scope+    , module Effectful.OpenTelemetry.Protocol.Transport+    , module Proto3.Wire.Class+    )+where++import Effectful.OpenTelemetry.Protocol.AnyValue+import Effectful.OpenTelemetry.Protocol.Attributes+import Effectful.OpenTelemetry.Protocol.Effect+import Effectful.OpenTelemetry.Protocol.Resource (Resource (Resource))+import Effectful.OpenTelemetry.Protocol.Scope (Scope (Scope))+import Effectful.OpenTelemetry.Protocol.Transport+import Proto3.Wire.Class
+ src/Effectful/OpenTelemetry/Protocol/AnyValue.hs view
@@ -0,0 +1,71 @@+module Effectful.OpenTelemetry.Protocol.AnyValue where++import Control.Monad (mzero)+import Data.Aeson.KeyMap qualified as KeyMap+import Data.Aeson.Types+    ( FromJSON (..)+    , ToJSON (..)+    , Value (..)+    , object+    , withObject+    , (.!=)+    , (.:)+    , (.:?)+    , (.=)+    )+import Data.Coerce (coerce)+import Data.Int (Int64)+import Data.Scientific qualified as Scientific+import Data.Text qualified as Text+import Text.Read qualified as Text+import Prelude++-- | Wrapper for 'Value' used for OTEL-specific JSON encoding.+--+-- See <https://opentelemetry.io/docs/specs/otel/common/#anyvalue the OpenTelemetry spec>.+newtype AnyValue = AnyValue Value+    deriving stock (Show, Eq)++instance ToJSON AnyValue where+    toJSON (AnyValue Null) = object []+    toJSON (AnyValue (Bool b)) = object ["boolValue" .= b]+    toJSON (AnyValue (String s)) = object ["stringValue" .= s]+    toJSON (AnyValue (Number d)) = case Scientific.toBoundedInteger @Int64 d of+        Just i -> object ["intValue" .= String (Text.show i)]+        Nothing -> object ["doubleValue" .= d]+    toJSON (AnyValue (Array a)) = object ["arrayValue" .= object ["values" .= (AnyValue <$> a)]]+    toJSON (AnyValue (Object o)) =+        object+            [ "kvlistValue"+                .= object+                    [ "values"+                        .= ( (\(k, v) -> object ["key" .= k, "value" .= AnyValue v])+                                <$> KeyMap.toList o+                           )+                    ]+            ]++instance FromJSON AnyValue where+    parseJSON = withObject "AnyValue" \o -> case KeyMap.toList o of+        [] -> pure $ AnyValue Null+        [("boolValue", Bool b)] -> pure . AnyValue $ Bool b+        [("stringValue", String s)] -> pure . AnyValue $ String s+        [("intValue", String s)]+            | Just n <- Text.readMaybe (Text.unpack s) -> pure . AnyValue . Number $ n+        [("doubleValue", Number n)] -> pure . AnyValue $ Number n+        [("doubleValue", String s)]+            | Just n <- Text.readMaybe (Text.unpack s) -> pure . AnyValue . Number $ n+        [("arrayValue", Object o')] ->+            AnyValue . Array . fmap coerce+                <$> (mapM (parseJSON @AnyValue) =<< o' .:? "values" .!= mempty)+        [("kvlistValue", Object o')] -> do+            AnyValue . Object . KeyMap.fromList+                <$> ( mapM+                        ( withObject "KVPair" $ \kvPair -> do+                            k <- kvPair .:? "key" .!= ""+                            AnyValue v <- kvPair .: "value"+                            pure (k, v)+                        )+                        =<< o' .:? "values" .!= mempty+                    )+        _ -> mzero
+ src/Effectful/OpenTelemetry/Protocol/Attributes.hs view
@@ -0,0 +1,70 @@+module Effectful.OpenTelemetry.Protocol.Attributes where++import Data.Aeson.KeyMap qualified as KeyMap+import Data.Aeson.Types+    ( FromJSON (..)+    , Object+    , Pair+    , ToJSON (..)+    , object+    , withArray+    , withObject+    , (.!=)+    , (.:)+    , (.:?)+    , (.=)+    )+import Data.Vector qualified as Vector+import Effectful.OpenTelemetry.Protocol.AnyValue (AnyValue (..))+import GHC.Generics (Generic)+import GHC.IsList (IsList)+import GHC.IsList qualified as GHC+import Prettyprinter.Extra (PrettyAnn (..))+import Prettyprinter.Extra qualified as Pretty+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | Key-value metadata that can be attached to telemetry.+--+-- See <https://opentelemetry.io/docs/specs/semconv/general/attributes/ the OpenTelemetry spec>.+newtype Attributes = Attributes Object+    deriving stock (Generic)+    deriving newtype (Monoid, Semigroup, Show, Eq)++instance IsList Attributes where+    type Item Attributes = Pair+    fromList = fromList+    toList = toList++null :: Attributes -> Bool+null (Attributes o) = Prelude.null o++fromList :: [Pair] -> Attributes+fromList = Attributes . KeyMap.fromList++toList :: Attributes -> [Pair]+toList (Attributes o) = KeyMap.toList o++instance ToJSON Attributes where+    toJSON =+        toJSON+            . fmap (\(k, v) -> object ["key" .= k, "value" .= AnyValue v])+            . toList++instance FromJSON Attributes where+    parseJSON =+        withArray "Attributes" $+            fmap (Attributes . KeyMap.fromList . Vector.toList)+                . mapM+                    ( withObject "KeyValue" \o -> do+                        k <- o .:? "key" .!= ""+                        AnyValue v <- parseJSON =<< o .: "value"+                        pure (k, v)+                    )++instance {-# OVERLAPPING #-} Proto.EncodeField Attributes where+    encodeField n (Attributes o) =+        foldMap (Proto.encodeField n) $ KeyMap.toList o++instance PrettyAnn ann Attributes where+    prettyAnn = Pretty.unwords . fmap prettyAnn . toList
+ src/Effectful/OpenTelemetry/Protocol/Effect.hs view
@@ -0,0 +1,46 @@+module Effectful.OpenTelemetry.Protocol.Effect {-# WARNING in "x-unstable-interface" "This is an unstable interface." #-} where++import Effectful+import Effectful.Concurrent (Concurrent)+import Effectful.Concurrent.STM (atomically, flushTQueue, newTQueueIO)+import Effectful.Dispatch.Static+import Effectful.OpenTelemetry.Exporter.STM qualified as Exporter+import Effectful.OpenTelemetry.Exporter.Type (Exporter (..), withExporter)+import Prelude++data OTLP a :: Effect++type instance DispatchOf (OTLP _) = 'Static 'WithSideEffects++newtype instance StaticRep (OTLP a) = OTLP {sink :: a -> IO ()}++-- | Run the 'OTLP' effect with a custom 'Exporter'.+-- Passes items to the exporter synchronously as they are emitted, without batching or retrying.+runOTLPWith :: forall a es r. (IOE :> es) => Exporter es a -> Eff (OTLP a ': es) r -> Eff es r+runOTLPWith Exporter{..} eff = withExporter \sink -> evalStaticRep OTLP{sink} eff++-- | Run the 'OTLP' effect, collecting telemetry in-memory rather than sending to a collector.+runInMemoryOTLP+    :: forall a es r+     . (IOE :> es, Concurrent :> es)+    => Eff (OTLP a ': es) r+    -> Eff es (r, [a])+runInMemoryOTLP eff = do+    queue <- newTQueueIO+    r <- runOTLPWith (Exporter.tqueue queue) eff+    items <- atomically $ flushTQueue queue+    pure (r, items)++-- | Run the 'OTLP' effect as a no-op action.+runNoOTLP :: forall a es r. (IOE :> es) => Eff (OTLP a ': es) r -> Eff es r+runNoOTLP = runOTLPWith mempty++export :: (OTLP a :> es) => a -> Eff es ()+export item = do+    sink <- exportIO+    unsafeEff_ $ sink item++exportIO :: forall a es. (OTLP a :> es) => Eff es (a -> IO ())+exportIO = do+    OTLP{..} <- getStaticRep+    pure sink
+ src/Effectful/OpenTelemetry/Protocol/Environment.hs view
@@ -0,0 +1,249 @@+{-# LANGUAGE AllowAmbiguousTypes #-}+{-# OPTIONS_GHC -Wno-name-shadowing #-}++module Effectful.OpenTelemetry.Protocol.Environment {-# WARNING in "x-unstable-interface" "This is an unstable interface." #-} where++import Control.Applicative ((<|>))+import Control.Exception (Exception (..))+import Data.Aeson.Types (Pair, Value (String))+import Data.Functor ((<&>))+import Data.List (intercalate)+import Data.List.Extra (splitOn, trim)+import Data.Maybe (fromMaybe)+import Data.String (IsString (fromString))+import Data.Text qualified as Text+import Effectful+import Effectful.Environment (Environment, lookupEnv)+import Effectful.Error.Static (Error, runErrorNoCallStackWith, throwError)+import Effectful.Exception (throwIO)+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+import Effectful.OpenTelemetry.Protocol.Attributes qualified as Attributes+import Effectful.OpenTelemetry.Protocol.Exception+    ( SomeOTLPException (..)+    , otlpExceptionFromException+    , otlpExceptionToException+    )+import Effectful.OpenTelemetry.Protocol.Export (Request (..))+import Effectful.OpenTelemetry.Protocol.Export qualified as Export+import Effectful.OpenTelemetry.Protocol.GRPC qualified as GRPC+import Effectful.OpenTelemetry.Protocol.HTTP qualified as HTTP+import Effectful.OpenTelemetry.Protocol.Resource (Resource (Resource))+import Effectful.OpenTelemetry.Protocol.Transport+    ( Compression (..)+    , Encoding (..)+    , Protocol (..)+    , Transport (..)+    )+import Network.URI (URI (..), URIAuth (..), parseAbsoluteURI)+import Network.URI qualified as URI+import Numeric.Natural (Natural)+import Safe (readMay)+import Text.Read (readMaybe)+import Prelude++data OTLPConfig = OTLPConfig {variable :: String, value :: String}+    deriving stock (Show, Eq)++data ConfigError+    = ParseEnvConfigError OTLPConfig+    | InvalidPort String+    | InvalidTimeout String+    | UnknownProtocol String+    | UnknownExporter String+    | InvalidOTLPEndpointURI String+    | OTLPEndpointMissingAuthority URI+    | UnsupportedOTLPCompression String+    | InvalidResourceAttribute String+    deriving stock (Show, Eq)++instance Exception ConfigError where+    toException = otlpExceptionToException+    fromException = otlpExceptionFromException+    displayException (ParseEnvConfigError OTLPConfig{..}) =+        "Error parsing " <> variable <> ". Invalid value: " <> value+    displayException (InvalidPort port) = "Invalid port number: " <> port+    displayException (InvalidTimeout timeout) = "Invalid timeout: " <> timeout+    displayException (UnknownProtocol protocol) =+        "Unknown OTEL_EXPORTER_OTLP_PROTOCOL: " <> protocol+    displayException (UnknownExporter name) =+        "Unknown exporter in OTEL_*_EXPORTER: " <> name+    displayException (InvalidOTLPEndpointURI uri) =+        "Invalid OTLP endpoint URI: " <> uri+    displayException (OTLPEndpointMissingAuthority uri) =+        "OTLP endpoint missing authority: " <> show uri+    displayException (UnsupportedOTLPCompression compression) =+        "Unsupported OTLP compression: " <> compression+    displayException (InvalidResourceAttribute pair) =+        "Invalid OTEL_RESOURCE_ATTRIBUTES entry (expected key=value): " <> pair++isSdkDisabled :: (Environment :> es) => Eff es Bool+isSdkDisabled = lookupEnv "OTEL_SDK_DISABLED" <&> (== Just "true")++runConfigError :: Eff (Error ConfigError ': es) a -> Eff es a+runConfigError = runErrorNoCallStackWith (throwIO . SomeOTLPException @ConfigError)++-- | The name of the signal-specific @OTEL_EXPORTER_OTLP_\<SIGNAL\>_*@ variable+-- for signal @a@.+signalEnvName :: forall a. (Export.Request a) => String -> String+signalEnvName suffix = intercalate "_" [exporterEnvPrefix, transportSignalEnvName @a, suffix]++-- | The name of the signal-agnostic @OTEL_EXPORTER_OTLP_*@ variable.+globalEnvName :: String -> String+globalEnvName suffix = intercalate "_" [exporterEnvPrefix, suffix]++exporterEnvPrefix :: String+exporterEnvPrefix = "OTEL_EXPORTER_OTLP"++-- | Read an @OTEL_EXPORTER_OTLP_*@ setting for signal @a@, preferring the+-- signal-specific variable over the global one.+lookupSignalOrGlobal+    :: forall a es+     . (Export.Request a, Environment :> es)+    => String+    -> Eff es (Maybe String)+lookupSignalOrGlobal suffix = do+    signal <- lookupEnv $ signalEnvName @a suffix+    global <- lookupEnv $ globalEnvName suffix+    pure $ signal <|> global++-- | The configured 'Compression' for signal @a@, from+-- @OTEL_EXPORTER_OTLP_<SIGNAL>_COMPRESSION@ or @OTEL_EXPORTER_OTLP_COMPRESSION@.+-- Defaults to 'NoCompression'.+lookupCompression+    :: forall a es+     . (Export.Request a, Environment :> es, Error ConfigError :> es)+    => Eff es Compression+lookupCompression =+    lookupSignalOrGlobal @a "COMPRESSION" >>= \case+        Nothing -> pure NoCompression+        Just name -> maybe (throwError $ UnsupportedOTLPCompression name) pure $ readMaybe name++-- | Detect a 'Resource' from the standard OpenTelemetry environment variables:+-- @OTEL_SERVICE_NAME@ and @OTEL_RESOURCE_ATTRIBUTES@.+-- @OTEL_SERVICE_NAME@ takes precedence over a @service.name@ set via @OTEL_RESOURCE_ATTRIBUTES@.+--+-- See <https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#general-sdk-configuration the OpenTelemetry spec>.+detectResource :: (Environment :> es, Error ConfigError :> es) => Eff es Resource+detectResource = do+    fromAttributes <-+        lookupEnv "OTEL_RESOURCE_ATTRIBUTES"+            >>= \case+                Just (trim -> s) | not (null s) -> either throwError (pure . Resource) $ parseResourceAttributes s+                _ -> pure mempty+    fromServiceName <-+        lookupEnv "OTEL_SERVICE_NAME"+            <&> \case+                Just (trim -> s) | not (null s) -> Resource $ Attributes.fromList [("service.name", String $ Text.pack s)]+                _ -> mempty+    pure $ fromServiceName <> fromAttributes+  where+    parseResourceAttributes :: String -> Either ConfigError Attributes+    parseResourceAttributes =+        fmap Attributes.fromList+            . mapM parsePair+            . filter (not . null)+            . map trim+            . splitOn ","+    parsePair :: String -> Either ConfigError Pair+    parsePair s = case break (== '=') s of+        (key, '=' : value) -> Right (decode key, decode value)+        _ -> Left $ InvalidResourceAttribute s+    decode :: (IsString s) => String -> s+    decode = fromString . URI.unEscapeString . trim++lookupTransport+    :: forall a es+     . ( Export.Request a+       , Error ConfigError :> es+       , Environment :> es+       )+    => Eff es Transport+lookupTransport = do+    protocol <- lookupProtocol+    compression <- lookupCompression @a+    timeoutStr <- lookupWithFallback "TIMEOUT" "30000"+    timeoutMs <- maybe (throwError $ InvalidTimeout timeoutStr) pure $ readMay timeoutStr+    pure Transport{..}+  where+    lookupProtocol =+        lookupWithFallback "PROTOCOL" "http/protobuf" >>= \case+            "http/protobuf" -> do+                HTTP Proto <$> lookupHttpUri+            "http/json" -> do+                HTTP Json <$> lookupHttpUri+            "grpc" -> do+                uri <- lookupGrpcUri+                (host, port) <- case uri of+                    URI{uriAuthority = Just URIAuth{uriRegName, ..}}+                        | ':' : (readMay -> Just port) <- uriPort -> pure (uriRegName, port)+                        | otherwise -> throwError $ InvalidPort uriPort+                    _ -> throwError $ OTLPEndpointMissingAuthority uri+                pure . GRPC host port $ exportGrpcRPC @a+            other -> throwError $ UnknownProtocol other++    lookupEndpoint :: Eff es (Maybe URI, Maybe URI)+    lookupEndpoint = do+        signal <- parseOrFail =<< lookupEnv (signalEnvName @a "ENDPOINT")+        global <- parseOrFail =<< lookupEnv (globalEnvName "ENDPOINT")+        pure (signal, global)++    lookupHttpUri :: Eff es URI+    lookupHttpUri = do+        lookupEndpoint <&> \case+            (Just endpoint, _) -> rootPathIfEmpty endpoint+            (_, Just endpoint) -> appendExportPath @a endpoint+            _ -> appendExportPath @a HTTP.defaultEndpoint++    lookupGrpcUri :: Eff es URI+    lookupGrpcUri = do+        lookupEndpoint <&> \case+            (Just endpoint, _) -> endpoint+            (_, Just endpoint) -> endpoint+            _ -> GRPC.defaultEndpoint++    lookupWithFallback :: String -> String -> Eff es String+    lookupWithFallback suffix fallback =+        fromMaybe fallback <$> lookupSignalOrGlobal @a suffix++    rootPathIfEmpty :: URI -> URI+    rootPathIfEmpty uri+        | null (uriPath uri) = uri{uriPath = "/"}+        | otherwise = uri++    parseOrFail :: Maybe String -> Eff es (Maybe URI)+    parseOrFail Nothing = pure Nothing+    parseOrFail (Just s) = maybe (throwError $ InvalidOTLPEndpointURI s) (pure . pure) $ parseAbsoluteURI s++lookupExportConfig+    :: forall a es+     . ( Export.Request a+       , Environment :> es+       , Error ConfigError :> es+       )+    => Export.Config+    -> Eff es Export.Config+lookupExportConfig Export.Config{..} = do+    batch <-+        maybe (pure Nothing) (fmap Just . uncurry lookupBatchConfig) $+            (,) <$> batchEnvPrefix @a <*> batch+    exportTimeoutMs <- envNatural "EXPORT_TIMEOUT" exportTimeoutMs+    pure Export.Config{..}+  where+    envNatural :: String -> Natural -> Eff es Natural+    envNatural name fallback = lookupEnv variable >>= maybe (pure fallback) parse+      where+        variable = intercalate "_" ["OTEL", exportSignalEnvName @a, name]+        parse value = maybe (throwError $ ParseEnvConfigError OTLPConfig{..}) pure $ readMay value++    lookupBatchConfig :: String -> Export.BatchConfig -> Eff es Export.BatchConfig+    lookupBatchConfig batchEnvPrefix Export.BatchConfig{..} = do+        maxQueueSize <- envNatural "MAX_QUEUE_SIZE" maxQueueSize+        scheduledDelayMs <- envNatural "SCHEDULE_DELAY" scheduledDelayMs+        maxBatchSize <- envNatural "MAX_EXPORT_BATCH_SIZE" maxBatchSize+        pure Export.BatchConfig{..}+      where+        envNatural :: String -> Natural -> Eff es Natural+        envNatural name fallback = lookupEnv variable >>= maybe (pure fallback) parse+          where+            variable = intercalate "_" ["OTEL", batchEnvPrefix, name]+            parse value = maybe (throwError $ ParseEnvConfigError OTLPConfig{..}) pure $ readMay value
+ src/Effectful/OpenTelemetry/Protocol/Exception.hs view
@@ -0,0 +1,61 @@+module Effectful.OpenTelemetry.Protocol.Exception where++import Data.Typeable (cast)+import Effectful.Exception+import Effectful.GrpcClient (GrpcError)+import Effectful.Http2Client (ClientError)+import Effectful.HttpClient (HttpException)+import Prelude++-- | Any 'Exception' thrown in @otel-effectful@ is encapsulated in a 'SomeOTLPException'.+data SomeOTLPException = forall e. (Exception e) => SomeOTLPException e++fromOTLPException :: (Exception e) => SomeOTLPException -> Maybe e+fromOTLPException (SomeOTLPException e) = fromException . toException $ e++instance Show SomeOTLPException where+    show (SomeOTLPException e) = show e++instance Exception SomeOTLPException where+    displayException (SomeOTLPException e) = displayException e+    backtraceDesired (SomeOTLPException e) = backtraceDesired e++otlpExceptionToException :: (Exception e) => e -> SomeException+otlpExceptionToException = toException . SomeOTLPException++otlpExceptionFromException :: (Exception e) => SomeException -> Maybe e+otlpExceptionFromException ex = do+    SomeOTLPException e <- fromException ex+    cast e++data ConnectionTimeout = ConnectionTimeout+    deriving stock (Show)++instance Exception ConnectionTimeout where+    toException = otlpExceptionToException+    fromException = otlpExceptionFromException+    displayException ConnectionTimeout = "connection timed out"++newtype OTLPHttpException = OTLPHttpException HttpException+    deriving newtype (Show)++instance Exception OTLPHttpException where+    toException = otlpExceptionToException+    fromException = otlpExceptionFromException+    displayException (OTLPHttpException e) = displayException e++newtype OTLPClientError = OTLPClientError ClientError+    deriving newtype (Show)++instance Exception OTLPClientError where+    toException = otlpExceptionToException+    fromException = otlpExceptionFromException+    displayException (OTLPClientError e) = displayException e++newtype OTLPGrpcError = OTLPGrpcError GrpcError+    deriving newtype (Show)++instance Exception OTLPGrpcError where+    toException = otlpExceptionToException+    fromException = otlpExceptionFromException+    displayException (OTLPGrpcError e) = displayException e
+ src/Effectful/OpenTelemetry/Protocol/Export.hs view
@@ -0,0 +1,63 @@+{-# LANGUAGE AllowAmbiguousTypes #-}++module Effectful.OpenTelemetry.Protocol.Export {-# WARNING in "x-unstable-interface" "This is an unstable interface." #-} where++import Data.Aeson qualified as Aeson+import Data.List qualified as List+import Effectful.OpenTelemetry.Protocol.Resource qualified as Resource+import Network.GRPC.HTTP2.Proto3Wire (RPC)+import Network.URI (URI (..))+import Network.URI qualified as URI+import Numeric.Natural (Natural)+import Proto3.Wire.Encode (MessageBuilder)+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | How to export a batch of signal @a@ via OTLP: HTTP paths, the gRPC RPC descriptor,+-- JSON\/proto encoders, the @OTEL_EXPORTER_OTLP_\<SIGNAL\>_*@ variable name, and batching defaults.+-- Implemented once per signal type (e.g. 'Effectful.OpenTelemetry.Tracing.Span.Span').+class Request a where+    exportHttpPathComponents :: [String]++    exportGrpcRPC :: RPC++    exportJson :: [Resource.Items a] -> Aeson.Value++    exportProto :: [Resource.Items a] -> MessageBuilder+    default exportProto :: (Proto.Encode a) => [Resource.Items a] -> MessageBuilder+    exportProto = foldMap (Proto.encodeField 1)++    transportSignalEnvName :: String+    exportSignalEnvName :: String+    exportSignalEnvName = transportSignalEnvName @a++    batchEnvPrefix :: Maybe String++    appendExportPath :: URI -> URI+    appendExportPath uri = uri{uriPath = base <> encodedPath}+      where+        base = List.dropWhileEnd (== '/') (uriPath uri)+        encodedPath =+            concatMap (('/' :) . URI.escapeURIString URI.isUnescapedInURIComponent) $+                exportHttpPathComponents @a++    defaultConfig :: Config++-- | Batch-processor / metric-reader configuration.+--+-- See <https://opentelemetry.io/docs/specs/otel/trace/sdk/#batching-processor>.+data Config = Config+    { batch :: Maybe BatchConfig+    , exportTimeoutMs :: Natural+    }+    deriving stock (Show, Eq)++data BatchConfig = BatchConfig+    { maxQueueSize :: Natural+    -- ^ Maximum number of items held in the queue before new items are dropped.+    , scheduledDelayMs :: Natural+    -- ^ Delay between export attempts in milliseconds.+    , maxBatchSize :: Natural+    -- ^ Maximum number of items forwarded to the exporter per call.+    }+    deriving stock (Show, Eq)
+ src/Effectful/OpenTelemetry/Protocol/GRPC.hs view
@@ -0,0 +1,94 @@+{-# 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
+ src/Effectful/OpenTelemetry/Protocol/HTTP.hs view
@@ -0,0 +1,55 @@+module Effectful.OpenTelemetry.Protocol.HTTP {-# WARNING in "x-unstable-interface" "This is an unstable interface." #-} where++import Codec.Compression.GZip qualified as GZip+import Control.Monad (void)+import Data.ByteString.Lazy (LazyByteString)+import Data.Maybe (maybeToList)+import Effectful+import Effectful.Error.Static (Error, throwError)+import Effectful.Exception (handle)+import Effectful.HttpClient+    ( HttpClient+    , HttpException+    , RequestBody (..)+    , ResponseTimeout+    , httpNoBody+    , method+    , requestBody+    , requestFromURI_+    , requestHeaders+    , responseTimeout+    )+import Effectful.OpenTelemetry.Protocol.Transport (Compression (..), Encoding)+import Effectful.OpenTelemetry.Protocol.Transport qualified as Transport+import Network.URI (URI)+import Network.URI.Static (uri)+import Prelude++-- | Default OTLP HTTP endpoint URI: @http:\/\/localhost:4318@.+defaultEndpoint :: URI+defaultEndpoint = [uri|http://localhost:4318|]++sendPayload+    :: (HttpClient :> es, Error HttpException :> es)+    => URI+    -> Encoding+    -> Compression+    -> ResponseTimeout+    -> LazyByteString+    -> Eff es ()+sendPayload endpoint encoding compression responseTimeout body =+    void+        . handle @HttpException throwError+        . httpNoBody+        $ (requestFromURI_ endpoint)+            { method = "POST"+            , requestBody = RequestBodyLBS (compress body)+            , requestHeaders =+                Transport.contentType encoding+                    : maybeToList (Transport.contentEncoding compression)+            , responseTimeout+            }+  where+    compress = case compression of+        NoCompression -> id+        GZip -> GZip.compress
+ src/Effectful/OpenTelemetry/Protocol/Resource.hs view
@@ -0,0 +1,50 @@+module Effectful.OpenTelemetry.Protocol.Resource where++import Data.Aeson.Types (ToJSON)+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+import Effectful.OpenTelemetry.Protocol.Scope (Scope)+import Effectful.OpenTelemetry.Protocol.Scope qualified as Scope+import GHC.Generics (Generic)+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | Represents the observed entity for which telemetry is produced.+-- Within OpenTelemetry, all signals are associated with a 'Resource',+-- enabling contextual correlation of data from the same source.+--+-- See <https://opentelemetry.io/docs/specs/otel/resource/ the OpenTelemetry spec>.+newtype Resource = Resource+    { attributes :: Attributes+    }+    deriving stock (Generic)+    deriving newtype (Semigroup, Monoid, Show, Eq)+    deriving anyclass (ToJSON)++-- https://github.com/open-telemetry/opentelemetry-proto/blob/main/opentelemetry/proto/resource/v1/resource.proto+instance Proto.Encode Resource where+    encode Resource{..} =+        mconcat+            [ Proto.encodeField 1 attributes+            -- 2: dropped_attributes_count: Not supported+            -- 3: entity_refs: Not supported+            ]++-- | A batch of scoped items grouped by 'Resource'.+data Items a = Items+    { resource :: Resource+    , scopeItems :: [Scope.Items a]+    }+    deriving stock (Generic)++-- NOTE: The JSON encoding for Items depends on the wrapped item type,+-- so it is defined in the Export.Request class.++instance (Proto.Encode a) => Proto.Encode (Items a) where+    encode Items{..} =+        mconcat+            [ Proto.encodeField 1 resource+            , foldMap (Proto.encodeField 2) scopeItems+            ]++wrapItems :: Resource -> Scope -> [a] -> [Items a]+wrapItems resource scope items = [Items{resource, scopeItems = [Scope.Items{scope, items}]}]
+ src/Effectful/OpenTelemetry/Protocol/Scope.hs view
@@ -0,0 +1,47 @@+module Effectful.OpenTelemetry.Protocol.Scope where++import Data.Aeson.Types (ToJSON)+import Data.Text (Text)+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+import GHC.Generics (Generic)+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | A logical unit of software with which emitted telemetry is associated.+-- It can represent a module, package, class, library, or framework;+-- any meaningful boundary that distinguishes one source of telemetry from another.+--+-- See <https://opentelemetry.io/docs/concepts/instrumentation-scope/ the OpenTelemetry spec>.+data Scope = Scope+    { name :: Text+    , version :: Text+    , attributes :: Attributes+    }+    deriving stock (Generic, Eq, Show)+    deriving anyclass (ToJSON)++instance Proto.Encode Scope where+    encode Scope{..} =+        mconcat+            [ Proto.encodeField 1 name+            , Proto.encodeField 2 version+            , Proto.encodeField 3 attributes+            -- 4: dropped_attributes_count: Not supported+            ]++-- | A batch of items grouped by 'Scope'.+data Items a = Items+    { scope :: Scope+    , items :: [a]+    }+    deriving stock (Generic)++-- NOTE: The JSON encoding for Items depends on the wrapped item type,+-- so it is defined in the Export.Request class.++instance (Proto.Encode a) => Proto.Encode (Items a) where+    encode Items{..} =+        mconcat+            [ Proto.encodeField 1 scope+            , foldMap (Proto.encodeField 2) items+            ]
+ src/Effectful/OpenTelemetry/Protocol/Transport.hs view
@@ -0,0 +1,42 @@+module Effectful.OpenTelemetry.Protocol.Transport {-# WARNING in "x-unstable-interface" "This is an unstable interface." #-} where++import Data.Char (toLower)+import Effectful.Http2Client (HostName, PortNumber)+import Network.GRPC.HTTP2.Proto3Wire (RPC)+import Network.HTTP.Types.Header (Header, hContentEncoding, hContentType)+import Network.URI (URI)+import Numeric.Natural+import Prelude++data Compression = NoCompression | GZip+    deriving stock (Show, Eq, Bounded, Enum)++instance Read Compression where+    readsPrec _ s =+        case toLower <$> s of+            "" -> [(NoCompression, "")]+            "none" -> [(NoCompression, "")]+            "nocompression" -> [(NoCompression, "")]+            "gzip" -> [(GZip, "")]+            _ -> []++contentEncoding :: Compression -> Maybe Header+contentEncoding NoCompression = Nothing+contentEncoding GZip = Just (hContentEncoding, "gzip")++data Encoding = Json | Proto+    deriving stock (Show, Eq)++contentType :: Encoding -> Header+contentType Json = (hContentType, "application/json")+contentType Proto = (hContentType, "application/x-protobuf")++data Protocol+    = HTTP Encoding URI+    | GRPC HostName PortNumber RPC++data Transport = Transport+    { protocol :: Protocol+    , compression :: Compression+    , timeoutMs :: Natural+    }
+ src/Effectful/OpenTelemetry/Timestamp.hs view
@@ -0,0 +1,42 @@+module Effectful.OpenTelemetry.Timestamp where++import Control.Monad (MonadPlus (mzero))+import Data.Aeson.Types (FromJSON (..), ToJSON (..), Value (..), withText)+import Data.Text qualified as Text+import Data.Time.Clock.POSIX (posixSecondsToUTCTime)+import Data.Time.Clock.System (SystemTime (..), getSystemTime)+import Data.Time.Format.ISO8601 (iso8601Show)+import Data.Word (Word64)+import GHC.Generics (Generic)+import Prettyprinter (Pretty (..))+import Prettyprinter.Extra (PrettyAnn (..))+import Proto3.Wire.Encode.Class qualified as Proto+import Text.Read (readMaybe)+import Prelude++-- | UNIX timestamp with nanosecond precision.+newtype Timestamp = Timestamp {nanos :: Word64}+    deriving stock (Generic)+    deriving newtype (Show, Eq, Ord)++epoch :: Timestamp+epoch = Timestamp 0++now :: IO Timestamp+now = do+    MkSystemTime{systemSeconds, systemNanoseconds} <- getSystemTime+    pure . Timestamp $+        fromIntegral systemSeconds * 1_000_000_000 + fromIntegral systemNanoseconds++instance FromJSON Timestamp where+    parseJSON = withText "Timestamp" $ maybe mzero (pure . Timestamp) . readMaybe . Text.unpack++instance ToJSON Timestamp where+    toJSON = String . Text.show++instance {-# OVERLAPPING #-} Proto.EncodeField Timestamp where+    encodeField n = Proto.fixed64 n . nanos++-- | yyyy-mm-ddThh:mm:ss[.sss]Z (ISO 8601:2004(E) sec. 4.3.2 extended format)+instance PrettyAnn ann Timestamp where+    prettyAnn Timestamp{..} = pretty . iso8601Show . posixSecondsToUTCTime $ fromIntegral nanos / 1_000_000_000
+ src/Effectful/OpenTelemetry/Tracing.hs view
@@ -0,0 +1,10 @@+module Effectful.OpenTelemetry.Tracing+    ( module Effectful.OpenTelemetry.Tracing.Effect+    , module Effectful.OpenTelemetry.Tracing.Span+    , module Effectful.OpenTelemetry.Tracing.Propagator+    )+where++import Effectful.OpenTelemetry.Tracing.Effect+import Effectful.OpenTelemetry.Tracing.Propagator+import Effectful.OpenTelemetry.Tracing.Span
+ src/Effectful/OpenTelemetry/Tracing/Effect.hs view
@@ -0,0 +1,386 @@+{-# LANGUAGE AllowAmbiguousTypes #-}+{-# LANGUAGE OverloadedLists #-}++module Effectful.OpenTelemetry.Tracing.Effect+    ( -- * Effect+      Tracing+    , inSpan+    , inOkSpan+    , inErrorSpan+    , withContext+    , addEvent+    , addEventNow+    , recordException+    , recordExceptionAt+    , currentContext+    , currentContextIO++      -- * Runners+    , runTracing+    , runHttpTracing+    , runGrpcTracing+    , runConsoleTracing+    , runInMemoryTracing+    , runTracingWith+    , runNoTracing+    )+where++import Control.Monad.Extra (ifM)+import Data.Aeson (ToJSON (..))+import Data.Functor (void, (<&>))+import Data.Sequence (Seq, (|>))+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Typeable (typeOf)+import Effectful+import Effectful.Concurrent (Concurrent)+import Effectful.Dispatch.Static+import Effectful.Environment (Environment)+import Effectful.Error.Static (Error, throwError, tryError)+import Effectful.Exception+    ( Exception+    , SomeException+    , displayException+    , mask+    , throwIO+    , try+    , uninterruptibleMask_+    )+import Effectful.Http2Client (HostName, PortNumber)+import Effectful.HttpClient (responseTimeoutDefault)+import Effectful.OpenTelemetry.Exporter (Exporter)+import Effectful.OpenTelemetry.Exporter qualified as Exporter+import Effectful.OpenTelemetry.Exporter.Console qualified as Console+import Effectful.OpenTelemetry.Exporter.Environment qualified as Exporter.Environment+import Effectful.OpenTelemetry.Exporter.OTLP (defaultGrpcSendTimeout)+import Effectful.OpenTelemetry.Protocol+    ( Attributes+    , Compression+    , Encoding+    , OTLP+    , Resource+    , Scope+    , exportIO+    , runInMemoryOTLP+    , runNoOTLP+    , runOTLPWith+    )+import Effectful.OpenTelemetry.Protocol.Environment qualified as Environment+import Effectful.OpenTelemetry.Protocol.Export qualified as Export+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Timestamp qualified as Timestamp+import Effectful.OpenTelemetry.Tracing.Span (Span (..))+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span (Context)+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span.Context+import Effectful.OpenTelemetry.Tracing.Span.Event (Event (..))+import Effectful.OpenTelemetry.Tracing.Span.Event qualified as Event+import Effectful.OpenTelemetry.Tracing.Span.Kind qualified as Span+import Effectful.OpenTelemetry.Tracing.Span.Status qualified as Span.Status+import Effectful.Retry (Retry)+import Effectful.Timeout (Timeout)+import GHC.Stack (whoCreated)+import Network.URI (URI)+import Prelude++data Tracing :: Effect++type instance DispatchOf Tracing = 'Static 'WithSideEffects++data instance StaticRep Tracing = Tracing+    { currentSpan :: Maybe Span.Context+    , events :: Seq Event+    , export :: Span -> IO ()+    }++data ExitCase e+    = ExitSuccess+    | ExitFailure e++-- | Execute an action with a given 'Span.Context'.+withContext :: (Tracing :> es) => Span.Context -> Eff es a -> Eff es a+withContext context = localStaticRep \tracing ->+    tracing{currentSpan = Just context, events = mempty}++-- | Execute an action within a named 'Span'.+-- If the action throws an exception, the 'Span' is finalised with 'Span.Status.Error'+-- and the exception is re-raised.+-- Otherwise, the 'Span' is exported with 'Span.Status.Unset'.+inSpan+    :: forall es a+     . (HasCallStack, Tracing :> es)+    => Text+    -> Span.Kind+    -> Attributes+    -> Eff es a+    -> Eff es a+inSpan name kind attributes action = do+    tracing <- getStaticRep @Tracing+    context <- unsafeEff_ $ Span.Context.new tracing.currentSpan+    startTime <- unsafeEff_ Timestamp.now+    withContext context $ finallyCaseIO @SomeException action \exitCase -> do+        endTime <- unsafeEff_ Timestamp.now+        tracing' <- getStaticRep @Tracing+        unsafeEff_ . tracing.export $+            Span+                { name+                , parentSpanId = tracing.currentSpan <&> (.spanId)+                , context+                , kind+                , startTime+                , endTime = Just endTime+                , attributes+                , events = tracing'.events+                , status =+                    case exitCase of+                        ExitSuccess -> Span.Status.unset+                        (ExitFailure e) -> Span.Status.fromException e+                }+  where+    finallyCaseIO+        :: forall e+         . (Exception e)+        => Eff es a+        -> (ExitCase e -> Eff es ())+        -> Eff es a+    finallyCaseIO thing after = mask \restore -> do+        try @e (restore thing) >>= \case+            Left e -> do+                recordException e mempty+                void . try @e . uninterruptibleMask_ . after $ ExitFailure e+                throwIO e+            Right y -> do+                void . try @e . uninterruptibleMask_ $ after ExitSuccess+                pure y++-- | Execute an action that may not fail within a named 'Span'.+-- Explicity marks the 'Span' with 'Span.Status.Ok'.+-- Does not export the 'Span' if the action raises an exception.+-- Use 'inErrorSpan' or 'inSpan' for actions that may fail.+inOkSpan+    :: forall a es+     . (Tracing :> es)+    => Text+    -> Span.Kind+    -> Attributes+    -> Eff es a+    -> Eff es a+inOkSpan name kind attributes action = do+    tracing <- getStaticRep @Tracing+    context <- unsafeEff_ $ Span.Context.new tracing.currentSpan+    startTime <- unsafeEff_ Timestamp.now+    (a, tracing') <- withContext context $ (,) <$> action <*> getStaticRep @Tracing+    endTime <- unsafeEff_ Timestamp.now+    unsafeEff_ . tracing.export $+        Span+            { name+            , parentSpanId = tracing.currentSpan <&> (.spanId)+            , context+            , kind+            , startTime+            , endTime = Just endTime+            , attributes+            , events = tracing'.events+            , status = Span.Status.ok+            }+    pure a++-- | Execute an action that may fail with the expected @Error e@ within a named 'Span'.+-- If the action fails with this expected @Error e@, the 'Span' is finalised with 'Span.Status.Error'+-- and the exception is re-raised.+-- Otherwise, exports the 'Span' with 'Span.Status.Unset'.+-- Does not export the 'Span' if the action fails with an unexpected exception (e.g. via 'throwIO').+-- If you need to export 'Span's in the presence of any exception, use 'inSpan' instead.+inErrorSpan+    :: forall e es a+     . (Exception e, Error e :> es, Tracing :> es)+    => Text+    -> Span.Kind+    -> Attributes+    -> Eff es a+    -> Eff es a+inErrorSpan name kind attributes action = do+    tracing <- getStaticRep @Tracing+    context <- unsafeEff_ $ Span.Context.new tracing.currentSpan+    startTime <- unsafeEff_ Timestamp.now+    withContext context . finallyCase action $ \exitCase -> do+        endTime <- unsafeEff_ Timestamp.now+        tracing' <- getStaticRep @Tracing+        unsafeEff_ . tracing.export $+            Span+                { name+                , parentSpanId = tracing.currentSpan <&> (.spanId)+                , context+                , kind+                , startTime+                , endTime = Just endTime+                , attributes+                , events = tracing'.events+                , status =+                    case exitCase of+                        ExitSuccess -> Span.Status.unset+                        (ExitFailure e) -> Span.Status.fromException e+                }+  where+    finallyCase :: Eff es a -> (ExitCase e -> Eff es ()) -> Eff es a+    finallyCase thing after = do+        tryError @e thing >>= \case+            Left (_, e) -> do+                recordException e mempty+                void . tryError @e . uninterruptibleMask_ . after $ ExitFailure e+                throwError e+            Right y -> do+                void . tryError @e . uninterruptibleMask_ . after $ ExitSuccess+                pure y++-- | Add an event to the current span. This is a no-op if not called in a span.+--+-- https://github.com/open-telemetry/opentelemetry-specification/blob/v1.57.0/specification/trace/api.md#add-events+addEvent :: (Tracing :> es) => Event -> Eff es ()+addEvent event = do+    tracing <- getStaticRep+    case tracing.currentSpan of+        Nothing -> pure ()+        Just context ->+            putStaticRep $ tracing{currentSpan = Just context, events = tracing.events |> event}++-- | Add an event that occurred at the current time to the current span. This is a no-op if not called in a span.+--+-- https://github.com/open-telemetry/opentelemetry-specification/blob/v1.57.0/specification/trace/api.md#add-events+addEventNow :: (Tracing :> es) => Text -> Attributes -> Eff es ()+addEventNow name attributes = addEvent =<< unsafeEff_ (Event.now name attributes)++-- | Add an exception to the current span. This is a no-op if not called in a span.+--+-- - https://github.com/open-telemetry/opentelemetry-specification/blob/v1.57.0/specification/trace/api.md#record-exception+-- - https://github.com/open-telemetry/opentelemetry-specification/blob/v1.57.0/specification/trace/exceptions.md+recordExceptionAt :: (Exception e, Tracing :> es) => Timestamp -> e -> Attributes -> Eff es ()+recordExceptionAt time e attributes = do+    stackTrace <- unsafeEff_ $ whoCreated e+    addEvent+        Event+            { time+            , name = "exception"+            , attributes =+                [ ("exception.type", toJSON . Text.pack . show $ typeOf e)+                , ("exception.message", toJSON . Text.pack $ displayException e)+                , ("exception.stacktrace", toJSON . Text.unlines $ Text.pack <$> stackTrace)+                ]+                    <> attributes+            }++-- | Add an exception event that occurred at the current time to the current span. This is a no-op if not called in a span.+--+-- - https://github.com/open-telemetry/opentelemetry-specification/blob/v1.57.0/specification/trace/api.md#record-exception+-- - https://github.com/open-telemetry/opentelemetry-specification/blob/v1.57.0/specification/trace/exceptions.md+recordException :: (Exception e, Tracing :> es) => e -> Attributes -> Eff es ()+recordException e attributes = do+    time <- unsafeEff_ Timestamp.now+    recordExceptionAt time e attributes++currentContext :: (Tracing :> es) => Eff es (Maybe Span.Context)+currentContext = do+    Tracing{..} <- getStaticRep+    pure currentSpan++currentContextIO :: (Tracing :> es) => Eff es (IO (Maybe Span.Context))+currentContextIO = unsafeEff $ pure . unEff currentContext++runTracingState+    :: (OTLP Span :> es, IOE :> es)+    => Eff (Tracing ': es) a+    -> Eff es a+runTracingState eff = do+    export <- exportIO+    evalStaticRep Tracing{currentSpan = Nothing, events = mempty, ..} eff++-- | Run the 'Tracing' effect, sending telemetry to an exporter.+-- Reads the configuration from <https://opentelemetry.io/docs/specs/otel/configuration/sdk-environment-variables/#general-sdk-configuration the standard environment variables>.+-- Delegates to 'runNoTracing' if @OTEL_SDK_DISABLED = true@.+runTracing+    :: ( IOE :> es+       , Concurrent :> es+       , Environment :> es+       , Retry :> es+       , Timeout :> es+       )+    => Resource+    -> Scope+    -> Eff (Tracing ': es) a+    -> Eff es a+runTracing resource scope eff =+    ifM Environment.isSdkDisabled (runNoTracing eff) do+        exporter <- Environment.runConfigError $ Exporter.Environment.lookup @Span resource scope+        runTracingWith exporter eff++-- | Run the 'Tracing' effect, sending telemetry to a collector at the given 'URI'+-- over HTTP with the given 'Encoding'.+runHttpTracing+    :: (IOE :> es, Concurrent :> es, Retry :> es, Timeout :> es)+    => Resource+    -> Scope+    -> Encoding+    -> Compression+    -> URI+    -> Eff (Tracing ': es) a+    -> Eff es a+runHttpTracing resource scope encoding compression endpoint =+    runTracingWith $+        Exporter.http+            @Span+            resource+            scope+            encoding+            endpoint+            (Export.defaultConfig @Span)+            compression+            responseTimeoutDefault++-- | Run the 'Tracing' effect, sending telemetry to a gRPC collector.+runGrpcTracing+    :: (IOE :> es, Concurrent :> es, Retry :> es, Timeout :> es)+    => Resource+    -> Scope+    -> HostName+    -> PortNumber+    -> Compression+    -> Eff (Tracing ': es) a+    -> Eff es a+runGrpcTracing resource scope host port compression =+    runTracingWith $+        Exporter.grpc @Span+            resource+            scope+            host+            port+            (Export.exportGrpcRPC @Span)+            (Export.defaultConfig @Span)+            compression+            defaultGrpcSendTimeout++-- | Run the 'Tracing' effect, printing telemetry to the console rather than sending to a collector.+runConsoleTracing :: (IOE :> es) => Eff (Tracing ': es) a -> Eff es a+runConsoleTracing = runTracingWith Console.stdout++-- | Run the 'Tracing' effect, collecting telemetry in-memory rather than sending to a collector.+runInMemoryTracing+    :: (IOE :> es, Concurrent :> es)+    => Eff (Tracing ': es) a+    -> Eff es (a, [Span])+runInMemoryTracing = runInMemoryOTLP @Span . runTracingState . inject++-- | Run the 'Tracing' effect with a custom 'Exporter'.+-- Passes 'Span's to the sink synchronously as they are finalised, without batching or retrying.+runTracingWith+    :: (IOE :> es)+    => Exporter es Span+    -> Eff (Tracing ': es) a+    -> Eff es a+runTracingWith exporter = runOTLPWith exporter . runTracingState . inject++-- | Run the 'Tracing' effect as a no-op action.+runNoTracing+    :: (IOE :> es)+    => Eff (Tracing ': es) a+    -> Eff es a+runNoTracing = runNoOTLP @Span . runTracingState . inject
+ src/Effectful/OpenTelemetry/Tracing/Propagator.hs view
@@ -0,0 +1,78 @@+module Effectful.OpenTelemetry.Tracing.Propagator+    ( -- * Propagator+      Propagator (..)++      -- * HTTP header propagators+    , w3cTraceContext+    )+where++import Control.Applicative ((<|>))+import Data.ByteString (ByteString)+import Data.ByteString qualified as ByteString+import Data.List qualified as List+import Data.Text (Text)+import Effectful+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span (Context)+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span.Context+import Effectful.OpenTelemetry.Tracing.Trace.State qualified as Trace.State+import Network.HTTP.Types (Header, HeaderName, RequestHeaders, ResponseHeaders)+import Prelude++-- | Propagates context information across process boundaries and+-- defines the restrictions imposed by a specific transport.+--+-- [@es@] An arbitrary effect stack.+-- [@context@] A cross-cutting concern context, for example a 'Span.Context'.+-- [@inboundCarrier@] Medium used by the propagator to read values from.+-- [@outboundCarrier@] Medium used by the propagator to write values to.+--+-- See <https://opentelemetry.io/docs/specs/otel/context/api-propagators/ the OpenTelemetry spec>.+data Propagator es context inboundCarrier outboundCarrier = Propagator+    { propagatorNames :: [Text]+    -- ^ The header (or field) names this propagator is responsible for.+    , extract :: inboundCarrier -> context -> Eff es context+    -- ^ Read propagation fields out of the inbound carrier, returning a+    -- context derived from the one passed in. If nothing valid can be parsed,+    -- the original context is returned unchanged.+    , inject :: context -> outboundCarrier -> Eff es outboundCarrier+    -- ^ Write the context's propagation fields into the outbound carrier.+    }++-- | The <https://www.w3.org/TR/trace-context/ W3C Trace Context> propagator,+-- carrying the @traceparent@ HTTP header.+--+-- See <https://www.w3.org/TR/trace-context/#traceparent-header the traceparent header specification>.+w3cTraceContext :: Propagator es (Maybe Span.Context) RequestHeaders ResponseHeaders+w3cTraceContext =+    Propagator+        { propagatorNames = ["traceparent", "tracestate"]+        , extract = \headers context -> pure . (<|> context) $ do+            context' <- Span.Context.fromTraceparent =<< List.lookup hTraceparent headers+            pure+                context'+                    { Span.Context.traceState =+                        foldMap (Trace.State.fromByteString . snd)+                            . filter ((hTracestate ==) . fst)+                            $ headers+                    }+        , inject = \context headers ->+            pure $ case context of+                Nothing -> headers+                Just ctx ->+                    addHeader hTracestate (Trace.State.toByteString ctx.traceState)+                        . setHeader hTraceparent (Span.Context.toTraceparent ctx)+                        $ headers+        }+  where+    setHeader :: HeaderName -> ByteString -> [Header] -> [Header]+    setHeader name value = addHeader name value . filter ((name /=) . fst)++    addHeader :: HeaderName -> ByteString -> [Header] -> [Header]+    addHeader name value+        | ByteString.null value = id+        | otherwise = ((name, value) :)++hTraceparent, hTracestate :: HeaderName+hTraceparent = "traceparent"+hTracestate = "tracestate"
+ src/Effectful/OpenTelemetry/Tracing/Span.hs view
@@ -0,0 +1,132 @@+module Effectful.OpenTelemetry.Tracing.Span where++import Data.Aeson.Types (Pair, ToJSON (..), Value, object, (.=))+import Data.Function ((&))+import Data.Maybe (catMaybes)+import Data.Sequence (Seq)+import Data.Text (Text)+import Data.Text qualified as Text+import Effectful.OpenTelemetry.Protocol.Attributes (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.Timestamp (Timestamp (..))+import Effectful.OpenTelemetry.Tracing.Span.Context (Context (..))+import Effectful.OpenTelemetry.Tracing.Span.Event (Event)+import Effectful.OpenTelemetry.Tracing.Span.ID qualified as Span (ID)+import Effectful.OpenTelemetry.Tracing.Span.Kind (Kind (..))+import Effectful.OpenTelemetry.Tracing.Span.Status+import GHC.Generics (Generic)+import Network.GRPC.HTTP2.Proto3Wire (RPC (..))+import Prettyprinter (Pretty (..))+import Prettyprinter.Extra (PrettyAnn (..))+import Prettyprinter.Extra qualified as Pretty+import Prettyprinter.Render.Terminal (AnsiStyle)+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | A single unit of work or operation in a trace. Spans are the building blocks of traces.+--+-- See <https://opentelemetry.io/docs/concepts/signals/traces/#spans the OpenTelemetry spec>.+data Span = Span+    { name :: Text+    , parentSpanId :: Maybe Span.ID+    , context :: Context+    , kind :: Kind+    , startTime :: Timestamp+    , endTime :: Maybe Timestamp+    , attributes :: Attributes+    , events :: Seq Event+    , status :: Status+    }+    deriving stock (Generic, Eq, Show)++instance ToJSON Span where+    toJSON Span{context = Context{..}, ..} =+        object $+            [ "traceId" .= traceId+            , "spanId" .= spanId+            , "traceState" .= traceState+            , "flags" .= traceFlags+            , "name" .= name+            , "kind" .= kind+            , "startTimeUnixNano" .= startTime+            , "attributes" .= attributes+            , "events" .= events+            , "status" .= status+            ]+                <> catMaybes @Pair+                    [ ("parentSpanId" .=) <$> parentSpanId+                    , ("endTimeUnixNano" .=) <$> endTime+                    ]++instance Proto.Encode Span where+    encode Span{context = Context{..}, ..} =+        mconcat+            [ Proto.encodeField 1 traceId+            , Proto.encodeField 2 spanId+            , Proto.encodeField 3 traceState+            , foldMap (Proto.encodeField 4) parentSpanId+            , Proto.encodeField 5 name+            , Proto.encodeField 6 kind+            , Proto.encodeField 7 startTime+            , foldMap (Proto.encodeField 8) endTime+            , Proto.encodeField 9 attributes+            , -- 10: dropped_attributes_count: Not supported+              Proto.encodeField 11 events+            , -- 12: dropped_events_count: Not supported+              -- TODO: 13: links+              -- 14: dropped_links_count: Not supported+              Proto.encodeField 15 status+            , Proto.encodeField 16 traceFlags+            ]++instance PrettyAnn AnsiStyle Span where+    prettyAnn Span{..} =+        Pretty.unwords . filter (not . Pretty.null) $+            [ prettyAnn startTime+            , endTime & maybe mempty \(nanos -> subtract startTime.nanos -> duration) ->+                pretty . (<> "ms") . Text.show $ duration `div` 1_000_000+            , pretty name+            , prettyAnn attributes+            , prettyAnn status+            ]++instance Export.Request Span where+    exportHttpPathComponents = ["v1", "traces"]+    exportGrpcRPC =+        RPC+            { pkg = "opentelemetry.proto.collector.trace.v1"+            , srv = "TraceService"+            , meth = "Export"+            }+    exportJson = object . pure . ("resourceSpans" .=) . fmap resourceSpans+      where+        resourceSpans :: Resource.Items Span -> Value+        resourceSpans Resource.Items{resource, scopeItems} =+            object+                [ "resource" .= resource+                , "scopeSpans" .= fmap scopeSpans scopeItems+                ]+        scopeSpans :: Scope.Items Span -> Value+        scopeSpans Scope.Items{..} =+            object+                [ "scope" .= scope+                , "spans" .= items+                ]++    transportSignalEnvName = "TRACES"+    batchEnvPrefix = Just "BSP"++    -- https://opentelemetry.io/docs/specs/otel/trace/sdk/#batching-processor+    defaultConfig =+        Export.Config+            { batch =+                Just+                    Export.BatchConfig+                        { maxQueueSize = 2_048+                        , scheduledDelayMs = 5_000+                        , maxBatchSize = 512+                        }+            , exportTimeoutMs = 30_000+            }
+ src/Effectful/OpenTelemetry/Tracing/Span/Context.hs view
@@ -0,0 +1,89 @@+module Effectful.OpenTelemetry.Tracing.Span.Context where++import Control.Monad ((>=>))+import Data.Aeson qualified as Aeson+import Data.ByteString (ByteString)+import Data.ByteString.Builder qualified as Builder+import Data.ByteString.Builder.Extra qualified as Builder+import Data.ByteString.Lazy qualified as LazyByteString+import Data.Either.Extra (eitherToMaybe)+import Data.List qualified as List+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text+import Effectful.OpenTelemetry.Tracing.Span.ID qualified as Span (ID)+import Effectful.OpenTelemetry.Tracing.Span.ID qualified as Span.ID+import Effectful.OpenTelemetry.Tracing.Trace.Flags qualified as Trace (Flags)+import Effectful.OpenTelemetry.Tracing.Trace.Flags qualified as Trace.Flags+import Effectful.OpenTelemetry.Tracing.Trace.ID qualified as Trace (ID)+import Effectful.OpenTelemetry.Tracing.Trace.ID qualified as Trace.ID+import Effectful.OpenTelemetry.Tracing.Trace.State qualified as Trace (State)+import GHC.Generics (Generic)+import Prelude++-- | A propagation mechanism that carries execution-scoped values across API boundaries.+--+-- See the <https://opentelemetry.io/docs/concepts/signals/traces/#span-context OpenTelemetry spec>.+data Context = Context+    { traceId :: Trace.ID+    , spanId :: Span.ID+    , traceFlags :: Trace.Flags+    , traceState :: Trace.State+    }+    deriving stock (Generic, Eq, Show)++-- | Create a new 'Context' with an optional parent and a new span 'Span.ID'.+-- If a parent is given, the new 'Context' inherits the parent's trace 'Trace.ID', 'Trace.Flags',+-- and 'Trace.State'.+-- Otherwise it gets a fresh trace 'Trace.ID', the 'Trace.Flags.sampled' and 'Trace.Flags.random'+-- trace 'Trace.Flags', and empty trace 'Trace.State'.+new :: Maybe Context -> IO Context+new parent = do+    traceId <- maybe Trace.ID.new (pure . traceId) parent+    spanId <- Span.ID.new+    pure+        Context+            { traceFlags =+                let flags =+                        maybe+                            mempty+                                { Trace.Flags.sampled = True+                                , Trace.Flags.random = True+                                }+                            traceFlags+                            parent+                 in flags{Trace.Flags.remote = Trace.Flags.IsNotRemote}+            , traceState = maybe mempty traceState parent+            , ..+            }++-- | Decode a 'Context' with an empty 'Trace.State' from a @traceparent@ header.+fromTraceparent :: ByteString -> Maybe Context+fromTraceparent =+    eitherToMaybe . Text.decodeUtf8' >=> \t ->+        case Aeson.String <$> Text.splitOn "-" t of+            [ "00"+                , Aeson.fromJSON -> Aeson.Success traceId+                , Aeson.fromJSON -> Aeson.Success spanId+                , Aeson.fromJSON -> Aeson.Success traceFlags+                ] ->+                    Just+                        Context+                            { traceState = mempty+                            , traceFlags = traceFlags{Trace.Flags.remote = Trace.Flags.IsRemote}+                            , ..+                            }+            _ -> Nothing++-- | Encode a 'Context' as a @traceparent@ header.+-- Drops the 'traceState'.+toTraceparent :: Context -> ByteString+toTraceparent Context{..} =+    LazyByteString.toStrict+        . Builder.toLazyByteStringWith (Builder.untrimmedStrategy 64 64) mempty+        . mconcat+        . List.intersperse "-"+        $ [ "00"+          , Builder.byteStringHex $ Trace.ID.toBytes traceId+          , Builder.byteStringHex $ Span.ID.toBytes spanId+          , Builder.word8HexFixed $ Trace.Flags.toWord8 traceFlags+          ]
+ src/Effectful/OpenTelemetry/Tracing/Span/Event.hs view
@@ -0,0 +1,42 @@+module Effectful.OpenTelemetry.Tracing.Span.Event where++import Data.Aeson.Types (ToJSON (..), object, (.=))+import Data.Text (Text)+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+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++-- | A structured log message or annotation on a @Span@,+-- typically used to denote a meaningful, singular point in time during the @Span@'s duration.+--+-- See <https://opentelemetry.io/docs/concepts/signals/traces/#span-events the OpenTelemetry spec>.+data Event = Event+    { time :: Timestamp+    , name :: Text+    , attributes :: Attributes+    }+    deriving stock (Generic, Eq, Show)++now :: Text -> Attributes -> IO Event+now name attributes = do+    time <- Timestamp.now+    pure Event{..}++instance ToJSON Event where+    toJSON Event{..} =+        object+            [ "timeUnixNano" .= time+            , "name" .= name+            , "attributes" .= attributes+            ]++instance Proto.Encode Event where+    encode Event{..} =+        mconcat+            [ Proto.encodeField 1 time+            , Proto.encodeField 2 name+            , Proto.encodeField 3 attributes+            ]
+ src/Effectful/OpenTelemetry/Tracing/Span/ID.hs view
@@ -0,0 +1,61 @@+module Effectful.OpenTelemetry.Tracing.Span.ID where++import Data.Aeson.Types (FromJSON (..), ToJSON (..), withText)+import Data.ByteString (ByteString)+import Data.ByteString qualified as ByteString+import Data.ByteString.Base16 qualified as Base16+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text+import System.Random.Extra (uniformIO)+import GHC.Generics (Generic)+import Proto3.Wire.Encode.Class qualified as Proto+import System.Random (Random (..), Uniform)+import System.Random.Stateful (Uniform (..), uniformByteStringM)+import Prelude hiding (length)++-- | A globally unique identifier of a span.+newtype ID = ID ByteString+    deriving stock (Generic)+    deriving newtype (Eq, Show)++length :: Int+length = 8++-- | Generate a random 'ID'.+new :: IO ID+new = uniformIO++-- | Encode an 'ID' to base-16.+-- The result is always 16 characters long.+toHex :: ID -> Text+toHex (ID bs) = Text.justifyRight (length * 2) '0' . Text.decodeUtf8 . Base16.encode $ bs++toBytes :: ID -> ByteString+toBytes (ID bs) = bs++fromBytes :: ByteString -> Either String ID+fromBytes bs+    | ByteString.length bs /= length = Left $ "Span ID must have exactly " <> show length <> " bytes"+    | ByteString.all (== 0) bs = Left "Span ID must not be all zero"+    | otherwise = Right $ ID bs++instance Uniform ID where+    uniformM g = do+        bs <- uniformByteStringM length g+        either (const $ uniformM g) pure $ fromBytes bs++instance Random ID where+    randomR _ = random++instance FromJSON ID where+    parseJSON = withText "ID" $ either fail pure . fromHex+      where+        fromHex :: Text -> Either String ID+        fromHex t = Base16.decode (Text.encodeUtf8 t) >>= fromBytes++instance ToJSON ID where+    toJSON = toJSON . toHex++instance {-# OVERLAPPING #-} Proto.EncodeField ID where+    encodeField n = Proto.byteString n . toBytes
+ src/Effectful/OpenTelemetry/Tracing/Span/Kind.hs view
@@ -0,0 +1,59 @@+module Effectful.OpenTelemetry.Tracing.Span.Kind where++import Data.Aeson.Types (ToJSON (..), Value (..))+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | Clarifies the relationship between spans that are correlated via parent/child relationships+-- or span links.+-- Describes two independent properties that benefit tracing systems during analysis:+--+-- - Whether a span represents an outgoing call to a remote service ('Client', 'Producer')+--   or a processing of an incoming request initiated externally ('Server', 'Consumer').+-- - Whether a span represents a request/response operation ('Client', 'Server')+--   or a deferred execution ('Producer', 'Consumer').+--+-- See <https://opentelemetry.io/docs/concepts/signals/traces/#span-kind the OpenTelemetry spec>+data Kind+    = -- | Indicates that a span represents an internal operation within an application,+      -- as opposed to an operation with remote parents or children.+      Internal+    | -- | Indicates that a span covers server-side handling of a remote request+      -- while the client awaits a response.+      Server+    | -- | Indicates that the span describes a request to a remote service+      -- where the client awaits a response.+      -- When the context of a 'Client' span is propagated, 'Client' span usually becomes+      -- a parent of a remote 'Server' span.+      Client+    | -- | Indicates that the span describes the initiation or scheduling of a local+      -- or remote operation. This initiating span often ends before the correlated 'Consumer'+      -- span, possibly even before the 'Consumer' span starts.+      --+      -- In messaging scenarios with batching, tracing individual messages requires a new+      -- 'Producer' span per message to be created.+      Producer+    | -- | Indicates that the span represents the processing of an operation initiated+      -- by a producer, where the producer does not wait for the outcome.+      Consumer+    deriving stock (Show, Eq, Bounded)++instance Enum Kind where+    fromEnum Internal = 1+    fromEnum Server = 2+    fromEnum Client = 3+    fromEnum Producer = 4+    fromEnum Consumer = 5++    toEnum 1 = Internal+    toEnum 2 = Server+    toEnum 3 = Client+    toEnum 4 = Producer+    toEnum 5 = Consumer+    toEnum _ = error "Enum.Kind.toEnum: bad argument"++instance ToJSON Kind where+    toJSON = Number . fromIntegral . fromEnum++instance {-# OVERLAPPING #-} Proto.EncodeField Kind where+    encodeField n = Proto.int32 n . fromIntegral . fromEnum
+ src/Effectful/OpenTelemetry/Tracing/Span/Status.hs view
@@ -0,0 +1,81 @@+module Effectful.OpenTelemetry.Tracing.Span.Status where++import Data.Aeson.Types (ToJSON (..), Value (..))+import Data.Text (Text)+import Data.Text qualified as Text+import Effectful.Exception (Exception, displayException)+import GHC.Generics (Generic)+import Prettyprinter (Pretty (..), annotate)+import Prettyprinter.Extra (PrettyAnn (..))+import Prettyprinter.Extra qualified as Pretty+import Prettyprinter.Render.Terminal (AnsiStyle, Color (..), color, colorDull)+import Proto3.Wire (fromProtoEnum)+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | The final status of a 'Span'.+-- See <https://opentelemetry.io/docs/specs/otel/trace/api/#set-status the OpenTelemetry spec>.+data Status = Status+    { message :: Text+    -- ^ A human readable message, typically an error message.+    , code :: Code+    }+    deriving stock (Generic, Eq, Show)+    deriving anyclass (ToJSON)++data Code+    = -- | The operation the span tracked successfully completed without an error.+      Unset+    | -- | The span has been explicitly marked as successful.+      Ok+    | -- | Some error occurred in the operation the span tracked.+      Error+    deriving stock (Show, Eq, Enum, Bounded)++instance ToJSON Code where+    toJSON = Number . fromIntegral . fromEnum++instance PrettyAnn AnsiStyle Code where+    prettyAnn Unset = mempty+    prettyAnn Ok = annotate (colorDull Green) "OK"+    prettyAnn Error = annotate (color Red) "ERROR"++instance Proto.Encode Status where+    encode Status{..} =+        mconcat+            [ Proto.encodeField 2 message+            , Proto.encodeField 3 code+            ]++instance {-# OVERLAPPING #-} Proto.EncodeField Code where+    encodeField n = Proto.encodeField n . fromProtoEnum++instance PrettyAnn AnsiStyle Status where+    prettyAnn Status{..} =+        Pretty.unwords . filter (not . Pretty.null) $+            [ prettyAnn code+            , pretty message+            ]++fromException :: (Exception e) => e -> Status+fromException e =+    Status+        { message = Text.pack . displayException $ e+        , code = Error+        }++-- Create a 'Status' with an 'Ok' code.+ok :: Status+ok =+    Status+        { message = mempty+        , code = Ok+        }++-- Create a 'Status' with an 'Unset' code.+unset :: Status+unset =+    Status+        { message = mempty+        , code = Unset+        }
+ src/Effectful/OpenTelemetry/Tracing/Trace/Flags.hs view
@@ -0,0 +1,101 @@+module Effectful.OpenTelemetry.Tracing.Trace.Flags+    ( Flags (..)+    , Remote (..)+    , toHex+    , toWord8+    , toWord32+    , fromWord32+    )+where++import Control.Monad (mzero)+import Data.Aeson.Types (FromJSON (..), ToJSON (..), Value (..))+import Data.Bits (Bits (..))+import Data.Bool (bool)+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Word (Word32, Word8)+import GHC.Generics (Generic)+import Numeric (readHex, showHex)+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude++-- | Describes whether a span's parent is remote.+data Remote+    = Unknown -- 0x+    | IsNotRemote -- 10+    | IsRemote -- 11+    deriving stock (Show, Eq, Bounded, Enum)++-- | Per-trace flags propagated alongside the 'Trace.ID' and 'Span.ID'.+data Flags = Flags+    { sampled :: Bool+    -- ^ The caller may have recorded trace data. When unset, the+    -- caller did not record trace data out-of-band.+    , random :: Bool+    -- ^ At least the right-most 7 bytes of the 'Trace.ID' have been+    -- selected randomly (or pseudo-randomly) with uniform distribution.+    , remote :: Remote+    }+    deriving stock (Generic, Show, Eq)++sampledBit, randomBit, hasIsRemoteBit, isRemoteBit :: Int+sampledBit = 0+randomBit = 1+hasIsRemoteBit = 8+isRemoteBit = 9++toWord32 :: Flags -> Word32+toWord32 Flags{..} =+    foldr+        (.|.)+        zeroBits+        [ setBit' sampledBit sampled+        , setBit' randomBit random+        , setBit' hasIsRemoteBit $ remote /= Unknown+        , setBit' isRemoteBit $ remote == IsRemote+        ]+  where+    setBit' :: Int -> Bool -> Word32+    setBit' = bool zeroBits . bit++fromWord32 :: Word32 -> Flags+fromWord32 bits =+    Flags+        { sampled = testBit bits sampledBit+        , random = testBit bits randomBit+        , remote = case (testBit bits hasIsRemoteBit, testBit bits isRemoteBit) of+            (False, _) -> Unknown+            (_, False) -> IsNotRemote+            (_, True) -> IsRemote+        }++instance Monoid Flags where+    mempty = fromWord32 zeroBits++instance Semigroup Flags where+    a <> b = fromWord32 $ toWord32 a .|. toWord32 b++toWord8 :: Flags -> Word8+toWord8 Flags{..} = setBit' sampledBit sampled .|. setBit' randomBit random+  where+    setBit' :: Int -> Bool -> Word8+    setBit' = bool zeroBits . bit++-- | Encode the @trace-flags@ field as the two hex characters it is specified as.+toHex :: Flags -> Text+toHex flags = Text.justifyRight 2 '0' . Text.pack $ showHex (toWord8 flags) ""++instance FromJSON Flags where+    parseJSON = \case+        Number n -> pure . fromWord32 . round $ n+        String s -> case readHex (Text.unpack s) of+            [(w, "")] -> pure $ fromWord32 w+            _ -> mzero+        _ -> mzero++instance ToJSON Flags where+    toJSON = Number . fromIntegral . toWord32++instance {-# OVERLAPPING #-} Proto.EncodeField Flags where+    encodeField n = Proto.encodeField n . toWord32
+ src/Effectful/OpenTelemetry/Tracing/Trace/ID.hs view
@@ -0,0 +1,61 @@+module Effectful.OpenTelemetry.Tracing.Trace.ID where++import Data.Aeson.Types (FromJSON (..), ToJSON (..), withText)+import Data.ByteString (ByteString)+import Data.ByteString qualified as ByteString+import Data.ByteString.Base16 qualified as Base16+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text+import System.Random.Extra (uniformIO)+import GHC.Generics (Generic)+import Proto3.Wire.Encode.Class qualified as Proto+import System.Random (Random (..), Uniform)+import System.Random.Stateful (Uniform (..), uniformByteStringM)+import Prelude hiding (length)++-- | A globally unique identifier of a trace.+newtype ID = ID ByteString+    deriving stock (Generic)+    deriving newtype (Eq, Show)++length :: Int+length = 16++-- | Generate a random 'ID'.+new :: IO ID+new = uniformIO++-- | Encode an 'ID' to base-16.+-- The result is always 32 characters long.+toHex :: ID -> Text+toHex (ID bs) = Text.justifyRight (length * 2) '0' . Text.decodeUtf8 . Base16.encode $ bs++toBytes :: ID -> ByteString+toBytes (ID bs) = bs++fromBytes :: ByteString -> Either String ID+fromBytes bs+    | ByteString.length bs /= length = Left $ "Trace ID must have exactly " <> show length <> " bytes"+    | ByteString.all (== 0) bs = Left "Trace ID must not be all zero"+    | otherwise = Right $ ID bs++instance Uniform ID where+    uniformM g = do+        bs <- uniformByteStringM length g+        either (const $ uniformM g) pure $ fromBytes bs++instance Random ID where+    randomR _ = random++instance FromJSON ID where+    parseJSON = withText "ID" $ either fail pure . fromHex+      where+        fromHex :: Text -> Either String ID+        fromHex t = Base16.decode (Text.encodeUtf8 t) >>= fromBytes++instance ToJSON ID where+    toJSON = toJSON . toHex++instance {-# OVERLAPPING #-} Proto.EncodeField ID where+    encodeField n = Proto.byteString n . toBytes
+ src/Effectful/OpenTelemetry/Tracing/Trace/State.hs view
@@ -0,0 +1,152 @@+module Effectful.OpenTelemetry.Tracing.Trace.State+    ( Key+    , key+    , Value+    , value+    , State+    , fromList+    , fromText+    , toText+    , fromByteString+    , toByteString+    , toList+    , lookup+    , insert+    , delete+    )+where++import Data.Aeson (ToJSON (..))+import Data.ByteString (ByteString)+import Data.Char qualified as Char+import Data.Foldable qualified as Foldable+import Data.List qualified as List+import Data.List.Extra qualified as List+import Data.Maybe (fromJust, mapMaybe)+import Data.Sequence (Seq)+import Data.Sequence qualified as Seq+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text+import GHC.IsList (IsList)+import GHC.IsList qualified as GHC+import Proto3.Wire.Encode.Class qualified as Proto+import Prelude hiding (lookup)++-- | W3C TraceState Key+--+-- See <https://www.w3.org/TR/trace-context/#key the W3C specification>.+newtype Key = Key Text+    deriving newtype (Show, Eq, Ord)++-- | W3C TraceState Value+--+-- See <https://www.w3.org/TR/trace-context/#value the W3C specification>.+newtype Value = Value Text+    deriving newtype (Show, Eq)++-- | W3C TraceState+--+-- May hold at most 32 unique 'Key's.+-- The 'Semigroup' operation for 'State' prefers values from the left operand.+newtype State = State (Seq (Key, Value))+    deriving newtype (Show, Eq, Monoid)++instance Semigroup State where+    State left <> State right =+        State . Seq.take 32 $ left <> Seq.filter (not . isInLeft . fst) right+      where+        isInLeft = (`elem` (fst <$> left))++instance IsList State where+    type Item State = (Key, Value)+    fromList = fromList+    toList = toList++instance ToJSON State where+    toJSON = toJSON . toText++instance {-# OVERLAPPING #-} Proto.EncodeField State where+    encodeField n = Proto.encodeField n . toText++toList :: State -> [(Key, Value)]+toList (State s) = Foldable.toList s++fromList :: [(Key, Value)] -> State+fromList = State . Seq.fromList . List.take 32 . List.nubOrdOn fst++-- | Parse a 'State' from a comma-separated list of 'key'@=@'value' pairs.+-- Drops any key-value pairs that cannot be parsed.+-- Evaluates to an empty state if the 'Text' contains no parseable key-value pairs.+fromText :: Text -> State+fromText = fromList . mapMaybe (kv . Text.strip) . Text.splitOn ","+  where+    kv :: Text -> Maybe (Key, Value)+    kv t = case Text.breakOn "=" t of+        ( key -> Just k+            , Text.stripPrefix "=" -> fromJust -> value -> Just v+            ) -> Just (k, v)+        _ -> Nothing++toText :: State -> Text+toText = Text.intercalate "," . fmap (\(Key k, Value v) -> k <> "=" <> v) . toList++toByteString :: State -> ByteString+toByteString = Text.encodeUtf8 . toText++-- | Parse a 'State' from a comma-separated list of 'key'@=@'value' pairs.+-- Drops any key-value pairs that cannot be parsed.+-- Evaluates to an empty state if the 'ByteString' contains no parseable key-value pairs.+fromByteString :: ByteString -> State+fromByteString = either (const mempty) fromText . Text.decodeUtf8'++-- | Construct a 'Key', or 'Nothing' on invalid input.+key :: Text -> Maybe Key+key (Text.strip -> t)+    | isValidKey = Just (Key t)+    | otherwise = Nothing+  where+    isValidKey+        | Text.length t > 256 = False+        | otherwise = case Text.splitOn "@" t of+            [simpleKey] -> isValidSimpleKey simpleKey+            [tenant, system] -> isValidTenantKey tenant system+            _ -> False++    isValidTenantKey tenant system =+        not (Text.null tenant)+            && Text.length tenant <= 241+            && Text.length system <= 14+            && isValidSimpleKey tenant+            && isValidSimpleKey system++    isValidSimpleKey k = case Text.uncons k of+        Just (firstChar, rest) -> Char.isAsciiLower firstChar && Text.all isSimpleChar rest+        Nothing -> False++    isSimpleChar c = Char.isAsciiLower c || Char.isDigit c || c `elem` ['_', '-', '*', '/']++-- | Construct a 'Value', or 'Nothing' on invalid input.+value :: Text -> Maybe Value+value (Text.strip -> t)+    | isValidValue = Just (Value t)+    | otherwise = Nothing+  where+    isValidValue = Text.length t < 256 && Text.all isValidChar t+    isValidChar c =+        let code = fromEnum c+         in code >= 0x20 && code <= 0x7E && c /= ',' && c /= '='++-- | Lookup a 'Value' by 'Key'.+lookup :: Key -> State -> Maybe Value+lookup k = List.lookup k . toList++-- | Insert or update a 'Key'-'Value' pair.+-- If the key already exists, it is moved to the left-most position in 'State'.+-- If the 'State' contains 32 keys before insertion, the right-most key will be removed.+insert :: Key -> Value -> State -> State+insert k v s = State (pure (k, v)) <> s++-- | Remove a 'Value' by 'Key' if it exists.+delete :: Key -> State -> State+delete k (State s) = State $ Seq.filter ((k /=) . fst) s
+ src/Prettyprinter/Extra.hs view
@@ -0,0 +1,39 @@+{-# LANGUAGE CPP #-}+{-# OPTIONS_GHC -Wno-orphans #-}++module Prettyprinter.Extra where++import Data.Aeson qualified as Aeson+import Data.Aeson.Key qualified as Aeson.Key+import Data.Aeson.Text qualified as Aeson+import Data.List qualified as List+import Prettyprinter.Internal (Doc (..), Pretty (..))+#if MIN_VERSION_prettyprinter(1,7,2)+import Prettyprinter.Internal (PrettyAnn (prettyAnn))+#endif+import Prelude++intercalate :: Doc ann -> [Doc ann] -> Doc ann+intercalate x = mconcat . List.intersperse x++unwords :: [Doc ann] -> Doc ann+unwords = intercalate " "++null :: Doc ann -> Bool+null Empty = True+null _ = False++#if !MIN_VERSION_prettyprinter(1,7,2)+class PrettyAnn ann a where+    prettyAnn :: a -> Doc ann+#endif++instance PrettyAnn ann Aeson.Key where+    prettyAnn = pretty . Aeson.Key.toText++instance PrettyAnn ann Aeson.Value where+    prettyAnn (Aeson.String t) = pretty t+    prettyAnn v = pretty $ Aeson.encodeToLazyText v++instance {-# OVERLAPPING #-} PrettyAnn ann (Aeson.Key, Aeson.Value) where+    prettyAnn (key, value) = prettyAnn key <> "=" <> prettyAnn value
+ src/Proto3/Wire/Encode/Class.hs view
@@ -0,0 +1,70 @@+{-# LANGUAGE FlexibleInstances #-}+{-# LANGUAGE UndecidableInstances #-}+{-# OPTIONS_GHC -Wno-orphans #-}++module Proto3.Wire.Encode.Class+    ( module Proto3.Wire.Encode.Class+    , module Proto3.Wire.Encode+    , module Proto3.Wire.Types+    )+where++import Data.Aeson.Key (Key)+import Data.Aeson.Key qualified as Key+import Data.Aeson.KeyMap qualified as KeyMap+import Data.Aeson.Types (Value (..))+import Data.Int (Int32)+import Data.Scientific (toBoundedInteger, toRealFloat)+import Data.Sequence (Seq)+import Data.Text (Text)+import Data.Text.Encoding qualified as Text+import Data.Vector qualified as Vector+import Data.Word (Word32)+import Proto3.Wire.Class (ProtoEnum)+import Proto3.Wire.Encode hiding (text)+import Proto3.Wire.Types (FieldNumber)+import Prelude hiding (span)++class Encode a where+    encode :: a -> MessageBuilder++class EncodeField a where+    encodeField :: FieldNumber -> a -> MessageBuilder++instance (Encode a) => EncodeField a where+    encodeField n = embedded n . encode++instance Encode Value where+    encode (String s) = encodeField 1 s+    encode (Bool b) = bool 2 b+    encode (Number n) = maybe (double 4 $ toRealFloat n) (int64 3) $ toBoundedInteger n+    encode (Array a) = embedded 5 . encodeField 1 . Vector.toList $ a+    encode (Object o) = embedded 6 . encodeField 1 . KeyMap.toList $ o+    encode Null = mempty++instance Encode (Key.Key, Value) where+    encode (k, v) =+        mconcat+            [ encodeField 1 k+            , encodeField 2 v+            ]++instance (Bounded a, Enum a) => ProtoEnum a++instance {-# OVERLAPPING #-} EncodeField Word32 where+    encodeField = fixed32++instance {-# OVERLAPPING #-} EncodeField Int32 where+    encodeField = int32++instance {-# OVERLAPPING #-} EncodeField Text where+    encodeField n = byteString n . Text.encodeUtf8++instance {-# OVERLAPPING #-} EncodeField Key where+    encodeField n = encodeField n . Key.toText++instance {-# OVERLAPPING #-} (EncodeField a) => EncodeField [a] where+    encodeField n = foldMap (encodeField n)++instance {-# OVERLAPPING #-} (EncodeField a) => EncodeField (Seq a) where+    encodeField n = foldMap (encodeField n)
+ src/System/Random/Extra.hs view
@@ -0,0 +1,34 @@+-- |+-- Module      : System.Random.Extra+-- Copyright   : (c) 2026 Institute for Digital Autonomy+-- License     : EUPL-1.2+-- Maintainer  : IDA+--+-- Thread-safe, low-contention generation of random values.+module System.Random.Extra where++import Control.Concurrent (myThreadId, threadCapability)+import Data.IORef (IORef, atomicModifyIORef', newIORef)+import Data.Tuple (swap)+import Data.Vector (Vector)+import Data.Vector qualified as Vector+import GHC.Conc (getNumCapabilities)+import System.IO.Unsafe (unsafePerformIO)+import System.Random (StdGen, newStdGen)+import System.Random.Stateful (Uniform, runStateGen, uniformM)+import Prelude++-- | One 'StdGen' per capability, each seeded independently from the global one.+generators :: Vector (IORef StdGen)+generators = unsafePerformIO do+    capabilities <- getNumCapabilities+    Vector.replicateM capabilities $ newIORef =<< newStdGen+{-# NOINLINE generators #-}++-- | Draw a uniformly random value from the current capability's generator.+uniformIO :: (Uniform a) => IO a+uniformIO = do+    (capability, _pinned) <- threadCapability =<< myThreadId+    let generator = generators Vector.! (capability `mod` Vector.length generators)+    atomicModifyIORef' generator \g -> swap $ runStateGen g uniformM+{-# INLINE uniformIO #-}
+ test/Arbitrary.hs view
@@ -0,0 +1,262 @@+{-# LANGUAGE StandaloneDeriving #-}+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-orphans #-}++module Arbitrary where++import Data.Aeson (toJSON)+import Data.Aeson.Key qualified as Key+import Data.Functor.Syntax ((<$$>), (<&&>))+import Data.List.Extra (nubOrd)+import Data.Scientific+import Data.Text (Text)+import Data.Text qualified as Text+import Effectful.Exception (AssertionFailed (..))+import Effectful.OpenTelemetry.Logging.Severity (Severity)+import Effectful.OpenTelemetry.Metrics.Measurement (AggregationTemporality)+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import Effectful.OpenTelemetry.Protocol.AnyValue (AnyValue (..))+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+import Effectful.OpenTelemetry.Protocol.Attributes qualified as Attributes+import Effectful.OpenTelemetry.Protocol.Transport (Compression (..))+import Effectful.OpenTelemetry.Timestamp (Timestamp (..))+import Effectful.OpenTelemetry.Tracing.Span.Context+import Effectful.OpenTelemetry.Tracing.Span.Event (Event (..))+import Effectful.OpenTelemetry.Tracing.Span.ID qualified as Span (ID)+import Effectful.OpenTelemetry.Tracing.Span.Kind qualified as Span (Kind)+import Effectful.OpenTelemetry.Tracing.Span.Status qualified as Status (Code)+import Effectful.OpenTelemetry.Tracing.Trace.Flags qualified as Trace (Flags (Flags))+import Effectful.OpenTelemetry.Tracing.Trace.Flags qualified as Trace.Flags+import Effectful.OpenTelemetry.Tracing.Trace.ID qualified as Trace (ID)+import Effectful.OpenTelemetry.Tracing.Trace.State qualified as Trace (State)+import Effectful.OpenTelemetry.Tracing.Trace.State qualified as Trace.State+import Effectful.QuickCheck+import GHC.IsList (fromList)+import System.Random (Random (..))+import Prelude++instance Random Scientific where+    random g = (scientific c e, g2)+      where+        (c, g1) = random g+        (e, g2) = random g1++    randomR (lo, hi) g = (fromFloatDigits x, g')+      where+        x :: Double+        (x, g') = randomR (toRealFloat lo, toRealFloat hi) g++instance Arbitrary Scientific where+    arbitrary = frequency [(70, getSmall <$> arbitrary), (30, getLarge <$> arbitrary)]++    shrink s = uncurry scientific <$> shrink (coefficient s, base10Exponent s)++instance {-# OVERLAPPING #-} Arbitrary (NonNegative Scientific) where+    arbitrary = NonNegative <$> ((scientific . getNonNegative <$> arbitrary) <*> arbitrary)++instance {-# OVERLAPPING #-} Arbitrary (Positive Scientific) where+    arbitrary = Positive <$> ((scientific . getPositive <$> arbitrary) <*> arbitrary)++-- | Integral values in the 'Int' range, floating values in the 'Float' range.+instance {-# OVERLAPPING #-} Arbitrary (Small Scientific) where+    arbitrary =+        Small+            <$> oneof+                [ fromIntegral . getSmall <$> (arbitrary :: Gen (Small Int))+                , fromFloatDigits <$> (arbitrary :: Gen Float)+                ]++instance {-# OVERLAPPING #-} Arbitrary (Large Scientific) where+    arbitrary = fmap Large $ scientific <$> arbitrary <*> arbitrary++newtype HistogramSamples = HistogramSamples [Scientific]+    deriving newtype (Eq)+    deriving stock (Show)++instance Arbitrary HistogramSamples where+    arbitrary = do+        n <- chooseInt (4, 12)+        HistogramSamples <$> vectorOf n (getSmall . getPositive <$> arbitrary)++newtype HistogramBounds = HistogramBounds [Scientific]+    deriving newtype (Eq)+    deriving stock (Show)++instance Arbitrary HistogramBounds where+    arbitrary = do+        n <- chooseInt (1, 5)+        HistogramBounds . nubOrd <$> vectorOf n (getSmall . getPositive <$> arbitrary)++newtype HistogramObservations = HistogramObservations [Scientific]+    deriving newtype (Eq)+    deriving stock (Show)++instance Arbitrary HistogramObservations where+    arbitrary = do+        n <- chooseInt (3, 20)+        HistogramObservations <$> vectorOf n (getSmall . getPositive <$> arbitrary)++data BucketedObservations = BucketedObservations+    { bounds :: [Scientific]+    , observations :: [[Scientific]]+    }+    deriving stock (Show, Eq)++instance Arbitrary BucketedObservations where+    arbitrary = do+        lo <- choose (1.0, 50.0)+        hi <- choose (lo + 1.0, lo + 500.0)+        below <- chooseInt (1, 5)+        between <- chooseInt (1, 5)+        above <- chooseInt (1, 5)+        belowSamples <- vectorOf below $ choose (lo - 100.0, lo - 0.001)+        betweenSamples <- vectorOf between $ choose (lo + 0.001, hi - 0.001)+        aboveSamples <- vectorOf above $ choose (hi + 0.001, hi + 100.0)+        pure+            BucketedObservations+                { bounds = [lo, hi]+                , observations = [belowSamples, betweenSamples, aboveSamples]+                }++data ConcurrentObservation = ConcurrentObservation+    { repeats :: Int+    , value :: Scientific+    }+    deriving stock (Show, Eq)++instance Arbitrary ConcurrentObservation where+    arbitrary = do+        repeats <- chooseInt (50, 500)+        Positive (Small value) <- arbitrary+        pure ConcurrentObservation{..}++deriving via String instance Eq AssertionFailed++instance Arbitrary Attributes where+    arbitrary = Attributes.fromList <$> resize 5 arbitrary++instance Arbitrary AnyValue where+    arbitrary = AnyValue <$> resize 5 arbitrary++instance Arbitrary Text where+    arbitrary = Text.pack <$> arbitrary+    shrink = fmap Text.pack . shrink . Text.unpack++instance Arbitrary Metadata where+    arbitrary = do+        attributes <- do+            n <- chooseInt (0, 3)+            Attributes.fromList <$> vectorOf n do+                key <- do+                    c <- elements $ ['a' .. 'z'] <> ['A' .. 'Z']+                    Marker rest <- arbitrary+                    pure . Key.fromText $ Text.cons c rest+                Marker (toJSON -> value) <- arbitrary+                pure (key, value)+        description <- getMarker <$$> arbitrary+        unit <- arbitrary <&&> \(Marker u) -> "{" <> u <> "}"+        pure Metadata{..}++instance Arbitrary Trace.ID where+    arbitrary = chooseAny++instance Arbitrary Span.ID where+    arbitrary = chooseAny++instance Arbitrary Span.Kind where+    arbitrary = arbitraryBoundedEnum++instance Arbitrary Status.Code where+    arbitrary = arbitraryBoundedEnum++instance Arbitrary Compression where+    arbitrary = arbitraryBoundedEnum++instance Arbitrary Severity where+    arbitrary = arbitraryBoundedEnum++instance Arbitrary AggregationTemporality where+    arbitrary = arbitraryBoundedEnum++instance Arbitrary Trace.Flags.Remote where+    arbitrary = arbitraryBoundedEnum++instance Arbitrary Trace.Flags where+    arbitrary = do+        sampled <- arbitrary+        random <- arbitrary+        remote <- arbitrary+        pure Trace.Flags{..}++instance Arbitrary Context where+    arbitrary = do+        traceId <- arbitrary+        spanId <- arbitrary+        traceFlags <- arbitrary+        traceState <- arbitrary+        pure Context{..}++instance Arbitrary Trace.State.Key where+    arbitrary = oneof [simpleKey, tenantKey]+      where+        simpleKey = do+            firstChar <- elements ['a' .. 'z']+            len <- chooseInt (0, 30)+            rest <- vectorOf len genSimpleChar+            maybe (error "Simple key generation failed") pure . Trace.State.key $ Text.pack (firstChar : rest)++        tenantKey = do+            tenant <- boundedSimpleKey 241+            system <- boundedSimpleKey 14+            maybe (error "Tenant key generation failed") pure . Trace.State.key $ tenant <> "@" <> system++        boundedSimpleKey maxLen = do+            firstChar <- elements ['a' .. 'z']+            len <- chooseInt (0, maxLen - 1)+            rest <- vectorOf len genSimpleChar+            pure $ Text.pack (firstChar : rest)++        genSimpleChar :: Gen Char+        genSimpleChar = elements $ ['a' .. 'z'] ++ ['0' .. '9'] ++ ['_', '-', '*', '/']++instance Arbitrary Trace.State.Value where+    arbitrary = do+        len <- chooseInt (0, 255)+        chars <- vectorOf len $ elements [c | c <- ['\x20' .. '\x7E'], c /= ',', c /= '=']+        maybe (error "Trace.State.Value generation failed") pure . Trace.State.value $ Text.pack chars++instance Arbitrary Trace.State where+    arbitrary = fmap fromList . flip vectorOf arbitrary =<< chooseInt (0, 32)++-- | Small alphanumeric 'Text', good enough for keys, values and names that+-- don't need to exercise anything beyond "some text arrived intact".+newtype Marker = Marker {getMarker :: Text}+    deriving newtype (Eq)+    deriving stock (Show)++instance Arbitrary Marker where+    arbitrary = Marker . Text.pack <$> vectorOf 8 (elements alphabet)+      where+        alphabet = ['a' .. 'z'] <> ['A' .. 'Z'] <> ['0' .. '9']++newtype MarkerAttributes = MarkerAttributes Attributes+    deriving newtype (Eq)+    deriving stock (Show)++instance Arbitrary MarkerAttributes where+    arbitrary = do+        n <- chooseInt (0, 3)+        MarkerAttributes . Attributes.fromList <$> vectorOf n do+            Marker (Key.fromText -> key) <- arbitrary+            Marker (toJSON -> value) <- arbitrary+            pure (key, value)++instance Arbitrary Timestamp where+    arbitrary = Timestamp <$> arbitrary++instance Arbitrary Event where+    arbitrary = do+        time <- arbitrary+        Marker name <- arbitrary+        MarkerAttributes attributes <- arbitrary+        pure Event{..}
+ test/Effectful/OpenTelemetry/Exporter/DeadSpec.hs view
@@ -0,0 +1,64 @@+{-# OPTIONS_GHC -Wno-type-defaults #-}++module Effectful.OpenTelemetry.Exporter.DeadSpec (spec) where++import Data.Maybe (isJust)+import Effectful+import Effectful.Concurrent (Concurrent, runConcurrent)+import Effectful.Environment (Environment, runEnvironment)+import Effectful.Hspec+import Effectful.Http2Client (HostName, PortNumber)+import Effectful.OpenTelemetry.Protocol (Compression (..), Encoding (..), Resource (..), Scope (..))+import Effectful.OpenTelemetry.Protocol.Attributes qualified as Attributes+import Effectful.OpenTelemetry.Protocol.Resource qualified as Resource+import Effectful.OpenTelemetry.Protocol.Scope qualified as Scope+import Effectful.OpenTelemetry.Tracing+import Effectful.OpenTelemetry.Tracing.Span.Kind qualified as Span.Kind+import Effectful.Retry (Retry, runRetry)+import Effectful.Timeout (Timeout, runTimeout)+import GHC.Stack (HasCallStack)+import Network.URI (URI)+import Network.URI.Static (uri)+import System.Timeout (timeout)+import Prelude++deadEndpoint :: URI+deadEndpoint = [uri|http://127.0.0.1:1/|]++deadHost :: HostName+deadHost = "127.0.0.1"++deadPort :: PortNumber+deadPort = 1++resource :: Resource+resource = Resource{attributes = Attributes.fromList [("service.name", "negative-path")]}++scope :: Scope+scope =+    Scope+        { name = "negative"+        , version = "0.0.0"+        , attributes = mempty+        }++-- | Fail rather than hang if a dead collector ever blocks the caller.+completes+    :: (HasCallStack, Hspec :> es, IOE :> es)+    => Eff '[Environment, Timeout, Retry, Concurrent, IOE] ()+    -> Eff es ()+completes act = do+    done <-+        liftIO . timeout 30_000_000 . runEff . runConcurrent . runRetry . runTimeout . runEnvironment $+            act+    done `shouldSatisfy` isJust++spec :: (HasCallStack, IOE :> es, Hspec :> es) => Eff es ()+spec = describe "connection refusals do not block 'inSpan'" do+    it "HTTP/JSON" . completes . runHttpTracing resource scope Json NoCompression deadEndpoint $ doomed+    it "HTTP/Protobuf" . completes . runHttpTracing resource scope Proto NoCompression deadEndpoint $+        doomed+    it "gRPC" . completes . runGrpcTracing resource scope deadHost deadPort NoCompression $ doomed+  where+    doomed :: (Tracing :> es) => Eff es ()+    doomed = inSpan "doomed" Span.Kind.Internal mempty $ pure ()
+ test/Effectful/OpenTelemetry/Exporter/Grafana/Loki.hs view
@@ -0,0 +1,119 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-orphans #-}++module Effectful.OpenTelemetry.Exporter.Grafana.Loki where++import Arbitrary+import Control.Monad ((>=>))+import Data.Aeson+import Data.Aeson.KeyMap (KeyMap)+import Data.Aeson.KeyMap qualified as KeyMap+import Data.Aeson.Types (Parser)+import Data.Text (Text)+import Data.Text qualified as Text+import Effectful+import Effectful.Concurrent (Concurrent)+import Effectful.Environment (Environment)+import Effectful.HUnit (HUnit)+import Effectful.Hspec+import Effectful.HttpClient (httpLbs, parseRequest_, responseBody, runHttpClientTls)+import Effectful.OpenTelemetry.Exporter.Grafana.Polling (checkReady, pollOrFail)+import Effectful.OpenTelemetry.Logging (Logging, log)+import Effectful.OpenTelemetry.Logging.LogRecord (LogRecord (..))+import Effectful.OpenTelemetry.Logging.Severity (Severity)+import Effectful.OpenTelemetry.Metrics (Metrics)+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes (..))+import Effectful.OpenTelemetry.Protocol.Transport (Protocol (..))+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.OpenTelemetry.Tracing (Tracing, inSpan)+import Effectful.OpenTelemetry.Tracing.Span.Context (Context (..))+import Effectful.OpenTelemetry.Tracing.Span.Kind qualified as Kind+import Effectful.Retry (Retry)+import Effectful.Timeout (Timeout)+import GHC.Generics (Generic)+import GHC.Stack (HasCallStack)+import Network.HTTP.Client (responseStatus)+import Network.HTTP.Types.Status (statusCode)+import Network.URI (URI (..), escapeURIString, isUnreserved)+import Text.Read (readMaybe)+import Util+import Prelude hiding (log)++-- HACK: Loki does not include the timestamp in the log records+data Stream = Stream+    { stream :: KeyMap Value+    , values :: [(Text, Text)]+    }+    deriving stock (Generic)+    deriving anyclass (FromJSON)++instance {-# OVERLAPPING #-} FromJSON [LogRecord] where+    parseJSON = parseJSON @[Stream] >=> fmap concat . traverse \Stream{..} -> traverse (parseEntry stream) values+      where+        parseEntry :: KeyMap Value -> (Text, Text) -> Parser LogRecord+        parseEntry labels (timeUnixNano, lineText) = do+            Object line <- either fail pure $ eitherDecodeStrictText lineText+            observedTimestamp <- parseJSON @Timestamp $ String timeUnixNano+            let o = line <> labels+            traceId <- o .:? "traceid"+            spanId <- o .:? "spanid"+            traceFlags <- o .:? "flags" .!= mempty+            let context = case (traceId, spanId) of+                    (Just traceId', Just spanId') ->+                        Just Context{traceId = traceId', spanId = spanId', traceFlags, traceState = mempty}+                    _ -> Nothing+            severity <- o .: "level"+            body <- o .:? "body" .!= Null+            attributes <- Attributes <$> (o .:? "attributes" .!= mempty)+            eventName <- o .:? "name"+            pure LogRecord{timestamp = observedTimestamp, ..}++instance FromJSON Severity where+    parseJSON = withText "Severity" $ maybe (fail "invalid severity") pure . readMaybe . Text.unpack++fetchLines :: (IOE :> es) => URI -> Text -> Eff es (Either String [LogRecord])+fetchLines baseUri needle = do+    let q = "{service_name=\"otel-effectful-test\"} |= `" <> Text.unpack needle <> "`"+        url =+            show+                baseUri+                    { uriPath = "/loki/api/v1/query_range"+                    , uriQuery = "?query=" <> escapeURIString isUnreserved q+                    }+    runHttpClientTls $ do+        resp <- httpLbs $ parseRequest_ url+        let code = statusCode $ responseStatus resp+        pure $ case code of+            200 -> case decode (responseBody resp) of+                Just (Object obj)+                    | Just (Object d) <- KeyMap.lookup "data" obj+                    , Just results <- KeyMap.lookup "result" d ->+                        case fromJSON results of+                            Success [] -> Left "Loki responded with no logs"+                            Success (records :: [LogRecord]) -> pure records+                            Error err -> Left $ "Loki response failed to parse: " <> err+                _ -> Left "Loki response was not the expected shape"+            _ -> Left $ "Loki returned HTTP " <> show code++spec+    :: (HasCallStack, IOE :> es, HUnit :> es, Hspec :> es, Retry :> es, Timeout :> es)+    => Protocol+    -> URI+    -> ( forall a+          . Eff '[Metrics, Logging, Tracing, Environment, Timeout, Retry, Concurrent, IOE] a+         -> Eff es (a, [LogRecord])+       )+    -> Eff es ()+spec (protocolLabel -> label) lokiUri runTest = describe "Loki" . parallel $ do+    prop "log line round-trip" \(Marker marker) (MarkerAttributes attrs) severity -> do+        checkReady "Loki" lokiUri+        let needle = label <> "-log-" <> marker+            spanName = needle <> "-span"+            eventName = "log.event." <> marker+        (_, capturedLogs) <-+            runTest . inSpan spanName Kind.Internal mempty $+                log severity (String needle) attrs (Just eventName)+        theLog <- only capturedLogs+        -- Alloy does not forward eventName+        let expected = theLog{eventName = Nothing}+        pollOrFail (fetchLines lokiUri needle) (`shouldBe` pure expected)
+ test/Effectful/OpenTelemetry/Exporter/Grafana/Mimir.hs view
@@ -0,0 +1,212 @@+module Effectful.OpenTelemetry.Exporter.Grafana.Mimir where++import Arbitrary+import Control.Applicative ((<|>))+import Data.Aeson (Key, Value (..), decode, toJSON)+import Data.Aeson.Key qualified as Key+import Data.Aeson.KeyMap qualified as KeyMap+import Data.Char (isAlphaNum, isAscii)+import Data.Foldable (toList)+import Data.Functor ((<&>))+import Data.Maybe (listToMaybe)+import Data.Scientific (Scientific)+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Word (Word64)+import Effectful+import Effectful.Concurrent (Concurrent)+import Effectful.Environment (Environment)+import Effectful.HUnit (HUnit, assertFailure)+import Effectful.Hspec+import Effectful.HttpClient (httpLbs, parseRequest_, responseBody, runHttpClientTls)+import Effectful.OpenTelemetry.Exporter.Grafana.Polling (checkReady, pollOrFail)+import Effectful.OpenTelemetry.Logging (Logging)+import Effectful.OpenTelemetry.Metrics (Metrics)+import Effectful.OpenTelemetry.Metrics.Counter qualified as Counter+import Effectful.OpenTelemetry.Metrics.Gauge qualified as Gauge+import Effectful.OpenTelemetry.Metrics.Histogram qualified as Histogram+import Effectful.OpenTelemetry.Metrics.Measurement+    ( HistogramDataPoint (..)+    , Measurement (..)+    , NumberDataPoint (..)+    )+import Effectful.OpenTelemetry.Metrics.Metadata (Metadata (..))+import Effectful.OpenTelemetry.Metrics.Metadata qualified as Metadata+import Effectful.OpenTelemetry.Metrics.UpDownCounter qualified as UpDownCounter+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+import Effectful.OpenTelemetry.Protocol.Attributes qualified as Attributes+import Effectful.OpenTelemetry.Protocol.Transport (Protocol (..))+import Effectful.OpenTelemetry.Tracing (Tracing)+import Effectful.QuickCheck (NonNegative (..), Small (..))+import Effectful.Retry (Retry)+import Effectful.Timeout (Timeout)+import GHC.Stack (HasCallStack)+import Network.HTTP.Client (responseStatus)+import Network.HTTP.Types.Status (statusCode)+import Network.URI (URI (..))+import Text.Read (readMaybe)+import Util+import Prelude++data NumberSample = NumberSample+    { value :: Scientific+    , attributes :: Attributes+    }+    deriving stock (Eq, Show)++numberSample :: Measurement -> Maybe NumberSample+numberSample m = numberDataPoint m <&> \NumberDataPoint{..} -> NumberSample{..}++extractNumberSample :: Value -> Maybe NumberSample+extractNumberSample v = do+    Object res <- Just v+    value <- extractInstantValue v+    Object metric <- KeyMap.lookup "metric" res+    pure+        NumberSample+            { value+            , attributes =+                Attributes.fromList+                    [ (k, val)+                    | (k, val) <- KeyMap.toList metric+                    , k `notElem` reservedLabels+                    ]+            }++data HistogramStats = HistogramStats+    { count :: Word64+    , total :: Maybe Scientific+    , buckets :: Int+    }+    deriving stock (Eq, Show)++histogramStats :: Measurement -> Maybe HistogramStats+histogramStats m =+    histogramDataPoint m <&> \dp ->+        HistogramStats+            { count = dp.count+            , total = dp.sum+            , buckets = length dp.explicitBounds + 1+            }++fetchInstant :: (IOE :> es) => (Value -> Maybe a) -> URI -> Text -> Eff es (Either String a)+fetchInstant extract baseUri metricName = do+    let url =+            show+                baseUri+                    { uriPath = "/prometheus/api/v1/query"+                    , uriQuery = "?query=" <> Text.unpack metricName+                    }+    runHttpClientTls $ do+        resp <- httpLbs $ parseRequest_ url+        let code = statusCode $ responseStatus resp+        pure $ case code of+            200 -> case decode (responseBody resp) of+                Just (Object obj)+                    | Just (Object d) <- KeyMap.lookup "data" obj+                    , Just (Array results) <- KeyMap.lookup "result" d ->+                        case listToMaybe (toList results) >>= extract of+                            Just v -> Right v+                            Nothing -> Left $ "Mimir: no samples for " <> Text.unpack metricName+                _ -> Left "Mimir response was not the expected shape"+            _ -> Left $ "Mimir returned HTTP " <> show code++fetchMetric :: (IOE :> es) => URI -> Text -> Eff es (Either String Scientific)+fetchMetric = fetchInstant extractInstantValue++fetchNumberSample :: (IOE :> es) => URI -> Text -> Eff es (Either String NumberSample)+fetchNumberSample = fetchInstant extractNumberSample++extractInstantValue :: Value -> Maybe Scientific+extractInstantValue v = do+    Object res <- Just v+    Array a <- KeyMap.lookup "value" res+    case toList a of+        [_, String txt] -> readMaybe $ Text.unpack txt+        _ -> Nothing++reservedLabels :: [Key]+reservedLabels = ["__name__", "job", "instance"]++hasDescriptionAndUnit :: Metadata -> Measurement -> Bool+hasDescriptionAndUnit metadata m = m.description == metadata.description && m.unit == metadata.unit++hasMetadata :: Metadata -> Measurement -> Bool+hasMetadata metadata m =+    hasDescriptionAndUnit metadata m+        && (numberAttributes <|> histogramAttributes) == Just metadata.attributes+  where+    numberAttributes = numberDataPoint m <&> (.attributes)+    histogramAttributes = histogramDataPoint m <&> (.attributes)++sanitiseName :: Text -> Text+sanitiseName = Text.map \c -> if isAsciiAlphaNum c || c == '_' then c else '_'+  where+    isAsciiAlphaNum c = isAscii c && isAlphaNum c++spec+    :: forall es+     . (HasCallStack, IOE :> es, HUnit :> es, Hspec :> es, Retry :> es, Timeout :> es)+    => Protocol+    -> URI+    -> ( forall a+          . Eff '[Metrics, Logging, Tracing, Environment, Timeout, Retry, Concurrent, IOE] a+         -> Eff es (a, [Measurement])+       )+    -> Eff es ()+spec (protocolLabel -> label) mimirUri runTest = describe "Mimir" . parallel $ do+    let instrumentName kind marker = sanitiseName $ label <> "_" <> kind <> "_" <> marker+        numberRoundTrip metadata selector act = do+            checkReady "Mimir" mimirUri+            m <- only . snd =<< runTest act+            m `shouldSatisfy` hasDescriptionAndUnit metadata+            sample <- maybe (assertFailure "measurement had no numeric value") pure $ numberSample m+            pollOrFail (fetchNumberSample mimirUri selector) (`shouldBe` sample)++    prop "counter round-trip" \(Marker marker) metadata (NonNegative (Small value)) ->+        let name = instrumentName "counter" marker+         in numberRoundTrip metadata (name <> "_total") do+                c <- Counter.new name metadata+                Counter.add c value++    prop "up-down counter round-trip" \(Marker marker) metadata (Small value) ->+        let name = instrumentName "updown" marker+         in numberRoundTrip metadata name do+                c <- UpDownCounter.new name metadata+                UpDownCounter.add c value++    prop "gauge round-trip" \(Marker marker) metadata (NonNegative (Small value)) ->+        let name = instrumentName "gauge" marker+         in numberRoundTrip metadata name do+                g <- Gauge.new name metadata+                Gauge.set g value++    prop "instrument attributes as labels" \(Marker marker) (Marker labelValue) metadata (NonNegative (Small increment)) ->+        let name = instrumentName "labelled" marker+            labelKey = "attr_kind" :: Text+            metadata' = metadata{Metadata.attributes = Attributes.fromList [(Key.fromText labelKey, toJSON labelValue)]}+            selector = name <> "_total{" <> labelKey <> "=\"" <> labelValue <> "\"}"+         in numberRoundTrip metadata' selector do+                c <- Counter.new name metadata'+                Counter.add c increment++    prop "histogram round-trip" \(Marker marker) metadata (HistogramSamples samples) (HistogramBounds bounds) -> do+        checkReady "Mimir" mimirUri+        let base = sanitiseName $ label <> "_hist_" <> marker+            countSeries = base <> "_count"+            sumSeries = base <> "_sum"+        m <-+            only . snd =<< runTest do+                h <- Histogram.newWithBounds base bounds metadata+                mapM_ (Histogram.record h) samples+        m `shouldSatisfy` hasMetadata metadata+        stats <- maybe (assertFailure "measurement had no histogram data") pure $ histogramStats m+        total <- maybe (assertFailure "measurement had no sum") pure stats.total+        pollOrFail (fetchMetric mimirUri countSeries) (`shouldBe` fromIntegral stats.count)+        pollOrFail (fetchMetric mimirUri sumSeries) \fetched -> abs (fetched - total) `shouldSatisfy` (< 0.01)+        pollOrFail+            (fetchMetric mimirUri $ "count(" <> base <> "_bucket)")+            (`shouldBe` fromIntegral stats.buckets)+        pollOrFail+            (fetchMetric mimirUri $ "max(" <> base <> "_bucket)")+            (`shouldBe` fromIntegral stats.count)
+ test/Effectful/OpenTelemetry/Exporter/Grafana/Polling.hs view
@@ -0,0 +1,100 @@+module Effectful.OpenTelemetry.Exporter.Grafana.Polling where++import Control.Monad (unless)+import Data.Either (isLeft)+import Data.Functor (void)+import Data.Maybe (fromMaybe)+import Effectful+import Effectful.Exception (bracket, catchSync)+import Effectful.HUnit (HUnit, assertFailure)+import Effectful.Hspec+import Effectful.HttpClient (httpLbs, parseRequest_, runHttpClientTls)+import Effectful.Retry+    ( Retry+    , capDelay+    , constantDelay+    , exponentialBackoff+    , limitRetriesByCumulativeDelay+    , retrying+    )+import Effectful.Timeout (Timeout, timeout)+import GHC.Stack (HasCallStack)+import Network.HTTP.Client (responseStatus)+import Network.HTTP.Types.Status (statusIsSuccessful)+import Network.Simple.TCP (closeSock, connectSock)+import Network.URI (URI (..), URIAuth (..))+import Prelude++pollRetrying+    :: forall es a+     . (Retry :> es, Timeout :> es)+    => Eff es (Either String a)+    -> Eff es (Either String a)+pollRetrying =+    retrying+        ( limitRetriesByCumulativeDelay 30_000_000+            . capDelay 2_000_000+            $ exponentialBackoff 250_000+        )+        (const $ pure . isLeft)+        . const+        . attempt+  where+    attempt :: Eff es (Either String a) -> Eff es (Either String a)+    attempt = fmap (fromMaybe $ Left "attempt timed out") . timeout 5_000_000++pollOrFail+    :: (HasCallStack, Retry :> es, HUnit :> es, Timeout :> es)+    => Eff es (Either String a)+    -> (a -> Eff es r)+    -> Eff es r+pollOrFail fetch onFound =+    pollRetrying fetch >>= \case+        Left reason -> assertFailure reason+        Right a -> onFound a++-- | Skip the surrounded example via 'pendingWith' when nothing is accepting+-- connections, for a collector with no @\/ready@ endpoint of its own.+withListening+    :: (IOE :> es, Hspec :> es)+    => String+    -> URI+    -> Eff es ()+    -> Eff es ()+withListening backend backendUri body = do+    listening <- isListening backendUri+    if listening+        then body+        else pendingWith (backend <> " is not listening at " <> show backendUri)++-- | Like 'withListening', but as a statement rather than a wrapper+checkListening :: (IOE :> es, Hspec :> es) => String -> URI -> Eff es ()+checkListening backend backendUri = do+    listening <- isListening backendUri+    unless listening $ pendingWith (backend <> " is not listening at " <> show backendUri)++-- | 'checkListening', plus waiting for the backend to report readiness.+checkReady :: (IOE :> es, Hspec :> es, Retry :> es) => String -> URI -> Eff es ()+checkReady backend backendUri = do+    checkListening backend backendUri+    waitReady backendUri++isListening :: (IOE :> es) => URI -> Eff es Bool+isListening baseUri = case baseUri of+    URI{uriAuthority = Just URIAuth{uriRegName, ..}}+        | ':' : port <- uriPort -> do+            (bracket (connectSock uriRegName port) (closeSock . fst) . const $ pure True)+                `catchSync` const (pure False)+    _ -> pure False++-- | Query the backend readiness via the standard Grafana-stack @\/ready@ endpoint.+waitReady :: (IOE :> es, Retry :> es) => URI -> Eff es ()+waitReady baseUri = do+    let url = show baseUri{uriPath = "/ready"}+    void+        . retrying+            (limitRetriesByCumulativeDelay 30_000_000 $ constantDelay 1_000_000)+            (const $ pure . not . statusIsSuccessful)+        . const+        . runHttpClientTls+        $ responseStatus <$> httpLbs (parseRequest_ url)
+ test/Effectful/OpenTelemetry/Exporter/Grafana/Tempo.hs view
@@ -0,0 +1,241 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-orphans #-}++module Effectful.OpenTelemetry.Exporter.Grafana.Tempo where++import Arbitrary+import Data.Aeson+import Data.Aeson.Key qualified as Key+import Data.Aeson.KeyMap qualified as KeyMap+import Data.Aeson.Types (Parser, typeMismatch)+import Data.ByteString (ByteString)+import Data.ByteString.Base16 qualified as Base16+import Data.ByteString.Base64 qualified as Base64+import Data.Foldable (for_)+import Data.Functor.Syntax ((<&&>))+import Data.Maybe (listToMaybe)+import Data.Sequence qualified as Seq+import Data.Text (Text)+import Data.Text qualified as Text+import Data.Text.Encoding qualified as Text+import Data.Tuple.Extra (thd3)+import Effectful+import Effectful.Concurrent (Concurrent)+import Effectful.Environment (Environment)+import Effectful.Exception (catchSync, throwIO)+import Effectful.HUnit (HUnit, assertFailure)+import Effectful.Hspec+import Effectful.HttpClient (httpLbs, parseRequest_, responseBody, runHttpClientTls)+import Effectful.OpenTelemetry.Exporter.Grafana.Polling (checkReady, pollOrFail)+import Effectful.OpenTelemetry.Logging (Logging)+import Effectful.OpenTelemetry.Metrics (Metrics)+import Effectful.OpenTelemetry.Protocol (Resource, Scope)+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+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.Timestamp (Timestamp (..))+import Effectful.OpenTelemetry.Tracing+    ( Tracing+    , addEvent+    , currentContext+    , inOkSpan+    , inSpan+    , withContext+    )+import Effectful.OpenTelemetry.Tracing.Span (Span (..))+import Effectful.OpenTelemetry.Tracing.Span.Context (Context (..))+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span.Context+import Effectful.OpenTelemetry.Tracing.Span.Event (Event (..))+import Effectful.OpenTelemetry.Tracing.Span.ID qualified as Span.ID+import Effectful.OpenTelemetry.Tracing.Span.Kind (Kind (..))+import Effectful.OpenTelemetry.Tracing.Span.Kind qualified as Kind+import Effectful.OpenTelemetry.Tracing.Span.Status (Code, Status (..))+import Effectful.OpenTelemetry.Tracing.Span.Status qualified as Status+import Effectful.OpenTelemetry.Tracing.Trace.ID qualified as Trace.ID+import Effectful.OpenTelemetry.Tracing.Trace.State qualified as Trace.State+import Effectful.QuickCheck (arbitrary, generate)+import Effectful.Retry (Retry)+import Effectful.Timeout (Timeout)+import GHC.Stack (HasCallStack)+import Network.HTTP.Client (responseStatus)+import Network.HTTP.Types.Status (statusCode)+import Network.URI (URI (..))+import Text.Read (readMaybe)+import Util+import Prelude++instance FromJSON (Resource.Items Span) where+    parseJSON = withObject "batch" \batch -> do+        resource <- batch .:? "resource" .!= mempty+        scopeItems <- batch .:? "scopeSpans" .!= []+        pure Resource.Items{resource, scopeItems}++instance FromJSON Resource.Resource where+    parseJSON = withObject "resource" $ fmap Resource.Resource . attributesAt "attributes"++instance FromJSON (Scope.Items Span) where+    parseJSON = withObject "scopeSpans" \ss -> do+        scope <- ss .:? "scope" .!= emptyScope+        items <- ss .:? "spans" .!= []+        pure Scope.Items{scope, items}+      where+        emptyScope = Scope.Scope{name = "", version = "", attributes = mempty}++instance FromJSON Scope.Scope where+    parseJSON = withObject "scope" \sc -> do+        name <- sc .:? "name" .!= ""+        version <- sc .:? "version" .!= ""+        attributes <- attributesAt "attributes" sc+        pure Scope.Scope{..}++instance FromJSON Span where+    parseJSON = withObject "Span" \o -> do+        context <- do+            traceId <- parseId Trace.ID.fromBytes =<< o .: "traceId"+            spanId <- parseId Span.ID.fromBytes =<< o .: "spanId"+            traceFlags <- o .:? "flags" .!= mempty+            traceState <- maybe mempty Trace.State.fromText <$> o .:? "traceState"+            pure Context{..}+        name <- o .:? "name" .!= ""+        parentSpanId <- traverse (parseId Span.ID.fromBytes) =<< o .:? "parentSpanId"+        kind <- o .: "kind"+        startTime <- parseTimestamp =<< o .: "startTimeUnixNano"+        endTime <- mapM parseTimestamp =<< o .:? "endTimeUnixNano"+        attributes <- attributesAt "attributes" o+        events <- Seq.fromList <$> o .:? "events" .!= []+        status <- o .:? "status" .!= Status.unset+        pure Span{..}+      where+        parseId :: (ByteString -> Either String a) -> Text.Text -> Parser a+        parseId fromBytes t@(Text.encodeUtf8 -> bs) =+            case (fromBytes =<< Base16.decode bs, fromBytes =<< Base64.decode bs) of+                (Right a, _) -> pure a+                (_, Right a) -> pure a+                (Left hexErr, Left b64Err) ->+                    fail $ "neither valid hex (" <> hexErr <> ") nor base64 (" <> b64Err <> "): " <> Text.unpack t++instance FromJSON Event where+    parseJSON = withObject "event" \e -> do+        name <- e .: "name"+        time <- parseTimestamp =<< e .: "timeUnixNano"+        attributes <- attributesAt "attributes" e+        pure Event{..}++instance FromJSON Status where+    parseJSON = withObject "status" \st -> do+        message <- st .:? "message" .!= ""+        code <- st .:? "code" .!= Status.Unset+        pure Status{..}++instance FromJSON Code where+    parseJSON (Number n) = case round n :: Int of+        0 -> pure Status.Unset+        1 -> pure Status.Ok+        2 -> pure Status.Error+        n' -> fail $ "unknown status code: " <> show n'+    parseJSON (String s) = case s of+        "STATUS_CODE_UNSET" -> pure Status.Unset+        "STATUS_CODE_OK" -> pure Status.Ok+        "STATUS_CODE_ERROR" -> pure Status.Error+        _ -> fail $ "unknown status code: " <> Text.unpack s+    parseJSON v = typeMismatch "Status.Code" v++instance FromJSON Kind where+    parseJSON (Number n) = case round n :: Int of+        1 -> pure Internal+        2 -> pure Server+        3 -> pure Client+        4 -> pure Producer+        5 -> pure Consumer+        n' -> fail $ "unknown span kind: " <> show n'+    parseJSON (String s) = case s of+        "SPAN_KIND_INTERNAL" -> pure Internal+        "SPAN_KIND_SERVER" -> pure Server+        "SPAN_KIND_CLIENT" -> pure Client+        "SPAN_KIND_PRODUCER" -> pure Producer+        "SPAN_KIND_CONSUMER" -> pure Consumer+        _ -> fail $ "unknown span kind: " <> Text.unpack s+    parseJSON v = typeMismatch "Kind" v++parseTimestamp :: Value -> Parser Timestamp+parseTimestamp (Number n) = pure . Timestamp $ round n+parseTimestamp (String t) = maybe (fail "invalid timestamp") (pure . Timestamp) . readMaybe $ Text.unpack t+parseTimestamp v = typeMismatch "Timestamp" v++attributesAt :: Key.Key -> Object -> Parser Attributes+attributesAt k o = o .:? k .!= mempty++fetchTrace :: (IOE :> es) => URI -> Text -> Eff es (Either String [(Resource, Scope, Span)])+fetchTrace baseUri traceIdHex = do+    let url = show baseUri{uriPath = "/api/traces/" <> Text.unpack traceIdHex}+    runHttpClientTls $ do+        resp <- httpLbs $ parseRequest_ url+        let code = statusCode $ responseStatus resp+        pure $ case code of+            200 -> case decode (responseBody resp) of+                Just (Object o) -> case KeyMap.lookup "batches" o of+                    Just batches -> case fromJSON batches of+                        Success (batchItems :: [Resource.Items Span]) -> case flattenBatches batchItems of+                            [] -> Left "Tempo returned trace with no spans"+                            ss -> Right ss+                        Error err -> Left $ "Tempo response failed to parse: " <> err+                    Nothing -> Left "Tempo response had no batches field"+                _ -> Left "Tempo response was not a JSON object"+            _ -> Left $ "Tempo returned HTTP " <> show code++flattenBatches :: [Resource.Items a] -> [(Resource, Scope, a)]+flattenBatches batchItems =+    [ (resource, scope, s)+    | Resource.Items{resource, scopeItems} <- batchItems+    , Scope.Items{scope, items} <- scopeItems+    , s <- items+    ]++spec+    :: forall es+     . (HasCallStack, IOE :> es, HUnit :> es, Hspec :> es, Retry :> es, Timeout :> es)+    => Protocol+    -> URI+    -> Resource+    -> Scope+    -> ( forall a+          . Eff '[Metrics, Logging, Tracing, Environment, Timeout, Retry, Concurrent, IOE] a+         -> Eff es (a, [Span])+       )+    -> Eff es ()+spec (protocolLabel -> label) tempoUri resource scope runTest = describe "Tempo" . parallel $ do+    Marker suffix <- liftIO $ generate arbitrary+    let nameFor :: Text -> Text+        nameFor kind = label <> "-" <> kind <> "-" <> suffix+        roundTrip act check = do+            checkReady "Tempo" tempoUri+            ((), captured) <- runTest act+            firstSpan <- maybe (assertFailure "no in-memory spans") pure $ listToMaybe captured+            pollOrFail (fetchTrace tempoUri $ Trace.ID.toHex firstSpan.context.traceId) (check captured)++    prop "span round-trip" \(Marker marker) spanKind (MarkerAttributes attrs) traceState events -> do+        let innerSpanName = label <> "-trace-" <> marker+            outerSpanName = innerSpanName <> "-outer"+        roundTrip+            ( inSpan outerSpanName Kind.Internal mempty do+                ctx <- currentContext <&&> \ctx -> ctx{Span.Context.traceState}+                maybe id withContext ctx . inSpan innerSpanName spanKind attrs $+                    mapM_ @[] addEvent events+            )+            \captured fetched -> do+                captured `shouldBe` (thd3 <$> fetched)+                for_ fetched \(tempoResource, tempoScope, _) -> do+                    tempoScope `shouldBe` scope+                    tempoResource `shouldBe` resource++    it "Ok status round-trip" $+        roundTrip (inOkSpan (nameFor "ok") Kind.Internal mempty $ pure ()) \captured fetched ->+            (thd3 <$> fetched) `shouldBe` captured++    it "Error status round-trip" $+        roundTrip+            ( (inSpan (nameFor "error") Kind.Internal mempty . throwIO . userError $ "boom-" <> Text.unpack suffix)+                `catchSync` const (pure ())+            )+            \captured fetched -> (thd3 <$> fetched) `shouldBe` captured
+ test/Effectful/OpenTelemetry/Exporter/GrafanaSpec.hs view
@@ -0,0 +1,318 @@+{-# 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
+ test/Effectful/OpenTelemetry/ExporterSpec.hs view
@@ -0,0 +1,36 @@+module Effectful.OpenTelemetry.ExporterSpec (spec) where++import Data.Foldable (for_)+import Effectful+import Effectful.Concurrent (runConcurrent)+import Effectful.Concurrent.STM (atomically, flushTQueue, newTQueueIO, writeTQueue)+import Effectful.Hspec+import Effectful.OpenTelemetry.Exporter qualified as Exporter+import Prelude++spec :: (IOE :> es, Hspec :> es) => Eff es ()+spec = runConcurrent do+    prop "associativity" \(xs :: [Int]) -> do+        q <- newTQueueIO+        let tag t = Exporter.atomically $ writeTQueue q . (t,)+            a = tag 'a'+            b = tag 'b'+            c = tag 'c'+            expected = [(t, x) | x <- xs, t <- "abc"]+            actual = atomically $ flushTQueue q++        Exporter.withExporter ((a <> b) <> c) $ liftIO . for_ xs+        actual `shouldReturn` expected++        Exporter.withExporter (a <> (b <> c)) $ liftIO . for_ xs+        actual `shouldReturn` expected++    prop "identity" \(xs :: [Int]) -> do+        let composeWith f = do+                q <- newTQueueIO+                let e = f $ Exporter.tqueue q+                Exporter.withExporter e $ liftIO . for_ xs+                atomically $ flushTQueue q+        composeWith id `shouldReturn` xs+        composeWith (mempty <>) `shouldReturn` xs+        composeWith (<> mempty) `shouldReturn` xs
+ test/Effectful/OpenTelemetry/Logging/SeveritySpec.hs view
@@ -0,0 +1,18 @@+module Effectful.OpenTelemetry.Logging.SeveritySpec (spec) where++import Arbitrary ()+import Data.Char (toLower, toUpper)+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Logging.Severity (Severity)+import Effectful.QuickCheck ((.&&.), (===))+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = parallel do+    prop "toEnum/fromEnum round trip" \(s :: Severity) ->+        toEnum (fromEnum s) === s+    prop "read/show round trip" \(s :: Severity) ->+        read (show s) === s+    prop "read is case-insensitive" \(s :: Severity) ->+        (read (map toLower (show s)) === s) .&&. (read (map toUpper (show s)) === s)
+ test/Effectful/OpenTelemetry/LoggingSpec.hs view
@@ -0,0 +1,92 @@+{-# OPTIONS_GHC -Wno-name-shadowing #-}+{-# OPTIONS_GHC -Wno-type-defaults #-}++module Effectful.OpenTelemetry.LoggingSpec (spec) where++import Arbitrary (Marker (..))+import Data.Aeson.Types (Value (..))+import Data.Functor ((<&>))+import Data.Maybe (isJust)+import Effectful+import Effectful.Concurrent (runConcurrent)+import Effectful.Concurrent.Async (replicateConcurrently_)+import Effectful.HUnit (HUnit)+import Effectful.Hspec+import Effectful.OpenTelemetry.Logging+import Effectful.OpenTelemetry.Logging.Severity qualified as Severity+import Effectful.OpenTelemetry.Timestamp (Timestamp (..))+import Effectful.OpenTelemetry.Tracing+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span.Context+import Effectful.OpenTelemetry.Tracing.Span.Kind qualified as Span.Kind+import Util+import Prelude hiding (log)++info :: (Logging :> es) => Value -> Eff es ()+info v = log Severity.Info v mempty Nothing++spec :: (IOE :> es, HUnit :> es, Hspec :> es) => Eff es ()+spec = runConcurrent . parallel $ do+    prop "emits a single log record" \severity body attributes name -> do+        record <-+            only . snd =<< (runNoTracing . runInMemoryLogging) do+                log severity body attributes (getMarker <$> name)+        record.severity `shouldBe` severity+        record.body `shouldBe` body+        record.attributes `shouldBe` attributes+        record.eventName `shouldBe` (getMarker <$> name)+        record.timestamp.nanos `shouldSatisfy` (> 0)++    prop "emits multiple log records in order" \body1 body2 body3 -> do+        (_, records) <- runNoTracing . runInMemoryLogging $ do+            log Severity.Info body1 mempty Nothing+            log Severity.Debug body2 mempty Nothing+            log Severity.Error body3 mempty Nothing+        (records <&> (.body)) `shouldBe` [body1, body2, body3]++    it "has no span context outside of a span" do+        record <- only . snd =<< (runNoTracing . runInMemoryLogging) do info $ String "test"+        record.context `shouldBe` Nothing++    it "correlates with innermost span" do+        ((_, records), spans) <- runInMemoryTracing . runInMemoryLogging $ do+            inSpan "outer" Span.Kind.Internal mempty+                . inSpan "inner" Span.Kind.Internal mempty+                . info+                $ String "test"+        record <- only records+        case filter (("inner" ==) . (.name)) spans of+            [inner] -> (record.context <&> (.spanId)) `shouldBe` Just inner.context.spanId+            _ -> expectationFailure "expected exactly one inner span"++    it "different logs correlate with their enclosing spans" do+        ((_, records), spans) <- runInMemoryTracing . runInMemoryLogging $ do+            inSpan "span-a" Span.Kind.Internal mempty . info $ String "test1"+            inSpan "span-b" Span.Kind.Internal mempty . info $ String "test2"+        case (records, filter (("span-a" ==) . (.name)) spans, filter (("span-b" ==) . (.name)) spans) of+            ([recordA, recordB], [spanA], [spanB]) -> do+                (recordA.context <&> (.spanId)) `shouldBe` Just spanA.context.spanId+                (recordB.context <&> (.spanId)) `shouldBe` Just spanB.context.spanId+            _ -> expectationFailure "expected two records, one span-a and one span-b"++    it "concurrent log emissions are all collected" do+        let n = 50 :: Int+        (_, records) <-+            runConcurrent+                . runNoTracing+                . runInMemoryLogging+                . replicateConcurrently_ n+                . info+                $ String "test"+        length records `shouldBe` n++    it "concurrent logging within spans preserves context" do+        (_, records) <-+            runConcurrent+                . runNoTracing+                . runInMemoryLogging+                . inSpan "parent" Span.Kind.Internal mempty+                . replicateConcurrently_ 10+                . info+                $ String "test"+        length records `shouldBe` 10+        mapM_ ((`shouldSatisfy` isJust) . (.context)) records
+ test/Effectful/OpenTelemetry/Metrics/MeasurementSpec.hs view
@@ -0,0 +1,12 @@+module Effectful.OpenTelemetry.Metrics.MeasurementSpec (spec) where++import Arbitrary ()+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Metrics.Measurement (AggregationTemporality)+import Effectful.QuickCheck ((===))+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = prop "toEnum/fromEnum round trip" \(t :: AggregationTemporality) ->+    toEnum (fromEnum t) === t
+ test/Effectful/OpenTelemetry/MetricsSpec.hs view
@@ -0,0 +1,331 @@+{-# OPTIONS_GHC -Wno-type-defaults #-}++module Effectful.OpenTelemetry.MetricsSpec (spec) where++import Arbitrary+import Control.Monad (forM_)+import Data.Either (isLeft)+import Data.Foldable (for_)+import Data.IORef (atomicModifyIORef', modifyIORef', newIORef, readIORef, writeIORef)+import Data.List qualified as List+import Data.List.Extra (nubOrd)+import Data.Maybe (isJust, mapMaybe)+import Data.Text (Text)+import Effectful+import Effectful.Concurrent (runConcurrent, threadDelay)+import Effectful.Concurrent.Async (replicateConcurrently_)+import Effectful.Exception (SomeException, throwIO, try)+import Effectful.HUnit (HUnit, assertFailure)+import Effectful.Hspec+import Effectful.OpenTelemetry.Exporter qualified as Exporter+import Effectful.OpenTelemetry.Metrics+import Effectful.OpenTelemetry.Metrics.Counter qualified as Counter+import Effectful.OpenTelemetry.Metrics.Gauge qualified as Gauge+import Effectful.OpenTelemetry.Metrics.Histogram qualified as Histogram+import Effectful.OpenTelemetry.Metrics.Metadata+import Effectful.OpenTelemetry.Metrics.ObservableCounter qualified as ObservableCounter+import Effectful.OpenTelemetry.Metrics.ObservableGauge qualified as ObservableGauge+import Effectful.OpenTelemetry.Metrics.ObservableUpDownCounter qualified as ObservableUpDownCounter+import Effectful.OpenTelemetry.Metrics.RTS qualified as Metrics.RTS+import Effectful.OpenTelemetry.Metrics.Sum qualified as Sum+import Effectful.OpenTelemetry.Metrics.UpDownCounter qualified as UpDownCounter+import Effectful.OpenTelemetry.Timestamp (Timestamp (..))+import Effectful.QuickCheck (NonNegative (..), Small (..))+import Util+import Prelude++getLastMeasurement :: (HUnit :> es) => [Measurement] -> Eff es Measurement+getLastMeasurement [] = assertFailure "Expected at least 1 measurement"+getLastMeasurement ms = pure $ last ms++named :: Text -> [Measurement] -> Maybe Measurement+named n = List.find ((n ==) . (.name))++spec :: (IOE :> es, HUnit :> es, Hspec :> es) => Eff es ()+spec = runConcurrent . parallel $ do+    describe "Counter" do+        it "has initial value zero" do+            m <-+                only . snd =<< runInMemoryMetrics do+                    _ <- Counter.new "requests" mempty+                    pure ()+            m.name `shouldBe` "requests"+            numberValue m `shouldBe` Just 0++        prop "increments monotonically" \(NonNegative (Small a)) (NonNegative (Small b)) -> do+            m <-+                only . snd =<< runInMemoryMetrics do+                    c <- Counter.new "ops" mempty+                    Counter.add c a+                    Counter.add c b+            numberValue m `shouldBe` Just (a + b)+            Sum.Payload{isMonotonic} <- maybe (assertFailure "not a Sum") pure $ toMetric m+            isMonotonic `shouldBe` True++        prop "carries its metadata" \metadata -> do+            m <-+                only . snd =<< runInMemoryMetrics do+                    c <- Counter.new "documented" metadata+                    Counter.add c 1+            m.description `shouldBe` metadata.description+            m.unit `shouldBe` metadata.unit+            dp <- maybe (assertFailure "no number data point") pure $ numberDataPoint m+            dp.attributes `shouldBe` metadata.attributes++    describe "UpDownCounter" do+        it "has initial value zero" do+            m <-+                only . snd =<< runInMemoryMetrics do+                    _ <- UpDownCounter.new "requests" mempty+                    pure ()+            m.name `shouldBe` "requests"+            numberValue m `shouldBe` Just 0++        prop "increments non-monotonically" \(Small a) (Small b) -> do+            m <-+                only . snd =<< runInMemoryMetrics do+                    c <- UpDownCounter.new "ops" mempty+                    UpDownCounter.add c a+                    UpDownCounter.add c b+            numberValue m `shouldBe` Just (a + b)+            Sum.Payload{isMonotonic} <- maybe (assertFailure "not a Sum") pure $ toMetric m+            isMonotonic `shouldBe` False++        prop "carries its metadata" \metadata -> do+            m <-+                only . snd =<< runInMemoryMetrics do+                    c <- UpDownCounter.new "documented" metadata+                    UpDownCounter.add c 1+            m.description `shouldBe` metadata.description+            m.unit `shouldBe` metadata.unit+            dp <- maybe (assertFailure "no number data point") pure $ numberDataPoint m+            dp.attributes `shouldBe` metadata.attributes++    describe "data point attributes" do+        prop "carries its attributes" \attrs -> do+            m <-+                only . snd =<< runInMemoryMetrics do+                    c <-+                        Counter.new+                            "labelled"+                            Metadata{attributes = attrs, description = Nothing, unit = Nothing}+                    Counter.add c 1+            dp <- maybe (assertFailure "no number data point") pure $ numberDataPoint m+            dp.attributes `shouldBe` attrs++    describe "Gauge" do+        it "never set, never exported" do+            (_, measurements) <- runInMemoryMetrics do+                _ <- Gauge.new "temperature" mempty+                pure ()+            measurements `shouldSatisfy` null++        prop "last value wins" \(NonNegative a) (NonNegative b) (NonNegative c) -> do+            m <-+                getLastMeasurement . snd =<< runInMemoryMetrics do+                    g <- Gauge.new "memory" mempty+                    Gauge.set g a+                    Gauge.set g b+                    Gauge.set g c+            numberValue m `shouldBe` Just c++    describe "Multiple instruments" do+        prop "are independent" \(NonNegative (Small a)) (NonNegative v) -> do+            (_, measurements) <- runInMemoryMetrics do+                c <- UpDownCounter.new "req_count" mempty+                g <- Gauge.new "req_latency" mempty+                UpDownCounter.add c a+                Gauge.set g v+            length measurements `shouldBe` 2+            (numberValue =<< named "req_count" measurements) `shouldBe` Just a+            (numberValue =<< named "req_latency" measurements) `shouldBe` Just v++    describe "Concurrency" do+        prop "up-down counter increments are atomic" \ConcurrentObservation{repeats} (NonNegative (Small value)) -> do+            m <-+                only . snd =<< (runConcurrent . runInMemoryMetrics) do+                    c <- UpDownCounter.new "concurrent" mempty+                    replicateConcurrently_ repeats $ UpDownCounter.add c value+            numberValue m `shouldBe` Just (fromIntegral repeats * value)++        it "gauge sets converge" do+            m <-+                only . snd =<< (runConcurrent . runInMemoryMetrics) do+                    g <- Gauge.new "contended" mempty+                    replicateConcurrently_ 50 $ Gauge.set g 42+            numberValue m `shouldBe` Just 42++        it "instrument creation is safe" do+            m <-+                only . snd =<< (runConcurrent . runInMemoryMetrics) do+                    replicateConcurrently_ 20 do+                        c <- UpDownCounter.new "shared" mempty+                        UpDownCounter.add c 1+                    pure ()+            m.name `shouldBe` "shared"++    describe "Histogram" do+        it "starts empty" do+            m <-+                only . snd =<< runInMemoryMetrics do+                    _ <- Histogram.new "latency_ms" mempty+                    pure ()+            m.name `shouldBe` "latency_ms"+            dp <- maybe (assertFailure "no histogram data point") pure $ histogramDataPoint m+            dp.count `shouldBe` 0+            dp.sum `shouldBe` Nothing++        prop "bounds and stats match" \(HistogramBounds bounds) (HistogramObservations samples) -> do+            m <-+                only . snd =<< runInMemoryMetrics do+                    h <- Histogram.newWithBounds "weights" bounds mempty+                    mapM_ (Histogram.record h) samples+            dp <- maybe (assertFailure "no histogram data point") pure $ histogramDataPoint m+            dp.explicitBounds `shouldBe` List.sort bounds+            dp.count `shouldBe` fromIntegral (length samples)+            case dp.sum of+                Just s -> abs (s - Prelude.sum samples) `shouldSatisfy` (< 1e-9)+                Nothing -> expectationFailure "no sum reported"+            dp.minValue `shouldBe` Just (minimum samples)+            dp.maxValue `shouldBe` Just (maximum samples)++        prop "carries its attributes" \attributes -> do+            m <-+                only . snd =<< runInMemoryMetrics do+                    h <- Histogram.new "attributed" Metadata{attributes, description = Nothing, unit = Nothing}+                    mapM_ (Histogram.record h) [1, 2, 3]+            dp <- maybe (assertFailure "no histogram data point") pure $ histogramDataPoint m+            dp.attributes `shouldBe` attributes++        prop "gauges and histograms carry metadata too" \metadataG metadataH (NonNegative v) -> do+            (_, measurements) <- runInMemoryMetrics do+                g <- Gauge.new "temperature" metadataG+                Gauge.set g v+                _ <- Histogram.new "latency" metadataH+                pure ()+            length measurements `shouldBe` 2+            tempMeasurement <-+                maybe (assertFailure "temperature not found") pure $ named "temperature" measurements+            histMeasurement <- maybe (assertFailure "latency not found") pure $ named "latency" measurements+            tempMeasurement.description `shouldBe` metadataG.description+            tempMeasurement.unit `shouldBe` metadataG.unit+            histMeasurement.description `shouldBe` metadataH.description+            histMeasurement.unit `shouldBe` metadataH.unit++        it "start precedes observation" do+            (_, measurements) <- runInMemoryMetrics do+                c <- Counter.new "started_total" mempty+                Counter.add c 1+                g <- Gauge.new "started_gauge" mempty+                Gauge.set g 1+                h <- Histogram.new "started_hist" mempty+                Histogram.record h 1+            let numberWindow m = (\dp -> (dp.startTime, dp.time)) <$> numberDataPoint m+                histogramWindow m = (\dp -> (dp.startTime, dp.time)) <$> histogramDataPoint m+                windows :: [(Timestamp, Timestamp)]+                windows = mapMaybe numberWindow measurements <> mapMaybe histogramWindow measurements+            length windows `shouldBe` length measurements+            windows `shouldSatisfy` all \(start, time) -> start.nanos > 0 && start <= time++        prop "bucket placement" \BucketedObservations{..} -> do+            m <-+                only . snd =<< runInMemoryMetrics do+                    h <- Histogram.newWithBounds "buckets" bounds mempty+                    mapM_ (Histogram.record h) $ mconcat observations+            dp <- maybe (assertFailure "no histogram data point") pure $ histogramDataPoint m+            dp.bucketCounts `shouldBe` (fromIntegral . length <$> observations)++        prop "concurrent observations are atomic" \ConcurrentObservation{..} -> do+            m <-+                only . snd =<< (runConcurrent . runInMemoryMetrics) do+                    h <- Histogram.new "concurrent_obs" mempty+                    replicateConcurrently_ repeats $ Histogram.record h value+            dp <- maybe (assertFailure "no histogram data point") pure $ histogramDataPoint m+            dp.count `shouldBe` fromIntegral repeats++    describe "asynchronous instruments" do+        prop "reports under its own name" \(NonNegative v) -> do+            m <-+                only . snd =<< runInMemoryMetrics do+                    _ <- ObservableGauge.new "obs.ratio" mempty . pure $ Just v+                    pure ()+            numberValue m `shouldBe` Just v++        prop "each kind reports correctly" \(NonNegative (Small a)) (NonNegative (Small b)) (NonNegative c) -> do+            (_, measurements) <- runInMemoryMetrics do+                _ <- ObservableCounter.new "obs.total" mempty $ pure a+                _ <- ObservableUpDownCounter.new "obs.size" mempty $ pure b+                _ <- ObservableGauge.new "obs.ratio" mempty . pure $ Just c+                pure ()+            (numberValue =<< named "obs.total" measurements) `shouldBe` Just a+            (numberValue =<< named "obs.size" measurements) `shouldBe` Just b+            (numberValue =<< named "obs.ratio" measurements) `shouldBe` Just c+            Sum.Payload{isMonotonic = totalMonotonic} <-+                maybe (assertFailure "obs.total not found") pure $ toMetric =<< named "obs.total" measurements+            Sum.Payload{isMonotonic = sizeMonotonic} <-+                maybe (assertFailure "obs.size not found") pure $ toMetric =<< named "obs.size" measurements+            totalMonotonic `shouldBe` True+            sizeMonotonic `shouldBe` False++        it "observes lazily" do+            ref <- liftIO $ newIORef 1+            m <-+                only . snd =<< runInMemoryMetrics do+                    _ <- ObservableGauge.new "obs.late" mempty $ Just <$> readIORef ref+                    liftIO $ writeIORef ref 42+            numberValue m `shouldBe` Just 42++        it "one callback per export" do+            ref <- liftIO $ newIORef (0 :: Int)+            _ <- runInMemoryMetrics do+                _ <- ObservableGauge.new "obs.counted" mempty do+                    calls <- atomicModifyIORef' ref \n -> (n + 1, n + 1)+                    pure . Just $ fromIntegral calls+                pure ()+            calls <- liftIO $ readIORef ref+            calls `shouldBe` 1++        it "nothing observed, nothing reported" do+            (_, measurements) <- runInMemoryMetrics do+                _ <- ObservableGauge.new "obs.empty" mempty $ pure Nothing+                pure ()+            measurements `shouldSatisfy` null++    describe "runMetricsWith" do+        it "samples once, even if the action raises" do+            sunk <- liftIO $ newIORef []+            (result :: Either SomeException ()) <-+                try @SomeException . runMetricsWith Nothing (Exporter.fromIO \m -> modifyIORef' sunk (<> [m])) $ do+                    c <- Counter.new "with.ops" mempty+                    Counter.add c 5+                    liftIO (length <$> readIORef sunk) `shouldReturn` 0+                    throwIO $ userError "boom"+            result `shouldSatisfy` isLeft+            m <- only =<< liftIO (readIORef sunk)+            numberValue m `shouldBe` Just 5++        it "samples repeatedly when given an interval" do+            sunk <- liftIO $ newIORef []+            runMetricsWith (Just 20) (Exporter.fromIO \m -> modifyIORef' sunk (<> [m])) do+                c <- Counter.new "with.periodic" mempty+                forM_ [1 .. 5 :: Int] . const $ Counter.add c 1 >> threadDelay 50_000+            measurements <- liftIO (readIORef sunk)+            mapMaybe numberValue measurements `shouldBe` [1, 2, 3, 4, 5]++    describe "RTS metrics" do+        it "reports readings" do+            (_, measurements) <- runInMemoryMetrics Metrics.RTS.register+            measurements `shouldSatisfy` any (("rts.allocated_bytes" ==) . (.name))+            measurements `shouldSatisfy` any (("rts.gc.allocated_bytes" ==) . (.name))++        it "observed at the same time" do+            (_, measurements) <- runInMemoryMetrics Metrics.RTS.register+            let times = mapMaybe (fmap (.time) . numberDataPoint) measurements+            times `shouldSatisfy` (1 ==) . length . nubOrd++        it "readings have metadata" do+            (_, measurements) <- runInMemoryMetrics Metrics.RTS.register+            for_ measurements \m -> do+                m.name `shouldNotBe` ""+                m.description `shouldSatisfy` isJust+                m.unit `shouldSatisfy` isJust
+ test/Effectful/OpenTelemetry/Protocol/AnyValueSpec.hs view
@@ -0,0 +1,11 @@+module Effectful.OpenTelemetry.Protocol.AnyValueSpec (spec) where++import Arbitrary ()+import Data.Aeson (Result (..), fromJSON, toJSON)+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Protocol.AnyValue (AnyValue (..))+import Effectful.QuickCheck ((===))++spec :: (Hspec :> es) => Eff es ()+spec = prop "toJSON/fromJSON round trip" \v -> fromJSON @AnyValue (toJSON v) === Success v
+ test/Effectful/OpenTelemetry/Protocol/AttributesSpec.hs view
@@ -0,0 +1,21 @@+{- HLINT ignore "Monoid law, left identity" -}+{- HLINT ignore "Monoid law, right identity" -}++module Effectful.OpenTelemetry.Protocol.AttributesSpec (spec) where++import Arbitrary ()+import Data.Aeson qualified as Aeson+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Protocol.Attributes (Attributes)+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = parallel do+    prop "associativity" \(a :: Attributes) b c ->+        (a <> b) <> c `shouldBe` a <> (b <> c)+    prop "identity" \(a :: Attributes) -> do+        mempty <> a `shouldBe` a+        a <> mempty `shouldBe` a+    prop "toJSON/fromJSON round-trip" \(attrs :: Attributes) ->+        (Aeson.fromJSON . Aeson.toJSON) attrs `shouldBe` Aeson.Success attrs
+ test/Effectful/OpenTelemetry/Protocol/EffectSpec.hs view
@@ -0,0 +1,19 @@+module Effectful.OpenTelemetry.Protocol.EffectSpec (spec) where++import Control.Monad (forM_)+import Effectful+import Effectful.Concurrent (runConcurrent)+import Effectful.Concurrent.STM (atomically, flushTQueue, newTQueueIO)+import Effectful.Hspec+import Effectful.OpenTelemetry.Exporter.STM qualified as Exporter+import Effectful.OpenTelemetry.Protocol.Effect (export, runOTLPWith)+import Prelude++spec :: (IOE :> es, Hspec :> es) => Eff es ()+spec = runConcurrent do+    prop "exports each item individually in the correct order" \xs -> do+        q <- newTQueueIO+        runOTLPWith (Exporter.tqueue @Int q) $+            forM_ @[] xs \x -> do+                export x+                atomically (flushTQueue q) `shouldReturn` [x]
+ test/Effectful/OpenTelemetry/Protocol/ExceptionSpec.hs view
@@ -0,0 +1,32 @@+module Effectful.OpenTelemetry.Protocol.ExceptionSpec where++import Control.Exception (throwIO, try)+import Data.Either (isLeft)+import Effectful+import Effectful.Dispatch.Static (unsafeEff_)+import Effectful.Hspec (Hspec, it, shouldReturn, shouldSatisfy)+import Effectful.OpenTelemetry.Protocol.Environment (ConfigError (..))+import Effectful.OpenTelemetry.Protocol.Exception+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = do+    it "catch wrapped exception" $+        unsafeEff_+            ( try @ConfigError @()+                ( throwIO+                    . SomeOTLPException+                    . UnsupportedOTLPCompression+                    $ "foo"+                )+            )+            `shouldReturn` Left (UnsupportedOTLPCompression "foo")+    it "catch wrapper exception" $ do+        result <-+            unsafeEff_+                . try @SomeOTLPException @()+                . throwIO+                . SomeOTLPException+                . UnsupportedOTLPCompression+                $ "foo"+        result `shouldSatisfy` isLeft
+ test/Effectful/OpenTelemetry/Protocol/TransportSpec.hs view
@@ -0,0 +1,12 @@+module Effectful.OpenTelemetry.Protocol.TransportSpec (spec) where++import Arbitrary ()+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Protocol.Transport (Compression)+import Effectful.QuickCheck ((===))+import Text.Read (readMaybe)+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = prop "show/read round trip" \(c :: Compression) -> readMaybe (show c) === Just c
+ test/Effectful/OpenTelemetry/TimestampSpec.hs view
@@ -0,0 +1,11 @@+module Effectful.OpenTelemetry.TimestampSpec (spec) where++import Arbitrary ()+import Data.Aeson (Result (..), fromJSON, toJSON)+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Timestamp (Timestamp)+import Effectful.QuickCheck ((===))++spec :: (Hspec :> es) => Eff es ()+spec = prop "toJSON/fromJSON round trip" \(t :: Timestamp) -> fromJSON (toJSON t) === Success t
+ test/Effectful/OpenTelemetry/Tracing/PropagatorSpec.hs view
@@ -0,0 +1,57 @@+module Effectful.OpenTelemetry.Tracing.PropagatorSpec (spec) where++import Arbitrary ()+import Data.ByteString (ByteString)+import Data.ByteString qualified as ByteString+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Tracing.Propagator qualified as Propagator+import Effectful.OpenTelemetry.Tracing.Span.Context (Context)+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span.Context+import Effectful.OpenTelemetry.Tracing.Trace.Flags qualified as Trace.Flags+import Prelude++spec :: (IOE :> es, Hspec :> es) => Eff es ()+spec = do+    describe "w3cTraceContext" do+        prop "round-trips a context through headers, marking it remote" \(ctx :: Context) -> do+            headers <- Propagator.w3cTraceContext.inject (Just ctx) mempty+            Propagator.w3cTraceContext.extract headers Nothing >>= \case+                Nothing -> expectationFailure "no context extracted"+                Just extracted -> do+                    extracted.traceId `shouldBe` ctx.traceId+                    extracted.spanId `shouldBe` ctx.spanId+                    extracted.traceState `shouldBe` ctx.traceState+                    Trace.Flags.sampled extracted.traceFlags `shouldBe` Trace.Flags.sampled ctx.traceFlags+                    Trace.Flags.remote extracted.traceFlags `shouldBe` Trace.Flags.IsRemote++        it "keeps a child in its parent's trace" do+            parent <- liftIO $ Span.Context.new Nothing+            child <- liftIO . Span.Context.new $ Just parent+            child.traceId `shouldBe` parent.traceId+            child.spanId `shouldNotBe` parent.spanId++        it "round-trips the example specification header" do+            let traceparent :: ByteString+                traceparent = "00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01"+            case Span.Context.fromTraceparent traceparent of+                Nothing -> expectationFailure "the example specification header did not parse"+                Just ctx -> do+                    headers <- Propagator.w3cTraceContext.inject (Just ctx) mempty+                    lookup "traceparent" headers `shouldBe` Just traceparent++        prop "injects a header of the specified shape" \(ctx :: Context) ->+            case ByteString.split 0x2d (Span.Context.toTraceparent ctx) of+                [version, traceId, spanId, flags] ->+                    ByteString.length+                        <$> [version, traceId, spanId, flags]+                            `shouldBe` [2, 32, 16, 2]+                _ -> expectationFailure "a traceparent has four dash-separated fields"++        it "extracts nothing from headers that carry no context" do+            Propagator.w3cTraceContext.extract mempty Nothing `shouldReturn` Nothing++        it "leaves an existing context alone when none is inbound" do+            ctx <- liftIO $ Span.Context.new Nothing+            Propagator.w3cTraceContext.extract mempty (Just ctx)+                `shouldReturn` Just ctx
+ test/Effectful/OpenTelemetry/Tracing/Span/ContextSpec.hs view
@@ -0,0 +1,19 @@+module Effectful.OpenTelemetry.Tracing.Span.ContextSpec where++import Arbitrary ()+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Tracing.Span.Context+import Effectful.OpenTelemetry.Tracing.Trace.Flags qualified as Trace.Flags+import Effectful.OpenTelemetry.Tracing.Trace.StateSpec ()+import Effectful.QuickCheck ((===))+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = prop "toTraceparent/fromTraceparerent round trip" \c ->+    fromTraceparent (toTraceparent c)+        === Just+            c+                { traceState = mempty+                , traceFlags = c.traceFlags{Trace.Flags.remote = Trace.Flags.IsRemote}+                }
+ test/Effectful/OpenTelemetry/Tracing/Span/IDSpec.hs view
@@ -0,0 +1,14 @@+module Effectful.OpenTelemetry.Tracing.Span.IDSpec (spec) where++import Arbitrary ()+import Data.Aeson (Result (..), fromJSON, toJSON)+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Tracing.Span.ID (ID, fromBytes, toBytes)+import Effectful.QuickCheck ((===))+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = parallel do+    prop "toBytes/fromBytes round trip" \i -> fromBytes (toBytes i) === Right i+    prop "toJSON/fromJSON round trip" \i -> fromJSON @ID (toJSON i) === Success i
+ test/Effectful/OpenTelemetry/Tracing/Span/KindSpec.hs view
@@ -0,0 +1,11 @@+module Effectful.OpenTelemetry.Tracing.Span.KindSpec (spec) where++import Arbitrary ()+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Tracing.Span.Kind (Kind)+import Effectful.QuickCheck ((===))+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = prop "fromEnum/toEnum round trip" \(k :: Kind) -> toEnum (fromEnum k) === k
+ test/Effectful/OpenTelemetry/Tracing/Trace/FlagsSpec.hs view
@@ -0,0 +1,28 @@+{-# OPTIONS_GHC -Wno-orphans #-}+{-# OPTIONS_GHC -Wno-term-variable-capture #-}++{- HLINT ignore "Monoid law, left identity" -}+{- HLINT ignore "Monoid law, right identity" -}++module Effectful.OpenTelemetry.Tracing.Trace.FlagsSpec where++import Arbitrary ()+import Data.Aeson (Result (..), fromJSON, toJSON)+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Tracing.Trace.Flags+import Effectful.OpenTelemetry.Tracing.Trace.Flags qualified as Trace+import Effectful.OpenTelemetry.Tracing.Trace.StateSpec ()+import Effectful.QuickCheck ((.&&.), (===))+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = parallel do+    prop "associativity" \(a :: Trace.Flags) b c ->+        (a <> b) <> c === a <> (b <> c)+    prop "identity" \(a :: Trace.Flags) ->+        (mempty <> a === a) .&&. (a <> mempty === a)+    prop "toWord32/fromWord32 round trip" \flags ->+        fromWord32 (toWord32 flags) === flags+    prop "toJSON/fromJSON round trip" \(flags :: Trace.Flags) ->+        fromJSON (toJSON flags) === Success flags
+ test/Effectful/OpenTelemetry/Tracing/Trace/IDSpec.hs view
@@ -0,0 +1,14 @@+module Effectful.OpenTelemetry.Tracing.Trace.IDSpec (spec) where++import Arbitrary ()+import Data.Aeson (Result (..), fromJSON, toJSON)+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Tracing.Trace.ID (ID, fromBytes, toBytes)+import Effectful.QuickCheck ((===))+import Prelude++spec :: (Hspec :> es) => Eff es ()+spec = parallel do+    prop "toBytes/fromBytes round trip" \i -> fromBytes (toBytes i) === Right i+    prop "toJSON/fromJSON round trip" \i -> fromJSON @ID (toJSON i) === Success i
+ test/Effectful/OpenTelemetry/Tracing/Trace/StateSpec.hs view
@@ -0,0 +1,59 @@+{-# OPTIONS_GHC -Wno-orphans #-}+{-# OPTIONS_GHC -Wno-term-variable-capture #-}++{- HLINT ignore "Monoid law, left identity" -}+{- HLINT ignore "Monoid law, right identity" -}++module Effectful.OpenTelemetry.Tracing.Trace.StateSpec where++import Arbitrary ()+import Data.List qualified as List+import Data.List.Extra qualified as List+import Effectful+import Effectful.Hspec+import Effectful.OpenTelemetry.Tracing.Trace.State+import Effectful.OpenTelemetry.Tracing.Trace.State qualified as State+import Effectful.QuickCheck+import Prelude++isValid :: State -> Property es+isValid (toList -> s) = List.nubOrdOn fst s === s .&&. List.length s <= 32++spec :: (Hspec :> es) => Eff es ()+spec = parallel do+    prop "construction" isValid+    prop "associativity" \(a :: State) b c -> (a <> b) <> c === a <> (b <> c)+    prop "identity" \(a :: State) -> (mempty <> a === a) .&&. (a <> mempty === a)+    prop "toList/fromList round trip" \s -> fromList (toList s) === s+    prop "toText/fromText round trip" \s -> fromText (toText s) === s+    prop "toByteString/fromByteString round trip" \s -> fromByteString (toByteString s) === s+    prop "concatenation with self" \s -> isValid (s <> s)+    prop "concatenation" \s1 s2 -> isValid (s1 <> s2)+    prop "lookup" \s (NonNegative (Small i)) ->+        let+            l = toList s+            (k, v) = l !! i+         in+            length l > i ==> State.lookup k s === Just v+    prop "insert" \s (k, v) ->+        let+            s' = insert k v s+         in+            counterexample (show s') $+                isValid s' .&&. case toList s' of+                    (k', v') : _ -> (k, v) === (k', v')+                    _ -> property False+    prop "delete" \s (NonNegative (Small i)) ->+        let+            l = toList s+            (k, _) = l !! i+            s' = State.delete k s+            l' = toList s'+            (a, b) = splitAt i l+         in+            (length l > i) ==> l' === a <> drop 1 b+    prop "left-bias" $ \k v1 v2 ->+        let left = insert k v1 mempty+            right = insert k v2 mempty+            combined = left <> right+         in State.lookup k combined === Just v1
+ test/Effectful/OpenTelemetry/TracingSpec.hs view
@@ -0,0 +1,413 @@+{- HLINT ignore "Use head" -}+{-# LANGUAGE OverloadedLists #-}++module Effectful.OpenTelemetry.TracingSpec where++import Arbitrary (Marker (..))+import Control.Monad (forM_)+import Data.Either (isLeft)+import Data.Functor ((<&>))+import Data.Text qualified as Text+import Effectful+import Effectful.Concurrent (runConcurrent)+import Effectful.Concurrent.Async (concurrently_, replicateConcurrently_)+import Effectful.Error.Static (runErrorNoCallStack, throwError)+import Effectful.Exception (AssertionFailed (..), throwIO, try)+import Effectful.HUnit (HUnit)+import Effectful.Hspec+import Effectful.OpenTelemetry.Tracing+import Effectful.OpenTelemetry.Tracing.Span.Context qualified as Span.Context+import Effectful.OpenTelemetry.Tracing.Span.Event (Event (..))+import Effectful.OpenTelemetry.Tracing.Span.Kind qualified as Span.Kind+import Effectful.OpenTelemetry.Tracing.Span.Status qualified as Span.Status+import Util+import Prelude++spec :: (IOE :> es, HUnit :> es, Hspec :> es) => Eff es ()+spec = runConcurrent . parallel $ do+    describe "inSpan" do+        prop "exports a span" \(Marker name) spanKind a -> do+            s <- only . snd =<< (runInMemoryTracing . inSpan name spanKind a . pure) ()+            s.name `shouldBe` name+            s.kind `shouldBe` spanKind+            s.attributes `shouldBe` a+            s.endTime `shouldSatisfy` (>= Just s.startTime)++        it "no parent, unique ids" do+            (_, spans) <- runInMemoryTracing do+                inSpan "a" Span.Kind.Internal mempty $ pure ()+                inSpan "b" Span.Kind.Internal mempty $ pure ()+                inSpan "c" Span.Kind.Internal mempty $ pure ()+            mapM_ ((`shouldBe` Nothing) . (.parentSpanId)) spans+            (spans <&> (.context.spanId)) `shouldSatisfy` allUnique++        describe "nesting" do+            it "parent-child chain" do+                (_, spans) <- runInMemoryTracing do+                    inSpan "root" Span.Kind.Internal mempty+                        . inSpan "middle" Span.Kind.Internal mempty+                        . inSpan "leaf" Span.Kind.Internal mempty+                        $ pure ()+                length spans `shouldBe` 3+                let leaf = spans !! 0+                    middle = spans !! 1+                    root = spans !! 2+                leaf.parentSpanId `shouldBe` Just middle.context.spanId+                middle.parentSpanId `shouldBe` Just root.context.spanId+                root.parentSpanId `shouldBe` Nothing++            it "shared parent" do+                (_, spans) <- runInMemoryTracing do+                    inSpan "parent" Span.Kind.Internal mempty do+                        inSpan "child-a" Span.Kind.Internal mempty $ pure ()+                        inSpan "child-b" Span.Kind.Internal mempty $ pure ()+                length spans `shouldBe` 3+                let childA = spans !! 0+                    childB = spans !! 1+                    parent = spans !! 2+                childA.parentSpanId `shouldBe` Just parent.context.spanId+                childB.parentSpanId `shouldBe` Just parent.context.spanId++            it "restores parent context" do+                (ctx, _) <- runInMemoryTracing do+                    inSpan "outer" Span.Kind.Internal mempty do+                        inSpan "inner" Span.Kind.Internal mempty $ pure ()+                        currentContext+                ctx `shouldSatisfy` \case+                    Just _ -> True+                    Nothing -> False++        describe "concurrency" do+            it "no loss, unique ids" do+                let n = 500 :: Int+                (_, spans) <-+                    runInMemoryTracing+                        . runConcurrent+                        . replicateConcurrently_ n+                        . inSpan "parallel" Span.Kind.Internal mempty+                        $ pure ()+                length spans `shouldBe` n+                (spans <&> (.context.spanId)) `shouldSatisfy` allUnique++            it "concurrent nesting" do+                (_, spans) <-+                    runInMemoryTracing+                        . runConcurrent+                        . inSpan "root" Span.Kind.Internal mempty+                        $ concurrently_+                            (inSpan "branch-a" Span.Kind.Internal mempty $ pure ())+                            (inSpan "branch-b" Span.Kind.Internal mempty $ pure ())+                length spans `shouldBe` 3+                case filter (("root" ==) . (.name)) spans of+                    [root] -> do+                        let branches = filter ((/= "root") . (.name)) spans+                        mapM_ ((`shouldBe` Just root.context.spanId) . (.parentSpanId)) branches+                    _ -> expectationFailure "expected exactly one root span"++        describe "exception handling" do+            prop "Error status on exception" \(Marker msg) -> do+                (result, spans) <-+                    runInMemoryTracing+                        . try @AssertionFailed+                        . inSpan @_ @() "item" Span.Kind.Internal mempty+                        . throwIO+                        $ AssertionFailed (Text.unpack msg)+                result `shouldSatisfy` isLeft+                length spans `shouldBe` 1+                forM_ spans \s -> do+                    s.status.code `shouldBe` Span.Status.Error+                    s.status.message `shouldSatisfy` Text.isPrefixOf msg+                    length s.events `shouldBe` 1+                    fmap (.name) s.events `shouldBe` ["exception"]++        describe "events" do+            prop "records events in order" \(Marker n1) a1 (Marker n2) a2 -> do+                s <-+                    only . snd =<< runInMemoryTracing do+                        inSpan "test-span" Span.Kind.Internal mempty do+                            addEventNow n1 a1+                            addEventNow n2 a2+                            pure ()+                fmap (\e -> (e.name, e.attributes)) s.events `shouldBe` [(n1, a1), (n2, a2)]++            it "events go to innermost span" do+                (_, spans) <- runInMemoryTracing do+                    inSpan "parent" Span.Kind.Internal mempty+                        . inSpan "child" Span.Kind.Internal mempty+                        $ do+                            addEventNow "event" mempty+                            pure ()+                length spans `shouldBe` 2+                let child = spans !! 0+                fmap (.name) child.events `shouldBe` ["event"]++    describe "inErrorSpan" do+        prop "exports a span" \(Marker name) spanKind a -> do+            s <-+                only . snd+                    =<< ( runInMemoryTracing+                            . runErrorNoCallStack @AssertionFailed+                            . inErrorSpan @AssertionFailed @_ @_ name spanKind a+                            . pure+                        )+                        ()+            s.name `shouldBe` name+            s.kind `shouldBe` spanKind+            s.attributes `shouldBe` a+            s.endTime `shouldSatisfy` (>= Just s.startTime)++        it "no parent, unique ids" do+            (_, spans) <- runInMemoryTracing . runErrorNoCallStack @AssertionFailed $ do+                inErrorSpan @AssertionFailed @_ @() "a" Span.Kind.Internal mempty $ pure ()+                inErrorSpan @AssertionFailed @_ @() "b" Span.Kind.Internal mempty $ pure ()+                inErrorSpan @AssertionFailed @_ @() "c" Span.Kind.Internal mempty $ pure ()+            mapM_ ((`shouldBe` Nothing) . (.parentSpanId)) spans+            (spans <&> (.context.spanId)) `shouldSatisfy` allUnique++        describe "nesting" do+            it "parent-child chain" do+                (_, spans) <- runInMemoryTracing . runErrorNoCallStack @AssertionFailed $ do+                    inErrorSpan @AssertionFailed @_ @() "root" Span.Kind.Internal mempty+                        . inErrorSpan @AssertionFailed @_ @() "middle" Span.Kind.Internal mempty+                        . inErrorSpan @AssertionFailed @_ @() "leaf" Span.Kind.Internal mempty+                        $ pure ()+                length spans `shouldBe` 3+                let leaf = spans !! 0+                    middle = spans !! 1+                    root = spans !! 2+                leaf.parentSpanId `shouldBe` Just middle.context.spanId+                middle.parentSpanId `shouldBe` Just root.context.spanId+                root.parentSpanId `shouldBe` Nothing++            it "shared parent" do+                (_, spans) <- runInMemoryTracing . runErrorNoCallStack @AssertionFailed $ do+                    inErrorSpan @AssertionFailed @_ @() "parent" Span.Kind.Internal mempty do+                        inErrorSpan @AssertionFailed @_ @() "child-a" Span.Kind.Internal mempty $ pure ()+                        inErrorSpan @AssertionFailed @_ @() "child-b" Span.Kind.Internal mempty $ pure ()+                length spans `shouldBe` 3+                let childA = spans !! 0+                    childB = spans !! 1+                    parent = spans !! 2+                childA.parentSpanId `shouldBe` Just parent.context.spanId+                childB.parentSpanId `shouldBe` Just parent.context.spanId++            it "restores parent context" do+                (ctx, _) <- runInMemoryTracing . runErrorNoCallStack @AssertionFailed $ do+                    inErrorSpan @AssertionFailed @_ @_ "outer" Span.Kind.Internal mempty do+                        inErrorSpan @AssertionFailed @_ @() "inner" Span.Kind.Internal mempty $ pure ()+                        currentContext+                ctx `shouldSatisfy` \case+                    Right (Just _) -> True+                    _ -> False++        describe "concurrency" do+            it "no loss, unique ids" do+                let n = 500 :: Int+                (_, spans) <-+                    runInMemoryTracing+                        . runErrorNoCallStack @AssertionFailed+                        . runConcurrent+                        . replicateConcurrently_ n+                        . inErrorSpan @AssertionFailed @_ @() "parallel" Span.Kind.Internal mempty+                        $ pure ()+                length spans `shouldBe` n+                (spans <&> (.context.spanId)) `shouldSatisfy` allUnique++            it "concurrent nesting" do+                (_, spans) <-+                    runInMemoryTracing+                        . runErrorNoCallStack @AssertionFailed+                        . runConcurrent+                        . inErrorSpan @AssertionFailed @_ @() "root" Span.Kind.Internal mempty+                        $ concurrently_+                            (inSpan "branch-a" Span.Kind.Internal mempty $ pure ())+                            (inSpan "branch-b" Span.Kind.Internal mempty $ pure ())+                length spans `shouldBe` 3+                case filter (("root" ==) . (.name)) spans of+                    [root] -> do+                        let branches = filter ((/= "root") . (.name)) spans+                        mapM_ ((`shouldBe` Just root.context.spanId) . (.parentSpanId)) branches+                    _ -> expectationFailure "expected exactly one root span"++        describe "exception handling" do+            prop "Error status on exception" \(Marker msg) -> do+                (result, spans) <-+                    runInMemoryTracing+                        . runErrorNoCallStack @AssertionFailed+                        . inErrorSpan+                            @AssertionFailed+                            @_+                            @()+                            "item"+                            Span.Kind.Internal+                            mempty+                        . throwError+                        $ AssertionFailed (Text.unpack msg)+                result `shouldSatisfy` isLeft+                length spans `shouldBe` 1+                forM_ spans \s -> do+                    s.status.code `shouldBe` Span.Status.Error+                    s.status.message `shouldBe` msg+                    length s.events `shouldBe` 1+                    fmap (.name) s.events `shouldBe` ["exception"]++            it "Ok status without exception" do+                (_, spans) <-+                    runInMemoryTracing+                        . runErrorNoCallStack @AssertionFailed+                        . inErrorSpan+                            @AssertionFailed+                            @_+                            @()+                            "item"+                            Span.Kind.Internal+                            mempty+                        $ pure ()+                length spans `shouldBe` 1+                forM_ (spans <&> (.status)) \s -> do+                    s.code `shouldBe` Span.Status.Unset+                    s.message `shouldBe` mempty++        describe "events" do+            prop "records events in order" \(Marker n1) a1 (Marker n2) a2 -> do+                s <-+                    only . snd =<< runInMemoryTracing do+                        runErrorNoCallStack @AssertionFailed+                            . inErrorSpan @AssertionFailed @_ @() "test-span" Span.Kind.Internal mempty+                            $ do+                                addEventNow n1 a1+                                addEventNow n2 a2+                                pure ()+                fmap (\e -> (e.name, e.attributes)) s.events `shouldBe` [(n1, a1), (n2, a2)]++            it "events go to innermost span" do+                (_, spans) <-+                    runInMemoryTracing+                        . runErrorNoCallStack @AssertionFailed+                        . inErrorSpan @AssertionFailed @_ @() "parent" Span.Kind.Internal mempty+                        . inErrorSpan @AssertionFailed @_ @() "child" Span.Kind.Internal mempty+                        $ do+                            addEventNow "event" mempty+                            pure ()+                length spans `shouldBe` 2+                let child = spans !! 0+                fmap (.name) child.events `shouldBe` ["event"]++    describe "inOkSpan" do+        prop "exports a span" \(Marker name) spanKind a -> do+            s <- only . snd =<< (runInMemoryTracing . inOkSpan name spanKind a . pure) ()+            s.name `shouldBe` name+            s.kind `shouldBe` spanKind+            s.attributes `shouldBe` a+            s.endTime `shouldSatisfy` (>= Just s.startTime)++        it "no parent, unique ids" do+            (_, spans) <- runInMemoryTracing do+                inOkSpan "a" Span.Kind.Internal mempty $ pure ()+                inOkSpan "b" Span.Kind.Internal mempty $ pure ()+                inOkSpan "c" Span.Kind.Internal mempty $ pure ()+            mapM_ ((`shouldBe` Nothing) . (.parentSpanId)) spans+            (spans <&> (.context.spanId)) `shouldSatisfy` allUnique++        describe "nesting" do+            it "parent-child chain" do+                (_, spans) <- runInMemoryTracing do+                    inOkSpan "root" Span.Kind.Internal mempty+                        . inOkSpan "middle" Span.Kind.Internal mempty+                        . inOkSpan "leaf" Span.Kind.Internal mempty+                        $ pure ()+                length spans `shouldBe` 3+                let leaf = spans !! 0+                    middle = spans !! 1+                    root = spans !! 2+                leaf.parentSpanId `shouldBe` Just middle.context.spanId+                middle.parentSpanId `shouldBe` Just root.context.spanId+                root.parentSpanId `shouldBe` Nothing++            it "shared parent" do+                (_, spans) <- runInMemoryTracing do+                    inOkSpan "parent" Span.Kind.Internal mempty do+                        inOkSpan "child-a" Span.Kind.Internal mempty $ pure ()+                        inOkSpan "child-b" Span.Kind.Internal mempty $ pure ()+                length spans `shouldBe` 3+                let childA = spans !! 0+                    childB = spans !! 1+                    parent = spans !! 2+                childA.parentSpanId `shouldBe` Just parent.context.spanId+                childB.parentSpanId `shouldBe` Just parent.context.spanId++            it "restores parent context" do+                (ctx, _) <- runInMemoryTracing do+                    inOkSpan "outer" Span.Kind.Internal mempty do+                        inOkSpan "inner" Span.Kind.Internal mempty $ pure ()+                        currentContext+                ctx `shouldSatisfy` \case+                    Just _ -> True+                    Nothing -> False++        describe "concurrency" do+            it "no loss, unique ids" do+                let n = 500 :: Int+                (_, spans) <-+                    runInMemoryTracing+                        . runConcurrent+                        . replicateConcurrently_ n+                        . inOkSpan "parallel" Span.Kind.Internal mempty+                        $ pure ()+                length spans `shouldBe` n+                (spans <&> (.context.spanId)) `shouldSatisfy` allUnique++            it "concurrent nesting" do+                (_, spans) <-+                    runInMemoryTracing+                        . runConcurrent+                        . inOkSpan "root" Span.Kind.Internal mempty+                        $ concurrently_+                            (inOkSpan "branch-a" Span.Kind.Internal mempty $ pure ())+                            (inOkSpan "branch-b" Span.Kind.Internal mempty $ pure ())+                length spans `shouldBe` 3+                case filter (("root" ==) . (.name)) spans of+                    [root] -> do+                        let branches = filter ((/= "root") . (.name)) spans+                        mapM_ ((`shouldBe` Just root.context.spanId) . (.parentSpanId)) branches+                    _ -> expectationFailure "expected exactly one root span"++        describe "exception handling" do+            it "no span on exception" do+                (_, spans) <-+                    runInMemoryTracing+                        . try @AssertionFailed+                        . inOkSpan "item" Span.Kind.Internal mempty+                        . throwIO+                        $ AssertionFailed "bla"+                length spans `shouldBe` 0++            it "Ok status without exception" do+                (_, spans) <-+                    runInMemoryTracing . inOkSpan "item" Span.Kind.Internal mempty $+                        pure ()+                length spans `shouldBe` 1+                forM_ (spans <&> (.status)) \s -> do+                    s.code `shouldBe` Span.Status.Ok+                    s.message `shouldBe` mempty++        describe "events" do+            prop "records events in order" \(Marker n1) a1 (Marker n2) a2 -> do+                s <-+                    only . snd =<< runInMemoryTracing do+                        inOkSpan "test-span" Span.Kind.Internal mempty do+                            addEventNow n1 a1+                            addEventNow n2 a2+                            pure ()+                fmap (\e -> (e.name, e.attributes)) s.events `shouldBe` [(n1, a1), (n2, a2)]++            it "events go to innermost span" do+                (_, spans) <- runInMemoryTracing do+                    inSpan "parent" Span.Kind.Internal mempty+                        . inOkSpan "child" Span.Kind.Internal mempty+                        $ do+                            addEventNow "event" mempty+                            pure ()+                length spans `shouldBe` 2+                let child = spans !! 0+                fmap (.name) child.events `shouldBe` ["event"]
+ test/Main.hs view
@@ -0,0 +1,1 @@+{-# OPTIONS_GHC -F -pgmF hspec-effectful-discover #-}
+ test/Util.hs view
@@ -0,0 +1,50 @@+module Util where++import Control.Applicative ((<|>))+import Data.Scientific (Scientific)+import Data.String (IsString (..))+import Effectful+import Effectful.HUnit+import Effectful.OpenTelemetry.Metrics.Gauge qualified as Gauge+import Effectful.OpenTelemetry.Metrics.Histogram qualified as Histogram+import Effectful.OpenTelemetry.Metrics.Measurement+    ( HistogramDataPoint+    , Measurement+    , NumberDataPoint (..)+    , toMetric+    )+import Effectful.OpenTelemetry.Metrics.Sum qualified as Sum+import Effectful.OpenTelemetry.Protocol.Transport+import GHC.Stack (HasCallStack)+import Prelude++protocolLabel :: (IsString s) => Protocol -> s+protocolLabel (HTTP Json _) = "HTTP/JSON"+protocolLabel (HTTP Proto _) = "HTTP/Protobuf"+protocolLabel GRPC{} = "gRPC"++only :: (HasCallStack, HUnit :> es, Show a) => [a] -> Eff es a+only [x] = pure x+only xs = assertFailure $ "expected exactly one item, got: " <> show xs++allUnique :: (Eq a) => [a] -> Bool+allUnique [] = True+allUnique (x : xs) = x `notElem` xs && allUnique xs++numberDataPoint :: Measurement -> Maybe NumberDataPoint+numberDataPoint m = fromSum <|> fromGauge+  where+    fromSum = do+        Sum.Payload{dataPoints = [dp]} <- toMetric m+        pure dp+    fromGauge = do+        Gauge.Payload{dataPoints = [dp]} <- toMetric m+        pure dp++numberValue :: Measurement -> Maybe Scientific+numberValue m = (.value) <$> numberDataPoint m++histogramDataPoint :: Measurement -> Maybe HistogramDataPoint+histogramDataPoint m = do+    Histogram.Payload{dataPoints = [dp]} <- toMetric m+    pure dp