0ac3b54d89
Build Library / Build Library (push) Successful in 5m54s
Refactored Progress Bar to remove multiprogress bar and only draw one instance. #36 is still open and not fixed with this push, but I believe this is the way to fix the issue. Also implemented an Arc Mutex on the progress bar so it can be controlled via different threads.
339 lines
13 KiB
Rust
339 lines
13 KiB
Rust
use crate::modules::datatypes::*;
|
|
use crate::prelude::*;
|
|
use crossbeam::channel::unbounded;
|
|
use std::sync::{Arc, Mutex};
|
|
use std::thread;
|
|
#[pyfunction]
|
|
pub fn pull_policy_exec_histories(
|
|
py: Python<'_>,
|
|
py_self: Py<PyAny>,
|
|
policy_names: String,
|
|
exec_types: String,
|
|
days: i64,
|
|
) -> Py<PyString> {
|
|
let headers: HeaderMap = match PyData::convert(py, &py_self, true) {
|
|
ExtractedValues::Headers(h) => h,
|
|
ExtractedValues::BaseUrl(_) => std::process::abort(),
|
|
};
|
|
let base_url: String = match PyData::convert(py, &py_self, false) {
|
|
ExtractedValues::Headers(_) => std::process::abort(),
|
|
ExtractedValues::BaseUrl(b) => b,
|
|
};
|
|
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);
|
|
std::process::abort();
|
|
}
|
|
};
|
|
rt.block_on(async {
|
|
let _ = init_tracer();
|
|
});
|
|
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()
|
|
)
|
|
.into();
|
|
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 = ApiResponse {
|
|
error: "Success".to_string(),
|
|
response: ExecHistories {
|
|
exechistories: vec![],
|
|
},
|
|
};
|
|
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) => {
|
|
println!("Failed to write to: {:?}: {}", &writeable_filepath, e);
|
|
std::process::abort();
|
|
}
|
|
}
|
|
let mut checkpoint_number: String = SkipBack::find_checkpoint(days).to_string();
|
|
let progress_bar = Arc::new(Mutex::new(ProgressBar::new(100)));
|
|
progress_bar
|
|
.lock()
|
|
.unwrap()
|
|
.set_draw_target(ProgressDrawTarget::stderr());
|
|
progress_bar.lock().unwrap().set_style(
|
|
ProgressStyle::default_bar()
|
|
.template("Total Completion: {spinner:.green} [{elapsed_precise}] [{bar:40.green/blue}] {pos}/{len} {message}")
|
|
.unwrap(),
|
|
);
|
|
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(
|
|
"info",
|
|
vec![KeyValue::new(
|
|
"Client Built Successfully",
|
|
format!("{:?}", client_result),
|
|
)],
|
|
);
|
|
client_result
|
|
}
|
|
Err(client_result) => {
|
|
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: chrono::NaiveDateTime = Local::now().naive_local() - Duration::days(days);
|
|
let (tx, rx) = unbounded::<Vec<Group>>();
|
|
let pb_clone = progress_bar.clone();
|
|
thread::spawn(move || {
|
|
let mut seen: HashMap<(String, String, String), Group> = if writeable_filepath.exists()
|
|
{
|
|
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(),
|
|
response: ExecHistories {
|
|
exechistories: vec![],
|
|
},
|
|
});
|
|
existing
|
|
.response
|
|
.exechistories
|
|
.into_iter()
|
|
.map(|entry: Group| {
|
|
(
|
|
(
|
|
entry.sha256.clone(),
|
|
entry.filename.clone(),
|
|
entry.hostname.clone(),
|
|
),
|
|
entry,
|
|
)
|
|
})
|
|
.collect()
|
|
} else {
|
|
HashMap::new()
|
|
};
|
|
while let Ok(parsed_responses) = rx.recv() {
|
|
for executions in parsed_responses {
|
|
if executions.checkpoint.is_empty() || executions.datetime.is_empty() {
|
|
continue;
|
|
}
|
|
let history_date: NaiveDate = 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: (String, String, String) = (
|
|
executions.sha256.clone(),
|
|
executions.filename.clone(),
|
|
executions.hostname.clone(),
|
|
);
|
|
seen.entry(key).or_insert(executions.clone());
|
|
}
|
|
}
|
|
let final_response: ApiResponse = ApiResponse {
|
|
error: "Success".to_string(),
|
|
response: ExecHistories {
|
|
exechistories: seen.values().cloned().collect(),
|
|
},
|
|
};
|
|
let data_write: String = serde_json::to_string_pretty(&final_response).unwrap();
|
|
match fs::write(&writeable_filepath, data_write) {
|
|
Ok(_) => {}
|
|
Err(e) => {
|
|
println!("Failed to write to: {:?}: {}", &writeable_filepath, e);
|
|
}
|
|
}
|
|
}
|
|
});
|
|
let mut first_date: Option<NaiveDate> = None;
|
|
tracer.in_span("Airlock Data Retreival", |cx| {
|
|
pb_clone
|
|
.lock()
|
|
.unwrap()
|
|
.enable_steady_tick(std::time::Duration::from_millis(100));
|
|
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 {
|
|
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: Vec<Group> = 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(
|
|
&last_item.datetime.replace(" +0000 UTC", ""),
|
|
"%Y-%m-%dT%H:%M:%SZ",
|
|
)
|
|
{
|
|
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;
|
|
pb_clone.lock().unwrap().set_position(percentage);
|
|
}
|
|
}
|
|
}
|
|
});
|
|
progress_bar
|
|
.lock()
|
|
.unwrap()
|
|
.finish_with_message("All Checkpoints Complete");
|
|
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);
|
|
std::process::abort();
|
|
}
|
|
};
|
|
drop(tx);
|
|
shutdown_tracer_provider();
|
|
return_data.to_string()
|
|
});
|
|
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> {
|
|
Client::builder()
|
|
.danger_accept_invalid_certs(true)
|
|
.default_headers(headers)
|
|
.timeout(std::time::Duration::from_secs(300))
|
|
.build()
|
|
}
|
|
|
|
#[tokio::main]
|
|
async fn history_logging(
|
|
base_url: &String,
|
|
exec_types: &String,
|
|
checkpoint_number: &String,
|
|
policy_names: &String,
|
|
client: &Client,
|
|
) -> ApiResponse {
|
|
let payload = format!(
|
|
r#"{{
|
|
"type": {},
|
|
"checkpoint": "{}",
|
|
"policy": ["{}"]
|
|
}}"#,
|
|
exec_types, checkpoint_number, policy_names
|
|
);
|
|
let res: Result<reqwest::Response, reqwest::Error> = client
|
|
.post(format!("{}/v1/logging/exechistories", base_url))
|
|
.body(payload)
|
|
.send()
|
|
.await;
|
|
match res {
|
|
Ok(res) => {
|
|
let first_response: ApiResponse = serde_json::from_str(&res.text().await.unwrap())
|
|
.expect("Failed to retrieve response from API");
|
|
return first_response;
|
|
}
|
|
Err(_res) => {
|
|
let failed_response: ApiResponse = ApiResponse {
|
|
error: "Failed".to_string(),
|
|
response: ExecHistories {
|
|
exechistories: vec![],
|
|
},
|
|
};
|
|
return failed_response;
|
|
}
|
|
}
|
|
}
|
|
|
|
pub fn get_base_directory() -> PathBuf {
|
|
let home = env::var_os("HOME")
|
|
.map(PathBuf::from)
|
|
.or_else(|| env::var_os("USERPROFILE").map(PathBuf::from))
|
|
.expect("Could not find Home Directory");
|
|
let os = std::env::consts::OS;
|
|
match os {
|
|
"windows" => {
|
|
let appdata = env::var_os("APPDATA")
|
|
.map(PathBuf::from)
|
|
.unwrap_or_else(|| home.join("AppData").join("Roaming"));
|
|
appdata.join("Loxide")
|
|
}
|
|
_ => home.join(".local").join("share").join("Loxide"),
|
|
}
|
|
}
|
|
|
|
fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> {
|
|
let cfg: TelemetryConfig = TelemetryConfig::load();
|
|
if !cfg.TELEMETRY {
|
|
global::set_tracer_provider(NoopTracerProvider::new());
|
|
return Ok(None);
|
|
}
|
|
let endpoint: String = cfg.TELEM_URL.unwrap_or_default();
|
|
let tracer: sdktrace::Tracer =
|
|
opentelemetry_otlp::new_pipeline()
|
|
.tracing()
|
|
.with_exporter(
|
|
opentelemetry_otlp::new_exporter()
|
|
.tonic()
|
|
.with_endpoint(endpoint),
|
|
)
|
|
.with_trace_config(sdktrace::config().with_resource(Resource::new(vec![
|
|
KeyValue::new("service.name", "LoxideLibs"),
|
|
])))
|
|
.install_simple()
|
|
.unwrap();
|
|
Ok(Some(tracer))
|
|
}
|