diff --git a/crates/rds-cli/src/desktop.rs b/crates/rds-cli/src/desktop.rs index 2a5d9a3..9e2f4db 100644 --- a/crates/rds-cli/src/desktop.rs +++ b/crates/rds-cli/src/desktop.rs @@ -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, + /// Internal, non-owning telemetry observation for the native flight recorder. + #[arg(skip)] + pub diagnostic_health: Option, } /// Shared bounded credential loading for CLI and native application clients. @@ -93,6 +96,9 @@ mod control; #[cfg(feature = "desktop")] mod diagnostics; +#[cfg(feature = "desktop")] +mod diagnostic_storage; + #[cfg(feature = "desktop")] mod native { use super::*; @@ -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 { @@ -336,20 +340,6 @@ mod native { } } - async fn persist_incident(directory: Option<&std::path::Path>, bytes: Option>) { - 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 { if let Ok(pinned) = rds_net::parse_target(original) { anyhow::ensure!( diff --git a/crates/rds-cli/src/desktop/diagnostic_storage.rs b/crates/rds-cli/src/desktop/diagnostic_storage.rs new file mode 100644 index 0000000..5f656d7 --- /dev/null +++ b/crates/rds-cli/src/desktop/diagnostic_storage.rs @@ -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, + pub last_incident_success_elapsed_ms: Option, + pub last_error_kind: Option<&'static str>, +} + +pub(super) struct Storage { + directory: Option, + pending: VecDeque>, + health: Health, +} + +impl Storage { + pub(super) fn new(directory: Option) -> 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>) { + 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>, + incident: Option>, + 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); + } +} diff --git a/crates/rds-cli/src/desktop/diagnostics.rs b/crates/rds-cli/src/desktop/diagnostics.rs index 596cea2..03a6ec0 100644 --- a/crates/rds-cli/src/desktop/diagnostics.rs +++ b/crates/rds-cli/src/desktop/diagnostics.rs @@ -19,12 +19,47 @@ pub(super) struct Recorder { previous_reconnects: Option, previous_slow_acks: Option, previous_slow_clipboards: Option, + previous_repairs: Option, + previous_log_errors: u64, + previous_storage_errors: u64, cooldown_until_ms: u64, } impl Recorder { - pub(super) fn observe(&mut self, snapshot: &ViewerSnapshot) -> Option> { - let value = serde_json::to_value(snapshot).ok()?; + #[cfg(test)] + fn observe(&mut self, snapshot: &ViewerSnapshot) -> Option> { + self.observe_with_health( + snapshot, + None, + &super::diagnostic_storage::Health::default(), + ) + .1 + } + + pub(super) fn observe_with_health( + &mut self, + snapshot: &ViewerSnapshot, + telemetry: Option, + storage: &super::diagnostic_storage::Health, + ) -> (serde_json::Value, Option>) { + let mut value = serde_json::to_value(snapshot).unwrap_or(serde_json::Value::Null); + if let Some(object) = value.as_object_mut() { + object.insert( + "diagnostics".into(), + serde_json::json!({"telemetry":telemetry,"storage":storage}), + ); + } + let incident = self.observe_value(snapshot, value.clone(), telemetry, storage); + (value, incident) + } + + fn observe_value( + &mut self, + snapshot: &ViewerSnapshot, + value: serde_json::Value, + telemetry: Option, + storage: &super::diagnostic_storage::Health, + ) -> Option> { self.history.push_back(value.clone()); if self.history.len() > HISTORY { self.history.pop_front(); @@ -41,20 +76,26 @@ impl Recorder { .previous_slow_clipboards .is_some_and(|previous| snapshot.report.slow_clipboard_transfers > previous); self.previous_slow_clipboards = Some(snapshot.report.slow_clipboard_transfers); - if let Some(incident) = &mut self.pending { - // Keep memory bounded even if the diagnostic timer runs rapidly. - if incident.snapshots.len() < HISTORY + 6 { - incident.snapshots.push(value); - } - if snapshot.elapsed_ms.saturating_sub(incident.triggered_ms) >= AFTER_MS { - self.cooldown_until_ms = snapshot.elapsed_ms.saturating_add(COOLDOWN_MS); - return self.finish(false); - } - return None; - } - if snapshot.elapsed_ms < self.cooldown_until_ms { - return None; + let repaired = self + .previous_repairs + .is_some_and(|n| snapshot.report.video_repair_requests > n); + self.previous_repairs = Some(snapshot.report.video_repair_requests); + let log_errors = telemetry.map(|h| { + h.telemetry_dropped_total + .saturating_add(h.telemetry_oversize_total) + .saturating_add(h.telemetry_write_errors_total) + }); + let lost_logs = log_errors.is_some_and(|n| n > self.previous_log_errors); + if let Some(n) = log_errors { + self.previous_log_errors = n; } + let storage_errors = storage + .snapshot_write_errors_total + .saturating_add(storage.incident_write_errors_total) + .saturating_add(storage.incident_queue_evicted_total) + .saturating_add(storage.incident_oversize_total); + let failed_storage = storage_errors > self.previous_storage_errors; + self.previous_storage_errors = storage_errors; let mut reasons = Vec::new(); if slow_ack { // Completed stalls can fall entirely between periodic snapshots. @@ -80,6 +121,11 @@ impl Recorder { if snapshot.decoded_frame_age_ms.is_some_and(|ms| ms >= 3000) { reasons.push("decoded_video_stalled"); } + if snapshot.status == "Connected" + && snapshot.control_echo_age_ms.is_some_and(|ms| ms >= 3000) + { + reasons.push("control_echo_stalled"); + } // Occlusion is a normal OS policy, not a frozen visible screen. if !snapshot.occluded && snapshot.render_stage != "surface occluded" @@ -93,6 +139,39 @@ impl Recorder { if reconnected { reasons.push("desktop_reconnected"); } + if repaired { + reasons.push("video_repair_requested"); + } + if lost_logs { + reasons.push("telemetry_records_lost"); + } + if failed_storage { + reasons.push("diagnostic_storage_failed"); + } + if let Some(incident) = &mut self.pending { + // Later control loss/recovery must not disappear behind the first + // symptom. The reason vocabulary and snapshot count remain bounded. + for reason in reasons { + if !incident.reasons.contains(&reason) { + incident.reasons.push(reason); + } + } + if incident.snapshots.len() < HISTORY + 6 { + incident.snapshots.push(value); + } + if snapshot.elapsed_ms.saturating_sub(incident.triggered_ms) >= AFTER_MS { + self.cooldown_until_ms = snapshot.elapsed_ms.saturating_add(COOLDOWN_MS); + return self.finish(false); + } + return None; + } + // A new recovery or evidence-loss event is distinct from repeated age + // samples and must survive the ordinary symptom cooldown. + if snapshot.elapsed_ms < self.cooldown_until_ms + && !(reconnected || repaired || lost_logs || failed_storage) + { + return None; + } if !reasons.is_empty() { tracing::warn!( ?reasons, @@ -256,4 +335,76 @@ mod tests { serde_json::json!(["visible_submission_stalled"]) ); } + + #[test] + fn later_control_loss_and_recovery_survive_an_open_window_and_cooldown() { + let mut recorder = Recorder::default(); + recorder.observe(&snapshot(0)); + let mut s = snapshot(2000); + s.decoded_frame_age_ms = Some(4000); + recorder.observe(&s); + s.elapsed_ms = 4000; + s.control_echo_age_ms = Some(5000); + s.report.video_repair_requests = 1; + recorder.observe(&s); + s.elapsed_ms = 12_000; + s.report.reconnects = 1; + let data: serde_json::Value = + serde_json::from_slice(&recorder.observe(&s).unwrap()).unwrap(); + assert_eq!( + data["reasons"], + serde_json::json!([ + "decoded_video_stalled", + "control_echo_stalled", + "video_repair_requested", + "desktop_reconnected" + ]) + ); + s.elapsed_ms = 14_000; + s.report.reconnects = 2; + recorder.observe(&s); // Still in ordinary age-symptom cooldown. + let data: serde_json::Value = + serde_json::from_slice(&recorder.finish(true).unwrap()).unwrap(); + assert!( + data["reasons"] + .as_array() + .unwrap() + .contains(&serde_json::json!("desktop_reconnected")) + ); + } + + #[test] + fn evidence_loss_is_recorded_independently_of_the_log_sink() { + let mut recorder = Recorder::default(); + let storage = super::super::diagnostic_storage::Health::default(); + recorder.observe_with_health(&snapshot(0), Some(rds_observe::Health::default()), &storage); + let telemetry = rds_observe::Health { + telemetry_write_errors_total: 1, + telemetry_dropped_total: 2, + ..Default::default() + }; + let failed_storage = super::super::diagnostic_storage::Health { + incident_write_errors_total: 1, + incident_queue_pending: 1, + ..Default::default() + }; + let (value, _) = + recorder.observe_with_health(&snapshot(2000), Some(telemetry), &failed_storage); + assert_eq!( + value["diagnostics"]["telemetry"]["telemetry_dropped_total"], + 2 + ); + assert_eq!(value["diagnostics"]["storage"]["incident_queue_pending"], 1); + let data: serde_json::Value = + serde_json::from_slice(&recorder.finish(true).unwrap()).unwrap(); + assert_eq!( + data["reasons"], + serde_json::json!(["telemetry_records_lost", "diagnostic_storage_failed"]) + ); + let (value, _) = recorder.observe_with_health(&snapshot(4000), None, &storage); + assert!( + value["diagnostics"]["telemetry"].is_null(), + "unavailable counters must not pretend to be zero" + ); + } } diff --git a/crates/rds-cli/src/logging.rs b/crates/rds-cli/src/logging.rs index 6327080..8a00c8b 100644 --- a/crates/rds-cli/src/logging.rs +++ b/crates/rds-cli/src/logging.rs @@ -227,15 +227,51 @@ pub fn viewer_snapshot(path: &Path, bytes: &[u8]) -> io::Result<()> { "unsafe existing viewer snapshot", )); } + publish_private(path, true, |file| file.write_all(bytes)) +} + +/// An incomplete write never occupies the final JSON name. Failed writes clean +/// their owned temporary inode; new incidents use no-replace publication. +fn publish_private( + path: &Path, + replace: bool, + write: impl FnOnce(&mut PrivateLog) -> io::Result<()>, +) -> io::Result<()> { let stamp = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map_err(io::Error::other)? .as_nanos(); let temporary = path.with_extension(format!("{stamp}.pending")); let mut file = PrivateLog::create(&temporary)?; - file.write_all(bytes)?; + let metadata = file.file.metadata()?; + struct Pending { + path: std::path::PathBuf, + dev: u64, + ino: u64, + } + impl Drop for Pending { + fn drop(&mut self) { + if let Ok(metadata) = std::fs::symlink_metadata(&self.path) + && metadata.is_file() + && metadata.dev() == self.dev + && metadata.ino() == self.ino + { + let _ = std::fs::remove_file(&self.path); + } + } + } + let pending = Pending { + path: temporary, + dev: metadata.dev(), + ino: metadata.ino(), + }; + write(&mut file)?; file.flush()?; - std::fs::rename(temporary, path) + if replace { + std::fs::rename(&pending.path, path) + } else { + std::fs::hard_link(&pending.path, path) + } } /// Keep ten completed fault windows independently of ordinary log rotation. @@ -248,9 +284,7 @@ pub fn viewer_incident(directory: &Path, bytes: &[u8]) -> io::Result<()> { .map_err(io::Error::other)? .as_nanos(); let path = directory.join(format!("incident-{stamp}-{}.json", std::process::id())); - let mut file = PrivateLog::create(&path)?; - file.write_all(bytes)?; - file.flush()?; + publish_private(&path, false, |file| file.write_all(bytes))?; let mut owned = Vec::new(); for entry in std::fs::read_dir(directory)? { let entry = entry?; @@ -291,6 +325,31 @@ mod tests { use super::*; use std::os::unix::fs::{DirBuilderExt, symlink}; + #[test] + fn a_failed_partial_publication_leaves_no_json_or_pending_file() { + let root = Path::new("/tmp") + .canonicalize() + .unwrap() + .join(format!("rds-partial-{}", std::process::id())); + std::fs::DirBuilder::new() + .mode(0o700) + .create(&root) + .unwrap(); + let path = root.join("incident.json"); + let failed = publish_private(&path, false, |file| { + file.write_all(b"{\"partial\":")?; + Err(io::Error::from_raw_os_error(28)) + }); + assert_eq!(failed.unwrap_err().kind(), io::ErrorKind::StorageFull); + assert_eq!(std::fs::read_dir(&root).unwrap().count(), 0); + publish_private(&path, false, |file| file.write_all(b"{\"complete\":true}")).unwrap(); + assert!(publish_private(&path, false, |file| file.write_all(b"overwrite")).is_err()); + assert_eq!(std::fs::read(&path).unwrap(), b"{\"complete\":true}"); + assert_eq!(std::fs::metadata(&path).unwrap().nlink(), 1); + assert_eq!(std::fs::read_dir(&root).unwrap().count(), 1); + std::fs::remove_dir_all(root).unwrap(); + } + #[test] fn incident_retention_is_bounded_and_preserves_unowned_entries() { let root = Path::new("/tmp") diff --git a/crates/rds-cli/src/viewer_main.rs b/crates/rds-cli/src/viewer_main.rs index 4f97ea0..23fe1a6 100644 --- a/crates/rds-cli/src/viewer_main.rs +++ b/crates/rds-cli/src/viewer_main.rs @@ -43,6 +43,7 @@ async fn main() -> std::process::ExitCode { })(); match setup { Ok(telemetry) => { + cli.options.diagnostic_health = Some(telemetry.health_observer()); let previous = std::panic::take_hook(); std::panic::set_hook(Box::new(move |info| { tracing::error!(panic=%info,"native viewer panic"); diff --git a/crates/rds-observe/src/lib.rs b/crates/rds-observe/src/lib.rs index a973eb3..bfabf8f 100644 --- a/crates/rds-observe/src/lib.rs +++ b/crates/rds-observe/src/lib.rs @@ -20,7 +20,7 @@ pub mod admin; mod output; use output::{Buffer, Output}; -pub use output::{ConsolePause, Health, Shutdown}; +pub use output::{ConsolePause, Health, HealthObserver, Shutdown}; const TARGET: &str = "rds_telemetry"; /// Each queued record, including its newline, is at most this size. @@ -203,6 +203,11 @@ impl Telemetry { self.output.health() } + /// Observe output loss independently of the potentially failing log sink. + pub fn health_observer(&self) -> HealthObserver { + self.output.health_observer() + } + /// This heartbeat means the main future is being polled, not readiness or /// proof that a remote device is usable. The product owns readiness events. pub async fn run(&self, future: impl Future>) -> Result { diff --git a/crates/rds-observe/src/output.rs b/crates/rds-observe/src/output.rs index 9c7a287..b6985e1 100644 --- a/crates/rds-observe/src/output.rs +++ b/crates/rds-observe/src/output.rs @@ -6,7 +6,7 @@ use std::io::{self, Write}; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; -use std::sync::{Arc, mpsc}; +use std::sync::{Arc, Weak, mpsc}; use std::thread::JoinHandle; use std::time::Duration; @@ -21,6 +21,18 @@ pub struct Health { pub telemetry_write_errors_total: u64, } +/// Read-only counters without retaining the writer, queue or output worker. +#[derive(Clone)] +pub struct HealthObserver(Weak); + +impl HealthObserver { + /// `None` means the counter lifetime has ended, not zero errors. Presence + /// alone does not establish writer readiness or successful delivery. + pub fn snapshot(&self) -> Option { + self.0.upgrade().map(|counters| counters.health()) + } +} + #[derive(Default)] struct Counters { dropped: AtomicU64, @@ -167,6 +179,9 @@ impl Output { pub fn health(&self) -> Health { self.sink.health() } + pub fn health_observer(&self) -> HealthObserver { + HealthObserver(Arc::downgrade(&self.sink.counters)) + } pub async fn pause(&self) -> io::Result { let (ready, acknowledged) = tokio::sync::oneshot::channel(); diff --git a/crates/rds-observe/src/tests.rs b/crates/rds-observe/src/tests.rs index 2f8d9ca..48798b3 100644 --- a/crates/rds-observe/src/tests.rs +++ b/crates/rds-observe/src/tests.rs @@ -28,6 +28,33 @@ impl Capture { } } +#[test] +fn health_observer_is_non_owning_and_reports_sink_failure() { + struct Broken; + impl Write for Broken { + fn write(&mut self, _: &[u8]) -> io::Result { + Err(io::Error::other("synthetic sink failure")) + } + fn flush(&mut self) -> io::Result<()> { + Ok(()) + } + } + let (subscriber, telemetry) = subscriber( + Service::Cli, + Config::new(Format::Json, "off").unwrap(), + Broken, + ) + .unwrap(); + let observer = telemetry.health_observer(); + let dispatch = tracing::Dispatch::new(subscriber); + tracing::dispatcher::with_default(&dispatch, || emit(Event::ProcessStarted)); + let result = telemetry.shutdown(); + assert!(result.drained); + assert_eq!(observer.snapshot().unwrap().telemetry_write_errors_total, 1); + drop(dispatch); + assert!(observer.snapshot().is_none()); +} + #[test] fn json_is_allowlisted_without_invoking_private_formatters() { struct MustNotFormat; diff --git a/docs/native-viewer.md b/docs/native-viewer.md index 33bcb92..f06958b 100644 --- a/docs/native-viewer.md +++ b/docs/native-viewer.md @@ -677,3 +677,30 @@ holds, pulses an intentional client repeat and restores the original setting when its last controller releases. This prevents a late KeyUp from manufacturing letters. Update both native viewer and serving agent together. See the [regression receipt](reports/rds-interactive-repair-20261004.md) for scope limits. + + +### Diagnostic evidence health + +Native live snapshots and incident windows append a `diagnostics` object. +`telemetry` reports dropped, oversized and output-write-error counters through +a weak observer independent of the log sink; unavailable observation is `null`, +not a fabricated zero. `storage` reports snapshot/incident write success/error, +last success time, four-window retry backlog, eviction and oversize counters. +Counters in a snapshot describe completed persistence work before that snapshot. +The native application wires this observation automatically; embedded/ordinary +CLI desktop callers without the internal observer retain `null` telemetry. + +The recorder merges reasons found during its ten-second post-trigger window. +Control-echo silence of at least3seconds, media repair, logging loss and storage +failure supplement existing input/video reasons. New recovery/evidence-loss +transitions bypass ordinary age-symptom cooldown. Four incident payloads of at +most256KiB are retained for retry, with one write attempt per diagnostic tick; +a new window evicts the oldest only at that explicit capacity. An in-flight +attempt may copy one further payload. Final shutdown is best effort: unsaved +memory is not durable, and a permanently full filesystem cannot store evidence. + +Snapshot replacement is atomic; new incident publication is complete and +no-replace. Failed publication removes only its temporary inode. Neither +boundary claims fsync/power-loss durability. Bounded file retention and the +log queue can still lose records; health counts make that uncertainty visible. +Content, key identity, clipboard payload and typed characters are not added. diff --git a/docs/observability.md b/docs/observability.md index d8c93d7..58fe0a4 100644 --- a/docs/observability.md +++ b/docs/observability.md @@ -6,6 +6,10 @@ workers. Operational event filtering and export allowlists remain independent; retaining a span does not enable its filtered diagnostic messages. Scope: W10.1/W10.2, implemented incrementally alongside correctness work. +The native [evidence-preservation increment](reports/rds-incident-evidence-20261005.md) +adds non-owning log-health observation, bounded incident retry and atomic private +publication. Its counters are local metadata; it does not assert collector or +OpenObserve delivery. The native desktop follow-up separates an idle image from actual outstanding presentation work and records the managed input queue/control-write boundaries in private metadata diagnostics. See diff --git a/docs/reports/rds-incident-evidence-20261005.md b/docs/reports/rds-incident-evidence-20261005.md new file mode 100644 index 0000000..3691786 --- /dev/null +++ b/docs/reports/rds-incident-evidence-20261005.md @@ -0,0 +1,30 @@ +# Native diagnostic evidence preservation — 2026-10-05 + +This O4/W10.1 and W6.7 increment repairs observed diagnostic failure boundaries. +It changes no remote wire, input replay, bitrate, path default or server lifecycle. +Private installed transport and host observations remain in the estate. + +## Corrected boundaries + +- A window previously retained only its initial reasons. Later control silence, + repair and reconnection now merge into its bounded reason set; distinct recovery + or evidence-loss transitions survive the ordinary age-symptom cooldown. +- A non-owning output-counter observer makes dropped/oversized records and sink + write failures visible in live snapshots and incident history. Missing counters + remain null. The observer retains no writer, queue or worker. +- Storage health reports completed success/error, backlog/eviction, oversize and + last-success clocks. Four at-most256KiB windows survive failed writes and retry + once per diagnostic tick. Retry capacity and deliberate loss remain explicit. +- A failed partial snapshot/incident write leaves no completed JSON name. Private + temporary publication cleans its own inode; incident publication is no-replace + and preserves existing entries. No power-loss/lossless guarantee is claimed. + +## Regression and validation + +22 CLI library tests and30 observation library tests pass, including synthetic +ENOSPC retention/recovery, retry-capacity overflow, partial JSON cleanup, +existing-file preservation, later control/recovery reason merging and non-owning +sink-failure observation. Strict expanded native workspace clippy passes. +An initial new test borrowed a subscriber where an owned subscriber was required; +it was corrected to use a retained Dispatch. Full workspace/CI and installed +native qualification remain separate evidence to record before completion. diff --git a/docs/research.md b/docs/research.md index bfd373e..fd737c9 100644 --- a/docs/research.md +++ b/docs/research.md @@ -1,5 +1,33 @@ # Deep research: remote + sync, all-Rust, minimum latency +## 2026-10-05 relay loss and evidence preservation + +[Iroh1.2 documentation](https://docs.rs/iroh/1.2.0/iroh/#relay-servers) +confirms its relay carrier uses TCP, while validated direct connectivity uses +QUIC. [Upstream issue4319](https://github.com/n0-computer/iroh/issues/4319) +reports an approximately30-second reachability gap after home-relay loss, +including existing connections despite multiple configured relay hints. +This is a reported analogous failure, not proof of the same defect/version or +a reason to assume a second relay is already a warm failover path. Real route, +network-change and relay actor evidence are required for installed diagnosis. +Allowing direct paths preserves the relay as fallback; it does not authorize +an OS routing or VPN change or guarantee physical-network availability. + +The [Noq stopped contract](https://docs.rs/noq/1.3.0/noq/struct.SendStream.html#method.stopped) +still distinguishes transport receipt from application processing. RDS retains +its existing exact payload receipts and input semantics. [Tokio's bounded-channel +contract](https://docs.rs/tokio/latest/tokio/sync/mpsc/) supports bounded observation +outside the interactive path; retries cannot grow an unbounded diagnostics queue. + +The flight recorder now retains later reasons in the same window, captures +control silence and repair transitions, and exposes log/storage loss counters. +A non-owning health observer avoids waiting on a failed log sink to discover +its failures. Four pending incident windows survive retryable storage failures; +overflow is explicit. Complete private JSON publication prevents a partial +write from masquerading as a completed incident. These boundaries do not provide +lossless audit or power-loss durability. See +[the increment receipt](reports/rds-incident-evidence-20261005.md). + ## 2026-10-05 bounded paste/control observation QUIC [stream flow control](https://www.rfc-editor.org/rfc/rfc9000.html#section-4)