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 +12/−0
- LICENCE +287/−0
- Setup.hs +2/−0
- otel-effectful.cabal +210/−0
- src/Effectful/OpenTelemetry.hs +173/−0
- src/Effectful/OpenTelemetry/Exporter.hs +12/−0
- src/Effectful/OpenTelemetry/Exporter/Console.hs +49/−0
- src/Effectful/OpenTelemetry/Exporter/Environment.hs +54/−0
- src/Effectful/OpenTelemetry/Exporter/OTLP.hs +295/−0
- src/Effectful/OpenTelemetry/Exporter/STM.hs +21/−0
- src/Effectful/OpenTelemetry/Exporter/Type.hs +25/−0
- src/Effectful/OpenTelemetry/Logging.hs +10/−0
- src/Effectful/OpenTelemetry/Logging/Effect.hs +196/−0
- src/Effectful/OpenTelemetry/Logging/LogRecord.hs +128/−0
- src/Effectful/OpenTelemetry/Logging/Severity.hs +123/−0
- src/Effectful/OpenTelemetry/Metrics.hs +24/−0
- src/Effectful/OpenTelemetry/Metrics/Counter.hs +77/−0
- src/Effectful/OpenTelemetry/Metrics/Effect.hs +253/−0
- src/Effectful/OpenTelemetry/Metrics/Gauge.hs +94/−0
- src/Effectful/OpenTelemetry/Metrics/Histogram.hs +176/−0
- src/Effectful/OpenTelemetry/Metrics/Instrument.hs +29/−0
- src/Effectful/OpenTelemetry/Metrics/Measurement.hs +236/−0
- src/Effectful/OpenTelemetry/Metrics/Metadata.hs +34/−0
- src/Effectful/OpenTelemetry/Metrics/ObservableCounter.hs +57/−0
- src/Effectful/OpenTelemetry/Metrics/ObservableGauge.hs +57/−0
- src/Effectful/OpenTelemetry/Metrics/ObservableUpDownCounter.hs +57/−0
- src/Effectful/OpenTelemetry/Metrics/RTS.hs +242/−0
- src/Effectful/OpenTelemetry/Metrics/Sum.hs +60/−0
- src/Effectful/OpenTelemetry/Metrics/UpDownCounter.hs +64/−0
- src/Effectful/OpenTelemetry/Protocol.hs +18/−0
- src/Effectful/OpenTelemetry/Protocol/AnyValue.hs +71/−0
- src/Effectful/OpenTelemetry/Protocol/Attributes.hs +70/−0
- src/Effectful/OpenTelemetry/Protocol/Effect.hs +46/−0
- src/Effectful/OpenTelemetry/Protocol/Environment.hs +249/−0
- src/Effectful/OpenTelemetry/Protocol/Exception.hs +61/−0
- src/Effectful/OpenTelemetry/Protocol/Export.hs +63/−0
- src/Effectful/OpenTelemetry/Protocol/GRPC.hs +94/−0
- src/Effectful/OpenTelemetry/Protocol/HTTP.hs +55/−0
- src/Effectful/OpenTelemetry/Protocol/Resource.hs +50/−0
- src/Effectful/OpenTelemetry/Protocol/Scope.hs +47/−0
- src/Effectful/OpenTelemetry/Protocol/Transport.hs +42/−0
- src/Effectful/OpenTelemetry/Timestamp.hs +42/−0
- src/Effectful/OpenTelemetry/Tracing.hs +10/−0
- src/Effectful/OpenTelemetry/Tracing/Effect.hs +386/−0
- src/Effectful/OpenTelemetry/Tracing/Propagator.hs +78/−0
- src/Effectful/OpenTelemetry/Tracing/Span.hs +132/−0
- src/Effectful/OpenTelemetry/Tracing/Span/Context.hs +89/−0
- src/Effectful/OpenTelemetry/Tracing/Span/Event.hs +42/−0
- src/Effectful/OpenTelemetry/Tracing/Span/ID.hs +61/−0
- src/Effectful/OpenTelemetry/Tracing/Span/Kind.hs +59/−0
- src/Effectful/OpenTelemetry/Tracing/Span/Status.hs +81/−0
- src/Effectful/OpenTelemetry/Tracing/Trace/Flags.hs +101/−0
- src/Effectful/OpenTelemetry/Tracing/Trace/ID.hs +61/−0
- src/Effectful/OpenTelemetry/Tracing/Trace/State.hs +152/−0
- src/Prettyprinter/Extra.hs +39/−0
- src/Proto3/Wire/Encode/Class.hs +70/−0
- src/System/Random/Extra.hs +34/−0
- test/Arbitrary.hs +262/−0
- test/Effectful/OpenTelemetry/Exporter/DeadSpec.hs +64/−0
- test/Effectful/OpenTelemetry/Exporter/Grafana/Loki.hs +119/−0
- test/Effectful/OpenTelemetry/Exporter/Grafana/Mimir.hs +212/−0
- test/Effectful/OpenTelemetry/Exporter/Grafana/Polling.hs +100/−0
- test/Effectful/OpenTelemetry/Exporter/Grafana/Tempo.hs +241/−0
- test/Effectful/OpenTelemetry/Exporter/GrafanaSpec.hs +318/−0
- test/Effectful/OpenTelemetry/ExporterSpec.hs +36/−0
- test/Effectful/OpenTelemetry/Logging/SeveritySpec.hs +18/−0
- test/Effectful/OpenTelemetry/LoggingSpec.hs +92/−0
- test/Effectful/OpenTelemetry/Metrics/MeasurementSpec.hs +12/−0
- test/Effectful/OpenTelemetry/MetricsSpec.hs +331/−0
- test/Effectful/OpenTelemetry/Protocol/AnyValueSpec.hs +11/−0
- test/Effectful/OpenTelemetry/Protocol/AttributesSpec.hs +21/−0
- test/Effectful/OpenTelemetry/Protocol/EffectSpec.hs +19/−0
- test/Effectful/OpenTelemetry/Protocol/ExceptionSpec.hs +32/−0
- test/Effectful/OpenTelemetry/Protocol/TransportSpec.hs +12/−0
- test/Effectful/OpenTelemetry/TimestampSpec.hs +11/−0
- test/Effectful/OpenTelemetry/Tracing/PropagatorSpec.hs +57/−0
- test/Effectful/OpenTelemetry/Tracing/Span/ContextSpec.hs +19/−0
- test/Effectful/OpenTelemetry/Tracing/Span/IDSpec.hs +14/−0
- test/Effectful/OpenTelemetry/Tracing/Span/KindSpec.hs +11/−0
- test/Effectful/OpenTelemetry/Tracing/Trace/FlagsSpec.hs +28/−0
- test/Effectful/OpenTelemetry/Tracing/Trace/IDSpec.hs +14/−0
- test/Effectful/OpenTelemetry/Tracing/Trace/StateSpec.hs +59/−0
- test/Effectful/OpenTelemetry/TracingSpec.hs +413/−0
- test/Main.hs +1/−0
- test/Util.hs +50/−0
+ 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