Merge pull request 'Fixed Telemetry Exporter' (#57) from fix/TelemetryExporter into master

Reviewed-on: brotoskyj/AirlockTools#57
This commit was merged in pull request #57.
This commit is contained in:
James Brotosky
2026-01-08 15:21:22 -05:00
5 changed files with 20 additions and 7 deletions
+1 -1
View File
@@ -26,7 +26,7 @@ dependencies = [
[[package]]
name = "airlock_libs"
version = "7.3.0"
version = "7.4.0"
dependencies = [
"chrono",
"crossbeam",
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "airlock_libs"
version = "7.3.0"
version = "7.4.0"
edition = "2024"
[dependencies]
+1 -1
View File
@@ -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" }
+16 -3
View File
@@ -1,3 +1,5 @@
use opentelemetry::trace::SpanContext;
use crate::modules::datatypes::*;
use crate::prelude::*;
@@ -107,9 +109,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<Group>)>();
let (tx, rx) = unbounded::<(SpanContext, Vec<Group>)>();
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();
@@ -138,7 +141,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(
@@ -250,7 +261,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(
@@ -282,6 +293,7 @@ pub fn pull_policy_exec_histories(
}
}
}
tracer_provider.force_flush();
});
progress_bar
.lock()
@@ -294,6 +306,7 @@ pub fn pull_policy_exec_histories(
std::process::abort();
}
};
tracer_provider.force_flush();
tracer_provider
.shutdown()
.expect("Failed to Shutdown Tracer Provdier");