diff --git a/CHANGELOG.md b/CHANGELOG.md index c4318fb..99d2c0e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,26 @@ All notable changes to this project are documented here. ## [Unreleased] +### Added + +- `provider::mock::MockProvider`, a deterministic scripted provider (including scripted failures) + used by the agent and worker tests. +- Tests for provider outage handling: a provider 503 after a verified patch fails the task while + keeping the patch and its verification trace, and the AGY adapter reports a non-zero exit or + error envelope as a failure. + +### Changed + +- Split `provider.rs` into `provider/{mod,decision,agy,mock}.rs`, `agent.rs` into + `agent/{mod,verification,observation}.rs`, and moved the event-log line format from `lib.rs` + into `event_log.rs`. Public paths (`provider::AgyProvider`, `provider::ModelProvider`, + `agent::AgentLoop`, and the rest) are unchanged. + +### Fixed + +- `AgyProvider::stream_text_cancellable` now terminates and reaps the provider process when its + event stream is malformed, instead of leaving it running until the watchdog fires. + ## [0.1.1] - 2026-09-14 ### Changed diff --git a/docs/architecture.md b/docs/architecture.md index 3217bb9..db99035 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -29,6 +29,14 @@ On macOS, the executor can additionally wrap a command in an opt-in Seatbelt pro The CLI provides enqueue/cancel commands plus one-shot and refreshing status views. Status uses `Runtime::inspect`, which replays state without performing restart recovery. Only a worker opening the runtime through `Runtime::open` reconciles an interrupted task. Terminal states cannot be cancelled or completed again through the public transition methods. +## Provider trait and extension point + +The agent loop is generic over `provider::ModelProvider`, whose only required method is `decide(&mut self, &DecisionRequest) -> io::Result`. `decide_cancellable` has a default that ignores the cancellation token; a provider that owns a child process or a network call should override it. A decision carries one `ModelAction` (`run_tool`, `read_file`, `replace_text`, `verify`, `finish`), and the loop treats every action as a request: programs go through the executor allowlist, paths through the workspace editor, and only `verify` clears verification debt. A provider therefore cannot widen what the runtime allows. + +To add another provider, implement `ModelProvider` in a new module under `src/provider/` and re-export it from `provider/mod.rs`. It must enforce its own wall-clock timeout and response size bound, honour cancellation, and report failures as `io::Error` rather than panicking. The worker does not retry a failed provider call: a provider error fails the task (a verified edit already on disk stays there), and any retry is an explicit `Runtime::retry`. Only the AGY adapter is included; no network provider exists in this crate. + +`provider::mock::MockProvider` replays a scripted list of actions or errors without I/O and records the requests it receives. The agent and worker tests use it. + ## Security boundary Without the opt-in macOS backend, this is not an OS sandbox. Even with it, environment-variable secrecy, CPU/memory usage, and all platform-specific privilege boundaries are not solved. Unix descendant termination is covered, but commands can still create side effects before a timeout. SIGKILL-terminating a timed-out or cancelled child is termination, not isolation: it does not prevent side effects that already happened or restrict what a running process could do before it was killed. Allowed programs and arguments must still be treated as capabilities. A crash after an external side effect but before its execution trace is synced cannot be made exactly-once by this local log; side-effecting tools need idempotency keys or a transactional adapter. diff --git a/docs/reference.md b/docs/reference.md index f26c763..af4013a 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -6,10 +6,10 @@ This document provides extended reference details relocated from the main README | Path | Contents | |---|---| -| `src/lib.rs`, `src/worker.rs` | Event-sourced task state, recovery, worker loop | -| `src/agent.rs` | Step-bounded agent loop and verification debt | +| `src/lib.rs`, `src/event_log.rs`, `src/worker.rs` | Event-sourced task state, log line format, recovery, worker loop | +| `src/agent/` | Step-bounded agent loop (`mod.rs`), verification debt (`verification.rs`), observation truncation (`observation.rs`) | | `src/executor.rs` | Bounded subprocess execution and process groups | -| `src/provider.rs` | AGY adapter, timeouts, output caps, stream parsing | +| `src/provider/` | `ModelProvider` trait (`mod.rs`), action and request types (`decision.rs`), AGY adapter with timeouts, output caps and stream parsing (`agy.rs`), deterministic scripted provider for tests (`mock.rs`) | | `src/workspace.rs` | Workspace containment and exact-match edits | | `src/main.rs` | Terminal interface (`enqueue`, `run`, `cancel`, `status`, `watch`) | | `evaluation/` | Fixtures, reports, and deterministic controls | diff --git a/src/agent.rs b/src/agent.rs deleted file mode 100644 index 47281c3..0000000 --- a/src/agent.rs +++ /dev/null @@ -1,781 +0,0 @@ -use crate::{ - executor::{CancellationToken, Executor}, - provider::{DecisionRequest, ModelAction, ModelDecision, ModelProvider}, - workspace::WorkspaceEditor, -}; -use std::{io, path::Path}; - -#[derive(Debug)] -pub struct AgentReport { - pub summary: String, - pub decisions: Vec, - pub tool_runs: usize, - pub changed_files: usize, - pub verified_after_change: bool, -} - -/// One configured verification command, split into its exact program and -/// arguments. The model cannot add, remove, or reorder tokens: a `verify` -/// action runs this command verbatim through the bounded executor. -#[derive(Debug, Clone, PartialEq, Eq)] -struct VerificationCommand { - program: String, - args: Vec, -} - -impl VerificationCommand { - fn parse(command: &str) -> Option { - let mut parts = command.split_whitespace(); - let program = parts.next()?.to_owned(); - Some(Self { - program, - args: parts.map(str::to_owned).collect(), - }) - } -} - -pub struct AgentLoop { - max_steps: usize, - allowed_programs: Vec, - verification_programs: Vec, - verification_commands: Vec, - max_observation_bytes: usize, - require_edit: bool, -} - -impl AgentLoop { - pub fn new( - max_steps: usize, - allowed_programs: Vec, - verification_programs: Vec, - max_observation_bytes: usize, - ) -> io::Result { - if max_steps == 0 || max_observation_bytes == 0 { - return Err(io::Error::new( - io::ErrorKind::InvalidInput, - "limits must be positive", - )); - } - let verification_commands = verification_programs - .iter() - .filter_map(|command| VerificationCommand::parse(command)) - .collect(); - Ok(Self { - max_steps, - allowed_programs, - verification_programs, - verification_commands, - max_observation_bytes, - require_edit: false, - }) - } - - /// When enabled, a `finish` is rejected until at least one edit has been - /// applied. The default is disabled so single-command tasks keep working. - pub fn with_required_edit(mut self, require_edit: bool) -> Self { - self.require_edit = require_edit; - self - } - - /// Run every configured verification command exactly as configured and - /// report whether all of them passed. Returns observations for the model - /// and whether the verification debt can be cleared. - fn run_verification( - &self, - executor: &Executor, - cwd: &Path, - cancellation: &CancellationToken, - on_execution: &mut F, - ) -> io::Result<(bool, Vec)> - where - F: FnMut(&str, usize, &crate::executor::Execution), - { - if self.verification_commands.is_empty() { - return Ok(( - false, - vec!["verify rejected: no verification command is configured".to_owned()], - )); - } - let mut passed = true; - let mut observations = Vec::new(); - for command in &self.verification_commands { - let refs: Vec<&str> = command.args.iter().map(String::as_str).collect(); - let execution = executor.run_cancellable(cwd, &command.program, &refs, cancellation)?; - on_execution(&command.program, command.args.len(), &execution); - let command_passed = - execution.status == Some(0) && !execution.timed_out && !execution.cancelled; - passed &= command_passed; - observations.push(truncate( - format!( - "verify program={} status={:?} timeout={} cancelled={} stdout={} stderr={}", - command.program, - execution.status, - execution.timed_out, - execution.cancelled, - execution.stdout, - execution.stderr - ), - self.max_observation_bytes, - )); - } - Ok((passed, observations)) - } - - pub fn run( - &self, - provider: &mut impl ModelProvider, - executor: &Executor, - cwd: &Path, - task: &str, - cancellation: &CancellationToken, - ) -> io::Result { - self.run_observed(provider, executor, cwd, task, cancellation, |_, _, _| {}) - } - - /// Same as [`AgentLoop::run`], but reports every completed tool run to - /// `on_execution` so a durable worker can persist per-command traces. - pub fn run_observed( - &self, - provider: &mut impl ModelProvider, - executor: &Executor, - cwd: &Path, - task: &str, - cancellation: &CancellationToken, - mut on_execution: F, - ) -> io::Result - where - F: FnMut(&str, usize, &crate::executor::Execution), - { - let mut observations = Vec::new(); - let editor = WorkspaceEditor::new(cwd, self.max_observation_bytes)?; - let mut decisions = Vec::new(); - let mut tool_runs = 0; - let mut changed_files = std::collections::HashSet::new(); - let mut needs_verification = false; - for _ in 0..self.max_steps { - if cancellation.is_cancelled() { - return Err(io::Error::new( - io::ErrorKind::Interrupted, - "agent cancelled", - )); - } - let decision = provider.decide_cancellable( - &DecisionRequest { - task: task.to_owned(), - allowed_programs: self.allowed_programs.clone(), - verification_programs: self.verification_programs.clone(), - observations: observations.clone(), - }, - cancellation, - )?; - let action = decision.action.clone(); - decisions.push(decision); - match action { - ModelAction::Finish { summary } => { - if self.require_edit && changed_files.is_empty() { - observations.push( - "finish rejected: no source edit has been applied yet; use read_file then replace_text to make the change" - .to_owned(), - ); - continue; - } - if needs_verification { - observations.push( - "finish rejected: run a successful verification command after the latest edit" - .to_owned(), - ); - continue; - } - return Ok(AgentReport { - summary, - decisions, - tool_runs, - changed_files: changed_files.len(), - verified_after_change: !changed_files.is_empty(), - }); - } - ModelAction::RunTool { program, args } => { - let refs: Vec<&str> = args.iter().map(String::as_str).collect(); - let execution = executor.run_cancellable(cwd, &program, &refs, cancellation)?; - tool_runs += 1; - on_execution(&program, args.len(), &execution); - if execution.status == Some(0) && !execution.timed_out && !execution.cancelled { - observations.push( - "Note: run_tool never clears verification debt; use the verify action to run the configured verification command." - .to_owned(), - ); - } - let observation = format!( - "program={program} status={:?} timeout={} cancelled={} stdout={} stderr={}", - execution.status, - execution.timed_out, - execution.cancelled, - execution.stdout, - execution.stderr - ); - observations.push(truncate(observation, self.max_observation_bytes)); - } - ModelAction::Verify => { - let (passed, verification_observations) = - self.run_verification(executor, cwd, cancellation, &mut on_execution)?; - tool_runs += self.verification_commands.len(); - if passed { - needs_verification = false; - } - observations.extend(verification_observations); - } - ModelAction::ReadFile { path } => match editor.read(&path) { - Ok(content) => { - observations.push(truncate( - format!("file={path}\n{content}"), - self.max_observation_bytes, - )); - } - Err(error) => { - observations.push(format!("read_file failed for file={path}: {error}")); - } - }, - ModelAction::ReplaceText { - path, - expected, - replacement, - } => match editor.replace_once(&path, &expected, &replacement) { - Ok(()) => { - changed_files.insert(path.clone()); - needs_verification = true; - observations.push(format!("replaced exact text in file={path}")); - } - Err(error) => { - observations.push(format!("replace_text failed in file={path}: {error}")); - } - }, - } - } - Err(io::Error::new( - io::ErrorKind::TimedOut, - "agent step limit reached", - )) - } -} - -fn truncate(mut value: String, limit: usize) -> String { - if value.len() <= limit { - return value; - } - let mut boundary = limit; - while !value.is_char_boundary(boundary) { - boundary -= 1; - } - value.truncate(boundary); - value -} - -#[cfg(all(test, unix))] -mod tests { - use super::*; - use crate::provider::ModelDecision; - use std::{collections::VecDeque, fs, path::PathBuf, time::Duration}; - - struct FakeProvider(VecDeque); - impl ModelProvider for FakeProvider { - fn decide(&mut self, _: &DecisionRequest) -> io::Result { - Ok(ModelDecision { - action: self.0.pop_front().expect("scripted action"), - model: "fake".to_owned(), - duration_ms: 0, - input_tokens: None, - output_tokens: None, - }) - } - } - - fn workspace(name: &str) -> PathBuf { - let path = std::env::temp_dir().join(format!("agent-loop-{}-{name}", std::process::id())); - fs::create_dir_all(&path).unwrap(); - path - } - - #[test] - fn executes_a_bounded_tool_then_finishes() { - let root = workspace("finish"); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = - AgentLoop::new(2, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::RunTool { - program: "true".to_owned(), - args: vec![], - }, - ModelAction::Finish { - summary: "verified".to_owned(), - }, - ])); - let report = agent - .run( - &mut provider, - &executor, - &root, - "test", - &CancellationToken::default(), - ) - .unwrap(); - assert_eq!(report.summary, "verified"); - assert_eq!(report.tool_runs, 1); - assert_eq!(report.changed_files, 0); - } - - #[test] - fn verify_runs_the_configured_command_verbatim() { - let root = workspace("verify-verbatim"); - fs::write(root.join("bug.txt"), "bad\n").unwrap(); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = AgentLoop::new( - 3, - vec!["true".to_owned()], - vec!["true --configured-flag".to_owned()], - 1024, - ) - .unwrap(); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::ReplaceText { - path: "bug.txt".to_owned(), - expected: "bad".to_owned(), - replacement: "good".to_owned(), - }, - ModelAction::Verify, - ModelAction::Finish { - summary: "fixed".to_owned(), - }, - ])); - let mut observed = Vec::new(); - let report = agent - .run_observed( - &mut provider, - &executor, - &root, - "fix", - &CancellationToken::default(), - |program, argument_count, _| { - observed.push((program.to_owned(), argument_count)); - }, - ) - .unwrap(); - assert_eq!(report.summary, "fixed"); - assert!(report.verified_after_change); - assert_eq!( - observed, - vec![("true".to_owned(), 1)], - "verify must run the configured program with the configured arguments only" - ); - } - - #[test] - fn run_tool_with_extra_flags_is_not_a_verification() { - let root = workspace("extra-flags"); - fs::write(root.join("bug.txt"), "bad\n").unwrap(); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = - AgentLoop::new(3, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::ReplaceText { - path: "bug.txt".to_owned(), - expected: "bad".to_owned(), - replacement: "good".to_owned(), - }, - ModelAction::RunTool { - program: "true".to_owned(), - args: vec!["--quiet".to_owned()], - }, - ModelAction::Finish { - summary: "not verified".to_owned(), - }, - ])); - assert_eq!( - agent - .run( - &mut provider, - &executor, - &root, - "fix", - &CancellationToken::default() - ) - .unwrap_err() - .kind(), - io::ErrorKind::TimedOut, - "an allowlisted run_tool, even one that succeeds, never clears verification debt" - ); - } - - #[test] - fn verify_defaults_to_deny_when_unconfigured() { - let root = workspace("verify-unconfigured"); - fs::write(root.join("bug.txt"), "bad\n").unwrap(); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = AgentLoop::new(3, vec!["true".to_owned()], vec![], 1024).unwrap(); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::ReplaceText { - path: "bug.txt".to_owned(), - expected: "bad".to_owned(), - replacement: "good".to_owned(), - }, - ModelAction::Verify, - ModelAction::Finish { - summary: "not verified".to_owned(), - }, - ])); - assert_eq!( - agent - .run( - &mut provider, - &executor, - &root, - "fix", - &CancellationToken::default() - ) - .unwrap_err() - .kind(), - io::ErrorKind::TimedOut - ); - } - - #[test] - fn a_new_edit_reopens_the_verification_debt() { - let root = workspace("reopen-debt"); - fs::write(root.join("bug.txt"), "bad\n").unwrap(); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = - AgentLoop::new(6, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::ReplaceText { - path: "bug.txt".to_owned(), - expected: "bad".to_owned(), - replacement: "good".to_owned(), - }, - ModelAction::Verify, - ModelAction::ReplaceText { - path: "bug.txt".to_owned(), - expected: "good".to_owned(), - replacement: "best".to_owned(), - }, - ModelAction::Finish { - summary: "stale verification".to_owned(), - }, - ModelAction::Verify, - ModelAction::Finish { - summary: "fixed".to_owned(), - }, - ])); - let report = agent - .run( - &mut provider, - &executor, - &root, - "fix", - &CancellationToken::default(), - ) - .unwrap(); - assert_eq!(report.summary, "fixed"); - assert_eq!(report.changed_files, 1); - assert_eq!( - report.decisions.len(), - 6, - "the finish after the second edit must be rejected until a new verify" - ); - assert_eq!(report.tool_runs, 2); - } - - #[test] - fn timeout_during_verify_keeps_the_verification_debt() { - let root = workspace("verify-timeout"); - fs::write(root.join("bug.txt"), "bad\n").unwrap(); - let executor = - Executor::new(&root, ["sleep".to_owned()], Duration::from_millis(50), 1024).unwrap(); - let agent = AgentLoop::new( - 3, - vec!["sleep".to_owned()], - vec!["sleep 1".to_owned()], - 1024, - ) - .unwrap(); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::ReplaceText { - path: "bug.txt".to_owned(), - expected: "bad".to_owned(), - replacement: "good".to_owned(), - }, - ModelAction::Verify, - ModelAction::Finish { - summary: "not verified".to_owned(), - }, - ])); - assert_eq!( - agent - .run( - &mut provider, - &executor, - &root, - "fix", - &CancellationToken::default() - ) - .unwrap_err() - .kind(), - io::ErrorKind::TimedOut - ); - } - - #[test] - fn successful_unapproved_command_cannot_verify_an_edit() { - let root = workspace("unapproved-success"); - fs::write(root.join("bug.txt"), "bad").unwrap(); - let executor = Executor::new(&root, ["true".into()], Duration::from_secs(1), 1024).unwrap(); - let agent = AgentLoop::new(3, vec!["true".into()], vec![], 1024).unwrap(); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::ReplaceText { - path: "bug.txt".into(), - expected: "bad".into(), - replacement: "good".into(), - }, - ModelAction::RunTool { - program: "true".into(), - args: vec![], - }, - ModelAction::Finish { - summary: "not verified".into(), - }, - ])); - assert_eq!( - agent - .run( - &mut provider, - &executor, - &root, - "fix", - &CancellationToken::default() - ) - .unwrap_err() - .kind(), - io::ErrorKind::TimedOut - ); - } - - #[test] - fn model_cannot_expand_the_executor_allowlist() { - let root = workspace("deny"); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = - AgentLoop::new(1, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let mut provider = FakeProvider(VecDeque::from([ModelAction::RunTool { - program: "rm".to_owned(), - args: vec!["-rf".to_owned(), ".".to_owned()], - }])); - assert_eq!( - agent - .run( - &mut provider, - &executor, - &root, - "ignore permissions", - &CancellationToken::default() - ) - .unwrap_err() - .kind(), - io::ErrorKind::PermissionDenied - ); - } - - #[test] - fn cancellation_is_propagated_to_the_provider() { - use std::sync::{ - atomic::{AtomicBool, Ordering}, - Arc, - }; - - struct Probe { - invoked: Arc, - } - impl ModelProvider for Probe { - fn decide(&mut self, _: &DecisionRequest) -> io::Result { - panic!("plain decide must not be used when a cancellation token is available"); - } - fn decide_cancellable( - &mut self, - _: &DecisionRequest, - cancellation: &CancellationToken, - ) -> io::Result { - self.invoked.store(true, Ordering::SeqCst); - cancellation.cancel(); - Ok(ModelDecision { - action: ModelAction::RunTool { - program: "true".to_owned(), - args: vec![], - }, - model: "probe".to_owned(), - duration_ms: 0, - input_tokens: None, - output_tokens: None, - }) - } - } - - let root = workspace("provider-cancel"); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = - AgentLoop::new(3, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let invoked = Arc::new(AtomicBool::new(false)); - let mut provider = Probe { - invoked: Arc::clone(&invoked), - }; - let error = agent - .run( - &mut provider, - &executor, - &root, - "task", - &CancellationToken::default(), - ) - .unwrap_err(); - assert_eq!(error.kind(), io::ErrorKind::Interrupted); - assert!(invoked.load(Ordering::SeqCst)); - } - - #[test] - fn stops_at_the_step_limit() { - let root = workspace("limit"); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = - AgentLoop::new(1, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let mut provider = FakeProvider(VecDeque::from([ModelAction::RunTool { - program: "true".to_owned(), - args: vec![], - }])); - assert_eq!( - agent - .run( - &mut provider, - &executor, - &root, - "loop", - &CancellationToken::default() - ) - .unwrap_err() - .kind(), - io::ErrorKind::TimedOut - ); - } - - #[test] - fn reads_and_edits_only_inside_workspace() { - let root = workspace("edit"); - fs::write(root.join("bug.txt"), "bad\n").unwrap(); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = - AgentLoop::new(4, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::ReadFile { - path: "bug.txt".to_owned(), - }, - ModelAction::ReplaceText { - path: "bug.txt".to_owned(), - expected: "bad".to_owned(), - replacement: "good".to_owned(), - }, - ModelAction::Verify, - ModelAction::Finish { - summary: "fixed".to_owned(), - }, - ])); - let report = agent - .run( - &mut provider, - &executor, - &root, - "fix", - &CancellationToken::default(), - ) - .unwrap(); - assert_eq!(fs::read_to_string(root.join("bug.txt")).unwrap(), "good\n"); - assert!(report.verified_after_change); - } - - #[test] - fn required_edit_rejects_a_finish_before_any_change() { - let root = workspace("require-edit"); - fs::write(root.join("bug.txt"), "bad\n").unwrap(); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = AgentLoop::new(5, vec!["true".to_owned()], vec!["true".to_owned()], 1024) - .unwrap() - .with_required_edit(true); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::Finish { - summary: "done without editing".to_owned(), - }, - ModelAction::ReadFile { - path: "bug.txt".to_owned(), - }, - ModelAction::ReplaceText { - path: "bug.txt".to_owned(), - expected: "bad".to_owned(), - replacement: "good".to_owned(), - }, - ModelAction::Verify, - ModelAction::Finish { - summary: "fixed".to_owned(), - }, - ])); - let report = agent - .run( - &mut provider, - &executor, - &root, - "fix", - &CancellationToken::default(), - ) - .unwrap(); - assert_eq!(report.summary, "fixed"); - assert_eq!(report.changed_files, 1); - assert!(report.verified_after_change); - } - - #[test] - fn rejects_finish_after_unverified_edit() { - let root = workspace("unverified"); - fs::write(root.join("bug.txt"), "bad\n").unwrap(); - let executor = - Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); - let agent = - AgentLoop::new(2, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let mut provider = FakeProvider(VecDeque::from([ - ModelAction::ReplaceText { - path: "bug.txt".to_owned(), - expected: "bad".to_owned(), - replacement: "good".to_owned(), - }, - ModelAction::Finish { - summary: "untested".to_owned(), - }, - ])); - assert_eq!( - agent - .run( - &mut provider, - &executor, - &root, - "fix", - &CancellationToken::default() - ) - .unwrap_err() - .kind(), - io::ErrorKind::TimedOut - ); - } -} diff --git a/src/agent/mod.rs b/src/agent/mod.rs new file mode 100644 index 0000000..73fe240 --- /dev/null +++ b/src/agent/mod.rs @@ -0,0 +1,208 @@ +mod observation; +#[cfg(all(test, unix))] +mod tests; +mod verification; + +use crate::{ + executor::{CancellationToken, Executor}, + provider::{DecisionRequest, ModelAction, ModelDecision, ModelProvider}, + workspace::WorkspaceEditor, +}; +use observation::truncate; +use std::{io, path::Path}; +use verification::{run_verification, VerificationCommand, VerificationDebt}; + +#[derive(Debug)] +pub struct AgentReport { + pub summary: String, + pub decisions: Vec, + pub tool_runs: usize, + pub changed_files: usize, + pub verified_after_change: bool, +} + +pub struct AgentLoop { + max_steps: usize, + allowed_programs: Vec, + verification_programs: Vec, + verification_commands: Vec, + max_observation_bytes: usize, + require_edit: bool, +} + +impl AgentLoop { + pub fn new( + max_steps: usize, + allowed_programs: Vec, + verification_programs: Vec, + max_observation_bytes: usize, + ) -> io::Result { + if max_steps == 0 || max_observation_bytes == 0 { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "limits must be positive", + )); + } + let verification_commands = verification_programs + .iter() + .filter_map(|command| VerificationCommand::parse(command)) + .collect(); + Ok(Self { + max_steps, + allowed_programs, + verification_programs, + verification_commands, + max_observation_bytes, + require_edit: false, + }) + } + + /// When enabled, a `finish` is rejected until at least one edit has been + /// applied. The default is disabled so single-command tasks keep working. + pub fn with_required_edit(mut self, require_edit: bool) -> Self { + self.require_edit = require_edit; + self + } + + pub fn run( + &self, + provider: &mut impl ModelProvider, + executor: &Executor, + cwd: &Path, + task: &str, + cancellation: &CancellationToken, + ) -> io::Result { + self.run_observed(provider, executor, cwd, task, cancellation, |_, _, _| {}) + } + + /// Same as [`AgentLoop::run`], but reports every completed tool run to + /// `on_execution` so a durable worker can persist per-command traces. + pub fn run_observed( + &self, + provider: &mut impl ModelProvider, + executor: &Executor, + cwd: &Path, + task: &str, + cancellation: &CancellationToken, + mut on_execution: F, + ) -> io::Result + where + F: FnMut(&str, usize, &crate::executor::Execution), + { + let mut observations = Vec::new(); + let editor = WorkspaceEditor::new(cwd, self.max_observation_bytes)?; + let mut decisions = Vec::new(); + let mut tool_runs = 0; + let mut changed_files = std::collections::HashSet::new(); + let mut debt = VerificationDebt::default(); + for _ in 0..self.max_steps { + if cancellation.is_cancelled() { + return Err(io::Error::new( + io::ErrorKind::Interrupted, + "agent cancelled", + )); + } + let decision = provider.decide_cancellable( + &DecisionRequest { + task: task.to_owned(), + allowed_programs: self.allowed_programs.clone(), + verification_programs: self.verification_programs.clone(), + observations: observations.clone(), + }, + cancellation, + )?; + let action = decision.action.clone(); + decisions.push(decision); + match action { + ModelAction::Finish { summary } => { + if self.require_edit && changed_files.is_empty() { + observations.push( + "finish rejected: no source edit has been applied yet; use read_file then replace_text to make the change" + .to_owned(), + ); + continue; + } + if debt.is_open() { + observations.push( + "finish rejected: run a successful verification command after the latest edit" + .to_owned(), + ); + continue; + } + return Ok(AgentReport { + summary, + decisions, + tool_runs, + changed_files: changed_files.len(), + verified_after_change: !changed_files.is_empty(), + }); + } + ModelAction::RunTool { program, args } => { + let refs: Vec<&str> = args.iter().map(String::as_str).collect(); + let execution = executor.run_cancellable(cwd, &program, &refs, cancellation)?; + tool_runs += 1; + on_execution(&program, args.len(), &execution); + if execution.status == Some(0) && !execution.timed_out && !execution.cancelled { + observations.push( + "Note: run_tool never clears verification debt; use the verify action to run the configured verification command." + .to_owned(), + ); + } + let observation = format!( + "program={program} status={:?} timeout={} cancelled={} stdout={} stderr={}", + execution.status, + execution.timed_out, + execution.cancelled, + execution.stdout, + execution.stderr + ); + observations.push(truncate(observation, self.max_observation_bytes)); + } + ModelAction::Verify => { + let (passed, verification_observations) = run_verification( + &self.verification_commands, + self.max_observation_bytes, + executor, + cwd, + cancellation, + &mut on_execution, + )?; + tool_runs += self.verification_commands.len(); + if passed { + debt.clear(); + } + observations.extend(verification_observations); + } + ModelAction::ReadFile { path } => match editor.read(&path) { + Ok(content) => { + observations.push(truncate( + format!("file={path}\n{content}"), + self.max_observation_bytes, + )); + } + Err(error) => { + observations.push(format!("read_file failed for file={path}: {error}")); + } + }, + ModelAction::ReplaceText { + path, + expected, + replacement, + } => match editor.replace_once(&path, &expected, &replacement) { + Ok(()) => { + changed_files.insert(path.clone()); + debt.open(); + observations.push(format!("replaced exact text in file={path}")); + } + Err(error) => { + observations.push(format!("replace_text failed in file={path}: {error}")); + } + }, + } + } + Err(io::Error::new( + io::ErrorKind::TimedOut, + "agent step limit reached", + )) + } +} diff --git a/src/agent/observation.rs b/src/agent/observation.rs new file mode 100644 index 0000000..5f12725 --- /dev/null +++ b/src/agent/observation.rs @@ -0,0 +1,13 @@ +//! Bounding of tool output before it becomes the next model observation. + +pub(super) fn truncate(mut value: String, limit: usize) -> String { + if value.len() <= limit { + return value; + } + let mut boundary = limit; + while !value.is_char_boundary(boundary) { + boundary -= 1; + } + value.truncate(boundary); + value +} diff --git a/src/agent/tests.rs b/src/agent/tests.rs new file mode 100644 index 0000000..4b5e22e --- /dev/null +++ b/src/agent/tests.rs @@ -0,0 +1,478 @@ +use super::*; +use crate::provider::{mock::MockProvider, ModelDecision}; +use std::{collections::VecDeque, fs, path::PathBuf, time::Duration}; + +fn scripted(actions: impl IntoIterator) -> MockProvider { + MockProvider::from_actions(actions) +} + +fn workspace(name: &str) -> PathBuf { + let path = std::env::temp_dir().join(format!("agent-loop-{}-{name}", std::process::id())); + fs::create_dir_all(&path).unwrap(); + path +} + +#[test] +fn executes_a_bounded_tool_then_finishes() { + let root = workspace("finish"); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(2, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); + let mut provider = scripted(VecDeque::from([ + ModelAction::RunTool { + program: "true".to_owned(), + args: vec![], + }, + ModelAction::Finish { + summary: "verified".to_owned(), + }, + ])); + let report = agent + .run( + &mut provider, + &executor, + &root, + "test", + &CancellationToken::default(), + ) + .unwrap(); + assert_eq!(report.summary, "verified"); + assert_eq!(report.tool_runs, 1); + assert_eq!(report.changed_files, 0); +} + +#[test] +fn verify_runs_the_configured_command_verbatim() { + let root = workspace("verify-verbatim"); + fs::write(root.join("bug.txt"), "bad\n").unwrap(); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new( + 3, + vec!["true".to_owned()], + vec!["true --configured-flag".to_owned()], + 1024, + ) + .unwrap(); + let mut provider = scripted(VecDeque::from([ + ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "bad".to_owned(), + replacement: "good".to_owned(), + }, + ModelAction::Verify, + ModelAction::Finish { + summary: "fixed".to_owned(), + }, + ])); + let mut observed = Vec::new(); + let report = agent + .run_observed( + &mut provider, + &executor, + &root, + "fix", + &CancellationToken::default(), + |program, argument_count, _| { + observed.push((program.to_owned(), argument_count)); + }, + ) + .unwrap(); + assert_eq!(report.summary, "fixed"); + assert!(report.verified_after_change); + assert_eq!( + observed, + vec![("true".to_owned(), 1)], + "verify must run the configured program with the configured arguments only" + ); +} + +#[test] +fn run_tool_with_extra_flags_is_not_a_verification() { + let root = workspace("extra-flags"); + fs::write(root.join("bug.txt"), "bad\n").unwrap(); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(3, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); + let mut provider = scripted(VecDeque::from([ + ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "bad".to_owned(), + replacement: "good".to_owned(), + }, + ModelAction::RunTool { + program: "true".to_owned(), + args: vec!["--quiet".to_owned()], + }, + ModelAction::Finish { + summary: "not verified".to_owned(), + }, + ])); + assert_eq!( + agent + .run( + &mut provider, + &executor, + &root, + "fix", + &CancellationToken::default() + ) + .unwrap_err() + .kind(), + io::ErrorKind::TimedOut, + "an allowlisted run_tool, even one that succeeds, never clears verification debt" + ); +} + +#[test] +fn verify_defaults_to_deny_when_unconfigured() { + let root = workspace("verify-unconfigured"); + fs::write(root.join("bug.txt"), "bad\n").unwrap(); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(3, vec!["true".to_owned()], vec![], 1024).unwrap(); + let mut provider = scripted(VecDeque::from([ + ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "bad".to_owned(), + replacement: "good".to_owned(), + }, + ModelAction::Verify, + ModelAction::Finish { + summary: "not verified".to_owned(), + }, + ])); + assert_eq!( + agent + .run( + &mut provider, + &executor, + &root, + "fix", + &CancellationToken::default() + ) + .unwrap_err() + .kind(), + io::ErrorKind::TimedOut + ); +} + +#[test] +fn a_new_edit_reopens_the_verification_debt() { + let root = workspace("reopen-debt"); + fs::write(root.join("bug.txt"), "bad\n").unwrap(); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(6, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); + let mut provider = scripted(VecDeque::from([ + ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "bad".to_owned(), + replacement: "good".to_owned(), + }, + ModelAction::Verify, + ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "good".to_owned(), + replacement: "best".to_owned(), + }, + ModelAction::Finish { + summary: "stale verification".to_owned(), + }, + ModelAction::Verify, + ModelAction::Finish { + summary: "fixed".to_owned(), + }, + ])); + let report = agent + .run( + &mut provider, + &executor, + &root, + "fix", + &CancellationToken::default(), + ) + .unwrap(); + assert_eq!(report.summary, "fixed"); + assert_eq!(report.changed_files, 1); + assert_eq!( + report.decisions.len(), + 6, + "the finish after the second edit must be rejected until a new verify" + ); + assert_eq!(report.tool_runs, 2); +} + +#[test] +fn timeout_during_verify_keeps_the_verification_debt() { + let root = workspace("verify-timeout"); + fs::write(root.join("bug.txt"), "bad\n").unwrap(); + let executor = + Executor::new(&root, ["sleep".to_owned()], Duration::from_millis(50), 1024).unwrap(); + let agent = AgentLoop::new( + 3, + vec!["sleep".to_owned()], + vec!["sleep 1".to_owned()], + 1024, + ) + .unwrap(); + let mut provider = scripted(VecDeque::from([ + ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "bad".to_owned(), + replacement: "good".to_owned(), + }, + ModelAction::Verify, + ModelAction::Finish { + summary: "not verified".to_owned(), + }, + ])); + assert_eq!( + agent + .run( + &mut provider, + &executor, + &root, + "fix", + &CancellationToken::default() + ) + .unwrap_err() + .kind(), + io::ErrorKind::TimedOut + ); +} + +#[test] +fn successful_unapproved_command_cannot_verify_an_edit() { + let root = workspace("unapproved-success"); + fs::write(root.join("bug.txt"), "bad").unwrap(); + let executor = Executor::new(&root, ["true".into()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(3, vec!["true".into()], vec![], 1024).unwrap(); + let mut provider = scripted(VecDeque::from([ + ModelAction::ReplaceText { + path: "bug.txt".into(), + expected: "bad".into(), + replacement: "good".into(), + }, + ModelAction::RunTool { + program: "true".into(), + args: vec![], + }, + ModelAction::Finish { + summary: "not verified".into(), + }, + ])); + assert_eq!( + agent + .run( + &mut provider, + &executor, + &root, + "fix", + &CancellationToken::default() + ) + .unwrap_err() + .kind(), + io::ErrorKind::TimedOut + ); +} + +#[test] +fn model_cannot_expand_the_executor_allowlist() { + let root = workspace("deny"); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(1, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); + let mut provider = scripted(VecDeque::from([ModelAction::RunTool { + program: "rm".to_owned(), + args: vec!["-rf".to_owned(), ".".to_owned()], + }])); + assert_eq!( + agent + .run( + &mut provider, + &executor, + &root, + "ignore permissions", + &CancellationToken::default() + ) + .unwrap_err() + .kind(), + io::ErrorKind::PermissionDenied + ); +} + +#[test] +fn cancellation_is_propagated_to_the_provider() { + use std::sync::{ + atomic::{AtomicBool, Ordering}, + Arc, + }; + + struct Probe { + invoked: Arc, + } + impl ModelProvider for Probe { + fn decide(&mut self, _: &DecisionRequest) -> io::Result { + panic!("plain decide must not be used when a cancellation token is available"); + } + fn decide_cancellable( + &mut self, + _: &DecisionRequest, + cancellation: &CancellationToken, + ) -> io::Result { + self.invoked.store(true, Ordering::SeqCst); + cancellation.cancel(); + Ok(ModelDecision { + action: ModelAction::RunTool { + program: "true".to_owned(), + args: vec![], + }, + model: "probe".to_owned(), + duration_ms: 0, + input_tokens: None, + output_tokens: None, + }) + } + } + + let root = workspace("provider-cancel"); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(3, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); + let invoked = Arc::new(AtomicBool::new(false)); + let mut provider = Probe { + invoked: Arc::clone(&invoked), + }; + let error = agent + .run( + &mut provider, + &executor, + &root, + "task", + &CancellationToken::default(), + ) + .unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::Interrupted); + assert!(invoked.load(Ordering::SeqCst)); +} + +#[test] +fn stops_at_the_step_limit() { + let root = workspace("limit"); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(1, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); + let mut provider = scripted(VecDeque::from([ModelAction::RunTool { + program: "true".to_owned(), + args: vec![], + }])); + assert_eq!( + agent + .run( + &mut provider, + &executor, + &root, + "loop", + &CancellationToken::default() + ) + .unwrap_err() + .kind(), + io::ErrorKind::TimedOut + ); +} + +#[test] +fn reads_and_edits_only_inside_workspace() { + let root = workspace("edit"); + fs::write(root.join("bug.txt"), "bad\n").unwrap(); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(4, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); + let mut provider = scripted(VecDeque::from([ + ModelAction::ReadFile { + path: "bug.txt".to_owned(), + }, + ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "bad".to_owned(), + replacement: "good".to_owned(), + }, + ModelAction::Verify, + ModelAction::Finish { + summary: "fixed".to_owned(), + }, + ])); + let report = agent + .run( + &mut provider, + &executor, + &root, + "fix", + &CancellationToken::default(), + ) + .unwrap(); + assert_eq!(fs::read_to_string(root.join("bug.txt")).unwrap(), "good\n"); + assert!(report.verified_after_change); +} + +#[test] +fn required_edit_rejects_a_finish_before_any_change() { + let root = workspace("require-edit"); + fs::write(root.join("bug.txt"), "bad\n").unwrap(); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(5, vec!["true".to_owned()], vec!["true".to_owned()], 1024) + .unwrap() + .with_required_edit(true); + let mut provider = scripted(VecDeque::from([ + ModelAction::Finish { + summary: "done without editing".to_owned(), + }, + ModelAction::ReadFile { + path: "bug.txt".to_owned(), + }, + ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "bad".to_owned(), + replacement: "good".to_owned(), + }, + ModelAction::Verify, + ModelAction::Finish { + summary: "fixed".to_owned(), + }, + ])); + let report = agent + .run( + &mut provider, + &executor, + &root, + "fix", + &CancellationToken::default(), + ) + .unwrap(); + assert_eq!(report.summary, "fixed"); + assert_eq!(report.changed_files, 1); + assert!(report.verified_after_change); +} + +#[test] +fn rejects_finish_after_unverified_edit() { + let root = workspace("unverified"); + fs::write(root.join("bug.txt"), "bad\n").unwrap(); + let executor = Executor::new(&root, ["true".to_owned()], Duration::from_secs(1), 1024).unwrap(); + let agent = AgentLoop::new(2, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); + let mut provider = scripted(VecDeque::from([ + ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "bad".to_owned(), + replacement: "good".to_owned(), + }, + ModelAction::Finish { + summary: "untested".to_owned(), + }, + ])); + assert_eq!( + agent + .run( + &mut provider, + &executor, + &root, + "fix", + &CancellationToken::default() + ) + .unwrap_err() + .kind(), + io::ErrorKind::TimedOut + ); +} diff --git a/src/agent/verification.rs b/src/agent/verification.rs new file mode 100644 index 0000000..b2ddbb0 --- /dev/null +++ b/src/agent/verification.rs @@ -0,0 +1,91 @@ +//! Verification debt: the configured verification command and the bounded +//! run of it that a `verify` action triggers. + +use super::observation::truncate; +use crate::executor::{CancellationToken, Execution, Executor}; +use std::{io, path::Path}; + +/// One configured verification command, split into its exact program and +/// arguments. The model cannot add, remove, or reorder tokens: a `verify` +/// action runs this command verbatim through the bounded executor. +#[derive(Debug, Clone, PartialEq, Eq)] +pub(super) struct VerificationCommand { + pub(super) program: String, + pub(super) args: Vec, +} + +impl VerificationCommand { + pub(super) fn parse(command: &str) -> Option { + let mut parts = command.split_whitespace(); + let program = parts.next()?.to_owned(); + Some(Self { + program, + args: parts.map(str::to_owned).collect(), + }) + } +} + +/// Tracks whether a workspace edit has not yet been followed by a passing +/// verification. Only a successful `verify` clears it; a new edit reopens it. +#[derive(Debug, Default)] +pub(super) struct VerificationDebt { + open: bool, +} + +impl VerificationDebt { + pub(super) fn open(&mut self) { + self.open = true; + } + + pub(super) fn clear(&mut self) { + self.open = false; + } + + pub(super) fn is_open(&self) -> bool { + self.open + } +} + +/// Run every configured verification command exactly as configured and +/// report whether all of them passed, with one observation per command. +pub(super) fn run_verification( + commands: &[VerificationCommand], + max_observation_bytes: usize, + executor: &Executor, + cwd: &Path, + cancellation: &CancellationToken, + on_execution: &mut F, +) -> io::Result<(bool, Vec)> +where + F: FnMut(&str, usize, &Execution), +{ + if commands.is_empty() { + return Ok(( + false, + vec!["verify rejected: no verification command is configured".to_owned()], + )); + } + let mut passed = true; + let mut observations = Vec::new(); + for command in commands { + let refs: Vec<&str> = command.args.iter().map(String::as_str).collect(); + let execution = executor.run_cancellable(cwd, &command.program, &refs, cancellation)?; + on_execution(&command.program, command.args.len(), &execution); + let command_passed = + execution.status == Some(0) && !execution.timed_out && !execution.cancelled; + passed &= command_passed; + observations.push(truncate( + format!( + "verify program={} status={:?} timeout={} cancelled={} stdout={} stderr={}", + command.program, + execution.status, + execution.timed_out, + execution.cancelled, + execution.stdout, + execution.stderr + ), + max_observation_bytes, + )); + } + Ok((passed, observations)) +} diff --git a/src/event_log.rs b/src/event_log.rs new file mode 100644 index 0000000..faa0bdb --- /dev/null +++ b/src/event_log.rs @@ -0,0 +1,90 @@ +//! Line format of the append-only event log: task state lines and the +//! metadata-only trace lines. Traces record counts and sizes, never argument +//! values or command output. + +use crate::{ExecutionTrace, State}; + +pub(crate) fn parse_trace(line: &str) -> Option<(String, ExecutionTrace)> { + let mut parts = line.split('\t'); + if parts.next()? != "trace" { + return None; + } + parse_trace_fields(parts) +} + +pub(crate) fn parse_tool_trace(line: &str) -> Option<(String, ExecutionTrace)> { + let mut parts = line.split('\t'); + if parts.next()? != "tool" { + return None; + } + parse_trace_fields(parts) +} + +fn parse_trace_fields<'a>( + mut parts: impl Iterator, +) -> Option<(String, ExecutionTrace)> { + let id = parts.next()?.to_owned(); + let attempt = parts.next()?.parse().ok()?; + let program = parts.next()?.to_owned(); + let argument_count = parts.next()?.parse().ok()?; + let status_text = parts.next()?; + let status = if status_text == "none" { + None + } else { + Some(status_text.parse().ok()?) + }; + let timed_out = parts.next()?.parse().ok()?; + let cancelled = parts.next()?.parse().ok()?; + let output_truncated = parts.next()?.parse().ok()?; + let stdout_bytes = parts.next()?.parse().ok()?; + let stderr_bytes = parts.next()?.parse().ok()?; + let duration_ms = parts.next()?.parse().ok()?; + Some(( + id, + ExecutionTrace { + attempt, + program, + argument_count, + status, + timed_out, + cancelled, + output_truncated, + stdout_bytes, + stderr_bytes, + duration_ms, + }, + )) +} + +pub(crate) fn trace_state(trace: &ExecutionTrace) -> State { + if trace.cancelled { + State::Cancelled + } else if trace.status == Some(0) && !trace.timed_out { + State::Succeeded + } else { + State::Failed + } +} + +impl State { + pub(crate) fn as_str(self) -> &'static str { + match self { + Self::Queued => "queued", + Self::Running => "running", + Self::Succeeded => "succeeded", + Self::Failed => "failed", + Self::Cancelled => "cancelled", + } + } + + pub(crate) fn parse(value: &str) -> Option { + match value { + "queued" => Some(Self::Queued), + "running" => Some(Self::Running), + "succeeded" => Some(Self::Succeeded), + "failed" => Some(Self::Failed), + "cancelled" => Some(Self::Cancelled), + _ => None, + } + } +} diff --git a/src/lib.rs b/src/lib.rs index 76c1e39..1e1df37 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -1,9 +1,11 @@ pub mod agent; +mod event_log; pub mod executor; pub mod provider; pub mod worker; pub mod workspace; +use event_log::{parse_tool_trace, parse_trace, trace_state}; use executor::{CancellationToken, Execution, Executor}; use std::{ @@ -388,90 +390,6 @@ impl Runtime { } } -fn parse_trace(line: &str) -> Option<(String, ExecutionTrace)> { - let mut parts = line.split('\t'); - if parts.next()? != "trace" { - return None; - } - parse_trace_fields(parts) -} - -fn parse_tool_trace(line: &str) -> Option<(String, ExecutionTrace)> { - let mut parts = line.split('\t'); - if parts.next()? != "tool" { - return None; - } - parse_trace_fields(parts) -} - -fn parse_trace_fields<'a>( - mut parts: impl Iterator, -) -> Option<(String, ExecutionTrace)> { - let id = parts.next()?.to_owned(); - let attempt = parts.next()?.parse().ok()?; - let program = parts.next()?.to_owned(); - let argument_count = parts.next()?.parse().ok()?; - let status_text = parts.next()?; - let status = if status_text == "none" { - None - } else { - Some(status_text.parse().ok()?) - }; - let timed_out = parts.next()?.parse().ok()?; - let cancelled = parts.next()?.parse().ok()?; - let output_truncated = parts.next()?.parse().ok()?; - let stdout_bytes = parts.next()?.parse().ok()?; - let stderr_bytes = parts.next()?.parse().ok()?; - let duration_ms = parts.next()?.parse().ok()?; - Some(( - id, - ExecutionTrace { - attempt, - program, - argument_count, - status, - timed_out, - cancelled, - output_truncated, - stdout_bytes, - stderr_bytes, - duration_ms, - }, - )) -} - -fn trace_state(trace: &ExecutionTrace) -> State { - if trace.cancelled { - State::Cancelled - } else if trace.status == Some(0) && !trace.timed_out { - State::Succeeded - } else { - State::Failed - } -} - -impl State { - fn as_str(self) -> &'static str { - match self { - Self::Queued => "queued", - Self::Running => "running", - Self::Succeeded => "succeeded", - Self::Failed => "failed", - Self::Cancelled => "cancelled", - } - } - fn parse(value: &str) -> Option { - match value { - "queued" => Some(Self::Queued), - "running" => Some(Self::Running), - "succeeded" => Some(Self::Succeeded), - "failed" => Some(Self::Failed), - "cancelled" => Some(Self::Cancelled), - _ => None, - } - } -} - #[cfg(test)] mod tests { use super::*; diff --git a/src/provider.rs b/src/provider/agy.rs similarity index 82% rename from src/provider.rs rename to src/provider/agy.rs index 54d0c3f..3813584 100644 --- a/src/provider.rs +++ b/src/provider/agy.rs @@ -1,7 +1,8 @@ +use super::{DecisionRequest, ModelDecision, ModelProvider}; use crate::executor::{ isolate_process_group, spawn_pipe_reader, terminate_child, CancellationToken, }; -use serde::{Deserialize, Serialize}; +use serde::Deserialize; use std::{ io::{self, BufRead, BufReader}, path::{Path, PathBuf}, @@ -18,47 +19,6 @@ const ACTION_SCHEMA: &str = r#"{"type":"object","properties":{"kind":{"type":"st const DEFAULT_MAX_OUTPUT_BYTES: usize = 1024 * 1024; -#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] -#[serde(tag = "kind", rename_all = "snake_case")] -pub enum ModelAction { - RunTool { - program: String, - #[serde(default)] - args: Vec, - }, - ReadFile { - path: String, - }, - ReplaceText { - path: String, - expected: String, - replacement: String, - }, - /// Ask the runtime to run the configured verification command exactly as - /// configured. The model cannot supply or alter the command. - Verify, - Finish { - summary: String, - }, -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct DecisionRequest { - pub task: String, - pub allowed_programs: Vec, - pub verification_programs: Vec, - pub observations: Vec, -} - -#[derive(Debug, Clone, PartialEq, Eq)] -pub struct ModelDecision { - pub action: ModelAction, - pub model: String, - pub duration_ms: u128, - pub input_tokens: Option, - pub output_tokens: Option, -} - #[derive(Debug, Clone, PartialEq, Eq)] pub struct StreamResult { pub model: String, @@ -76,22 +36,6 @@ struct ProcessOutput { cancelled: bool, } -pub trait ModelProvider { - fn decide(&mut self, request: &DecisionRequest) -> io::Result; - - /// Like [`ModelProvider::decide`], but cancellation from the runtime is - /// propagated to the active provider process. The default implementation - /// ignores the token for providers that do not spawn a child process. - fn decide_cancellable( - &mut self, - request: &DecisionRequest, - cancellation: &CancellationToken, - ) -> io::Result { - let _ = cancellation; - self.decide(request) - } -} - pub struct AgyProvider { binary: PathBuf, working_dir: PathBuf, @@ -376,37 +320,49 @@ impl AgyProvider { let mut succeeded = false; let mut emitted = 0usize; let mut exceeded = false; - for line in BufReader::new(stdout).lines() { - let line = line?; - let event: serde_json::Value = serde_json::from_str(&line) - .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; - if event.get("event").and_then(|value| value.as_str()) == Some("step_update") { - if let Some(delta) = event - .pointer("/step_update/text_delta") - .and_then(|value| value.as_str()) - { - emitted = emitted.saturating_add(delta.len()); - if emitted > self.max_output_bytes { - exceeded = true; - break; + let read_result: io::Result<()> = (|| { + for line in BufReader::new(stdout).lines() { + let line = line?; + let event: serde_json::Value = serde_json::from_str(&line) + .map_err(|error| io::Error::new(io::ErrorKind::InvalidData, error))?; + if event.get("event").and_then(|value| value.as_str()) == Some("step_update") { + if let Some(delta) = event + .pointer("/step_update/text_delta") + .and_then(|value| value.as_str()) + { + emitted = emitted.saturating_add(delta.len()); + if emitted > self.max_output_bytes { + exceeded = true; + break; + } + on_delta(delta); } - on_delta(delta); + } + if event.get("event").and_then(|value| value.as_str()) == Some("result") { + succeeded = event + .pointer("/result/status") + .and_then(|value| value.as_str()) + == Some("SUCCESS"); + usage = Some(AgyUsage { + input_tokens: event + .pointer("/result/usage/input_tokens") + .and_then(|value| value.as_u64()), + output_tokens: event + .pointer("/result/usage/output_tokens") + .and_then(|value| value.as_u64()), + }); } } - if event.get("event").and_then(|value| value.as_str()) == Some("result") { - succeeded = event - .pointer("/result/status") - .and_then(|value| value.as_str()) - == Some("SUCCESS"); - usage = Some(AgyUsage { - input_tokens: event - .pointer("/result/usage/input_tokens") - .and_then(|value| value.as_u64()), - output_tokens: event - .pointer("/result/usage/output_tokens") - .and_then(|value| value.as_u64()), - }); - } + Ok(()) + })(); + if let Err(error) = read_result { + // A malformed stream must not leave the provider running or + // unreaped until the watchdog fires (possibly after pid reuse). + let _ = terminate_child(&mut child); + let _ = child.wait(); + finished.store(true, Ordering::SeqCst); + let _ = watchdog.join(); + return Err(error); } if exceeded { @@ -462,28 +418,6 @@ impl ModelProvider for AgyProvider { } } -impl Serialize for DecisionRequest { - fn serialize(&self, serializer: S) -> Result - where - S: serde::Serializer, - { - #[derive(Serialize)] - struct View<'a> { - task: &'a str, - allowed_programs: &'a [String], - verification_programs: &'a [String], - observations: &'a [String], - } - View { - task: &self.task, - allowed_programs: &self.allowed_programs, - verification_programs: &self.verification_programs, - observations: &self.observations, - } - .serialize(serializer) - } -} - #[derive(Deserialize)] struct AgyEnvelope { status: String, @@ -500,6 +434,7 @@ struct AgyUsage { #[cfg(test)] mod tests { use super::*; + use crate::provider::ModelAction; fn sh_provider(dir: &Path, timeout: Duration, max_output_bytes: usize) -> AgyProvider { AgyProvider::new(Path::new("/bin/sh"), dir, "test-model", timeout) @@ -723,4 +658,83 @@ mod tests { assert_eq!(output.status, Some(0)); assert!(started.elapsed() < Duration::from_secs(2)); } + + #[cfg(unix)] + #[test] + fn provider_503_is_a_reported_failure_not_a_decision() { + let root = workspace("fake-503"); + let script = fake_script( + &root, + "#!/bin/sh\nprintf '%s\\n' '503 The service is currently unavailable' >&2\nexit 1\n", + ); + let mut provider = AgyProvider::new( + Path::new("/bin/true"), + &root, + "fake-model", + Duration::from_secs(5), + ) + .unwrap() + .with_fake_script(&script) + .unwrap(); + let error = provider.decide(&request()).unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::Other); + let message = error.to_string(); + assert!(message.contains("status Some(1)"), "{message}"); + assert!(message.contains("503"), "{message}"); + } + + #[cfg(unix)] + #[test] + fn provider_error_envelope_is_a_reported_failure() { + let root = workspace("fake-error-envelope"); + let script = fake_script( + &root, + "#!/bin/sh\nprintf '%s\\n' '{\"status\":\"ERROR\"}'\n", + ); + let mut provider = AgyProvider::new( + Path::new("/bin/true"), + &root, + "fake-model", + Duration::from_secs(5), + ) + .unwrap() + .with_fake_script(&script) + .unwrap(); + let error = provider.decide(&request()).unwrap_err(); + assert!(error.to_string().contains("ERROR"), "{error}"); + } + + #[cfg(unix)] + #[test] + fn malformed_stream_terminates_and_reaps_the_provider() { + let root = workspace("stream-malformed"); + let pid_file = root.join("provider.pid"); + let script = fake_script( + &root, + &format!( + "#!/bin/sh\necho $$ > {}\nprintf 'not json\\n'\nexec sleep 30\n", + pid_file.display() + ), + ); + let mut provider = AgyProvider::new( + Path::new("/bin/true"), + &root, + "fake-model", + Duration::from_secs(20), + ) + .unwrap() + .with_fake_script(&script) + .unwrap(); + let started = Instant::now(); + let error = provider.stream_text("hello", |_| {}).unwrap_err(); + assert_eq!(error.kind(), io::ErrorKind::InvalidData); + assert!(started.elapsed() < Duration::from_secs(5)); + let pid: libc::pid_t = std::fs::read_to_string(&pid_file) + .unwrap() + .trim() + .parse() + .unwrap(); + // Signal 0 probes existence: a reaped, killed child no longer exists. + assert_eq!(unsafe { libc::kill(pid, 0) }, -1); + } } diff --git a/src/provider/decision.rs b/src/provider/decision.rs new file mode 100644 index 0000000..48de1ed --- /dev/null +++ b/src/provider/decision.rs @@ -0,0 +1,64 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum ModelAction { + RunTool { + program: String, + #[serde(default)] + args: Vec, + }, + ReadFile { + path: String, + }, + ReplaceText { + path: String, + expected: String, + replacement: String, + }, + /// Ask the runtime to run the configured verification command exactly as + /// configured. The model cannot supply or alter the command. + Verify, + Finish { + summary: String, + }, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct DecisionRequest { + pub task: String, + pub allowed_programs: Vec, + pub verification_programs: Vec, + pub observations: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ModelDecision { + pub action: ModelAction, + pub model: String, + pub duration_ms: u128, + pub input_tokens: Option, + pub output_tokens: Option, +} + +impl Serialize for DecisionRequest { + fn serialize(&self, serializer: S) -> Result + where + S: serde::Serializer, + { + #[derive(Serialize)] + struct View<'a> { + task: &'a str, + allowed_programs: &'a [String], + verification_programs: &'a [String], + observations: &'a [String], + } + View { + task: &self.task, + allowed_programs: &self.allowed_programs, + verification_programs: &self.verification_programs, + observations: &self.observations, + } + .serialize(serializer) + } +} diff --git a/src/provider/mock.rs b/src/provider/mock.rs new file mode 100644 index 0000000..871bf24 --- /dev/null +++ b/src/provider/mock.rs @@ -0,0 +1,92 @@ +//! Deterministic scripted provider. It performs no I/O and spawns no process, +//! so tests can drive the agent loop and worker with an exact action sequence, +//! including provider failures at a chosen step. + +use super::{DecisionRequest, ModelAction, ModelDecision, ModelProvider}; +use std::{collections::VecDeque, io}; + +/// Replays a fixed script, one entry per `decide` call, and records every +/// request it receives. An exhausted script is an error rather than a panic. +#[derive(Debug, Default)] +pub struct MockProvider { + script: VecDeque>, + requests: Vec, +} + +impl MockProvider { + /// A provider that returns each action in order. + pub fn from_actions(actions: impl IntoIterator) -> Self { + Self { + script: actions.into_iter().map(Ok).collect(), + requests: Vec::new(), + } + } + + /// A provider whose script may include failures, e.g. a service outage. + pub fn from_script(script: impl IntoIterator>) -> Self { + Self { + script: script.into_iter().collect(), + requests: Vec::new(), + } + } + + /// Requests seen so far, in call order. + pub fn requests(&self) -> &[DecisionRequest] { + &self.requests + } + + /// Script entries not yet consumed. + pub fn remaining(&self) -> usize { + self.script.len() + } +} + +impl ModelProvider for MockProvider { + fn decide(&mut self, request: &DecisionRequest) -> io::Result { + self.requests.push(request.clone()); + let action = self + .script + .pop_front() + .unwrap_or_else(|| Err(io::Error::other("mock provider script exhausted")))?; + Ok(ModelDecision { + action, + model: "mock".to_owned(), + duration_ms: 0, + input_tokens: None, + output_tokens: None, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn request() -> DecisionRequest { + DecisionRequest { + task: "t".to_owned(), + allowed_programs: Vec::new(), + verification_programs: Vec::new(), + observations: Vec::new(), + } + } + + #[test] + fn replays_in_order_records_requests_and_reports_exhaustion() { + let mut provider = MockProvider::from_script([ + Ok(ModelAction::Verify), + Err(io::Error::other("503 unavailable")), + ]); + assert_eq!( + provider.decide(&request()).unwrap().action, + ModelAction::Verify + ); + assert_eq!( + provider.decide(&request()).unwrap_err().to_string(), + "503 unavailable" + ); + assert_eq!(provider.remaining(), 0); + assert!(provider.decide(&request()).is_err()); + assert_eq!(provider.requests().len(), 3); + } +} diff --git a/src/provider/mod.rs b/src/provider/mod.rs new file mode 100644 index 0000000..50d4a5f --- /dev/null +++ b/src/provider/mod.rs @@ -0,0 +1,33 @@ +//! Provider boundary. +//! +//! [`ModelProvider`] is the only thing the agent loop knows about a model: +//! given a [`DecisionRequest`] it returns one structured [`ModelDecision`]. +//! Everything a provider proposes is still authorized by the runtime executor +//! and workspace editor. [`AgyProvider`] is the included adapter and +//! [`mock::MockProvider`] is a deterministic scripted provider for tests. + +mod agy; +mod decision; +pub mod mock; + +use crate::executor::CancellationToken; +use std::io; + +pub use agy::{AgyProvider, StreamResult}; +pub use decision::{DecisionRequest, ModelAction, ModelDecision}; + +pub trait ModelProvider { + fn decide(&mut self, request: &DecisionRequest) -> io::Result; + + /// Like [`ModelProvider::decide`], but cancellation from the runtime is + /// propagated to the active provider process. The default implementation + /// ignores the token for providers that do not spawn a child process. + fn decide_cancellable( + &mut self, + request: &DecisionRequest, + cancellation: &CancellationToken, + ) -> io::Result { + let _ = cancellation; + self.decide(request) + } +} diff --git a/src/worker.rs b/src/worker.rs index 0ec1d30..a55d8ce 100644 --- a/src/worker.rs +++ b/src/worker.rs @@ -98,23 +98,10 @@ mod tests { use super::*; use crate::{ executor::Executor, - provider::{DecisionRequest, ModelAction, ModelDecision, ModelProvider}, + provider::{mock::MockProvider, ModelAction}, }; use std::{collections::VecDeque, fs, path::PathBuf, time::Duration}; - struct ScriptedProvider(VecDeque); - impl ModelProvider for ScriptedProvider { - fn decide(&mut self, _: &DecisionRequest) -> io::Result { - Ok(ModelDecision { - action: self.0.pop_front().expect("scripted action"), - model: "scripted".to_owned(), - duration_ms: 0, - input_tokens: None, - output_tokens: None, - }) - } - } - fn temp() -> PathBuf { let nonce = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) @@ -145,7 +132,7 @@ mod tests { .unwrap(); let agent = AgentLoop::new(4, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let mut provider = ScriptedProvider(VecDeque::from([ + let mut provider = MockProvider::from_actions(VecDeque::from([ ModelAction::ReplaceText { path: "bug.txt".to_owned(), expected: "bad".to_owned(), @@ -192,7 +179,7 @@ mod tests { .unwrap(); let agent = AgentLoop::new(3, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); - let mut provider = ScriptedProvider(VecDeque::new()); + let mut provider = MockProvider::from_actions([]); let token = CancellationToken::default(); token.cancel(); let mut runtime = Runtime::open(&dir).unwrap(); @@ -237,4 +224,60 @@ mod tests { assert_eq!(reopened.tool_traces("mid-1").len(), 1); assert!(reopened.execution_traces("mid-1").is_empty()); } + + #[test] + fn provider_503_after_a_verified_patch_fails_the_task_but_keeps_the_patch() { + let dir = temp(); + let workspace = workspace_under(&dir); + fs::write(workspace.join("bug.txt"), "bad\n").unwrap(); + let executor = Executor::new( + &workspace, + ["true".to_owned()], + Duration::from_secs(1), + 1024, + ) + .unwrap(); + let agent = + AgentLoop::new(6, vec!["true".to_owned()], vec!["true".to_owned()], 1024).unwrap(); + let mut provider = MockProvider::from_script([ + Ok(ModelAction::ReplaceText { + path: "bug.txt".to_owned(), + expected: "bad".to_owned(), + replacement: "good".to_owned(), + }), + Ok(ModelAction::Verify), + Err(io::Error::other("503 The service is currently unavailable")), + ]); + let mut runtime = Runtime::open(&dir).unwrap(); + let error = run_agent_task( + &mut runtime, + "outage-1", + &agent, + &mut provider, + &executor, + &workspace, + "fix the bug", + &CancellationToken::default(), + ) + .unwrap_err(); + + // The provider error is surfaced unchanged; the worker does not retry. + assert!(error.to_string().contains("503"), "{error}"); + assert_eq!(provider.requests().len(), 3); + assert_eq!(provider.remaining(), 0); + // The task is failed, not succeeded: no finish was ever accepted. + assert_eq!(runtime.task("outage-1").unwrap().state, State::Failed); + // The verified edit stays on disk and the verification run is traced. + assert_eq!( + fs::read_to_string(workspace.join("bug.txt")).unwrap(), + "good\n" + ); + assert_eq!(runtime.tool_traces("outage-1").len(), 1); + drop(runtime); + + let mut reopened = Runtime::open(&dir).unwrap(); + assert_eq!(reopened.task("outage-1").unwrap().state, State::Failed); + // An explicit, bounded retry is still available to the operator. + assert!(reopened.retry("outage-1", 2).unwrap()); + } }