From ad50f2487b81e71d1cf0eea59c1154c8fc224535 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sun, 30 Aug 2026 19:13:35 -0400 Subject: [PATCH] feat(observability): preserve cross-client trace continuity Signed-off-by: Yordis Prieto --- Cargo.lock | 129 ++++- .../trogon/eventstore/client-spans.yaml | 27 + .../registry/rust/observability.rs.j2 | 1 + .../templates/registry/rust/weaver.yaml | 13 +- trogon-eventstore/Cargo.toml | 1 + trogon-eventstore/src/batch.rs | 6 +- trogon-eventstore/src/commands.rs | 22 +- trogon-eventstore/src/observability.rs | 490 ++++++++++++++++- .../src/observability/generated.rs | 8 + trogon-eventstore/tests/compatibility.rs | 494 ++++++++++++++++++ 10 files changed, 1170 insertions(+), 21 deletions(-) create mode 100644 trogon-eventstore/tests/compatibility.rs diff --git a/Cargo.lock b/Cargo.lock index 8c321b4..00b022d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1352,6 +1352,36 @@ dependencies = [ "thiserror", ] +[[package]] +name = "opentelemetry-otlp" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9966929966d17620d7c316c643ba62631826e10021409357772d5eea84f62c35" +dependencies = [ + "http", + "opentelemetry", + "opentelemetry-proto", + "opentelemetry_sdk", + "prost 0.14.4", + "thiserror", + "tokio", + "tonic 0.14.6", + "tonic-types", +] + +[[package]] +name = "opentelemetry-proto" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "56d658ba1faf63f7b9c492cfbe6e0ec365440a16132d3270c1065f7b33f1b638" +dependencies = [ + "opentelemetry", + "opentelemetry_sdk", + "prost 0.14.4", + "tonic 0.14.6", + "tonic-prost", +] + [[package]] name = "opentelemetry_sdk" version = "0.32.1" @@ -1537,7 +1567,17 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2796faa41db3ec313a31f7624d9286acf277b52de526150b7e69f3debf891ee5" dependencies = [ "bytes", - "prost-derive", + "prost-derive 0.13.5", +] + +[[package]] +name = "prost" +version = "0.14.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "528ac67416ff8646872a3c02cad9cc4ee5dc9f9540c9b10771855c95cb2e5ae1" +dependencies = [ + "bytes", + "prost-derive 0.14.4", ] [[package]] @@ -1553,8 +1593,8 @@ dependencies = [ "once_cell", "petgraph", "prettyplease", - "prost", - "prost-types", + "prost 0.13.5", + "prost-types 0.13.5", "regex", "syn 2.0.119", "tempfile", @@ -1573,13 +1613,35 @@ dependencies = [ "syn 2.0.119", ] +[[package]] +name = "prost-derive" +version = "0.14.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf" +dependencies = [ + "anyhow", + "itertools", + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "prost-types" version = "0.13.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "52c2c1bf36ddb1a1c396b3601a3cec27c2462e45f07c386894ec3ccf5332bd16" dependencies = [ - "prost", + "prost 0.13.5", +] + +[[package]] +name = "prost-types" +version = "0.14.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f94967dc7688f3054c7fac87473ffae4cc4c3904800e2d9f5b857246d8963b0a" +dependencies = [ + "prost 0.14.4", ] [[package]] @@ -2533,7 +2595,7 @@ dependencies = [ "hyper-util", "percent-encoding", "pin-project", - "prost", + "prost 0.13.5", "rustls-native-certs", "socket2 0.5.10", "tokio", @@ -2545,6 +2607,32 @@ dependencies = [ "tracing", ] +[[package]] +name = "tonic" +version = "0.14.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ac2a5518c70fa84342385732db33fb3f44bc4cc748936eb5833d2df34d6445ef" +dependencies = [ + "async-trait", + "base64", + "bytes", + "http", + "http-body", + "http-body-util", + "hyper", + "hyper-timeout", + "hyper-util", + "percent-encoding", + "pin-project", + "sync_wrapper", + "tokio", + "tokio-stream", + "tower", + "tower-layer", + "tower-service", + "tracing", +] + [[package]] name = "tonic-build" version = "0.13.1" @@ -2554,11 +2642,33 @@ dependencies = [ "prettyplease", "proc-macro2", "prost-build", - "prost-types", + "prost-types 0.13.5", "quote", "syn 2.0.119", ] +[[package]] +name = "tonic-prost" +version = "0.14.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "50849f68853be452acf590cde0b146665b8d507b3b8af17261df47e02c209ea0" +dependencies = [ + "bytes", + "prost 0.14.4", + "tonic 0.14.6", +] + +[[package]] +name = "tonic-types" +version = "0.14.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73ab1b02061f83d519bba3caa167f88f261ef05720ab8ebc954ade70de3348e8" +dependencies = [ + "prost 0.14.4", + "prost-types 0.14.4", + "tonic 0.14.6", +] + [[package]] name = "tower" version = "0.5.3" @@ -2690,9 +2800,10 @@ dependencies = [ "names", "nom", "opentelemetry", + "opentelemetry-otlp", "opentelemetry_sdk", - "prost", - "prost-types", + "prost 0.13.5", + "prost-types 0.13.5", "rand 0.9.5", "reqwest", "rustls", @@ -2704,7 +2815,7 @@ dependencies = [ "tokio", "tokio-rustls", "toml", - "tonic", + "tonic 0.13.1", "tonic-build", "tower", "tracing", diff --git a/otel/semconv/registry/trogon/eventstore/client-spans.yaml b/otel/semconv/registry/trogon/eventstore/client-spans.yaml index 9354d5a..6c8abd5 100644 --- a/otel/semconv/registry/trogon/eventstore/client-spans.yaml +++ b/otel/semconv/registry/trogon/eventstore/client-spans.yaml @@ -3,6 +3,11 @@ groups: type: attribute_group brief: TrogonEventStore client attributes. attributes: + - id: trogon.eventstore.event.type + type: string + stability: development + brief: The type of event received by the client. + examples: [order-created] - id: trogon.eventstore.batch.correlation_id type: string stability: development @@ -28,6 +33,28 @@ groups: requirement_level: conditionally_required: If and only if the operation failed. + - id: span.trogon.eventstore.client.receive + type: span + span_kind: client + stability: development + brief: Receives an event delivered by a subscription. + attributes: + - ref: messaging.system + requirement_level: required + - ref: messaging.operation.name + requirement_level: required + - ref: messaging.operation.type + requirement_level: required + - ref: messaging.destination.name + requirement_level: required + - ref: messaging.message.id + requirement_level: recommended + - ref: messaging.consumer.group.name + requirement_level: + conditionally_required: If the event was delivered by a persistent subscription. + - ref: trogon.eventstore.event.type + requirement_level: recommended + - id: span.trogon.eventstore.client.append_to_stream type: span extends: span.trogon.eventstore.client diff --git a/otel/semconv/templates/registry/rust/observability.rs.j2 b/otel/semconv/templates/registry/rust/observability.rs.j2 index bb70c25..430ea24 100644 --- a/otel/semconv/templates/registry/rust/observability.rs.j2 +++ b/otel/semconv/templates/registry/rust/observability.rs.j2 @@ -8,6 +8,7 @@ pub(crate) const {{ attribute.name | screaming_snake_case }}: &str = "{{ attribu {% endfor %} pub(crate) const CLIENT_SPAN_KIND: SpanKind = SpanKind::{{ ctx.base.span_kind | pascal_case }}; +pub(crate) const RECEIVE_SPAN_KIND: SpanKind = SpanKind::{{ ctx.receive.span_kind | pascal_case }}; pub(crate) mod operation { use super::ClientOperation; diff --git a/otel/semconv/templates/registry/rust/weaver.yaml b/otel/semconv/templates/registry/rust/weaver.yaml index 48e9288..48e167b 100644 --- a/otel/semconv/templates/registry/rust/weaver.yaml +++ b/otel/semconv/templates/registry/rust/weaver.yaml @@ -7,12 +7,23 @@ templates: filter: > { base: [semconv_grouped_spans[].spans[] | select(.id == "span.trogon.eventstore.client")][0], - operations: [semconv_grouped_spans[].spans[] | select(.id | startswith("span.trogon.eventstore.client."))], + receive: [semconv_grouped_spans[].spans[] | select(.id == "span.trogon.eventstore.client.receive")][0], + operations: [semconv_grouped_spans[].spans[] | select( + (.id | startswith("span.trogon.eventstore.client.")) and + .annotations.code_generation.operation_name != null + )], attributes: [semconv_grouped_spans[].spans[].attributes[] | select( .name == "db.collection.name" or .name == "db.operation.name" or .name == "db.system.name" or .name == "error.type" or + .name == "messaging.consumer.group.name" or + .name == "messaging.destination.name" or + .name == "messaging.message.id" or + .name == "messaging.operation.name" or + .name == "messaging.operation.type" or + .name == "messaging.system" or + .name == "trogon.eventstore.event.type" or .name == "trogon.eventstore.batch.correlation_id" )] | unique_by(.name) } diff --git a/trogon-eventstore/Cargo.toml b/trogon-eventstore/Cargo.toml index 77befe6..3cf4895 100755 --- a/trogon-eventstore/Cargo.toml +++ b/trogon-eventstore/Cargo.toml @@ -69,6 +69,7 @@ name = "integration" [dev-dependencies] names = "0.14" +opentelemetry-otlp = { version = "0.32", default-features = false, features = ["grpc-tonic", "trace"] } opentelemetry_sdk = { version = "0.32", default-features = false, features = ["testing", "trace"] } serde = { version = "1", features = ["derive"] } testcontainers = "0.23" diff --git a/trogon-eventstore/src/batch.rs b/trogon-eventstore/src/batch.rs index 16c1c32..7edec9a 100644 --- a/trogon-eventstore/src/batch.rs +++ b/trogon-eventstore/src/batch.rs @@ -21,6 +21,7 @@ pub(crate) struct Req { pub(crate) stream_name: String, pub(crate) events: Vec, pub(crate) expected_revision: StreamState, + pub(crate) context: Context, } impl Req { @@ -30,6 +31,7 @@ impl Req { stream_name, events, expected_revision, + context: Context::current(), } } } @@ -200,7 +202,7 @@ mod tests { DB_COLLECTION_NAME, DB_OPERATION_NAME, TROGON_EVENTSTORE_BATCH_CORRELATION_ID, }; use opentelemetry::global; - use opentelemetry::trace::{Status, noop::NoopTracerProvider}; + use opentelemetry::trace::{Status, TraceContextExt, noop::NoopTracerProvider}; use opentelemetry_sdk::trace::{InMemorySpanExporter, SdkTracerProvider}; #[tokio::test] @@ -223,6 +225,7 @@ mod tests { panic!("expected an inbound batch request"); }; let correlation_id = req.id.to_string(); + let captured_span_id = req.context.span().span_context().span_id(); sender .send(Ok(BatchWriteResult::new( "stream".to_string(), @@ -239,6 +242,7 @@ mod tests { .iter() .find(|span| span.name == "batch_append_to_stream stream") .expect("batch append client span"); + assert_eq!(span.span_context.span_id(), captured_span_id); assert!(span.attributes.iter().any(|attribute| { attribute.key.as_str() == TROGON_EVENTSTORE_BATCH_CORRELATION_ID && attribute.value.to_string() == correlation_id diff --git a/trogon-eventstore/src/commands.rs b/trogon-eventstore/src/commands.rs index c05ae31..35c48c6 100644 --- a/trogon-eventstore/src/commands.rs +++ b/trogon-eventstore/src/commands.rs @@ -41,9 +41,11 @@ use crate::{ fn convert_event_data_to_batch_proposed_message( event: EventData, + context: &opentelemetry::Context, ) -> streams::batch_append_req::ProposedMessage { use streams::batch_append_req; + let event = crate::observability::inject_event_context(event, context); let id = event.id_opt.unwrap_or_else(uuid::Uuid::new_v4).into(); let custom_metadata = event.custom_metadata.unwrap_or_default(); @@ -148,6 +150,7 @@ pub async fn append_to_stream( use streams::AppendReq; use streams::append_req::{self, Content}; + let context = opentelemetry::Context::current(); let stream_identifier = Some(StreamIdentifier { stream_name: stream.into_stream_name(), }); @@ -163,7 +166,7 @@ pub async fn append_to_stream( yield header; for event in events { - yield event.into(); + yield crate::observability::inject_event_context(event, &context).into(); } }; @@ -233,10 +236,11 @@ pub async fn batch_append( let expected_stream_position = Some(expected_stream_position); + let context = &req.context; let proposed_messages: Vec = req .events .into_iter() - .map(convert_event_data_to_batch_proposed_message) + .map(|event| convert_event_data_to_batch_proposed_message(event, context)) .collect(); let deadline = common_operation_options @@ -759,6 +763,7 @@ impl Subscription { use streams::read_req::options::stream_options::RevisionOption; use streams::read_req::options::{self, StreamOption}; + let receive = crate::observability::SubscriptionReceive::start(); loop { if let Some(mut stream) = self.stream.take() { match stream.try_next().await { @@ -806,6 +811,7 @@ impl Subscription { } } + receive.complete(None, &event); return Ok(SubscriptionEvent::EventAppeared(event)); } @@ -1321,6 +1327,7 @@ pub async fn subscribe_to_persistent_subscription>( use persistent::read_req::{self, Options, options::StreamOption}; let handle = connection.current_selected_node().await?; + let group_name = group_name.as_ref().to_owned(); if to_all && !handle.supports_feature(Features::PERSISTENT_SUBSCRIPITON_TO_ALL) { return Err(crate::Error::UnsupportedFeature); @@ -1349,7 +1356,7 @@ pub async fn subscribe_to_persistent_subscription>( let req_options = Options { stream_option, - group_name: group_name.as_ref().to_string(), + group_name: group_name.clone(), buffer_size: options.buffer_size as i32, uuid_option: Some(uuid_option), }; @@ -1376,6 +1383,7 @@ pub async fn subscribe_to_persistent_subscription>( ack_sender: sender, channel_id, inner: resp.into_inner(), + group_name, }), } } @@ -1385,10 +1393,12 @@ pub struct PersistentSubscription { ack_sender: mpsc::Sender, channel_id: uuid::Uuid, inner: tonic::Streaming, + group_name: String, } impl PersistentSubscription { pub async fn next_subscription_event(&mut self) -> crate::Result { + let receive = crate::observability::SubscriptionReceive::start(); match self.inner.try_next().await { Err(status) => { if let Some("persistent-subscription-dropped") = status @@ -1408,7 +1418,11 @@ impl PersistentSubscription { } Ok(resp) => { if let Some(content) = resp.and_then(|r| r.content) { - return Ok(content.into()); + let event = content.into(); + if let PersistentSubscriptionEvent::EventAppeared { event, .. } = &event { + receive.complete(Some(&self.group_name), event); + } + return Ok(event); } unreachable!() diff --git a/trogon-eventstore/src/observability.rs b/trogon-eventstore/src/observability.rs index 72da08a..2c19940 100644 --- a/trogon-eventstore/src/observability.rs +++ b/trogon-eventstore/src/observability.rs @@ -1,17 +1,29 @@ mod generated; -use generated::CLIENT_SPAN_KIND; +use bytes::Bytes; +use generated::{CLIENT_SPAN_KIND, RECEIVE_SPAN_KIND}; pub(crate) use generated::{ DB_COLLECTION_NAME, DB_OPERATION_NAME, DB_SYSTEM_NAME, ERROR_TYPE, - TROGON_EVENTSTORE_BATCH_CORRELATION_ID, + MESSAGING_CONSUMER_GROUP_NAME, MESSAGING_DESTINATION_NAME, MESSAGING_MESSAGE_ID, + MESSAGING_OPERATION_NAME, MESSAGING_OPERATION_TYPE, MESSAGING_SYSTEM, + TROGON_EVENTSTORE_BATCH_CORRELATION_ID, TROGON_EVENTSTORE_EVENT_TYPE, }; -use opentelemetry::trace::{FutureExt, Status, TraceContextExt, Tracer}; +use opentelemetry::propagation::{Extractor, Injector}; +use opentelemetry::trace::{FutureExt, Link, Span as _, Status, TraceContextExt, Tracer}; use opentelemetry::{Context, InstrumentationScope, KeyValue, global}; +use serde_json::{Map, Value}; use std::borrow::Cow; use std::future::Future; +use std::time::SystemTime; + +use crate::{EventData, ResolvedEvent}; pub(crate) use generated::operation; +const TRACE_PARENT: &str = "traceparent"; +const TRACE_STATE: &str = "tracestate"; +const RECEIVE_OPERATION: &str = "receive"; + #[cfg(test)] pub(crate) static TEST_GLOBALS: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); @@ -106,6 +118,149 @@ fn start_client_operation(operation: impl Into) -> Context { Context::current_with_span(span) } +struct EventMetadataCarrier(Map); + +impl Injector for EventMetadataCarrier { + fn set(&mut self, key: &str, value: String) { + if !is_persisted_trace_field(key) { + return; + } + + self.0 + .retain(|existing, _| !existing.eq_ignore_ascii_case(key)); + self.0.insert(key.to_owned(), Value::String(value)); + } +} + +fn is_persisted_trace_field(key: &str) -> bool { + key.eq_ignore_ascii_case(TRACE_PARENT) || key.eq_ignore_ascii_case(TRACE_STATE) +} + +impl Extractor for EventMetadataCarrier { + fn get(&self, key: &str) -> Option<&str> { + self.0.iter().find_map(|(existing, value)| { + existing + .eq_ignore_ascii_case(key) + .then(|| value.as_str()) + .flatten() + }) + } + + fn keys(&self) -> Vec<&str> { + self.0.keys().map(String::as_str).collect() + } +} + +pub(crate) fn inject_event_context(mut event: EventData, context: &Context) -> EventData { + if !context.span().span_context().is_valid() { + return event; + } + + let mut propagation_fields = EventMetadataCarrier(Map::new()); + global::get_text_map_propagator(|propagator| { + propagator.inject_context(context, &mut propagation_fields) + }); + if propagation_fields.0.is_empty() { + return event; + } + + let metadata = match event.custom_metadata.as_deref() { + None | Some([]) => Map::new(), + Some(metadata) => match serde_json::from_slice::(metadata) { + Ok(Value::Object(metadata)) => metadata, + _ => return event, + }, + }; + + let mut carrier = EventMetadataCarrier(metadata); + for (key, value) in propagation_fields.0 { + if let Value::String(value) = value { + carrier.set(&key, value); + } + } + + if let Ok(metadata) = serde_json::to_vec(&Value::Object(carrier.0)) { + event.custom_metadata = Some(Bytes::from(metadata)); + } + + event +} + +fn extract_event_context(metadata: &[u8]) -> Context { + if metadata.is_empty() { + return Context::new(); + } + + let Ok(Value::Object(metadata)) = serde_json::from_slice(metadata) else { + return Context::new(); + }; + + let carrier = EventMetadataCarrier(metadata); + global::get_text_map_propagator(|propagator| propagator.extract(&carrier)) +} + +pub(crate) struct SubscriptionReceive { + parent: Context, + started_at: SystemTime, +} + +impl SubscriptionReceive { + pub(crate) fn start() -> Self { + Self { + parent: Context::current(), + started_at: SystemTime::now(), + } + } + + pub(crate) fn complete( + self, + consumer_group_name: Option<&str>, + resolved_event: &ResolvedEvent, + ) { + let Some(delivered_event) = resolved_event + .event + .as_ref() + .or(resolved_event.link.as_ref()) + else { + return; + }; + let original_event = resolved_event.get_original_event(); + let destination = original_event.stream_id(); + let mut attributes = vec![ + KeyValue::new(MESSAGING_SYSTEM, "trogoneventstore"), + KeyValue::new(MESSAGING_OPERATION_NAME, RECEIVE_OPERATION), + KeyValue::new(MESSAGING_OPERATION_TYPE, RECEIVE_OPERATION), + KeyValue::new(MESSAGING_DESTINATION_NAME, destination.to_owned()), + KeyValue::new(MESSAGING_MESSAGE_ID, original_event.id.to_string()), + KeyValue::new( + TROGON_EVENTSTORE_EVENT_TYPE, + delivered_event.event_type.clone(), + ), + ]; + if let Some(consumer_group_name) = consumer_group_name { + attributes.push(KeyValue::new( + MESSAGING_CONSUMER_GROUP_NAME, + consumer_group_name.to_owned(), + )); + } + + let message_context = extract_event_context(&delivered_event.custom_metadata); + let message_span_context = message_context.span().span_context().clone(); + let tracer = global::tracer_with_scope(instrumentation_scope()); + let mut builder = tracer + .span_builder(format!("receive {destination}")) + .with_kind(RECEIVE_SPAN_KIND) + .with_start_time(self.started_at) + .with_attributes(attributes); + if message_span_context.is_valid() { + builder = builder.with_links(vec![Link::with_context(message_span_context)]); + } + + let mut span = builder.start_with_context(&tracer, &self.parent); + span.end(); + } +} + fn error_type(error: &crate::Error) -> &'static str { match error { crate::Error::ServerError(_) => "server_error", @@ -167,12 +322,26 @@ where mod tests { use super::{ DB_COLLECTION_NAME, DB_OPERATION_NAME, DB_SYSTEM_NAME, ERROR_TYPE, INSTRUMENTATION_SCOPE, - client_operation, operation::APPEND_TO_STREAM, + MESSAGING_CONSUMER_GROUP_NAME, MESSAGING_DESTINATION_NAME, MESSAGING_MESSAGE_ID, + MESSAGING_OPERATION_NAME, MESSAGING_OPERATION_TYPE, MESSAGING_SYSTEM, SubscriptionReceive, + TROGON_EVENTSTORE_EVENT_TYPE, client_operation, extract_event_context, + inject_event_context, operation::APPEND_TO_STREAM, }; - use opentelemetry::Context; + use bytes::Bytes; + use chrono::Utc; + use opentelemetry::baggage::BaggageExt; use opentelemetry::global; - use opentelemetry::trace::{SpanKind, Status, TraceContextExt, noop::NoopTracerProvider}; + use opentelemetry::propagation::TextMapCompositePropagator; + use opentelemetry::trace::{ + SpanContext, SpanId, SpanKind, Status, TraceContextExt, TraceFlags, TraceId, TraceState, + noop::NoopTracerProvider, + }; + use opentelemetry::{Context, KeyValue}; + use opentelemetry_sdk::propagation::{BaggagePropagator, TraceContextPropagator}; use opentelemetry_sdk::trace::{InMemorySpanExporter, SdkTracerProvider}; + use serde_json::json; + + use crate::{EventData, Position, RecordedEvent, ResolvedEvent}; #[tokio::test] async fn client_operation_emits_a_semantic_database_client_span() { @@ -247,6 +416,315 @@ mod tests { global::set_tracer_provider(NoopTracerProvider::new()); } + #[tokio::test] + async fn event_context_round_trips_without_replacing_custom_metadata() { + let _guard = super::TEST_GLOBALS.lock().await; + global::set_text_map_propagator(TextMapCompositePropagator::new(vec![ + Box::new(TraceContextPropagator::new()), + Box::new(BaggagePropagator::new()), + ])); + let context = remote_context().with_baggage([KeyValue::new("tenant", "secret")]); + let event = EventData::json("order-created", &json!({"orderId": "42"})) + .unwrap() + .metadata_as_json(&json!({ + "tenant": "north", + "TraceParent": "00-00000000000000000000000000000001-0000000000000001-00" + })) + .unwrap(); + + let event = inject_event_context(event, &context); + let metadata: serde_json::Value = + serde_json::from_slice(event.custom_metadata.as_ref().unwrap()).unwrap(); + let object = metadata.as_object().unwrap(); + + assert_eq!(object.get("tenant"), Some(&json!("north"))); + assert!(!object.keys().any(|key| key.eq_ignore_ascii_case("baggage"))); + assert_eq!( + object + .keys() + .filter(|key| key.eq_ignore_ascii_case("traceparent")) + .count(), + 1 + ); + let extracted = extract_event_context(event.custom_metadata.as_ref().unwrap()); + assert_eq!( + extracted.span().span_context().trace_id(), + context.span().span_context().trace_id() + ); + assert_eq!( + extracted.span().span_context().span_id(), + context.span().span_context().span_id() + ); + + reset_propagator(); + } + + #[tokio::test] + async fn event_context_injection_preserves_unsupported_metadata_and_empty_contexts() { + let _guard = super::TEST_GLOBALS.lock().await; + global::set_text_map_propagator(TraceContextPropagator::new()); + for unsupported in [ + b"not-json".as_slice(), + b"[]".as_slice(), + br#""scalar""#.as_slice(), + b"42".as_slice(), + ] { + let original = Bytes::copy_from_slice(unsupported); + let event = EventData::binary("binary", Bytes::new()).metadata(original.clone()); + let event = inject_event_context(event, &remote_context()); + assert_eq!(event.custom_metadata.as_ref(), Some(&original)); + } + + let original = Bytes::from_static(br#"{"tenant":"north"}"#); + let event = EventData::binary("binary", Bytes::new()).metadata(original.clone()); + let event = inject_event_context(event, &Context::new()); + assert_eq!(event.custom_metadata.as_ref(), Some(&original)); + + reset_propagator(); + } + + #[tokio::test] + async fn subscription_event_emits_a_receive_span_with_ambient_parent_and_message_link() { + let _guard = super::TEST_GLOBALS.lock().await; + global::set_text_map_propagator(TraceContextPropagator::new()); + let exporter = InMemorySpanExporter::default(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + global::set_tracer_provider(provider.clone()); + let message_context = remote_context(); + let ambient_context = alternate_remote_context(); + let event = inject_event_context( + EventData::json("order-created", &json!({"orderId": "42"})).unwrap(), + &message_context, + ); + let recorded = recorded_event(event.custom_metadata.unwrap()); + let message_id = recorded.id.to_string(); + let resolved = ResolvedEvent { + event: Some(recorded), + link: None, + commit_position: None, + }; + + let _ambient = ambient_context.clone().attach(); + let receive = SubscriptionReceive::start(); + tokio::time::sleep(std::time::Duration::from_millis(10)).await; + receive.complete(Some("billing"), &resolved); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let span = spans + .iter() + .find(|span| span.name == "receive orders") + .expect("receive span"); + assert_eq!(span.span_kind, SpanKind::Client); + assert!( + span.end_time.duration_since(span.start_time).unwrap() + >= std::time::Duration::from_millis(10) + ); + assert_eq!( + span.parent_span_id, + ambient_context.span().span_context().span_id() + ); + assert_eq!(span.links.links.len(), 1); + assert_eq!( + span.links.links[0].span_context, + message_context.span().span_context().clone() + ); + assert_eq!(attribute(span, MESSAGING_SYSTEM), Some("trogoneventstore")); + assert_eq!(attribute(span, MESSAGING_OPERATION_NAME), Some("receive")); + assert_eq!(attribute(span, MESSAGING_OPERATION_TYPE), Some("receive")); + assert_eq!(attribute(span, MESSAGING_DESTINATION_NAME), Some("orders")); + assert_eq!( + attribute(span, MESSAGING_MESSAGE_ID), + Some(message_id.as_str()) + ); + assert_eq!( + attribute(span, MESSAGING_CONSUMER_GROUP_NAME), + Some("billing") + ); + assert_eq!( + attribute(span, TROGON_EVENTSTORE_EVENT_TYPE), + Some("order-created") + ); + + global::set_tracer_provider(NoopTracerProvider::new()); + reset_propagator(); + } + + #[tokio::test] + async fn subscription_event_emits_receive_telemetry_without_persisted_context() { + let _guard = super::TEST_GLOBALS.lock().await; + global::set_text_map_propagator(TraceContextPropagator::new()); + let exporter = InMemorySpanExporter::default(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + global::set_tracer_provider(provider.clone()); + let ambient_context = alternate_remote_context(); + let resolved = ResolvedEvent { + event: Some(recorded_event(Bytes::new())), + link: None, + commit_position: None, + }; + + let _ambient = ambient_context.clone().attach(); + SubscriptionReceive::start().complete(None, &resolved); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let span = spans + .iter() + .find(|span| span.name == "receive orders") + .expect("receive span"); + assert_eq!(span.span_kind, SpanKind::Client); + assert_eq!( + span.parent_span_id, + ambient_context.span().span_context().span_id() + ); + assert!(span.links.links.is_empty()); + + global::set_tracer_provider(NoopTracerProvider::new()); + reset_propagator(); + } + + #[tokio::test] + async fn resolved_subscription_event_uses_delivered_event_type_with_original_identity() { + let _guard = super::TEST_GLOBALS.lock().await; + global::set_text_map_propagator(TraceContextPropagator::new()); + let exporter = InMemorySpanExporter::default(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + global::set_tracer_provider(provider.clone()); + let message_context = remote_context(); + let event = inject_event_context( + EventData::json("order-created", &json!({"orderId": "42"})).unwrap(), + &message_context, + ); + let mut link = recorded_event(Bytes::new()); + link.stream_id_raw = Bytes::from_static(b"$ce-orders"); + link.event_type = "$>".to_owned(); + let link_id = link.id.to_string(); + let resolved = ResolvedEvent { + event: Some(recorded_event(event.custom_metadata.unwrap())), + link: Some(link), + commit_position: None, + }; + + SubscriptionReceive::start().complete(None, &resolved); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let span = spans + .iter() + .find(|span| span.name == "receive $ce-orders") + .expect("receive span"); + assert_eq!( + attribute(span, MESSAGING_DESTINATION_NAME), + Some("$ce-orders") + ); + assert_eq!( + attribute(span, MESSAGING_MESSAGE_ID), + Some(link_id.as_str()) + ); + assert_eq!( + attribute(span, TROGON_EVENTSTORE_EVENT_TYPE), + Some("order-created") + ); + + global::set_tracer_provider(NoopTracerProvider::new()); + reset_propagator(); + } + + #[tokio::test] + async fn link_only_subscription_event_emits_receive_telemetry_from_the_link() { + let _guard = super::TEST_GLOBALS.lock().await; + global::set_text_map_propagator(TraceContextPropagator::new()); + let exporter = InMemorySpanExporter::default(); + let provider = SdkTracerProvider::builder() + .with_simple_exporter(exporter.clone()) + .build(); + global::set_tracer_provider(provider.clone()); + let message_context = remote_context(); + let event = inject_event_context( + EventData::json("$>", &json!({"link": "0@orders"})).unwrap(), + &message_context, + ); + let mut link = recorded_event(event.custom_metadata.unwrap()); + link.stream_id_raw = Bytes::from_static(b"$ce-orders"); + link.event_type = "$>".to_owned(); + let link_id = link.id.to_string(); + let resolved = ResolvedEvent { + event: None, + link: Some(link), + commit_position: None, + }; + + SubscriptionReceive::start().complete(None, &resolved); + provider.force_flush().unwrap(); + + let spans = exporter.get_finished_spans().unwrap(); + let span = spans + .iter() + .find(|span| span.name == "receive $ce-orders") + .expect("receive span"); + assert_eq!( + attribute(span, MESSAGING_DESTINATION_NAME), + Some("$ce-orders") + ); + assert_eq!( + attribute(span, MESSAGING_MESSAGE_ID), + Some(link_id.as_str()) + ); + assert_eq!(attribute(span, TROGON_EVENTSTORE_EVENT_TYPE), Some("$>")); + assert_eq!(span.links.links.len(), 1); + assert_eq!( + span.links.links[0].span_context, + message_context.span().span_context().clone() + ); + + global::set_tracer_provider(NoopTracerProvider::new()); + reset_propagator(); + } + + fn remote_context() -> Context { + remote_context_with("58406520a006649127e371903a2de979", "58406520a0066491") + } + + fn alternate_remote_context() -> Context { + remote_context_with("68406520a006649127e371903a2de978", "68406520a0066492") + } + + fn remote_context_with(trace_id: &str, span_id: &str) -> Context { + Context::new().with_remote_span_context(SpanContext::new( + TraceId::from_hex(trace_id).unwrap(), + SpanId::from_hex(span_id).unwrap(), + TraceFlags::SAMPLED, + true, + TraceState::default(), + )) + } + + fn recorded_event(custom_metadata: Bytes) -> RecordedEvent { + RecordedEvent { + stream_id_raw: Bytes::from_static(b"orders"), + id: uuid::Uuid::new_v4(), + revision: 0, + event_type: "order-created".to_owned(), + data: Bytes::new(), + metadata: Default::default(), + custom_metadata, + is_json: true, + position: Position::start(), + created: Utc::now(), + } + } + + fn reset_propagator() { + global::set_text_map_propagator(TextMapCompositePropagator::new(Vec::new())); + } + fn attribute<'a>(span: &'a opentelemetry_sdk::trace::SpanData, key: &str) -> Option<&'a str> { span.attributes .iter() diff --git a/trogon-eventstore/src/observability/generated.rs b/trogon-eventstore/src/observability/generated.rs index a19b71e..2295ade 100644 --- a/trogon-eventstore/src/observability/generated.rs +++ b/trogon-eventstore/src/observability/generated.rs @@ -7,10 +7,18 @@ pub(crate) const DB_COLLECTION_NAME: &str = "db.collection.name"; pub(crate) const DB_OPERATION_NAME: &str = "db.operation.name"; pub(crate) const DB_SYSTEM_NAME: &str = "db.system.name"; pub(crate) const ERROR_TYPE: &str = "error.type"; +pub(crate) const MESSAGING_CONSUMER_GROUP_NAME: &str = "messaging.consumer.group.name"; +pub(crate) const MESSAGING_DESTINATION_NAME: &str = "messaging.destination.name"; +pub(crate) const MESSAGING_MESSAGE_ID: &str = "messaging.message.id"; +pub(crate) const MESSAGING_OPERATION_NAME: &str = "messaging.operation.name"; +pub(crate) const MESSAGING_OPERATION_TYPE: &str = "messaging.operation.type"; +pub(crate) const MESSAGING_SYSTEM: &str = "messaging.system"; pub(crate) const TROGON_EVENTSTORE_BATCH_CORRELATION_ID: &str = "trogon.eventstore.batch.correlation_id"; +pub(crate) const TROGON_EVENTSTORE_EVENT_TYPE: &str = "trogon.eventstore.event.type"; pub(crate) const CLIENT_SPAN_KIND: SpanKind = SpanKind::Client; +pub(crate) const RECEIVE_SPAN_KIND: SpanKind = SpanKind::Client; pub(crate) mod operation { use super::ClientOperation; diff --git a/trogon-eventstore/tests/compatibility.rs b/trogon-eventstore/tests/compatibility.rs new file mode 100644 index 0000000..1b44d29 --- /dev/null +++ b/trogon-eventstore/tests/compatibility.rs @@ -0,0 +1,494 @@ +use std::io::Write; +use std::path::{Path, PathBuf}; +use std::time::Duration; + +use opentelemetry::propagation::TextMapCompositePropagator; +use opentelemetry::trace::{FutureExt, TraceContextExt, Tracer}; +use opentelemetry::{Context, global}; +use opentelemetry_otlp::WithExportConfig; +use opentelemetry_sdk::Resource; +use opentelemetry_sdk::propagation::{BaggagePropagator, TraceContextPropagator}; +use opentelemetry_sdk::trace::SdkTracerProvider; +use serde::{Deserialize, Serialize}; +use trogon_eventstore::{ + AppendToStreamOptions, Client, EventData, PersistentSubscriptionEvent, + PersistentSubscriptionOptions, RecordedEvent, StreamPosition, StreamState, + SubscribeToStreamOptions, SubscriptionEvent, +}; + +const EVENT_TYPE: &str = "trogon-compatibility"; +const PRODUCER: &str = "rust"; +const SERVICE_NAME: &str = "trogon-eventstore-client-rust"; +const TIMEOUT: Duration = Duration::from_secs(60); + +#[derive(Debug, Deserialize, PartialEq, Serialize)] +#[serde(rename_all = "camelCase")] +struct CompatibilityPayload { + producer: String, + run_id: String, +} + +#[derive(Debug, Serialize)] +#[serde(rename_all = "camelCase")] +struct CompatibilityResult { + command: String, + stream: String, + group: Option, + producer: String, + run_id: String, + event_id: Option, +} + +struct CompatibilityOptions { + command: CompatibilityCommand, + uri: String, + run_id: RunId, + stream: StreamName, + group: Option, + ready_file: Option, + otlp_endpoint: String, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum CompatibilityCommand { + Write, + BatchWrite, + Read, + Subscribe, + CreatePersistentSubscription, + ConsumePersistentSubscription, +} + +struct StreamName(String); +struct GroupName(String); +struct RunId(String); +struct ReadyFile(PathBuf); + +#[tokio::test(flavor = "multi_thread")] +#[ignore = "requires an external server and OTLP collector"] +async fn cross_client_compatibility() -> eyre::Result<()> { + let options = CompatibilityOptions::load()?; + let provider = configure_telemetry(options.otlp_endpoint.clone())?; + let tracer = global::tracer(SERVICE_NAME); + let root = tracer.start(format!("compatibility {}", options.command.as_str())); + let context = Context::current_with_span(root); + + let result = + tokio::time::timeout(TIMEOUT, execute(options).with_context(context.clone())).await; + context.span().end(); + provider.force_flush()?; + + let result = result.map_err(|_| eyre::eyre!("compatibility command timed out"))??; + println!("{}", serde_json::to_string(&result)?); + Ok(()) +} + +async fn execute(options: CompatibilityOptions) -> eyre::Result { + let settings = options.uri.parse()?; + let client = Client::new(settings)?; + + match options.command { + CompatibilityCommand::Write => write(&client, options).await, + CompatibilityCommand::BatchWrite => batch_write(&client, options).await, + CompatibilityCommand::Read => read(&client, options).await, + CompatibilityCommand::Subscribe => subscribe(&client, options).await, + CompatibilityCommand::CreatePersistentSubscription => { + create_persistent_subscription(&client, options).await + } + CompatibilityCommand::ConsumePersistentSubscription => { + consume_persistent_subscription(&client, options).await + } + } +} + +async fn write( + client: &Client, + options: CompatibilityOptions, +) -> eyre::Result { + let (payload, event_id, event) = compatibility_event(&options)?; + let append_options = AppendToStreamOptions::default().stream_state(StreamState::NoStream); + client + .append_to_stream(options.stream.as_str(), &append_options, event) + .await?; + + Ok(options.result(payload, Some(event_id))) +} + +async fn batch_write( + client: &Client, + options: CompatibilityOptions, +) -> eyre::Result { + let (payload, event_id, event) = compatibility_event(&options)?; + let batch = client.batch_append(&Default::default()).await?; + batch + .append_to_stream(options.stream.as_str(), StreamState::NoStream, vec![event]) + .await?; + + Ok(options.result(payload, Some(event_id))) +} + +fn compatibility_event( + options: &CompatibilityOptions, +) -> eyre::Result<(CompatibilityPayload, uuid::Uuid, EventData)> { + let payload = CompatibilityPayload { + producer: PRODUCER.to_owned(), + run_id: options.run_id.as_str().to_owned(), + }; + let event_id = uuid::Uuid::new_v4(); + let event = EventData::json(EVENT_TYPE, &payload)?.id(event_id); + + Ok((payload, event_id, event)) +} + +async fn read(client: &Client, options: CompatibilityOptions) -> eyre::Result { + let mut read = client + .read_stream(options.stream.as_str(), &Default::default()) + .await?; + + while let Some(event) = read.next().await? { + if let Some(payload) = + matching_payload(event.get_original_event(), options.run_id.as_str())? + { + return Ok(options.result(payload, Some(event.get_original_event().id))); + } + } + + Err(missing_event(&options)) +} + +async fn subscribe( + client: &Client, + options: CompatibilityOptions, +) -> eyre::Result { + let subscription_options = + SubscribeToStreamOptions::default().start_from(StreamPosition::Start); + let mut subscription = client + .subscribe_to_stream(options.stream.as_str(), &subscription_options) + .await; + loop { + if let SubscriptionEvent::Confirmed(_) = subscription.next_subscription_event().await? { + break; + } + } + options.signal_ready()?; + + loop { + let event = subscription.next().await?; + if let Some(payload) = + matching_payload(event.get_original_event(), options.run_id.as_str())? + { + return Ok(options.result(payload, Some(event.get_original_event().id))); + } + } +} + +async fn create_persistent_subscription( + client: &Client, + options: CompatibilityOptions, +) -> eyre::Result { + let group = options.required_group()?; + let subscription_options = + PersistentSubscriptionOptions::default().start_from(StreamPosition::Start); + client + .create_persistent_subscription(options.stream.as_str(), group, &subscription_options) + .await?; + let payload = CompatibilityPayload { + producer: PRODUCER.to_owned(), + run_id: options.run_id.as_str().to_owned(), + }; + + Ok(options.result(payload, None)) +} + +async fn consume_persistent_subscription( + client: &Client, + options: CompatibilityOptions, +) -> eyre::Result { + let group = options.required_group()?; + let mut subscription = client + .subscribe_to_persistent_subscription(options.stream.as_str(), group, &Default::default()) + .await?; + + loop { + if let PersistentSubscriptionEvent::Confirmed(_) = + subscription.next_subscription_event().await? + { + break; + } + } + options.signal_ready()?; + + loop { + match subscription.next_subscription_event().await? { + PersistentSubscriptionEvent::EventAppeared { event, .. } => { + subscription.ack(&event).await?; + if let Some(payload) = + matching_payload(event.get_original_event(), options.run_id.as_str())? + { + return Ok(options.result(payload, Some(event.get_original_event().id))); + } + } + PersistentSubscriptionEvent::Confirmed(_) => {} + } + } +} + +fn matching_payload( + event: &RecordedEvent, + run_id: &str, +) -> eyre::Result> { + if event.event_type != EVENT_TYPE { + return Ok(None); + } + + let payload: CompatibilityPayload = event.as_json()?; + if payload.run_id != run_id { + return Ok(None); + } + if payload.producer.trim().is_empty() { + return Err(eyre::eyre!("compatibility event producer is required")); + } + + Ok(Some(payload)) +} + +fn missing_event(options: &CompatibilityOptions) -> eyre::Report { + eyre::eyre!( + "stream {} does not contain a {EVENT_TYPE} event for run {}", + options.stream.as_str(), + options.run_id.as_str() + ) +} + +impl CompatibilityOptions { + fn load() -> eyre::Result { + Self::load_with(|name| std::env::var(name).ok()) + } + + fn load_with(read_environment: impl Fn(&str) -> Option) -> eyre::Result { + let command = CompatibilityCommand::parse(&required_value( + &read_environment, + "TROGON_EVENTSTORE_COMMAND", + )?)?; + let group = command + .requires_group() + .then(|| required_value(&read_environment, "TROGON_EVENTSTORE_GROUP").map(GroupName)) + .transpose()?; + let ready_file = optional_value(&read_environment, "TROGON_EVENTSTORE_READY_FILE") + .map(|value| ReadyFile(PathBuf::from(value))); + + Ok(Self { + command, + uri: required_value(&read_environment, "TROGON_EVENTSTORE_URI")?, + run_id: RunId(required_value( + &read_environment, + "TROGON_EVENTSTORE_RUN_ID", + )?), + stream: StreamName(required_value( + &read_environment, + "TROGON_EVENTSTORE_STREAM", + )?), + group, + ready_file, + otlp_endpoint: required_value(&read_environment, "OTEL_EXPORTER_OTLP_ENDPOINT")?, + }) + } + + fn required_group(&self) -> eyre::Result<&str> { + self.group.as_ref().map(GroupName::as_str).ok_or_else(|| { + eyre::eyre!( + "TROGON_EVENTSTORE_GROUP is required for {}", + self.command.as_str() + ) + }) + } + + fn signal_ready(&self) -> eyre::Result<()> { + if let Some(ready_file) = &self.ready_file { + ready_file.signal()?; + } + + Ok(()) + } + + fn result( + self, + payload: CompatibilityPayload, + event_id: Option, + ) -> CompatibilityResult { + let CompatibilityOptions { + command, + stream, + group, + .. + } = self; + CompatibilityResult { + command: command.as_str().to_owned(), + stream: stream.0, + group: group.map(|group| group.0), + producer: payload.producer, + run_id: payload.run_id, + event_id, + } + } +} + +impl CompatibilityCommand { + fn parse(value: &str) -> eyre::Result { + match value { + "write" => Ok(Self::Write), + "batch-write" => Ok(Self::BatchWrite), + "read" => Ok(Self::Read), + "subscribe" => Ok(Self::Subscribe), + "create-persistent-subscription" => Ok(Self::CreatePersistentSubscription), + "consume-persistent-subscription" => Ok(Self::ConsumePersistentSubscription), + _ => Err(eyre::eyre!("unsupported compatibility command {value}")), + } + } + + const fn as_str(self) -> &'static str { + match self { + Self::Write => "write", + Self::BatchWrite => "batch-write", + Self::Read => "read", + Self::Subscribe => "subscribe", + Self::CreatePersistentSubscription => "create-persistent-subscription", + Self::ConsumePersistentSubscription => "consume-persistent-subscription", + } + } + + const fn requires_group(self) -> bool { + matches!( + self, + Self::CreatePersistentSubscription | Self::ConsumePersistentSubscription + ) + } +} + +impl StreamName { + fn as_str(&self) -> &str { + &self.0 + } +} + +impl GroupName { + fn as_str(&self) -> &str { + &self.0 + } +} + +impl RunId { + fn as_str(&self) -> &str { + &self.0 + } +} + +impl ReadyFile { + fn signal(&self) -> eyre::Result<()> { + let mut file = std::fs::OpenOptions::new() + .write(true) + .create_new(true) + .open(&self.0)?; + file.write_all(b"ready\n")?; + Ok(()) + } + + #[cfg(test)] + fn path(&self) -> &Path { + &self.0 + } +} + +fn configure_telemetry(endpoint: String) -> eyre::Result { + global::set_text_map_propagator(TextMapCompositePropagator::new(vec![ + Box::new(TraceContextPropagator::new()), + Box::new(BaggagePropagator::new()), + ])); + let exporter = opentelemetry_otlp::SpanExporter::builder() + .with_tonic() + .with_endpoint(endpoint) + .build()?; + let resource = Resource::builder().with_service_name(SERVICE_NAME).build(); + let provider = SdkTracerProvider::builder() + .with_resource(resource) + .with_batch_exporter(exporter) + .build(); + global::set_tracer_provider(provider.clone()); + + Ok(provider) +} + +fn required_value( + read_environment: &impl Fn(&str) -> Option, + name: &str, +) -> eyre::Result { + optional_value(read_environment, name) + .ok_or_else(|| eyre::eyre!("required environment variable {name} is not set or is empty")) +} + +fn optional_value( + read_environment: &impl Fn(&str) -> Option, + name: &str, +) -> Option { + read_environment(name).and_then(|value| { + let value = value.trim(); + if value.is_empty() { + None + } else { + Some(value.to_owned()) + } + }) +} + +#[cfg(test)] +mod tests { + use super::{CompatibilityCommand, CompatibilityOptions}; + use std::collections::HashMap; + use std::path::Path; + + #[test] + fn options_load_an_optional_ready_file() { + let mut environment = base_environment("subscribe"); + environment.insert( + "TROGON_EVENTSTORE_READY_FILE", + " /tmp/trogon-ready ".to_owned(), + ); + + let options = CompatibilityOptions::load_with(|name| environment.get(name).cloned()) + .expect("compatibility options"); + + assert_eq!(options.command, CompatibilityCommand::Subscribe); + assert_eq!( + options.ready_file.as_ref().map(|file| file.path()), + Some(Path::new("/tmp/trogon-ready")) + ); + assert!(options.group.is_none()); + } + + #[test] + fn persistent_commands_require_a_group() { + let environment = base_environment("consume-persistent-subscription"); + + let error = CompatibilityOptions::load_with(|name| environment.get(name).cloned()) + .err() + .expect("missing group error"); + + assert!(error.to_string().contains("TROGON_EVENTSTORE_GROUP")); + } + + fn base_environment(command: &str) -> HashMap<&'static str, String> { + HashMap::from([ + ("TROGON_EVENTSTORE_COMMAND", command.to_owned()), + ( + "TROGON_EVENTSTORE_URI", + "esdb://localhost:2113?tls=false".to_owned(), + ), + ("TROGON_EVENTSTORE_RUN_ID", "run-1".to_owned()), + ("TROGON_EVENTSTORE_STREAM", "compatibility-1".to_owned()), + ( + "OTEL_EXPORTER_OTLP_ENDPOINT", + "http://localhost:4317".to_owned(), + ), + ]) + } +}