diff --git a/kernel/relayflowd/src/engine/drive.rs b/kernel/relayflowd/src/engine/drive.rs index 2005b9e05..7eb27f0eb 100644 --- a/kernel/relayflowd/src/engine/drive.rs +++ b/kernel/relayflowd/src/engine/drive.rs @@ -1,4 +1,8 @@ -use std::{collections::BTreeSet, thread, time::Duration}; +use std::{ + collections::{BTreeMap, BTreeSet}, + thread, + time::Duration, +}; use anyhow::{Context, Result, bail}; use relayflowd_core::{ @@ -50,6 +54,12 @@ impl Engine { let mut backpressured = false; let mut handoff_failed = false; let mut skipped_dispatches = BTreeSet::new(); + // The batch was elected at one instant but runs serially, so a + // start journaled after a deterministic peer ran is late by that + // peer's runtime. Stamp each start when it is journaled and move + // its lease deadline — and its dispatch's — by the same delay, so + // neither the attempt's time nor its lease is spent in advance. + let mut start_delays = BTreeMap::new(); for action in actions { match action { Action::Append(mut entry) => { @@ -106,6 +116,8 @@ impl Engine { self.prepare_start_entry(&state, &mut entry)?; self.assign_executor(&state, &mut entry)?; self.route_start(&mut journal, &state, &mut entry)?; + // After placement, which can itself be slow. + self.stamp_start(&mut entry, &mut start_delays); self.append(&mut journal, &entry)?; Ok(()) })(); @@ -194,6 +206,11 @@ impl Engine { if skipped_dispatches.remove(&(step.id.clone(), attempt)) { continue; } + let lease_deadline_ms = lease_deadline_ms.saturating_add( + start_delays + .remove(&(step.id.clone(), attempt)) + .unwrap_or(0), + ); let input = self.resolve_step_input(&mut journal, &step, attempt); if !matches!(input, Ok(Some(_))) { if let Some(dispatcher) = &self.dispatcher { @@ -362,6 +379,32 @@ impl Engine { } } + /// Stamp a start with the moment it is journaled and move its lease + /// deadline by the delay since the batch was elected, recording the delay + /// so the batch's dispatch for the same attempt moves with it. Anything + /// but an attempt start is left alone. + fn stamp_start( + &self, + entry: &mut relayflowd_core::JournalEntry, + start_delays: &mut BTreeMap<(String, u32), i64>, + ) { + if entry.entry_type != relayflowd_core::EntryType::StepAttemptStarted { + return; + } + let now_ms = self.clock.now_ms(); + let delay = now_ms.saturating_sub(entry.at_ms); + if delay <= 0 { + return; + } + entry.at_ms = now_ms; + if let Some(deadline) = entry.payload["lease_deadline_ms"].as_i64() { + entry.payload["lease_deadline_ms"] = deadline.saturating_add(delay).into(); + } + if let (Some(step_id), Some(attempt)) = (entry.step_id.clone(), entry.attempt) { + start_delays.insert((step_id, attempt), delay); + } + } + /// Journal a dispatch that could not be honored as a declared `worker_error` /// completion of the attempt. It runs through the same completion path as /// any other rejected attempt, so retry policy and the failure taxonomy diff --git a/kernel/relayflowd/tests/late_start.rs b/kernel/relayflowd/tests/late_start.rs new file mode 100644 index 000000000..b0b8a71e2 --- /dev/null +++ b/kernel/relayflowd/tests/late_start.rs @@ -0,0 +1,130 @@ +//! A batch is elected at one instant and driven serially, so a start journaled +//! after a slow deterministic peer must be stamped when it is journaled: its +//! wall clock is its own, and its full lease — journaled and dispatched — is +//! still ahead of it. + +use std::sync::{Arc, Mutex}; + +use relayflowd::worker::{DispatchOutcome, JournalObserver, StepDispatch, StepDispatcher}; +use relayflowd::{Engine, RunStatus}; +use relayflowd_core::{EntryType, JournalEntry, RunSpec, StepType}; +use serde_json::json; +use tempfile::tempdir; + +/// The kernel's default lease for a step that declares none. +const DEFAULT_LEASE_MS: i64 = 30_000; + +#[derive(Default)] +struct Recorder { + run_id: Mutex>, + dispatches: Mutex>, +} + +impl JournalObserver for Recorder { + fn appended(&self, entry: &JournalEntry) { + if entry.entry_type == EntryType::RunSpawned { + *self.run_id.lock().unwrap() = Some(entry.run_id.clone()); + } + } +} + +impl StepDispatcher for Recorder { + fn executor(&self, _step_type: StepType) -> Option { + Some("late-start-test".to_owned()) + } + + fn available(&self, step_type: StepType) -> bool { + step_type == StepType::Llm + } + + fn dispatch(&self, dispatch: StepDispatch) -> anyhow::Result { + self.dispatches.lock().unwrap().push(dispatch); + Ok(DispatchOutcome::Dispatched) + } +} + +fn entry_of(entries: &[JournalEntry], entry_type: EntryType, step_id: &str) -> JournalEntry { + entries + .iter() + .find(|entry| entry.entry_type == entry_type && entry.step_id.as_deref() == Some(step_id)) + .cloned() + .unwrap_or_else(|| panic!("no {entry_type:?} for {step_id}")) +} + +#[test] +fn a_start_after_a_slow_deterministic_peer_is_stamped_when_journaled() { + let spec = RunSpec::parse(&json!({ + "name": "late-start", + "steps": [ + {"id": "slow", "type": "deterministic", "command": "sleep 0.4"}, + {"id": "quick", "type": "deterministic", "command": "true"} + ] + })) + .unwrap(); + let directory = tempdir().unwrap(); + let engine = Engine::new(directory.path()); + let outcome = engine.start(spec, "test", None).unwrap(); + assert_eq!(outcome.status, RunStatus::Completed); + let entries = engine + .journal_entries(&outcome.run_id, 1, usize::MAX) + .unwrap(); + + let slow_done = entry_of(&entries, EntryType::StepCompleted, "slow"); + let quick_start = entry_of(&entries, EntryType::StepAttemptStarted, "quick"); + let quick_done = entry_of(&entries, EntryType::StepCompleted, "quick"); + // Ordering, not a wall-clock bound: the start is journaled after the peer + // finished, so the step's own span cannot contain the peer's runtime. + assert!( + quick_start.at_ms >= slow_done.at_ms, + "quick started at {} before slow completed at {}", + quick_start.at_ms, + slow_done.at_ms + ); + assert_eq!( + quick_done.payload["spend"]["wallclock_ms"].as_i64().unwrap(), + quick_done.at_ms - quick_start.at_ms, + ); + for step in ["slow", "quick"] { + let start = entry_of(&entries, EntryType::StepAttemptStarted, step); + assert_eq!( + start.payload["lease_deadline_ms"].as_i64().unwrap() - start.at_ms, + DEFAULT_LEASE_MS, + "{step}'s lease must run from its own start" + ); + } +} + +#[test] +fn a_dispatch_after_a_slow_deterministic_peer_keeps_its_full_lease() { + let spec = RunSpec::parse(&json!({ + "name": "late-dispatch", + "steps": [ + {"id": "slow", "type": "deterministic", "command": "sleep 0.4"}, + {"id": "model", "type": "llm", "prompt": "p", "model": "stub"} + ] + })) + .unwrap(); + let directory = tempdir().unwrap(); + let recorder = Arc::new(Recorder::default()); + let engine = Engine::with_runtime(directory.path(), recorder.clone(), recorder.clone()); + engine.start(spec, "test", None).unwrap(); + let run_id = recorder.run_id.lock().unwrap().clone().unwrap(); + let entries = engine.journal_entries(&run_id, 1, usize::MAX).unwrap(); + + let slow_done = entry_of(&entries, EntryType::StepCompleted, "slow"); + let start = entry_of(&entries, EntryType::StepAttemptStarted, "model"); + let dispatch = recorder + .dispatches + .lock() + .unwrap() + .iter() + .find(|dispatch| dispatch.step_id == "model") + .cloned() + .unwrap(); + assert!(start.at_ms >= slow_done.at_ms); + assert_eq!( + dispatch.lease_deadline_ms, + start.payload["lease_deadline_ms"].as_i64().unwrap() + ); + assert_eq!(dispatch.lease_deadline_ms - start.at_ms, DEFAULT_LEASE_MS); +}