try debug difference trace_id otel by use async for gateway and attach it to span in surb_reply

This commit is contained in:
Floriane TUERNAL SABOTINOV
2025-08-29 18:19:52 +02:00
parent c5590bddc5
commit 572fc331b0
2 changed files with 113 additions and 118 deletions
@@ -34,6 +34,7 @@ use nym_node_metrics::events::MetricsEvent;
use nym_sphinx::DestinationAddressBytes;
use nym_task::ShutdownToken;
use opentelemetry::trace::TraceContextExt;
use opentelemetry::TraceId;
use opentelemetry_sdk::trace::{IdGenerator, RandomIdGenerator};
use rand::CryptoRng;
use std::net::SocketAddr;
@@ -45,6 +46,7 @@ use tokio::time::timeout;
use tokio_tungstenite::tungstenite::{protocol::Message, Error as WsError};
use tracing::{debug, error, info, info_span, instrument, warn};
use tracing_opentelemetry::OpenTelemetrySpanExt;
use tracing::Instrument;
#[derive(Debug, Error)]
pub(crate) enum InitialAuthenticationError {
@@ -851,7 +853,7 @@ impl<R, S> FreshHandler<R, S> {
debug!("failed to reply with protocol version: {err}")
}
}
#[instrument(skip_all)]
pub(crate) async fn handle_initial_client_request(
&mut self,
request: ClientControlRequest,
@@ -860,118 +862,118 @@ impl<R, S> FreshHandler<R, S> {
S: AsyncRead + AsyncWrite + Unpin + Send,
R: CryptoRng + RngCore + Send,
{
let span = if let ClientControlRequest::AuthenticateV2(ref auth_req) = request {
if let Some(ref trace_id) = auth_req.debug_trace_id {
warn!("RAW TRACE ID: {trace_id:?}");
let trace_id = opentelemetry::trace::TraceId::from_hex(&trace_id)
.expect("Invalid trace ID format");
warn!("🫂TraceID: {trace_id}🫂");
let distributed_context = if let ClientControlRequest::AuthenticateV2(ref auth_req) = request {
if let Some(ref trace_id_str) = auth_req.debug_trace_id {
warn!("Received debug trace ID: {trace_id_str}");
// We don't need to try and preserve the SpanID, just the TraceID (right?) so
// just making a new SpanID for the moment
let id_generator = RandomIdGenerator::default();
let span_id = id_generator.new_span_id();
match TraceId::from_hex(&trace_id_str) {
Ok(trace_id) => {
warn!("Parsed debug trace ID: {trace_id}");
let span_context = opentelemetry::trace::SpanContext::new(
trace_id,
span_id,
opentelemetry::trace::TraceFlags::SAMPLED,
true, // is_remote = true since this comes from another service
Default::default(),
);
let id_gen = RandomIdGenerator::default();
let span_id = id_gen.new_span_id();
let remote_context =
opentelemetry::Context::current().with_remote_span_context(span_context);
let span_context = opentelemetry::trace::SpanContext::new(
trace_id,
span_id,
opentelemetry::trace::TraceFlags::SAMPLED,
true,
Default::default(),
);
let _context_guard = remote_context.clone().attach();
let span = info_span!(
"authenticate_v2",
trace_id = %trace_id
);
span.set_parent(remote_context.clone());
let remote_context = opentelemetry::Context::current()
.with_remote_span_context(span_context);
Some(span)
let span = info_span!(
"handle_initial_client_request_distributed",
trace_id = %trace_id,
service = "nym-node"
);
span.set_parent(remote_context);
warn!("Distributed context established");
Some(span)
}
Err(err) => {
warn!("Failed to parse debug trace ID: {err}");
None
}
}
} else {
warn!("AuthenticateV2 request but no trace_id provided");
warn!("Received debug trace ID without context");
None
}
} else {
warn!("Not an AuthenticateV2 request");
None
};
// Probably a nicer way to do this but for now just match
let _guard = match &span {
Some(s) => {
warn!("ENTERED SPAN");
Some(s.enter())
}
None => {
warn!("COULDN'T ENTER SPAN");
None
}
};
// we can handle stateless client requests without prior authentication, like `ClientControlRequest::SupportedProtocol`
let auth_result = match request {
ClientControlRequest::Authenticate {
protocol_version,
address,
enc_address,
iv,
debug_trace_id: None,
} => {
self.handle_legacy_authenticate(protocol_version, address, enc_address, iv)
.await
}
ClientControlRequest::AuthenticateV2(req) => self.handle_authenticate_v2(req).await,
ClientControlRequest::RegisterHandshakeInitRequest {
protocol_version,
data,
} => self.handle_register(protocol_version, data).await,
ClientControlRequest::SupportedProtocol { .. } => {
self.handle_reply_supported_protocol_request().await;
return Ok(None);
}
_ => {
debug!("received an invalid client request");
return Err(InitialAuthenticationError::InvalidRequest);
}
};
let auth_result = match auth_result {
Ok(res) => res,
Err(err) => {
match &err {
InitialAuthenticationError::StorageError(inner_storage) => {
debug!("authentication failure due to storage issue: {inner_storage}")
let handle_request = async {
let auth_result = match request {
ClientControlRequest::Authenticate {
protocol_version,
address,
enc_address,
iv,
debug_trace_id: None,
} => {
self.handle_legacy_authenticate(protocol_version, address, enc_address, iv)
.await
}
other => debug!("authentication failure: {other}"),
ClientControlRequest::AuthenticateV2(req) => self.handle_authenticate_v2(req).await,
ClientControlRequest::RegisterHandshakeInitRequest {
protocol_version,
data,
} => self.handle_register(protocol_version, data).await,
ClientControlRequest::SupportedProtocol { .. } => {
self.handle_reply_supported_protocol_request().await;
return Ok(None);
}
_ => {
debug!("received an invalid client request");
return Err(InitialAuthenticationError::InvalidRequest);
}
};
let auth_result = match auth_result {
Ok(res) => res,
Err(err) => {
match &err {
InitialAuthenticationError::StorageError(inner_storage) => {
debug!("authentication failure due to storage issue: {inner_storage}")
}
other => debug!("authentication failure: {other}"),
}
self.send_and_forget_error_response(&err).await;
return Err(err);
}
};
// try to send auth response back to the client
if let Err(source) = self
.send_websocket_message(auth_result.server_response)
.await
{
debug!("failed to send authentication response: {source}");
return Err(InitialAuthenticationError::ResponseSendFailure {
source: Box::new(source),
});
}
self.send_and_forget_error_response(&err).await;
return Err(err);
}
let Some(client_details) = auth_result.client_details else {
// honestly, it's been so long I don't remember under what conditions its possible (if at all)
// to have empty client details
warn!("could not establish client details");
return Err(InitialAuthenticationError::EmptyClientDetails);
};
Ok(Some(client_details))
};
// try to send auth response back to the client
if let Err(source) = self
.send_websocket_message(auth_result.server_response)
.await
{
debug!("failed to send authentication response: {source}");
return Err(InitialAuthenticationError::ResponseSendFailure {
source: Box::new(source),
});
match distributed_context {
Some(span) => handle_request.instrument(span).await,
None => handle_request.await,
}
let Some(client_details) = auth_result.client_details else {
// honestly, it's been so long I don't remember under what conditions its possible (if at all)
// to have empty client details
warn!("could not establish client details");
return Err(InitialAuthenticationError::EmptyClientDetails);
};
Ok(Some(client_details))
}
#[instrument(skip_all)]
+16 -23
View File
@@ -3,39 +3,32 @@ use nym_sdk::mixnet::{
// StoragePaths,
};
use opentelemetry::trace::TraceContextExt;
use opentelemetry::trace::Tracer;
use opentelemetry::Context;
use opentelemetry::global;
// use opentelemetry::trace::Tracer;
// use opentelemetry::Context;
// use opentelemetry::global;
use tracing::warn;
use tracing::{instrument, info_span};
use tracing_opentelemetry::OpenTelemetrySpanExt;
// use std::path::PathBuf;
// use tempfile::TempDir;
#[tokio::main]
#[instrument(name = "sdk-example-surb-reply", skip_all)]
async fn main() {
let _guard = nym_bin_common::opentelemetry::setup_tracing_logger("sdk-example-surb-reply".to_string()).unwrap();
nym_bin_common::opentelemetry::setup_tracing_logger("sdk-example-surb-reply".to_string()).unwrap();
let tracer = global::tracer("sdk-example-surb-reply");
let span = tracer.start("test_span");
let cx = Context::current_with_span(span);
let _guard = cx.clone().attach();
let main_span = info_span!("surb_example_session");
let _main_span_enter = main_span.enter();
let trace_id = cx.span().span_context().trace_id();
warn!("Main trace ID: {}", trace_id);
let span = info_span!(
"surb_reply_example session",
trace_id = %trace_id.to_string()
);
let _enter = span.enter();
let otel_context = opentelemetry::Context::current();
warn!("OTEL CONTEXT: {:?}", otel_context);
let span = otel_context.span();
let context = span.span_context();
let trace_id = context.trace_id();
warn!("TRACE_ID: {:?}", trace_id);
let current_span = tracing::Span::current();
let otel_context = current_span.context();
let binding = otel_context.span();
let span_context = binding.span_context();
let trace_id = span_context.trace_id();
warn!("Starting the SURB reply example - trace id: {}", trace_id);
warn!("Otel context: {:?}", otel_context);
warn!("trace id: {}", trace_id);
// // Specify some config options
// let config_dir: PathBuf = TempDir::new().unwrap().path().to_path_buf();