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]]
name = "airlock_libs"
version = "5.1.0"
version = "5.1.2"
dependencies = [
"chrono",
"crossbeam",
"flexi_logger",
"indicatif",
"log",
"mongodb",
"opentelemetry 0.18.0",
"opentelemetry-otlp",
@@ -799,6 +801,19 @@ version = "0.4.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
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]]
name = "fnv"
version = "1.0.7"
@@ -1590,9 +1605,9 @@ dependencies = [
[[package]]
name = "log"
version = "0.4.28"
version = "0.4.29"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "34080505efa8e45a4b816c349525ebe327ceaa8559756f0356cba97ef3bf7432"
checksum = "5e5032e24019045c762d3c0f28f5b6b8bbf38563a65908389bf7978758920897"
dependencies = [
"value-bag",
]
+3 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "airlock_libs"
version = "5.1.0"
version = "5.1.2"
edition = "2024"
[lib]
@@ -26,6 +26,8 @@ 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"
log = "0.4.29"
flexi_logger = "0.31.7"
[package.metadata.maturin]
generate-abi-stubs = true
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "maturin"
[project]
name = "airlock_libs"
version = "5.1.0"
version = "5.1.2"
description = "Airlock Digital API Wrapper"
readme = "README.md"
license = { text = "AGPL-3.0-only" }
+41 -36
View File
@@ -1,7 +1,5 @@
use std::thread;
use crossbeam::channel::unbounded;
use crate::modules::datatypes::*;
use crate::prelude::*;
#[pyfunction]
@@ -16,12 +14,12 @@ pub fn pull_policy_exec_histories(
ExtractedValues::Headers(h) => h,
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::BaseUrl(b) => b,
};
let handle = std::thread::spawn(move || {
let rt = match tokio::runtime::Runtime::new() {
let handle: thread::JoinHandle<String> = std::thread::spawn(move || {
let rt: tokio::runtime::Runtime = match tokio::runtime::Runtime::new() {
Ok(rt) => rt,
Err(e) => {
println!("Failed to build Tokio Runtime: {:?}", e);
@@ -31,8 +29,8 @@ pub fn pull_policy_exec_histories(
rt.block_on(async {
let _ = init_tracer();
});
let tracer = global::tracer("global_tracer");
let _cx = Context::new();
let tracer: global::BoxedTracer = global::tracer("global_tracer");
let _cx: Context = Context::new();
let file_path: PathBuf = format!(
"{}\\cache\\chunkinator.json",
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(),
response: ExecHistories {
exechistories: vec![],
},
};
let writeable_filepath = file_path.clone();
let data_write = serde_json::to_string_pretty(&data).expect("Failed to serialize");
let writeable_filepath: PathBuf = file_path.clone();
let data_write: String = serde_json::to_string_pretty(&data).expect("Failed to serialize");
match fs::write(writeable_filepath.clone(), data_write) {
Ok(_) => {}
Err(e) => {
@@ -74,17 +72,17 @@ pub fn pull_policy_exec_histories(
}
}
let mut checkpoint_number: String = SkipBack::find_checkpoint(days).to_string();
let multi_progress = MultiProgress::new();
multi_progress.set_draw_target(ProgressDrawTarget::stdout());
let progress_bar = multi_progress.add(ProgressBar::new(100));
let multi_progress: MultiProgress = MultiProgress::new();
multi_progress.set_draw_target(ProgressDrawTarget::stderr());
let progress_bar: ProgressBar = 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} {message}")
.unwrap(),
);
progress_bar.enable_steady_tick(std::time::Duration::from_millis(100));
let client = tracer.in_span("Building HTTP Client", |cx| {
let client_result = build_client(headers);
let client: Client = tracer.in_span("Building HTTP Client", |cx| {
let client_result: Result<Client, reqwest::Error> = build_client(headers);
match client_result {
Ok(client_result) => {
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>>();
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 contents: String = fs::read_to_string(&writeable_filepath).unwrap_or_default();
let existing: ApiResponse =
serde_json::from_str(&contents).unwrap_or(ApiResponse {
error: "Success".to_string(),
@@ -128,7 +126,7 @@ pub fn pull_policy_exec_histories(
.response
.exechistories
.into_iter()
.map(|entry| {
.map(|entry: Group| {
(
(
entry.sha256.clone(),
@@ -147,7 +145,7 @@ pub fn pull_policy_exec_histories(
if executions.checkpoint.is_empty() || executions.datetime.is_empty() {
continue;
}
let history_date = match NaiveDate::parse_from_str(
let history_date: NaiveDate = match NaiveDate::parse_from_str(
&executions.datetime.replace(" +0000 UTC", ""),
"%Y-%m-%dT%H:%M:%SZ",
) {
@@ -155,7 +153,7 @@ pub fn pull_policy_exec_histories(
Err(_) => continue,
};
if history_date >= cutoff.into() {
let key = (
let key: (String, String, String) = (
executions.sha256.clone(),
executions.filename.clone(),
executions.hostname.clone(),
@@ -163,13 +161,13 @@ pub fn pull_policy_exec_histories(
seen.entry(key).or_insert(executions.clone());
}
}
let final_response = ApiResponse {
let final_response: ApiResponse = ApiResponse {
error: "Success".to_string(),
response: ExecHistories {
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) {
Ok(_) => {}
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| {
let span = cx.span();
let span: opentelemetry::trace::SpanRef<'_> = cx.span();
span.set_attribute(Key::new("Days").string(days.to_string()));
span.set_attribute(KeyValue::new("Policy Name", policy_names.clone()));
loop {
@@ -197,7 +196,7 @@ pub fn pull_policy_exec_histories(
));
results
});
let parsed_responses = execution_histories.response.exechistories;
let parsed_responses: Vec<Group> = execution_histories.response.exechistories;
if parsed_responses.is_empty() {
break;
}
@@ -209,16 +208,22 @@ pub fn pull_policy_exec_histories(
"%Y-%m-%dT%H:%M:%SZ",
)
{
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).round() as u64;
progress_bar.set_position(percentage_diff);
progress_bar.set_message("Total Percent Complete");
if first_date.is_none() {
first_date = Some(last_date);
}
if let Some(base_date) = first_date {
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");
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,
Err(e) => {
println!("Failed to read data from: {:?}: {}", &file_path, e);
@@ -229,8 +234,8 @@ pub fn pull_policy_exec_histories(
shutdown_tracer_provider();
return_data.to_string()
});
let gil_value = handle.join().unwrap();
Python::attach(|py| PyString::new(py, &gil_value).into())
let gil_value: String = handle.join().unwrap();
Python::attach(|py: Python<'_>| PyString::new(py, &gil_value).into())
}
fn build_client(headers: HeaderMap) -> Result<reqwest::Client, reqwest::Error> {
@@ -257,7 +262,7 @@ async fn history_logging(
}}"#,
exec_types, checkpoint_number, policy_names
);
let res = client
let res: Result<reqwest::Response, reqwest::Error> = client
.post(format!("{}/v1/logging/exechistories", base_url))
.body(payload)
.send()
@@ -298,13 +303,13 @@ pub fn get_base_directory() -> PathBuf {
}
fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> {
let cfg = TelemetryConfig::load();
let cfg: TelemetryConfig = TelemetryConfig::load();
if !cfg.TELEMETRY {
global::set_tracer_provider(NoopTracerProvider::new());
return Ok(None);
}
let endpoint = cfg.TELEM_URL.unwrap_or_default();
let tracer =
let endpoint: String = cfg.TELEM_URL.unwrap_or_default();
let tracer: sdktrace::Tracer =
opentelemetry_otlp::new_pipeline()
.tracing()
.with_exporter(
+1 -1
View File
@@ -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.1.0
airlock_libs==5.1.2