From 6a76c2ded7d824cb3ab8330b4665730bcd0a5db1 Mon Sep 17 00:00:00 2001 From: brotoskyj Date: Thu, 8 Jan 2026 15:20:54 -0500 Subject: [PATCH] Fixed Telemetry Exporter Closes #56 --- airlock_libs/Cargo.lock | 2 +- airlock_libs/Cargo.toml | 2 +- airlock_libs/pyproject.toml | 2 +- airlock_libs/src/services.rs | 19 ++++++++++++++++--- requirements.txt | 2 +- 5 files changed, 20 insertions(+), 7 deletions(-) diff --git a/airlock_libs/Cargo.lock b/airlock_libs/Cargo.lock index e14b90e..93b9cee 100644 --- a/airlock_libs/Cargo.lock +++ b/airlock_libs/Cargo.lock @@ -26,7 +26,7 @@ dependencies = [ [[package]] name = "airlock_libs" -version = "7.3.0" +version = "7.4.0" dependencies = [ "chrono", "crossbeam", diff --git a/airlock_libs/Cargo.toml b/airlock_libs/Cargo.toml index 8827af5..f238c5b 100644 --- a/airlock_libs/Cargo.toml +++ b/airlock_libs/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "airlock_libs" -version = "7.3.0" +version = "7.4.0" edition = "2024" [dependencies] diff --git a/airlock_libs/pyproject.toml b/airlock_libs/pyproject.toml index 5060fe8..6a0c9f9 100644 --- a/airlock_libs/pyproject.toml +++ b/airlock_libs/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "maturin" [project] name = "airlock_libs" -version = "7.3.0" +version = "7.4.0" description = "Airlock Digital API Wrapper" readme = "README.md" license = { text = "AGPL-3.0-only" } diff --git a/airlock_libs/src/services.rs b/airlock_libs/src/services.rs index 3023f45..bc3d9cd 100644 --- a/airlock_libs/src/services.rs +++ b/airlock_libs/src/services.rs @@ -1,3 +1,5 @@ +use opentelemetry::trace::SpanContext; + use crate::modules::datatypes::*; use crate::prelude::*; @@ -106,9 +108,10 @@ pub fn pull_policy_exec_histories( }); let cutoff: chrono::NaiveDateTime = Local::now().naive_local() - chrono::Duration::days(days); - let (tx, rx) = unbounded::<(Context, Vec)>(); + let (tx, rx) = unbounded::<(SpanContext, Vec)>(); let pb_clone = progress_bar.clone(); thread::spawn(move || { + let tracer = global::tracer("Loxide"); let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists() { let contents: String = fs::read_to_string(&writeable_filepath).unwrap_or_default(); @@ -137,7 +140,15 @@ pub fn pull_policy_exec_histories( } else { HashMap::new() }; - while let Ok((cx, parsed_responses)) = rx.recv() { + while let Ok((parent_spancontext, parsed_responses)) = rx.recv() { + let parent_ctx = Context::new().with_remote_span_context(parent_spancontext); + let span = tracer.build_with_context( + tracer + .span_builder("Deduplicate and Write") + .with_kind(trace::SpanKind::Consumer), + &parent_ctx, + ); + let cx = Context::current_with_span(span); cx.span().add_event( "Received Data from Producer", vec![KeyValue::new( @@ -249,7 +260,7 @@ pub fn pull_policy_exec_histories( if parsed_responses.is_empty() { break; } - match tx.send((cx.clone(), parsed_responses.clone())) { + match tx.send((cx.span().span_context().clone(), parsed_responses.clone())) { Ok(_) => {} Err(e) => { cx.span().add_event( @@ -281,6 +292,7 @@ pub fn pull_policy_exec_histories( } } } + tracer_provider.force_flush(); }); progress_bar .lock() @@ -293,6 +305,7 @@ pub fn pull_policy_exec_histories( std::process::abort(); } }; + tracer_provider.force_flush(); tracer_provider .shutdown() .expect("Failed to Shutdown Tracer Provdier"); diff --git a/requirements.txt b/requirements.txt index d3c4959..28908a5 100644 --- a/requirements.txt +++ b/requirements.txt @@ -23,4 +23,4 @@ pyperclip==1.11.0 # Custom/Private packages --extra-index-url https://git.racooncity.org/api/packages/brotoskyj/pypi/simple/ -airlock_libs==7.3.0 \ No newline at end of file +airlock_libs==7.4.0 \ No newline at end of file -- 2.34.1