diff --git a/.env.example b/.env.example index 64107dbc..598e59a3 100644 --- a/.env.example +++ b/.env.example @@ -60,13 +60,6 @@ TINYCOMPUTER_LAB_SELF_EMAIL= # TINYCOMPUTER_BROWSER_ENDPOINT=http://127.0.0.1:9222 # The travel fixture's address for browser_fixture. # TINYCOMPUTER_FIXTURE_URL=http://127.0.0.1:8000 -# The attested module the bus runners (scripts/lab, task_live, task_fixture) -# load; scripts/build-module builds it and prints this path. -# TINYCOMPUTER_MODULE=target/lab/libtinycomputer.dylib -# Where scripts/build-module installs it instead of /lab: the loader -# refuses a module under a directory another user can write, so the Docker -# runners use $HOME/.local/lib/tinycomputer. -# TINYCOMPUTER_MODULE_DIR= # The planner model task_live asks for. # TINYCOMPUTER_PLANNER_MODEL= # The reasoning model task_live rescues a failed step with (default diff --git a/crates/tinycomputer-browser/src/error/error_tests.rs b/crates/tinycomputer-browser/src/error/error_tests.rs index 26e68e7c..520ef5b9 100644 --- a/crates/tinycomputer-browser/src/error/error_tests.rs +++ b/crates/tinycomputer-browser/src/error/error_tests.rs @@ -149,6 +149,8 @@ fn a_stale_ref_envelope_matches_the_desktop_recovery() { .suggestion .is_some_and(|s| s.contains("BrowserSnapshot")) ); + // agent-browser could not resolve the ref, so nothing reached the page: + // retrying with a fresh ref cannot repeat an effect. assert_eq!(envelope.disposition.retry, RetryDisposition::Safe); } @@ -158,6 +160,10 @@ fn a_timeout_may_have_reached_the_page() { assert_eq!(envelope.code, "TIMEOUT"); assert_eq!(envelope.disposition.delivery, DeliveryDisposition::Unknown); assert!(envelope.suggestion.is_none()); + // A click may already have landed: inspect before repeating it. + let hint = envelope.recovery.expect("a timeout has a way out"); + assert_eq!(hint.strategy, "inspect_state_then_retry_original"); + assert!(hint.requires_fresh_snapshot); } #[test] @@ -173,12 +179,36 @@ fn a_refused_navigation_says_not_to_retry() { .suggestion .is_some_and(|s| s.starts_with("do not retry")) ); + // The domain filter refuses before any navigation reaches the page. assert_eq!( envelope.disposition.delivery, DeliveryDisposition::NotDelivered ); } +#[test] +fn only_failures_decided_before_the_page_claim_nothing_was_delivered() { + let before_the_page = |error: &Error| { + matches!( + error, + Error::NoSuchSession { .. } + | Error::NoSuchOutput { .. } + | Error::StaleRef { .. } + | Error::BlockedByPolicy { .. } + ) + }; + for error in every_variant() { + let expected = before_the_page(&error); + let name = error.to_string(); + let delivery = error.envelope().disposition.delivery; + assert_eq!( + delivery == DeliveryDisposition::NotDelivered, + expected, + "{name}" + ); + } +} + #[test] fn a_lost_connection_suggests_a_new_session() { let envelope = Error::connection_lost("closed").envelope(); diff --git a/crates/tinycomputer-browser/src/error/mod.rs b/crates/tinycomputer-browser/src/error/mod.rs index 0f4c7bd8..8f57d105 100644 --- a/crates/tinycomputer-browser/src/error/mod.rs +++ b/crates/tinycomputer-browser/src/error/mod.rs @@ -160,9 +160,12 @@ impl Error { /// The code is [`errors::code`] of the wire name — the desktop's spelling /// where the meaning is shared — the recovery hint is /// [`errors::recovery`]'s, and the full wire name rides in - /// `details.name` for a host that matches on it. A failure refused before - /// the browser was asked anything is marked not delivered, so a caller - /// knows retrying it cannot repeat an effect. + /// `details.name` for a host that matches on it. Only a failure decided + /// before anything reaches the page is marked not delivered, so a caller + /// knows retrying it cannot repeat an effect: an unknown session or + /// output (local lookups), an unresolvable ref, or a refused origin + /// (rejected inside agent-browser before any input is sent). Every other + /// failure's delivery stays unknown. /// /// # Examples /// @@ -218,17 +221,21 @@ impl Error { } } - /// Whether this failure was decided before any command reached the - /// browser. + /// Whether this failure is always decided before anything reaches the + /// page, so repeating the call cannot repeat an effect: a session or an + /// output looked up locally and not found, a ref agent-browser could not + /// resolve, and a destination its domain filter refused — both of those + /// are rejected inside agent-browser before any input is sent to the page. + /// + /// Invalid input and a limit can also come back after work was done — a + /// capture that turned out too large — so their delivery stays unknown. fn refused_before_delivery(&self) -> bool { matches!( self, - Self::InvalidInput { .. } - | Self::NoSuchSession { .. } + Self::NoSuchSession { .. } + | Self::NoSuchOutput { .. } | Self::StaleRef { .. } | Self::BlockedByPolicy { .. } - | Self::NoSuchOutput { .. } - | Self::LimitExceeded { .. } ) } diff --git a/crates/tinycomputer-browser/src/reply/mod.rs b/crates/tinycomputer-browser/src/reply/mod.rs index efa6334a..f423325c 100644 --- a/crates/tinycomputer-browser/src/reply/mod.rs +++ b/crates/tinycomputer-browser/src/reply/mod.rs @@ -82,7 +82,12 @@ pub(crate) fn classify(message: &str) -> Error { { return Error::invalid_input(message); } - if lower.starts_with("evaluation error") || lower.contains("dialog is blocking the page") { + // A dialog in front of the page passes once someone answers it, so the + // command is not wrong — the page is not ready for it yet. + if lower.contains("dialog is blocking the page") { + return Error::not_actionable(message); + } + if lower.starts_with("evaluation error") { return Error::page(message); } Error::failed(message) diff --git a/crates/tinycomputer-browser/src/reply/reply_tests.rs b/crates/tinycomputer-browser/src/reply/reply_tests.rs index 968fe15c..104bc229 100644 --- a/crates/tinycomputer-browser/src/reply/reply_tests.rs +++ b/crates/tinycomputer-browser/src/reply/reply_tests.rs @@ -91,7 +91,7 @@ fn engine_messages_map_to_what_the_caller_should_do() { }), ( "A JavaScript confirm dialog is blocking the page: \"Leave?\"", - |e| matches!(e, Error::PageError { .. }), + |e| matches!(e, Error::NotActionable { .. }), ), ("something unexpected", |e| { matches!(e, Error::ModuleFailed { .. }) diff --git a/crates/tinycomputer-browser/src/sessions/mod.rs b/crates/tinycomputer-browser/src/sessions/mod.rs index f70f6711..9aba478e 100644 --- a/crates/tinycomputer-browser/src/sessions/mod.rs +++ b/crates/tinycomputer-browser/src/sessions/mod.rs @@ -6,7 +6,7 @@ use std::collections::HashMap; use std::path::PathBuf; -use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; use serde_json::{Value, json}; @@ -32,6 +32,20 @@ pub struct Browser { outputs: Mutex, scratch: PathBuf, counter: AtomicU64, + /// Sessions being launched: they count against [`MAX_SESSIONS`] from + /// the moment the limit is checked, so two concurrent opens cannot both + /// see room for the last slot. + opening: AtomicUsize, +} + +/// One reserved launch slot, given back when the launch ends — inserted, +/// failed, or dropped mid-launch by a cancelled caller. +struct Reservation<'a>(&'a AtomicUsize); + +impl Drop for Reservation<'_> { + fn drop(&mut self) { + self.0.fetch_sub(1, Ordering::SeqCst); + } } struct Session { @@ -102,6 +116,7 @@ impl Browser { outputs: Mutex::new(OutputStore::default()), scratch, counter: AtomicU64::new(0), + opening: AtomicUsize::new(0), } } @@ -112,11 +127,18 @@ impl Browser { /// [`Error::LimitExceeded`] when [`MAX_SESSIONS`] are open, and whatever /// the engine reports when the browser cannot be launched or reached. pub async fn open_session(&self, options: SessionOptions) -> Result { - if self.lock_sessions()?.len() >= MAX_SESSIONS { - return Err(Error::LimitExceeded { - message: format!("at most {MAX_SESSIONS} browser sessions may be open"), - }); - } + // Checked and reserved under the table's lock, so the check and the + // claim are one step for every concurrent caller. + let reservation = { + let sessions = self.lock_sessions()?; + if sessions.len() + self.opening.load(Ordering::SeqCst) >= MAX_SESSIONS { + return Err(Error::LimitExceeded { + message: format!("at most {MAX_SESSIONS} browser sessions may be open"), + }); + } + self.opening.fetch_add(1, Ordering::SeqCst); + Reservation(&self.opening) + }; let id = SessionId::new(format!("s-{}", self.next())); let mut session = Session { engine: self.launcher.open(id.as_str()), @@ -140,8 +162,12 @@ impl Browser { session.info.url = page.url; session.info.title = page.title; let info = session.info.clone(); - self.lock_sessions()? - .insert(id, Arc::new(tokio::sync::Mutex::new(session))); + // The slot moves from the reservation to the table in one step under + // the lock, so no concurrent check ever counts this launch twice. + let mut sessions = self.lock_sessions()?; + sessions.insert(id, Arc::new(tokio::sync::Mutex::new(session))); + drop(reservation); + drop(sessions); Ok(info) } diff --git a/crates/tinycomputer-browser/src/sessions/sessions_tests/lifecycle_tests.rs b/crates/tinycomputer-browser/src/sessions/sessions_tests/lifecycle_tests.rs index 759c56ec..b81ddbbc 100644 --- a/crates/tinycomputer-browser/src/sessions/sessions_tests/lifecycle_tests.rs +++ b/crates/tinycomputer-browser/src/sessions/sessions_tests/lifecycle_tests.rs @@ -81,3 +81,64 @@ fn the_default_scratch_space_is_private_to_the_process() { let browser = Browser::new(Arc::new(Fake::new())); assert!(format!("{browser:?}").contains(&std::process::id().to_string())); } + +/// An engine that yields before every reply, so concurrent launches +/// interleave the way real ones do. +#[derive(Debug)] +struct Slow; + +impl crate::engine::Engine for Slow { + fn execute(&mut self, command: serde_json::Value) -> crate::engine::Reply<'_> { + Box::pin(async move { + tokio::task::yield_now().await; + crate::fake::default_reply(&command) + }) + } +} + +impl crate::engine::Launcher for Slow { + fn open(&self, _session: &str) -> Box { + Box::new(Slow) + } +} + +#[tokio::test] +async fn concurrent_opens_never_exceed_the_cap() { + let browser = Arc::new(Browser::with_scratch(Arc::new(Slow), scratch("race"))); + let mut opens = tokio::task::JoinSet::new(); + for _ in 0..MAX_SESSIONS + 4 { + let browser = browser.clone(); + opens.spawn(async move { + browser + .open_session(SessionOptions { + endpoint: Some("ws://127.0.0.1:9222".to_owned()), + ..SessionOptions::default() + }) + .await + }); + } + let mut refused = 0; + while let Some(opened) = opens.join_next().await { + if matches!(opened.unwrap(), Err(Error::LimitExceeded { .. })) { + refused += 1; + } + } + assert_eq!(refused, 4); + assert_eq!(browser.list_sessions().await.unwrap().len(), MAX_SESSIONS); +} + +#[tokio::test] +async fn a_failed_launch_gives_its_slot_back() { + let fake = + Fake::scripted(|command| (command["action"] == "launch").then(|| failure("Chrome exited"))); + let browser = Browser::with_scratch(Arc::new(fake), scratch("slot")); + // Every launch fails; none may be refused for want of a slot, which is + // what leaked reservations would cause after MAX_SESSIONS attempts. + for _ in 0..MAX_SESSIONS + 2 { + let error = browser + .open_session(SessionOptions::default()) + .await + .unwrap_err(); + assert!(!matches!(error, Error::LimitExceeded { .. }), "{error:?}"); + } +} diff --git a/crates/tinycomputer-bus/src/browser/errors/errors_tests.rs b/crates/tinycomputer-bus/src/browser/errors/errors_tests.rs index 8c268a6d..894eb8cd 100644 --- a/crates/tinycomputer-bus/src/browser/errors/errors_tests.rs +++ b/crates/tinycomputer-bus/src/browser/errors/errors_tests.rs @@ -129,3 +129,29 @@ fn a_refused_navigation_offers_no_way_out() { } assert!(recovery(NO_SUCH_SESSION).is_some_and(|hint| !hint.retryable)); } + +#[test] +fn a_timeout_says_inspect_before_retrying() { + let hint = recovery(TIMEOUT).expect("a timeout has a way out"); + assert_eq!(hint.strategy, "inspect_state_then_retry_original"); + assert!(hint.retryable && hint.requires_fresh_snapshot); +} + +#[test] +fn a_page_error_says_change_the_request_not_repeat_it() { + let hint = recovery(PAGE_ERROR).expect("a page error has a way out"); + assert_eq!(hint.strategy, "inspect_state_then_revise_request"); + assert!(!hint.retryable && hint.requires_fresh_snapshot); +} + +#[test] +fn every_recoverable_name_has_a_recovery_hint() { + for name in NAMES { + if is_agent_recoverable(name) { + assert!( + recovery(name).is_some(), + "{name} is recoverable without a hint" + ); + } + } +} diff --git a/crates/tinycomputer-bus/src/browser/errors/mod.rs b/crates/tinycomputer-bus/src/browser/errors/mod.rs index ba840236..2787cb14 100644 --- a/crates/tinycomputer-bus/src/browser/errors/mod.rs +++ b/crates/tinycomputer-bus/src/browser/errors/mod.rs @@ -157,6 +157,13 @@ pub fn code(name: &str) -> &'static str { /// (`refresh_snapshot_then_retry_original`), so an agent recovering from a /// stale ref does not care which surface it was on. /// +/// A timeout, which may have come after the action reached the page, is +/// `inspect_state_then_retry_original`, never a blind retry: a click or a +/// form submission may already have happened, and repeating it could +/// duplicate the effect. A page error — a thrown script, a rejected command — +/// is `inspect_state_then_revise_request` and not retryable: the same request +/// fails the same way again. +/// /// # Examples /// /// ``` @@ -176,8 +183,10 @@ pub fn recovery(name: &str) -> Option { match name { STALE_REF => Some(hint("refresh_snapshot_then_retry_original", true, true)), NO_SUCH_ELEMENT => Some(hint("refresh_snapshot_then_choose_again", true, true)), - NOT_ACTIONABLE => Some(hint("inspect_state_then_retry_original", true, true)), - TIMEOUT => Some(hint("retry_original", true, false)), + NOT_ACTIONABLE | TIMEOUT => Some(hint("inspect_state_then_retry_original", true, true)), + // A thrown script or a rejected command fails the same way again: + // look at the page and change the request rather than repeat it. + PAGE_ERROR => Some(hint("inspect_state_then_revise_request", false, true)), INVALID_INPUT => Some(hint("fix_request_then_retry", false, false)), NO_SUCH_SESSION => Some(hint("open_session_then_retry_original", false, false)), NO_SUCH_OUTPUT => Some(hint("capture_again_then_read", false, false)), diff --git a/crates/tinycomputer-bus/src/browser/names/mod.rs b/crates/tinycomputer-bus/src/browser/names/mod.rs index e60c47e6..b24a4658 100644 --- a/crates/tinycomputer-bus/src/browser/names/mod.rs +++ b/crates/tinycomputer-bus/src/browser/names/mod.rs @@ -92,7 +92,8 @@ pub mod methods { pub const SCREENSHOT: &str = "BrowserScreenshot"; /// Reads one chunk of a held output: a screenshot this interface took, - /// or one a task view or report names. + /// or one a task's `TaskReport.artifacts` names (task views never carry + /// one). /// /// Takes a [`crate::browser::ReadOutputRequest`]; its `data` is an /// [`crate::browser::OutputChunk`]. diff --git a/crates/tinycomputer-bus/src/browser/names/names_tests.rs b/crates/tinycomputer-bus/src/browser/names/names_tests.rs index 6706b190..bd372c3e 100644 --- a/crates/tinycomputer-bus/src/browser/names/names_tests.rs +++ b/crates/tinycomputer-bus/src/browser/names/names_tests.rs @@ -10,8 +10,10 @@ use super::{INTERFACE, METHODS, OBJECT_PATH, methods}; #[test] fn browser_members_share_the_module_interface() { - assert_eq!(INTERFACE, crate::names::INTERFACE); - assert_eq!(OBJECT_PATH, crate::names::OBJECT_PATH); + // Literals, not the aliases they are defined from: a change to the + // published identity must be made here on purpose. + assert_eq!(INTERFACE, "ai.tinyhumans.tinycomputer.Desktop"); + assert_eq!(OBJECT_PATH, "/ai/tinyhumans/tinycomputer/Desktop"); } #[test] diff --git a/crates/tinycomputer-bus/src/catalogue/catalogue_tests.rs b/crates/tinycomputer-bus/src/catalogue/catalogue_tests.rs index b0e0dcbb..253ced91 100644 --- a/crates/tinycomputer-bus/src/catalogue/catalogue_tests.rs +++ b/crates/tinycomputer-bus/src/catalogue/catalogue_tests.rs @@ -36,6 +36,49 @@ fn every_browser_member_is_in_the_browser_family_and_nothing_else_is() { } } +#[test] +fn every_task_and_browser_name_has_a_catalogue_entry() { + for name in crate::agent::names::METHODS { + assert_eq!( + member(name).map(|entry| entry.family), + Some(Family::Task), + "{name}" + ); + } + for name in crate::browser::names::METHODS { + assert_eq!( + member(name).map(|entry| entry.family), + Some(Family::Browser), + "{name}" + ); + } +} + +#[test] +fn the_flow_family_is_exactly_the_flow_members() { + use crate::names::methods; + let flow = [ + methods::RESOLVE_INTENT, + methods::RUN_GOAL, + methods::RUN_FLOW, + methods::VALIDATE_FLOW, + methods::FLOW_GUIDE, + ]; + for name in flow { + assert_eq!( + member(name).map(|entry| entry.family), + Some(Family::Flow), + "{name}" + ); + } + let catalogued = MEMBERS + .iter() + .filter(|entry| entry.family == Family::Flow) + .map(|entry| entry.name) + .collect::>(); + assert_eq!(catalogued, flow); +} + #[test] fn task_confidentiality_agrees_with_the_task_names() { for entry in MEMBERS.iter().filter(|entry| entry.family == Family::Task) { @@ -53,6 +96,12 @@ fn every_summary_is_one_sentence() { for entry in MEMBERS { assert!(entry.summary.ends_with('.'), "{} summary", entry.name); assert!(entry.summary.len() < 120, "{} summary is long", entry.name); + let body = &entry.summary[..entry.summary.len() - 1]; + assert!( + !body.contains(". ") && !body.contains("? ") && !body.contains("! "), + "{} summary is more than one sentence", + entry.name + ); } } diff --git a/crates/tinycomputer-engine/Cargo.toml b/crates/tinycomputer-engine/Cargo.toml index bb062b84..d6863d0e 100644 --- a/crates/tinycomputer-engine/Cargo.toml +++ b/crates/tinycomputer-engine/Cargo.toml @@ -40,7 +40,9 @@ default = [] planner = ["dep:tinyinference-llm"] [dev-dependencies] -tokio = { workspace = true } +# `test-util` for paused time: task tests step past the controller's +# capture and await timeouts without waiting for them in real time. +tokio = { workspace = true, features = ["test-util"] } [lints] workspace = true diff --git a/crates/tinycomputer-engine/src/lib.rs b/crates/tinycomputer-engine/src/lib.rs index 4ab4c30d..48b57f79 100644 --- a/crates/tinycomputer-engine/src/lib.rs +++ b/crates/tinycomputer-engine/src/lib.rs @@ -49,7 +49,9 @@ pub use planner::{ }; pub use rescue::{Briefing, Guidance, MAX_RESCUE_STEPS, MAX_RESCUES, Rescuer, SCREEN_CHARS}; pub use shape::{Harvest, RECORDS_CHARS, Shaper}; -pub use task::{FlowFuture, FlowRunner, MAX_AWAIT_MS, MAX_TASKS, Tasks, TextFuture, capabilities}; +pub use task::{ + CaptureFuture, FlowFuture, FlowRunner, MAX_AWAIT_MS, MAX_TASKS, Tasks, TextFuture, capabilities, +}; pub use tinycomputer_bus::DesktopResponse; use tinycomputer_desktop::Desktop; pub use workspace::{BROWSER, Workspace}; diff --git a/crates/tinycomputer-engine/src/task/artifact.rs b/crates/tinycomputer-engine/src/task/artifact.rs new file mode 100644 index 00000000..d7808467 --- /dev/null +++ b/crates/tinycomputer-engine/src/task/artifact.rs @@ -0,0 +1,54 @@ +//! The screenshot a stopped run leaves: taken before the task's surfaces are +//! released, and kept in `TaskReport.artifacts` only. +//! +//! It never goes on a status. A `TaskView` also travels through `AwaitTask` +//! and `ListTasks`, which are not confidential, and an output handle is a +//! bearer token for `BrowserReadOutput`: a screenshot of a filled traveller +//! or payment form must only be reachable through the confidential report. + +use std::time::Duration; + +use tinycomputer_bus::agent::TaskStatus; + +use super::FlowRunner; +use super::store::Cell; + +/// How long a capture may take before the task goes on without it: a +/// screenshot is evidence, never a reason to hold a task's final state back. +pub(super) const CAPTURE_TIMEOUT: Duration = Duration::from_secs(10); + +/// `status`, unchanged, after keeping a screenshot of the task's surface for +/// `TaskReport.artifacts` — when the status is one a caller acts on (a +/// checkpoint, an approval, a person's turn, or the end) and the runner can +/// take one. `needs_input` and `needs_plan` take none: both are decided +/// before a run starts, so nothing on screen is the task's doing yet. +pub(super) async fn captured( + cell: &Cell, + runner: &dyn FlowRunner, + status: TaskStatus, +) -> TaskStatus { + if !matches!( + status, + TaskStatus::Running | TaskStatus::NeedsInput { .. } | TaskStatus::NeedsPlan { .. } + ) { + let _kept = capture(cell, runner).await; + } + status +} + +/// Takes a screenshot of the task's surface, keeps it for the report, and +/// returns it; `None` when the runner has none to give in time. +pub(super) async fn capture( + cell: &Cell, + runner: &dyn FlowRunner, +) -> Option { + let id = cell.view.borrow().id.clone(); + let shot = tokio::time::timeout(CAPTURE_TIMEOUT, runner.capture(&id)) + .await + .ok() + .flatten()?; + if let Ok(mut state) = cell.state.lock() { + state.artifacts.push(shot.clone()); + } + Some(shot) +} diff --git a/crates/tinycomputer-engine/src/task/budget.rs b/crates/tinycomputer-engine/src/task/budget.rs index 38e302df..cdb6039c 100644 --- a/crates/tinycomputer-engine/src/task/budget.rs +++ b/crates/tinycomputer-engine/src/task/budget.rs @@ -6,6 +6,7 @@ use std::collections::BTreeMap; use tinycomputer_bus::RunFlowRequest; use tinycomputer_bus::agent::{TaskConstraints, TaskStatus}; +use super::artifact::captured; use super::brief::brief; use super::names::fact_names; use super::publish::{publish, stopped_summary}; @@ -74,7 +75,9 @@ pub(super) fn run_request( /// Ends a task outright with `status`, without a flow run to interpret: its /// time budget ran out before a run of it could even start, or a run of it /// had to be cut off mid-flight to keep from spending past what remains. -pub(super) fn stop_task(cell: &Cell, runner: &dyn FlowRunner, status: TaskStatus) { +pub(super) async fn stop_task(cell: &Cell, runner: &dyn FlowRunner, status: TaskStatus) { + // A task cut off by its time budget still leaves its last screen. + let status = captured(cell, runner, status).await; let summary = stopped_summary(&status); publish(cell, status, &summary); runner.release(&cell.view.borrow().id); diff --git a/crates/tinycomputer-engine/src/task/controller.rs b/crates/tinycomputer-engine/src/task/controller.rs index 7c54a9de..f9ed1381 100644 --- a/crates/tinycomputer-engine/src/task/controller.rs +++ b/crates/tinycomputer-engine/src/task/controller.rs @@ -315,7 +315,7 @@ impl Tasks { flow: Some(state.flow.clone()), steps: state.steps.clone(), records: records(&state.reads), - artifacts: Vec::new(), + artifacts: state.artifacts.clone(), learned: state.learned.clone(), trace: state.exchanges.clone(), rescues: state.rescues.clone(), diff --git a/crates/tinycomputer-engine/src/task/describe.rs b/crates/tinycomputer-engine/src/task/describe.rs index 0936e55c..31673942 100644 --- a/crates/tinycomputer-engine/src/task/describe.rs +++ b/crates/tinycomputer-engine/src/task/describe.rs @@ -193,11 +193,10 @@ fn members() -> Vec { "id": {"type": "string"}, "trace": { "type": "boolean", - "default": true, - "description": "include every Jev exchange StartTask.trace recorded; false keeps the report small" + "description": "include every Jev exchange StartTask.trace recorded; false keeps the report small. Required: a confidential body of only {\"id\"} is refused by the TinyBus client as a stream handle" } }), - &["id"], + &["id", "trace"], ), "TaskReport", ), diff --git a/crates/tinycomputer-engine/src/task/drive.rs b/crates/tinycomputer-engine/src/task/drive.rs index cbbb8d2c..e2220654 100644 --- a/crates/tinycomputer-engine/src/task/drive.rs +++ b/crates/tinycomputer-engine/src/task/drive.rs @@ -7,6 +7,7 @@ use std::time::{Duration, Instant}; use tinycomputer_bus::agent::TaskStatus; +use super::artifact::{capture, captured}; use super::budget::{elapsed_budget_failed, run_request, stop_task}; use super::human::human_wall; use super::interpret::{Next, finished, run_outcome}; @@ -96,7 +97,7 @@ pub(super) async fn drive(cell: Arc, runner: Arc, runs: Ve return; }; if time_left == Some(0) { - stop_task(&cell, runner.as_ref(), elapsed_budget_failed()); + stop_task(&cell, runner.as_ref(), elapsed_budget_failed()).await; return; } let id = cell.view.borrow().id.clone(); @@ -104,7 +105,7 @@ pub(super) async fn drive(cell: Arc, runner: Arc, runs: Ve let run_call = runner.run(&id, &constraints, request); let reply = if let Some(ms) = time_left { let Ok(reply) = tokio::time::timeout(Duration::from_millis(ms), run_call).await else { - stop_task(&cell, runner.as_ref(), elapsed_budget_failed()); + stop_task(&cell, runner.as_ref(), elapsed_budget_failed()).await; return; }; reply @@ -152,6 +153,7 @@ pub(super) async fn drive(cell: Arc, runner: Arc, runs: Ve continue; } }; + let status = captured(&cell, runner.as_ref(), status).await; let summary = redacted.redact(&stopped_summary(&status)); let ended = matches!( status, @@ -170,6 +172,8 @@ pub(super) async fn drive(cell: Arc, runner: Arc, runs: Ve /// Every run finished: the task is done, with what it read, shaped as its /// `output` asks when it asks. pub(super) async fn finish(cell: &Cell, runner: &dyn FlowRunner) { + // The last look at the surface, before shaping and before release. + let _ = capture(cell, runner).await; let (answer, records, harvest) = { let Ok(state) = cell.state.lock() else { return; diff --git a/crates/tinycomputer-engine/src/task/mod.rs b/crates/tinycomputer-engine/src/task/mod.rs index 0b257f00..fe15bdae 100644 --- a/crates/tinycomputer-engine/src/task/mod.rs +++ b/crates/tinycomputer-engine/src/task/mod.rs @@ -33,8 +33,10 @@ //! `store`; the background run in `drive`, capped by `budget` and briefed by //! `brief`; answering a paused task in `resume`; rescuing a failed step in //! `recovery`; spotting a wall only a person can pass in `human`; and what -//! the caller sees in `publish`. +//! the caller sees in `publish`; the screenshots a stopped run leaves in +//! `artifact`. +mod artifact; mod brief; mod budget; mod controller; @@ -53,6 +55,7 @@ use std::future::Future; use std::pin::Pin; use tinycomputer_bus::agent::{TaskConstraints, TaskId}; +use tinycomputer_bus::browser::OutputRef; use tinycomputer_bus::{DesktopResponse, RunFlowRequest}; pub use controller::Tasks; @@ -65,6 +68,10 @@ pub type FlowFuture = Pin + Send>>; /// The future [`FlowRunner::visible_text`] returns. pub type TextFuture = Pin> + Send>>; +/// The future [`FlowRunner::capture`] returns: a held screenshot, if the +/// task's surface could take one. +pub type CaptureFuture = Pin> + Send>>; + /// Runs a task's flows, on surfaces that live as long as the task. pub trait FlowRunner: Send + Sync + 'static { /// Runs `request` for `task` within `constraints`, returning `RunFlow`'s @@ -83,6 +90,20 @@ pub trait FlowRunner: Send + Sync + 'static { Box::pin(async { Vec::new() }) } + /// A screenshot of the task's surface as it stands, held for the caller + /// to read. Taken whenever a run stops on its own — at a checkpoint, + /// before an approval, at a person's turn, finished, failed, or cut off + /// by its time budget — before the task's surfaces are released, so a + /// finished or failed task still leaves one. Not for `needs_input` or + /// `needs_plan`: those are decided before a run starts, with nothing on + /// screen yet that the task did. + /// `CancelTask` releases at once without one: the caller chose to stop, + /// and can take its own with `BrowserScreenshot` first. `None` by + /// default, and whenever the surface cannot take one. + fn capture(&self, _task: &TaskId) -> CaptureFuture { + Box::pin(async { None }) + } + /// Lets go of whatever the task held, once it has ended. fn release(&self, _task: &TaskId) {} } diff --git a/crates/tinycomputer-engine/src/task/store.rs b/crates/tinycomputer-engine/src/task/store.rs index db082f36..458cadfa 100644 --- a/crates/tinycomputer-engine/src/task/store.rs +++ b/crates/tinycomputer-engine/src/task/store.rs @@ -9,6 +9,7 @@ use tinycomputer_bus::agent::{ AgentResponse, Rescue, StartTaskRequest, TaskBudget, TaskConstraints, TaskId, TaskOutput, TaskStatus, TaskView, }; +use tinycomputer_bus::browser::OutputRef; use tinycomputer_bus::{FLOW_GUIDE, Flow, GroundingHint, JevExchange, StepReport}; use tinycomputer_core::Facts; use tokio::sync::watch; @@ -56,6 +57,8 @@ pub(super) struct State { pub(super) rescues: Vec, /// The shape the caller wants the answer in, if any. pub(super) output: Option, + /// Screenshots taken each time a run stopped, oldest first. + pub(super) artifacts: Vec, } /// A task's cumulative spend against its [`TaskBudget`], across every run. @@ -111,6 +114,7 @@ impl Tasks { spent: Spent::default(), rescues: Vec::new(), output: request.output.clone(), + artifacts: Vec::new(), }), worker: Mutex::new(None), rescuer: self.rescuer.clone(), diff --git a/crates/tinycomputer-engine/src/task/task_tests.rs b/crates/tinycomputer-engine/src/task/task_tests.rs index b3024d30..6800b48d 100644 --- a/crates/tinycomputer-engine/src/task/task_tests.rs +++ b/crates/tinycomputer-engine/src/task/task_tests.rs @@ -8,6 +8,7 @@ #![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] mod approval_tests; +mod artifact_tests; mod describe_tests; mod errors_tests; mod human_tests; @@ -45,6 +46,12 @@ struct Script { screen: Mutex>, /// Tasks let go of, in order. released: Mutex>, + /// The screenshot `capture` hands back, if any. + shot: Mutex>, + /// A capture that never answers, like a hung surface. + stuck: std::sync::atomic::AtomicBool, + /// `capture` and `release` calls, in the order they arrived. + events: Mutex>, } impl FlowRunner for Script { @@ -69,7 +76,17 @@ impl FlowRunner for Script { Box::pin(async move { texts }) } + fn capture(&self, _task: &TaskId) -> super::CaptureFuture { + self.events.lock().unwrap().push("capture"); + let shot = self.shot.lock().unwrap().clone(); + if self.stuck.load(std::sync::atomic::Ordering::SeqCst) { + return Box::pin(std::future::pending()); + } + Box::pin(async move { shot }) + } + fn release(&self, task: &TaskId) { + self.events.lock().unwrap().push("release"); self.released.lock().unwrap().push(task.clone()); } } diff --git a/crates/tinycomputer-engine/src/task/task_tests/artifact_tests.rs b/crates/tinycomputer-engine/src/task/task_tests/artifact_tests.rs new file mode 100644 index 00000000..63d688f5 --- /dev/null +++ b/crates/tinycomputer-engine/src/task/task_tests/artifact_tests.rs @@ -0,0 +1,202 @@ +//! Tests for the screenshot a stopped run leaves: on the status that has room +//! for one, and in the report, taken before the task's surfaces are released. + +use super::*; + +use tinycomputer_bus::browser::{OutputId, OutputRef}; + +fn shot() -> OutputRef { + OutputRef { + id: OutputId::new("o-1"), + total_bytes: 3, + sha256: "abc".to_owned(), + media_type: "image/png".to_owned(), + width: 1, + height: 1, + } +} + +fn with_shot(replies: Vec) -> (Tasks, Arc