From 01eb0e32d62b77ec07015b76c05044e87d382df2 Mon Sep 17 00:00:00 2001 From: kjgbot Date: Wed, 2 Sep 2026 20:37:23 +0200 Subject: [PATCH] feat(kernel): add durable run cancellation Session-Id: 01a06343-9355-7393-95d7-d2fb2d972c73 --- kernel/DESIGN.md | 9 +- kernel/relayflowd-core/src/entry.rs | 9 + kernel/relayflowd-core/src/lib.rs | 1 + kernel/relayflowd-core/src/machine.rs | 7 + kernel/relayflowd-core/src/machine/cancel.rs | 100 +++++++ .../relayflowd-core/src/machine/recovery.rs | 6 + kernel/relayflowd-core/src/machine/tests.rs | 95 ++++++- kernel/relayflowd-core/src/state.rs | 11 +- kernel/relayflowd/src/engine.rs | 37 ++- kernel/relayflowd/src/engine/remote.rs | 3 + kernel/relayflowd/src/lib.rs | 2 +- kernel/relayflowd/src/main.rs | 30 +- kernel/relayflowd/src/server.rs | 20 +- kernel/relayflowd/src/server/cancel.rs | 35 +++ kernel/relayflowd/src/server/client.rs | 61 ++++ kernel/relayflowd/src/server/session.rs | 10 + kernel/relayflowd/tests/crash_resume.rs | 100 ++++++- .../tests/crash_resume/concurrency.rs | 81 +++++- .../20260902-v2-run-cancel-evidence.md | 268 ++++++++++++++++++ sdk/src/index.ts | 2 + sdk/src/journal-client.ts | 5 + sdk/src/protocol.ts | 7 + sdk/tests/journal-client-loopback.ts | 1 + sdk/tests/journal-client.test.ts | 19 ++ sdk/tests/live-kernel.test.ts | 44 +++ 25 files changed, 944 insertions(+), 19 deletions(-) create mode 100644 kernel/relayflowd-core/src/machine/cancel.rs create mode 100644 kernel/relayflowd/src/server/cancel.rs create mode 100644 ops/reviews/20260902-v2-run-cancel-evidence.md diff --git a/kernel/DESIGN.md b/kernel/DESIGN.md index 2b0d8b7e..3bdd525f 100644 --- a/kernel/DESIGN.md +++ b/kernel/DESIGN.md @@ -38,6 +38,12 @@ First entry of segment 1. Payload: (int, stamped per segment thereafter via `epoch.summary`), `created_by` (client identity string). +### 1.1a `run.cancel.requested` +Durable operator intent, appended before cancellation closes any work. Payload: +`requested_by` (client identity string). Resume treats this entry as a one-way +state transition: no new work may start, active attempts and waits close with +`completionReason: canceled`, and exactly one terminal canceled fact follows. + ### 1.2 `step.attempt.started` One per attempt. Payload: @@ -330,7 +336,7 @@ kernel/ │ └── registry.rs # relayflowd.sqlite3 run index (rebuildable) └── relayflowd/ # the binary └── src/ - ├── main.rs # CLI: run | resume | serve + ├── main.rs # CLI: run | resume | cancel | serve ├── engine.rs # drives machine.rs Actions against journal + executors ├── exec_det.rs # deterministic steps: spawn, capture, timeout ├── server.rs # unix socket, protocol v0 (§5), worker dispatch @@ -367,6 +373,7 @@ Minimal verb set for gate 1: | `hello` | `{protocol: 0, client}` → `{protocol: 0, server}` | handshake; version mismatch is a hard error | | `run.start` | `{spec}` → `{run_id}` | validate spec (zero-agent flows are legal), create run file, append `run.spawned`, begin scheduling | | `run.resume` | `{run_id}` → `{run_id, state}` | §3 memoized resume | +| `run.cancel` | `{run_id}` → `{run_id, status, completion_reason}` | append durable intent, close active leases, and append the terminal canceled fact; repeated calls return the existing outcome | | `run.get` | `{run_id}` → `{status, steps, budget}` | snapshot for legibility | | `run.watch` | `{run_id}` → stream of `{event: "entry", data: Entry}` | every appended entry, pushed | | `worker.attach` | `{worker_id, step_types: ["llm","agent"], pins}` → `{}` | connection becomes a worker; agent workers **must** supply opaque initial workspace revisions/stream offsets (refused otherwise) and receive `step.dispatch` events with pins plus recovery context. A step whose declared surfaces no attached worker holds parks — it is not dispatched | diff --git a/kernel/relayflowd-core/src/entry.rs b/kernel/relayflowd-core/src/entry.rs index 6cd6594f..f26fe716 100644 --- a/kernel/relayflowd-core/src/entry.rs +++ b/kernel/relayflowd-core/src/entry.rs @@ -9,6 +9,8 @@ use crate::spec::{RecoveryMode, StepType}; pub enum EntryType { #[serde(rename = "run.spawned")] RunSpawned, + #[serde(rename = "run.cancel.requested")] + RunCancelRequested, #[serde(rename = "event.received")] EventReceived, #[serde(rename = "subscription.registered")] @@ -53,6 +55,7 @@ impl EntryType { pub fn as_str(self) -> &'static str { match self { Self::RunSpawned => "run.spawned", + Self::RunCancelRequested => "run.cancel.requested", Self::EventReceived => "event.received", Self::SubscriptionRegistered => "subscription.registered", Self::SubscriptionMatched => "subscription.matched", @@ -75,6 +78,7 @@ impl EntryType { pub fn parse(value: &str) -> Option { Some(match value { "run.spawned" => Self::RunSpawned, + "run.cancel.requested" => Self::RunCancelRequested, "event.received" => Self::EventReceived, "subscription.registered" => Self::SubscriptionRegistered, "subscription.matched" => Self::SubscriptionMatched, @@ -140,6 +144,11 @@ pub struct RunSpawnedPayload { pub created_by: String, } +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +pub struct RunCancelRequestedPayload { + pub requested_by: String, +} + #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] pub struct AttemptStartedPayload { pub step_type: StepType, diff --git a/kernel/relayflowd-core/src/lib.rs b/kernel/relayflowd-core/src/lib.rs index 0b9f3e2b..006f54ef 100644 --- a/kernel/relayflowd-core/src/lib.rs +++ b/kernel/relayflowd-core/src/lib.rs @@ -21,6 +21,7 @@ pub use journal::{Journal, JournalError, MemoryJournal}; pub use machine::{ Action, AttemptResult, RecoveryInstruction, abandonment_actions, carried_pins_for, completion_actions, next_actions, recovery_actions, recovery_actions_filtered, + request_cancel_action, }; pub use spec::*; pub use state::{RunState, StateError, StepRuntime, StepState}; diff --git a/kernel/relayflowd-core/src/machine.rs b/kernel/relayflowd-core/src/machine.rs index 861eb411..6e97ca95 100644 --- a/kernel/relayflowd-core/src/machine.rs +++ b/kernel/relayflowd-core/src/machine.rs @@ -17,6 +17,10 @@ use crate::{ const LEASE_DURATION_MS: i64 = 30_000; +mod cancel; +use cancel::cancel_run_actions; +pub use cancel::request_cancel_action; + #[derive(Debug, Clone, PartialEq)] pub enum Action { Append(JournalEntry), @@ -87,6 +91,9 @@ pub fn next_actions(state: &RunState, now_ms: i64) -> Vec { if state.completion.is_some() { return Vec::new(); } + if state.cancel_requested.is_some() { + return cancel_run_actions(state, now_ms); + } if let Some(failed_step_id) = state.failed_step() { return complete_run_actions( state, diff --git a/kernel/relayflowd-core/src/machine/cancel.rs b/kernel/relayflowd-core/src/machine/cancel.rs new file mode 100644 index 00000000..d68a542a --- /dev/null +++ b/kernel/relayflowd-core/src/machine/cancel.rs @@ -0,0 +1,100 @@ +use serde_json::Value; + +use super::{Action, complete_run_actions, retry_wait_id}; +use crate::{ + Budget, CompletionReason, Disposition, EntryType, JournalEntry, RunCancelRequestedPayload, + RunCompletionReason, RunState, StepCompletedPayload, StepState, WaitCompletedPayload, + WaitCompletionReason, +}; + +/// Persist cancellation intent before any work is closed. Repeating the +/// request is a no-op both while cancellation is in progress and after the +/// terminal fact exists. +pub fn request_cancel_action( + state: &RunState, + requested_by: impl Into, + now_ms: i64, +) -> Option { + if state.completion.is_some() || state.cancel_requested.is_some() { + return None; + } + Some(Action::Append(JournalEntry::new( + EntryType::RunCancelRequested, + state.run_id.clone(), + None, + None, + now_ms, + RunCancelRequestedPayload { + requested_by: requested_by.into(), + }, + ))) +} + +pub(super) fn cancel_run_actions(state: &RunState, now_ms: i64) -> Vec { + let mut actions = Vec::new(); + for step in &state.spec.steps { + let runtime = &state.steps[&step.id]; + match &runtime.state { + StepState::Running { attempt, .. } => { + actions.push(Action::Append(JournalEntry::new( + EntryType::StepCompleted, + state.run_id.clone(), + Some(step.id.clone()), + Some(*attempt), + now_ms, + StepCompletedPayload { + completion_reason: CompletionReason::Canceled, + disposition: Disposition::StepDone, + output: Value::Null, + verification: None, + end_pins: None, + effects: Vec::new(), + trajectory_tail: None, + budget: Budget::default(), + completed_by: "kernel".to_owned(), + next_attempt_at_ms: None, + }, + ))); + } + StepState::Waiting { wait_id } | StepState::NeedsHuman { wait_id } => { + actions.push(Action::Append(JournalEntry::new( + EntryType::WaitCompleted, + state.run_id.clone(), + Some(step.id.clone()), + Some(runtime.attempts), + now_ms, + WaitCompletedPayload { + wait_id: wait_id.clone(), + completion_reason: WaitCompletionReason::Canceled, + result: Value::Null, + }, + ))); + } + StepState::Backoff { + attempt, + wake_at_ms, + } => { + actions.push(Action::Append(JournalEntry::new( + EntryType::WaitCompleted, + state.run_id.clone(), + Some(step.id.clone()), + Some(*attempt), + now_ms, + WaitCompletedPayload { + wait_id: retry_wait_id(&state.run_id, &step.id, *attempt, *wake_at_ms), + completion_reason: WaitCompletionReason::Canceled, + result: Value::Null, + }, + ))); + } + StepState::Pending | StepState::Runnable | StepState::Done { .. } => {} + } + } + actions.extend(complete_run_actions( + state, + RunCompletionReason::Canceled, + None, + now_ms, + )); + actions +} diff --git a/kernel/relayflowd-core/src/machine/recovery.rs b/kernel/relayflowd-core/src/machine/recovery.rs index a84a8bb5..85e697ba 100644 --- a/kernel/relayflowd-core/src/machine/recovery.rs +++ b/kernel/relayflowd-core/src/machine/recovery.rs @@ -27,6 +27,12 @@ pub fn recovery_actions_filtered( now_ms: i64, lease_is_active: &dyn Fn(&str, u32) -> bool, ) -> Vec { + // Durable cancellation intent wins over crash classification. The cancel + // path closes the same lease as canceled; recovery must not get there + // first and rewrite the reason merely because the process restarted. + if state.cancel_requested.is_some() { + return Vec::new(); + } let mut actions = Vec::new(); for spec in &state.spec.steps { let runtime = &state.steps[&spec.id]; diff --git a/kernel/relayflowd-core/src/machine/tests.rs b/kernel/relayflowd-core/src/machine/tests.rs index e9b6987f..31aa0e3b 100644 --- a/kernel/relayflowd-core/src/machine/tests.rs +++ b/kernel/relayflowd-core/src/machine/tests.rs @@ -1,7 +1,7 @@ use serde_json::json; use super::*; -use crate::{entry::AttemptStartedPayload, state::RunState}; +use crate::{Clock, SimClock, entry::AttemptStartedPayload, state::RunState}; fn retrying_spec() -> crate::RunSpec { serde_json::from_value(json!({ @@ -133,6 +133,91 @@ fn successful_memo_is_never_scheduled_again() { ); } +#[test] +fn cancel_request_closes_the_active_lease_before_the_terminal_fact() { + let clock = SimClock::new(10); + let spec = crate::RunSpec::parse(&json!({ + "steps": [{"id": "model", "type": "llm", "prompt": "answer"}] + })) + .unwrap(); + let fresh = RunState::fold("run", spec.clone(), &[]).unwrap(); + let Action::Append(started) = next_actions(&fresh, clock.now_ms()).remove(0) else { + panic!("the attempt lease must be durable"); + }; + let running = RunState::fold("run", spec.clone(), std::slice::from_ref(&started)).unwrap(); + clock.advance(10); + let Action::Append(requested) = + request_cancel_action(&running, "operator", clock.now_ms()).unwrap() + else { + panic!("cancel must first persist its request"); + }; + assert_eq!(requested.entry_type, EntryType::RunCancelRequested); + + let canceling = RunState::fold("run", spec, &[started, requested]).unwrap(); + clock.advance(10); + let actions = next_actions(&canceling, clock.now_ms()); + let Action::Append(closed) = &actions[0] else { + panic!("the active attempt must be closed first"); + }; + let closed: StepCompletedPayload = serde_json::from_value(closed.payload.clone()).unwrap(); + assert_eq!(closed.completion_reason, CompletionReason::Canceled); + assert_eq!(closed.disposition, Disposition::StepDone); + let Action::Append(terminal) = &actions[1] else { + panic!("the run fact must follow the lease closure"); + }; + assert_eq!(terminal.entry_type, EntryType::RunCompleted); + let terminal: RunCompletedPayload = serde_json::from_value(terminal.payload.clone()).unwrap(); + assert_eq!(terminal.completion_reason, RunCompletionReason::Canceled); +} + +#[test] +fn repeated_cancel_request_is_idempotent() { + let spec = retrying_spec(); + let state = RunState::fold("run", spec.clone(), &[]).unwrap(); + let Action::Append(requested) = request_cancel_action(&state, "operator", 10).unwrap() else { + panic!(); + }; + let canceling = RunState::fold("run", spec.clone(), std::slice::from_ref(&requested)).unwrap(); + assert!(request_cancel_action(&canceling, "operator", 11).is_none()); + let terminal_entries = next_actions(&canceling, 12) + .into_iter() + .filter_map(|action| match action { + Action::Append(entry) => Some(entry), + _ => None, + }) + .collect::>(); + let terminal = RunState::fold("run", spec, &[requested, terminal_entries[0].clone()]).unwrap(); + assert!(request_cancel_action(&terminal, "operator", 13).is_none()); +} + +#[test] +fn durable_cancel_request_outranks_crash_recovery() { + let clock = SimClock::new(10); + let spec = crate::RunSpec::parse(&json!({ + "steps": [{"id": "model", "type": "llm", "prompt": "answer"}] + })) + .unwrap(); + let fresh = RunState::fold("run", spec.clone(), &[]).unwrap(); + let Action::Append(started) = next_actions(&fresh, clock.now_ms()).remove(0) else { + panic!("the attempt lease must be durable"); + }; + let running = RunState::fold("run", spec.clone(), std::slice::from_ref(&started)).unwrap(); + clock.advance(10); + let Action::Append(requested) = + request_cancel_action(&running, "operator", clock.now_ms()).unwrap() + else { + panic!("the cancel request must be durable"); + }; + let canceling = RunState::fold("run", spec, &[started, requested]).unwrap(); + + assert!(recovery_actions(&canceling, clock.now_ms()).is_empty()); + let Action::Append(closed) = &next_actions(&canceling, clock.now_ms())[0] else { + panic!("cancellation must close the active lease"); + }; + let closed: StepCompletedPayload = serde_json::from_value(closed.payload.clone()).unwrap(); + assert_eq!(closed.completion_reason, CompletionReason::Canceled); +} + #[test] fn crashed_attempt_does_not_consume_an_iteration() { // max_iterations 2: crash attempt 1, verification-fail the replacement @@ -146,7 +231,7 @@ fn crashed_attempt_does_not_consume_an_iteration() { // kill -9 between steps: attempt 1 is Running with no result. Recovery // must record the dead attempt as a retry, not a consumed iteration. - let state = RunState::fold("run", spec.clone(), &[started.clone()]).unwrap(); + let state = RunState::fold("run", spec.clone(), std::slice::from_ref(&started)).unwrap(); assert_eq!(state.steps["hello"].semantic_executions, 0); let recovery = recovery_actions(&state, 1_000); let Action::Append(crashed) = &recovery[0] else { @@ -276,7 +361,7 @@ fn reset_recovery_dispatches_the_original_pinned_revision() { let spec = agent_spec("reset"); let pinned = workspace_pins("rev-clean"); let started = started_agent(&spec, pinned.clone()); - let running = RunState::fold("run", spec.clone(), &[started.clone()]).unwrap(); + let running = RunState::fold("run", spec.clone(), std::slice::from_ref(&started)).unwrap(); let recovered = recovery_actions(&running, 20); let entries = vec![ started, @@ -311,7 +396,7 @@ fn inspect_recovery_injects_the_dirty_pin_completion_reason_and_tail() { let clean = workspace_pins("rev-clean"); let dirty = workspace_pins("rev-dirty"); let started = started_agent(&spec, clean); - let running = RunState::fold("run", spec.clone(), &[started.clone()]).unwrap(); + let running = RunState::fold("run", spec.clone(), std::slice::from_ref(&started)).unwrap(); let result = AttemptResult { output: Value::Null, budget: Budget::default(), @@ -359,7 +444,7 @@ fn inspect_recovery_injects_the_dirty_pin_completion_reason_and_tail() { fn manual_recovery_parks_needs_human_and_never_redispatches() { let spec = agent_spec("manual"); let started = started_agent(&spec, workspace_pins("rev-clean")); - let running = RunState::fold("run", spec.clone(), &[started.clone()]).unwrap(); + let running = RunState::fold("run", spec.clone(), std::slice::from_ref(&started)).unwrap(); let recovered = recovery_actions(&running, 20); let Action::Append(wait) = &recovered[1] else { panic!( diff --git a/kernel/relayflowd-core/src/state.rs b/kernel/relayflowd-core/src/state.rs index c26f424f..3409909f 100644 --- a/kernel/relayflowd-core/src/state.rs +++ b/kernel/relayflowd-core/src/state.rs @@ -6,8 +6,8 @@ use thiserror::Error; use crate::{ entry::{ Budget, CompletionReason, Disposition, EntryType, EpochSummaryPayload, JournalEntry, Pins, - RunCompletedPayload, RunCompletionReason, SleepUntilPayload, StepCompletedPayload, - WaitCompletedPayload, WaitCompletionReason, + RunCancelRequestedPayload, RunCompletedPayload, RunCompletionReason, SleepUntilPayload, + StepCompletedPayload, WaitCompletedPayload, WaitCompletionReason, }, spec::{RunSpec, StepKind, StepType}, }; @@ -67,6 +67,9 @@ pub struct RunState { pub memo: BTreeMap, pub budget: Budget, pub completion: Option, + /// Durable cancellation intent. Once present, scheduling can only close + /// live work and append the terminal canceled fact. + pub cancel_requested: Option, /// Appendix A rule 6 chain head: the last successful agent completion. pub current_pins: Option, } @@ -103,6 +106,7 @@ impl RunState { memo: BTreeMap::new(), budget: Budget::default(), completion: None, + cancel_requested: None, current_pins: None, }; @@ -164,6 +168,9 @@ impl RunState { let payload: RunCompletedPayload = decode(entry)?; state.completion = Some(payload.completion_reason); } + EntryType::RunCancelRequested => { + state.cancel_requested = Some(decode(entry)?); + } EntryType::RunSpawned | EntryType::EventReceived | EntryType::SubscriptionRegistered diff --git a/kernel/relayflowd/src/engine.rs b/kernel/relayflowd/src/engine.rs index c1f7ec0a..7484235a 100644 --- a/kernel/relayflowd/src/engine.rs +++ b/kernel/relayflowd/src/engine.rs @@ -6,7 +6,7 @@ use std::{ use anyhow::{Context, Result, anyhow, bail}; use relayflowd_core::{ Clock, EntryType, Journal, JournalEntry, RunSpawnedPayload, RunSpec, RunState, StepKind, - recovery_actions_filtered, + recovery_actions_filtered, request_cancel_action, }; use relayflowd_journal::{Registry, SqliteJournal}; use sha2::{Digest, Sha256}; @@ -35,6 +35,14 @@ pub struct DriveOptions { pub pause_before_completion: bool, } +#[derive(Debug, Clone, Default)] +#[doc(hidden)] +pub struct CancelOptions { + /// Crash-injection seam: pause after intent is durable but before facts + /// close the run. Production callers leave this false. + pub pause_after_request: bool, +} + pub struct Engine { data_dir: PathBuf, clock: C, @@ -186,6 +194,33 @@ impl Engine { Ok(snapshot) } + pub fn cancel(&self, run_id: &str, requested_by: &str) -> Result { + self.cancel_with_options(run_id, requested_by, CancelOptions::default()) + } + + #[doc(hidden)] + pub fn cancel_with_options( + &self, + run_id: &str, + requested_by: &str, + options: CancelOptions, + ) -> Result { + let mut journal = self.open_run(run_id)?; + let spec = journal.run_spec().context("read run spec")?; + let state = self.load_state(&journal, spec.clone())?; + if let Some(reason) = state.completion { + return Ok(outcome_from_state(&state, reason)); + } + let requested = request_cancel_action(&state, requested_by, self.clock.now_ms()); + if let Some(action) = requested { + self.persist_only(&mut journal, action)?; + if options.pause_after_request { + std::thread::sleep(std::time::Duration::from_secs(300)); + } + } + self.drive(journal, spec, DriveOptions::default()) + } + pub fn journal_entries( &self, run_id: &str, diff --git a/kernel/relayflowd/src/engine/remote.rs b/kernel/relayflowd/src/engine/remote.rs index 38af96cf..9b7ca635 100644 --- a/kernel/relayflowd/src/engine/remote.rs +++ b/kernel/relayflowd/src/engine/remote.rs @@ -41,6 +41,9 @@ impl Engine { let mut journal = self.open_run(run_id)?; let spec = journal.run_spec().context("read run spec")?; let state = self.load_state(&journal, spec.clone())?; + if state.cancel_requested.is_some() || state.completion.is_some() { + bail!("run {run_id} no longer accepts step completions") + } let step = spec .step(step_id) .cloned() diff --git a/kernel/relayflowd/src/lib.rs b/kernel/relayflowd/src/lib.rs index a90a98c3..014bd359 100644 --- a/kernel/relayflowd/src/lib.rs +++ b/kernel/relayflowd/src/lib.rs @@ -5,6 +5,6 @@ pub mod server; pub mod worker; pub use engine::{ - DriveOptions, Engine, OutOfBandCompletion, RunOutcome, RunSnapshot, RunStatus, + CancelOptions, DriveOptions, Engine, OutOfBandCompletion, RunOutcome, RunSnapshot, RunStatus, StepSnapshot, StepStatus, }; diff --git a/kernel/relayflowd/src/main.rs b/kernel/relayflowd/src/main.rs index 9ec51d1b..fc427474 100644 --- a/kernel/relayflowd/src/main.rs +++ b/kernel/relayflowd/src/main.rs @@ -2,7 +2,7 @@ use std::path::PathBuf; use anyhow::{Result, bail}; use clap::{Parser, Subcommand}; -use relayflowd::{DriveOptions, Engine, RunStatus, engine::read_spec, server}; +use relayflowd::{CancelOptions, DriveOptions, Engine, RunStatus, engine::read_spec, server}; #[derive(Debug, Parser)] #[command( @@ -43,6 +43,13 @@ enum Command { #[arg(long, hide = true)] stop_after: Option, }, + /// Durably cancel a run and close any active lease. + Cancel { + run_id: String, + /// Test/debug boundary: pause after intent is durable, before facts. + #[arg(long, hide = true)] + pause_after_request: bool, + }, /// Serve journal protocol v0 over a Unix socket. Serve, } @@ -97,6 +104,27 @@ fn main() -> Result<()> { bail!("run {} failed", outcome.run_id); } } + Command::Cancel { + run_id, + pause_after_request, + } => { + let outcome = if !pause_after_request { + if let Some(outcome) = server::cancel_via_socket(&cli.data_dir, &run_id)? { + outcome + } else { + engine.cancel(&run_id, "cli")? + } + } else { + engine.cancel_with_options( + &run_id, + "cli", + CancelOptions { + pause_after_request: true, + }, + )? + }; + println!("{}", serde_json::to_string(&outcome)?); + } Command::Serve => server::serve(&cli.data_dir)?, } Ok(()) diff --git a/kernel/relayflowd/src/server.rs b/kernel/relayflowd/src/server.rs index 115b171e..491463bc 100644 --- a/kernel/relayflowd/src/server.rs +++ b/kernel/relayflowd/src/server.rs @@ -7,6 +7,8 @@ use serde_json::{Value, json}; use crate::{Engine, OutOfBandCompletion}; +#[cfg(unix)] +mod cancel; #[cfg(unix)] pub mod liveness; #[cfg(unix)] @@ -20,7 +22,7 @@ use session::{ProtocolHub, SharedWriter, write_frame}; /// human and for the next attempt's prompt, not a transcript store. const TRAJECTORY_TAIL_MAX_BYTES: usize = 16 * 1024; mod client; -pub use client::resume_via_socket; +pub use client::{cancel_via_socket, resume_via_socket}; mod wire; use wire::*; @@ -191,11 +193,17 @@ fn handle_request( let _guard = lock.lock().expect("run lock"); // Live resume: attempts with a valid, heartbeating lease on this // hub stay running; only genuinely dead attempts are recovered. - to_value( - engine - .resume_live(¶ms.run_id, hub.as_ref()) - .map_err(internal_error)?, - ) + let outcome = engine + .resume_live(¶ms.run_id, hub.as_ref()) + .map_err(internal_error)?; + if outcome.completion_reason.is_some() { + hub.finish_run(¶ms.run_id); + } + to_value(outcome) + } + "run.cancel" => { + let params: RunIdParams = decode_params(request.params)?; + cancel::handle(data_dir, hub, &engine, ¶ms.run_id) } "run.get" => { let params: RunIdParams = decode_params(request.params)?; diff --git a/kernel/relayflowd/src/server/cancel.rs b/kernel/relayflowd/src/server/cancel.rs new file mode 100644 index 00000000..b86ba9c4 --- /dev/null +++ b/kernel/relayflowd/src/server/cancel.rs @@ -0,0 +1,35 @@ +use std::{path::Path, sync::Arc}; + +use serde_json::Value; + +use super::{ProtocolHub, ProtocolResult, internal_error, to_value}; +use crate::Engine; + +pub(super) fn handle( + data_dir: &Path, + hub: &Arc, + engine: &Engine, + run_id: &str, +) -> ProtocolResult { + let registry = relayflowd_journal::Registry::open(data_dir.join("relayflowd.sqlite3")) + .map_err(|error| internal_error(error.into()))?; + if registry + .lookup(run_id) + .map_err(|error| internal_error(error.into()))? + .is_none() + { + return Err(("run_not_found", format!("run {run_id} does not exist"))); + } + + // The same lock guards completion: whichever request acquires it first + // becomes the one linear history represented by the journal. + let lock = hub.run_lock(run_id); + let _guard = lock.lock().expect("run lock"); + let outcome = engine + .cancel(run_id, "protocol-v0") + .map_err(internal_error)?; + if outcome.completion_reason.is_some() { + hub.finish_run(run_id); + } + to_value(outcome) +} diff --git a/kernel/relayflowd/src/server/client.rs b/kernel/relayflowd/src/server/client.rs index e880fed4..b8ff5a84 100644 --- a/kernel/relayflowd/src/server/client.rs +++ b/kernel/relayflowd/src/server/client.rs @@ -11,6 +11,57 @@ use serde_json::json; use super::Response; use crate::Engine; +#[cfg(unix)] +fn lifecycle_request_via_socket( + data_dir: &Path, + request_id: &str, + verb: &str, + run_id: &str, +) -> Result> { + use std::{ + io::{BufRead, BufReader, Write}, + os::unix::net::UnixStream, + }; + let socket = data_dir.join("relayflowd.sock"); + if !socket.exists() { + return Ok(None); + } + let mut connection = match UnixStream::connect(&socket) { + Ok(connection) => connection, + Err(error) + if matches!( + error.kind(), + io::ErrorKind::NotFound | io::ErrorKind::ConnectionRefused + ) => + { + return Ok(None); + } + Err(error) => return Err(error.into()), + }; + serde_json::to_writer( + &mut connection, + &json!({"id": request_id, "verb": verb, "params": {"run_id": run_id}}), + )?; + connection.write_all(b"\n")?; + connection.flush()?; + for line in BufReader::new(connection).lines() { + let response: Response = serde_json::from_str(&line?)?; + if response.id != request_id { + continue; + } + if !response.ok { + let error = response.error.context("protocol error omitted detail")?; + anyhow::bail!("{}: {}", error.code, error.message); + } + return Ok(Some(serde_json::from_value( + response + .result + .context("protocol response omitted result")?, + )?)); + } + anyhow::bail!("relayflowd serve closed before {verb} replied") +} + /// How long after a lease deadline the CLI still waits before calling the /// attempt dead. The reconciler sweeps expired leases and journals their /// completion; this covers the gap between the deadline passing and that sweep @@ -104,6 +155,11 @@ pub fn resume_via_socket(data_dir: &Path, run_id: &str) -> Result Result> { + lifecycle_request_via_socket(data_dir, "cli-cancel", "run.cancel", run_id) +} + /// Wait out an attempt a worker is running out of band. The loop follows the /// lease: while the worker keeps renewing it the attempt is alive however long /// it takes, and the wait ends when the run reaches a terminal state, leaves @@ -138,6 +194,11 @@ pub fn resume_via_socket(_data_dir: &Path, _run_id: &str) -> Result Result> { + Ok(None) +} + #[cfg(test)] mod tests { use super::*; diff --git a/kernel/relayflowd/src/server/session.rs b/kernel/relayflowd/src/server/session.rs index e72823ad..58d440bd 100644 --- a/kernel/relayflowd/src/server/session.rs +++ b/kernel/relayflowd/src/server/session.rs @@ -247,6 +247,16 @@ impl ProtocolHub { .remove(key); } + /// Release every in-memory lease after the journal has durably made the + /// run terminal. Calling this repeatedly is intentionally harmless. + pub fn finish_run(&self, run_id: &str) { + self.sessions + .lock() + .expect("protocol sessions lock") + .assignments + .retain(|(assigned_run, _, _), _| assigned_run != run_id); + } + /// Assignments whose (heartbeat-renewed) lease deadline has passed. The /// worker may still hold an open socket — a hung worker is exactly the /// case the expiry reconciler exists for. Assignments are NOT removed diff --git a/kernel/relayflowd/tests/crash_resume.rs b/kernel/relayflowd/tests/crash_resume.rs index f1b40700..e956b1ad 100644 --- a/kernel/relayflowd/tests/crash_resume.rs +++ b/kernel/relayflowd/tests/crash_resume.rs @@ -18,7 +18,10 @@ use std::{ fs, io::Write, os::unix::net::UnixStream, os::unix::process::CommandExt, process::Command, }; -use relayflowd_core::{CompletionReason, EntryType, StepCompletedPayload}; +use relayflowd::{RunOutcome, RunStatus}; +use relayflowd_core::{ + CompletionReason, EntryType, RunCompletedPayload, RunCompletionReason, StepCompletedPayload, +}; use serde_json::{Value, json}; use support::{ @@ -147,3 +150,98 @@ fn sigkill_under_serve_resumes_the_socket_started_run() { CompletionReason::Crashed | CompletionReason::LeaseExpired )); } + +#[test] +fn sigkill_after_cancel_request_resumes_to_one_canceled_fact() { + let fixture = Fixture::hello("cancel-request"); + let interrupted = Command::new(env!("CARGO_BIN_EXE_relayflowd")) + .args(["--data-dir", fixture.data_dir.to_str().unwrap(), "run"]) + .arg(&fixture.spec_path) + .args(["--stop-after", "1"]) + .output() + .unwrap(); + assert!( + interrupted.status.success(), + "initial run failed: {interrupted:?}" + ); + let run_id = support::only_run_id(&fixture.data_dir); + + let mut cancel = Command::new(env!("CARGO_BIN_EXE_relayflowd")) + .args([ + "--data-dir", + fixture.data_dir.to_str().unwrap(), + "cancel", + &run_id, + "--pause-after-request", + ]) + .process_group(0) + .spawn() + .unwrap(); + wait_until("durable cancel request", || { + journal_entries(&fixture.data_dir).is_some_and(|entries| { + entries + .iter() + .any(|entry| entry.entry_type == EntryType::RunCancelRequested) + }) + }); + kill_process_group(&mut cancel); + let before_resume = journal_entries(&fixture.data_dir).unwrap(); + assert_eq!( + before_resume + .iter() + .filter(|entry| entry.entry_type == EntryType::RunCancelRequested) + .count(), + 1 + ); + assert!( + !before_resume + .iter() + .any(|entry| entry.entry_type == EntryType::RunCompleted) + ); + + let resumed = Command::new(env!("CARGO_BIN_EXE_relayflowd")) + .args([ + "--data-dir", + fixture.data_dir.to_str().unwrap(), + "resume", + &run_id, + ]) + .output() + .unwrap(); + let outcome: RunOutcome = serde_json::from_slice(&resumed.stdout).unwrap(); + assert_eq!(outcome.status, RunStatus::Failed); + assert_eq!( + outcome.completion_reason, + Some(RunCompletionReason::Canceled) + ); + + let repeated = Command::new(env!("CARGO_BIN_EXE_relayflowd")) + .args([ + "--data-dir", + fixture.data_dir.to_str().unwrap(), + "cancel", + &run_id, + ]) + .output() + .unwrap(); + assert!( + repeated.status.success(), + "repeated cancel failed: {repeated:?}" + ); + let entries = journal_entries(&fixture.data_dir).unwrap(); + assert_eq!( + entries + .iter() + .filter(|entry| entry.entry_type == EntryType::RunCancelRequested) + .count(), + 1 + ); + let terminal = entries + .iter() + .filter(|entry| entry.entry_type == EntryType::RunCompleted) + .map(|entry| serde_json::from_value::(entry.payload.clone()).unwrap()) + .collect::>(); + assert_eq!(terminal.len(), 1); + assert_eq!(terminal[0].completion_reason, RunCompletionReason::Canceled); + assert_eq!(fs::read_to_string(&fixture.marker).unwrap(), "first\n"); +} diff --git a/kernel/relayflowd/tests/crash_resume/concurrency.rs b/kernel/relayflowd/tests/crash_resume/concurrency.rs index 5393e3e3..6a14b586 100644 --- a/kernel/relayflowd/tests/crash_resume/concurrency.rs +++ b/kernel/relayflowd/tests/crash_resume/concurrency.rs @@ -7,7 +7,9 @@ use std::{ thread, }; -use relayflowd_core::{CompletionReason, EntryType, StepCompletedPayload}; +use relayflowd_core::{ + CompletionReason, EntryType, RunCompletedPayload, RunCompletionReason, StepCompletedPayload, +}; use serde_json::json; use super::{ @@ -124,6 +126,83 @@ fn live_resume_leaves_an_active_lease_running() { assert_eq!(done["status"], "completed"); } +#[test] +fn cancel_closes_the_lease_and_rejects_a_late_completion() { + let fixture = LlmFixture::new("cancel-late-completion", false); + let _server = ServerGuard::start(&fixture); + let mut worker = attached_worker(&fixture, "cancel-stub"); + let run_id = start_run(&fixture); + let dispatch = worker.event("step.dispatch").unwrap(); + let mut control = ProtocolClient::connect(&fixture.data_dir.join("relayflowd.sock")); + + let canceled = control + .request("run.cancel", json!({"run_id": run_id})) + .unwrap(); + assert_eq!(canceled["completion_reason"], "canceled"); + let repeated = control + .request("run.cancel", json!({"run_id": run_id})) + .unwrap(); + assert_eq!(repeated, canceled); + assert!(complete(&mut worker, &dispatch, json!({"answer": 4})).is_err()); + + let entries = journal_entries(&fixture.data_dir).unwrap(); + assert_eq!( + entries + .iter() + .filter(|entry| entry.entry_type == EntryType::RunCancelRequested) + .count(), + 1 + ); + let completions = model_completions(&fixture); + assert_eq!(completions.len(), 1); + assert_eq!(completions[0].completion_reason, CompletionReason::Canceled); +} + +#[test] +fn cancel_and_completion_race_has_one_terminal_fact() { + let fixture = LlmFixture::new("cancel-completion-race", false); + let _server = ServerGuard::start(&fixture); + let mut worker = attached_worker(&fixture, "race-stub"); + let run_id = start_run(&fixture); + let dispatch = worker.event("step.dispatch").unwrap(); + let socket = fixture.data_dir.join("relayflowd.sock"); + let barrier = Arc::new(Barrier::new(2)); + + let cancel_barrier = barrier.clone(); + let cancel_run = run_id.clone(); + let cancel = thread::spawn(move || { + let mut control = ProtocolClient::connect(&socket); + cancel_barrier.wait(); + control.request("run.cancel", json!({"run_id": cancel_run})) + }); + let complete_barrier = barrier.clone(); + let completion = thread::spawn(move || { + complete_barrier.wait(); + complete(&mut worker, &dispatch, json!({"answer": 4})) + }); + let cancel = cancel.join().unwrap(); + let completion = completion.join().unwrap(); + + let entries = journal_entries(&fixture.data_dir).unwrap(); + let terminal = entries + .iter() + .filter(|entry| entry.entry_type == EntryType::RunCompleted) + .map(|entry| serde_json::from_value::(entry.payload.clone()).unwrap()) + .collect::>(); + assert_eq!(terminal.len(), 1); + match terminal[0].completion_reason { + RunCompletionReason::Canceled => { + assert!(cancel.is_ok()); + assert!(completion.is_err()); + } + RunCompletionReason::Success => { + assert!(completion.is_ok()); + assert_eq!(cancel.unwrap()["completion_reason"], "success"); + } + reason => panic!("unexpected race result: {reason:?}"), + } +} + fn model_completions(fixture: &LlmFixture) -> Vec { journal_entries(&fixture.data_dir) .unwrap_or_default() diff --git a/ops/reviews/20260902-v2-run-cancel-evidence.md b/ops/reviews/20260902-v2-run-cancel-evidence.md new file mode 100644 index 00000000..caae8842 --- /dev/null +++ b/ops/reviews/20260902-v2-run-cancel-evidence.md @@ -0,0 +1,268 @@ +# V2 run cancellation — captured evidence + +Scope: `feat/v2-run-cancel`. Commands ran from this worktree on 2026-09-02. +Rust commands ran from `kernel/`; SDK commands ran from `sdk/`. + +## Mutation check: cancellation priority removed + +The specific three-line `cancel_requested` priority branch in +`relayflowd-core/src/machine.rs` was removed with `apply_patch`, the focused +core and real-process/socket tests were run, and the branch was then restored +byte-for-byte with `apply_patch`. + +```text +$ CARGO_TARGET_DIR=/tmp/flows-v2-cancel-target cargo test -p relayflowd-core cancel_request -- --nocapture + +running 2 tests +test machine::tests::repeated_cancel_request_is_idempotent ... ok +test machine::tests::cancel_request_closes_the_active_lease_before_the_terminal_fact ... FAILED + +thread 'machine::tests::cancel_request_closes_the_active_lease_before_the_terminal_fact' panicked at relayflowd-core/src/machine/tests.rs:159:42: +index out of bounds: the len is 0 but the index is 0 + +failures: + machine::tests::cancel_request_closes_the_active_lease_before_the_terminal_fact + +test result: FAILED. 1 passed; 1 failed; 0 ignored; 0 measured; 26 filtered out +MUTATION_CORE_EXIT=101 + +$ CARGO_TARGET_DIR=/tmp/flows-v2-cancel-target cargo test -p relayflowd cancel_ -- --nocapture + +running 3 tests + +thread 'concurrency::cancel_closes_the_lease_and_rejects_a_late_completion' panicked at relayflowd/tests/crash_resume/concurrency.rs:141:5: +assertion `left == right` failed + left: Null + right: "canceled" +test concurrency::cancel_closes_the_lease_and_rejects_a_late_completion ... FAILED +test concurrency::cancel_and_completion_race_has_one_terminal_fact ... ok + +thread 'sigkill_after_cancel_request_resumes_to_one_canceled_fact' panicked at relayflowd/tests/crash_resume.rs:212:5: +assertion `left == right` failed + left: Completed + right: Failed +test sigkill_after_cancel_request_resumes_to_one_canceled_fact ... FAILED + +failures: + concurrency::cancel_closes_the_lease_and_rejects_a_late_completion + sigkill_after_cancel_request_resumes_to_one_canceled_fact + +test result: FAILED. 1 passed; 2 failed; 0 ignored; 0 measured; 19 filtered out +MUTATION_LIVE_EXIT=101 +``` + +## Restored focused pass + +```text +$ CARGO_TARGET_DIR=/tmp/flows-v2-cancel-target cargo test -p relayflowd-core cancel_request -- --nocapture + +running 2 tests +test machine::tests::repeated_cancel_request_is_idempotent ... ok +test machine::tests::cancel_request_closes_the_active_lease_before_the_terminal_fact ... ok + +test result: ok. 2 passed; 0 failed; 0 ignored; 0 measured; 26 filtered out + +running 0 tests + +test result: ok. 0 passed; 0 failed; 0 ignored; 0 measured; 5 filtered out +RESTORED_CORE_EXIT=0 + +$ CARGO_TARGET_DIR=/tmp/flows-v2-cancel-target cargo test -p relayflowd cancel_ -- --nocapture + +running 3 tests +test concurrency::cancel_closes_the_lease_and_rejects_a_late_completion ... ok +test sigkill_after_cancel_request_resumes_to_one_canceled_fact ... ok +test concurrency::cancel_and_completion_race_has_one_terminal_fact ... ok + +test result: ok. 3 passed; 0 failed; 0 ignored; 0 measured; 19 filtered out +RESTORED_LIVE_EXIT=0 +``` + +The final focused kernel pass also pins that a durable cancel request outranks +crash recovery for an active lease: + +```text +$ CARGO_TARGET_DIR=/tmp/flows-v2-cancel-target cargo test -p relayflowd-core cancel -- --nocapture + Finished `test` profile [unoptimized + debuginfo] target(s) in 0.60s + Running unittests src/lib.rs (/tmp/flows-v2-cancel-target/debug/deps/relayflowd_core-eb215c870b5b8c71) + +running 3 tests +test machine::tests::repeated_cancel_request_is_idempotent ... ok +test machine::tests::cancel_request_closes_the_active_lease_before_the_terminal_fact ... ok +test machine::tests::durable_cancel_request_outranks_crash_recovery ... ok + +test result: ok. 3 passed; 0 failed; 0 ignored; 0 measured; 26 filtered out; finished in 0.00s + + Running tests/spec_parity.rs (/tmp/flows-v2-cancel-target/debug/deps/spec_parity-25916cc116a7848d) + +running 0 tests + +test result: ok. 0 passed; 0 failed; 0 ignored; 0 measured; 5 filtered out; finished in 0.00s + +FINAL_CORE_CANCEL_EXIT=0 +``` + +## Full Rust regression and lint + +```text +$ CARGO_TARGET_DIR=/tmp/flows-v2-cancel-target cargo clippy --workspace --all-targets -- -D warnings && CARGO_TARGET_DIR=/tmp/flows-v2-cancel-target cargo test --workspace --quiet + Checking relayflowd-core v0.1.0 (/Users/khaliqgant/AgentWorkforce/flows-v2-cancel-wt/kernel/relayflowd-core) + Finished `dev` profile [unoptimized + debuginfo] target(s) in 56.94s + +running 22 tests +...................... +test result: ok. 22 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 0 tests +test result: ok. 0 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 22 tests +...................... +test result: ok. 22 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 1 test +. +test result: ok. 1 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 1 test +. +test result: ok. 1 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 3 tests +... +test result: ok. 3 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 29 tests +............................. +test result: ok. 29 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 5 tests +..... +test result: ok. 5 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 17 tests +................. +test result: ok. 17 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 0 tests +test result: ok. 0 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 0 tests +test result: ok. 0 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out + +running 0 tests +test result: ok. 0 passed; 0 failed; 0 ignored; 0 measured; 0 filtered out +``` + +## SDK typecheck/build and full regression + +`npm` itself hung before printing its version on this host. The worktree had no +dependency directory, so the command used the already-installed dependency +tree from the sibling `flows-132-direct-input-wt` worktree (temporarily linked +as `sdk/node_modules`, then unlinked). It ran this worktree's compiler inputs, +config, source, built CLI, tests, and freshly-built cancellation-capable daemon. + +```text +$ ./node_modules/.bin/tsc --noEmit && ./node_modules/.bin/tsc && node scripts/make-cli-executable.mjs && RELAYFLOWD_BIN=/tmp/flows-v2-cancel-target/debug/relayflowd ./node_modules/.bin/vitest run + + RUN v2.1.9 /Users/khaliqgant/AgentWorkforce/flows-v2-cancel-wt/sdk + + ✓ tests/preflight.test.ts (14 tests) 31ms + ✓ tests/validate.test.ts (36 tests) 70ms + ✓ tests/backlog-picker.test.ts (14 tests) 317ms + ✓ tests/journal-client.test.ts (14 tests) 172ms + ✓ tests/cli-hn-monitor.test.ts (16 tests) 179ms + ✓ tests/work-package-consumer.test.ts (13 tests) 963ms + ✓ tests/hn-poller.test.ts (6 tests) 55ms + ✓ tests/deterministic-llm.test.ts (5 tests) 137ms + ✓ tests/dir-watcher-poller.test.ts (6 tests) 40ms + ✓ tests/hello-deterministic.test.ts (5 tests) 164ms + ✓ tests/work-package-validator.test.ts (7 tests) 60ms + ✓ tests/backlog-picker-flow.test.ts (6 tests) 2582ms + ✓ tests/parse-json-output.test.ts (7 tests) 12ms + ✓ tests/cli.test.ts (50 tests) 2522ms + ✓ tests/spec-parity.test.ts (15 tests) 180ms + ✓ tests/bin.test.ts (7 tests) 2404ms + ✓ tests/live-kernel.test.ts (18 tests) 56671ms + + Test Files 17 passed (17) + Tests 239 passed (239) + Duration 58.63s (transform 1.69s, setup 0ms, collect 5.49s, tests 66.56s, environment 9ms, prepare 6.70s) +``` + +After the final kernel ordering change, the full suite was rerun with file +parallelism disabled. This avoids exhausting the existing five-second timeout +in an unrelated live CLI case under parallel load. The analyzer skip variable +is the repository's documented allowance for a non-gate run; in this captured +run the real analyzer was available and executed successfully anyway. + +```text +$ RELAYFLOWS_ALLOW_ANALYZER_SKIP=1 RELAYFLOWD_BIN=/tmp/flows-v2-cancel-target/debug/relayflowd ./node_modules/.bin/vitest run --no-file-parallelism + + RUN v2.1.9 /Users/khaliqgant/AgentWorkforce/flows-v2-cancel-wt/sdk + + ✓ tests/live-kernel.test.ts (18 tests) 68166ms + ✓ tests/cli.test.ts (50 tests) 1985ms + ✓ tests/journal-client.test.ts (14 tests) 167ms + ✓ tests/validate.test.ts (36 tests) 79ms + ✓ tests/cli-hn-monitor.test.ts (16 tests) 93ms + ✓ tests/preflight.test.ts (14 tests) 27ms + ✓ tests/backlog-picker.test.ts (14 tests) 628ms + ✓ tests/backlog-picker-flow.test.ts (6 tests) 1824ms + ✓ tests/work-package-consumer.test.ts (13 tests) 724ms + ✓ tests/deterministic-llm.test.ts (5 tests) 46ms + ✓ tests/bin.test.ts (7 tests) 1267ms + ✓ tests/hn-poller.test.ts (6 tests) 20ms + ✓ tests/dir-watcher-poller.test.ts (6 tests) 15ms + ✓ tests/hello-deterministic.test.ts (5 tests) 66ms + ✓ tests/work-package-validator.test.ts (7 tests) 10ms + ✓ tests/spec-parity.test.ts (15 tests) 91ms + ✓ tests/parse-json-output.test.ts (7 tests) 7ms + + Test Files 17 passed (17) + Tests 239 passed (239) + Duration 89.31s (transform 1.02s, setup 0ms, collect 2.93s, tests 75.22s, environment 8ms, prepare 3.47s) +``` + +The parallel rerun immediately before that was not counted as a pass. It +reported one pre-existing five-second timeout under load; the same test passed +alone in 2.17 seconds: + +```text +$ RELAYFLOWS_ALLOW_ANALYZER_SKIP=1 RELAYFLOWD_BIN=/tmp/flows-v2-cancel-target/debug/relayflowd ./node_modules/.bin/vitest run + + Test Files 1 failed | 16 passed (17) + Tests 1 failed | 238 passed (239) + Duration 74.79s (transform 3.32s, setup 0ms, collect 11.22s, tests 92.34s, environment 15ms, prepare 7.39s) + +$ RELAYFLOWD_BIN=/tmp/flows-v2-cancel-target/debug/relayflowd ./node_modules/.bin/vitest run tests/live-kernel.test.ts -t 'runs rung \(a\), parks rung \(b\)' + + ✓ tests/live-kernel.test.ts (18 tests | 17 skipped) 2174ms + ✓ built flows CLI against live relayflowd > runs rung (a), parks rung (b), and keeps JSON report-shaped 2170ms + + Test Files 1 passed (1) + Tests 1 passed | 17 skipped (18) + Duration 3.68s (transform 302ms, setup 0ms, collect 419ms, tests 2.17s, environment 0ms, prepare 191ms) +``` + +The focused full-stack client/server cancellation case was also run alone +after the final kernel change: + +```text +$ RELAYFLOWD_BIN=/tmp/flows-v2-cancel-target/debug/relayflowd ./node_modules/.bin/vitest run tests/live-kernel.test.ts -t 'cancels over the real socket' + + RUN v2.1.9 /Users/khaliqgant/AgentWorkforce/flows-v2-cancel-wt/sdk + + ✓ tests/live-kernel.test.ts (18 tests | 17 skipped) 182ms + + Test Files 1 passed (1) + Tests 1 passed | 17 skipped (18) + Duration 1.75s (transform 293ms, setup 0ms, collect 408ms, tests 182ms, environment 0ms, prepare 245ms) +``` + +## Scoped formatting + +```text +$ rustfmt --edition 2024 --check --config skip_children=true kernel/relayflowd-core/src/entry.rs kernel/relayflowd-core/src/lib.rs kernel/relayflowd-core/src/machine.rs kernel/relayflowd-core/src/machine/cancel.rs kernel/relayflowd-core/src/machine/recovery.rs kernel/relayflowd-core/src/machine/tests.rs kernel/relayflowd-core/src/state.rs kernel/relayflowd/src/engine.rs kernel/relayflowd/src/engine/remote.rs kernel/relayflowd/src/lib.rs kernel/relayflowd/src/main.rs kernel/relayflowd/src/server.rs kernel/relayflowd/src/server/cancel.rs kernel/relayflowd/src/server/client.rs kernel/relayflowd/src/server/session.rs kernel/relayflowd/tests/crash_resume.rs kernel/relayflowd/tests/crash_resume/concurrency.rs +RUSTFMT_SCOPED_EXIT=0 +``` diff --git a/sdk/src/index.ts b/sdk/src/index.ts index f6ba9445..79144817 100644 --- a/sdk/src/index.ts +++ b/sdk/src/index.ts @@ -89,6 +89,8 @@ export type { RunGetResult, RunOutcome, RunCompletionReason, + RunCancelParams, + RunCancelResult, RunResumeParams, RunResumeResult, RunStatus, diff --git a/sdk/src/journal-client.ts b/sdk/src/journal-client.ts index 5082d022..b26e73f1 100644 --- a/sdk/src/journal-client.ts +++ b/sdk/src/journal-client.ts @@ -184,6 +184,11 @@ export class JournalClient extends EventEmitter { return this.request('run.resume', { run_id: runId }, null); } + /** Durably request cancellation and return the terminal run fact. */ + runCancel(runId: string): Promise { + return this.request('run.cancel', { run_id: runId }, null); + } + /** Snapshot for legibility. */ runGet(runId: string): Promise { return this.request('run.get', { run_id: runId }); diff --git a/sdk/src/protocol.ts b/sdk/src/protocol.ts index 22837ed3..3ec7f306 100644 --- a/sdk/src/protocol.ts +++ b/sdk/src/protocol.ts @@ -44,6 +44,7 @@ export type Verb = | 'hello' | 'run.start' | 'run.resume' + | 'run.cancel' | 'run.get' | 'run.watch' | 'worker.attach' @@ -93,6 +94,11 @@ export interface RunResumeParams { } export type RunResumeResult = RunOutcome; +export interface RunCancelParams { + run_id: string; +} +export type RunCancelResult = RunOutcome; + export interface RunGetParams { run_id: string; } @@ -340,6 +346,7 @@ export interface VerbContract { hello: { params: HelloParams; result: HelloResult }; 'run.start': { params: RunStartParams; result: RunStartResult }; 'run.resume': { params: RunResumeParams; result: RunResumeResult }; + 'run.cancel': { params: RunCancelParams; result: RunCancelResult }; 'run.get': { params: RunGetParams; result: RunGetResult }; 'run.watch': { params: RunWatchParams; result: RunWatchResult }; 'worker.attach': { params: WorkerAttachParams; result: WorkerAttachResult }; diff --git a/sdk/tests/journal-client-loopback.ts b/sdk/tests/journal-client-loopback.ts index c7ad8c20..84df2333 100644 --- a/sdk/tests/journal-client-loopback.ts +++ b/sdk/tests/journal-client-loopback.ts @@ -27,6 +27,7 @@ export interface LoopbackHandlers { hello?: (ctx: FrameCtx) => void; 'run.start'?: (ctx: FrameCtx, params: Record) => void; 'run.resume'?: (ctx: FrameCtx, params: Record) => void; + 'run.cancel'?: (ctx: FrameCtx, params: Record) => void; 'run.get'?: (ctx: FrameCtx, params: Record) => void; 'run.watch'?: (ctx: FrameCtx, params: Record) => void; 'worker.attach'?: (ctx: FrameCtx, params: Record) => void; diff --git a/sdk/tests/journal-client.test.ts b/sdk/tests/journal-client.test.ts index c6be6930..527088de 100644 --- a/sdk/tests/journal-client.test.ts +++ b/sdk/tests/journal-client.test.ts @@ -61,6 +61,14 @@ describe('JournalClient: protocol v0 over unix socket', () => { 'run.resume': (ctx, params) => { sendResult(ctx, { run_id: params.run_id, status: 'completed', completed_steps: 2 }); }, + 'run.cancel': (ctx, params) => { + sendResult(ctx, { + run_id: params.run_id, + status: 'failed', + completion_reason: 'canceled', + completed_steps: 1, + }); + }, 'run.get': (ctx) => { sendResult(ctx, { status: 'running', steps: [], budget: { tokens_in: 0, tokens_out: 0, dollars: '0' } }); }, @@ -222,6 +230,17 @@ steps: expect(read.messages).toEqual([{ offset: 7 }]); }); + it('cancels a run through the typed lifecycle surface', async () => { + client = new JournalClient(path, { requestTimeoutMs: 2000 }); + await client.connect(); + await expect(client.runCancel('run-01')).resolves.toEqual({ + run_id: 'run-01', + status: 'failed', + completion_reason: 'canceled', + completed_steps: 1, + }); + }); + it('fails closed when the server returns an error', async () => { const errorPath = sockPath(); const errorServer = startLoopback(errorPath, { diff --git a/sdk/tests/live-kernel.test.ts b/sdk/tests/live-kernel.test.ts index 5ae7a9b0..3e7f1b22 100644 --- a/sdk/tests/live-kernel.test.ts +++ b/sdk/tests/live-kernel.test.ts @@ -204,6 +204,50 @@ steps: expect(completed.stderr).not.toContain('protocol_error'); }); + it('cancels over the real socket and rejects the lease holder after closure', async () => { + const dataDir = temporaryDirectory('flows-live-cancel-'); + await startDaemon(dataDir); + const worker = await connectClient(dataDir); + await worker.hello('live-cancel-worker'); + const dispatched = eventOnce(worker, 'step.dispatch'); + await worker.workerAttach('live-cancel-worker', ['llm']); + const control = await connectClient(dataDir); + await control.hello('live-cancel-control'); + const started = await control.runStart(toKernelSpec(compileYaml(` +version: '0.1.0' +steps: + - id: model + type: llm + prompt: answer +`))); + const lease = await dispatched; + + const canceled = await control.runCancel(started.run_id); + expect(canceled).toMatchObject({ + run_id: started.run_id, + status: 'failed', + completion_reason: 'canceled', + }); + await expect(control.runCancel(started.run_id)).resolves.toEqual(canceled); + await expect(worker.stepComplete( + lease.run_id, + lease.step_id, + lease.attempt, + lease.idempotency_key, + 'success', + { output: { answer: 4 } }, + )).rejects.toMatchObject({ code: 'lease_conflict' }); + + const entries = (await control.journalRead(started.run_id, 1)).entries as { + entry_type: string; + payload: { completionReason?: string }; + }[]; + expect(entries.filter((entry) => entry.entry_type === 'run.cancel.requested')).toHaveLength(1); + expect(entries.filter((entry) => entry.entry_type === 'run.completed')).toEqual([ + expect.objectContaining({ payload: expect.objectContaining({ completionReason: 'canceled' }) }), + ]); + }); + it('runs an agent CLI end to end through the SDK worker', async () => { const directory = temporaryDirectory('flows-live-agent-worker-'); const dataDir = join(directory, 'data');