RustImplementation #49

Merged
mysticmomba merged 33 commits from RustImplementation into master 2025-12-17 14:19:20 -05:00
6 changed files with 413 additions and 935 deletions
Showing only changes of commit 630e0a3cdf - Show all commits
+355 -884
View File
File diff suppressed because it is too large Load Diff
+11 -13
View File
@@ -1,33 +1,31 @@
[package] [package]
name = "airlock_libs" name = "signoz_test"
version = "5.2.1" version = "6.0.0"
edition = "2024" edition = "2024"
[lib]
crate-type = ["cdylib"]
[dependencies] [dependencies]
chrono = "0.4.42" chrono = "0.4.42"
indicatif = "0.18.2" indicatif = "0.18.2"
mongodb = "3.3.0" mongodb = "3.3.0"
opentelemetry = { version = "0.18.0", features = ["rt-tokio", "metrics", "trace"] } opentelemetry = { version = "0.27.0", features = ["logs", "metrics", "trace"] }
opentelemetry-otlp = { version = "0.11.0", features = ["trace", "metrics"] } opentelemetry-otlp = { version = "0.27.0", features = ["trace", "metrics", "grpc-tonic", "http-proto", "tls", "reqwest-client", "reqwest-rustls"] }
opentelemetry-semantic-conventions = { version = "0.10.0" } opentelemetry-semantic-conventions = { version = "0.27.0" }
opentelemetry-proto = { version = "0.1.0"} opentelemetry-proto = { version = "0.27.0"}
pyo3 = { version = "0.27.0", features = ["extension-module", "generate-import-lib"] } pyo3 = { version = "0.27.0", features = ["extension-module", "generate-import-lib"] }
reqwest = { version = "0.12.24", features = ["json", "native-tls"] } reqwest = { version = "0.12.24", features = ["json", "native-tls", "rustls-tls"] }
serde = "1.0.228" serde = "1.0.228"
serde-pyobject = "0.8.0" serde-pyobject = "0.8.0"
serde_json = "1.0.145" serde_json = "1.0.145"
tokio = { version = "1.48.0", features = ["full"] } tokio = { version = "1.48.0", features = ["full"] }
tonic = { version = "0.8.2", features = ["tls-roots"] } tonic = { version = "0.12.3", 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"] }
crossbeam = "0.8.4" crossbeam = "0.8.4"
log = "0.4.29" log = "0.4.29"
flexi_logger = "0.31.7" flexi_logger = "0.31.7"
opentelemetry-appender-log = "0.27.0"
opentelemetry_sdk = { version = "0.27.0", features = ["rt-tokio", "trace"] }
[package.metadata.maturin] [package.metadata.maturin]
generate-abi-stubs = true generate-abi-stubs = true
@@ -40,4 +38,4 @@ codegen-units = 1
panic = 'abort' panic = 'abort'
strip = true strip = true
debug-assertions = false debug-assertions = false
overflow-checks = false overflow-checks = true
+1 -1
View File
@@ -4,7 +4,7 @@ build-backend = "maturin"
[project] [project]
name = "airlock_libs" name = "airlock_libs"
version = "5.2.1" version = "6.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" }
+6 -6
View File
@@ -1,15 +1,15 @@
pub use chrono::{Duration, Local, NaiveDate}; pub use chrono::{Duration, Local, NaiveDate};
pub use indicatif::{MultiProgress, ProgressBar, ProgressDrawTarget, ProgressStyle}; pub use indicatif::{MultiProgress, ProgressBar, ProgressDrawTarget, ProgressStyle};
pub use mongodb::bson::oid::ObjectId; pub use mongodb::bson::oid::ObjectId;
pub use opentelemetry::global::shutdown_tracer_provider; pub use opentelemetry::global::GlobalTracerProvider;
pub use opentelemetry::sdk::Resource;
pub use opentelemetry::trace::noop::NoopTracerProvider; pub use opentelemetry::trace::noop::NoopTracerProvider;
pub use opentelemetry::trace::{Status, TraceContextExt, TraceError}; pub use opentelemetry::trace::{Status, TraceContextExt, Tracer};
pub use opentelemetry::{Context, KeyValue, sdk::trace as sdktrace, trace::Tracer}; pub use opentelemetry::*;
pub use opentelemetry::{Key, global}; pub use opentelemetry_otlp::ExportConfig;
pub use opentelemetry_otlp::WithExportConfig; pub use opentelemetry_otlp::WithExportConfig;
pub use opentelemetry_sdk::Resource;
pub use opentelemetry_sdk::trace::{Config, TracerProvider};
pub use pyo3::{prelude::*, types::PyString}; pub use pyo3::{prelude::*, types::PyString};
pub use pyo3_async_runtimes::async_std;
pub use reqwest::{ pub use reqwest::{
Client, Client,
header::{HeaderMap, HeaderName, HeaderValue}, header::{HeaderMap, HeaderName, HeaderValue},
+39 -30
View File
@@ -1,8 +1,10 @@
use crate::modules::datatypes::*; use crate::modules::datatypes::*;
use crate::prelude::*; use crate::prelude::*;
use crossbeam::channel::unbounded; use crossbeam::channel::unbounded;
use opentelemetry_otlp::WithTonicConfig;
use std::sync::{Arc, Mutex}; use std::sync::{Arc, Mutex};
use std::thread; use std::thread;
use tonic::transport::{Channel, ClientTlsConfig};
#[pyfunction] #[pyfunction]
pub fn pull_policy_exec_histories( pub fn pull_policy_exec_histories(
py: Python<'_>, py: Python<'_>,
@@ -11,9 +13,10 @@ pub fn pull_policy_exec_histories(
exec_types: String, exec_types: String,
days: i64, days: i64,
) -> Py<PyString> { ) -> Py<PyString> {
let data = PyData::extract_data(py, &py_self); println!();
let headers = data.headers; let data: PyData = PyData::extract_data(py, &py_self);
let base_url = data.base_url; let headers: HeaderMap = data.headers;
let base_url: String = data.base_url;
let handle: thread::JoinHandle<String> = std::thread::spawn(move || { let handle: thread::JoinHandle<String> = std::thread::spawn(move || {
let rt: tokio::runtime::Runtime = match tokio::runtime::Runtime::new() { let rt: tokio::runtime::Runtime = match tokio::runtime::Runtime::new() {
Ok(rt) => rt, Ok(rt) => rt,
@@ -22,10 +25,9 @@ pub fn pull_policy_exec_histories(
std::process::abort(); std::process::abort();
} }
}; };
rt.block_on(async { let tracer_provider = rt.block_on(async { init_tracer() });
let _ = init_tracer(); global::set_tracer_provider(tracer_provider.clone());
}); let tracer: global::BoxedTracer = global::tracer("tracer");
let tracer: global::BoxedTracer = global::tracer("global_tracer");
let _cx: Context = Context::new(); let _cx: Context = Context::new();
let file_path: PathBuf = format!( let file_path: PathBuf = format!(
"{}\\cache\\chunkinator.json", "{}\\cache\\chunkinator.json",
@@ -106,7 +108,8 @@ pub fn pull_policy_exec_histories(
} }
} }
}); });
let cutoff: chrono::NaiveDateTime = Local::now().naive_local() - Duration::days(days); let cutoff: chrono::NaiveDateTime =
Local::now().naive_local() - chrono::Duration::days(days);
let (tx, rx) = unbounded::<Vec<Group>>(); let (tx, rx) = unbounded::<Vec<Group>>();
let pb_clone = progress_bar.clone(); let pb_clone = progress_bar.clone();
thread::spawn(move || { thread::spawn(move || {
@@ -181,7 +184,10 @@ pub fn pull_policy_exec_histories(
.unwrap() .unwrap()
.enable_steady_tick(std::time::Duration::from_millis(100)); .enable_steady_tick(std::time::Duration::from_millis(100));
let span: opentelemetry::trace::SpanRef<'_> = 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(Key::new("Days"));
//span.set_attribute(KeyValue::new("Policy Name", policy_names.clone()));
span.set_attribute(KeyValue::new("Days", days));
span.set_attribute(KeyValue::new("Policy Name", policy_names.clone())); span.set_attribute(KeyValue::new("Policy Name", policy_names.clone()));
loop { loop {
let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| { let execution_histories = tracer.in_span(checkpoint_number.to_string(), |cx| {
@@ -237,8 +243,10 @@ pub fn pull_policy_exec_histories(
std::process::abort(); std::process::abort();
} }
}; };
tracer_provider
.shutdown()
.expect("Failed to Shutdown Tracer Provdier");
drop(tx); drop(tx);
shutdown_tracer_provider();
return_data.to_string() return_data.to_string()
}); });
let gil_value: String = handle.join().unwrap(); let gil_value: String = handle.join().unwrap();
@@ -313,25 +321,26 @@ pub fn get_base_directory() -> PathBuf {
} }
} }
fn init_tracer() -> Result<Option<sdktrace::Tracer>, TraceError> { fn init_tracer() -> opentelemetry_sdk::trace::TracerProvider {
let cfg: TelemetryConfig = TelemetryConfig::load(); let cfg: TelemetryConfig = TelemetryConfig::load();
if !cfg.TELEMETRY { let endpoint = cfg.TELEM_URL.unwrap_or_default().clone();
global::set_tracer_provider(NoopTracerProvider::new()); let channel_endpoint = endpoint.clone();
return Ok(None); let channel = Channel::from_shared(channel_endpoint.clone())
} .unwrap()
let endpoint: String = cfg.TELEM_URL.unwrap_or_default(); .tls_config(ClientTlsConfig::new().with_native_roots())
let tracer: sdktrace::Tracer = .unwrap()
opentelemetry_otlp::new_pipeline() .connect_lazy();
.tracing() let exporter = opentelemetry_otlp::SpanExporter::builder()
.with_exporter( .with_tonic()
opentelemetry_otlp::new_exporter() .with_endpoint(endpoint.clone())
.tonic() .with_channel(channel)
.with_endpoint(endpoint), .build()
) .expect("Failed to build exporter");
.with_trace_config(sdktrace::config().with_resource(Resource::new(vec![ opentelemetry_sdk::trace::TracerProvider::builder()
KeyValue::new("service.name", "LoxideLibs"), .with_simple_exporter(exporter)
]))) .with_resource(Resource::new(vec![KeyValue::new(
.install_simple() "service.name",
.unwrap(); "LoxideLibs",
Ok(Some(tracer)) )]))
.build()
} }
+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.2.1 airlock_libs==6.0.0