From 3d3f50f693400f7712794fff8b94a03052220d40 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Tue, 18 Aug 2026 13:24:15 -0600 Subject: [PATCH] feat(observer): record durable progress and signal empty windows MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The per-segment .metadata sidecar now carries has_durable_media, durable_byte_count, and last_durable_write_at alongside start_timestamp, refreshed by the production tick from a scan of the open segment directory. start_timestamp is write-once: the refresh sources it from ObserverState::segment_start_wall, never a read-modify-write of the sidecar. last_durable_write_at is omitted, never null, until media is observed. A boundary that never started video and produced no media emits one observe.stream_silent with reason: empty_window, latched per observer instance and cleared by a rotation that keeps media. The latch is guarded by a per-window video_started flag so the unhealthy-stream watchdog path — which already emitted via stop_video after the backend unlinked a header-only webm — is not double-counted. Legacy 39-byte {"start_timestamp":…} sidecars still recover unchanged. --- crates/solstone-linux/src/observer.rs | 281 +++++++++++++++++++++++++- crates/solstone-linux/src/recovery.rs | 108 ++++++++-- crates/solstone-linux/src/run.rs | 69 ++++++- crates/solstone-linux/src/segment.rs | 4 +- 4 files changed, 427 insertions(+), 35 deletions(-) diff --git a/crates/solstone-linux/src/observer.rs b/crates/solstone-linux/src/observer.rs index ec72f07..0499257 100644 --- a/crates/solstone-linux/src/observer.rs +++ b/crates/solstone-linux/src/observer.rs @@ -8,7 +8,7 @@ use crate::{ chunking::{DrainedChunk, HitGate}, config::Config, encoding::{AudioOutputPlan, audio_output_plan}, - recovery::write_segment_metadata, + recovery::{SegmentProgress, scan_segment_progress, write_segment_metadata}, segment::{clamp_duration, finalize_segment_dir, segment_key, timestamp_parts}, streams::is_healthy_file_size, }; @@ -144,6 +144,8 @@ pub struct Registration { pub health: HealthBeacon, } +pub const EMPTY_WINDOW_REASON: &str = "empty_window"; + #[derive(Clone, Debug, PartialEq)] pub struct StreamSilentEvent { pub connector: String, @@ -154,6 +156,7 @@ pub struct StreamSilentEvent { pub duration_seconds: i64, pub host: String, pub platform: String, + pub reason: Option<&'static str>, } #[derive(Clone, Debug, PartialEq, Eq)] pub struct SegmentCompletedEvent { @@ -235,6 +238,9 @@ pub struct ObserverState { pub cached_is_active: bool, pub cached_activity: ActivityState, pub current_streams: Vec, + // Survives stop_video: the watchdog must not look like "video never started". + pub video_started: bool, + pub empty_window_signalled: bool, pub hit_gate: HitGate, pub frames: Vec, pub capture_stats: CaptureStats, @@ -290,6 +296,8 @@ where cached_is_active: false, cached_activity: ActivityState::default(), current_streams: vec![], + video_started: false, + empty_window_signalled: false, hit_gate: HitGate::default(), frames: vec![], capture_stats: CaptureStats { @@ -364,6 +372,20 @@ where self.publish(); return Ok(()); } + if let Some(dir) = self.state.segment_dir.as_ref() { + let (has_durable_media, durable_byte_count) = scan_segment_progress(dir); + let last_durable_write_at = + has_durable_media.then(|| self.backends.clock.wall_seconds()); + write_segment_metadata( + dir, + self.state.segment_start_wall, + SegmentProgress { + has_durable_media, + durable_byte_count, + last_durable_write_at, + }, + ); + } let target = self.probe_status().unwrap_or(self.state.mode); if self.state.mode == Mode::Screencast && !self.backends.video.is_healthy() { self.stop_video()?; @@ -392,6 +414,27 @@ where self.write_gated_audio()?; self.state.frames.clear(); self.state.hit_gate = HitGate::default(); + // shutdown and finish_paused_segment deliberately do not signal empty_window. + if let Some(dir) = self.state.segment_dir.as_ref() { + let (has_media, _) = scan_segment_progress(dir); + if has_media { + self.state.empty_window_signalled = false; + } else if !self.state.video_started && !self.state.empty_window_signalled { + // !video_started keeps the watchdog path from double-counting: x11/portal + // unlink header-only webms inside stop() before emit_silent, so emptiness + // alone is not sufficient. + self.emit_silent( + StoppedStream { + node_id: 0, + connector: String::new(), + position: String::new(), + file_bytes: 0, + }, + Some(EMPTY_WINDOW_REASON), + ); + self.state.empty_window_signalled = true; + } + } self.finalize_segment()?; self.state.segment_is_muted = self.state.cached_is_muted; self.state.mode = target; @@ -451,6 +494,7 @@ where return Err(ObserverError::VideoStart("no streams available".into())); } self.state.current_streams = streams; + self.state.video_started = true; } Ok(()) } @@ -464,10 +508,11 @@ where .join(&self.config.stream) .join(format!("{time}.incomplete")); fs::create_dir_all(&dir)?; - write_segment_metadata(&dir, wall); + write_segment_metadata(&dir, wall, SegmentProgress::default()); self.state.segment_start_wall = wall; self.state.segment_start_mono = self.backends.clock.monotonic_seconds(); self.state.segment_dir = Some(dir.clone()); + self.state.video_started = false; Ok(dir) } fn finalize_segment(&mut self) -> Result<(), ObserverError> { @@ -503,11 +548,11 @@ where .into_iter() .filter(|s| !is_healthy_file_size(Some(s.file_bytes))) { - self.emit_silent(s) + self.emit_silent(s, None) } Ok(()) } - fn emit_silent(&mut self, s: StoppedStream) { + fn emit_silent(&mut self, s: StoppedStream, reason: Option<&'static str>) { let segment_dir = self .state .segment_dir @@ -527,6 +572,7 @@ where duration_seconds, host: self.host.clone(), platform: self.platform.clone(), + reason, }); } fn refresh_stats(&mut self) { @@ -1491,12 +1537,15 @@ pub(crate) mod tests { let mut f = fixture(true); initialize(&mut f); f.wall.set(f.wall.get() - 10.0); - f.observer.emit_silent(StoppedStream { - node_id: 1, - connector: "x".into(), - position: "p".into(), - file_bytes: 1, - }); + f.observer.emit_silent( + StoppedStream { + node_id: 1, + connector: "x".into(), + position: "p".into(), + file_bytes: 1, + }, + None, + ); let s = f.events.silent.borrow(); assert_eq!(s[0].segment_dir, ""); assert_eq!(s[0].duration_seconds, -10) @@ -1696,4 +1745,216 @@ pub(crate) mod tests { ); } } + + fn sidecar(dir: &Path) -> Value { + serde_json::from_str(&fs::read_to_string(dir.join(".metadata")).unwrap()).unwrap() + } + + // AC1: open sidecar is default progress with last_durable_write_at omitted. + #[test] + fn open_metadata_is_default_progress() { + let mut f = fixture(false); + f.observer.backends.activity.0.push_back(Ok(idle())); + initialize(&mut f); + let meta = sidecar(f.observer.state.segment_dir.as_ref().unwrap()); + assert_eq!(meta["start_timestamp"], 1_700_000_000.0); + assert_eq!(meta["has_durable_media"], false); + assert_eq!(meta["durable_byte_count"], 0); + assert!(meta.get("last_durable_write_at").is_none()); + } + + // AC2: production tick refresh sees a test-planted file. Do not call the writer. + #[test] + fn tick_refresh_stamps_observer_wall() { + let mut f = fixture(false); + f.observer.backends.activity.0.push_back(Ok(idle())); + initialize(&mut f); + let dir = f.observer.state.segment_dir.clone().unwrap(); + fs::write(dir.join("planted.bin"), b"abcd").unwrap(); + f.observer.backends.activity.0.push_back(Ok(idle())); + f.observer.tick().unwrap(); + let meta = sidecar(&dir); + assert_eq!(meta["has_durable_media"], true); + assert_eq!(meta["durable_byte_count"], 4); + assert_eq!(meta["last_durable_write_at"], 1_700_000_000.0); + assert_eq!(meta["start_timestamp"], 1_700_000_000.0); + } + + // AC3: a tick with no media omits last_durable_write_at. + #[test] + fn tick_refresh_without_media_omits_write_at() { + let mut f = fixture(false); + f.observer.backends.activity.0.push_back(Ok(idle())); + initialize(&mut f); + let dir = f.observer.state.segment_dir.clone().unwrap(); + f.observer.backends.activity.0.push_back(Ok(idle())); + f.observer.tick().unwrap(); + let meta = sidecar(&dir); + assert_eq!(meta["has_durable_media"], false); + assert_eq!(meta["durable_byte_count"], 0); + assert!(meta.get("last_durable_write_at").is_none()); + } + + // AC4: start_timestamp is write-once from state; last_durable_write_at is the refresh wall. + #[test] + fn tick_refresh_keeps_open_start_timestamp() { + let mut f = fixture(false); + f.observer.backends.activity.0.push_back(Ok(idle())); + initialize(&mut f); + let dir = f.observer.state.segment_dir.clone().unwrap(); + fs::write(dir.join("planted.bin"), b"x").unwrap(); + f.wall.set(1_700_000_050.0); + f.observer.backends.activity.0.push_back(Ok(idle())); + f.observer.tick().unwrap(); + let meta = sidecar(&dir); + assert_eq!(meta["start_timestamp"], 1_700_000_000.0); + assert_eq!(meta["last_durable_write_at"], 1_700_000_050.0); + assert_eq!( + crate::recovery::read_segment_start(&dir), + Some(1_700_000_000.0) + ); + } + + // AC6: empty-window latch across consecutive empty rotations and a media reset. + #[test] + fn empty_window_latch_across_rotations() { + let mut f = fixture(false); + f.observer.backends.activity.0.push_back(Ok(idle())); + initialize(&mut f); + let first = f.observer.state.segment_dir.clone().unwrap(); + f.observer.backends.activity.0.push_back(Ok(idle())); + f.wall.set(1_700_000_001.0); + f.mono.set(300.0); + f.observer.tick().unwrap(); + assert!(!first.exists()); + { + let silent = f.events.silent.borrow(); + assert_eq!(silent.len(), 1); + assert_eq!(silent[0].reason, Some(EMPTY_WINDOW_REASON)); + assert_eq!(silent[0].connector, ""); + assert_eq!(silent[0].position, ""); + assert_eq!(silent[0].node_id, 0); + assert_eq!(silent[0].file_bytes, 0); + assert!(silent[0].segment_dir.ends_with(".incomplete")); + assert_eq!(silent[0].duration_seconds, 1); + assert_eq!(silent[0].host, "host"); + assert_eq!(silent[0].platform, "linux"); + } + assert!(f.events.completed.borrow().is_empty()); + + let second = f.observer.state.segment_dir.clone().unwrap(); + f.observer.backends.activity.0.push_back(Ok(idle())); + f.wall.set(1_700_000_002.0); + f.mono.set(600.0); + f.observer.tick().unwrap(); + assert!(!second.exists()); + assert_eq!(f.events.silent.borrow().len(), 1); + assert!(f.events.completed.borrow().is_empty()); + + let media = f.observer.state.segment_dir.clone().unwrap(); + fs::write(media.join("kept.bin"), b"keep").unwrap(); + f.observer.backends.activity.0.push_back(Ok(idle())); + f.wall.set(1_700_000_003.0); + f.mono.set(900.0); + f.observer.tick().unwrap(); + assert_eq!(f.events.silent.borrow().len(), 1); + assert_eq!(f.events.completed.borrow().len(), 1); + assert!(!f.observer.state.empty_window_signalled); + + let fourth = f.observer.state.segment_dir.clone().unwrap(); + f.observer.backends.activity.0.push_back(Ok(idle())); + f.wall.set(1_700_000_004.0); + f.mono.set(1_200.0); + f.observer.tick().unwrap(); + assert!(!fourth.exists()); + assert_eq!(f.events.silent.borrow().len(), 2); + assert_eq!( + f.events.silent.borrow()[1].reason, + Some(EMPTY_WINDOW_REASON) + ); + assert_eq!(f.events.completed.borrow().len(), 1); + } + + // AC7: unhealthy silent is unchanged and is not joined by empty_window. + #[test] + fn unhealthy_silent_is_not_empty_window() { + let mut f = fixture(false); + initialize(&mut f); + f.observer.backends.video.stopped = vec![ + StoppedStream { + node_id: 1, + connector: "x".into(), + position: "left".into(), + file_bytes: 2047, + }, + StoppedStream { + node_id: 2, + connector: "y".into(), + position: "right".into(), + file_bytes: 2048, + }, + ]; + f.wall.set(f.wall.get() + 400.0); + f.observer.handle_boundary(Mode::Idle).unwrap(); + let silent = f.events.silent.borrow(); + assert_eq!(silent.len(), 1); + assert_eq!(silent[0].connector, "x"); + assert_eq!(silent[0].position, "left"); + assert_eq!(silent[0].node_id, 1); + assert_eq!(silent[0].file_bytes, 2047); + assert!(silent[0].reason.is_none()); + assert!( + silent + .iter() + .all(|event| event.reason != Some(EMPTY_WINDOW_REASON)) + ); + } + + // AC8: Idle→Idle boundary that keeps planted media completes and does not signal empty_window. + #[test] + fn idle_boundary_with_media_completes_without_empty_window() { + let mut f = fixture(false); + f.observer.backends.activity.0.push_back(Ok(idle())); + initialize(&mut f); + fs::write( + f.observer + .state + .segment_dir + .as_ref() + .unwrap() + .join("kept.bin"), + b"keep", + ) + .unwrap(); + f.observer.backends.activity.0.push_back(Ok(idle())); + f.mono.set(300.0); + f.observer.tick().unwrap(); + assert!(f.events.silent.borrow().is_empty()); + assert_eq!(f.events.completed.borrow().len(), 1); + } + + // Watchdog already emitted (or would have) via stop_video; video_started blocks empty_window. + #[test] + fn watchdog_stop_does_not_also_emit_empty_window() { + let mut f = fixture(false); + initialize(&mut f); + fs::remove_file( + f.observer + .state + .segment_dir + .as_ref() + .unwrap() + .join("screen.webm"), + ) + .unwrap(); + f.observer.backends.video.healthy = false; + f.observer.tick().unwrap(); + assert!( + f.events + .silent + .borrow() + .iter() + .all(|event| event.reason != Some(EMPTY_WINDOW_REASON)) + ); + } } diff --git a/crates/solstone-linux/src/recovery.rs b/crates/solstone-linux/src/recovery.rs index 0e93bb8..655e430 100644 --- a/crates/solstone-linux/src/recovery.rs +++ b/crates/solstone-linux/src/recovery.rs @@ -3,7 +3,8 @@ use crate::segment::clamp_duration; use claxon::{FlacReader, FlacReaderOptions}; -use serde_json::{Value, json}; +use serde::Serialize; +use serde_json::Value; use std::{ fs::{self, File, FileTimes}, path::{Path, PathBuf}, @@ -49,8 +50,27 @@ fn stream_duration(samples: Option, sample_rate: u32) -> Option { Some(samples? as f64 / f64::from(sample_rate)) } -pub fn write_segment_metadata(segment_dir: &Path, start_timestamp: f64) { - let Ok(mut text) = serde_json::to_string(&json!({"start_timestamp": start_timestamp})) else { +#[derive(Clone, Copy, Debug, Default, PartialEq, Serialize)] +pub struct SegmentProgress { + pub has_durable_media: bool, + pub durable_byte_count: u64, + #[serde(skip_serializing_if = "Option::is_none")] + pub last_durable_write_at: Option, +} + +#[derive(Serialize)] +struct SegmentMetadataFile { + start_timestamp: f64, + #[serde(flatten)] + progress: SegmentProgress, +} + +pub fn write_segment_metadata(segment_dir: &Path, start_timestamp: f64, progress: SegmentProgress) { + let document = SegmentMetadataFile { + start_timestamp, + progress, + }; + let Ok(mut text) = serde_json::to_string(&document) else { tracing::warn!("Failed to write segment metadata"); return; }; @@ -60,6 +80,34 @@ pub fn write_segment_metadata(segment_dir: &Path, start_timestamp: f64) { } } +// Existence of a non-`.metadata` regular file, not byte_count > 0 — a 0-byte leftover +// is media, matching finalize_segment. Not shared with finalize_segment (boolean any-file +// after deleting `.metadata`) or recover_segment (every non-`.metadata` entry, including dirs). +pub fn scan_segment_progress(segment_dir: &Path) -> (bool, u64) { + let Ok(entries) = fs::read_dir(segment_dir) else { + return (false, 0); + }; + let mut has_durable_media = false; + let mut durable_byte_count = 0; + for entry in entries.flatten() { + let path = entry.path(); + if path + .file_name() + .is_some_and(|name| name == METADATA_FILENAME) + { + continue; + } + if !path.is_file() { + continue; + } + has_durable_media = true; + durable_byte_count += fs::metadata(&path) + .map(|metadata| metadata.len()) + .unwrap_or(0); + } + (has_durable_media, durable_byte_count) +} + pub fn read_segment_start(segment_dir: &Path) -> Option { let value: Value = serde_json::from_str(&fs::read_to_string(segment_dir.join(METADATA_FILENAME)).ok()?) @@ -333,7 +381,7 @@ mod tests { fn metadata_duration() { let t = tempfile::tempdir().unwrap(); let path = incomplete(t.path(), "140000", 1000.0); - write_segment_metadata(&path, 940.0); + write_segment_metadata(&path, 940.0, SegmentProgress::default()); age(&path, 1000.0); recover_incomplete_segments(t.path(), 300, 1000.0, &NoMedia); assert!(path.with_file_name("140000_60").exists()); @@ -343,7 +391,7 @@ mod tests { fn metadata_ceiling() { let t = tempfile::tempdir().unwrap(); let path = incomplete(t.path(), "140000", 1000.0); - write_segment_metadata(&path, 0.0); + write_segment_metadata(&path, 0.0, SegmentProgress::default()); age(&path, 1000.0); recover_incomplete_segments(t.path(), 60, 1000.0, &NoMedia); assert!(path.with_file_name("140000_60").exists()); @@ -366,7 +414,7 @@ mod tests { path.join("audio.flac"), ) .unwrap(); - write_segment_metadata(&path, 0.0); + write_segment_metadata(&path, 0.0, SegmentProgress::default()); age(&path, 1000.0); recover_incomplete_segments(t.path(), 300, 1000.0, &ClaxonMediaDurationProbe); assert!(path.with_file_name("140000_4").exists()); @@ -376,7 +424,7 @@ mod tests { fn webm_uses_ceiling() { let t = tempfile::tempdir().unwrap(); let path = incomplete(t.path(), "140000", 1000.0); - write_segment_metadata(&path, 0.0); + write_segment_metadata(&path, 0.0, SegmentProgress::default()); age(&path, 1000.0); recover_incomplete_segments(t.path(), 60, 1000.0, &NoMedia); assert!(path.with_file_name("140000_60").exists()); @@ -427,7 +475,7 @@ mod tests { fn metadata_removed() { let t = tempfile::tempdir().unwrap(); let path = incomplete(t.path(), "140000", 1000.0); - write_segment_metadata(&path, 940.0); + write_segment_metadata(&path, 940.0, SegmentProgress::default()); age(&path, 1000.0); recover_incomplete_segments(t.path(), 300, 1000.0, &NoMedia); assert!( @@ -453,7 +501,7 @@ mod tests { assert_eq!(filesystem_duration(0.5, 0.0, 300), 1); let t = tempfile::tempdir().unwrap(); let path = incomplete(t.path(), "140000", 1000.0); - write_segment_metadata(&path, 999.5); + write_segment_metadata(&path, 999.5, SegmentProgress::default()); age(&path, 1000.0); recover_incomplete_segments(t.path(), 300, 1000.0, &FixedMedia(0.5)); assert!(path.with_file_name("140000_1").exists()); @@ -463,13 +511,13 @@ mod tests { fn rename_failure_continues() { let t = tempfile::tempdir().unwrap(); let first = incomplete(t.path(), "120000", 1000.0); - write_segment_metadata(&first, 940.0); + write_segment_metadata(&first, 940.0, SegmentProgress::default()); age(&first, 1000.0); let collision = first.with_file_name("120000_60"); fs::create_dir(&collision).unwrap(); fs::write(collision.join("occupied"), b"x").unwrap(); let second = incomplete(t.path(), "130000", 1000.0); - write_segment_metadata(&second, 940.0); + write_segment_metadata(&second, 940.0, SegmentProgress::default()); age(&second, 1000.0); assert_eq!( recover_incomplete_segments(t.path(), 300, 1000.0, &NoMedia), @@ -490,7 +538,7 @@ mod tests { let path = incomplete(t.path(), "140000", 1000.0); fs::remove_file(path.join("screen.webm")).unwrap(); fs::create_dir(path.join("nested")).unwrap(); - write_segment_metadata(&path, 940.0); + write_segment_metadata(&path, 940.0, SegmentProgress::default()); assert!(recover_segment(&path, 300, 1000.0, &NoMedia)); } // AC: media candidates use max duration and swallow individual failures. @@ -548,7 +596,7 @@ mod tests { let segment = actual_day.join("archon/140000.incomplete"); fs::create_dir_all(&segment).unwrap(); fs::write(segment.join("screen.webm"), b"x").unwrap(); - write_segment_metadata(&segment, 940.0); + write_segment_metadata(&segment, 940.0, SegmentProgress::default()); age(&segment, 1000.0); fs::create_dir(&captures).unwrap(); symlink(&actual_day, captures.join("20260403")).unwrap(); @@ -559,4 +607,38 @@ mod tests { ); assert!(actual_day.join("archon/140000_60").exists()); } + + // Existence, not byte_count > 0: a 0-byte leftover is durable media (D3). + #[test] + fn scan_zero_byte_file_is_media() { + let t = tempfile::tempdir().unwrap(); + fs::write(t.path().join(METADATA_FILENAME), b"{}\n").unwrap(); + fs::write(t.path().join("leftover.bin"), b"").unwrap(); + assert_eq!(scan_segment_progress(t.path()), (true, 0)); + } + + // AC5: in-flight legacy sidecars from the previous writer still recover. + #[test] + fn legacy_sidecar_reads_and_recovers() { + const LEGACY: &[u8] = b"{\"start_timestamp\":1700000000.00000000}"; + assert_eq!(LEGACY.len(), 39); + + let t = tempfile::tempdir().unwrap(); + let now = 1_700_000_060.0; + let path = incomplete(t.path(), "140000", now); + fs::write(path.join(METADATA_FILENAME), LEGACY).unwrap(); + age(&path, now); + assert_eq!(read_segment_start(&path), Some(1_700_000_000.0)); + assert_eq!(recover_incomplete_segments(t.path(), 300, now, &NoMedia), 1); + assert!(path.with_file_name("140000_60").exists()); + + let extra = t.path().join("extra"); + fs::create_dir(&extra).unwrap(); + fs::write( + extra.join(METADATA_FILENAME), + br#"{"start_timestamp":1700000000.00000000,"unknown":true,"also":1}"#, + ) + .unwrap(); + assert_eq!(read_segment_start(&extra), Some(1_700_000_000.0)); + } } diff --git a/crates/solstone-linux/src/run.rs b/crates/solstone-linux/src/run.rs index e2376b9..e72464a 100644 --- a/crates/solstone-linux/src/run.rs +++ b/crates/solstone-linux/src/run.rs @@ -665,6 +665,22 @@ impl SyncWake for SyncTrigger { } } +fn stream_silent_fields(event: &StreamSilentEvent) -> Map { + let mut fields = Map::new(); + fields.insert("connector".into(), json!(event.connector)); + fields.insert("position".into(), json!(event.position)); + fields.insert("node_id".into(), json!(event.node_id)); + fields.insert("file_bytes".into(), json!(event.file_bytes)); + fields.insert("segment_dir".into(), json!(event.segment_dir)); + fields.insert("duration_seconds".into(), json!(event.duration_seconds)); + fields.insert("host".into(), json!(event.host)); + fields.insert("platform".into(), json!(event.platform)); + if let Some(reason) = event.reason { + fields.insert("reason".into(), json!(reason)); + } + fields +} + struct UploadEventSink { client: Arc, sync: W, @@ -675,16 +691,9 @@ impl EventSink for UploadEventSink { self.client.enqueue_status(fields); } fn stream_silent(&mut self, event: StreamSilentEvent) { - let mut fields = Map::new(); - fields.insert("connector".into(), json!(event.connector)); - fields.insert("position".into(), json!(event.position)); - fields.insert("node_id".into(), json!(event.node_id)); - fields.insert("file_bytes".into(), json!(event.file_bytes)); - fields.insert("segment_dir".into(), json!(event.segment_dir)); - fields.insert("duration_seconds".into(), json!(event.duration_seconds)); - fields.insert("host".into(), json!(event.host)); - fields.insert("platform".into(), json!(event.platform)); - let _ = self.client.enqueue_stream_silent(fields); + let _ = self + .client + .enqueue_stream_silent(stream_silent_fields(&event)); } fn segment_completed(&mut self, event: SegmentCompletedEvent) { tracing::debug!(key = %event.key, "segment completed"); @@ -1000,6 +1009,7 @@ mod tests { duration_seconds: 1, host: "h".into(), platform: "linux".into(), + reason: None, }); assert_eq!(count.load(Ordering::Acquire), 0); sink.segment_completed(SegmentCompletedEvent { @@ -1011,6 +1021,45 @@ mod tests { client.stop(Duration::from_secs(1)).await; } + // AC9: reason is omitted when None and present only when Some. + #[test] + fn stream_silent_omits_reason_unless_some() { + let none = stream_silent_fields(&StreamSilentEvent { + connector: "c".into(), + position: "p".into(), + node_id: 1, + file_bytes: 0, + segment_dir: "s".into(), + duration_seconds: 1, + host: "h".into(), + platform: "linux".into(), + reason: None, + }); + assert!(none.get("reason").is_none()); + assert_eq!(none["connector"], json!("c")); + assert_eq!(none["position"], json!("p")); + assert_eq!(none["node_id"], json!(1)); + assert_eq!(none["file_bytes"], json!(0)); + assert_eq!(none["segment_dir"], json!("s")); + assert_eq!(none["duration_seconds"], json!(1)); + assert_eq!(none["host"], json!("h")); + assert_eq!(none["platform"], json!("linux")); + + let some = stream_silent_fields(&StreamSilentEvent { + connector: "c".into(), + position: "p".into(), + node_id: 1, + file_bytes: 0, + segment_dir: "s".into(), + duration_seconds: 1, + host: "h".into(), + platform: "linux".into(), + reason: Some(crate::observer::EMPTY_WINDOW_REASON), + }); + assert_eq!(some["reason"], json!("empty_window")); + assert_eq!(some["connector"], json!("c")); + } + // AC: 8 — desktop tasks stop before final observer work, walker join, and sender stop. #[tokio::test] async fn shutdown_order_includes_linked_owner_last() { diff --git a/crates/solstone-linux/src/segment.rs b/crates/solstone-linux/src/segment.rs index ea8cd60..e8860ae 100644 --- a/crates/solstone-linux/src/segment.rs +++ b/crates/solstone-linux/src/segment.rs @@ -44,7 +44,7 @@ pub fn finalize_segment_dir(incomplete: &Path, key: &str) -> io::Result #[cfg(test)] mod tests { use super::*; - use crate::recovery::{read_segment_start, write_segment_metadata}; + use crate::recovery::{SegmentProgress, read_segment_start, write_segment_metadata}; // observer.py::_get_timestamp_parts shape contract. #[test] @@ -87,7 +87,7 @@ mod tests { #[test] fn metadata_round_trip() { let t = tempfile::tempdir().unwrap(); - write_segment_metadata(t.path(), 1234.5); + write_segment_metadata(t.path(), 1234.5, SegmentProgress::default()); assert_eq!(read_segment_start(t.path()), Some(1234.5)); } } -- 2.51.2