Skip to main content

bambu_rs/
park.rs

1//! Live "parked frame per layer" capture — the reusable I/O runner for the smooth
2//! timelapse. The pure detector lives in [`crate::core::park`]; this owns the I/O and is
3//! a first-class library module so the CLI, the server, and other consumers all drive it
4//! the same way (not locked behind the `server` feature).
5//!
6//! One ffmpeg per stream camera opens the camera's MJPEG stream and tees a tiny gray
7//! rawvideo (read here, fed to the detector) plus full-resolution JPEGs to a ring dir, so
8//! the emitted preview is full-res without decoding JPEGs in Rust. On each emitted park
9//! the chosen ring JPEG is copied atomically to `latest_park.jpg` (what a dashboard
10//! shows) plus `park_NNNNNN.jpg`, and a line is appended to `parks.jsonl`.
11//!
12//! ffmpeg is the only external tool (already a dependency for the plain recorder); there
13//! is no python3 runtime dependency. The frame-reading + write-mapping logic is injected
14//! so it's unit-tested with an in-memory stream and fake ring files — the ffmpeg spawn
15//! itself is the thin, on-device-verified seam, and it reports progress through a
16//! callback rather than any server type, so every caller adapts it to its own output.
17
18use std::io::{Read, Write};
19use std::path::{Path, PathBuf};
20use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
21use std::sync::{Arc, Mutex};
22use std::time::Duration;
23
24use crate::core::park::{
25    LiveParkDetector, Park, ParkTuning, SegmentPick, SegmentSelector, SelectTuning,
26};
27
28/// Detection decode size (tiny grayscale): enough signal for the left-zone park, cheap
29/// to score per frame. Matches the Python miner's default.
30pub const DECODE_W: usize = 64;
31pub const DECODE_H: usize = 36;
32
33/// ffmpeg argv (after the program name) for the live park tee: one MJPEG stream in, two
34/// frame-synced outputs — a tiny gray rawvideo to stdout (`pipe:1`, for detection) and
35/// full-res JPEGs to `ring_dir/full_%09d.jpg` (index-aligned to the gray frames). Pure,
36/// so the command shape is unit-tested without ffmpeg.
37pub fn live_park_args(
38    stream_url: &str,
39    ring_dir: &Path,
40    fps: f64,
41    w: usize,
42    h: usize,
43) -> Vec<String> {
44    let ring = ring_dir.join("full_%09d.jpg");
45    vec![
46        "-v".into(),
47        "error".into(),
48        "-f".into(),
49        "mpjpeg".into(),
50        "-i".into(),
51        stream_url.into(),
52        "-filter_complex".into(),
53        format!("[0:v]fps={fps},split=2[full][det];[det]scale={w}:{h},format=gray[gray]"),
54        "-map".into(),
55        "[gray]".into(),
56        "-f".into(),
57        "rawvideo".into(),
58        "pipe:1".into(),
59        "-map".into(),
60        "[full]".into(),
61        "-start_number".into(),
62        "0".into(),
63        "-q:v".into(),
64        "3".into(),
65        ring.display().to_string(),
66    ]
67}
68
69/// The full-res ring JPEG path ffmpeg writes for gray frame `idx`.
70pub fn ring_jpeg_path(ring_dir: &Path, idx: u64) -> PathBuf {
71    ring_dir.join(format!("full_{idx:09}.jpg"))
72}
73
74/// Read exactly `buf.len()` bytes (one gray frame) from a blocking reader. Returns false
75/// on EOF or error (a short/partial read at stream end counts as EOF).
76fn read_full(reader: &mut dyn Read, buf: &mut [u8]) -> bool {
77    let mut filled = 0;
78    while filled < buf.len() {
79        match reader.read(&mut buf[filled..]) {
80            Ok(0) => return false,
81            Ok(n) => filled += n,
82            Err(_) => return false,
83        }
84    }
85    true
86}
87
88/// Writes emitted parks to disk. A `replace` park (a stronger close pair — same layer)
89/// reuses the previous index and overwrites, so the timelapse keeps exactly one frame
90/// per park; `latest_park.jpg` is updated atomically (temp + rename) so a polling
91/// dashboard never reads a half-written file.
92pub struct ParkWriter {
93    out_dir: PathBuf,
94    emitted: u64,
95}
96
97impl ParkWriter {
98    pub fn new(out_dir: PathBuf) -> Self {
99        Self {
100            out_dir,
101            emitted: 0,
102        }
103    }
104
105    /// Number of distinct park frames written so far (a `replace` doesn't increment it).
106    pub fn emitted(&self) -> u64 {
107        self.emitted
108    }
109
110    /// Copy `ring_jpeg` to `latest_park.jpg` (atomic) and `park_NNNNNN.jpg`, and append a
111    /// line to `parks.jsonl`. Returns the index written and whether it REPLACED the
112    /// previous park (a stronger close pair, same layer — overwritten, not a new frame).
113    pub fn write(&mut self, park: &Park, ring_jpeg: &Path) -> std::io::Result<ParkWritten> {
114        let replaced = park.replace && self.emitted > 0;
115        let index = if replaced {
116            self.emitted - 1
117        } else {
118            self.emitted
119        };
120
121        let tmp = self.out_dir.join("latest_park.jpg.tmp");
122        std::fs::copy(ring_jpeg, &tmp)?;
123        std::fs::rename(&tmp, self.out_dir.join("latest_park.jpg"))?;
124        std::fs::copy(ring_jpeg, self.out_dir.join(format!("park_{index:06}.jpg")))?;
125
126        let line = serde_json::json!({
127            "n": index, "idx": park.idx, "t": park.t, "left_mass": park.left_mass,
128            "sharpness": park.sharpness, "confidence": park.confidence, "replace": replaced,
129        });
130        let mut f = std::fs::OpenOptions::new()
131            .create(true)
132            .append(true)
133            .open(self.out_dir.join("parks.jsonl"))?;
134        writeln!(f, "{line}")?;
135
136        if !replaced {
137            self.emitted += 1;
138        }
139        Ok(ParkWritten { index, replaced })
140    }
141}
142
143/// Result of [`ParkWriter::write`]: the index of the `park_NNNNNN.jpg` written, and
144/// whether it overwrote the previous park (a same-layer supersession) vs. added a new one.
145pub struct ParkWritten {
146    pub index: u64,
147    pub replaced: bool,
148}
149
150/// What happened to one emitted park — the live progress hook's event. A `Replaced` is
151/// NOT a new layer: it overwrote the previous park with a stronger frame, so it must not
152/// be counted as an additional park.
153#[derive(Debug, Clone, Copy, PartialEq, Eq)]
154pub enum ParkEvent {
155    /// A new distinct park frame was written (`park_NNNNNN.jpg` added).
156    Written,
157    /// A stronger close pair superseded the previous park (overwritten, same layer).
158    Replaced,
159    /// The park's ring JPEG never arrived / the write failed — the frame is lost.
160    Dropped,
161}
162
163/// Outcome of one detection run over a gray stream.
164#[derive(Debug, Default, PartialEq, Eq)]
165pub struct ParkRunStats {
166    /// Gray frames read from the stream.
167    pub frames: u64,
168    /// Distinct parks written (one `park_NNNNNN.jpg` each; a replace does NOT count here).
169    pub parks: u64,
170    /// Parks superseded by a stronger close pair (overwrote an existing frame).
171    pub replaced: u64,
172    /// Parks whose full-res ring JPEG never arrived / failed to write — dropped with a
173    /// count rather than silently, so a missing layer is visible.
174    pub dropped: u64,
175}
176
177/// Resolve one emitted park to disk: await its ring JPEG (briefly), write it, and
178/// tally/report the outcome. `on_emit` fires per park so the caller can update live
179/// progress without waiting for the run to end — distinguishing a new park from a
180/// supersession so neither over-counts.
181fn emit_park(
182    park: &Park,
183    ring_jpeg: &dyn Fn(u64) -> PathBuf,
184    await_ring: &dyn Fn(&Path) -> bool,
185    writer: &mut ParkWriter,
186    on_emit: &mut dyn FnMut(ParkEvent),
187    stats: &mut ParkRunStats,
188) {
189    let path = ring_jpeg(park.idx);
190    let event = if await_ring(&path) {
191        match writer.write(park, &path) {
192            Ok(w) if w.replaced => ParkEvent::Replaced,
193            Ok(_) => ParkEvent::Written,
194            Err(_) => ParkEvent::Dropped,
195        }
196    } else {
197        ParkEvent::Dropped
198    };
199    match event {
200        ParkEvent::Written => stats.parks += 1,
201        ParkEvent::Replaced => stats.replaced += 1,
202        ParkEvent::Dropped => stats.dropped += 1,
203    }
204    on_emit(event);
205}
206
207/// Drive the detector over a gray rawvideo stream, writing each emitted park. Pure of
208/// ffmpeg: the gray `reader`, the per-idx ring path, the "is the ring JPEG ready yet?"
209/// wait, `cancel`, and the `on_emit` progress hook are all injected, so this is
210/// unit-tested with an in-memory stream and pre-created fake ring files. `await_ring`
211/// lets the real worker poll briefly for the JPEG (it can lag the gray by a tick) while
212/// a test returns instantly; `on_emit` fires a [`ParkEvent`] per park for live status.
213#[allow(clippy::too_many_arguments)]
214pub fn detect_stream(
215    reader: &mut dyn Read,
216    detector: &mut LiveParkDetector,
217    frame_size: usize,
218    writer: &mut ParkWriter,
219    ring_jpeg: &dyn Fn(u64) -> PathBuf,
220    await_ring: &dyn Fn(&Path) -> bool,
221    cancel: &dyn Fn() -> bool,
222    on_emit: &mut dyn FnMut(ParkEvent),
223) -> ParkRunStats {
224    let mut stats = ParkRunStats::default();
225    let mut buf = vec![0u8; frame_size];
226    let mut idx: u64 = 0;
227    loop {
228        if cancel() || !read_full(reader, &mut buf) {
229            break;
230        }
231        if let Some(park) = detector.push(&buf, idx) {
232            emit_park(&park, ring_jpeg, await_ring, writer, on_emit, &mut stats);
233        }
234        idx += 1;
235        stats.frames += 1;
236    }
237    if let Some(park) = detector.flush() {
238        emit_park(&park, ring_jpeg, await_ring, writer, on_emit, &mut stats);
239    }
240    stats
241}
242
243/// Resolve one segment pick to disk: it carries no left-mass/sharpness (the median-subtract
244/// selector only yields an offset + confidence), and its `t` records the print LAYER rather
245/// than a timestamp — what a per-layer timelapse keys on. A pick is always a fresh frame
246/// (one per layer), never a same-layer supersession, so `replace` is false.
247fn emit_segment_pick(
248    pick: &SegmentPick,
249    ring_jpeg: &dyn Fn(u64) -> PathBuf,
250    await_ring: &dyn Fn(&Path) -> bool,
251    writer: &mut ParkWriter,
252    on_emit: &mut dyn FnMut(ParkEvent),
253    stats: &mut ParkRunStats,
254) {
255    let park = Park {
256        idx: pick.idx,
257        t: pick.layer as f64,
258        left_mass: 0.0,
259        sharpness: 0.0,
260        confidence: pick.confidence,
261        replace: false,
262    };
263    emit_park(&park, ring_jpeg, await_ring, writer, on_emit, stats);
264}
265
266/// Drive the dense-stream [`SegmentSelector`] over a gray rawvideo stream, writing the
267/// picked frame per print layer. The continuous stream alternative to [`detect_stream`]'s
268/// camera-only miner: instead of detecting the park by image change (which misses the
269/// brief native park), it segments the stream by LAYER (`layer_of(idx)` — the test maps
270/// idx→layer, the server reads a live atomic fed by MQTT) and median-subtract-selects the
271/// parked frame within each layer's first `window_ms`. Frame capture time is modeled from
272/// the fixed ffmpeg cadence (`idx / fps`), matching the device model. Pure of ffmpeg the
273/// same way `detect_stream` is — `reader`, `ring_jpeg`, `await_ring`, `cancel`, `on_emit`
274/// are injected, so it's unit-tested with an in-memory modeled stream and fake ring files.
275#[allow(clippy::too_many_arguments)]
276pub fn detect_stream_segmented(
277    reader: &mut dyn Read,
278    sel: &mut SegmentSelector,
279    frame_size: usize,
280    fps: f64,
281    layer_of: &dyn Fn(u64) -> i64,
282    writer: &mut ParkWriter,
283    ring_jpeg: &dyn Fn(u64) -> PathBuf,
284    await_ring: &dyn Fn(&Path) -> bool,
285    cancel: &dyn Fn() -> bool,
286    on_emit: &mut dyn FnMut(ParkEvent),
287) -> ParkRunStats {
288    let mut stats = ParkRunStats::default();
289    let mut buf = vec![0u8; frame_size];
290    let mut idx: u64 = 0;
291    loop {
292        if cancel() || !read_full(reader, &mut buf) {
293            break;
294        }
295        let layer = layer_of(idx);
296        let t_ms = (idx as f64 * 1000.0 / fps) as u64;
297        if let Some(pick) = sel.push(layer, idx, t_ms, buf.clone()) {
298            emit_segment_pick(&pick, ring_jpeg, await_ring, writer, on_emit, &mut stats);
299        }
300        idx += 1;
301        stats.frames += 1;
302    }
303    if let Some(pick) = sel.finish() {
304        emit_segment_pick(&pick, ring_jpeg, await_ring, writer, on_emit, &mut stats);
305    }
306    stats
307}
308
309/// One stream camera selected for live park detection: its id, the MJPEG stream URL
310/// ffmpeg opens, and its per-camera tuning (framing is camera-specific, so the tuning is
311/// too — there are no shared defaults).
312#[derive(Clone)]
313pub struct ParkCapture {
314    pub id: String,
315    pub stream_url: String,
316    pub tuning: ParkTuning,
317}
318
319/// Delete all but the newest `keep` ring JPEGs (`full_<idx>.jpg`) to bound disk during a
320/// long print — ffmpeg writes one full-res JPEG per gray frame. The newest are kept, so
321/// an in-flight park (its idx is within a few frames of the latest) is never pruned out
322/// from under [`detect_stream`]'s `await_ring`. Returns how many were removed.
323pub fn prune_ring(ring_dir: &Path, keep: usize) -> usize {
324    let mut entries: Vec<(u64, PathBuf)> = match std::fs::read_dir(ring_dir) {
325        Ok(rd) => rd
326            .flatten()
327            .filter_map(|e| {
328                let p = e.path();
329                let idx = p
330                    .file_name()
331                    .and_then(|f| f.to_str())
332                    .and_then(|f| f.strip_prefix("full_"))
333                    .and_then(|f| f.strip_suffix(".jpg"))
334                    .and_then(|f| f.parse::<u64>().ok())?;
335                Some((idx, p))
336            })
337            .collect(),
338        Err(_) => return 0,
339    };
340    if entries.len() <= keep {
341        return 0;
342    }
343    entries.sort_by_key(|(idx, _)| *idx);
344    let remove = entries.len() - keep;
345    entries
346        .into_iter()
347        .take(remove)
348        .filter(|(_, p)| std::fs::remove_file(p).is_ok())
349        .count()
350}
351
352/// Seconds of full-res ring JPEGs to retain (sliding window) — well above the few-frame
353/// park lag, so an in-flight park's JPEG is always still present when it's written out.
354const RING_KEEP_SECONDS: f64 = 60.0;
355
356/// A spawned park-tee ffmpeg: the gray rawvideo stdout to read, plus the handles to shut it
357/// down cleanly. Shared by [`run_park_camera`] and [`run_segment_camera`] — both open the
358/// identical MJPEG tee ([`live_park_args`]) and differ only in HOW they consume the gray
359/// stream (image-change detection vs. layer segmenting).
360struct ParkFfmpeg {
361    stdout: std::process::ChildStdout,
362    child: Arc<Mutex<std::process::Child>>,
363    done: Arc<AtomicBool>,
364    aux: std::thread::JoinHandle<()>,
365}
366
367impl ParkFfmpeg {
368    /// Tell the aux thread to stop, join it, and reap ffmpeg. The caller removes the ring
369    /// dir afterward (it owns the path).
370    fn shutdown(self) {
371        self.done.store(true, Ordering::Relaxed);
372        let _ = self.aux.join();
373        let _ = self.child.lock().unwrap().wait();
374    }
375}
376
377/// Spawn the park-tee ffmpeg for one camera, teeing into `ring`, plus the aux thread that
378/// KILLS ffmpeg the instant `cancel` is set (the consuming read blocks until bytes arrive,
379/// so it can't observe `cancel` itself) and prunes the ring on a sliding window. `id` only
380/// labels the spawn error. This is the thin, on-device-verified I/O seam both runners share.
381fn spawn_park_ffmpeg(
382    id: &str,
383    stream_url: &str,
384    ring: &Path,
385    fps: f64,
386    w: usize,
387    h: usize,
388    cancel: &Arc<AtomicBool>,
389) -> Result<ParkFfmpeg, String> {
390    let mut child = std::process::Command::new("ffmpeg")
391        .args(live_park_args(stream_url, ring, fps, w, h))
392        .stdin(std::process::Stdio::null())
393        .stdout(std::process::Stdio::piped())
394        .stderr(std::process::Stdio::null())
395        .spawn()
396        .map_err(|e| format!("park {id}: ffmpeg spawn failed: {e}"))?;
397    let stdout = child.stdout.take().expect("piped stdout");
398    let child = Arc::new(Mutex::new(child));
399    let done = Arc::new(AtomicBool::new(false));
400
401    let aux = {
402        let (cancel, done, child, ring) = (
403            cancel.clone(),
404            done.clone(),
405            child.clone(),
406            ring.to_path_buf(),
407        );
408        let keep = (RING_KEEP_SECONDS * fps).max(40.0) as usize;
409        std::thread::spawn(move || {
410            let mut ticks = 0u32;
411            loop {
412                if cancel.load(Ordering::Relaxed) {
413                    let _ = child.lock().unwrap().kill();
414                    return;
415                }
416                if done.load(Ordering::Relaxed) {
417                    return;
418                }
419                ticks += 1;
420                if ticks.is_multiple_of(4) {
421                    prune_ring(&ring, keep);
422                }
423                std::thread::sleep(Duration::from_millis(500));
424            }
425        })
426    };
427
428    Ok(ParkFfmpeg {
429        stdout,
430        child,
431        done,
432        aux,
433    })
434}
435
436/// Run live park detection for ONE stream camera until the stream ends or `cancel` is
437/// set, writing its output into `cam_dir` (`latest_park.jpg` + `park_NNNNNN.jpg` +
438/// `parks.jsonl`, with a transient `.ring` subdir). The caller picks the dir — the server
439/// uses one per camera under the run dir; the CLI passes its `--out` directly. Spawns one
440/// ffmpeg ([`live_park_args`]) that opens the MJPEG stream and tees a tiny gray rawvideo
441/// (read here → the detector) plus full-res JPEGs to the ring; on each park the chosen
442/// JPEG is written to `latest_park.jpg`. An aux thread KILLS ffmpeg the instant `cancel`
443/// is set — the main read blocks until bytes arrive, so it can't observe `cancel` itself —
444/// and prunes the ring on a sliding window.
445///
446/// Reports each park live via `on_park` ([`ParkEvent`]); returns the run's
447/// [`ParkRunStats`], or an error string if it couldn't even start (no server/CLI type
448/// leaks in, so any caller adapts it). `stats.frames == 0` on return means the stream
449/// produced nothing — the caller decides how to surface that. This is the thin,
450/// on-device-verified I/O seam; the detector, argv, writer, and prune are the unit-tested
451/// pure pieces.
452pub fn run_park_camera(
453    cap: &ParkCapture,
454    cam_dir: &Path,
455    w: usize,
456    h: usize,
457    cancel: &Arc<AtomicBool>,
458    on_park: &mut dyn FnMut(ParkEvent),
459) -> Result<ParkRunStats, String> {
460    let ring = cam_dir.join(".ring");
461    std::fs::create_dir_all(&ring)
462        .map_err(|e| format!("park {}: create {}: {e}", cap.id, ring.display()))?;
463
464    let mut ff = spawn_park_ffmpeg(
465        &cap.id,
466        &cap.stream_url,
467        &ring,
468        cap.tuning.fps,
469        w,
470        h,
471        cancel,
472    )?;
473
474    let mut det = LiveParkDetector::new(w, h, &cap.tuning);
475    let mut writer = ParkWriter::new(cam_dir.to_path_buf());
476    let ring_path = ring.clone();
477    let read_cancel = cancel.clone();
478    let wait_cancel = cancel.clone();
479    let stats = detect_stream(
480        &mut ff.stdout,
481        &mut det,
482        w * h,
483        &mut writer,
484        &|idx| ring_jpeg_path(&ring_path, idx),
485        &|p| await_ring(p, &wait_cancel),
486        &|| read_cancel.load(Ordering::Relaxed),
487        on_park,
488    );
489
490    ff.shutdown();
491    let _ = std::fs::remove_dir_all(&ring);
492    Ok(stats)
493}
494
495/// One stream camera for dense-stream segmented capture: its id, the MJPEG stream URL, the
496/// capture fps (the gray cadence), the per-layer accumulation `window_ms`, and the
497/// median-subtract select knobs. Like [`ParkCapture`] these are all camera-specific
498/// (framing is) — there are no shared defaults.
499#[derive(Clone)]
500pub struct SegmentCapture {
501    pub id: String,
502    pub stream_url: String,
503    pub fps: f64,
504    pub window_ms: u64,
505    pub select_tuning: SelectTuning,
506}
507
508/// Run dense-stream segmented capture for ONE stream camera until the stream ends or
509/// `cancel` is set, writing the picked frame per print layer into `cam_dir` (the same
510/// `latest_park.jpg` + `park_NNNNNN.jpg` + `parks.jsonl` layout as [`run_park_camera`], so
511/// the dashboard reads it identically). The continuous-stream alternative to
512/// `run_park_camera`: instead of detecting the park by image change (which misses the brief
513/// native park), it segments the stream by the live print LAYER and median-subtract-selects
514/// the parked frame within each layer's window. `current_layer` is the live layer the
515/// server updates from MQTT `layer_num`; the runner just reads it per frame (a `< 0`
516/// "unknown" value still segments by whatever it reads). The ffmpeg spawn + aux are shared
517/// with `run_park_camera`; the segmenting and selection are the unit-tested pure pieces.
518pub fn run_segment_camera(
519    cap: &SegmentCapture,
520    cam_dir: &Path,
521    w: usize,
522    h: usize,
523    current_layer: &Arc<AtomicI64>,
524    cancel: &Arc<AtomicBool>,
525    on_park: &mut dyn FnMut(ParkEvent),
526) -> Result<ParkRunStats, String> {
527    let ring = cam_dir.join(".ring");
528    std::fs::create_dir_all(&ring)
529        .map_err(|e| format!("park {}: create {}: {e}", cap.id, ring.display()))?;
530
531    let mut ff = spawn_park_ffmpeg(&cap.id, &cap.stream_url, &ring, cap.fps, w, h, cancel)?;
532
533    let mut sel = SegmentSelector::new(w, h, cap.window_ms, cap.select_tuning);
534    let mut writer = ParkWriter::new(cam_dir.to_path_buf());
535    let ring_path = ring.clone();
536    let read_cancel = cancel.clone();
537    let wait_cancel = cancel.clone();
538    let layer = current_layer.clone();
539    let stats = detect_stream_segmented(
540        &mut ff.stdout,
541        &mut sel,
542        w * h,
543        cap.fps,
544        &|_idx| layer.load(Ordering::Relaxed),
545        &mut writer,
546        &|idx| ring_jpeg_path(&ring_path, idx),
547        &|p| await_ring(p, &wait_cancel),
548        &|| read_cancel.load(Ordering::Relaxed),
549        on_park,
550    );
551
552    ff.shutdown();
553    let _ = std::fs::remove_dir_all(&ring);
554    Ok(stats)
555}
556
557/// Poll for a ring JPEG to appear (it can lag its gray frame by a tick), up to ~500ms,
558/// bailing early if cancelled. Bounded — a JPEG that never arrives is dropped, never
559/// waited on forever (which would stall draining ffmpeg's stdout).
560fn await_ring(path: &Path, cancel: &Arc<AtomicBool>) -> bool {
561    for _ in 0..10 {
562        if path.exists() {
563            return true;
564        }
565        if cancel.load(Ordering::Relaxed) {
566            return false;
567        }
568        std::thread::sleep(Duration::from_millis(50));
569    }
570    path.exists()
571}
572
573#[cfg(test)]
574mod tests {
575    use super::*;
576    use std::io::Cursor;
577
578    #[test]
579    fn live_park_args_tees_synced_gray_and_full_outputs() {
580        let args = live_park_args("http://cam/stream", Path::new("/ring"), 4.0, 64, 36);
581        let joined = args.join(" ");
582        assert!(
583            joined.contains("-f mpjpeg -i http://cam/stream"),
584            "{joined}"
585        );
586        assert!(
587            joined.contains("split=2[full][det]"),
588            "one input, two outputs: {joined}"
589        );
590        assert!(
591            joined.contains("scale=64:36,format=gray"),
592            "detection stream: {joined}"
593        );
594        assert!(joined.contains("rawvideo pipe:1"), "{joined}");
595        assert!(
596            joined.trim_end().ends_with("/ring/full_%09d.jpg"),
597            "full-res ring: {joined}"
598        );
599        assert!(
600            joined.contains("fps=4,"),
601            "fps renders ffmpeg-friendly: {joined}"
602        );
603    }
604
605    fn tmp(tag: &str) -> PathBuf {
606        let d = std::env::temp_dir().join(format!("bambu-park-{tag}-{}", std::process::id()));
607        let _ = std::fs::remove_dir_all(&d);
608        std::fs::create_dir_all(&d).unwrap();
609        d
610    }
611
612    fn park(idx: u64, replace: bool) -> Park {
613        Park {
614            idx,
615            t: idx as f64 / 4.0,
616            left_mass: 9000.0,
617            sharpness: 500.0,
618            confidence: 0.9,
619            replace,
620        }
621    }
622
623    #[test]
624    fn writer_writes_latest_and_indexed_and_jsonl() {
625        let dir = tmp("writer");
626        let ring = dir.join("full_000000003.jpg");
627        std::fs::write(&ring, b"JPEGA").unwrap();
628        let mut w = ParkWriter::new(dir.clone());
629        let wr = w.write(&park(3, false), &ring).unwrap();
630        assert_eq!(wr.index, 0);
631        assert!(!wr.replaced);
632        assert_eq!(w.emitted(), 1);
633        assert_eq!(
634            std::fs::read(dir.join("latest_park.jpg")).unwrap(),
635            b"JPEGA"
636        );
637        assert_eq!(
638            std::fs::read(dir.join("park_000000.jpg")).unwrap(),
639            b"JPEGA"
640        );
641        let jl = std::fs::read_to_string(dir.join("parks.jsonl")).unwrap();
642        assert_eq!(jl.lines().count(), 1);
643        assert!(jl.contains("\"n\":0"), "{jl}");
644        let _ = std::fs::remove_dir_all(&dir);
645    }
646
647    #[test]
648    fn a_replace_park_overwrites_the_previous_index() {
649        let dir = tmp("replace");
650        let r1 = dir.join("full_000000003.jpg");
651        let r2 = dir.join("full_000000005.jpg");
652        std::fs::write(&r1, b"WEAK").unwrap();
653        std::fs::write(&r2, b"STRONG").unwrap();
654        let mut w = ParkWriter::new(dir.clone());
655        w.write(&park(3, false), &r1).unwrap(); // n=0, emitted -> 1
656        let wr = w.write(&park(5, true), &r2).unwrap(); // replace: reuse n=0, emitted stays 1
657        assert_eq!(wr.index, 0, "replace reuses the previous index");
658        assert!(wr.replaced, "flagged as a replacement");
659        assert_eq!(w.emitted(), 1, "replace does not add a frame");
660        assert_eq!(
661            std::fs::read(dir.join("park_000000.jpg")).unwrap(),
662            b"STRONG",
663            "overwritten"
664        );
665        assert_eq!(
666            std::fs::read(dir.join("latest_park.jpg")).unwrap(),
667            b"STRONG"
668        );
669        assert_eq!(
670            std::fs::read_to_string(dir.join("parks.jsonl"))
671                .unwrap()
672                .lines()
673                .count(),
674            2
675        );
676        let _ = std::fs::remove_dir_all(&dir);
677    }
678
679    // ── synthetic gray frames (mirrors the core detector's fixtures) ──
680    const W: usize = 48;
681    const H: usize = 24;
682
683    fn cfg() -> ParkTuning {
684        ParkTuning {
685            fps: 3.0,
686            left_frac: 0.33,
687            ema_seconds: 6.0,
688            abs_floor: 1500.0,
689            mad_k: 6.0,
690            merge_gap_s: 1.2,
691            max_island_s: 3.0,
692            min_sep_s: 3.0,
693            candidate_frac: 0.75,
694            warmup_s: 0.5,
695            baseline_s: 20.0,
696        }
697    }
698
699    fn cframe(head_x: usize) -> Vec<u8> {
700        let (bg, fix, obj, head) = (200u8, 40u8, 110u8, 25u8);
701        let mut img = vec![bg; W * H];
702        for y in 0..H {
703            let row = y * W;
704            img[row] = fix;
705            img[row + 1] = fix;
706            for x in (W / 2 - 3)..(W / 2 + 3) {
707                img[row + x] = obj;
708            }
709            for x in 0..W {
710                if (x as i64 - head_x as i64).unsigned_abs() <= 4 {
711                    img[row + x] = head;
712                }
713            }
714        }
715        img
716    }
717
718    /// Like [`cframe`] but the head is drawn only on the top half of rows → a fainter
719    /// island (less left-mass), e.g. a blurred travel move that a real park supersedes.
720    fn cframe_weak(head_x: usize) -> Vec<u8> {
721        let (bg, fix, obj, head) = (200u8, 40u8, 110u8, 25u8);
722        let mut img = vec![bg; W * H];
723        for y in 0..H {
724            let row = y * W;
725            img[row] = fix;
726            img[row + 1] = fix;
727            for x in (W / 2 - 3)..(W / 2 + 3) {
728                img[row + x] = obj;
729            }
730            if y >= H / 2 {
731                continue;
732            }
733            for x in 0..W {
734                if (x as i64 - head_x as i64).unsigned_abs() <= 4 {
735                    img[row + x] = head;
736                }
737            }
738        }
739        img
740    }
741
742    #[test]
743    fn detect_stream_writes_one_park_and_maps_the_ring_index() {
744        let dir = tmp("detect");
745        let ring_dir = dir.join(".ring");
746        std::fs::create_dir_all(&ring_dir).unwrap();
747        for idx in 8..=10 {
748            std::fs::write(ring_jpeg_path(&ring_dir, idx), format!("RING{idx}")).unwrap();
749        }
750
751        let mut bytes = Vec::new();
752        for _ in 0..8 {
753            bytes.extend_from_slice(&cframe(W / 2));
754        }
755        for _ in 0..3 {
756            bytes.extend_from_slice(&cframe(6));
757        }
758        for _ in 0..10 {
759            bytes.extend_from_slice(&cframe(W / 2));
760        }
761
762        let mut det = LiveParkDetector::new(W, H, &cfg());
763        let mut writer = ParkWriter::new(dir.clone());
764        let rd = ring_dir.clone();
765        let mut emitted = 0u32;
766        let stats = detect_stream(
767            &mut Cursor::new(bytes),
768            &mut det,
769            W * H,
770            &mut writer,
771            &|idx| ring_jpeg_path(&rd, idx),
772            &|p| p.exists(),
773            &|| false,
774            &mut |ev| {
775                if ev == ParkEvent::Written {
776                    emitted += 1;
777                }
778            },
779        );
780        assert_eq!(stats.parks, 1, "{stats:?}");
781        assert_eq!(stats.dropped, 0);
782        assert_eq!(emitted, 1, "on_emit fired for the written park");
783        assert_eq!(stats.frames, 21, "all frames read");
784        assert!(dir.join("latest_park.jpg").exists());
785        assert!(dir.join("park_000000.jpg").exists());
786        let _ = std::fs::remove_dir_all(&dir);
787    }
788
789    #[test]
790    fn detect_stream_counts_a_replacement_as_replaced_not_a_new_park() {
791        // a weak island then a stronger one <min_sep later (the core replace scenario):
792        // the second SUPERSEDES the first, so it must count as Replaced, not a 2nd park.
793        let dir = tmp("replace-stream");
794        let ring_dir = dir.join(".ring");
795        std::fs::create_dir_all(&ring_dir).unwrap();
796        for idx in 0..30 {
797            std::fs::write(ring_jpeg_path(&ring_dir, idx), b"R").unwrap();
798        }
799        let mut bytes = Vec::new();
800        for _ in 0..8 {
801            bytes.extend_from_slice(&cframe(W / 2));
802        }
803        for _ in 0..2 {
804            bytes.extend_from_slice(&cframe_weak(6));
805        }
806        for _ in 0..5 {
807            bytes.extend_from_slice(&cframe(W / 2));
808        }
809        for _ in 0..2 {
810            bytes.extend_from_slice(&cframe(6));
811        }
812        for _ in 0..8 {
813            bytes.extend_from_slice(&cframe(W / 2));
814        }
815
816        let mut det = LiveParkDetector::new(W, H, &cfg());
817        let mut writer = ParkWriter::new(dir.clone());
818        let rd = ring_dir.clone();
819        let (mut written, mut replaced) = (0u32, 0u32);
820        let stats = detect_stream(
821            &mut Cursor::new(bytes),
822            &mut det,
823            W * H,
824            &mut writer,
825            &|idx| ring_jpeg_path(&rd, idx),
826            &|p| p.exists(),
827            &|| false,
828            &mut |ev| match ev {
829                ParkEvent::Written => written += 1,
830                ParkEvent::Replaced => replaced += 1,
831                ParkEvent::Dropped => {}
832            },
833        );
834        assert_eq!(stats.parks, 1, "one distinct park on disk: {stats:?}");
835        assert_eq!(
836            stats.replaced, 1,
837            "the stronger pair superseded it: {stats:?}"
838        );
839        assert_eq!(written, 1, "live count: one new park, not two");
840        assert_eq!(replaced, 1);
841        assert_eq!(writer.emitted(), 1, "one park_*.jpg on disk");
842        let _ = std::fs::remove_dir_all(&dir);
843    }
844
845    #[test]
846    fn detect_stream_stops_promptly_on_cancel() {
847        let dir = tmp("cancel");
848        let mut det = LiveParkDetector::new(W, H, &cfg());
849        let mut writer = ParkWriter::new(dir.clone());
850        let bytes = cframe(W / 2).repeat(5);
851        let stats = detect_stream(
852            &mut Cursor::new(bytes),
853            &mut det,
854            W * H,
855            &mut writer,
856            &|idx| PathBuf::from(format!("/nope/{idx}")),
857            &|_| true,
858            &|| true, // already cancelled
859            &mut |_| {},
860        );
861        assert_eq!(
862            stats,
863            ParkRunStats::default(),
864            "cancel before reading any frame"
865        );
866        let _ = std::fs::remove_dir_all(&dir);
867    }
868
869    #[test]
870    fn prune_ring_keeps_only_the_newest_jpegs() {
871        let dir = tmp("prune");
872        for idx in 0..10u64 {
873            std::fs::write(ring_jpeg_path(&dir, idx), b"x").unwrap();
874        }
875        std::fs::write(dir.join("notes.txt"), b"x").unwrap(); // a non-ring file
876        assert_eq!(prune_ring(&dir, 3), 7);
877        for idx in 0..7u64 {
878            assert!(!ring_jpeg_path(&dir, idx).exists(), "old {idx} pruned");
879        }
880        for idx in 7..10u64 {
881            assert!(ring_jpeg_path(&dir, idx).exists(), "newest {idx} kept");
882        }
883        assert!(dir.join("notes.txt").exists(), "non-ring file untouched");
884        assert_eq!(
885            prune_ring(&dir, 10),
886            0,
887            "nothing to prune when under the cap"
888        );
889        let _ = std::fs::remove_dir_all(&dir);
890    }
891
892    // ── dense-stream segmented capture: software device model ──
893    // Mirrors core::park's segment_tests model (bright bed + static center object + a dark
894    // head parked FAR-LEFT for a brief, per-layer-VARIABLE dwell). The I/O wiring under test
895    // here is detect_stream_segmented: it must compute each frame's time from the fps
896    // cadence, segment by the injected per-frame layer, and write the PICKED stream index's
897    // ring JPEG out — no ffmpeg, camera, or printer needed.
898    const SEG_CENTER: i64 = 24;
899
900    fn sframe(head_x: i64) -> Vec<u8> {
901        let (bg, obj, head) = (200u8, 110u8, 25u8);
902        let head_hw: i64 = 5;
903        let mut img = vec![bg; W * H];
904        for y in 0..H {
905            let row = y * W;
906            for px in
907                img[row + (SEG_CENTER as usize - 3)..row + (SEG_CENTER as usize + 3)].iter_mut()
908            {
909                *px = obj; // static dark print object at center
910            }
911            for x in 0..W {
912                if (x as i64 - head_x).abs() <= head_hw {
913                    img[row + x] = head;
914                }
915            }
916        }
917        img
918    }
919
920    fn seg_tuning() -> SelectTuning {
921        SelectTuning {
922            left_frac: 0.33,
923            min_outlier: 2.5,
924            min_left_density: 3.0,
925            select_candidate_frac: 0.6,
926            min_confidence: 0.40,
927        }
928    }
929
930    /// Build a modeled continuous gray stream over several layers (head over-print except a
931    /// brief far-left park at a per-layer delay), returning the bytes and the per-frame layer
932    /// signal that `detect_stream_segmented`'s `layer_of` closure replays.
933    fn sim_stream(
934        layers: &[(i64, u64)], // (layer, park_at_ms); each parks 500ms over a 4000ms layer
935        fps: f64,
936    ) -> (Vec<u8>, Vec<i64>) {
937        let dt = (1000.0 / fps) as u64;
938        let mut bytes = Vec::new();
939        let mut layer_of_frame = Vec::new();
940        for &(layer, park_at) in layers {
941            let mut t = 0u64;
942            while t < 4000 {
943                let parked = t >= park_at && t < park_at + 500;
944                bytes.extend_from_slice(&sframe(if parked { 6 } else { SEG_CENTER }));
945                layer_of_frame.push(layer);
946                t += dt;
947            }
948        }
949        (bytes, layer_of_frame)
950    }
951
952    #[test]
953    fn segmented_stream_writes_one_park_per_layer_mapping_the_ring_index() {
954        let dir = tmp("segmented");
955        let ring_dir = dir.join(".ring");
956        std::fs::create_dir_all(&ring_dir).unwrap();
957
958        // 3 layers, 10 fps, parks at WILDLY varying delays (what defeats a sparse grid).
959        let fps = 10.0;
960        let (bytes, layer_of_frame) = sim_stream(&[(0, 500), (1, 2300), (2, 900)], fps);
961        for idx in 0..layer_of_frame.len() as u64 {
962            std::fs::write(ring_jpeg_path(&ring_dir, idx), format!("RING{idx}")).unwrap();
963        }
964
965        let mut sel = SegmentSelector::new(W, H, 3000, seg_tuning());
966        let mut writer = ParkWriter::new(dir.clone());
967        let rd = ring_dir.clone();
968        let lf = layer_of_frame.clone();
969        let mut written = 0u32;
970        let stats = detect_stream_segmented(
971            &mut Cursor::new(bytes),
972            &mut sel,
973            W * H,
974            fps,
975            &|idx| lf[idx as usize],
976            &mut writer,
977            &|idx| ring_jpeg_path(&rd, idx),
978            &|p| p.exists(),
979            &|| false,
980            &mut |ev| {
981                if ev == ParkEvent::Written {
982                    written += 1;
983                }
984            },
985        );
986        assert_eq!(stats.parks, 3, "one park caught per layer: {stats:?}");
987        assert_eq!(stats.dropped, 0, "{stats:?}");
988        assert_eq!(written, 3, "on_emit fired for each written park");
989        assert_eq!(writer.emitted(), 3, "three park_*.jpg on disk");
990
991        // parks.jsonl records the LAYER in `t`, in layer order, with mapped ring indices.
992        let jl = std::fs::read_to_string(dir.join("parks.jsonl")).unwrap();
993        let layers: Vec<f64> = jl
994            .lines()
995            .map(|l| {
996                serde_json::from_str::<serde_json::Value>(l).unwrap()["t"]
997                    .as_f64()
998                    .unwrap()
999            })
1000            .collect();
1001        assert_eq!(layers, vec![0.0, 1.0, 2.0], "{jl}");
1002        let _ = std::fs::remove_dir_all(&dir);
1003    }
1004
1005    #[test]
1006    fn segmented_stream_skips_an_over_print_only_layer() {
1007        // Layer 0 never parks (head stays over-print) → no pick; layer 1 parks → one pick.
1008        let dir = tmp("segmented-skip");
1009        let ring_dir = dir.join(".ring");
1010        std::fs::create_dir_all(&ring_dir).unwrap();
1011        let fps = 10.0;
1012        let (bytes, lf) = sim_stream(&[(0, 99_999), (1, 800)], fps);
1013        for idx in 0..lf.len() as u64 {
1014            std::fs::write(ring_jpeg_path(&ring_dir, idx), b"R").unwrap();
1015        }
1016        let mut sel = SegmentSelector::new(W, H, 3000, seg_tuning());
1017        let mut writer = ParkWriter::new(dir.clone());
1018        let rd = ring_dir.clone();
1019        let lfc = lf.clone();
1020        let stats = detect_stream_segmented(
1021            &mut Cursor::new(bytes),
1022            &mut sel,
1023            W * H,
1024            fps,
1025            &|idx| lfc[idx as usize],
1026            &mut writer,
1027            &|idx| ring_jpeg_path(&rd, idx),
1028            &|p| p.exists(),
1029            &|| false,
1030            &mut |_| {},
1031        );
1032        assert_eq!(
1033            stats.parks, 1,
1034            "only the parked layer yields a frame: {stats:?}"
1035        );
1036        assert_eq!(writer.emitted(), 1);
1037        let _ = std::fs::remove_dir_all(&dir);
1038    }
1039
1040    #[test]
1041    fn segmented_stream_drops_a_pick_whose_ring_jpeg_never_arrived() {
1042        // The park is selected but its full-res JPEG never landed → counted as dropped, not
1043        // silently lost (a missing layer stays visible).
1044        let dir = tmp("segmented-drop");
1045        let ring_dir = dir.join(".ring");
1046        std::fs::create_dir_all(&ring_dir).unwrap(); // empty: no ring JPEGs at all
1047        let fps = 10.0;
1048        let (bytes, lf) = sim_stream(&[(0, 500)], fps);
1049        let mut sel = SegmentSelector::new(W, H, 3000, seg_tuning());
1050        let mut writer = ParkWriter::new(dir.clone());
1051        let rd = ring_dir.clone();
1052        let lfc = lf.clone();
1053        let stats = detect_stream_segmented(
1054            &mut Cursor::new(bytes),
1055            &mut sel,
1056            W * H,
1057            fps,
1058            &|idx| lfc[idx as usize],
1059            &mut writer,
1060            &|idx| ring_jpeg_path(&rd, idx),
1061            &|p| p.exists(),
1062            &|| false,
1063            &mut |_| {},
1064        );
1065        assert_eq!(stats.parks, 0, "{stats:?}");
1066        assert_eq!(
1067            stats.dropped, 1,
1068            "the picked frame's JPEG was missing: {stats:?}"
1069        );
1070        let _ = std::fs::remove_dir_all(&dir);
1071    }
1072}