diff --git a/airlock_libs/Cargo.lock b/airlock_libs/Cargo.lock index 66728f0..ca2d548 100644 --- a/airlock_libs/Cargo.lock +++ b/airlock_libs/Cargo.lock @@ -26,9 +26,10 @@ dependencies = [ [[package]] name = "airlock_libs" -version = "5.0.1" +version = "5.1.0" dependencies = [ "chrono", + "crossbeam", "indicatif", "mongodb", "opentelemetry 0.18.0", @@ -504,6 +505,19 @@ version = "1.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" +[[package]] +name = "crossbeam" +version = "0.8.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1137cd7e7fc0fb5d3c5a8678be38ec56e819125d8d7907411fe24ccb943faca8" +dependencies = [ + "crossbeam-channel", + "crossbeam-deque", + "crossbeam-epoch", + "crossbeam-queue", + "crossbeam-utils", +] + [[package]] name = "crossbeam-channel" version = "0.5.15" @@ -513,6 +527,16 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "crossbeam-deque" +version = "0.8.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9dd111b7b7f7d55b72c0a6ae361660ee5853c9af73f70c3c2ef6858b950e2e51" +dependencies = [ + "crossbeam-epoch", + "crossbeam-utils", +] + [[package]] name = "crossbeam-epoch" version = "0.9.18" @@ -522,6 +546,15 @@ dependencies = [ "crossbeam-utils", ] +[[package]] +name = "crossbeam-queue" +version = "0.3.12" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "0f58bbc28f91df819d0aa2a2c00cd19754769c2fad90579b3592b1c9ba7a3115" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crossbeam-utils" version = "0.8.21" diff --git a/airlock_libs/Cargo.toml b/airlock_libs/Cargo.toml index 5704114..36647a0 100644 --- a/airlock_libs/Cargo.toml +++ b/airlock_libs/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "airlock_libs" -version = "5.0.1" +version = "5.1.0" edition = "2024" [lib] @@ -25,6 +25,7 @@ tracing = "0.1.41" tracing-subscriber = "0.3.20" tracing-opentelemetry = "0.32.0" pyo3-async-runtimes = { version = "0.27.0", features = ["async-std", "tokio"] } +crossbeam = "0.8.4" [package.metadata.maturin] generate-abi-stubs = true diff --git a/airlock_libs/pyproject.toml b/airlock_libs/pyproject.toml index 9ceacbc..02f03a9 100644 --- a/airlock_libs/pyproject.toml +++ b/airlock_libs/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "maturin" [project] name = "airlock_libs" -version = "5.0.1" +version = "5.1.0" description = "Airlock Digital API Wrapper" readme = "README.md" license = { text = "AGPL-3.0-only" } diff --git a/airlock_libs/src/lib.rs b/airlock_libs/src/lib.rs index b991583..a3b6e3d 100644 --- a/airlock_libs/src/lib.rs +++ b/airlock_libs/src/lib.rs @@ -1,7 +1,7 @@ use pyo3::prelude::*; pub mod modules; -pub mod services; pub mod prelude; +pub mod services; #[pymodule] fn airlock_libs(py: Python<'_>, m: &Bound) -> PyResult<()> { m.add_function(wrap_pyfunction!(services::pull_policy_exec_histories, py)?)?; diff --git a/airlock_libs/src/modules/datatypes.rs b/airlock_libs/src/modules/datatypes.rs index c27af30..5df65f6 100644 --- a/airlock_libs/src/modules/datatypes.rs +++ b/airlock_libs/src/modules/datatypes.rs @@ -109,4 +109,4 @@ impl SkipBack { let objectid_hex = format!("{}0000000000000000", hex_timestamp); ObjectId::parse_str(&objectid_hex).expect("Invalid ObjectId hex") } -} \ No newline at end of file +} diff --git a/airlock_libs/src/modules/mod.rs b/airlock_libs/src/modules/mod.rs index 58fd615..0d8681d 100644 --- a/airlock_libs/src/modules/mod.rs +++ b/airlock_libs/src/modules/mod.rs @@ -1 +1 @@ -pub mod datatypes; \ No newline at end of file +pub mod datatypes; diff --git a/airlock_libs/src/prelude.rs b/airlock_libs/src/prelude.rs index faa125a..8d87356 100644 --- a/airlock_libs/src/prelude.rs +++ b/airlock_libs/src/prelude.rs @@ -24,4 +24,4 @@ pub use std::{ io::{Read, Seek, SeekFrom}, path::PathBuf, str::FromStr, -}; \ No newline at end of file +}; diff --git a/airlock_libs/src/services.rs b/airlock_libs/src/services.rs index e38c7e0..ce95090 100644 --- a/airlock_libs/src/services.rs +++ b/airlock_libs/src/services.rs @@ -1,5 +1,9 @@ -use crate::prelude::*; +use std::thread; + +use crossbeam::channel::unbounded; + use crate::modules::datatypes::*; +use crate::prelude::*; #[pyfunction] pub fn pull_policy_exec_histories( py: Python<'_>, @@ -34,7 +38,6 @@ pub fn pull_policy_exec_histories( get_base_directory().display() ) .into(); - let writeable_filepath = file_path.clone(); if !&file_path.exists() { if let Some(parent_dir) = &file_path.parent() && !parent_dir.exists() @@ -61,6 +64,7 @@ pub fn pull_policy_exec_histories( exechistories: vec![], }, }; + let writeable_filepath = file_path.clone(); let data_write = serde_json::to_string_pretty(&data).expect("Failed to serialize"); match fs::write(writeable_filepath.clone(), data_write) { Ok(_) => {} @@ -75,7 +79,7 @@ pub fn pull_policy_exec_histories( let progress_bar = multi_progress.add(ProgressBar::new(100)); progress_bar.set_style( ProgressStyle::default_bar() - .template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len}") + .template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len} {message}") .unwrap(), ); progress_bar.enable_steady_tick(std::time::Duration::from_millis(100)); @@ -108,80 +112,41 @@ pub fn pull_policy_exec_histories( } }); let cutoff = Local::now().naive_local() - Duration::days(days); - let mut f = match File::open(&writeable_filepath) { - Ok(f) => f, - Err(e) => { - println!("Failed to Access {:?}: {}", &writeable_filepath, e); - std::process::abort(); - } - }; - tracer.in_span("Airlock Data Retreival", |cx| { - let span = cx.span(); - span.set_attribute(Key::new("Days").string(days.to_string())); - span.set_attribute(KeyValue::new("Policy Name", policy_names.clone())); - loop { - match f.seek(SeekFrom::Start(0)) { - Ok(_) => {} - Err(e) => { - println!("Failed to seek start of {:?}: {}", f, e); - std::process::abort(); - } - } - let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| { - let results: ApiResponse = history_logging( - &base_url, - &exec_types, - &checkpoint_number, - &policy_names, - &client, - ); - cx.span().set_attribute(KeyValue::new( - "items_in_response", - results.response.exechistories.len().to_string(), - )); - results - }); - let parsed_responses = execution_histories.response.exechistories; - if parsed_responses.is_empty() { - break; - } - let mut seen: HashMap<(String, String, String), Group> = - if writeable_filepath.exists() { - let mut contents = String::new(); - f.read_to_string(&mut contents).unwrap(); - let existing_data: ApiResponse = - serde_json::from_str(&contents).unwrap_or(ApiResponse { - error: "Success".to_string(), - response: ExecHistories { - exechistories: vec![], - }, - }); - existing_data - .response - .exechistories - .into_iter() - .map(|entry| { - ( - ( - entry.sha256.clone(), - entry.filename.clone(), - entry.hostname.clone(), - ), - entry, - ) - }) - .collect() - } else { - HashMap::new() - }; - for (index, executions) in parsed_responses.iter().enumerate() { + let (tx, rx) = unbounded::>(); + thread::spawn(move || { + let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists() + { + let contents = fs::read_to_string(&writeable_filepath).unwrap_or_default(); + let existing: ApiResponse = + serde_json::from_str(&contents).unwrap_or(ApiResponse { + error: "Success".to_string(), + response: ExecHistories { + exechistories: vec![], + }, + }); + existing + .response + .exechistories + .into_iter() + .map(|entry| { + ( + ( + entry.sha256.clone(), + entry.filename.clone(), + entry.hostname.clone(), + ), + entry, + ) + }) + .collect() + } else { + HashMap::new() + }; + while let Ok(parsed_responses) = rx.recv() { + for executions in parsed_responses { if executions.checkpoint.is_empty() || executions.datetime.is_empty() { continue; } - if index == parsed_responses.len() - 1 { - checkpoint_number = executions.checkpoint.clone(); - break; - } let history_date = match NaiveDate::parse_from_str( &executions.datetime.replace(" +0000 UTC", ""), "%Y-%m-%dT%H:%M:%SZ", @@ -211,7 +176,34 @@ pub fn pull_policy_exec_histories( println!("Failed to write to: {:?}: {}", &writeable_filepath, e); } } - if let Some(last_item) = &final_response.response.exechistories.last() + } + }); + tracer.in_span("Airlock Data Retreival", |cx| { + let span = cx.span(); + span.set_attribute(Key::new("Days").string(days.to_string())); + span.set_attribute(KeyValue::new("Policy Name", policy_names.clone())); + loop { + let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| { + let results: ApiResponse = history_logging( + &base_url, + &exec_types, + &checkpoint_number, + &policy_names, + &client, + ); + cx.span().set_attribute(KeyValue::new( + "items_in_response", + results.response.exechistories.len().to_string(), + )); + results + }); + let parsed_responses = execution_histories.response.exechistories; + if parsed_responses.is_empty() { + break; + } + tx.send(parsed_responses.clone()).unwrap(); + checkpoint_number = parsed_responses.last().unwrap().checkpoint.clone(); + if let Some(last_item) = parsed_responses.last() && let Ok(last_date) = NaiveDate::parse_from_str( &last_item.datetime.replace(" +0000 UTC", ""), "%Y-%m-%dT%H:%M:%SZ", @@ -219,20 +211,21 @@ pub fn pull_policy_exec_histories( { let date_diff = Local::now().naive_local().date() - last_date; let percentage_diff = - (days - date_diff.num_days()) as f64 / days as f64 * 100.0; - progress_bar.set_position(percentage_diff.round() as u64); + ((days - date_diff.num_days()) as f64 / days as f64 * 100.0).round() as u64; + progress_bar.set_position(percentage_diff); progress_bar.set_message("Total Percent Complete"); } } }); progress_bar.finish_with_message("All Checkpoints Complete"); - let return_data = match fs::read_to_string(&writeable_filepath) { + let return_data = match fs::read_to_string(file_path.clone()) { Ok(return_data) => return_data, Err(e) => { - println!("Failed to read data from: {:?}: {}", &writeable_filepath, e); + println!("Failed to read data from: {:?}: {}", &file_path, e); std::process::abort(); } }; + drop(tx); shutdown_tracer_provider(); return_data.to_string() }); @@ -325,4 +318,4 @@ fn init_tracer() -> Result, TraceError> { .install_simple() .unwrap(); Ok(Some(tracer)) -} \ No newline at end of file +} diff --git a/requirements.txt b/requirements.txt index 7c0d5ff..94ee402 100644 --- a/requirements.txt +++ b/requirements.txt @@ -11,4 +11,4 @@ urllib3==2.5.0 pyperclip==1.11.0 --extra-index-url https://git.racooncity.org/api/packages/brotoskyj/pypi/simple/ -airlock_libs==5.0.1 \ No newline at end of file +airlock_libs==5.1.0 \ No newline at end of file