qualia_client_core/
system_telemetry.rs1use std::sync::atomic::{AtomicBool, Ordering};
4use std::sync::mpsc::{SyncSender, TrySendError};
5use std::sync::{Mutex, OnceLock};
6use std::thread;
7use std::time::Duration;
8
9use serde::{Deserialize, Serialize};
10
11#[derive(Debug, Clone, Serialize, Deserialize)]
12pub struct SystemTelemetryEvent {
13 pub ram_used_mb: u32,
14 pub ram_total_mb: u32,
15 pub vram_used_mb: u32,
16 pub vram_total_mb: u32,
17 pub llm_memory_mb: u32,
18 pub kv_cache_mb: u32,
19 pub lifecycle: String,
20 pub status: String,
21 pub activation_in_progress: bool,
22}
23
24struct TelemetryBus {
25 senders: Mutex<Vec<SyncSender<SystemTelemetryEvent>>>,
26}
27
28impl TelemetryBus {
29 fn new() -> Self {
30 Self {
31 senders: Mutex::new(Vec::new()),
32 }
33 }
34
35 fn subscribe(&self, tx: SyncSender<SystemTelemetryEvent>) {
36 if let Ok(mut senders) = self.senders.lock() {
37 senders.push(tx);
38 }
39 }
40
41 fn publish(&self, event: SystemTelemetryEvent) {
42 if let Ok(mut senders) = self.senders.lock() {
43 senders.retain(|tx| match tx.try_send(event.clone()) {
44 Ok(_) => true,
45 Err(TrySendError::Full(_)) => true,
46 Err(TrySendError::Disconnected(_)) => false,
47 });
48 }
49 }
50}
51
52fn bus() -> &'static TelemetryBus {
53 static BUS: OnceLock<TelemetryBus> = OnceLock::new();
54 BUS.get_or_init(TelemetryBus::new)
55}
56
57static ACTIVATION_TICKER: AtomicBool = AtomicBool::new(false);
58static TICKER_STOP: AtomicBool = AtomicBool::new(false);
59
60pub fn subscribe_system_telemetry(tx: SyncSender<SystemTelemetryEvent>) {
61 bus().subscribe(tx);
62}
63
64pub fn publish_system_telemetry(event: SystemTelemetryEvent) {
65 bus().publish(event);
66}
67
68fn probe_vram_usage_mb() -> (u32, u32) {
69 #[cfg(target_os = "windows")]
70 {
71 if let Ok(memory) = qualia_core_db::directml_bridge::probe_best_adapter_memory() {
72 let used = memory.local_usage_bytes / (1024 * 1024);
73 let total = memory.local_budget_bytes / (1024 * 1024);
74 return (used as u32, total as u32);
75 }
76 }
77 (0, 0)
78}
79
80fn sample_event(status: &str, activation_in_progress: bool) -> SystemTelemetryEvent {
81 use sysinfo::System;
82
83 let mut sys = System::new_all();
84 sys.refresh_memory();
85 let (vram_used_mb, vram_total_mb) = probe_vram_usage_mb();
86 let ram_total_mb = (sys.total_memory() / (1024 * 1024)).min(u32::MAX as u64) as u32;
87 let ram_used_mb = (sys.used_memory() / (1024 * 1024)).min(u32::MAX as u64) as u32;
88 let llm_memory_mb = (crate::model_lifecycle::get_llm_memory_bytes() / (1024 * 1024))
89 .min(u32::MAX as u64) as u32;
90
91 SystemTelemetryEvent {
92 ram_used_mb,
93 ram_total_mb,
94 vram_used_mb,
95 vram_total_mb,
96 llm_memory_mb,
97 kv_cache_mb: crate::model_lifecycle::get_kv_cache_used_mb(),
98 lifecycle: crate::model_lifecycle::lifecycle_label(
99 crate::model_lifecycle::get_model_lifecycle_state(),
100 )
101 .to_string(),
102 status: status.to_string(),
103 activation_in_progress,
104 }
105}
106
107pub fn stop_activation_telemetry() {
108 TICKER_STOP.store(true, Ordering::Release);
109}
110
111pub fn start_activation_telemetry(status: impl Into<String>) {
113 if ACTIVATION_TICKER
114 .compare_exchange(false, true, Ordering::AcqRel, Ordering::Relaxed)
115 .is_err()
116 {
117 return;
118 }
119 TICKER_STOP.store(false, Ordering::Release);
120 let status = status.into();
121 thread::Builder::new()
122 .name("qualia-telemetry-ticker".into())
123 .spawn(move || {
124 bus().publish(sample_event(&status, true));
125 while !TICKER_STOP.load(Ordering::Acquire) {
126 bus().publish(sample_event(&status, true));
127 thread::sleep(Duration::from_millis(100));
128 }
129 bus().publish(sample_event("Model activation complete", false));
130 ACTIVATION_TICKER.store(false, Ordering::Release);
131 })
132 .ok();
133}
134
135pub fn publish_idle_telemetry() {
136 bus().publish(sample_event("Idle", false));
137}
138
139#[cfg(test)]
140mod tests {
141 use super::*;
142
143 #[test]
144 fn sample_event_has_nonzero_ram_total_on_host() {
145 let event = sample_event("test", false);
146 assert!(event.ram_total_mb > 0);
147 }
148}