From 75b798c38499ed82b3266b529d7e6abaaf4635e4 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Sun, 26 Jul 2026 21:36:03 -0600 Subject: [PATCH 1/7] fix(process): preserve complete command stdout --- src/process/managed.rs | 25 +++++++++++++++++++++---- 1 file changed, 21 insertions(+), 4 deletions(-) diff --git a/src/process/managed.rs b/src/process/managed.rs index ccec5a86..b3736a4b 100644 --- a/src/process/managed.rs +++ b/src/process/managed.rs @@ -340,10 +340,13 @@ impl ManagedProcess { .wait_for_completion(options.wait_timeout) .with_raw_output( DEFAULT_OUTPUT_EOF_TIMEOUT, - RawOutputOptions::symmetric(RawCollectionOptions::Bounded { - max_bytes: options.stderr_limit.bytes(), - overflow_behavior: CollectionOverflowBehavior::DropOldestData, - }), + RawOutputOptions { + stdout_collection_options: RawCollectionOptions::TrustedUnbounded, + stderr_collection_options: RawCollectionOptions::Bounded { + max_bytes: options.stderr_limit.bytes(), + overflow_behavior: CollectionOverflowBehavior::DropOldestData, + }, + }, ) .or_terminate(shutdown) .await?; @@ -861,6 +864,20 @@ mod tests { assert!(String::from_utf8_lossy(&output.stderr).contains("frame=")); } + #[tokio::test] + async fn managed_process_output_keeps_complete_stdout_when_stderr_is_bounded() { + let cmd = fixture_command("stdout-noise-stderr-ffmpeg-progress"); + let options = ManagedProcessOptions::default().with_stderr_limit(4); + + let output = ManagedProcess::spawn_with_options("output fixture", cmd, options) + .expect("spawn shell fixture") + .output() + .await + .expect("collect output"); + + assert!(String::from_utf8_lossy(&output.stdout).contains("stdout-noise")); + } + #[tokio::test] async fn managed_process_terminates_after_timeout() { let cmd = fixture_command("sleep-long"); From 1644670f4b0b784977ee76343510596a4d235279 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Tue, 28 Jul 2026 18:42:22 -0600 Subject: [PATCH 2/7] fix(worker): use system temp directory --- src/command/worker.rs | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 deletions(-) diff --git a/src/command/worker.rs b/src/command/worker.rs index e21f115e..9346945c 100644 --- a/src/command/worker.rs +++ b/src/command/worker.rs @@ -1194,9 +1194,7 @@ fn heartbeat_payload(path: &Path, active_video_id: Option) -> HeartbeatPayl } fn worker_job_input_dir(job_id: &str) -> PathBuf { - std::env::current_dir() - .expect("current working directory") - .join(format!("ab-av1-worker-{}", job_id)) + std::env::temp_dir().join(format!("ab-av1-worker-{}", job_id)) } #[derive(Debug, Deserialize, Serialize, PartialEq, Eq)] @@ -6150,6 +6148,14 @@ mod tests { assert_eq!(status, "job_assigned (job_id=job-123)"); } + #[test] + fn worker_job_input_dir_uses_system_temp_dir() { + assert_eq!( + worker_job_input_dir("job-123"), + std::env::temp_dir().join("ab-av1-worker-job-123") + ); + } + #[test] fn worker_job_lowering_uses_an_isolated_temp_dir_and_target_vmaf() { let job_dir = From 61648018dbea1931d404e6ce1014bd9db7deda71 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Thu, 30 Jul 2026 11:09:18 -0600 Subject: [PATCH 3/7] fix(worker): reconnect dead Phoenix channels --- src/command/worker.rs | 225 ++++++++++++++++++++++++++++++++++-------- 1 file changed, 186 insertions(+), 39 deletions(-) diff --git a/src/command/worker.rs b/src/command/worker.rs index 9346945c..e07d6710 100644 --- a/src/command/worker.rs +++ b/src/command/worker.rs @@ -47,6 +47,7 @@ const TRANSFER_CHUNK_TYPE: u8 = 1; const TRANSFER_CHUNK_HEADER_LEN: usize = 52; const MAX_TRANSFER_FRAME_BYTES: usize = 640 * 1024 * 1024; const HEARTBEAT_INTERVAL: Duration = Duration::from_secs(10); +const EVENT_ACK_TIMEOUT: Duration = Duration::from_secs(30); const HTTP_TRANSFER_PROGRESS_INTERVAL: Duration = Duration::from_millis(500); const HTTP_TRANSFER_IDLE_TIMEOUT: Duration = Duration::from_secs(60); static HEARTBEAT_SYSTEM: OnceLock> = OnceLock::new(); @@ -893,6 +894,13 @@ async fn run_worker_job_with_reporting( } } } + Some(WorkerPush::ChannelClosed) => { + debug!( + job_id = %job.assignment.job_id, + "worker channel closed during job" + ); + *worker = None; + } _ => {} } } @@ -1317,6 +1325,7 @@ enum WorkerPush { Cancel(CancelPayload), Control(ControlPayload), Started(TransferStartedPayload), + ChannelClosed, } #[derive(Debug)] @@ -1368,6 +1377,21 @@ struct MultiplexJob { struct PendingEventAck { job_id: String, name: &'static str, + sent_at: Instant, +} + +impl PendingEventAck { + fn new(job_id: impl Into, name: &'static str) -> Self { + Self { + job_id: job_id.into(), + name, + sent_at: Instant::now(), + } + } + + fn deadline(&self, timeout: Duration) -> Instant { + self.sent_at + timeout + } } struct MultiplexedWorker { @@ -2226,9 +2250,15 @@ impl ConnectedWorker { while let Some(message) = self.socket.next().await { match message.context("read websocket message")? { Message::Text(text) => { - if let Some(WorkerPush::Control(control)) = decode_worker_push(&text)? { - self.pending_control = Some(control.action); - continue; + match decode_worker_push(&text)? { + Some(WorkerPush::Control(control)) => { + self.pending_control = Some(control.action); + continue; + } + Some(WorkerPush::ChannelClosed) => { + bail!("worker channel closed while waiting for work"); + } + _ => {} } if let Some(reply) = decode_expected_reply(&text, &expected_ref, "pull_work")? { return reply; @@ -2358,6 +2388,9 @@ impl ConnectedWorker { Ok(PendingJobOutcome::Waiting) } }, + Some(WorkerPush::ChannelClosed) => { + bail!("worker channel closed while waiting for input") + } Some(_) => Ok(PendingJobOutcome::Waiting), None => Ok(PendingJobOutcome::Waiting), } @@ -2453,23 +2486,29 @@ impl ConnectedWorker { | Some(Ok(Message::Binary(_))) | Some(Ok(Message::Frame(_))) => {} Some(Ok(Message::Text(text))) => { - if let Some(WorkerPush::Control(control)) = decode_worker_push(&text)? { - match control.action { - ControlAction::Start | ControlAction::Resume => { - self.send_control_state(ControlState::Running, None).await?; - *control_state = WorkerControlState::Running; - return Ok(stopped); - } - ControlAction::Pause if !stopped => { - self.send_control_state(ControlState::Paused, None).await?; - *control_state = WorkerControlState::Paused; - } - ControlAction::Pause | ControlAction::Stop => { - self.send_control_state(ControlState::Stopped, None).await?; - *control_state = WorkerControlState::Stopped; - stopped = true; + match decode_worker_push(&text)? { + Some(WorkerPush::Control(control)) => { + match control.action { + ControlAction::Start | ControlAction::Resume => { + self.send_control_state(ControlState::Running, None).await?; + *control_state = WorkerControlState::Running; + return Ok(stopped); + } + ControlAction::Pause if !stopped => { + self.send_control_state(ControlState::Paused, None).await?; + *control_state = WorkerControlState::Paused; + } + ControlAction::Pause | ControlAction::Stop => { + self.send_control_state(ControlState::Stopped, None).await?; + *control_state = WorkerControlState::Stopped; + stopped = true; + } } } + Some(WorkerPush::ChannelClosed) => { + bail!("worker channel closed while worker stopped") + } + _ => {} } } Some(Ok(Message::Close(frame))) => bail!("websocket closed while worker stopped: {frame:?}"), @@ -3235,6 +3274,9 @@ async fn run_multiplexed_worker( let mut scheduler = WorkScheduler::new(config.crf_searches_per_encode()); let mut reconnect = tokio::time::interval(runtime.reconnect_base_delay); reconnect.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + let mut heartbeat = tokio::time::interval(HEARTBEAT_INTERVAL); + heartbeat.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + heartbeat.tick().await; let mut offline_deadline = None; loop { @@ -3255,6 +3297,7 @@ async fn run_multiplexed_worker( let mut connection_lost = false; if let Some(worker) = connection.as_mut() { + let event_ack_deadline = next_event_ack_deadline(&pending_acks, EVENT_ACK_TIMEOUT); tokio::select! { item = outputs.recv() => { connection_lost = !handle_multiplex_output( @@ -3279,6 +3322,17 @@ async fn run_multiplexed_worker( config.local_path.as_deref(), ).await?; } + _ = heartbeat.tick() => { + connection_lost = !send_multiplex_heartbeat( + worker, + &jobs, + &mut pending_acks, + ).await; + } + _ = wait_for_deadline(event_ack_deadline) => { + warn!("worker channel acknowledgement timed out; reconnecting"); + connection_lost = true; + } _ = reconnect.tick() => { no_work.clear(); } @@ -3307,6 +3361,9 @@ async fn run_multiplexed_worker( offline_deadline = None; reconnect = tokio::time::interval(runtime.reconnect_base_delay); reconnect.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + heartbeat = tokio::time::interval(HEARTBEAT_INTERVAL); + heartbeat.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + heartbeat.tick().await; } } Err(error) => { @@ -3314,7 +3371,7 @@ async fn run_multiplexed_worker( } } } - _ = wait_for_offline_deadline(offline_deadline) => { + _ = wait_for_deadline(offline_deadline) => { expire_offline_jobs(&mut jobs, &mut pending, &mut outputs).await?; bail!("coordinator unavailable until worker job lease expired"); } @@ -3337,7 +3394,7 @@ async fn run_multiplexed_worker( } } -async fn wait_for_offline_deadline(deadline: Option) { +async fn wait_for_deadline(deadline: Option) { match deadline { Some(deadline) => { tokio::time::sleep_until(tokio::time::Instant::from_std(deadline)).await; @@ -3346,6 +3403,44 @@ async fn wait_for_offline_deadline(deadline: Option) { } } +fn next_event_ack_deadline( + pending_acks: &HashMap, + timeout: Duration, +) -> Option { + pending_acks.values().map(|ack| ack.deadline(timeout)).min() +} + +async fn send_multiplex_heartbeat( + worker: &mut MultiplexedWorker, + jobs: &HashMap, + pending_acks: &mut HashMap, +) -> bool { + if pending_acks.values().any(|ack| ack.name == "heartbeat") { + return true; + } + + let active_job = jobs.values().find(|job| !job.finished); + let path = active_job.map_or(Path::new("."), |job| job.job.input_dir.as_path()); + let active_video_id = active_job.map(|job| job.job.assignment.video_id); + + match worker + .send_event(ClientEvent::Heartbeat(heartbeat_payload( + path, + active_video_id, + ))) + .await + { + Ok(reference) => { + pending_acks.insert(reference, PendingEventAck::new("", "heartbeat")); + true + } + Err(error) => { + debug!(error = %error, "multiplexed worker heartbeat failed"); + false + } + } +} + async fn expire_offline_jobs( jobs: &mut HashMap, pending: &mut HashMap)>, @@ -3563,13 +3658,7 @@ async fn send_multiplex_event( match worker.send_event(event).await { Ok(reference) => { if event_requires_ack(name) { - pending_acks.insert( - reference, - PendingEventAck { - job_id: job_id.to_owned(), - name, - }, - ); + pending_acks.insert(reference, PendingEventAck::new(job_id, name)); } true } @@ -3798,6 +3887,10 @@ async fn handle_multiplex_frame( .map_err(|_| anyhow!("worker input command channel closed"))?; } } + Some(WorkerPush::ChannelClosed) => { + debug!("worker Phoenix channel closed"); + return Ok(false); + } None => {} } Ok(true) @@ -4128,6 +4221,7 @@ fn decode_worker_push(text: &str) -> Result> { "received worker push" ); let push = match frame.3.as_str() { + "phx_error" | "phx_close" => WorkerPush::ChannelClosed, "cancel" => WorkerPush::Cancel( serde_json::from_value::(payload.clone()) .context("decode cancel push")?, @@ -4928,20 +5022,14 @@ mod tests { let mut pending = HashMap::new(); let mut pending_acks = HashMap::from([( "7".into(), - PendingEventAck { - job_id: "encode-123".into(), - name: "encode_completed", - }, + PendingEventAck::new("encode-123", "encode_completed"), )]); remove_finished_job_if_acknowledged("encode-123", &mut jobs, &mut pending, &pending_acks); assert!(jobs.contains_key("encode-123")); acknowledge_multiplex_event( - &PendingEventAck { - job_id: "encode-123".into(), - name: "control_state", - }, + &PendingEventAck::new("encode-123", "control_state"), &mut jobs, &mut pending, ); @@ -4964,6 +5052,32 @@ mod tests { assert!(!event_requires_ack("encode_progress")); } + #[test] + fn phoenix_channel_errors_are_connection_loss() -> Result<()> { + for event in ["phx_error", "phx_close"] { + let text = json!([null, null, CRF_SEARCH_TOPIC, event, {}]).to_string(); + + assert!(matches!( + decode_worker_push(&text)?, + Some(WorkerPush::ChannelClosed) + )); + } + + Ok(()) + } + + #[test] + fn pending_event_ack_sets_liveness_deadline() { + let timeout = Duration::from_secs(30); + let ack = PendingEventAck { + job_id: String::new(), + name: "heartbeat", + sent_at: Instant::now() - timeout, + }; + + assert!(ack.deadline(timeout) <= Instant::now()); + } + #[test] fn discarded_active_attempt_is_cancelled() -> Result<()> { let (command, mut commands) = mpsc::unbounded_channel(); @@ -4979,10 +5093,7 @@ mod tests { let mut pending = HashMap::new(); handle_multiplex_ack( - &PendingEventAck { - job_id: "encode-stale".into(), - name: "job_active", - }, + &PendingEventAck::new("encode-stale", "job_active"), &serde_json::json!({ "accepted": false, "discarded": true, @@ -5941,6 +6052,42 @@ mod tests { Ok(()) } + #[tokio::test(flavor = "current_thread")] + async fn multiplexed_channel_error_requests_reconnect() -> Result<()> { + let coordinator = FakeCoordinator::with_no_work_replies(0).await?; + let connected = + ConnectedWorker::connect(&coordinator.worker_config(WorkerTestConfig::continuous())) + .await?; + let mut worker = MultiplexedWorker::from_connected(connected); + let mut jobs = HashMap::new(); + let mut pending = HashMap::new(); + let mut no_work = HashMap::new(); + let mut pending_acks = HashMap::new(); + let mut scheduler = WorkScheduler::new(std::num::NonZeroUsize::new(1).unwrap()); + let (output, _outputs) = mpsc::unbounded_channel(); + let mut completed_pulls = 0; + let frame = + Message::Text(json!([null, null, CRF_SEARCH_TOPIC, "phx_error", {}]).to_string()); + + let connection_is_alive = handle_multiplex_frame( + &mut worker, + Some(Ok(frame)), + &mut jobs, + &mut pending, + &mut no_work, + &mut pending_acks, + &mut scheduler, + &output, + &mut completed_pulls, + None, + ) + .await?; + + assert!(!connection_is_alive); + coordinator.finish().await; + Ok(()) + } + #[tokio::test(flavor = "current_thread")] async fn worker_requests_input_resend_when_resumed_job_lacks_local_file() -> Result<()> { let job_id = "missing-input-resend"; From b15575e596e4296510dbc639256db82a0bd9f613 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Thu, 30 Jul 2026 11:54:57 -0600 Subject: [PATCH 4/7] refactor(worker): normalize connection loss --- src/command/worker.rs | 529 +++++++++++++++++++----------------------- 1 file changed, 233 insertions(+), 296 deletions(-) diff --git a/src/command/worker.rs b/src/command/worker.rs index e07d6710..7855c719 100644 --- a/src/command/worker.rs +++ b/src/command/worker.rs @@ -839,82 +839,54 @@ async fn run_worker_job_with_reporting( } } } - frame = current_worker.socket.next() => { + frame = current_worker.next_frame() => { match frame { - Some(Ok(Message::Ping(payload))) => { - current_worker - .socket - .send(Message::Pong(payload)) - .await - .context("send websocket pong")?; + Ok(WorkerFrame::Push(WorkerPush::Cancel(cancel))) + if cancel.job_id == job.assignment.job_id => + { + eprintln!( + "worker job {} canceled: {}", + cancel.job_id, cancel.reason + ); + return Err(anyhow!( + "worker job {} canceled: {}", + cancel.job_id, cancel.reason + )); } - Some(Ok(Message::Pong(_))) => {} - Some(Ok(Message::Text(text))) => { - match decode_worker_push(&text)? { - Some(WorkerPush::Cancel(cancel)) - if cancel.job_id == job.assignment.job_id => - { - eprintln!( - "worker job {} canceled: {}", - cancel.job_id, cancel.reason - ); - return Err(anyhow!( - "worker job {} canceled: {}", - cancel.job_id, cancel.reason - )); + Ok(WorkerFrame::Push(WorkerPush::Control(control))) + if control.video_id.is_none() + || control.video_id == Some(job.assignment.video_id) => + { + match control.action { + ControlAction::Pause => { + crate::process::managed::pause_active_processes()?; + paused = true; + current_worker.send_control_state( + ControlState::Paused, + Some(job.assignment.video_id), + ).await?; } - Some(WorkerPush::Control(control)) - if control.video_id.is_none() - || control.video_id == Some(job.assignment.video_id) => - { - match control.action { - ControlAction::Pause => { - crate::process::managed::pause_active_processes()?; - paused = true; - current_worker.send_control_state( - ControlState::Paused, - Some(job.assignment.video_id), - ).await?; - } - ControlAction::Resume | ControlAction::Start => { - crate::process::managed::resume_active_processes()?; - paused = false; - current_worker.send_control_state( - ControlState::Running, - Some(job.assignment.video_id), - ).await?; - } - ControlAction::Stop => { - crate::process::managed::resume_active_processes()?; - current_worker.send_control_state( - ControlState::Stopped, - None, - ).await?; - return Ok(WorkerJobOutcome::Stopped); - } - } + ControlAction::Resume | ControlAction::Start => { + crate::process::managed::resume_active_processes()?; + paused = false; + current_worker.send_control_state( + ControlState::Running, + Some(job.assignment.video_id), + ).await?; } - Some(WorkerPush::ChannelClosed) => { - debug!( - job_id = %job.assignment.job_id, - "worker channel closed during job" - ); - *worker = None; + ControlAction::Stop => { + crate::process::managed::resume_active_processes()?; + current_worker.send_control_state( + ControlState::Stopped, + None, + ).await?; + return Ok(WorkerJobOutcome::Stopped); } - _ => {} } } - Some(Ok(Message::Binary(_))) | Some(Ok(Message::Frame(_))) => {} - Some(Ok(Message::Close(frame))) => { - debug!(job_id = %job.assignment.job_id, ?frame, "worker socket closed during job"); - *worker = None; - } - Some(Err(error)) => { - debug!(job_id = %job.assignment.job_id, error = %error, "worker socket lost during job"); - *worker = None; - } - None => { - debug!(job_id = %job.assignment.job_id, "worker websocket ended during job"); + Ok(_) => {} + Err(error) => { + debug!(job_id = %job.assignment.job_id, %error, "worker connection lost during job"); *worker = None; } } @@ -1325,7 +1297,14 @@ enum WorkerPush { Cancel(CancelPayload), Control(ControlPayload), Started(TransferStartedPayload), - ChannelClosed, +} + +#[derive(Debug)] +enum WorkerFrame { + Push(WorkerPush), + Text(String), + Binary(Vec), + Ping(Vec), } #[derive(Debug)] @@ -1441,6 +1420,16 @@ impl MultiplexedWorker { .await .context("send websocket pong") } + + async fn next_frame(&mut self) -> Result { + loop { + match decode_worker_frame(self.reader.next().await)? { + Some(WorkerFrame::Ping(payload)) => self.send_pong(payload).await?, + Some(frame) => return Ok(frame), + None => {} + } + } + } } fn multiplex_event( @@ -2247,36 +2236,19 @@ impl ConnectedWorker { send_json(&mut self.socket, frame).await?; let expected_ref = request_ref.to_string(); - while let Some(message) = self.socket.next().await { - match message.context("read websocket message")? { - Message::Text(text) => { - match decode_worker_push(&text)? { - Some(WorkerPush::Control(control)) => { - self.pending_control = Some(control.action); - continue; - } - Some(WorkerPush::ChannelClosed) => { - bail!("worker channel closed while waiting for work"); - } - _ => {} - } + loop { + match self.next_frame().await? { + WorkerFrame::Push(WorkerPush::Control(control)) => { + self.pending_control = Some(control.action); + } + WorkerFrame::Text(text) => { if let Some(reply) = decode_expected_reply(&text, &expected_ref, "pull_work")? { return reply; } } - Message::Ping(payload) => { - self.socket - .send(Message::Pong(payload)) - .await - .context("send websocket pong")?; - } - Message::Close(frame) => { - bail!("websocket closed while waiting for work: {frame:?}") - } - Message::Pong(_) | Message::Binary(_) | Message::Frame(_) => {} + WorkerFrame::Push(_) | WorkerFrame::Binary(_) | WorkerFrame::Ping(_) => {} } } - bail!("websocket ended while waiting for work") } fn take_pending_control(&mut self) -> Option { @@ -2289,6 +2261,21 @@ impl ConnectedWorker { send_json(&mut self.socket, ClientFrame::new(request_ref, event)).await } + async fn next_frame(&mut self) -> Result { + loop { + match decode_worker_frame(self.socket.next().await)? { + Some(WorkerFrame::Ping(payload)) => { + self.socket + .send(Message::Pong(payload)) + .await + .context("send websocket pong")?; + } + Some(frame) => return Ok(frame), + None => {} + } + } + } + async fn send_control_state( &mut self, state: ControlState, @@ -2327,75 +2314,58 @@ impl ConnectedWorker { idle_delay: Duration, ) -> Result { tokio::select! { - frame = self.socket.next() => { + frame = self.next_frame() => { match frame { - Some(Ok(Message::Ping(payload))) => { - self.socket - .send(Message::Pong(payload)) - .await - .context("send websocket pong")?; + Ok(WorkerFrame::Push(WorkerPush::Cancel(cancel))) + if cancel.job_id == pending_job.job().assignment.job_id => + { + eprintln!( + "worker job {} canceled: {}", + cancel.job_id, cancel.reason + ); + Ok(PendingJobOutcome::Canceled) + } + Ok(WorkerFrame::Push(WorkerPush::Started(started))) + if started.transfer_id == pending_job.job().assignment.job_id => + { + pending_job.ensure_receiver(started.chunk_size_bytes)?; + debug!( + job_id = %started.transfer_id, + source_name = %started.source_name, + chunk_size_bytes = started.chunk_size_bytes, + size_bytes = started.size_bytes, + total_bytes = started.total_bytes, + total_chunks = started.total_chunks, + received_bytes = pending_job + .receiver + .as_ref() + .map(ChunkReceiver::received_bytes) + .unwrap_or_default(), + "transfer started" + ); Ok(PendingJobOutcome::Waiting) } - Some(Ok(Message::Pong(_))) => Ok(PendingJobOutcome::Waiting), - Some(Ok(Message::Text(text))) => { - match decode_worker_push(&text)? { - Some(WorkerPush::Cancel(cancel)) - if cancel.job_id == pending_job.job().assignment.job_id => - { - eprintln!( - "worker job {} canceled: {}", - cancel.job_id, cancel.reason - ); - Ok(PendingJobOutcome::Canceled) - } - Some(WorkerPush::Started(started)) - if started.transfer_id == pending_job.job().assignment.job_id => - { - pending_job.ensure_receiver(started.chunk_size_bytes)?; - debug!( - job_id = %started.transfer_id, - source_name = %started.source_name, - chunk_size_bytes = started.chunk_size_bytes, - size_bytes = started.size_bytes, - total_bytes = started.total_bytes, - total_chunks = started.total_chunks, - received_bytes = pending_job - .receiver - .as_ref() - .map(ChunkReceiver::received_bytes) - .unwrap_or_default(), - "transfer started" - ); - Ok(PendingJobOutcome::Waiting) - } - Some(WorkerPush::Control(control)) => match control.action { - ControlAction::Stop => { - self.send_control_state(ControlState::Stopped, None).await?; - Ok(PendingJobOutcome::Stopped) - } - ControlAction::Pause => { - self.send_control_state( - ControlState::Paused, - Some(pending_job.job.assignment.video_id), - ).await?; - Ok(PendingJobOutcome::Paused) - } - ControlAction::Resume | ControlAction::Start => { - self.send_control_state( - ControlState::Running, - Some(pending_job.job.assignment.video_id), - ).await?; - Ok(PendingJobOutcome::Waiting) - } - }, - Some(WorkerPush::ChannelClosed) => { - bail!("worker channel closed while waiting for input") - } - Some(_) => Ok(PendingJobOutcome::Waiting), - None => Ok(PendingJobOutcome::Waiting), + Ok(WorkerFrame::Push(WorkerPush::Control(control))) => match control.action { + ControlAction::Stop => { + self.send_control_state(ControlState::Stopped, None).await?; + Ok(PendingJobOutcome::Stopped) } - } - Some(Ok(Message::Binary(bytes))) => { + ControlAction::Pause => { + self.send_control_state( + ControlState::Paused, + Some(pending_job.job.assignment.video_id), + ).await?; + Ok(PendingJobOutcome::Paused) + } + ControlAction::Resume | ControlAction::Start => { + self.send_control_state( + ControlState::Running, + Some(pending_job.job.assignment.video_id), + ).await?; + Ok(PendingJobOutcome::Waiting) + } + }, + Ok(WorkerFrame::Binary(bytes)) => { let chunk = decode_binary_worker_push(&bytes)?; if let Some(chunk) = chunk && chunk.transfer_id == pending_job.job().assignment.job_id @@ -2439,14 +2409,12 @@ impl ConnectedWorker { } Ok(PendingJobOutcome::Waiting) } - Some(Ok(Message::Frame(_))) => { + Ok(WorkerFrame::Push(_) | WorkerFrame::Text(_) | WorkerFrame::Ping(_)) => { Ok(PendingJobOutcome::Waiting) } - Some(Ok(Message::Close(frame))) => { - bail!("websocket closed while waiting for worker input: {frame:?}") + Err(error) => { + Err(error).context("read worker websocket frame while waiting for input") } - Some(Err(error)) => Err(error).context("read websocket message"), - None => bail!("websocket ended while waiting for worker input"), } } _ = tokio::time::sleep(idle_delay) => { @@ -2478,42 +2446,29 @@ impl ConnectedWorker { let mut stopped = *control_state == WorkerControlState::Stopped; loop { tokio::select! { - frame = self.socket.next() => match frame { - Some(Ok(Message::Ping(payload))) => { - self.socket.send(Message::Pong(payload)).await.context("send websocket pong")?; - } - Some(Ok(Message::Pong(_))) - | Some(Ok(Message::Binary(_))) - | Some(Ok(Message::Frame(_))) => {} - Some(Ok(Message::Text(text))) => { - match decode_worker_push(&text)? { - Some(WorkerPush::Control(control)) => { - match control.action { - ControlAction::Start | ControlAction::Resume => { - self.send_control_state(ControlState::Running, None).await?; - *control_state = WorkerControlState::Running; - return Ok(stopped); - } - ControlAction::Pause if !stopped => { - self.send_control_state(ControlState::Paused, None).await?; - *control_state = WorkerControlState::Paused; - } - ControlAction::Pause | ControlAction::Stop => { - self.send_control_state(ControlState::Stopped, None).await?; - *control_state = WorkerControlState::Stopped; - stopped = true; - } - } + frame = self.next_frame() => match frame? { + WorkerFrame::Push(WorkerPush::Control(control)) => { + match control.action { + ControlAction::Start | ControlAction::Resume => { + self.send_control_state(ControlState::Running, None).await?; + *control_state = WorkerControlState::Running; + return Ok(stopped); + } + ControlAction::Pause if !stopped => { + self.send_control_state(ControlState::Paused, None).await?; + *control_state = WorkerControlState::Paused; } - Some(WorkerPush::ChannelClosed) => { - bail!("worker channel closed while worker stopped") + ControlAction::Pause | ControlAction::Stop => { + self.send_control_state(ControlState::Stopped, None).await?; + *control_state = WorkerControlState::Stopped; + stopped = true; } - _ => {} } } - Some(Ok(Message::Close(frame))) => bail!("websocket closed while worker stopped: {frame:?}"), - Some(Err(error)) => return Err(error).context("read websocket message while worker stopped"), - None => bail!("websocket ended while worker stopped"), + WorkerFrame::Push(_) + | WorkerFrame::Text(_) + | WorkerFrame::Binary(_) + | WorkerFrame::Ping(_) => {} }, _ = tokio::time::sleep(idle_delay) => { self.send_event(ClientEvent::Heartbeat(heartbeat_payload(Path::new("."), None))).await?; @@ -3308,9 +3263,8 @@ async fn run_multiplexed_worker( &mut pending_acks, ).await?; } - frame = worker.reader.next() => { + frame = worker.next_frame() => { connection_lost = !handle_multiplex_frame( - worker, frame, &mut jobs, &mut pending, @@ -3320,7 +3274,7 @@ async fn run_multiplexed_worker( &output, completed_pulls, config.local_path.as_deref(), - ).await?; + )?; } _ = heartbeat.tick() => { connection_lost = !send_multiplex_heartbeat( @@ -3700,9 +3654,8 @@ fn record_multiplex_event(state: &mut WorkerJobReportState, event: &ClientEvent) } #[allow(clippy::too_many_arguments)] -async fn handle_multiplex_frame( - worker: &mut MultiplexedWorker, - frame: Option>, +fn handle_multiplex_frame( + frame: Result, jobs: &mut HashMap, pending: &mut HashMap)>, no_work: &mut HashMap, @@ -3712,24 +3665,15 @@ async fn handle_multiplex_frame( completed_pulls: &mut usize, local_path: Option<&Path>, ) -> Result { - let Some(frame) = frame else { - return Ok(false); - }; let frame = match frame { Ok(frame) => frame, Err(error) => { - debug!(error = %error, "multiplexed worker websocket read failed"); + debug!(%error, "multiplexed worker connection lost"); return Ok(false); } }; match frame { - Message::Ping(payload) => worker.send_pong(payload).await.map(|()| true), - Message::Pong(_) | Message::Frame(_) => Ok(true), - Message::Close(frame) => { - debug!(?frame, "multiplexed worker websocket closed"); - Ok(false) - } - Message::Binary(bytes) => { + WorkerFrame::Binary(bytes) => { if let Some(chunk) = decode_binary_worker_push(&bytes)? && let Some(job) = jobs.get(&chunk.transfer_id) { @@ -3739,7 +3683,7 @@ async fn handle_multiplex_frame( } Ok(true) } - Message::Text(text) => { + WorkerFrame::Text(text) => { let mut event_ack = None; for (reference, ack) in pending_acks.iter() { if let Some(decoded) = decode_expected_reply::(&text, reference, ack.name)? { @@ -3853,48 +3797,45 @@ async fn handle_multiplex_frame( return Ok(true); } - match decode_worker_push(&text)? { - Some(WorkerPush::Cancel(cancel)) => { - if let Some(job) = jobs.get(&cancel.job_id) { - job.command - .send(JobCommand::Cancel(cancel)) - .map_err(|_| anyhow!("worker input command channel closed"))?; - } - } - Some(WorkerPush::Control(control)) => { - let jobs: Vec<_> = jobs - .values() - .filter(|job| { - control_targets_job( - &control, - &job.job.assignment.job_id, - job.job.assignment.video_id, - job.finished, - ) - }) - .collect(); + Ok(true) + } + WorkerFrame::Push(WorkerPush::Cancel(cancel)) => { + if let Some(job) = jobs.get(&cancel.job_id) { + job.command + .send(JobCommand::Cancel(cancel)) + .map_err(|_| anyhow!("worker input command channel closed"))?; + } + Ok(true) + } + WorkerFrame::Push(WorkerPush::Control(control)) => { + let jobs: Vec<_> = jobs + .values() + .filter(|job| { + control_targets_job( + &control, + &job.job.assignment.job_id, + job.job.assignment.video_id, + job.finished, + ) + }) + .collect(); - for job in jobs { - job.command - .send(JobCommand::Control(control.clone())) - .map_err(|_| anyhow!("worker control channel closed"))?; - } - } - Some(WorkerPush::Started(started)) => { - if let Some(job) = jobs.get(&started.transfer_id) { - job.command - .send(JobCommand::TransferStarted(started)) - .map_err(|_| anyhow!("worker input command channel closed"))?; - } - } - Some(WorkerPush::ChannelClosed) => { - debug!("worker Phoenix channel closed"); - return Ok(false); - } - None => {} + for job in jobs { + job.command + .send(JobCommand::Control(control.clone())) + .map_err(|_| anyhow!("worker control channel closed"))?; + } + Ok(true) + } + WorkerFrame::Push(WorkerPush::Started(started)) => { + if let Some(job) = jobs.get(&started.transfer_id) { + job.command + .send(JobCommand::TransferStarted(started)) + .map_err(|_| anyhow!("worker input command channel closed"))?; } Ok(true) } + WorkerFrame::Ping(_) => Ok(true), } } @@ -4201,6 +4142,25 @@ fn control_targets_job( } } +fn decode_worker_frame( + frame: Option>, +) -> Result> { + let frame = frame + .ok_or_else(|| anyhow!("worker websocket ended"))? + .context("read worker websocket frame")?; + + match frame { + Message::Text(text) => Ok(match decode_worker_push(&text)? { + Some(push) => Some(WorkerFrame::Push(push)), + None => Some(WorkerFrame::Text(text)), + }), + Message::Binary(bytes) => Ok(Some(WorkerFrame::Binary(bytes))), + Message::Ping(payload) => Ok(Some(WorkerFrame::Ping(payload))), + Message::Pong(_) | Message::Frame(_) => Ok(None), + Message::Close(frame) => bail!("worker websocket closed: {frame:?}"), + } +} + fn decode_worker_push(text: &str) -> Result> { let frame: ServerPushFrame = match serde_json::from_str(text) { Ok(frame) => frame, @@ -4221,7 +4181,9 @@ fn decode_worker_push(text: &str) -> Result> { "received worker push" ); let push = match frame.3.as_str() { - "phx_error" | "phx_close" => WorkerPush::ChannelClosed, + "phx_error" | "phx_close" => { + bail!("worker Phoenix channel closed: {}", frame.3) + } "cancel" => WorkerPush::Cancel( serde_json::from_value::(payload.clone()) .context("decode cancel push")?, @@ -4426,23 +4388,16 @@ where + Unpin, T: for<'de> Deserialize<'de>, { - while let Some(message) = reader.next().await { - match message.context("read websocket message")? { - Message::Text(text) => { + loop { + match decode_worker_frame(reader.next().await)? { + Some(WorkerFrame::Text(text)) => { if let Some(reply) = decode_expected_reply(&text, expected_ref, expected_event)? { return reply; } } - Message::Close(frame) => { - bail!("websocket closed before {expected_event} reply: {frame:?}") - } - Message::Ping(_) | Message::Pong(_) | Message::Binary(_) | Message::Frame(_) => { - continue; - } + Some(WorkerFrame::Push(_) | WorkerFrame::Binary(_) | WorkerFrame::Ping(_)) | None => {} } } - - bail!("websocket ended before {expected_event} reply") } fn decode_expected_reply( @@ -5053,17 +5008,15 @@ mod tests { } #[test] - fn phoenix_channel_errors_are_connection_loss() -> Result<()> { + fn worker_frames_normalize_every_connection_loss() { for event in ["phx_error", "phx_close"] { let text = json!([null, null, CRF_SEARCH_TOPIC, event, {}]).to_string(); - assert!(matches!( - decode_worker_push(&text)?, - Some(WorkerPush::ChannelClosed) - )); + assert!(decode_worker_frame(Some(Ok(Message::Text(text)))).is_err()); } - Ok(()) + assert!(decode_worker_frame(Some(Ok(Message::Close(None)))).is_err()); + assert!(decode_worker_frame(None).is_err()); } #[test] @@ -6013,13 +5966,8 @@ mod tests { Ok(()) } - #[tokio::test(flavor = "current_thread")] - async fn multiplexed_pull_error_requests_reconnect() -> Result<()> { - let coordinator = FakeCoordinator::with_no_work_replies(0).await?; - let connected = - ConnectedWorker::connect(&coordinator.worker_config(WorkerTestConfig::continuous())) - .await?; - let mut worker = MultiplexedWorker::from_connected(connected); + #[test] + fn multiplexed_pull_error_requests_reconnect() -> Result<()> { let mut jobs = HashMap::new(); let mut pending = HashMap::from([("3".into(), (JobKind::CrfSearch, None))]); let mut no_work = HashMap::new(); @@ -6032,9 +5980,8 @@ mod tests { ReplyBody::error(ErrorReplyPayload::new("unmatched topic")), ))?); - let connection_is_lost = handle_multiplex_frame( - &mut worker, - Some(Ok(frame)), + let connection_is_alive = handle_multiplex_frame( + decode_worker_frame(Some(Ok(frame))).map(|frame| frame.expect("text worker frame")), &mut jobs, &mut pending, &mut no_work, @@ -6043,22 +5990,15 @@ mod tests { &output, &mut completed_pulls, None, - ) - .await?; + )?; - assert!(!connection_is_lost); + assert!(!connection_is_alive); assert_eq!(completed_pulls, 0); - coordinator.finish().await; Ok(()) } - #[tokio::test(flavor = "current_thread")] - async fn multiplexed_channel_error_requests_reconnect() -> Result<()> { - let coordinator = FakeCoordinator::with_no_work_replies(0).await?; - let connected = - ConnectedWorker::connect(&coordinator.worker_config(WorkerTestConfig::continuous())) - .await?; - let mut worker = MultiplexedWorker::from_connected(connected); + #[test] + fn multiplexed_channel_error_requests_reconnect() -> Result<()> { let mut jobs = HashMap::new(); let mut pending = HashMap::new(); let mut no_work = HashMap::new(); @@ -6070,8 +6010,7 @@ mod tests { Message::Text(json!([null, null, CRF_SEARCH_TOPIC, "phx_error", {}]).to_string()); let connection_is_alive = handle_multiplex_frame( - &mut worker, - Some(Ok(frame)), + decode_worker_frame(Some(Ok(frame))).map(|frame| frame.expect("text worker frame")), &mut jobs, &mut pending, &mut no_work, @@ -6080,11 +6019,9 @@ mod tests { &output, &mut completed_pulls, None, - ) - .await?; + )?; assert!(!connection_is_alive); - coordinator.finish().await; Ok(()) } From b38ac4f0618bd2057561a47fea068c7d94d31718 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Sun, 2 Aug 2026 11:33:23 -0600 Subject: [PATCH 5/7] fix(worker): discard stale reconnect attempts --- src/command/worker.rs | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/command/worker.rs b/src/command/worker.rs index 7855c719..888ef753 100644 --- a/src/command/worker.rs +++ b/src/command/worker.rs @@ -3873,6 +3873,8 @@ fn handle_multiplex_ack( })) .map_err(|_| anyhow!("worker job command channel closed"))?; } + jobs.remove(&ack.job_id); + pending.retain(|_, (_, resend)| resend.as_deref() != Some(&ack.job_id)); return Ok(()); } @@ -5061,7 +5063,7 @@ mod tests { Ok(JobCommand::Cancel(CancelPayload { ref job_id, .. })) if job_id == "encode-stale" )); - assert!(jobs.contains_key("encode-stale")); + assert!(!jobs.contains_key("encode-stale")); Ok(()) } From 60e6fd6fb5fa7c9bde9ab745b9a49a5e4ce548f3 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Thu, 6 Aug 2026 10:54:59 -0600 Subject: [PATCH 6/7] fix(worker): report temp filesystem capacity --- src/command/worker.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/command/worker.rs b/src/command/worker.rs index 888ef753..d5ad6182 100644 --- a/src/command/worker.rs +++ b/src/command/worker.rs @@ -2471,7 +2471,7 @@ impl ConnectedWorker { | WorkerFrame::Ping(_) => {} }, _ = tokio::time::sleep(idle_delay) => { - self.send_event(ClientEvent::Heartbeat(heartbeat_payload(Path::new("."), None))).await?; + self.send_event(ClientEvent::Heartbeat(heartbeat_payload(&std::env::temp_dir(), None))).await?; } } } From f47659d5276e0b30ad881f2a90447731884b30f6 Mon Sep 17 00:00:00 2001 From: Mika Cohen Date: Thu, 6 Aug 2026 11:36:51 -0600 Subject: [PATCH 7/7] fix(worker): report typed CRF failures --- src/command/crf_search.rs | 12 +++-- src/command/crf_search/decision.rs | 49 ++++++++++++++++++- src/command/worker.rs | 77 +++++++++++++++++++++++++++++- 3 files changed, 132 insertions(+), 6 deletions(-) diff --git a/src/command/crf_search.rs b/src/command/crf_search.rs index f845ab28..051966a5 100644 --- a/src/command/crf_search.rs +++ b/src/command/crf_search.rs @@ -508,11 +508,11 @@ pub fn run( } } - let sample = Sample { - crf: args.crf, + let sample = Sample::new( + args.crf, q, - enc: sample_enc_output.context("no sample output?")?, - }; + sample_enc_output.context("no sample output?")?, + ); let score = output_search_score(&sample.enc, use_xpsnr); crf_attempts.push(sample.clone()); @@ -559,6 +559,10 @@ pub struct Sample { } impl Sample { + pub(crate) fn new(crf: f32, q: i64, enc: sample_encode::Output) -> Self { + Self { enc, crf, q } + } + pub fn print_attempt( &self, bar: &ProgressBar, diff --git a/src/command/crf_search/decision.rs b/src/command/crf_search/decision.rs index 7201038f..61ca083b 100644 --- a/src/command/crf_search/decision.rs +++ b/src/command/crf_search/decision.rs @@ -88,7 +88,7 @@ pub fn decide_next_transition( if lower_score >= min_score { Error::ensure_or_no_good_crf( lower.enc.encode_percent <= max_encoded_percent.get(), - sample, + lower, )?; Ok(SearchTransition::RunResultThenDone { run_result: sample.clone(), @@ -147,3 +147,50 @@ pub fn guess_progress(run: usize, sample_progress: f32, thorough: bool) -> f64 { let sample_progress = sample_progress.clamp(0.0, 1.0) as f64; (((run - 1) as f64 + sample_progress) * BAR_LEN as f64 / total_runs_guess).min(BAR_LEN as f64) } + +#[cfg(test)] +mod tests { + use super::*; + use crate::command::sample_encode; + use std::time::Duration; + + fn sample(q: i64, score: f32, encode_percent: f64) -> Sample { + Sample::new( + q as f32, + q, + sample_encode::Output { + vmaf_score: Some(score), + xpsnr_score: None, + predicted_encode_size: 1_000, + encode_percent, + predicted_encode_time: Duration::from_secs(60), + from_cache: false, + }, + ) + } + + #[test] + fn adjacent_quality_candidate_reports_the_oversized_candidate() { + let lower = sample(30, 96.0, 81.0); + let current = sample(31, 94.0, 50.0); + let decision = SearchDecision { + min_score: 95.0, + higher_tolerance: 1.0, + thorough: false, + cut_on_iter2: false, + run: 2, + min_q: 5, + max_q: 70, + use_xpsnr: false, + max_encoded_percent: MaxEncodedPercent::new(80.0).expect("valid percentage"), + }; + + let error = decide_next_transition(¤t, 94.0, &[lower], decision) + .expect_err("the quality candidate exceeds the size limit"); + + assert!(matches!( + error, + Error::NoGoodCrf { last } if last.enc.encode_percent == 81.0 + )); + } +} diff --git a/src/command/worker.rs b/src/command/worker.rs index d5ad6182..cbfa73d8 100644 --- a/src/command/worker.rs +++ b/src/command/worker.rs @@ -551,6 +551,26 @@ impl WorkerJob { .find_map(|cause| cause.downcast_ref::()) .and_then(ProcessExitError::code) .unwrap_or(crate::FAILURE_EXIT_CODE); + let no_good_crf = + error + .chain() + .find_map(|cause| match cause.downcast_ref::() { + Some(crf_search::Error::NoGoodCrf { last }) => Some(last), + Some(crf_search::Error::Other(_)) | None => None, + }); + let max_encoded_percent = self + .crf_search_config() + .ok() + .map(|config| config.max_encoded_percent.get()); + let category = match (no_good_crf, max_encoded_percent) { + (Some(last), Some(max)) if last.enc.encode_percent > max => "size_limits", + (Some(_), _) => "crf_optimization", + (None, _) => "process_failure", + }; + let argv = match self.assignment.job_type { + JobKind::CrfSearch => &self.assignment.crf_search_args, + JobKind::Encode => &self.assignment.encode_args, + }; FailureReportPayload { job_id: self.assignment.job_id.clone(), @@ -560,7 +580,7 @@ impl WorkerJob { JobKind::Encode => "encoding", } .into(), - category: "process_failure".into(), + category: category.into(), message: error.to_string(), code: format!("EXIT_{exit_code}"), context: json!({ @@ -568,6 +588,9 @@ impl WorkerJob { "source_name": self.assignment.source_name, "exit_code": exit_code, "error_chain": format!("{error:#}"), + "argv": argv, + "encode_percent": no_good_crf.map(|last| last.enc.encode_percent), + "max_encoded_percent": max_encoded_percent, }), retriable: false, stderr_excerpt: Some(format!("{error:#}")), @@ -4762,6 +4785,58 @@ mod tests { assert_eq!(payload.category, "process_failure"); assert_eq!(payload.code, "EXIT_254"); assert_eq!(payload.context["exit_code"], json!(254)); + assert_eq!(payload.context["argv"], json!(["crf-search", "movie.mkv"])); + } + + #[test] + fn worker_crf_failure_reports_why_no_crf_was_acceptable() { + let job = WorkerJob::new( + JobAssignedPayload { + status: WorkStatus::JobAssigned, + job_type: JobKind::CrfSearch, + job_id: "crf-search-123".into(), + video_id: 123, + source_name: "movie.mkv".into(), + size_bytes: 1_000, + chunk_size_bytes: 256, + target_vmaf: 95.0, + transfer: None, + output_transfer: None, + output_shared_path: None, + encode_args: Vec::new(), + crf_search_args: vec![ + "crf-search".into(), + "--max-encoded-percent".into(), + "80".into(), + "--input".into(), + "movie.mkv".into(), + ], + }, + PathBuf::from("/tmp/crf-search-123"), + PathBuf::from("/tmp/crf-search-123/movie.mkv"), + ); + let sample = crf_search::Sample::new( + 30.0, + 30, + sample_encode::Output { + vmaf_score: Some(95.0), + xpsnr_score: None, + predicted_encode_size: 1_000, + encode_percent: 81.0, + predicted_encode_time: Duration::from_secs(60), + from_cache: false, + }, + ); + + let error = anyhow!(crf_search::Error::Other(anyhow!( + crf_search::Error::NoGoodCrf { last: sample } + ))) + .context("worker CRF search failed"); + let payload = job.failure_payload(&error); + + assert_eq!(payload.category, "size_limits"); + assert_eq!(payload.context["encode_percent"], json!(81.0)); + assert_eq!(payload.context["max_encoded_percent"], json!(80.0)); } #[test]