Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 16 additions & 26 deletions crates/rds-cli/src/desktop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,9 @@ pub struct Options {
/// Controlled-marker descriptor for causal input-to-submit diagnostics.
#[arg(long, conflicts_with = "headless")]
pub diagnostic_visual_probe: Option<PathBuf>,
/// Internal, non-owning telemetry observation for the native flight recorder.
#[arg(skip)]
pub diagnostic_health: Option<rds_observe::HealthObserver>,
}

/// Shared bounded credential loading for CLI and native application clients.
Expand Down Expand Up @@ -93,6 +96,9 @@ mod control;
#[cfg(feature = "desktop")]
mod diagnostics;

#[cfg(feature = "desktop")]
mod diagnostic_storage;

#[cfg(feature = "desktop")]
mod native {
use super::*;
Expand Down Expand Up @@ -176,33 +182,31 @@ mod native {
let diagnostic_stop = stop.clone();
let diagnostic_view = handle.clone();
let diagnostic_dir = options.diagnostics_dir.clone();
let diagnostic_health = options.diagnostic_health.clone();
workers.spawn(async move {
let mut tick = tokio::time::interval(Duration::from_secs(2));
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
let path = diagnostic_dir.as_ref().map(|d| d.join(format!("state-{}.json",std::process::id())));
let mut recorder = diagnostics::Recorder::default();
let mut storage = diagnostic_storage::Storage::new(diagnostic_dir);
loop {
tokio::select! {
_ = diagnostic_stop.cancelled() => break,
_ = tick.tick() => {
diagnostic_view.heartbeat_ui();
let snapshot = diagnostic_view.snapshot();
let incident = recorder.observe(&snapshot);
if let Ok(json) = serde_json::to_string(&snapshot) {
let health = diagnostic_health.as_ref().and_then(rds_observe::HealthObserver::snapshot);
let (value, incident) = recorder.observe_with_health(&snapshot, health, &storage.health());
let json = serde_json::to_string(&value).ok();
if let Some(json) = &json {
tracing::info!(snapshot=%json, "viewer health");
if let Some(path) = &path {
let path = path.clone();
let result = tokio::task::spawn_blocking(move || -> std::io::Result<()> {
crate::logging::viewer_snapshot(&path,json.as_bytes())
}).await;
if !matches!(result,Ok(Ok(()))) { tracing::warn!(error=?result,"viewer snapshot write failed"); }
}
}
persist_incident(diagnostic_dir.as_deref(), incident).await;
storage.persist(json.map(String::into_bytes), incident, snapshot.elapsed_ms).await;
},
}
}
persist_incident(diagnostic_dir.as_deref(), recorder.finish(true)).await;
storage.persist(None, recorder.finish(true), diagnostic_view.snapshot().elapsed_ms).await;
let remaining = storage.health().incident_queue_pending;
if remaining > 0 { tracing::warn!(remaining, "viewer shutdown retains unsaved incident windows in memory only"); }
Ok(())
});
workers.spawn(async move {
Expand Down Expand Up @@ -336,20 +340,6 @@ mod native {
}
}

async fn persist_incident(directory: Option<&std::path::Path>, bytes: Option<Vec<u8>>) {
let (Some(directory), Some(bytes)) = (directory, bytes) else {
return;
};
let directory = directory.to_owned();
let result = tokio::task::spawn_blocking(move || {
crate::logging::viewer_incident(&directory, &bytes)
})
.await;
if !matches!(result, Ok(Ok(()))) {
tracing::warn!(error=?result, "viewer incident write failed");
}
}

fn reconnect_target(original: &str, authenticated: String) -> anyhow::Result<String> {
if let Ok(pinned) = rds_net::parse_target(original) {
anyhow::ensure!(
Expand Down
179 changes: 179 additions & 0 deletions crates/rds-cli/src/desktop/diagnostic_storage.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,179 @@
//! Bounded retry ownership for private native incident files. This worker is
//! independent of input/media; failure is observable without dumping contents.
use std::{collections::VecDeque, io, path::PathBuf};

const MAX_PENDING: usize = 4;
const MAX_INCIDENT_BYTES: usize = 256 * 1024;

#[derive(Clone, Default, serde::Serialize)]
pub(super) struct Health {
pub enabled: bool,
pub snapshot_writes_total: u64,
pub snapshot_write_errors_total: u64,
pub incident_writes_total: u64,
pub incident_write_errors_total: u64,
pub incident_queue_evicted_total: u64,
pub incident_oversize_total: u64,
pub incident_queue_pending: usize,
pub last_snapshot_success_elapsed_ms: Option<u64>,
pub last_incident_success_elapsed_ms: Option<u64>,
pub last_error_kind: Option<&'static str>,
}

pub(super) struct Storage {
directory: Option<PathBuf>,
pending: VecDeque<Vec<u8>>,
health: Health,
}

impl Storage {
pub(super) fn new(directory: Option<PathBuf>) -> Self {
Self {
health: Health {
enabled: directory.is_some(),
..Health::default()
},
directory,
pending: VecDeque::new(),
}
}

pub(super) fn health(&self) -> Health {
self.health.clone()
}

fn enqueue(&mut self, incident: Option<Vec<u8>>) {
let Some(bytes) = incident.filter(|_| self.directory.is_some()) else {
return;
};
if bytes.len() > MAX_INCIDENT_BYTES {
self.health.incident_oversize_total += 1;
return;
}
if self.pending.len() == MAX_PENDING {
self.pending.pop_front();
self.health.incident_queue_evicted_total += 1;
}
self.pending.push_back(bytes);
self.health.incident_queue_pending = self.pending.len();
}

fn complete_incident(&mut self, result: io::Result<()>, elapsed_ms: u64) {
match result {
Ok(()) => {
self.pending.pop_front();
self.health.incident_writes_total += 1;
self.health.last_incident_success_elapsed_ms = Some(elapsed_ms);
}
Err(error) => {
self.health.incident_write_errors_total += 1;
self.health.last_error_kind = Some(error_kind(&error));
tracing::warn!(
error_kind = error_kind(&error),
pending = self.pending.len(),
"viewer incident retained for retry after storage failure"
);
}
}
self.health.incident_queue_pending = self.pending.len();
}

pub(super) async fn persist(
&mut self,
snapshot: Option<Vec<u8>>,
incident: Option<Vec<u8>>,
elapsed_ms: u64,
) {
self.enqueue(incident);
let Some(directory) = self.directory.clone() else {
return;
};
if let Some(bytes) = snapshot {
let path = directory.join(format!("state-{}.json", std::process::id()));
let result =
tokio::task::spawn_blocking(move || crate::logging::viewer_snapshot(&path, &bytes))
.await;
match result.unwrap_or_else(|_| Err(io::Error::other("snapshot worker ended"))) {
Ok(()) => {
self.health.snapshot_writes_total += 1;
self.health.last_snapshot_success_elapsed_ms = Some(elapsed_ms);
}
Err(error) => {
self.health.snapshot_write_errors_total += 1;
self.health.last_error_kind = Some(error_kind(&error));
tracing::warn!(
error_kind = error_kind(&error),
"viewer snapshot write failed"
);
}
}
}
// One bounded attempt per diagnostic tick. A failed window remains
// owned; the next tick retries rather than silently discarding it.
if let Some(bytes) = self.pending.front().cloned() {
let result = tokio::task::spawn_blocking(move || {
crate::logging::viewer_incident(&directory, &bytes)
})
.await;
self.complete_incident(
result.unwrap_or_else(|_| Err(io::Error::other("incident worker ended"))),
elapsed_ms,
);
}
}
}

fn error_kind(error: &io::Error) -> &'static str {
match error.kind() {
io::ErrorKind::PermissionDenied => "permission_denied",
io::ErrorKind::StorageFull => "storage_full",
io::ErrorKind::NotFound => "not_found",
io::ErrorKind::AlreadyExists => "already_exists",
_ => "io_failure",
}
}

#[cfg(test)]
mod tests {
use super::*;

#[test]
fn disk_failure_retains_the_exact_incident_until_a_successful_retry() {
let mut storage = Storage::new(Some(PathBuf::from("/unused")));
storage.enqueue(Some(b"bounded synthetic evidence".to_vec()));
storage.complete_incident(Err(io::Error::from_raw_os_error(28)), 2000);
assert_eq!(
storage.pending.front().unwrap(),
b"bounded synthetic evidence"
);
assert_eq!(storage.health().incident_queue_pending, 1);
assert_eq!(storage.health().incident_write_errors_total, 1);
assert_eq!(storage.health().last_error_kind, Some("storage_full"));
assert_eq!(storage.health().incident_writes_total, 0);
storage.complete_incident(Ok(()), 4000);
assert!(storage.pending.is_empty());
assert_eq!(storage.health().incident_writes_total, 1);
assert_eq!(
storage.health().last_incident_success_elapsed_ms,
Some(4000)
);
}

#[test]
fn outage_backlog_is_bounded_and_evictions_are_explicit() {
let mut storage = Storage::new(Some(PathBuf::from("/unused")));
for n in 0..6 {
storage.enqueue(Some(vec![n; MAX_INCIDENT_BYTES]));
}
assert_eq!(storage.pending.len(), 4);
assert_eq!(storage.pending.front().unwrap()[0], 2);
assert_eq!(storage.health().incident_queue_evicted_total, 2);
storage.enqueue(Some(vec![0; MAX_INCIDENT_BYTES + 1]));
assert_eq!(storage.health().incident_oversize_total, 1);
assert_eq!(storage.pending.len(), 4);
let mut disabled = Storage::new(None);
disabled.enqueue(Some(vec![0; 16]));
assert!(disabled.pending.is_empty());
assert!(!disabled.health().enabled);
}
}
Loading
Loading