1use 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
28pub const DECODE_W: usize = 64;
31pub const DECODE_H: usize = 36;
32
33pub 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
69pub fn ring_jpeg_path(ring_dir: &Path, idx: u64) -> PathBuf {
71 ring_dir.join(format!("full_{idx:09}.jpg"))
72}
73
74fn 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
88pub 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 pub fn emitted(&self) -> u64 {
107 self.emitted
108 }
109
110 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
143pub struct ParkWritten {
146 pub index: u64,
147 pub replaced: bool,
148}
149
150#[derive(Debug, Clone, Copy, PartialEq, Eq)]
154pub enum ParkEvent {
155 Written,
157 Replaced,
159 Dropped,
161}
162
163#[derive(Debug, Default, PartialEq, Eq)]
165pub struct ParkRunStats {
166 pub frames: u64,
168 pub parks: u64,
170 pub replaced: u64,
172 pub dropped: u64,
175}
176
177fn 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#[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
243fn 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#[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#[derive(Clone)]
313pub struct ParkCapture {
314 pub id: String,
315 pub stream_url: String,
316 pub tuning: ParkTuning,
317}
318
319pub 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
352const RING_KEEP_SECONDS: f64 = 60.0;
355
356struct 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 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
377fn 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
436pub 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#[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
508pub 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
557fn 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(); let wr = w.write(&park(5, true), &r2).unwrap(); 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 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 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 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, &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(); 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 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; }
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 fn sim_stream(
934 layers: &[(i64, u64)], 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 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 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 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 let dir = tmp("segmented-drop");
1045 let ring_dir = dir.join(".ring");
1046 std::fs::create_dir_all(&ring_dir).unwrap(); 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}