Feature added - MPSC channel
Build Library / Build Library (push) Successful in 5m44s

closes #38
This commit is contained in:
brotoskyj
2025-12-03 17:54:57 -05:00
parent 6ab9de413f
commit 76bd3a6087
9 changed files with 115 additions and 88 deletions
+34 -1
View File
@@ -26,9 +26,10 @@ dependencies = [
[[package]] [[package]]
name = "airlock_libs" name = "airlock_libs"
version = "5.0.1" version = "5.1.0"
dependencies = [ dependencies = [
"chrono", "chrono",
"crossbeam",
"indicatif", "indicatif",
"mongodb", "mongodb",
"opentelemetry 0.18.0", "opentelemetry 0.18.0",
@@ -504,6 +505,19 @@ version = "1.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "790eea4361631c5e7d22598ecd5723ff611904e3344ce8720784c93e3d83d40b" 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]] [[package]]
name = "crossbeam-channel" name = "crossbeam-channel"
version = "0.5.15" version = "0.5.15"
@@ -513,6 +527,16 @@ dependencies = [
"crossbeam-utils", "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]] [[package]]
name = "crossbeam-epoch" name = "crossbeam-epoch"
version = "0.9.18" version = "0.9.18"
@@ -522,6 +546,15 @@ dependencies = [
"crossbeam-utils", "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]] [[package]]
name = "crossbeam-utils" name = "crossbeam-utils"
version = "0.8.21" version = "0.8.21"
+2 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "airlock_libs" name = "airlock_libs"
version = "5.0.1" version = "5.1.0"
edition = "2024" edition = "2024"
[lib] [lib]
@@ -25,6 +25,7 @@ tracing = "0.1.41"
tracing-subscriber = "0.3.20" tracing-subscriber = "0.3.20"
tracing-opentelemetry = "0.32.0" tracing-opentelemetry = "0.32.0"
pyo3-async-runtimes = { version = "0.27.0", features = ["async-std", "tokio"] } pyo3-async-runtimes = { version = "0.27.0", features = ["async-std", "tokio"] }
crossbeam = "0.8.4"
[package.metadata.maturin] [package.metadata.maturin]
generate-abi-stubs = true generate-abi-stubs = true
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "maturin"
[project] [project]
name = "airlock_libs" name = "airlock_libs"
version = "5.0.1" version = "5.1.0"
description = "Airlock Digital API Wrapper" description = "Airlock Digital API Wrapper"
readme = "README.md" readme = "README.md"
license = { text = "AGPL-3.0-only" } license = { text = "AGPL-3.0-only" }
+1 -1
View File
@@ -1,7 +1,7 @@
use pyo3::prelude::*; use pyo3::prelude::*;
pub mod modules; pub mod modules;
pub mod services;
pub mod prelude; pub mod prelude;
pub mod services;
#[pymodule] #[pymodule]
fn airlock_libs(py: Python<'_>, m: &Bound<PyModule>) -> PyResult<()> { fn airlock_libs(py: Python<'_>, m: &Bound<PyModule>) -> PyResult<()> {
m.add_function(wrap_pyfunction!(services::pull_policy_exec_histories, py)?)?; m.add_function(wrap_pyfunction!(services::pull_policy_exec_histories, py)?)?;
+1 -1
View File
@@ -109,4 +109,4 @@ impl SkipBack {
let objectid_hex = format!("{}0000000000000000", hex_timestamp); let objectid_hex = format!("{}0000000000000000", hex_timestamp);
ObjectId::parse_str(&objectid_hex).expect("Invalid ObjectId hex") ObjectId::parse_str(&objectid_hex).expect("Invalid ObjectId hex")
} }
} }
+1 -1
View File
@@ -1 +1 @@
pub mod datatypes; pub mod datatypes;
+1 -1
View File
@@ -24,4 +24,4 @@ pub use std::{
io::{Read, Seek, SeekFrom}, io::{Read, Seek, SeekFrom},
path::PathBuf, path::PathBuf,
str::FromStr, str::FromStr,
}; };
+73 -80
View File
@@ -1,5 +1,9 @@
use crate::prelude::*; use std::thread;
use crossbeam::channel::unbounded;
use crate::modules::datatypes::*; use crate::modules::datatypes::*;
use crate::prelude::*;
#[pyfunction] #[pyfunction]
pub fn pull_policy_exec_histories( pub fn pull_policy_exec_histories(
py: Python<'_>, py: Python<'_>,
@@ -34,7 +38,6 @@ pub fn pull_policy_exec_histories(
get_base_directory().display() get_base_directory().display()
) )
.into(); .into();
let writeable_filepath = file_path.clone();
if !&file_path.exists() { if !&file_path.exists() {
if let Some(parent_dir) = &file_path.parent() if let Some(parent_dir) = &file_path.parent()
&& !parent_dir.exists() && !parent_dir.exists()
@@ -61,6 +64,7 @@ pub fn pull_policy_exec_histories(
exechistories: vec![], exechistories: vec![],
}, },
}; };
let writeable_filepath = file_path.clone();
let data_write = serde_json::to_string_pretty(&data).expect("Failed to serialize"); let data_write = serde_json::to_string_pretty(&data).expect("Failed to serialize");
match fs::write(writeable_filepath.clone(), data_write) { match fs::write(writeable_filepath.clone(), data_write) {
Ok(_) => {} Ok(_) => {}
@@ -75,7 +79,7 @@ pub fn pull_policy_exec_histories(
let progress_bar = multi_progress.add(ProgressBar::new(100)); let progress_bar = multi_progress.add(ProgressBar::new(100));
progress_bar.set_style( progress_bar.set_style(
ProgressStyle::default_bar() 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(), .unwrap(),
); );
progress_bar.enable_steady_tick(std::time::Duration::from_millis(100)); 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 cutoff = Local::now().naive_local() - Duration::days(days);
let mut f = match File::open(&writeable_filepath) { let (tx, rx) = unbounded::<Vec<Group>>();
Ok(f) => f, thread::spawn(move || {
Err(e) => { let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists()
println!("Failed to Access {:?}: {}", &writeable_filepath, e); {
std::process::abort(); let contents = fs::read_to_string(&writeable_filepath).unwrap_or_default();
} let existing: ApiResponse =
}; serde_json::from_str(&contents).unwrap_or(ApiResponse {
tracer.in_span("Airlock Data Retreival", |cx| { error: "Success".to_string(),
let span = cx.span(); response: ExecHistories {
span.set_attribute(Key::new("Days").string(days.to_string())); exechistories: vec![],
span.set_attribute(KeyValue::new("Policy Name", policy_names.clone())); },
loop { });
match f.seek(SeekFrom::Start(0)) { existing
Ok(_) => {} .response
Err(e) => { .exechistories
println!("Failed to seek start of {:?}: {}", f, e); .into_iter()
std::process::abort(); .map(|entry| {
} (
} (
let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| { entry.sha256.clone(),
let results: ApiResponse = history_logging( entry.filename.clone(),
&base_url, entry.hostname.clone(),
&exec_types, ),
&checkpoint_number, entry,
&policy_names, )
&client, })
); .collect()
cx.span().set_attribute(KeyValue::new( } else {
"items_in_response", HashMap::new()
results.response.exechistories.len().to_string(), };
)); while let Ok(parsed_responses) = rx.recv() {
results for executions in parsed_responses {
});
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() {
if executions.checkpoint.is_empty() || executions.datetime.is_empty() { if executions.checkpoint.is_empty() || executions.datetime.is_empty() {
continue; continue;
} }
if index == parsed_responses.len() - 1 {
checkpoint_number = executions.checkpoint.clone();
break;
}
let history_date = match NaiveDate::parse_from_str( let history_date = match NaiveDate::parse_from_str(
&executions.datetime.replace(" +0000 UTC", ""), &executions.datetime.replace(" +0000 UTC", ""),
"%Y-%m-%dT%H:%M:%SZ", "%Y-%m-%dT%H:%M:%SZ",
@@ -211,7 +176,34 @@ pub fn pull_policy_exec_histories(
println!("Failed to write to: {:?}: {}", &writeable_filepath, e); 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( && let Ok(last_date) = NaiveDate::parse_from_str(
&last_item.datetime.replace(" +0000 UTC", ""), &last_item.datetime.replace(" +0000 UTC", ""),
"%Y-%m-%dT%H:%M:%SZ", "%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 date_diff = Local::now().naive_local().date() - last_date;
let percentage_diff = let percentage_diff =
(days - date_diff.num_days()) as f64 / days as f64 * 100.0; ((days - date_diff.num_days()) as f64 / days as f64 * 100.0).round() as u64;
progress_bar.set_position(percentage_diff.round() as u64); progress_bar.set_position(percentage_diff);
progress_bar.set_message("Total Percent Complete"); progress_bar.set_message("Total Percent Complete");
} }
} }
}); });
progress_bar.finish_with_message("All Checkpoints 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, Ok(return_data) => return_data,
Err(e) => { Err(e) => {
println!("Failed to read data from: {:?}: {}", &writeable_filepath, e); println!("Failed to read data from: {:?}: {}", &file_path, e);
std::process::abort(); std::process::abort();
} }
}; };
drop(tx);
shutdown_tracer_provider(); shutdown_tracer_provider();
return_data.to_string() return_data.to_string()
}); });
@@ -325,4 +318,4 @@ fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> {
.install_simple() .install_simple()
.unwrap(); .unwrap();
Ok(Some(tracer)) Ok(Some(tracer))
} }
+1 -1
View File
@@ -11,4 +11,4 @@ urllib3==2.5.0
pyperclip==1.11.0 pyperclip==1.11.0
--extra-index-url https://git.racooncity.org/api/packages/brotoskyj/pypi/simple/ --extra-index-url https://git.racooncity.org/api/packages/brotoskyj/pypi/simple/
airlock_libs==5.0.1 airlock_libs==5.1.0