use smbcloud_gresiq_sdk::{Environment, GresiqClient, GresiqCredentials};
use super::events::{InferenceEvent, ModelLoadedEvent};
const EMBEDDED_API_KEY_DEV: Option<&str> = option_env!("GRESIQ_API_KEY_DEV");
const EMBEDDED_API_SECRET_DEV: Option<&str> = option_env!("GRESIQ_API_SECRET_DEV");
const EMBEDDED_API_KEY_PRODUCTION: Option<&str> = option_env!("GRESIQ_API_KEY_PRODUCTION");
const EMBEDDED_API_SECRET_SECRET_PRODUCTION: Option<&str> =
option_env!("GRESIQ_API_SECRET_PRODUCTION");
#[derive(Debug, Clone)]
pub struct PulseClient {
inner: GresiqClient,
edge_id: String,
onde_app_id: Option<String>,
}
impl PulseClient {
pub fn disabled_by_env() -> bool {
matches!(
std::env::var("ONDE_DISABLE_PULSE")
.ok()
.as_deref()
.map(str::trim)
.map(str::to_ascii_lowercase)
.as_deref(),
Some("1") | Some("true") | Some("yes") | Some("on")
)
}
fn dual_write_enabled() -> bool {
!matches!(
std::env::var("ONDE_PULSE_DUAL_WRITE")
.ok()
.as_deref()
.map(str::trim)
.map(str::to_ascii_lowercase)
.as_deref(),
Some("0") | Some("false") | Some("no") | Some("off")
)
}
pub fn new(
environment: Environment,
edge_id: String,
onde_app_id: Option<String>,
) -> Option<Self> {
if Self::disabled_by_env() {
return None;
}
let (api_key, api_secret) = match environment {
Environment::Dev => (EMBEDDED_API_KEY_DEV?, EMBEDDED_API_SECRET_DEV?),
Environment::Production => (
EMBEDDED_API_KEY_PRODUCTION?,
EMBEDDED_API_SECRET_SECRET_PRODUCTION?,
),
};
let edge_id = if edge_id.is_empty() {
"onde-unknown".to_string()
} else {
edge_id
};
if tokio::runtime::Handle::try_current().is_err() {
log::warn!(
"pulse: no Tokio runtime available — \
deferring PulseClient creation"
);
return None;
}
let credentials = GresiqCredentials {
api_key,
api_secret,
};
let inner = match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
GresiqClient::from_credentials(environment, credentials)
})) {
Ok(client) => client,
Err(_) => {
log::warn!(
"pulse: GresiqClient::from_credentials panicked \
(likely missing Tokio reactor) — telemetry disabled"
);
return None;
}
};
Some(PulseClient {
inner,
edge_id,
onde_app_id,
})
}
pub async fn record_model_loaded(
&self,
model_id: String,
model_name: String,
load_duration_ms: u64,
) {
let client = self.clone();
let send_event = async move {
let event = ModelLoadedEvent {
edge_id: client.edge_id.clone(),
model_id,
model_name,
load_duration_ms,
onde_app_id: client.onde_app_id.clone(),
};
if let Err(error) = client.inner.insert("pulse/model_loaded", &event).await {
log::warn!("pulse: model_loaded failed: {}", error);
}
client.dual_write_model_loaded(&event).await;
};
if tokio::runtime::Handle::try_current().is_ok() {
send_event.await;
return;
}
let join = std::thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build();
match runtime {
Ok(runtime) => {
runtime.block_on(send_event);
}
Err(error) => {
log::warn!(
"pulse: could not create fallback runtime for model_loaded: {}",
error
);
}
}
});
if join.join().is_err() {
log::warn!("pulse: fallback model_loaded thread panicked");
}
}
pub fn record_inference(
&self,
model_id: String,
request_id: String,
duration_ms: u64,
status: String,
) {
let client = self.clone();
let send_event = async move {
let event = InferenceEvent {
edge_id: client.edge_id.clone(),
model_id,
request_id,
duration_ms,
status,
onde_app_id: client.onde_app_id.clone(),
};
if let Err(error) = client.inner.insert("pulse/inference_event", &event).await {
log::warn!("pulse: inference_event failed: {}", error);
}
client.dual_write_inference(&event).await;
};
if let Ok(handle) = tokio::runtime::Handle::try_current() {
handle.spawn(send_event);
return;
}
std::thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build();
match runtime {
Ok(runtime) => {
runtime.block_on(send_event);
}
Err(error) => {
log::warn!(
"pulse: could not create fallback runtime for inference_event: {}",
error
);
}
}
});
}
async fn dual_write_model_loaded(&self, event: &ModelLoadedEvent) {
if !Self::dual_write_enabled() {
return;
}
self.write_doc(
"pulse_edges",
Some(&event.edge_id),
&serde_json::json!({ "edge_id": event.edge_id }),
)
.await;
let model_key = event
.onde_app_id
.clone()
.unwrap_or_else(|| event.model_id.clone());
self.write_doc(
"pulse_models",
Some(&model_key),
&serde_json::json!({
"slug": event.model_id,
"model_id": event.model_id,
"model_name": event.model_name,
"onde_app_id": event.onde_app_id,
}),
)
.await;
let deployment_key = format!("{}:{}", event.edge_id, event.model_id);
self.write_doc(
"pulse_deployments",
Some(&deployment_key),
&serde_json::json!({
"edge_id": event.edge_id,
"model_id": event.model_id,
"load_duration_ms": event.load_duration_ms,
"onde_app_id": event.onde_app_id,
}),
)
.await;
}
async fn dual_write_inference(&self, event: &InferenceEvent) {
if !Self::dual_write_enabled() {
return;
}
self.write_doc(
"pulse_inference_events",
None,
&serde_json::json!({
"edge_id": event.edge_id,
"model_id": event.model_id,
"request_id": event.request_id,
"duration_ms": event.duration_ms,
"status": event.status,
"onde_app_id": event.onde_app_id,
}),
)
.await;
}
async fn write_doc(&self, collection: &str, key: Option<&str>, doc: &serde_json::Value) {
if let Err(error) = self.inner.upsert_document(collection, key, doc).await {
log::warn!("pulse: dual-write to {} failed: {}", collection, error);
}
}
}