Compare commits

...

1 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
6 changed files with 55 additions and 40 deletions
+18 -3
View File
@@ -26,11 +26,13 @@ dependencies = [
[[package]] [[package]]
name = "airlock_libs" name = "airlock_libs"
version = "5.1.1" 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.1" 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.1" 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" }
+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")
} }
} }
+31 -33
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::stderr()); 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) => {
@@ -180,7 +178,7 @@ pub fn pull_policy_exec_histories(
}); });
let mut first_date: Option<NaiveDate> = None; 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 {
@@ -198,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;
} }
@@ -214,9 +212,9 @@ pub fn pull_policy_exec_histories(
first_date = Some(last_date); first_date = Some(last_date);
} }
if let Some(base_date) = first_date { if let Some(base_date) = first_date {
let date_diff = last_date - base_date; let date_diff: chrono::TimeDelta = last_date - base_date;
let total_span = (Local::now().naive_local().date() - base_date).num_days(); let total_span: i64 = (Local::now().naive_local().date() - base_date).num_days();
let percentage = ((date_diff.num_days() as f64 / total_span as f64) * 100.0) let percentage: u64 = ((date_diff.num_days() as f64 / total_span as f64) * 100.0)
.clamp(0.0, 100.0) .clamp(0.0, 100.0)
.round() as u64; .round() as u64;
progress_bar.set_position(percentage); progress_bar.set_position(percentage);
@@ -225,7 +223,7 @@ pub fn pull_policy_exec_histories(
} }
}); });
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);
@@ -236,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> {
@@ -264,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()
@@ -305,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.1 airlock_libs==5.1.2