use std::io::Write;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, AtomicI64, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
use tokio::sync::watch;
use tokio::task::JoinHandle;
use super::camera::StreamOpen;
use super::stream_record::record_loop;
use crate::core::park::{Park, SelectTuning};
use crate::core::status::PrinterStatus;
use crate::core::timelapse::{ActivityAction, CaptureAction, CaptureSession, PrintActivitySession};
use crate::park::{
DECODE_H, DECODE_W, ParkCapture, ParkEvent, ParkRunStats, ParkWriter, SegmentCapture,
run_park_camera, run_segment_camera,
};
pub type ParkSpawn = Arc<
dyn Fn(ParkCapture, PathBuf, Arc<AtomicBool>, Arc<Mutex<TimelapseStatus>>) -> JoinHandle<()>
+ Send
+ Sync,
>;
pub fn real_park_spawn() -> ParkSpawn {
Arc::new(
|cap: ParkCapture, out_dir, cancel, status: Arc<Mutex<TimelapseStatus>>| {
tokio::task::spawn_blocking(move || {
let mut on_park = park_progress(cap.id.clone(), status.clone());
let cam_dir = out_dir.join(&cap.id);
let outcome =
run_park_camera(&cap, &cam_dir, DECODE_W, DECODE_H, &cancel, &mut on_park);
report_park_run(&cap.id, &cap.stream_url, outcome, &status);
})
},
)
}
pub type SegmentSpawn = Arc<
dyn Fn(
SegmentCapture,
PathBuf,
Arc<AtomicI64>,
Arc<AtomicBool>,
Arc<Mutex<TimelapseStatus>>,
) -> JoinHandle<()>
+ Send
+ Sync,
>;
pub fn real_segment_spawn() -> SegmentSpawn {
Arc::new(
|cap: SegmentCapture,
out_dir,
current_layer: Arc<AtomicI64>,
cancel,
status: Arc<Mutex<TimelapseStatus>>| {
tokio::task::spawn_blocking(move || {
let mut on_park = park_progress(cap.id.clone(), status.clone());
let cam_dir = out_dir.join(&cap.id);
let outcome = run_segment_camera(
&cap,
&cam_dir,
DECODE_W,
DECODE_H,
¤t_layer,
&cancel,
&mut on_park,
);
report_park_run(&cap.id, &cap.stream_url, outcome, &status);
})
},
)
}
fn park_progress(id: String, status: Arc<Mutex<TimelapseStatus>>) -> impl FnMut(ParkEvent) {
move |ev: ParkEvent| {
let mut s = status.lock().unwrap();
match ev {
ParkEvent::Written => s.frames += 1,
ParkEvent::Replaced => {}
ParkEvent::Dropped => {
s.failures += 1;
s.last_error = Some(format!("park {id}: a ring JPEG never arrived"));
}
}
}
}
fn report_park_run(
id: &str,
source: &str,
outcome: Result<ParkRunStats, String>,
status: &Arc<Mutex<TimelapseStatus>>,
) {
let mut s = status.lock().unwrap();
match outcome {
Ok(stats) if stats.frames == 0 => {
s.failures += 1;
s.last_error = Some(format!("park {id}: read 0 frames from {source}"));
}
Ok(_) => {}
Err(e) => {
s.failures += 1;
s.last_error = Some(e);
}
}
}
pub type FrameGrab = Arc<dyn Fn() -> Result<Vec<u8>, String> + Send + Sync>;
const MAX_STREAM_BYTES: u64 = 2 * 1024 * 1024 * 1024;
const MAX_STREAM_INPUT_BYTES: u64 = 64 * 1024 * 1024 * 1024;
pub enum PlainCapture {
Sample { id: String, grab: FrameGrab },
Stream { id: String, open: StreamOpen },
}
impl PlainCapture {
fn id(&self) -> &str {
match self {
PlainCapture::Sample { id, .. } | PlainCapture::Stream { id, .. } => id,
}
}
}
#[derive(Clone, Default)]
pub struct TimelapseStatus {
pub running: bool,
pub mode: &'static str,
pub cameras: Vec<String>,
pub every: u64,
pub interval_ms: Option<u64>,
pub frames: u64,
pub failures: u64,
pub current_layer: Option<i64>,
pub out_dir: Option<String>,
pub last_error: Option<String>,
}
impl TimelapseStatus {
pub fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"running": self.running,
"mode": self.mode,
"cameras": self.cameras,
"camera": self.cameras.first(),
"every": self.every,
"interval_ms": self.interval_ms,
"frames": self.frames,
"failures": self.failures,
"current_layer": self.current_layer,
"out_dir": self.out_dir,
"last_error": self.last_error,
})
}
}
#[derive(Default)]
struct Inner {
status: Arc<Mutex<TimelapseStatus>>,
handle: Option<JoinHandle<()>>,
cancel: Arc<AtomicBool>,
}
#[derive(Default)]
pub struct TimelapseManager {
smooth: Mutex<Inner>,
plain: Mutex<Inner>,
park: Mutex<Inner>,
segment: Mutex<Inner>,
}
impl TimelapseManager {
pub fn start_smooth(
&self,
cameras: Vec<(String, FrameGrab)>,
every: u64,
burst_offsets: Vec<u64>,
rx: watch::Receiver<PrinterStatus>,
out_dir: PathBuf,
) -> Result<(), String> {
self.start_smooth_with_select(cameras, every, burst_offsets, rx, out_dir, Vec::new())
}
pub fn start_smooth_with_select(
&self,
cameras: Vec<(String, FrameGrab)>,
every: u64,
burst_offsets: Vec<u64>,
rx: watch::Receiver<PrinterStatus>,
out_dir: PathBuf,
selects: Vec<Option<SelectTuning>>,
) -> Result<(), String> {
let every = every.max(1);
let burst_offsets = normalize_burst_offsets(burst_offsets);
let ids = cameras.iter().map(|(id, _)| id.clone()).collect();
start_slot(
&self.smooth,
cameras,
ids,
out_dir,
TimelapseStatus {
mode: "smooth",
every,
..Default::default()
},
move |status, cams, dir, cancel| {
tokio::spawn(run(
status,
rx,
cams,
dir,
every,
burst_offsets,
selects,
cancel,
))
},
)
}
pub fn start_plain(
&self,
cameras: Vec<PlainCapture>,
interval_ms: u64,
rx: watch::Receiver<PrinterStatus>,
out_dir: PathBuf,
) -> Result<(), String> {
let interval_ms = interval_ms.max(1);
let ids = cameras.iter().map(|c| c.id().to_string()).collect();
start_slot(
&self.plain,
cameras,
ids,
out_dir,
TimelapseStatus {
mode: "plain",
interval_ms: Some(interval_ms),
..Default::default()
},
move |status, caps, dir, cancel| {
tokio::spawn(run_plain(status, rx, caps, dir, interval_ms, cancel))
},
)
}
pub fn start_park(
&self,
cameras: Vec<ParkCapture>,
rx: watch::Receiver<PrinterStatus>,
out_dir: PathBuf,
spawn_worker: ParkSpawn,
) -> Result<(), String> {
let ids = cameras.iter().map(|c| c.id.clone()).collect();
start_slot(
&self.park,
cameras,
ids,
out_dir,
TimelapseStatus {
mode: "park",
..Default::default()
},
move |status, caps, dir, cancel| {
tokio::spawn(run_park(status, rx, caps, dir, cancel, spawn_worker))
},
)
}
pub fn start_segment(
&self,
cameras: Vec<SegmentCapture>,
rx: watch::Receiver<PrinterStatus>,
out_dir: PathBuf,
spawn_worker: SegmentSpawn,
) -> Result<(), String> {
let ids = cameras.iter().map(|c| c.id.clone()).collect();
start_slot(
&self.segment,
cameras,
ids,
out_dir,
TimelapseStatus {
mode: "segment",
..Default::default()
},
move |status, caps, dir, cancel| {
tokio::spawn(run_segment(status, rx, caps, dir, cancel, spawn_worker))
},
)
}
pub fn stop_smooth(&self) -> bool {
stop_slot(&self.smooth)
}
pub fn stop_plain(&self) -> bool {
stop_slot(&self.plain)
}
pub fn stop_park(&self) -> bool {
stop_slot(&self.park)
}
pub fn stop_segment(&self) -> bool {
stop_slot(&self.segment)
}
pub fn status_smooth(&self) -> TimelapseStatus {
self.smooth.lock().unwrap().status.lock().unwrap().clone()
}
pub fn status_plain(&self) -> TimelapseStatus {
self.plain.lock().unwrap().status.lock().unwrap().clone()
}
pub fn status_park(&self) -> TimelapseStatus {
self.park.lock().unwrap().status.lock().unwrap().clone()
}
pub fn status_segment(&self) -> TimelapseStatus {
self.segment.lock().unwrap().status.lock().unwrap().clone()
}
}
fn start_slot<C>(
inner: &Mutex<Inner>,
cameras: Vec<C>,
ids: Vec<String>,
out_dir: PathBuf,
init: TimelapseStatus,
spawn: impl FnOnce(Arc<Mutex<TimelapseStatus>>, Vec<C>, PathBuf, Arc<AtomicBool>) -> JoinHandle<()>,
) -> Result<(), String> {
let mut g = inner.lock().unwrap();
if g.status.lock().unwrap().running {
return Err(format!("a {} timelapse is already running", init.mode));
}
if ids.is_empty() {
return Err("no cameras to capture".to_string());
}
for id in &ids {
let dir = out_dir.join(id);
std::fs::create_dir_all(&dir).map_err(|e| format!("create {}: {e}", dir.display()))?;
}
let cancel = Arc::new(AtomicBool::new(false));
let status = Arc::new(Mutex::new(TimelapseStatus {
running: true,
cameras: ids,
out_dir: Some(out_dir.display().to_string()),
..init
}));
let handle = spawn(status.clone(), cameras, out_dir, cancel.clone());
g.status = status;
g.handle = Some(handle);
g.cancel = cancel;
Ok(())
}
fn stop_slot(inner: &Mutex<Inner>) -> bool {
let mut g = inner.lock().unwrap();
g.cancel.store(true, Ordering::Relaxed);
if let Some(h) = g.handle.take() {
h.abort();
}
let mut s = g.status.lock().unwrap();
let was = s.running;
s.running = false;
was
}
pub const DEFAULT_SMOOTH_BURST_MS: &[u64] = &[
100, 300, 500, 700, 900, 1100, 1300, 1500, 1700, 1900, 2100, 2300, 2500, 2700, 2900,
];
fn burst_frame_name(frame_no: u64, layer: i64, offset_ms: u64) -> String {
format!("frame_{frame_no:06}_layer_{layer:05}_t{offset_ms:04}.jpg")
}
fn normalize_burst_offsets(mut offsets: Vec<u64>) -> Vec<u64> {
offsets.sort_unstable();
offsets.dedup();
if offsets.is_empty() { vec![0] } else { offsets }
}
#[derive(Clone)]
struct BurstSink {
tx: tokio::sync::mpsc::Sender<(FrameGrab, PathBuf)>,
cameras: Arc<Vec<(String, FrameGrab)>>,
status: Arc<Mutex<TimelapseStatus>>,
out_dir: PathBuf,
cancel: Arc<AtomicBool>,
}
struct LiveSelect {
cam_dir: PathBuf,
tuning: SelectTuning,
writer: Arc<Mutex<ParkWriter>>,
}
const FINALIZE_MARGIN_MS: u64 = 800;
fn schedule_finalize(
live: &Arc<Vec<LiveSelect>>,
frame_no: u64,
layer: i64,
max_offset: u64,
cancel: &Arc<AtomicBool>,
) {
let live = live.clone();
let cancel = cancel.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(max_offset + FINALIZE_MARGIN_MS)).await;
if cancel.load(Ordering::Relaxed) {
return; }
for sel in live.iter() {
let cam_dir = sel.cam_dir.clone();
let tuning = sel.tuning;
let writer = sel.writer.clone();
let _ = tokio::task::spawn_blocking(move || {
if let Ok(Some((path, confidence))) =
crate::captures::select_layer_burst(&cam_dir, layer, &tuning)
{
let park = Park {
idx: frame_no,
t: layer as f64,
left_mass: 0.0,
sharpness: 0.0,
confidence,
replace: false,
};
let _ = writer.lock().unwrap().write(&park, &path);
}
})
.await;
}
});
}
fn schedule_burst(sink: &BurstSink, frame_no: u64, layer: i64, offsets: &[u64]) {
for &offset_ms in offsets {
let sink = sink.clone();
tokio::spawn(async move {
tokio::time::sleep(Duration::from_millis(offset_ms)).await;
if sink.cancel.load(Ordering::Relaxed) {
return; }
let name = burst_frame_name(frame_no, layer, offset_ms);
for (id, grab) in sink.cameras.iter() {
let path = sink.out_dir.join(id).join(&name);
if sink.tx.try_send((grab.clone(), path)).is_err() {
let mut s = sink.status.lock().unwrap();
s.failures += 1;
s.last_error = Some("capture fell behind — frame skipped".to_string());
}
}
});
}
}
#[allow(clippy::too_many_arguments)]
async fn run(
status: Arc<Mutex<TimelapseStatus>>,
mut rx: watch::Receiver<PrinterStatus>,
cameras: Vec<(String, FrameGrab)>,
out_dir: PathBuf,
every: u64,
burst_offsets: Vec<u64>,
selects: Vec<Option<SelectTuning>>,
cancel: Arc<AtomicBool>,
) {
let mut session = CaptureSession::new(every, true);
let live: Vec<LiveSelect> = cameras
.iter()
.enumerate()
.filter_map(|(i, (id, _))| {
selects.get(i).copied().flatten().map(|tuning| LiveSelect {
cam_dir: out_dir.join(id),
tuning,
writer: Arc::new(Mutex::new(ParkWriter::new(out_dir.join(id)))),
})
})
.collect();
let live = Arc::new(live);
let max_offset = burst_offsets.iter().copied().max().unwrap_or(0);
let cameras = Arc::new(cameras);
let bound = (4 * cameras.len() * burst_offsets.len().max(1)).max(8);
let (tx, mut jobs) = tokio::sync::mpsc::channel::<(FrameGrab, PathBuf)>(bound);
let wstatus = status.clone();
let worker = tokio::spawn(async move {
while let Some((grab, path)) = jobs.recv().await {
let res = tokio::task::spawn_blocking(move || grab()).await;
let mut s = wstatus.lock().unwrap();
match res {
Ok(Ok(bytes)) => match std::fs::write(&path, &bytes) {
Ok(()) => s.frames += 1,
Err(e) => {
s.failures += 1;
s.last_error = Some(format!("write {}: {e}", path.display()));
}
},
Ok(Err(e)) => {
s.failures += 1;
s.last_error = Some(e);
}
Err(_) => {
s.failures += 1;
s.last_error = Some("frame grab task failed".to_string());
}
}
}
});
let sink = BurstSink {
tx,
cameras,
status: status.clone(),
out_dir,
cancel,
};
loop {
let snap = rx.borrow_and_update().clone();
status.lock().unwrap().current_layer = snap.layer_num;
match session.observe(&snap) {
CaptureAction::Capture { frame_no, layer } => {
schedule_burst(&sink, frame_no, layer, &burst_offsets);
if !live.is_empty() {
schedule_finalize(&live, frame_no, layer, max_offset, &sink.cancel);
}
}
CaptureAction::Stop => break,
CaptureAction::Continue => {}
}
if rx.changed().await.is_err() {
break; }
}
drop(sink);
let _ = worker.await;
status.lock().unwrap().running = false;
}
fn live_mp4_args(out: &std::path::Path) -> Vec<String> {
vec![
"-y".into(),
"-f".into(),
"mpjpeg".into(),
"-i".into(),
"-".into(),
"-c:v".into(),
"libx264".into(),
"-pix_fmt".into(),
"yuv420p".into(),
"-movflags".into(),
"+faststart".into(),
out.display().to_string(),
]
}
fn spawn_stream_recorders(
streams: Vec<(String, StreamOpen)>,
out_dir: &std::path::Path,
status: &Arc<Mutex<TimelapseStatus>>,
cancel: &Arc<AtomicBool>,
) -> Vec<JoinHandle<()>> {
streams
.into_iter()
.map(|(id, open)| {
let dir = out_dir.join(&id);
let mp4 = dir.join("plain.mp4");
let mjpeg = dir.join("plain.mjpeg");
let cancel = cancel.clone();
let wstatus = status.clone();
tokio::task::spawn_blocking(move || {
let cancel_fn = || cancel.load(Ordering::Relaxed);
let backoff = |attempt: u32| {
let total_ms = (500u64 * u64::from(attempt)).min(5_000);
let mut slept = 0u64;
while slept < total_ms && !cancel.load(Ordering::Relaxed) {
std::thread::sleep(Duration::from_millis(50));
slept += 50;
}
};
let ffmpeg = std::process::Command::new("ffmpeg")
.args(live_mp4_args(&mp4))
.stdin(std::process::Stdio::piped())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn();
let (stats, target, encode_ok) = match ffmpeg {
Ok(mut child) => {
let mut stdin = child.stdin.take().expect("piped stdin");
let stats = record_loop(
&open,
&mut stdin,
&cancel_fn,
MAX_STREAM_INPUT_BYTES,
32 * 1024,
&backoff,
);
drop(stdin); let ok = child.wait().map(|s| s.success()).unwrap_or(false);
(stats, mp4, ok)
}
Err(_) => {
let file = match std::fs::File::create(&mjpeg) {
Ok(f) => f,
Err(e) => {
let mut s = wstatus.lock().unwrap();
s.failures += 1;
s.last_error = Some(format!("create {}: {e}", mjpeg.display()));
return;
}
};
let mut sink = std::io::BufWriter::new(file);
let stats = record_loop(
&open,
&mut sink,
&cancel_fn,
MAX_STREAM_BYTES,
32 * 1024,
&backoff,
);
let _ = sink.flush();
(stats, mjpeg, true)
}
};
let mut s = wstatus.lock().unwrap();
s.failures += u64::from(stats.failures);
if stats.bytes == 0 {
s.last_error = Some(format!(
"stream {id}: no data recorded ({})",
target.display()
));
} else if !encode_ok {
s.failures += 1;
s.last_error = Some(format!(
"stream {id}: ffmpeg failed to encode {} (missing libx264, or bad stream)",
target.display()
));
}
})
})
.collect()
}
async fn run_plain(
status: Arc<Mutex<TimelapseStatus>>,
mut rx: watch::Receiver<PrinterStatus>,
cameras: Vec<PlainCapture>,
out_dir: PathBuf,
interval_ms: u64,
cancel: Arc<AtomicBool>,
) {
let mut samples: Vec<(String, FrameGrab)> = Vec::new();
let mut streams: Vec<(String, StreamOpen)> = Vec::new();
for cap in cameras {
match cap {
PlainCapture::Sample { id, grab } => samples.push((id, grab)),
PlainCapture::Stream { id, open } => streams.push((id, open)),
}
}
let mut streams = streams;
let mut stream_workers: Vec<JoinHandle<()>> = Vec::new();
let mut activity = PrintActivitySession::new(true);
let bound = (4 * samples.len()).max(4);
let (tx, mut jobs) = tokio::sync::mpsc::channel::<(FrameGrab, PathBuf)>(bound);
let wstatus = status.clone();
let worker = tokio::spawn(async move {
while let Some((grab, path)) = jobs.recv().await {
let res = tokio::task::spawn_blocking(move || grab()).await;
let mut s = wstatus.lock().unwrap();
match res {
Ok(Ok(bytes)) => match std::fs::write(&path, &bytes) {
Ok(()) => s.frames += 1,
Err(e) => {
s.failures += 1;
s.last_error = Some(format!("write {}: {e}", path.display()));
}
},
Ok(Err(e)) => {
s.failures += 1;
s.last_error = Some(e);
}
Err(_) => {
s.failures += 1;
s.last_error = Some("frame grab task failed".to_string());
}
}
}
});
let mut frame_no: u64 = 0;
let mut ticker = tokio::time::interval(std::time::Duration::from_millis(interval_ms));
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
tokio::select! {
_ = ticker.tick() => {
let snap = rx.borrow().clone();
status.lock().unwrap().current_layer = snap.layer_num;
match activity.observe(&snap) {
ActivityAction::Capture => {
if !streams.is_empty() {
stream_workers = spawn_stream_recorders(
std::mem::take(&mut streams),
&out_dir,
&status,
&cancel,
);
}
frame_no += 1;
let name = format!("frame_{frame_no:06}.jpg");
for (id, grab) in &samples {
let path = out_dir.join(id).join(&name);
if tx.try_send((grab.clone(), path)).is_err() {
let mut s = status.lock().unwrap();
s.failures += 1;
s.last_error = Some("capture fell behind — frame skipped".to_string());
}
}
}
ActivityAction::Idle => {}
ActivityAction::Stop => break,
}
}
changed = rx.changed() => {
if changed.is_err() {
break;
}
let snap = rx.borrow().clone();
if activity.observe(&snap) == ActivityAction::Stop {
break;
}
}
}
}
drop(tx);
let _ = worker.await;
cancel.store(true, Ordering::Relaxed);
for w in stream_workers {
let _ = w.await;
}
status.lock().unwrap().running = false;
}
async fn run_park(
status: Arc<Mutex<TimelapseStatus>>,
mut rx: watch::Receiver<PrinterStatus>,
captures: Vec<ParkCapture>,
out_dir: PathBuf,
cancel: Arc<AtomicBool>,
spawn_worker: ParkSpawn,
) {
let mut activity = PrintActivitySession::new(true);
let mut pending = Some(captures); let mut workers: Vec<JoinHandle<()>> = Vec::new();
loop {
let snap = rx.borrow_and_update().clone();
status.lock().unwrap().current_layer = snap.layer_num;
match activity.observe(&snap) {
ActivityAction::Capture => {
if let Some(caps) = pending.take() {
for cap in caps {
workers.push(spawn_worker(
cap,
out_dir.clone(),
cancel.clone(),
status.clone(),
));
}
}
}
ActivityAction::Idle => {}
ActivityAction::Stop => break,
}
if rx.changed().await.is_err() {
break; }
}
cancel.store(true, Ordering::Relaxed);
for w in workers {
let _ = w.await;
}
status.lock().unwrap().running = false;
}
async fn run_segment(
status: Arc<Mutex<TimelapseStatus>>,
mut rx: watch::Receiver<PrinterStatus>,
captures: Vec<SegmentCapture>,
out_dir: PathBuf,
cancel: Arc<AtomicBool>,
spawn_worker: SegmentSpawn,
) {
let current_layer = Arc::new(AtomicI64::new(-1));
let mut activity = PrintActivitySession::new(true);
let mut pending = Some(captures); let mut workers: Vec<JoinHandle<()>> = Vec::new();
loop {
let snap = rx.borrow_and_update().clone();
if let Some(l) = snap.layer_num {
current_layer.store(l, Ordering::Relaxed);
}
status.lock().unwrap().current_layer = snap.layer_num;
match activity.observe(&snap) {
ActivityAction::Capture => {
if let Some(caps) = pending.take() {
for cap in caps {
workers.push(spawn_worker(
cap,
out_dir.clone(),
current_layer.clone(),
cancel.clone(),
status.clone(),
));
}
}
}
ActivityAction::Idle => {}
ActivityAction::Stop => break,
}
if rx.changed().await.is_err() {
break; }
}
cancel.store(true, Ordering::Relaxed);
for w in workers {
let _ = w.await;
}
status.lock().unwrap().running = false;
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::park::ParkTuning;
use crate::core::status::PrinterStatus;
use std::sync::atomic::AtomicUsize;
fn st(state: &str, layer: Option<i64>) -> PrinterStatus {
PrinterStatus {
gcode_state: Some(state.to_string()),
layer_num: layer,
..Default::default()
}
}
fn one(id: &str, grab: FrameGrab) -> Vec<(String, FrameGrab)> {
vec![(id.to_string(), grab)]
}
fn sample(id: &str, grab: FrameGrab) -> Vec<PlainCapture> {
vec![PlainCapture::Sample {
id: id.to_string(),
grab,
}]
}
#[tokio::test]
async fn runs_a_capture_from_a_watch_feed_writing_one_frame_per_layer() {
let dir = std::env::temp_dir().join(format!("bambu-tl-test-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (tx, rx) = watch::channel(st("IDLE", None));
let grab: FrameGrab = Arc::new(|| Ok(vec![0xff, 0xd8, 0xff, 0x42]));
let mgr = TimelapseManager::default();
mgr.start_smooth(one("ext-0", grab), 1, vec![0], rx, dir.clone())
.unwrap();
for s in [
st("RUNNING", Some(1)),
st("RUNNING", Some(2)),
st("RUNNING", Some(3)),
st("FINISH", Some(3)),
] {
tx.send(s).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(40)).await;
}
tokio::time::sleep(std::time::Duration::from_millis(80)).await;
let s = mgr.status_smooth();
assert!(
!s.running,
"capture should auto-stop when the print finishes"
);
assert_eq!(s.frames, 3, "one frame per advancing layer");
assert_eq!(s.failures, 0);
let n = std::fs::read_dir(dir.join("ext-0")).unwrap().count();
assert_eq!(n, 3, "three JPEG files written under the camera's subdir");
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn captures_every_camera_once_per_layer_into_per_camera_subdirs() {
let dir = std::env::temp_dir().join(format!("bambu-tl-multi-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (tx, rx) = watch::channel(st("IDLE", None));
let g: FrameGrab = Arc::new(|| Ok(vec![0xff, 0xd8, 0xff, 0x01]));
let mgr = TimelapseManager::default();
mgr.start_smooth(
vec![("ext-0".into(), g.clone()), ("ext-1".into(), g)],
1,
vec![0],
rx,
dir.clone(),
)
.unwrap();
for s in [
st("RUNNING", Some(1)),
st("RUNNING", Some(2)),
st("FINISH", Some(2)),
] {
tx.send(s).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(40)).await;
}
tokio::time::sleep(std::time::Duration::from_millis(80)).await;
let s = mgr.status_smooth();
assert_eq!(s.cameras, vec!["ext-0".to_string(), "ext-1".to_string()]);
assert_eq!(s.frames, 4, "2 layers × 2 cameras");
assert_eq!(s.failures, 0);
assert_eq!(std::fs::read_dir(dir.join("ext-0")).unwrap().count(), 2);
assert_eq!(std::fs::read_dir(dir.join("ext-1")).unwrap().count(), 2);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn start_twice_is_rejected_until_stopped() {
let dir = std::env::temp_dir().join(format!("bambu-tl-test2-{}", std::process::id()));
let (_tx, rx) = watch::channel(st("RUNNING", Some(0)));
let grab: FrameGrab = Arc::new(|| Ok(vec![0xff, 0xd8, 0xff]));
let mgr = TimelapseManager::default();
mgr.start_smooth(
one("ext-0", grab.clone()),
1,
vec![0],
rx.clone(),
dir.clone(),
)
.unwrap();
assert!(
mgr.start_smooth(one("ext-1", grab), 1, vec![0], rx, dir.clone())
.is_err()
);
assert!(mgr.stop_smooth());
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn start_with_no_cameras_is_rejected() {
let dir = std::env::temp_dir().join(format!("bambu-tl-empty-{}", std::process::id()));
let (_tx, rx) = watch::channel(st("RUNNING", Some(0)));
let mgr = TimelapseManager::default();
assert!(
mgr.start_smooth(vec![], 1, vec![0], rx, dir).is_err(),
"need at least one camera"
);
}
#[tokio::test]
async fn a_failing_grab_counts_failures_and_keeps_going() {
let dir = std::env::temp_dir().join(format!("bambu-tl-test3-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (tx, rx) = watch::channel(st("RUNNING", Some(0)));
let grab: FrameGrab = Arc::new(|| Err("camera offline".to_string()));
let mgr = TimelapseManager::default();
mgr.start_smooth(one("ext-0", grab), 1, vec![0], rx, dir.clone())
.unwrap();
for s in [
st("RUNNING", Some(1)),
st("RUNNING", Some(2)),
st("FINISH", Some(2)),
] {
tx.send(s).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(40)).await;
}
tokio::time::sleep(std::time::Duration::from_millis(80)).await;
let s = mgr.status_smooth();
assert!(s.failures >= 2, "grab failures are counted");
assert_eq!(s.frames, 0, "no files on failure");
assert!(s.last_error.is_some());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn burst_frame_name_tags_frame_layer_and_offset() {
assert_eq!(
super::burst_frame_name(1, 5, 800),
"frame_000001_layer_00005_t0800.jpg"
);
assert_eq!(
super::burst_frame_name(12, 240, 0),
"frame_000012_layer_00240_t0000.jpg"
);
}
#[test]
fn normalize_burst_offsets_sorts_dedups_and_defaults_empty() {
assert_eq!(
super::normalize_burst_offsets(vec![800, 400, 800, 600]),
vec![400, 600, 800]
);
assert_eq!(super::normalize_burst_offsets(vec![500, 500]), vec![500]);
assert_eq!(super::normalize_burst_offsets(vec![]), vec![0]);
}
#[tokio::test(start_paused = true)]
async fn burst_enqueues_one_grab_per_offset_at_its_due_time() {
let (tx, mut rx) = tokio::sync::mpsc::channel::<(FrameGrab, PathBuf)>(64);
let g: FrameGrab = Arc::new(|| Ok(vec![0xff, 0xd8, 0xff]));
let cameras = Arc::new(vec![("ext-0".to_string(), g)]);
let status = Arc::new(Mutex::new(TimelapseStatus::default()));
let cancel = Arc::new(AtomicBool::new(false));
let sink = super::BurstSink {
tx,
cameras,
status,
out_dir: std::path::PathBuf::from("/cap"),
cancel,
};
super::schedule_burst(&sink, 1, 5, &[10, 30]);
tokio::task::yield_now().await;
assert!(
rx.try_recv().is_err(),
"nothing is due before the first offset"
);
tokio::time::advance(Duration::from_millis(10)).await;
tokio::task::yield_now().await;
let (_g, p) = rx.try_recv().expect("first sample due at 10ms");
assert!(
p.ends_with("frame_000001_layer_00005_t0010.jpg"),
"{}",
p.display()
);
assert!(rx.try_recv().is_err(), "the 30ms sample is not due yet");
tokio::time::advance(Duration::from_millis(20)).await;
tokio::task::yield_now().await;
let (_g, p) = rx.try_recv().expect("second sample due at 30ms");
assert!(
p.ends_with("frame_000001_layer_00005_t0030.jpg"),
"{}",
p.display()
);
}
#[tokio::test(start_paused = true)]
async fn a_cancelled_burst_enqueues_nothing() {
let (tx, mut rx) = tokio::sync::mpsc::channel::<(FrameGrab, PathBuf)>(8);
let g: FrameGrab = Arc::new(|| Ok(vec![0xff, 0xd8, 0xff]));
let cameras = Arc::new(vec![("ext-0".to_string(), g)]);
let status = Arc::new(Mutex::new(TimelapseStatus::default()));
let cancel = Arc::new(AtomicBool::new(true)); let sink = super::BurstSink {
tx,
cameras,
status,
out_dir: std::path::PathBuf::from("/cap"),
cancel,
};
super::schedule_burst(&sink, 1, 5, &[10]);
tokio::task::yield_now().await;
tokio::time::advance(Duration::from_millis(20)).await;
tokio::task::yield_now().await;
assert!(
rx.try_recv().is_err(),
"a burst that fires after stop must not grab"
);
}
#[tokio::test]
async fn plain_samples_frames_on_an_interval_while_printing() {
let dir = std::env::temp_dir().join(format!("bambu-tl-plain-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (tx, rx) = watch::channel(st("IDLE", None));
let grab: FrameGrab = Arc::new(|| Ok(vec![0xff, 0xd8, 0xff, 0x09]));
let mgr = TimelapseManager::default();
mgr.start_plain(sample("ext-0", grab), 20, rx, dir.clone())
.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(60)).await;
assert_eq!(
mgr.status_plain().frames,
0,
"no sampling before the print is active"
);
tx.send(st("RUNNING", Some(1))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(150)).await;
let mid = mgr.status_plain().frames;
assert!(
mid >= 2,
"plain samples on its own clock while printing (got {mid})"
);
tx.send(st("FINISH", Some(1))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(80)).await;
assert!(
!mgr.status_plain().running,
"plain stops when the print finishes"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn plain_stream_recorder_starts_and_stops_with_the_print() {
use crate::server::camera::{OpenedCameraStream, StreamOpen};
let dir = std::env::temp_dir().join(format!("bambu-tl-stream-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (tx, rx) = watch::channel(st("RUNNING", Some(1)));
let open: StreamOpen = Arc::new(|| {
Ok(OpenedCameraStream {
content_type: "multipart/x-mixed-replace".to_string(),
reader: Box::new(std::io::Cursor::new(b"JPEGDATA".to_vec())),
})
});
let caps = vec![PlainCapture::Stream {
id: "ext-1".to_string(),
open,
}];
let mgr = TimelapseManager::default();
mgr.start_plain(caps, 20, rx, dir.clone()).unwrap();
assert!(dir.join("ext-1").is_dir(), "per-camera dir created");
tokio::time::sleep(std::time::Duration::from_millis(80)).await;
tx.send(st("FINISH", Some(1))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(250)).await;
assert!(
!mgr.status_plain().running,
"stream recorder stops cleanly when the print finishes"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn live_mp4_args_pipe_mpjpeg_stdin_to_h264() {
let args = super::live_mp4_args(std::path::Path::new("/cap/ext-1/plain.mp4"));
let joined = args.join(" ");
assert!(joined.contains("-f mpjpeg"), "{joined}");
assert!(
joined.contains("-i -"),
"reads the stream from stdin: {joined}"
);
assert!(joined.contains("libx264"));
assert!(joined.trim_end().ends_with("/cap/ext-1/plain.mp4"));
}
#[tokio::test]
async fn smooth_and_plain_run_concurrently_and_stop_independently() {
let dir = std::env::temp_dir().join(format!("bambu-tl-both-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (tx, rx) = watch::channel(st("IDLE", None));
let g: FrameGrab = Arc::new(|| Ok(vec![0xff, 0xd8, 0xff, 0x01]));
let mgr = TimelapseManager::default();
mgr.start_smooth(
one("ext-0", g.clone()),
1,
vec![0],
rx.clone(),
dir.join("smooth"),
)
.unwrap();
mgr.start_plain(sample("ext-0", g), 20, rx, dir.join("plain"))
.unwrap();
assert!(mgr.status_smooth().running && mgr.status_plain().running);
tx.send(st("RUNNING", Some(1))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(60)).await;
tx.send(st("RUNNING", Some(2))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(60)).await;
assert!(mgr.status_smooth().frames >= 1, "smooth captured layers");
assert!(mgr.status_plain().frames >= 2, "plain sampled its interval");
assert!(mgr.stop_smooth());
assert!(!mgr.status_smooth().running);
assert!(
mgr.status_plain().running,
"plain keeps running after smooth stops"
);
assert!(mgr.stop_plain());
let _ = std::fs::remove_dir_all(&dir);
}
fn park_cap(id: &str) -> ParkCapture {
ParkCapture {
id: id.to_string(),
stream_url: "http://cam/stream".to_string(),
tuning: ParkTuning {
fps: 4.0,
left_frac: 0.33,
ema_seconds: 30.0,
abs_floor: 1500.0,
mad_k: 6.0,
merge_gap_s: 1.2,
max_island_s: 3.0,
min_sep_s: 3.0,
candidate_frac: 0.75,
warmup_s: 4.0,
baseline_s: 90.0,
},
}
}
#[tokio::test]
async fn park_spawns_one_worker_per_camera_on_active_and_stops_at_finish() {
let dir = std::env::temp_dir().join(format!("bambu-park-slot-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (tx, rx) = watch::channel(st("IDLE", None));
let spawned = Arc::new(AtomicUsize::new(0));
let spawn_worker: ParkSpawn = {
let spawned = spawned.clone();
Arc::new(
move |_cap, _dir, cancel: Arc<AtomicBool>, status: Arc<Mutex<TimelapseStatus>>| {
spawned.fetch_add(1, Ordering::SeqCst);
tokio::task::spawn_blocking(move || {
status.lock().unwrap().frames += 1; while !cancel.load(Ordering::Relaxed) {
std::thread::sleep(std::time::Duration::from_millis(10));
}
})
},
)
};
let mgr = TimelapseManager::default();
mgr.start_park(
vec![park_cap("ext-0"), park_cap("ext-1")],
rx,
dir.clone(),
spawn_worker,
)
.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(40)).await;
assert_eq!(spawned.load(Ordering::SeqCst), 0, "no workers while idle");
tx.send(st("RUNNING", Some(1))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(60)).await;
assert_eq!(spawned.load(Ordering::SeqCst), 2, "one per camera");
tx.send(st("RUNNING", Some(2))).unwrap(); tokio::time::sleep(std::time::Duration::from_millis(40)).await;
assert_eq!(
spawned.load(Ordering::SeqCst),
2,
"spawned once, not per tick"
);
assert_eq!(mgr.status_park().frames, 2);
tx.send(st("FINISH", Some(2))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(150)).await;
assert!(
!mgr.status_park().running,
"park stops when the print finishes"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn park_with_no_cameras_is_rejected() {
let dir = std::env::temp_dir().join(format!("bambu-park-empty-{}", std::process::id()));
let (_tx, rx) = watch::channel(st("RUNNING", Some(0)));
let mgr = TimelapseManager::default();
let noop: ParkSpawn = Arc::new(|_, _, _, _| tokio::task::spawn_blocking(|| {}));
assert!(
mgr.start_park(vec![], rx, dir, noop).is_err(),
"need at least one camera"
);
}
fn segment_cap(id: &str) -> SegmentCapture {
SegmentCapture {
id: id.to_string(),
stream_url: "http://cam/stream".to_string(),
fps: 10.0,
window_ms: 3000,
select_tuning: SelectTuning {
left_frac: 0.33,
min_outlier: 2.5,
min_left_density: 3.0,
select_candidate_frac: 0.6,
min_confidence: 0.40,
},
}
}
#[tokio::test]
async fn segment_spawns_per_camera_and_feeds_the_live_layer() {
let dir = std::env::temp_dir().join(format!("bambu-seg-slot-{}", std::process::id()));
let _ = std::fs::remove_dir_all(&dir);
let (tx, rx) = watch::channel(st("IDLE", None));
let spawned = Arc::new(AtomicUsize::new(0));
let seen_layer = Arc::new(AtomicI64::new(i64::MIN));
let spawn_worker: SegmentSpawn = {
let (spawned, seen_layer) = (spawned.clone(), seen_layer.clone());
Arc::new(
move |_cap,
_dir,
current_layer: Arc<AtomicI64>,
cancel: Arc<AtomicBool>,
status: Arc<Mutex<TimelapseStatus>>| {
spawned.fetch_add(1, Ordering::SeqCst);
status.lock().unwrap().frames += 1; let seen_layer = seen_layer.clone();
tokio::task::spawn_blocking(move || {
while !cancel.load(Ordering::Relaxed) {
seen_layer
.store(current_layer.load(Ordering::Relaxed), Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(5));
}
})
},
)
};
let mgr = TimelapseManager::default();
mgr.start_segment(
vec![segment_cap("ext-0"), segment_cap("ext-1")],
rx,
dir.clone(),
spawn_worker,
)
.unwrap();
tokio::time::sleep(std::time::Duration::from_millis(40)).await;
assert_eq!(spawned.load(Ordering::SeqCst), 0, "no workers while idle");
tx.send(st("RUNNING", Some(7))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(40)).await;
tx.send(st("RUNNING", Some(8))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(40)).await;
assert_eq!(
spawned.load(Ordering::SeqCst),
2,
"one per camera, spawned once"
);
assert_eq!(mgr.status_segment().frames, 2);
assert_eq!(
seen_layer.load(Ordering::SeqCst),
8,
"the worker reads the latest MQTT layer through the shared atomic"
);
tx.send(st("FINISH", Some(8))).unwrap();
tokio::time::sleep(std::time::Duration::from_millis(120)).await;
assert!(
!mgr.status_segment().running,
"segment stops when the print finishes"
);
let _ = std::fs::remove_dir_all(&dir);
}
#[tokio::test]
async fn segment_with_no_cameras_is_rejected() {
let dir = std::env::temp_dir().join(format!("bambu-seg-empty-{}", std::process::id()));
let (_tx, rx) = watch::channel(st("RUNNING", Some(0)));
let mgr = TimelapseManager::default();
let noop: SegmentSpawn = Arc::new(|_, _, _, _, _| tokio::task::spawn_blocking(|| {}));
assert!(
mgr.start_segment(vec![], rx, dir, noop).is_err(),
"need at least one camera"
);
}
}