Compare commits

...

8 Commits

Author SHA1 Message Date
brotoskyj f3c1d97d28 Implemented Threaded Rust Process
Build Library / Build Library (push) Successful in 6m18s
It still locks the GIL, but this is the ground work for the non blocking GUI. Python objects are converted to Rust objects in the beginning of the function call so they can be sent to threads safely.
2025-12-02 10:53:25 -05:00
brotoskyj 99d7b5f74e Merge remote-tracking branch 'refs/remotes/origin/RustImplementation' into RustImplementation 2025-11-25 15:30:07 -05:00
brotoskyj 6327adeabf Refactor of Loxide Libs
Removed spaces, imports, and converters from code
2025-11-25 15:29:55 -05:00
brotoskyj b289e9324c Slight Changes
Build Library / Build Library (push) Successful in 5m15s
Refactor of Loxide Libs
Removed spaces, imports, and converters from code
2025-11-25 15:27:16 -05:00
brotoskyj a80c2ca1e1 Features
Build Library / Build Library (push) Successful in 4m46s
closes #34
Removed a lot of unwraps, most remaining unwraps will likely stay in the code, as it will be expected behavior to panic and crash rather than a system level abort/exit event
2025-11-24 11:36:50 -05:00
brotoskyj dd206eb272 Feature
closes #37
2025-11-24 10:35:38 -05:00
brotoskyj 87ab1e3b28 Telemetry Config Change
Build Library / Build Library (push) Successful in 4m54s
Added implementation instead of separate function to load Telemetry configuration
2025-11-21 14:00:09 -05:00
brotoskyj 303ecd8368 Removed a lot of unwraps
Build Library / Build Library (push) Successful in 5m2s
Removing the unwraps now has better error handling and will cause the program to crash with crash dumps instead of panic!
2025-11-21 12:20:11 -05:00
6 changed files with 820 additions and 499 deletions
+1
View File
@@ -1,3 +1,4 @@
/target /target
build.sh build.sh
pythontest.py pythontest.py
changelog.md
+536 -297
View File
File diff suppressed because it is too large Load Diff
+2 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "airlock_libs" name = "airlock_libs"
version = "3.1.2" version = "5.0.0"
edition = "2024" edition = "2024"
[lib] [lib]
@@ -24,6 +24,7 @@ tonic = { version = "0.8.2", features = ["tls-roots"] }
tracing = "0.1.41" 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"] }
[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 = "3.1.2" version = "5.0.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" }
+278 -198
View File
@@ -9,6 +9,7 @@ use opentelemetry::{Context, KeyValue, sdk::trace as sdktrace, trace::Tracer};
use opentelemetry::{Key, global}; use opentelemetry::{Key, global};
use opentelemetry_otlp::WithExportConfig; use opentelemetry_otlp::WithExportConfig;
use pyo3::{prelude::*, types::PyString}; use pyo3::{prelude::*, types::PyString};
use pyo3_async_runtimes::async_std;
use reqwest::{ use reqwest::{
Client, Client,
header::{HeaderMap, HeaderName, HeaderValue}, header::{HeaderMap, HeaderName, HeaderValue},
@@ -25,12 +26,34 @@ use std::{
str::FromStr, str::FromStr,
}; };
#[allow(non_snake_case)]
#[derive(Deserialize, Debug)] #[derive(Deserialize, Debug)]
struct TelemetryConfig { struct TelemetryConfig {
TELEMETRY: bool, TELEMETRY: bool,
TELEM_URL: Option<String>, TELEM_URL: Option<String>,
} }
impl TelemetryConfig {
pub fn load() -> Self {
let cfg_path = get_base_directory().join("config\\user_config.json");
if !cfg_path.exists() {
return Self {
TELEMETRY: false,
TELEM_URL: None,
};
}
match fs::read_to_string(&cfg_path) {
Ok(contents) => serde_json::from_str::<Self>(&contents).unwrap_or(Self {
TELEMETRY: false,
TELEM_URL: None,
}),
Err(_) => Self {
TELEMETRY: false,
TELEM_URL: None,
},
}
}
}
#[derive(Debug, Deserialize, Serialize)] #[derive(Debug, Deserialize, Serialize)]
struct ApiResponse { struct ApiResponse {
error: String, error: String,
@@ -66,6 +89,40 @@ struct Group {
localip: String, localip: String,
} }
enum ExtractedValues {
Headers(reqwest::header::HeaderMap),
BaseUrl(String),
}
trait Converter {
fn convert(py: Python<'_>, py_self: &Py<PyAny>, extract_headers: bool) -> ExtractedValues;
}
struct PyData;
impl Converter for PyData {
fn convert(py: Python<'_>, py_self: &Py<PyAny>, extract_headers: bool) -> ExtractedValues {
if extract_headers {
let headers = py_self.getattr(py, "headers").unwrap().to_string();
let headers_replace = headers.replace('\'', "\"");
let parsed: Value = serde_json::from_str(headers_replace.as_str()).unwrap();
let mut header_map = HeaderMap::new();
if let Some(obj) = parsed.as_object() {
for (_key, value) in obj {
if let Some(v) = value.as_str() {
let val = HeaderValue::from_str(v).unwrap();
header_map.insert(HeaderName::from_str("X-APIKey").unwrap(), val);
}
}
}
ExtractedValues::Headers(header_map)
} else {
let base_url = py_self.getattr(py, "base_url").unwrap().to_string();
ExtractedValues::BaseUrl(base_url)
}
}
}
#[pyfunction] #[pyfunction]
pub fn pull_policy_exec_histories( pub fn pull_policy_exec_histories(
py: Python<'_>, py: Python<'_>,
@@ -74,209 +131,254 @@ pub fn pull_policy_exec_histories(
exec_types: String, exec_types: String,
days: i64, days: i64,
) -> Py<PyString> { ) -> Py<PyString> {
let rt = tokio::runtime::Runtime::new().unwrap(); let headers: HeaderMap = match PyData::convert(py, &py_self, true) {
rt.block_on(async { ExtractedValues::Headers(h) => h,
let _ = init_tracer(); ExtractedValues::BaseUrl(_) => std::process::abort(),
});
let tracer = global::tracer("global_tracer");
let _cx = Context::new();
let file_path: PathBuf = format!(
"{}\\cache\\chunkinator.json",
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()
{
fs::create_dir_all(parent_dir).unwrap();
}
fs::File::create(file_path).unwrap();
}
let data = ApiResponse {
error: "Success".to_string(),
response: ExecHistories {
exechistories: vec![],
},
}; };
let data_write = serde_json::to_string_pretty(&data).expect("Failed to serialize"); let base_url = match PyData::convert(py, &py_self, false) {
fs::write(writeable_filepath.clone(), data_write).unwrap(); ExtractedValues::Headers(_) => std::process::abort(),
let mut checkpoint_number: String = skipback(days).to_string(); ExtractedValues::BaseUrl(b) => b,
let multi_progress = MultiProgress::new(); };
multi_progress.set_draw_target(ProgressDrawTarget::stdout()); let handle = std::thread::spawn(move || {
let progress_bar = multi_progress.add(ProgressBar::new(100)); let rt = match tokio::runtime::Runtime::new() {
progress_bar.set_style( Ok(rt) => rt,
ProgressStyle::default_bar() Err(e) => {
.template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len}") println!("Failed to build Tokio Runtime: {:?}", e);
.unwrap(), std::process::abort();
);
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(py, &py_self);
match client_result {
Ok(client_result) => {
cx.span().add_event(
"info",
vec![KeyValue::new(
"Client Built Successfully",
format!("{:?}", client_result),
)],
);
client_result
} }
Err(client_result) => { };
cx.span().add_event( rt.block_on(async {
"warn", let _ = init_tracer();
vec![KeyValue::new( });
"Client Failed to Build", let tracer = global::tracer("global_tracer");
format!("{:?}", &client_result), let _cx = Context::new();
)], let file_path: PathBuf = format!(
); "{}\\cache\\chunkinator.json",
cx.span() get_base_directory().display()
.set_status(Status::error("Client Failed to Build")); )
panic!("Failed to Build Client: {:?}", client_result); .into();
let writeable_filepath = file_path.clone();
if !&file_path.exists() {
if let Some(parent_dir) = &file_path.parent()
&& !parent_dir.exists()
{
match fs::create_dir_all(parent_dir) {
Ok(_) => {}
Err(e) => {
println!("Failed to Create Directory {:?}: {}", parent_dir, e);
std::process::abort();
}
}
}
match fs::File::create(&file_path) {
Ok(_) => {}
Err(e) => {
println!("Failed to Create Directory {:?}: {}", &file_path, e);
std::process::abort();
}
} }
} }
}); let data = ApiResponse {
let api: Py<PyAny> = py_self; error: "Success".to_string(),
let cutoff = Local::now().naive_local() - Duration::days(days); response: ExecHistories {
let mut f = File::open(&writeable_filepath).unwrap(); exechistories: vec![],
tracer.in_span("Airlock Data Retreival", |cx| { },
let span = cx.span(); };
span.set_attribute(Key::new("Days").string(days.to_string().to_string())); let data_write = serde_json::to_string_pretty(&data).expect("Failed to serialize");
loop { match fs::write(writeable_filepath.clone(), data_write) {
f.seek(SeekFrom::Start(0)).unwrap(); Ok(_) => {}
let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| { Err(e) => {
let results: ApiResponse = history_logging( println!("Failed to write to: {:?}: {}", &writeable_filepath, e);
py, std::process::abort();
&api,
&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 checkpoint_number: String = skipback(days).to_string();
let mut contents = String::new(); let multi_progress = MultiProgress::new();
f.read_to_string(&mut contents).unwrap(); multi_progress.set_draw_target(ProgressDrawTarget::stdout());
let existing_data: ApiResponse = let progress_bar = multi_progress.add(ProgressBar::new(100));
serde_json::from_str(&contents).unwrap_or(ApiResponse { progress_bar.set_style(
error: "Success".to_string(), ProgressStyle::default_bar()
response: ExecHistories { .template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len}")
exechistories: vec![], .unwrap(),
}, );
}); progress_bar.enable_steady_tick(std::time::Duration::from_millis(100));
existing_data let client = tracer.in_span("Building HTTP Client", |cx| {
.response let client_result = build_client(headers);
.exechistories match client_result {
.into_iter() Ok(client_result) => {
.map(|entry| { cx.span().add_event(
( "info",
( vec![KeyValue::new(
entry.sha256.clone(), "Client Built Successfully",
entry.filename.clone(), format!("{:?}", client_result),
entry.hostname.clone(), )],
), );
entry, client_result
)
})
.collect()
} else {
HashMap::new()
};
for (index, executions) in parsed_responses.iter().enumerate() {
if executions.checkpoint.is_empty() || executions.datetime.is_empty() {
continue;
} }
if index == parsed_responses.len() - 1 { Err(client_result) => {
checkpoint_number = executions.checkpoint.clone(); cx.span().add_event(
"warn",
vec![KeyValue::new(
"Client Failed to Build",
format!("{:?}", &client_result),
)],
);
cx.span()
.set_status(Status::error("Client Failed to Build"));
println!("Failed to Build Client: {:?}", client_result);
std::process::abort();
}
}
});
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; break;
} }
let history_date = match NaiveDate::parse_from_str( let mut seen: HashMap<(String, String, String), Group> =
&executions.datetime.replace(" +0000 UTC", ""), if writeable_filepath.exists() {
"%Y-%m-%dT%H:%M:%SZ", let mut contents = String::new();
) { f.read_to_string(&mut contents).unwrap();
Ok(date) => date, let existing_data: ApiResponse =
Err(_) => continue, 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() {
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",
) {
Ok(date) => date,
Err(_) => continue,
};
if history_date >= cutoff.into() {
let key = (
executions.sha256.clone(),
executions.filename.clone(),
executions.hostname.clone(),
);
seen.entry(key).or_insert(executions.clone());
}
}
let final_response = ApiResponse {
error: "Success".to_string(),
response: ExecHistories {
exechistories: seen.values().cloned().collect(),
},
}; };
if history_date >= cutoff.into() { let data_write = serde_json::to_string_pretty(&final_response).unwrap();
let key = ( match fs::write(&writeable_filepath, data_write) {
executions.sha256.clone(), Ok(_) => {}
executions.filename.clone(), Err(e) => {
executions.hostname.clone(), println!("Failed to write to: {:?}: {}", &writeable_filepath, e);
); }
seen.entry(key).or_insert(executions.clone()); }
if let Some(last_item) = &final_response.response.exechistories.last()
&& let Ok(last_date) = NaiveDate::parse_from_str(
&last_item.datetime.replace(" +0000 UTC", ""),
"%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;
progress_bar.set_position(percentage_diff.round() as u64);
progress_bar.set_message("Total Percent Complete");
} }
} }
let final_response = ApiResponse { });
error: "Success".to_string(), progress_bar.finish_with_message("All Checkpoints Complete");
response: ExecHistories { let return_data = match fs::read_to_string(&writeable_filepath) {
exechistories: seen.values().cloned().collect(), Ok(return_data) => return_data,
}, Err(e) => {
}; println!("Failed to read data from: {:?}: {}", &writeable_filepath, e);
let data_write = serde_json::to_string_pretty(&final_response).unwrap(); std::process::abort();
fs::write(&writeable_filepath, data_write).unwrap();
if let Some(last_item) = &final_response.response.exechistories.last()
&& let Ok(last_date) = NaiveDate::parse_from_str(
&last_item.datetime.replace(" +0000 UTC", ""),
"%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;
progress_bar.set_position(percentage_diff.round() as u64);
progress_bar.set_message("Total Percent Complete");
} }
} };
shutdown_tracer_provider();
return_data.to_string()
}); });
progress_bar.finish_with_message("All Checkpoints Complete"); let gil_value = handle.join().unwrap();
let return_data = fs::read_to_string(&writeable_filepath).unwrap(); Python::attach(|py| PyString::new(py, &gil_value).into())
shutdown_tracer_provider();
PyString::new(py, &return_data).into()
} }
fn build_client(py: Python<'_>, py_self: &Py<PyAny>) -> Result<reqwest::Client, reqwest::Error> { fn build_client(headers: HeaderMap) -> Result<reqwest::Client, reqwest::Error> {
let headers = py_self.getattr(py, "headers").unwrap().to_string();
let headers_replace = headers.replace('\'', "\"");
let parsed: Value = serde_json::from_str(headers_replace.as_str()).unwrap();
let mut header_map = HeaderMap::new();
if let Some(obj) = parsed.as_object() {
for (_key, value) in obj {
if let Some(v) = value.as_str() {
let val = HeaderValue::from_str(v).unwrap();
header_map.insert(HeaderName::from_str("X-APIKey").unwrap(), val);
}
}
}
Client::builder() Client::builder()
.danger_accept_invalid_certs(true) .danger_accept_invalid_certs(true)
.default_headers(header_map) .default_headers(headers)
.timeout(std::time::Duration::from_secs(300)) .timeout(std::time::Duration::from_secs(300))
.build() .build()
} }
#[tokio::main] #[tokio::main]
async fn history_logging( async fn history_logging(
py: Python<'_>, base_url: &String,
py_self: &Py<PyAny>,
exec_types: &String, exec_types: &String,
checkpoint_number: &String, checkpoint_number: &String,
policy_names: &String, policy_names: &String,
client: &Client, client: &Client,
) -> ApiResponse { ) -> ApiResponse {
let base_url = py_self.getattr(py, "base_url").unwrap().to_string();
let payload = format!( let payload = format!(
r#"{{ r#"{{
"type": {}, "type": {},
@@ -334,30 +436,8 @@ fn skipback(days: i64) -> ObjectId {
ObjectId::parse_str(&objectid_hex).expect("Invalid ObjectId hex") ObjectId::parse_str(&objectid_hex).expect("Invalid ObjectId hex")
} }
fn load_telemetry_config() -> TelemetryConfig {
let cfg_path = get_base_directory().join("config\\user_config.json");
if !cfg_path.exists() {
return TelemetryConfig {
TELEMETRY: false,
TELEM_URL: None,
};
}
match fs::read_to_string(&cfg_path) {
Ok(contents) => {
serde_json::from_str::<TelemetryConfig>(&contents).unwrap_or(TelemetryConfig {
TELEMETRY: false,
TELEM_URL: None,
})
}
Err(_) => TelemetryConfig {
TELEMETRY: false,
TELEM_URL: None,
},
}
}
fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> { fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> {
let cfg = load_telemetry_config(); let cfg = 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);
+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==3.1.2 airlock_libs==5.0.0