Compare commits

...

2 Commits

Author SHA1 Message Date
brotoskyj 0aabbfd36e Style Change
Build Library / Build Library (push) Successful in 6m25s
Cleaned up services.rs file and removed whitespaces
Specified Data Types were needed instead of allowing the compiler to select types
2025-12-05 11:12:29 -05:00
brotoskyj b19eeb6c96 Fixed Progress Bar
Build Library / Build Library (push) Successful in 6m22s
Changed date math so that the last date from the first response from the API is the minimum value for the progress bar. This gives true progress percentages from the first date to the last date on the last request. #36
2025-12-04 11:57:49 -05:00
6 changed files with 65 additions and 43 deletions
+18 -3
View File
@@ -26,11 +26,13 @@ dependencies = [
[[package]] [[package]]
name = "airlock_libs" name = "airlock_libs"
version = "5.1.0" version = "5.1.2"
dependencies = [ dependencies = [
"chrono", "chrono",
"crossbeam", "crossbeam",
"flexi_logger",
"indicatif", "indicatif",
"log",
"mongodb", "mongodb",
"opentelemetry 0.18.0", "opentelemetry 0.18.0",
"opentelemetry-otlp", "opentelemetry-otlp",
@@ -799,6 +801,19 @@ version = "0.4.2"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0ce7134b9999ecaf8bcd65542e436736ef32ddca1b3e06094cb6ec5755203b80" checksum = "0ce7134b9999ecaf8bcd65542e436736ef32ddca1b3e06094cb6ec5755203b80"
[[package]]
name = "flexi_logger"
version = "0.31.7"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "31e5335674a3a259527f97e9176a3767dcc9b220b8e29d643daeb2d6c72caf8b"
dependencies = [
"chrono",
"log",
"nu-ansi-term",
"regex",
"thiserror 2.0.17",
]
[[package]] [[package]]
name = "fnv" name = "fnv"
version = "1.0.7" version = "1.0.7"
@@ -1590,9 +1605,9 @@ dependencies = [
[[package]] [[package]]
name = "log" name = "log"
version = "0.4.28" version = "0.4.29"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "34080505efa8e45a4b816c349525ebe327ceaa8559756f0356cba97ef3bf7432" checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897"
dependencies = [ dependencies = [
"value-bag", "value-bag",
] ]
+3 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "airlock_libs" name = "airlock_libs"
version = "5.1.0" version = "5.1.2"
edition = "2024" edition = "2024"
[lib] [lib]
@@ -26,6 +26,8 @@ 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" crossbeam = "0.8.4"
log = "0.4.29"
flexi_logger = "0.31.7"
[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.1.0" version = "5.1.2"
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" }
+41 -36
View File
@@ -1,7 +1,5 @@
use std::thread; use std::thread;
use crossbeam::channel::unbounded; use crossbeam::channel::unbounded;
use crate::modules::datatypes::*; use crate::modules::datatypes::*;
use crate::prelude::*; use crate::prelude::*;
#[pyfunction] #[pyfunction]
@@ -16,12 +14,12 @@ pub fn pull_policy_exec_histories(
ExtractedValues::Headers(h) => h, ExtractedValues::Headers(h) => h,
ExtractedValues::BaseUrl(_) => std::process::abort(), ExtractedValues::BaseUrl(_) => std::process::abort(),
}; };
let base_url = match PyData::convert(py, &py_self, false) { let base_url: String = match PyData::convert(py, &py_self, false) {
ExtractedValues::Headers(_) => std::process::abort(), ExtractedValues::Headers(_) => std::process::abort(),
ExtractedValues::BaseUrl(b) => b, ExtractedValues::BaseUrl(b) => b,
}; };
let handle = std::thread::spawn(move || { let handle: thread::JoinHandle<String> = std::thread::spawn(move || {
let rt = match tokio::runtime::Runtime::new() { let rt: tokio::runtime::Runtime = match tokio::runtime::Runtime::new() {
Ok(rt) => rt, Ok(rt) => rt,
Err(e) => { Err(e) => {
println!("Failed to build Tokio Runtime: {:?}", e); println!("Failed to build Tokio Runtime: {:?}", e);
@@ -31,8 +29,8 @@ pub fn pull_policy_exec_histories(
rt.block_on(async { rt.block_on(async {
let _ = init_tracer(); let _ = init_tracer();
}); });
let tracer = global::tracer("global_tracer"); let tracer: global::BoxedTracer = global::tracer("global_tracer");
let _cx = Context::new(); let _cx: Context = Context::new();
let file_path: PathBuf = format!( let file_path: PathBuf = format!(
"{}\\cache\\chunkinator.json", "{}\\cache\\chunkinator.json",
get_base_directory().display() get_base_directory().display()
@@ -58,14 +56,14 @@ pub fn pull_policy_exec_histories(
} }
} }
} }
let data = ApiResponse { let data: ApiResponse = ApiResponse {
error: "Success".to_string(), error: "Success".to_string(),
response: ExecHistories { response: ExecHistories {
exechistories: vec![], exechistories: vec![],
}, },
}; };
let writeable_filepath = file_path.clone(); let writeable_filepath: PathBuf = file_path.clone();
let data_write = serde_json::to_string_pretty(&data).expect("Failed to serialize"); let data_write: String = 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(_) => {}
Err(e) => { Err(e) => {
@@ -74,17 +72,17 @@ pub fn pull_policy_exec_histories(
} }
} }
let mut checkpoint_number: String = SkipBack::find_checkpoint(days).to_string(); let mut checkpoint_number: String = SkipBack::find_checkpoint(days).to_string();
let multi_progress = MultiProgress::new(); let multi_progress: MultiProgress = MultiProgress::new();
multi_progress.set_draw_target(ProgressDrawTarget::stdout()); multi_progress.set_draw_target(ProgressDrawTarget::stderr());
let progress_bar = multi_progress.add(ProgressBar::new(100)); let progress_bar: ProgressBar = 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} {message}") .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));
let client = tracer.in_span("Building HTTP Client", |cx| { let client: Client = tracer.in_span("Building HTTP Client", |cx| {
let client_result = build_client(headers); let client_result: Result<Client, reqwest::Error> = build_client(headers);
match client_result { match client_result {
Ok(client_result) => { Ok(client_result) => {
cx.span().add_event( cx.span().add_event(
@@ -111,12 +109,12 @@ pub fn pull_policy_exec_histories(
} }
} }
}); });
let cutoff = Local::now().naive_local() - Duration::days(days); let cutoff: chrono::NaiveDateTime = Local::now().naive_local() - Duration::days(days);
let (tx, rx) = unbounded::<Vec<Group>>(); let (tx, rx) = unbounded::<Vec<Group>>();
thread::spawn(move || { thread::spawn(move || {
let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists() let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists()
{ {
let contents = fs::read_to_string(&writeable_filepath).unwrap_or_default(); let contents: String = fs::read_to_string(&writeable_filepath).unwrap_or_default();
let existing: ApiResponse = let existing: ApiResponse =
serde_json::from_str(&contents).unwrap_or(ApiResponse { serde_json::from_str(&contents).unwrap_or(ApiResponse {
error: "Success".to_string(), error: "Success".to_string(),
@@ -128,7 +126,7 @@ pub fn pull_policy_exec_histories(
.response .response
.exechistories .exechistories
.into_iter() .into_iter()
.map(|entry| { .map(|entry: Group| {
( (
( (
entry.sha256.clone(), entry.sha256.clone(),
@@ -147,7 +145,7 @@ pub fn pull_policy_exec_histories(
if executions.checkpoint.is_empty() || executions.datetime.is_empty() { if executions.checkpoint.is_empty() || executions.datetime.is_empty() {
continue; continue;
} }
let history_date = match NaiveDate::parse_from_str( let history_date: NaiveDate = 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",
) { ) {
@@ -155,7 +153,7 @@ pub fn pull_policy_exec_histories(
Err(_) => continue, Err(_) => continue,
}; };
if history_date >= cutoff.into() { if history_date >= cutoff.into() {
let key = ( let key: (String, String, String) = (
executions.sha256.clone(), executions.sha256.clone(),
executions.filename.clone(), executions.filename.clone(),
executions.hostname.clone(), executions.hostname.clone(),
@@ -163,13 +161,13 @@ pub fn pull_policy_exec_histories(
seen.entry(key).or_insert(executions.clone()); seen.entry(key).or_insert(executions.clone());
} }
} }
let final_response = ApiResponse { let final_response: ApiResponse = ApiResponse {
error: "Success".to_string(), error: "Success".to_string(),
response: ExecHistories { response: ExecHistories {
exechistories: seen.values().cloned().collect(), exechistories: seen.values().cloned().collect(),
}, },
}; };
let data_write = serde_json::to_string_pretty(&final_response).unwrap(); let data_write: String = serde_json::to_string_pretty(&final_response).unwrap();
match fs::write(&writeable_filepath, data_write) { match fs::write(&writeable_filepath, data_write) {
Ok(_) => {} Ok(_) => {}
Err(e) => { Err(e) => {
@@ -178,8 +176,9 @@ pub fn pull_policy_exec_histories(
} }
} }
}); });
let mut first_date: Option<NaiveDate> = None;
tracer.in_span("Airlock Data Retreival", |cx| { tracer.in_span("Airlock Data Retreival", |cx| {
let span = cx.span(); let span: opentelemetry::trace::SpanRef<'_> = cx.span();
span.set_attribute(Key::new("Days").string(days.to_string())); span.set_attribute(Key::new("Days").string(days.to_string()));
span.set_attribute(KeyValue::new("Policy Name", policy_names.clone())); span.set_attribute(KeyValue::new("Policy Name", policy_names.clone()));
loop { loop {
@@ -197,7 +196,7 @@ pub fn pull_policy_exec_histories(
)); ));
results results
}); });
let parsed_responses = execution_histories.response.exechistories; let parsed_responses: Vec<Group> = execution_histories.response.exechistories;
if parsed_responses.is_empty() { if parsed_responses.is_empty() {
break; break;
} }
@@ -209,16 +208,22 @@ pub fn pull_policy_exec_histories(
"%Y-%m-%dT%H:%M:%SZ", "%Y-%m-%dT%H:%M:%SZ",
) )
{ {
let date_diff = Local::now().naive_local().date() - last_date; if first_date.is_none() {
let percentage_diff = first_date = Some(last_date);
((days - date_diff.num_days()) as f64 / days as f64 * 100.0).round() as u64; }
progress_bar.set_position(percentage_diff); if let Some(base_date) = first_date {
progress_bar.set_message("Total Percent Complete"); let date_diff: chrono::TimeDelta = last_date - base_date;
let total_span: i64 = (Local::now().naive_local().date() - base_date).num_days();
let percentage: u64 = ((date_diff.num_days() as f64 / total_span as f64) * 100.0)
.clamp(0.0, 100.0)
.round() as u64;
progress_bar.set_position(percentage);
}
} }
} }
}); });
progress_bar.finish_with_message("All Checkpoints Complete"); progress_bar.finish_with_message("All Checkpoints Complete");
let return_data = match fs::read_to_string(file_path.clone()) { let return_data: String = 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: {:?}: {}", &file_path, e); println!("Failed to read data from: {:?}: {}", &file_path, e);
@@ -229,8 +234,8 @@ pub fn pull_policy_exec_histories(
shutdown_tracer_provider(); shutdown_tracer_provider();
return_data.to_string() return_data.to_string()
}); });
let gil_value = handle.join().unwrap(); let gil_value: String = handle.join().unwrap();
Python::attach(|py| PyString::new(py, &gil_value).into()) Python::attach(|py: Python<'_>| PyString::new(py, &gil_value).into())
} }
fn build_client(headers: HeaderMap) -> Result<reqwest::Client, reqwest::Error> { fn build_client(headers: HeaderMap) -> Result<reqwest::Client, reqwest::Error> {
@@ -257,7 +262,7 @@ async fn history_logging(
}}"#, }}"#,
exec_types, checkpoint_number, policy_names exec_types, checkpoint_number, policy_names
); );
let res = client let res: Result<reqwest::Response, reqwest::Error> = client
.post(format!("{}/v1/logging/exechistories", base_url)) .post(format!("{}/v1/logging/exechistories", base_url))
.body(payload) .body(payload)
.send() .send()
@@ -298,13 +303,13 @@ pub fn get_base_directory() -> PathBuf {
} }
fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> { fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> {
let cfg = TelemetryConfig::load(); let cfg: TelemetryConfig = TelemetryConfig::load();
if !cfg.TELEMETRY { if !cfg.TELEMETRY {
global::set_tracer_provider(NoopTracerProvider::new()); global::set_tracer_provider(NoopTracerProvider::new());
return Ok(None); return Ok(None);
} }
let endpoint = cfg.TELEM_URL.unwrap_or_default(); let endpoint: String = cfg.TELEM_URL.unwrap_or_default();
let tracer = let tracer: sdktrace::Tracer =
opentelemetry_otlp::new_pipeline() opentelemetry_otlp::new_pipeline()
.tracing() .tracing()
.with_exporter( .with_exporter(
+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.1.0 airlock_libs==5.1.2