Skip to content

Commit f309717

Browse files
authored
feat(observability): preserve cross-client trace continuity (#10)
Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent 9f54ab9 commit f309717

10 files changed

Lines changed: 1170 additions & 21 deletions

File tree

Cargo.lock

Lines changed: 120 additions & 9 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

otel/semconv/registry/trogon/eventstore/client-spans.yaml

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,11 @@ groups:
33
type: attribute_group
44
brief: TrogonEventStore client attributes.
55
attributes:
6+
- id: trogon.eventstore.event.type
7+
type: string
8+
stability: development
9+
brief: The type of event received by the client.
10+
examples: [order-created]
611
- id: trogon.eventstore.batch.correlation_id
712
type: string
813
stability: development
@@ -28,6 +33,28 @@ groups:
2833
requirement_level:
2934
conditionally_required: If and only if the operation failed.
3035

36+
- id: span.trogon.eventstore.client.receive
37+
type: span
38+
span_kind: client
39+
stability: development
40+
brief: Receives an event delivered by a subscription.
41+
attributes:
42+
- ref: messaging.system
43+
requirement_level: required
44+
- ref: messaging.operation.name
45+
requirement_level: required
46+
- ref: messaging.operation.type
47+
requirement_level: required
48+
- ref: messaging.destination.name
49+
requirement_level: required
50+
- ref: messaging.message.id
51+
requirement_level: recommended
52+
- ref: messaging.consumer.group.name
53+
requirement_level:
54+
conditionally_required: If the event was delivered by a persistent subscription.
55+
- ref: trogon.eventstore.event.type
56+
requirement_level: recommended
57+
3158
- id: span.trogon.eventstore.client.append_to_stream
3259
type: span
3360
extends: span.trogon.eventstore.client

otel/semconv/templates/registry/rust/observability.rs.j2

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@ pub(crate) const {{ attribute.name | screaming_snake_case }}: &str = "{{ attribu
88
{% endfor %}
99

1010
pub(crate) const CLIENT_SPAN_KIND: SpanKind = SpanKind::{{ ctx.base.span_kind | pascal_case }};
11+
pub(crate) const RECEIVE_SPAN_KIND: SpanKind = SpanKind::{{ ctx.receive.span_kind | pascal_case }};
1112

1213
pub(crate) mod operation {
1314
use super::ClientOperation;

otel/semconv/templates/registry/rust/weaver.yaml

Lines changed: 12 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,12 +7,23 @@ templates:
77
filter: >
88
{
99
base: [semconv_grouped_spans[].spans[] | select(.id == "span.trogon.eventstore.client")][0],
10-
operations: [semconv_grouped_spans[].spans[] | select(.id | startswith("span.trogon.eventstore.client."))],
10+
receive: [semconv_grouped_spans[].spans[] | select(.id == "span.trogon.eventstore.client.receive")][0],
11+
operations: [semconv_grouped_spans[].spans[] | select(
12+
(.id | startswith("span.trogon.eventstore.client.")) and
13+
.annotations.code_generation.operation_name != null
14+
)],
1115
attributes: [semconv_grouped_spans[].spans[].attributes[] | select(
1216
.name == "db.collection.name" or
1317
.name == "db.operation.name" or
1418
.name == "db.system.name" or
1519
.name == "error.type" or
20+
.name == "messaging.consumer.group.name" or
21+
.name == "messaging.destination.name" or
22+
.name == "messaging.message.id" or
23+
.name == "messaging.operation.name" or
24+
.name == "messaging.operation.type" or
25+
.name == "messaging.system" or
26+
.name == "trogon.eventstore.event.type" or
1627
.name == "trogon.eventstore.batch.correlation_id"
1728
)] | unique_by(.name)
1829
}

trogon-eventstore/Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@ name = "integration"
6969

7070
[dev-dependencies]
7171
names = "0.14"
72+
opentelemetry-otlp = { version = "0.32", default-features = false, features = ["grpc-tonic", "trace"] }
7273
opentelemetry_sdk = { version = "0.32", default-features = false, features = ["testing", "trace"] }
7374
serde = { version = "1", features = ["derive"] }
7475
testcontainers = "0.23"

trogon-eventstore/src/batch.rs

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,6 +21,7 @@ pub(crate) struct Req {
2121
pub(crate) stream_name: String,
2222
pub(crate) events: Vec<EventData>,
2323
pub(crate) expected_revision: StreamState,
24+
pub(crate) context: Context,
2425
}
2526

2627
impl Req {
@@ -30,6 +31,7 @@ impl Req {
3031
stream_name,
3132
events,
3233
expected_revision,
34+
context: Context::current(),
3335
}
3436
}
3537
}
@@ -200,7 +202,7 @@ mod tests {
200202
DB_COLLECTION_NAME, DB_OPERATION_NAME, TROGON_EVENTSTORE_BATCH_CORRELATION_ID,
201203
};
202204
use opentelemetry::global;
203-
use opentelemetry::trace::{Status, noop::NoopTracerProvider};
205+
use opentelemetry::trace::{Status, TraceContextExt, noop::NoopTracerProvider};
204206
use opentelemetry_sdk::trace::{InMemorySpanExporter, SdkTracerProvider};
205207

206208
#[tokio::test]
@@ -223,6 +225,7 @@ mod tests {
223225
panic!("expected an inbound batch request");
224226
};
225227
let correlation_id = req.id.to_string();
228+
let captured_span_id = req.context.span().span_context().span_id();
226229
sender
227230
.send(Ok(BatchWriteResult::new(
228231
"stream".to_string(),
@@ -239,6 +242,7 @@ mod tests {
239242
.iter()
240243
.find(|span| span.name == "batch_append_to_stream stream")
241244
.expect("batch append client span");
245+
assert_eq!(span.span_context.span_id(), captured_span_id);
242246
assert!(span.attributes.iter().any(|attribute| {
243247
attribute.key.as_str() == TROGON_EVENTSTORE_BATCH_CORRELATION_ID
244248
&& attribute.value.to_string() == correlation_id

0 commit comments

Comments
 (0)