From 66bb21ed88046853c91e56bdddaaf473e2cf3f7b Mon Sep 17 00:00:00 2001 From: brotoskyj Date: Wed, 17 Dec 2025 16:38:53 -0500 Subject: [PATCH 1/4] Styling - Progress Bar Changed progress bar indicators and added policy name to the progress bar closes #50 --- airlock_libs/Cargo.lock | 2 +- airlock_libs/Cargo.toml | 2 +- airlock_libs/pyproject.toml | 2 +- airlock_libs/src/services.rs | 6 ++++-- requirements.txt | 2 +- 5 files changed, 8 insertions(+), 6 deletions(-) diff --git a/airlock_libs/Cargo.lock b/airlock_libs/Cargo.lock index e2386f3..391d9b4 100644 --- a/airlock_libs/Cargo.lock +++ b/airlock_libs/Cargo.lock @@ -26,7 +26,7 @@ dependencies = [ [[package]] name = "airlock_libs" -version = "6.1.1" +version = "6.2.0" dependencies = [ "chrono", "crossbeam", diff --git a/airlock_libs/Cargo.toml b/airlock_libs/Cargo.toml index cf49ac1..888c19c 100644 --- a/airlock_libs/Cargo.toml +++ b/airlock_libs/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "airlock_libs" -version = "6.1.1" +version = "6.2.0" edition = "2024" [dependencies] diff --git a/airlock_libs/pyproject.toml b/airlock_libs/pyproject.toml index 5ed6be4..1c2fcd2 100644 --- a/airlock_libs/pyproject.toml +++ b/airlock_libs/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "maturin" [project] name = "airlock_libs" -version = "6.1.1" +version = "6.2.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 c567f41..3ddcf1e 100644 --- a/airlock_libs/src/services.rs +++ b/airlock_libs/src/services.rs @@ -73,8 +73,9 @@ pub fn pull_policy_exec_histories( .set_draw_target(ProgressDrawTarget::stderr()); progress_bar.lock().unwrap().set_style( ProgressStyle::default_bar() - .template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len} {message}") - .unwrap(), + .template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/cyan}] {pos}/{len} - Current Policy: {msg}") + .unwrap() + .progress_chars("#> ") ); let client: Client = tracer.in_span("Building HTTP Client", |cx| { let client_result: Result = build_client(headers); @@ -108,6 +109,7 @@ pub fn pull_policy_exec_histories( Local::now().naive_local() - chrono::Duration::days(days); let (tx, rx) = unbounded::>(); let pb_clone = progress_bar.clone(); + pb_clone.lock().unwrap().set_message(policy_names.clone()); thread::spawn(move || { let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists() { diff --git a/requirements.txt b/requirements.txt index 30d5f35..83a2071 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==6.1.1 \ No newline at end of file +airlock_libs==6.2.0 \ No newline at end of file From 729b45f52afed33f7c644f9cb149a0f05b993647 Mon Sep 17 00:00:00 2001 From: brotoskyj Date: Thu, 18 Dec 2025 10:14:58 -0500 Subject: [PATCH 2/4] Features - Optional Policy Name in Execution Histories closes #48 --- airlock_libs/Cargo.lock | 2 +- airlock_libs/Cargo.toml | 2 +- airlock_libs/pyproject.toml | 2 +- airlock_libs/src/services.rs | 29 ++++++++++++++++++++--------- requirements.txt | 2 +- 5 files changed, 24 insertions(+), 13 deletions(-) diff --git a/airlock_libs/Cargo.lock b/airlock_libs/Cargo.lock index 391d9b4..d263687 100644 --- a/airlock_libs/Cargo.lock +++ b/airlock_libs/Cargo.lock @@ -26,7 +26,7 @@ dependencies = [ [[package]] name = "airlock_libs" -version = "6.2.0" +version = "7.0.0" dependencies = [ "chrono", "crossbeam", diff --git a/airlock_libs/Cargo.toml b/airlock_libs/Cargo.toml index 888c19c..2e6b4f0 100644 --- a/airlock_libs/Cargo.toml +++ b/airlock_libs/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "airlock_libs" -version = "6.2.0" +version = "7.0.0" edition = "2024" [dependencies] diff --git a/airlock_libs/pyproject.toml b/airlock_libs/pyproject.toml index 1c2fcd2..95fce4f 100644 --- a/airlock_libs/pyproject.toml +++ b/airlock_libs/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "maturin" [project] name = "airlock_libs" -version = "6.2.0" +version = "7.0.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 3ddcf1e..657f9b3 100644 --- a/airlock_libs/src/services.rs +++ b/airlock_libs/src/services.rs @@ -5,7 +5,7 @@ use crate::prelude::*; pub fn pull_policy_exec_histories( py: Python<'_>, py_self: Py, - policy_names: String, + policy_names: Option, exec_types: String, days: i64, ) -> Py { @@ -73,9 +73,8 @@ pub fn pull_policy_exec_histories( .set_draw_target(ProgressDrawTarget::stderr()); progress_bar.lock().unwrap().set_style( ProgressStyle::default_bar() - .template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/cyan}] {pos}/{len} - Current Policy: {msg}") - .unwrap() - .progress_chars("#> ") + .template("Total - Policy Name: {msg}: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len}") + .unwrap().progress_chars("≫>›") ); let client: Client = tracer.in_span("Building HTTP Client", |cx| { let client_result: Result = build_client(headers); @@ -109,7 +108,6 @@ pub fn pull_policy_exec_histories( Local::now().naive_local() - chrono::Duration::days(days); let (tx, rx) = unbounded::>(); let pb_clone = progress_bar.clone(); - pb_clone.lock().unwrap().set_message(policy_names.clone()); thread::spawn(move || { let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists() { @@ -181,9 +179,18 @@ pub fn pull_policy_exec_histories( .lock() .unwrap() .enable_steady_tick(std::time::Duration::from_millis(100)); + pb_clone + .lock() + .unwrap() + .set_message(policy_names.clone().unwrap_or("Statistics".to_string())); let span: opentelemetry::trace::SpanRef<'_> = cx.span(); span.set_attribute(KeyValue::new("Days", days.to_string())); - span.set_attribute(KeyValue::new("Policy Name", policy_names.clone())); + span.set_attribute(KeyValue::new( + "Policy Name", + policy_names + .clone() + .unwrap_or("Statistics Monitoring".to_string()), + )); loop { let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| { let results: ApiResponse = history_logging( @@ -261,16 +268,20 @@ async fn history_logging( base_url: &String, exec_types: &String, checkpoint_number: &String, - policy_names: &String, + policy_names: &Option, client: &Client, ) -> ApiResponse { + let policy_json = match policy_names { + Some(name) => format!(r#"[ "{}" ]"#, name), // JSON array with one element + None => "[]".to_string(), // Empty JSON array + }; let payload = format!( r#"{{ "type": {}, "checkpoint": "{}", - "policy": ["{}"] + "policy": {} }}"#, - exec_types, checkpoint_number, policy_names + exec_types, checkpoint_number, policy_json ); let res: Result = client .post(format!("{}/v1/logging/exechistories", base_url)) diff --git a/requirements.txt b/requirements.txt index 83a2071..0357766 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==6.2.0 \ No newline at end of file +airlock_libs==7.0.0 \ No newline at end of file From 55557474227fb420fee794d6816626f0aab76744 Mon Sep 17 00:00:00 2001 From: brotoskyj Date: Thu, 18 Dec 2025 12:58:51 -0500 Subject: [PATCH 3/4] Styling - Progress Bar Changed glyphs in progress bar for better printing to console closes #51 --- airlock_libs/Cargo.lock | 2 +- airlock_libs/Cargo.toml | 2 +- airlock_libs/pyproject.toml | 2 +- airlock_libs/src/services.rs | 2 +- requirements.txt | 2 +- 5 files changed, 5 insertions(+), 5 deletions(-) diff --git a/airlock_libs/Cargo.lock b/airlock_libs/Cargo.lock index d263687..eaaadaa 100644 --- a/airlock_libs/Cargo.lock +++ b/airlock_libs/Cargo.lock @@ -26,7 +26,7 @@ dependencies = [ [[package]] name = "airlock_libs" -version = "7.0.0" +version = "7.1.0" dependencies = [ "chrono", "crossbeam", diff --git a/airlock_libs/Cargo.toml b/airlock_libs/Cargo.toml index 2e6b4f0..4833a01 100644 --- a/airlock_libs/Cargo.toml +++ b/airlock_libs/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "airlock_libs" -version = "7.0.0" +version = "7.1.0" edition = "2024" [dependencies] diff --git a/airlock_libs/pyproject.toml b/airlock_libs/pyproject.toml index 95fce4f..aee7519 100644 --- a/airlock_libs/pyproject.toml +++ b/airlock_libs/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "maturin" [project] name = "airlock_libs" -version = "7.0.0" +version = "7.1.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 657f9b3..767598d 100644 --- a/airlock_libs/src/services.rs +++ b/airlock_libs/src/services.rs @@ -74,7 +74,7 @@ pub fn pull_policy_exec_histories( progress_bar.lock().unwrap().set_style( ProgressStyle::default_bar() .template("Total - Policy Name: {msg}: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len}") - .unwrap().progress_chars("≫>›") + .unwrap().progress_chars("⣿⣦⣀") ); let client: Client = tracer.in_span("Building HTTP Client", |cx| { let client_result: Result = build_client(headers); diff --git a/requirements.txt b/requirements.txt index 0357766..6862175 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.0.0 \ No newline at end of file +airlock_libs==7.1.0 \ No newline at end of file From 72681218e73640bea162e32369ae56b9ec250251 Mon Sep 17 00:00:00 2001 From: brotoskyj Date: Wed, 7 Jan 2026 17:01:33 -0500 Subject: [PATCH 4/4] Fix for API Data Retrieval Performance Optimization - History Logging was spawning a new tokio runtime every function call, it now uses the global run time Documentation - Added a lot more span events for better logging --- airlock_libs/Cargo.lock | 2 +- airlock_libs/Cargo.toml | 2 +- airlock_libs/pyproject.toml | 2 +- airlock_libs/src/services.rs | 66 +++++++++++++++++++++++++++++++----- requirements.txt | 2 +- 5 files changed, 61 insertions(+), 13 deletions(-) diff --git a/airlock_libs/Cargo.lock b/airlock_libs/Cargo.lock index eaaadaa..e14b90e 100644 --- a/airlock_libs/Cargo.lock +++ b/airlock_libs/Cargo.lock @@ -26,7 +26,7 @@ dependencies = [ [[package]] name = "airlock_libs" -version = "7.1.0" +version = "7.3.0" dependencies = [ "chrono", "crossbeam", diff --git a/airlock_libs/Cargo.toml b/airlock_libs/Cargo.toml index 4833a01..8827af5 100644 --- a/airlock_libs/Cargo.toml +++ b/airlock_libs/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "airlock_libs" -version = "7.1.0" +version = "7.3.0" edition = "2024" [dependencies] diff --git a/airlock_libs/pyproject.toml b/airlock_libs/pyproject.toml index aee7519..5060fe8 100644 --- a/airlock_libs/pyproject.toml +++ b/airlock_libs/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "maturin" [project] name = "airlock_libs" -version = "7.1.0" +version = "7.3.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 767598d..3023f45 100644 --- a/airlock_libs/src/services.rs +++ b/airlock_libs/src/services.rs @@ -106,7 +106,7 @@ pub fn pull_policy_exec_histories( }); let cutoff: chrono::NaiveDateTime = Local::now().naive_local() - chrono::Duration::days(days); - let (tx, rx) = unbounded::>(); + let (tx, rx) = unbounded::<(Context, Vec)>(); let pb_clone = progress_bar.clone(); thread::spawn(move || { let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists() @@ -137,7 +137,14 @@ pub fn pull_policy_exec_histories( } else { HashMap::new() }; - while let Ok(parsed_responses) = rx.recv() { + while let Ok((cx, parsed_responses)) = rx.recv() { + cx.span().add_event( + "Received Data from Producer", + vec![KeyValue::new( + "Items to Process", + parsed_responses.len().to_string(), + )], + ); for executions in parsed_responses { if executions.checkpoint.is_empty() || executions.datetime.is_empty() { continue; @@ -166,11 +173,28 @@ pub fn pull_policy_exec_histories( }; let data_write: String = serde_json::to_string_pretty(&final_response).unwrap(); match fs::write(&writeable_filepath, data_write) { - Ok(_) => {} + Ok(_) => { + cx.span().add_event( + "Writing Data to File", + vec![KeyValue::new("Success", "Ok".to_string())], + ); + } Err(e) => { - println!("Failed to write to: {:?}: {}", &writeable_filepath, e); + cx.span().add_event( + "Writing Data to File", + vec![KeyValue::new("Failed", e.to_string())], + ); + cx.span() + .set_status(Status::error("Failed to Write to File")); } } + cx.span().add_event( + "Finished Deduplicating Data", + vec![KeyValue::new( + "Items Successfully Processed", + seen.len().to_string(), + )], + ); } }); let mut first_date: Option = None; @@ -193,13 +217,28 @@ pub fn pull_policy_exec_histories( )); loop { let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| { - let results: ApiResponse = history_logging( + cx.span().add_event( + "Retrieving Responses from API", + vec![KeyValue::new( + "Checkpoint Number", + checkpoint_number.to_string(), + )], + ); + let results: ApiResponse = rt.block_on(history_logging( &base_url, &exec_types, &checkpoint_number, &policy_names, &client, + )); + cx.span().add_event( + "Got Responses from API", + vec![KeyValue::new( + "Items in Response", + results.response.exechistories.len().to_string(), + )], ); + cx.span().set_status(Status::Ok); cx.span().set_attribute(KeyValue::new( "items_in_response", results.response.exechistories.len().to_string(), @@ -210,7 +249,16 @@ pub fn pull_policy_exec_histories( if parsed_responses.is_empty() { break; } - tx.send(parsed_responses.clone()).unwrap(); + match tx.send((cx.clone(), parsed_responses.clone())) { + Ok(_) => {} + Err(e) => { + cx.span().add_event( + "Failed to Send Items to Processor", + vec![KeyValue::new("Response from Processor", e.to_string())], + ); + cx.span().set_status(Status::error("Processor Failed")) + } + } checkpoint_number = parsed_responses.last().unwrap().checkpoint.clone(); if let Some(last_item) = parsed_responses.last() && let Ok(last_date) = NaiveDate::parse_from_str( @@ -263,7 +311,7 @@ fn build_client(headers: HeaderMap) -> Result { .build() } -#[tokio::main] +#[tracing::instrument(name = "history_logging")] async fn history_logging( base_url: &String, exec_types: &String, @@ -292,7 +340,7 @@ async fn history_logging( Ok(res) => { let first_response: ApiResponse = serde_json::from_str(&res.text().await.unwrap()) .expect("Failed to retrieve response from API"); - return first_response; + first_response } Err(_res) => { let failed_response: ApiResponse = ApiResponse { @@ -301,7 +349,7 @@ async fn history_logging( exechistories: vec![], }, }; - return failed_response; + failed_response } } } diff --git a/requirements.txt b/requirements.txt index 6862175..d3c4959 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.1.0 \ No newline at end of file +airlock_libs==7.3.0 \ No newline at end of file