返回 CodeWhale
vm.rs
1 //! The sandboxed QuickJS VM that executes Workflow scripts.
2 //!
3 //! Threading model (design §2.2): `rquickjs` contexts and every `'js` value
4 //! are `!Send`, so each run gets a dedicated OS thread with its own
5 //! current-thread tokio reactor. Host functions do no heavy work inline —
6 //! only `Send` data (JSON strings, [`TaskRequest`]s, oneshot replies) crosses
7 //! to the driver; conversion back into JS values happens on the VM thread
8 //! after the await resolves.
9 //!
10 //! Sandbox: the context registers only standard ECMAScript intrinsics plus
11 //! the Workflow globals (`task`, `parallel`, `pipeline`, `log`, `phase`,
12 //! `budget`, `args`). There is no module loader, no fs/net/process access,
13 //! and `Date`/`Math.random` are overridden to throw so recorded runs stay
14 //! deterministic for replay.
15
16 use std::cell::Cell;
17 use std::env;
18 use std::future::Future;
19 use std::rc::Rc;
20 use std::sync::atomic::{AtomicBool, Ordering};
21 use std::sync::{Arc, OnceLock};
22
23 use rquickjs::function::{Async, Func};
24 use rquickjs::{AsyncContext, AsyncRuntime, CatchResultExt, CaughtError, Ctx, Promise, Value};
25 use serde::Deserialize;
26 use tokio::runtime::Handle;
27 use tokio::sync::{OwnedSemaphorePermit, Semaphore, oneshot, watch};
28 use tokio::task::JoinSet;
29
30 use crate::driver::{ProgressEvent, TaskCompletion, TaskRequest, WorkflowDriver};
31 use crate::error::{TaskError, TaskErrorKind, WorkflowJsError};
32 use crate::schema::{
33 ReplyDecodeError, SCHEMA_REPAIR_MAX_ATTEMPTS, carried_raw, compile_schema, decode_reply,
34 repair_prompt,
35 };
36 use crate::{
37 CODEMODE_MAX_TOOL_CALLS, PARALLEL_MAX_ITEMS, ToolCallRequest, ToolInvoker,
38 WORKFLOW_LIFETIME_CAP, normalize_profile,
39 };
40
41 const DEFAULT_VM_MEMORY_LIMIT_BYTES: usize = 32 * 1024 * 1024;
42 const MIN_VM_MEMORY_LIMIT_BYTES: usize = 4 * 1024 * 1024;
43 const MAX_VM_MEMORY_LIMIT_BYTES: usize = 512 * 1024 * 1024;
44 const DEFAULT_VM_STACK_BYTES: usize = 1024 * 1024;
45 const MIN_VM_STACK_BYTES: usize = 128 * 1024;
46 const MAX_VM_STACK_BYTES: usize = 8 * 1024 * 1024;
47 const DEFAULT_VM_THREAD_STACK_BYTES: usize = 2 * 1024 * 1024;
48 const MIN_VM_THREAD_STACK_BYTES: usize = 512 * 1024;
49 const MAX_VM_THREAD_STACK_BYTES: usize = 16 * 1024 * 1024;
50 const DEFAULT_MAX_CONCURRENT_VMS: usize = 4;
51 const MAX_CONCURRENT_VMS: usize = 256;
52
53 const VM_MEMORY_LIMIT_MB_ENV: &str = "CODEWHALE_WORKFLOW_JS_MEMORY_LIMIT_MB";
54 const VM_STACK_KB_ENV: &str = "CODEWHALE_WORKFLOW_JS_STACK_KB";
55 const VM_THREAD_STACK_KB_ENV: &str = "CODEWHALE_WORKFLOW_JS_THREAD_STACK_KB";
56 const VM_MAX_CONCURRENT_ENV: &str = "CODEWHALE_WORKFLOW_JS_MAX_CONCURRENT";
57
58 /// Resource limits applied to the QuickJS runtime before any script runs.
59 ///
60 /// There is deliberately no wall-clock timeout here: cancellation (dropping
61 /// the run future, or the driver's cancel cascade) is the deadline mechanism.
62 #[derive(Debug, Clone, Copy)]
63 pub struct VmLimits {
64 /// QuickJS heap ceiling in bytes (default 32 MiB).
65 pub memory_limit_bytes: usize,
66 /// Maximum interpreter stack in bytes (default 1 MiB).
67 pub max_stack_bytes: usize,
68 }
69
70 impl Default for VmLimits {
71 fn default() -> Self {
72 Self::from_env()
73 }
74 }
75
76 impl VmLimits {
77 pub fn from_env() -> Self {
78 Self {
79 memory_limit_bytes: env_usize_bytes(
80 VM_MEMORY_LIMIT_MB_ENV,
81 1024 * 1024,
82 MIN_VM_MEMORY_LIMIT_BYTES,
83 MAX_VM_MEMORY_LIMIT_BYTES,
84 DEFAULT_VM_MEMORY_LIMIT_BYTES,
85 ),
86 max_stack_bytes: env_usize_bytes(
87 VM_STACK_KB_ENV,
88 1024,
89 MIN_VM_STACK_BYTES,
90 MAX_VM_STACK_BYTES,
91 DEFAULT_VM_STACK_BYTES,
92 ),
93 }
94 }
95 }
96
97 fn env_usize_bytes(name: &str, unit: usize, min: usize, max: usize, default: usize) -> usize {
98 env::var(name)
99 .ok()
100 .and_then(|raw| raw.parse::<usize>().ok())
101 .and_then(|value| value.checked_mul(unit))
102 .map(|bytes| bytes.clamp(min, max))
103 .unwrap_or(default)
104 }
105
106 fn max_concurrent_vms() -> usize {
107 env::var(VM_MAX_CONCURRENT_ENV)
108 .ok()
109 .and_then(|raw| raw.parse::<usize>().ok())
110 .map(|value| value.clamp(1, MAX_CONCURRENT_VMS))
111 .unwrap_or(DEFAULT_MAX_CONCURRENT_VMS)
112 }
113
114 fn vm_thread_stack_bytes() -> usize {
115 env_usize_bytes(
116 VM_THREAD_STACK_KB_ENV,
117 1024,
118 MIN_VM_THREAD_STACK_BYTES,
119 MAX_VM_THREAD_STACK_BYTES,
120 DEFAULT_VM_THREAD_STACK_BYTES,
121 )
122 }
123
124 fn vm_admission() -> &'static Arc<Semaphore> {
125 static ADMISSION: OnceLock<Arc<Semaphore>> = OnceLock::new();
126 ADMISSION.get_or_init(|| Arc::new(Semaphore::new(max_concurrent_vms())))
127 }
128
129 /// Executes Workflow scripts, one isolated QuickJS runtime per run.
130 ///
131 /// Every [`WorkflowVm::run_script`] call spins up a fresh interpreter on a
132 /// dedicated thread, so runs share nothing (globals, heap, interned atoms)
133 /// and a wedged script can never stall a sibling run.
134 #[derive(Debug, Clone, Default)]
135 pub struct WorkflowVm {
136 limits: VmLimits,
137 }
138
139 impl WorkflowVm {
140 /// A VM with the default [`VmLimits`].
141 pub fn new() -> Self {
142 Self::default()
143 }
144
145 /// A VM with explicit resource limits.
146 pub fn with_limits(limits: VmLimits) -> Self {
147 Self { limits }
148 }
149
150 /// Run one Workflow script to completion.
151 ///
152 /// * `source` is the script body; it is wrapped in an async function, so
153 /// top-level `await` and `return` both work. The returned value is the
154 /// script's `return` value, JSON-encoded (`undefined` becomes `null`).
155 /// * `args` is exposed verbatim to the script as the `args` global.
156 /// * `driver` executes `task()` spawns and receives progress events. A
157 /// driver instance is scoped to exactly one run: `cancel_all` is always
158 /// invoked at run teardown (success, script error, or cancellation), so
159 /// stray children never outlive the script that spawned them.
160 ///
161 /// Cancellation cascade (design §9): dropping the returned future cancels
162 /// the run — the interrupt handler aborts executing JS, pending `task()`
163 /// awaits resolve to errors, and `driver.cancel_all()` is invoked
164 /// immediately from the dropping thread.
165 pub async fn run_script(
166 &self,
167 source: &str,
168 args: serde_json::Value,
169 driver: Arc<dyn WorkflowDriver>,
170 ) -> Result<serde_json::Value, WorkflowJsError> {
171 self.run_script_with_cancel(source, args, driver, WorkflowRunCancel::new())
172 .await
173 }
174
175 /// Like [`Self::run_script`], but accepts an external cancel handle so the
176 /// host can interrupt the VM without dropping the run future.
177 pub async fn run_script_with_cancel(
178 &self,
179 source: &str,
180 args: serde_json::Value,
181 driver: Arc<dyn WorkflowDriver>,
182 cancel: WorkflowRunCancel,
183 ) -> Result<serde_json::Value, WorkflowJsError> {
184 self.run_inner(source, args, driver, None, cancel).await
185 }
186
187 /// Run a code-mode script: the same VM, plus the `tools.call()` host
188 /// binding backed by `invoker`. Workflow runs (no invoker) never see the
189 /// binding, so the Workflow sandbox keeps its documented host surface.
190 pub async fn run_tools_script(
191 &self,
192 source: &str,
193 args: serde_json::Value,
194 driver: Arc<dyn WorkflowDriver>,
195 invoker: Arc<dyn ToolInvoker>,
196 cancel: WorkflowRunCancel,
197 ) -> Result<serde_json::Value, WorkflowJsError> {
198 self.run_inner(source, args, driver, Some(invoker), cancel)
199 .await
200 }
201
202 async fn run_inner(
203 &self,
204 source: &str,
205 args: serde_json::Value,
206 driver: Arc<dyn WorkflowDriver>,
207 invoker: Option<Arc<dyn ToolInvoker>>,
208 cancel: WorkflowRunCancel,
209 ) -> Result<serde_json::Value, WorkflowJsError> {
210 let args_json = serde_json::to_string(&args)
211 .map_err(|err| WorkflowJsError::InvalidArgs(err.to_string()))?;
212 let cancel = cancel.0;
213 let (result_tx, result_rx) = oneshot::channel();
214 let mut guard = RunGuard {
215 cancel: cancel.clone(),
216 driver: driver.clone(),
217 armed: true,
218 };
219
220 let permit = vm_admission()
221 .clone()
222 .acquire_owned()
223 .await
224 .map_err(|_| WorkflowJsError::VmInit("VM admission gate closed".to_string()))?;
225 let limits = self.limits;
226 let source = source.to_string();
227 let thread_driver = driver.clone();
228 let thread_cancel = cancel.clone();
229 let thread_invoker = invoker.clone();
230 // Native admission belongs to the originating host executor. The VM
231 // thread carries only Send requests/replies, never the Child Engine's
232 // Rust future on its bounded interpreter stack.
233 let host_runtime = Handle::try_current().map_err(|_| {
234 WorkflowJsError::VmInit("workflow host Tokio runtime is unavailable".to_string())
235 })?;
236 let spawned = std::thread::Builder::new()
237 .name("workflow-js-vm".to_string())
238 .stack_size(vm_thread_stack_bytes())
239 .spawn(move || {
240 let _permit: OwnedSemaphorePermit = permit;
241 let outcome = vm_thread_main(
242 source,
243 args_json,
244 thread_driver.clone(),
245 thread_cancel,
246 thread_invoker,
247 limits,
248 host_runtime,
249 );
250 // Run teardown: this driver is scoped to one run, so any task
251 // still in flight is unreachable now — cancel the cascade.
252 thread_driver.cancel_all();
253 let _ = result_tx.send(outcome);
254 });
255 if let Err(err) = spawned {
256 guard.armed = false;
257 return Err(WorkflowJsError::VmInit(format!(
258 "failed to spawn VM thread: {err}"
259 )));
260 }
261
262 match result_rx.await {
263 Ok(outcome) => {
264 // The VM thread has already torn down and cancelled children.
265 guard.armed = false;
266 outcome
267 }
268 // VM thread panicked before reporting; leave the guard armed so
269 // its drop (right now, at return) cancels outstanding tasks.
270 Err(_) => Err(WorkflowJsError::VmTerminated(
271 "VM thread exited without reporting a result".to_string(),
272 )),
273 }
274 }
275 }
276
277 /// Cooperative cancel signal shared by the run future (guard side) and the VM
278 /// thread. The atomic flag feeds the QuickJS interrupt handler (sync, called
279 /// mid-bytecode); the watch channel wakes host futures parked on driver
280 /// completions.
281 #[derive(Clone)]
282 pub struct WorkflowRunCancel(CancelHandle);
283
284 impl WorkflowRunCancel {
285 #[must_use]
286 pub fn new() -> Self {
287 Self(CancelHandle::new())
288 }
289
290 pub fn cancel(&self) {
291 self.0.cancel();
292 }
293 }
294
295 impl Default for WorkflowRunCancel {
296 fn default() -> Self {
297 Self::new()
298 }
299 }
300
301 #[derive(Clone)]
302 struct CancelHandle {
303 flag: Arc<AtomicBool>,
304 tx: Arc<watch::Sender<bool>>,
305 }
306
307 impl CancelHandle {
308 fn new() -> Self {
309 let (tx, _rx) = watch::channel(false);
310 Self {
311 flag: Arc::new(AtomicBool::new(false)),
312 tx: Arc::new(tx),
313 }
314 }
315
316 fn cancel(&self) {
317 self.flag.store(true, Ordering::SeqCst);
318 self.tx.send_replace(true);
319 }
320
321 fn is_cancelled(&self) -> bool {
322 self.flag.load(Ordering::SeqCst)
323 }
324
325 async fn cancelled(&self) {
326 let mut rx = self.tx.subscribe();
327 let _ = rx.wait_for(|cancelled| *cancelled).await;
328 }
329
330 fn flag_arc(&self) -> Arc<AtomicBool> {
331 self.flag.clone()
332 }
333 }
334
335 /// Fires the cancel cascade if the caller drops the run future before the VM
336 /// reports a result.
337 struct RunGuard {
338 cancel: CancelHandle,
339 driver: Arc<dyn WorkflowDriver>,
340 armed: bool,
341 }
342
343 impl Drop for RunGuard {
344 fn drop(&mut self) {
345 if self.armed {
346 self.cancel.cancel();
347 self.driver.cancel_all();
348 }
349 }
350 }
351
352 fn vm_thread_main(
353 source: String,
354 args_json: String,
355 driver: Arc<dyn WorkflowDriver>,
356 cancel: CancelHandle,
357 invoker: Option<Arc<dyn ToolInvoker>>,
358 limits: VmLimits,
359 host_runtime: Handle,
360 ) -> Result<serde_json::Value, WorkflowJsError> {
361 let reactor = tokio::runtime::Builder::new_current_thread()
362 .enable_all()
363 .build()
364 .map_err(|err| WorkflowJsError::VmInit(format!("failed to build VM reactor: {err}")))?;
365 reactor.block_on(run_in_vm(
366 source,
367 args_json,
368 driver,
369 cancel,
370 invoker,
371 limits,
372 host_runtime,
373 ))
374 }
375
376 async fn run_in_vm(
377 source: String,
378 args_json: String,
379 driver: Arc<dyn WorkflowDriver>,
380 cancel: CancelHandle,
381 invoker: Option<Arc<dyn ToolInvoker>>,
382 limits: VmLimits,
383 host_runtime: Handle,
384 ) -> Result<serde_json::Value, WorkflowJsError> {
385 let runtime = AsyncRuntime::new().map_err(|err| WorkflowJsError::VmInit(err.to_string()))?;
386 runtime.set_memory_limit(limits.memory_limit_bytes).await;
387 runtime.set_max_stack_size(limits.max_stack_bytes).await;
388 let interrupt_flag = cancel.flag_arc();
389 runtime
390 .set_interrupt_handler(Some(Box::new(move || {
391 interrupt_flag.load(Ordering::Acquire)
392 })))
393 .await;
394 let context = AsyncContext::full(&runtime)
395 .await
396 .map_err(|err| WorkflowJsError::VmInit(err.to_string()))?;
397
398 let result = context
399 .async_with(async |ctx| {
400 run_in_ctx(
401 ctx,
402 source,
403 args_json,
404 driver,
405 invoker,
406 cancel,
407 host_runtime,
408 )
409 .await
410 })
411 .await;
412 drop(context);
413 runtime.run_gc().await;
414 result
415 }
416
417 async fn run_in_ctx(
418 ctx: Ctx<'_>,
419 source: String,
420 args_json: String,
421 driver: Arc<dyn WorkflowDriver>,
422 invoker: Option<Arc<dyn ToolInvoker>>,
423 cancel: CancelHandle,
424 host_runtime: Handle,
425 ) -> Result<serde_json::Value, WorkflowJsError> {
426 install_host(
427 &ctx,
428 driver,
429 invoker.clone(),
430 cancel.clone(),
431 &args_json,
432 host_runtime,
433 )?;
434 ctx.eval::<(), _>(prelude())
435 .catch(&ctx)
436 .map_err(|err| WorkflowJsError::VmInit(format!("prelude failed: {err}")))?;
437 if invoker.is_some() {
438 ctx.eval::<(), _>(CODEMODE_PRELUDE)
439 .catch(&ctx)
440 .map_err(|err| WorkflowJsError::VmInit(format!("codemode prelude failed: {err}")))?;
441 }
442
443 let desugared = desugar_export_default(&source);
444 let wrapped = format!("(async () => {{\n{desugared}\n}})()");
445 let promise = ctx
446 .eval::<Promise, _>(wrapped)
447 .catch(&ctx)
448 .map_err(|err| script_error(&cancel, err))?;
449 let value = promise
450 .into_future::<Value>()
451 .await
452 .catch(&ctx)
453 .map_err(|err| script_error(&cancel, err))?;
454 js_value_to_json(&ctx, value)
455 }
456
457 /// Rewrite the documented module-style authoring shape
458 /// (`export default async function (args) { ... }`) into the script form the
459 /// VM actually evals. Sources are wrapped in an async IIFE, where the
460 /// module-only `export` keyword is a syntax error, so without this every
461 /// imperative `export default` workflow (including the #4131 dogfood
462 /// fixtures) failed to parse. The default export is captured, invoked with
463 /// the `args` global when it is a function, and its result becomes the run
464 /// result; a non-function default export is returned as-is.
465 fn desugar_export_default(source: &str) -> String {
466 const EXPORT_DEFAULT: &str = "export default";
467 let Some(offset) = line_leading_export_default(source) else {
468 return source.to_string();
469 };
470 let mut out = source.to_string();
471 out.replace_range(
472 offset..offset + EXPORT_DEFAULT.len(),
473 "globalThis.__workflow_default =",
474 );
475 out.push('\n');
476 out.push_str(
477 ";{\n const __wf_default = globalThis.__workflow_default;\n delete globalThis.__workflow_default;\n if (typeof __wf_default === \"function\") {\n return await __wf_default(args);\n }\n if (__wf_default !== undefined) {\n return __wf_default;\n }\n}\n",
478 );
479 out
480 }
481
482 /// Return the byte offset of a line-leading `export default` token that is
483 /// actual JavaScript syntax, not text inside a string, template literal, or
484 /// comment. This intentionally recognizes only the documented authoring shape
485 /// instead of attempting to implement a general JavaScript module parser.
486 fn line_leading_export_default(source: &str) -> Option<usize> {
487 const EXPORT_DEFAULT: &[u8] = b"export default";
488 let bytes = source.as_bytes();
489 let mut idx = 0usize;
490 let mut quote = None;
491 let mut escaped = false;
492 let mut line_comment = false;
493 let mut block_comment = false;
494 let mut line_has_only_whitespace = true;
495
496 while idx < bytes.len() {
497 let byte = bytes[idx];
498
499 if line_comment {
500 if byte == b'\n' {
501 line_comment = false;
502 line_has_only_whitespace = true;
503 }
504 idx += 1;
505 continue;
506 }
507
508 if block_comment {
509 if byte == b'*' && bytes.get(idx + 1) == Some(&b'/') {
510 block_comment = false;
511 line_has_only_whitespace = false;
512 idx += 2;
513 continue;
514 }
515 if byte == b'\n' {
516 line_has_only_whitespace = true;
517 } else if !byte.is_ascii_whitespace() {
518 line_has_only_whitespace = false;
519 }
520 idx += 1;
521 continue;
522 }
523
524 if let Some(active_quote) = quote {
525 if byte == b'\n' {
526 line_has_only_whitespace = true;
527 escaped = false;
528 } else {
529 if !byte.is_ascii_whitespace() {
530 line_has_only_whitespace = false;
531 }
532 if escaped {
533 escaped = false;
534 } else if byte == b'\\' {
535 escaped = true;
536 } else if byte == active_quote {
537 quote = None;
538 }
539 }
540 idx += 1;
541 continue;
542 }
543
544 if byte == b'\n' {
545 line_has_only_whitespace = true;
546 idx += 1;
547 continue;
548 }
549 if line_has_only_whitespace && byte.is_ascii_whitespace() {
550 idx += 1;
551 continue;
552 }
553 if line_has_only_whitespace && bytes[idx..].starts_with(EXPORT_DEFAULT) {
554 return Some(idx);
555 }
556
557 line_has_only_whitespace = false;
558 if byte == b'/' && bytes.get(idx + 1) == Some(&b'/') {
559 line_comment = true;
560 idx += 2;
561 } else if byte == b'/' && bytes.get(idx + 1) == Some(&b'*') {
562 block_comment = true;
563 idx += 2;
564 } else {
565 if matches!(byte, b'\'' | b'"' | b'`') {
566 quote = Some(byte);
567 }
568 idx += 1;
569 }
570 }
571
572 None
573 }
574
575 fn script_error(cancel: &CancelHandle, err: CaughtError<'_>) -> WorkflowJsError {
576 if cancel.is_cancelled() {
577 WorkflowJsError::Cancelled
578 } else {
579 WorkflowJsError::Script(err.to_string())
580 }
581 }
582
583 fn js_value_to_json<'js>(
584 ctx: &Ctx<'js>,
585 value: Value<'js>,
586 ) -> Result<serde_json::Value, WorkflowJsError> {
587 if value.is_undefined() {
588 return Ok(serde_json::Value::Null);
589 }
590 let text = ctx
591 .json_stringify(value)
592 .map_err(|err| WorkflowJsError::ResultEncoding(err.to_string()))?;
593 match text {
594 None => Ok(serde_json::Value::Null),
595 Some(text) => {
596 let text = text
597 .to_string()
598 .map_err(|err| WorkflowJsError::ResultEncoding(err.to_string()))?;
599 serde_json::from_str(&text)
600 .map_err(|err| WorkflowJsError::ResultEncoding(err.to_string()))
601 }
602 }
603 }
604
605 /// Prelude fragment for code-mode runs only: captures `__codemode_call` into
606 /// the frozen `globalThis.tools` surface, then deletes the raw binding.
607 /// Workflow runs never eval this, so their documented host surface is
608 /// unchanged — `tools` simply does not exist there.
609 const CODEMODE_PRELUDE: &str = r#"
610 (() => {
611 const hostCall = __codemode_call;
612 globalThis.tools = Object.freeze({
613 call: async (tool, input) => {
614 const envelope = JSON.parse(
615 await hostCall(JSON.stringify({ tool, input: input === undefined ? {} : input }))
616 );
617 if (envelope.error !== undefined) {
618 const err = new Error(envelope.error);
619 err.kind = envelope.error_kind;
620 throw err;
621 }
622 return envelope.result;
623 },
624 });
625 try { delete globalThis.__codemode_call; } catch (_) { /* frozen shape */ }
626 })();
627 "#;
628
629 // Construct and poll native host futures on the captured host executor.
630 // JoinSet aborts on drop; explicit cancellation also joins that abort before
631 // the VM receives its existing typed cancellation error.
632 async fn call_on_host<T, F>(
633 runtime: &Handle,
634 cancel: &CancelHandle,
635 cancel_error: TaskError,
636 call: impl FnOnce() -> F + Send + 'static,
637 ) -> Result<T, TaskError>
638 where
639 T: Send + 'static,
640 F: Future<Output = T> + Send + 'static,
641 {
642 let mut pending = JoinSet::new();
643 let host_cancel = cancel.clone();
644 let host_error = cancel_error.clone();
645 pending.spawn_on(
646 async move {
647 if host_cancel.is_cancelled() {
648 Err(host_error)
649 } else {
650 Ok(call().await)
651 }
652 },
653 runtime,
654 );
655 tokio::select! {
656 biased;
657 _ = cancel.cancelled() => {
658 pending.shutdown().await;
659 Err(cancel_error)
660 }
661 outcome = pending.join_next() => match outcome {
662 Some(Ok(value)) => value,
663 Some(Err(error)) if error.is_panic() => {
664 std::panic::resume_unwind(error.into_panic())
665 }
666 _ => Err(TaskError::new(
667 TaskErrorKind::Driver,
668 "workflow host callback ended without a reply",
669 )),
670 },
671 }
672 }
673
674 /// The `tools.call()` host call. Infallible at the binding level, like
675 /// `task_host`: outcomes return through the `{result}` /
676 /// `{error, error_kind}` envelope so the prelude rethrows typed errors.
677 /// Gate refusals arrive as admission errors (nothing ran); a nested tool
678 /// that ran and failed arrives as an agent error (work failed).
679 async fn tools_call_host(
680 call_json: String,
681 invoker: Arc<dyn ToolInvoker>,
682 cancel: CancelHandle,
683 invoked: Rc<Cell<u64>>,
684 host_runtime: Handle,
685 ) -> String {
686 let outcome = tools_call_host_inner(call_json, invoker, cancel, invoked, host_runtime).await;
687 let envelope = match outcome {
688 Ok(result) => serde_json::json!({ "result": result }),
689 Err(TaskError { kind, message }) => {
690 serde_json::json!({ "error": message, "error_kind": kind.as_str() })
691 }
692 };
693 envelope.to_string()
694 }
695
696 async fn tools_call_host_inner(
697 call_json: String,
698 invoker: Arc<dyn ToolInvoker>,
699 cancel: CancelHandle,
700 invoked: Rc<Cell<u64>>,
701 host_runtime: Handle,
702 ) -> Result<serde_json::Value, TaskError> {
703 let admission = |message: String| TaskError::new(TaskErrorKind::Admission, message);
704 let request: ToolCallRequest = serde_json::from_str(&call_json).map_err(|err| {
705 admission(format!(
706 "tools.call(): expected {{\"tool\", \"input\"}} JSON: {err}"
707 ))
708 })?;
709 if request.tool.trim().is_empty() {
710 return Err(admission(
711 "tools.call(): `tool` must be a non-empty string".to_string(),
712 ));
713 }
714 if !request.input.is_object() {
715 return Err(admission(
716 "tools.call(): `input` must be a JSON object".to_string(),
717 ));
718 }
719 if invoked.get() >= CODEMODE_MAX_TOOL_CALLS {
720 return Err(admission(format!(
721 "tools.call(): per-run tool-call cap ({CODEMODE_MAX_TOOL_CALLS}) reached"
722 )));
723 }
724 if cancel.is_cancelled() {
725 return Err(TaskError::new(
726 TaskErrorKind::Cancelled,
727 "tools.call(): run cancelled".to_string(),
728 ));
729 }
730 invoked.set(invoked.get() + 1);
731 let response = call_on_host(
732 &host_runtime,
733 &cancel,
734 TaskError::new(TaskErrorKind::Cancelled, "tools.call(): run cancelled"),
735 move || async move { invoker.invoke(request).await },
736 )
737 .await?
738 .map_err(|err| TaskError::new(TaskErrorKind::from(&err), format!("tools.call(): {err}")))?;
739 if response.ok {
740 Ok(response.result)
741 } else {
742 let message = response
743 .result
744 .as_str()
745 .unwrap_or("tools.call(): tool failed without a message")
746 .to_string();
747 Err(TaskError::new(TaskErrorKind::Agent, message))
748 }
749 }
750
751 fn install_host(
752 ctx: &Ctx<'_>,
753 driver: Arc<dyn WorkflowDriver>,
754 invoker: Option<Arc<dyn ToolInvoker>>,
755 cancel: CancelHandle,
756 args_json: &str,
757 host_runtime: Handle,
758 ) -> Result<(), WorkflowJsError> {
759 let globals = ctx.globals();
760
761 let args_value: Value = ctx
762 .json_parse(args_json)
763 .map_err(|err| WorkflowJsError::InvalidArgs(err.to_string()))?;
764 globals.set("args", args_value).map_err(init_err)?;
765
766 // Per-run lifetime counter (design §4.3): counts spawn *attempts*, and the
767 // check + increment happen with no await in between so a parallel burst
768 // cannot slip past the cap on the single-threaded VM.
769 let spawned = Rc::new(Cell::new(0u64));
770
771 let task_driver = driver.clone();
772 let task_cancel = cancel.clone();
773 let task_runtime = host_runtime.clone();
774 globals
775 .set(
776 "__workflow_task",
777 Func::from(Async(move |opts_json: String| {
778 let driver = task_driver.clone();
779 let cancel = task_cancel.clone();
780 let spawned = spawned.clone();
781 let host_runtime = task_runtime.clone();
782 async move { task_host(opts_json, driver, cancel, spawned, host_runtime).await }
783 })),
784 )
785 .map_err(init_err)?;
786
787 // Code-mode runs only: per-run tool-call counter. Same single-threaded
788 // check+increment discipline as `spawned` above — no await between the
789 // cap check and the increment, so a burst cannot slip past it.
790 if let Some(invoker) = invoker {
791 let invoked = Rc::new(Cell::new(0u64));
792 let call_invoker = invoker.clone();
793 let call_cancel = cancel.clone();
794 globals
795 .set(
796 "__codemode_call",
797 Func::from(Async(move |call_json: String| {
798 let invoker = call_invoker.clone();
799 let cancel = call_cancel.clone();
800 let invoked = invoked.clone();
801 let host_runtime = host_runtime.clone();
802 async move {
803 tools_call_host(call_json, invoker, cancel, invoked, host_runtime).await
804 }
805 })),
806 )
807 .map_err(init_err)?;
808 }
809
810 let log_driver = driver.clone();
811 globals
812 .set(
813 "__workflow_log",
814 Func::from(move |message: String| {
815 log_driver.progress(ProgressEvent::Log { message });
816 }),
817 )
818 .map_err(init_err)?;
819
820 // Structured twin of the prelude's "every slot failed" breadcrumb (R9):
821 // a dead fan-out of script-thrown thunks leaves no task record behind,
822 // so the host needs a typed event — not a log line — to keep the run's
823 // terminal status honest.
824 let fanout_driver = driver.clone();
825 globals
826 .set(
827 "__workflow_every_slot_failed",
828 Func::from(move |construct: String, failed: u32, total: u32| {
829 fanout_driver.progress(ProgressEvent::FanoutAllSlotsFailed {
830 construct,
831 failed,
832 total,
833 });
834 }),
835 )
836 .map_err(init_err)?;
837
838 // Structured twin of the per-slot "dropped a failed slot as null"
839 // breadcrumb (R9): a PARTIALLY failed fan-out still resolves and leaves
840 // no task record for the dropped slot, so without this event the ledger
841 // cannot see the loss and records a clean Completed.
842 let fanout_driver = driver.clone();
843 globals
844 .set(
845 "__workflow_slot_dropped",
846 Func::from(move |construct: String, kind: String, slot: u32| {
847 fanout_driver.progress(ProgressEvent::FanoutSlotDropped {
848 construct,
849 kind,
850 slot,
851 });
852 }),
853 )
854 .map_err(init_err)?;
855
856 let phase_driver = driver.clone();
857 globals
858 .set(
859 "__workflow_phase",
860 Func::from(move |title: String| {
861 phase_driver.progress(ProgressEvent::Phase { title });
862 }),
863 )
864 .map_err(init_err)?;
865
866 // Budget reads are live driver snapshots (design §5.2). NaN encodes
867 // "no ceiling" for `total`; the prelude maps it to `null`.
868 let total_driver = driver.clone();
869 globals
870 .set(
871 "__workflow_budget_total",
872 Func::from(move || -> f64 {
873 match total_driver.budget().total {
874 Some(total) => total as f64,
875 None => f64::NAN,
876 }
877 }),
878 )
879 .map_err(init_err)?;
880
881 let spent_driver = driver.clone();
882 globals
883 .set(
884 "__workflow_budget_spent",
885 Func::from(move || -> f64 { spent_driver.budget().spent as f64 }),
886 )
887 .map_err(init_err)?;
888
889 globals
890 .set(
891 "__workflow_budget_remaining",
892 Func::from(move || -> f64 {
893 match driver.budget().remaining() {
894 Some(remaining) => remaining as f64,
895 None => f64::INFINITY,
896 }
897 }),
898 )
899 .map_err(init_err)?;
900
901 Ok(())
902 }
903
904 fn init_err(err: rquickjs::Error) -> WorkflowJsError {
905 WorkflowJsError::VmInit(err.to_string())
906 }
907
908 /// The `task()` host call. Everything that can go wrong is reported through
909 /// the JSON envelope (`{"error": ..., "error_kind": ...}`) so the prelude
910 /// re-throws it as a real JS `Error` with a script-side stack and a typed
911 /// [`TaskErrorKind`] on `.kind` (R9). The kind is assigned here, where the
912 /// failure actually happened — never re-derived from the message text.
913 async fn task_host(
914 opts_json: String,
915 driver: Arc<dyn WorkflowDriver>,
916 cancel: CancelHandle,
917 spawned: Rc<Cell<u64>>,
918 host_runtime: Handle,
919 ) -> String {
920 let outcome = task_host_inner(opts_json, driver, cancel, spawned, host_runtime).await;
921 let envelope = match outcome {
922 Ok(value) => serde_json::json!({ "value": value }),
923 Err(TaskError { kind, message }) => {
924 serde_json::json!({ "error": message, "error_kind": kind.as_str() })
925 }
926 };
927 envelope.to_string()
928 }
929
930 /// Best-effort `label`/`phase` from raw `task()` options, for rejection
931 /// receipts when the options never survived parsing.
932 fn task_identity_hint(opts_json: &str) -> (Option<String>, Option<String>) {
933 let value: serde_json::Value =
934 serde_json::from_str(opts_json).unwrap_or(serde_json::Value::Null);
935 let pluck = |key: &str| {
936 value
937 .get(key)
938 .and_then(serde_json::Value::as_str)
939 .map(str::trim)
940 .filter(|text| !text.is_empty())
941 .map(str::to_string)
942 };
943 (pluck("label"), pluck("phase"))
944 }
945
946 /// Record a pre-spawn `task()` rejection on the host ledger, then hand the
947 /// message back for the JS throw. Rejections that never reach `spawn_task`
948 /// would otherwise be invisible to the run record (#5035's surviving gap).
949 fn reject_task(driver: &Arc<dyn WorkflowDriver>, opts_json: &str, message: String) -> String {
950 let (label, phase) = task_identity_hint(opts_json);
951 driver.progress(ProgressEvent::TaskRejected {
952 label,
953 phase,
954 message: message.clone(),
955 });
956 message
957 }
958
959 /// Emit the terminal schema-failure receipt for a `task()` whose reply failed
960 /// `responseSchema` with no repair left to try, and hand the message back for
961 /// the JS throw. `note` (when present) names why a repair was skipped, so the
962 /// operator can tell "repair refused to run" from "repair also failed".
963 fn fail_schema(
964 driver: &Arc<dyn WorkflowDriver>,
965 task_id: String,
966 attempt: u32,
967 error: &ReplyDecodeError,
968 raw: String,
969 raw_truncated: bool,
970 note: Option<String>,
971 ) -> String {
972 let message = match note {
973 Some(note) => format!("{} (repair skipped: {note})", error.message()),
974 None => error.message().to_string(),
975 };
976 driver.progress(ProgressEvent::TaskSchemaValidationFailed {
977 task_id,
978 kind: error.kind().to_string(),
979 attempt,
980 message: message.clone(),
981 raw,
982 raw_truncated,
983 });
984 message
985 }
986
987 /// Build the repair request for `next_attempt` (#5583): the same child
988 /// identity, budget fields, and schema as the original, with a repair prompt
989 /// carrying the original task, the schema, the failed reply, and why it
990 /// failed, plus the wall clock the first attempt did not spend.
991 fn repair_request(
992 original: &TaskRequest,
993 next_attempt: u32,
994 error: &ReplyDecodeError,
995 failed_raw: &str,
996 wall_time_secs: Option<u64>,
997 ) -> TaskRequest {
998 let mut request = original.clone();
999 let schema = original
1000 .response_schema
1001 .as_ref()
1002 .expect("repair only runs when responseSchema is set");
1003 // The bracket prefix identifies the repair at a glance on progress
1004 // surfaces and lets tests script the repair reply by rule order.
1005 request.description = format!(
1006 "[schema repair {next_attempt}] {}",
1007 repair_prompt(&original.description, schema, failed_raw, error)
1008 );
1009 request.wall_time_secs = wall_time_secs;
1010 // Label and phase stay inherited so progress surfaces group the repair
1011 // with its task.
1012 request
1013 }
1014
1015 /// The one cancellation error: the run's deadline fired. Fatal in every
1016 /// `parallel()` / `pipeline()` mode — a cancelled run must never resolve into
1017 /// a slot value.
1018 fn cancelled_task() -> TaskError {
1019 TaskError::new(TaskErrorKind::Cancelled, "task(): run cancelled")
1020 }
1021
1022 /// A terminal `responseSchema` failure, already receipted by [`fail_schema`].
1023 fn schema_task(message: String) -> TaskError {
1024 TaskError::new(TaskErrorKind::Schema, message)
1025 }
1026
1027 async fn task_host_inner(
1028 opts_json: String,
1029 driver: Arc<dyn WorkflowDriver>,
1030 cancel: CancelHandle,
1031 spawned: Rc<Cell<u64>>,
1032 host_runtime: Handle,
1033 ) -> Result<serde_json::Value, TaskError> {
1034 let admission = |message: String| TaskError::new(TaskErrorKind::Admission, message);
1035 let workspace = driver.workspace_root();
1036 let request = parse_task_options(&opts_json, workspace.as_deref())
1037 .map_err(|message| admission(reject_task(&driver, &opts_json, message)))?;
1038 // Compile the schema before spawning so a malformed one fails fast
1039 // instead of burning a subagent.
1040 let validator = request
1041 .response_schema
1042 .as_ref()
1043 .map(compile_schema)
1044 .transpose()
1045 .map_err(|message| admission(reject_task(&driver, &opts_json, message)))?;
1046
1047 // Lifetime backstop (design §4.3) — checked and bumped before any await.
1048 if spawned.get() >= WORKFLOW_LIFETIME_CAP {
1049 return Err(admission(reject_task(
1050 &driver,
1051 &opts_json,
1052 format!(
1053 "task(): Workflow lifetime agent cap ({WORKFLOW_LIFETIME_CAP}) reached for this run"
1054 ),
1055 )));
1056 }
1057 // Fast-fail budget gate. The authoritative reservation lives in the
1058 // driver (design §5.3); this only stops obviously-doomed spawns early.
1059 let snapshot = driver.budget();
1060 if snapshot.exhausted() {
1061 return Err(TaskError::new(
1062 TaskErrorKind::Budget,
1063 reject_task(
1064 &driver,
1065 &opts_json,
1066 format!(
1067 "task(): budget exhausted ({} of {} tokens spent)",
1068 snapshot.spent,
1069 snapshot.total.unwrap_or(0)
1070 ),
1071 ),
1072 ));
1073 }
1074 if cancel.is_cancelled() {
1075 return Err(cancelled_task());
1076 }
1077
1078 // Bounded schema repair (#5583): after a failed `responseSchema` decode,
1079 // re-ask the same route before throwing. `None` is the default single
1080 // repair; `Some(0)` disables it. The first attempt is attempt 1, so the
1081 // task is schema-terminal once `attempt` reaches this ceiling.
1082 let last_attempt = 1 + request.schema_repair_attempts.unwrap_or(1);
1083 // The wall clock is shared across attempts: a repair inherits the time
1084 // the first attempt did not spend, not a fresh budget.
1085 let started = std::time::Instant::now();
1086 let mut wall_time_secs_left = request.wall_time_secs;
1087 let mut current = request.clone();
1088 let mut attempt: u32 = 0;
1089 loop {
1090 attempt += 1;
1091 spawned.set(spawned.get() + 1);
1092 // Admission can wait (a saturated concurrency gate, a routing call),
1093 // so it races the run's cancel exactly as the completion does below
1094 // and as `tools.call()` does: a cancel that only closed the driver's
1095 // gate used to be the one thing that could wake it.
1096 let task_driver = driver.clone();
1097 let task_request = current.clone();
1098 let spawned_task = call_on_host(
1099 &host_runtime,
1100 &cancel,
1101 cancelled_task(),
1102 move || async move { task_driver.spawn_task(task_request).await },
1103 )
1104 .await?
1105 .map_err(|err| TaskError::new(TaskErrorKind::from(&err), err.to_string()))?;
1106 let task_id = spawned_task.task_id;
1107 let completion_rx = spawned_task.completion;
1108 let completion = tokio::select! {
1109 _ = cancel.cancelled() => return Err(cancelled_task()),
1110 completion = completion_rx => completion.map_err(|_| {
1111 TaskError::new(
1112 TaskErrorKind::Driver,
1113 "task(): driver dropped the completion channel",
1114 )
1115 })?,
1116 };
1117
1118 let text = match completion {
1119 TaskCompletion::Completed { text } => text,
1120 TaskCompletion::Failed { message } => {
1121 return Err(TaskError::new(
1122 TaskErrorKind::Agent,
1123 format!("task(): subagent failed: {message}"),
1124 ));
1125 }
1126 TaskCompletion::Cancelled => {
1127 return Err(TaskError::new(
1128 TaskErrorKind::Cancelled,
1129 "task(): subagent cancelled",
1130 ));
1131 }
1132 TaskCompletion::BudgetExhausted { message } => {
1133 return Err(TaskError::new(
1134 TaskErrorKind::Budget,
1135 format!("task(): budget exhausted: {message}"),
1136 ));
1137 }
1138 };
1139 // Without a schema the raw text is the contract; with one, the decode
1140 // decides — and a failure may still be repaired.
1141 let Some(validator) = validator.as_ref() else {
1142 return Ok(serde_json::Value::String(text));
1143 };
1144 let error = match decode_reply(&text, validator) {
1145 Ok(value) => return Ok(value),
1146 Err(error) => error,
1147 };
1148 let (raw, raw_truncated) = carried_raw(&text);
1149 if attempt >= last_attempt {
1150 return Err(schema_task(fail_schema(
1151 &driver,
1152 task_id,
1153 attempt,
1154 &error,
1155 raw,
1156 raw_truncated,
1157 None,
1158 )));
1159 }
1160 // The attempt failed but a repair remains: record it as a receipt
1161 // (visible even when the repair succeeds), then re-run the admission
1162 // gates — a repair is a real child, not a free retry.
1163 driver.progress(ProgressEvent::TaskSchemaRepairAttempted {
1164 task_id: task_id.clone(),
1165 kind: error.kind().to_string(),
1166 attempt,
1167 message: error.message().to_string(),
1168 raw: raw.clone(),
1169 raw_truncated,
1170 });
1171 if spawned.get() >= WORKFLOW_LIFETIME_CAP {
1172 return Err(schema_task(fail_schema(
1173 &driver,
1174 task_id,
1175 attempt,
1176 &error,
1177 raw,
1178 raw_truncated,
1179 Some(format!(
1180 "workflow lifetime agent cap ({WORKFLOW_LIFETIME_CAP}) reached"
1181 )),
1182 )));
1183 }
1184 let snapshot = driver.budget();
1185 if snapshot.exhausted() {
1186 return Err(schema_task(fail_schema(
1187 &driver,
1188 task_id,
1189 attempt,
1190 &error,
1191 raw,
1192 raw_truncated,
1193 Some("budget exhausted".to_string()),
1194 )));
1195 }
1196 if cancel.is_cancelled() {
1197 return Err(cancelled_task());
1198 }
1199 if let Some(wall) = wall_time_secs_left {
1200 let remaining = wall.saturating_sub(started.elapsed().as_secs());
1201 if remaining == 0 {
1202 return Err(schema_task(fail_schema(
1203 &driver,
1204 task_id,
1205 attempt,
1206 &error,
1207 raw,
1208 raw_truncated,
1209 Some("no wall-time left from wallTimeSecs".to_string()),
1210 )));
1211 }
1212 wall_time_secs_left = Some(remaining);
1213 }
1214 current = repair_request(&request, attempt + 1, &error, &raw, wall_time_secs_left);
1215 }
1216 }
1217
1218 /// JS-facing option names for `task()` (design §3.3). Unknown fields are
1219 /// rejected so a typo (`responseschema`) fails loudly instead of being
1220 /// silently dropped. Every multi-word field also accepts its snake_case
1221 /// spelling, and the `agent` tool's `workspace_policy` name is accepted as an
1222 /// alias for worktree isolation — the two spawn surfaces are written by the
1223 /// same authors (often models), so a schema that runs on one must not be an
1224 /// unknown-field error on the other.
1225 #[derive(Debug, Deserialize)]
1226 #[serde(rename_all = "camelCase", deny_unknown_fields)]
1227 struct TaskOptions {
1228 #[serde(alias = "title")]
1229 description: Option<String>,
1230 prompt: Option<String>,
1231 #[serde(alias = "type", alias = "subagent_type")]
1232 subagent_type: Option<String>,
1233 /// Fleet role name (#4177). Preferred step identity field.
1234 role: Option<String>,
1235 profile: Option<String>,
1236 model: Option<String>,
1237 #[serde(alias = "model_strength")]
1238 model_strength: Option<String>,
1239 thinking: Option<String>,
1240 cwd: Option<String>,
1241 #[serde(default)]
1242 worktree: bool,
1243 /// `agent`-tool alias for worktree isolation: "shared" | "worktree".
1244 #[serde(default, alias = "workspace_policy")]
1245 workspace_policy: Option<String>,
1246 #[serde(alias = "write_authority")]
1247 write_authority: Option<String>,
1248 #[serde(default, alias = "write_roots")]
1249 write_roots: Vec<String>,
1250 #[serde(default, alias = "exact_files")]
1251 exact_files: Vec<String>,
1252 #[serde(default, alias = "coordination_contracts")]
1253 coordination_contracts: Vec<String>,
1254 #[serde(default)]
1255 dependencies: Vec<String>,
1256 #[serde(default)]
1257 acceptance: Vec<String>,
1258 #[serde(alias = "allowed_tools")]
1259 allowed_tools: Option<Vec<String>>,
1260 #[serde(alias = "max_depth")]
1261 max_depth: Option<u32>,
1262 #[serde(alias = "token_budget")]
1263 token_budget: Option<u64>,
1264 #[serde(alias = "max_steps")]
1265 max_steps: Option<u32>,
1266 #[serde(alias = "wall_time_secs")]
1267 wall_time_secs: Option<u64>,
1268 #[serde(alias = "response_schema")]
1269 response_schema: Option<serde_json::Value>,
1270 /// Bounded `responseSchema` repair attempts after a failed decode
1271 /// (#5583): re-ask the same route with the schema and the failed reply.
1272 /// Defaults to one; `0` disables; capped at
1273 /// [`SCHEMA_REPAIR_MAX_ATTEMPTS`] so repair stays a bounded recovery.
1274 #[serde(default, alias = "schema_repair_attempts")]
1275 schema_repair_attempts: Option<u32>,
1276 label: Option<String>,
1277 phase: Option<String>,
1278 }
1279
1280 fn parse_task_options(
1281 opts_json: &str,
1282 workspace: Option<&std::path::Path>,
1283 ) -> Result<TaskRequest, String> {
1284 let mut options: TaskOptions =
1285 serde_json::from_str(opts_json).map_err(|err| format!("task(): invalid options: {err}"))?;
1286 if let Some(policy) = options.workspace_policy.take() {
1287 match policy.trim().to_ascii_lowercase().as_str() {
1288 "worktree" => options.worktree = true,
1289 "shared" => {
1290 if options.worktree {
1291 return Err(
1292 "task(): workspacePolicy 'shared' conflicts with worktree: true"
1293 .to_string(),
1294 );
1295 }
1296 }
1297 other => {
1298 return Err(format!(
1299 "task(): workspacePolicy must be shared or worktree; got {other:?}"
1300 ));
1301 }
1302 }
1303 }
1304 let description = options
1305 .prompt
1306 .or(options.description)
1307 .filter(|description| !description.trim().is_empty())
1308 .ok_or_else(|| "task(): 'description' (or 'prompt') is required".to_string())?;
1309 let role = options
1310 .role
1311 .as_deref()
1312 .map(normalize_profile)
1313 .transpose()
1314 .map_err(|err| format!("task(): role: {err}"))?;
1315 let profile = options
1316 .profile
1317 .as_deref()
1318 .map(normalize_profile)
1319 .transpose()
1320 .map_err(|err| format!("task(): {err}"))?;
1321 options.write_roots = normalize_task_paths("writeRoots", options.write_roots, 32)?;
1322 options.exact_files = normalize_task_paths("exactFiles", options.exact_files, 32)?;
1323 let cwd = options
1324 .cwd
1325 .as_deref()
1326 .map(|cwd| normalize_task_cwd_in(cwd, workspace))
1327 .transpose()?;
1328 options.coordination_contracts =
1329 normalize_task_string_list("coordinationContracts", options.coordination_contracts, 16)?;
1330 options.dependencies = normalize_task_string_list("dependencies", options.dependencies, 8)?;
1331 options.acceptance = normalize_task_string_list("acceptance", options.acceptance, 8)?;
1332 let write_authority = options
1333 .write_authority
1334 .as_deref()
1335 .map(|value| value.trim().to_ascii_lowercase())
1336 .map(|value| match value.as_str() {
1337 "read_only" | "workspace_write" | "worktree_write" => Ok(value),
1338 _ => Err(format!(
1339 "task(): writeAuthority must be read_only, workspace_write, or worktree_write; got {value:?}"
1340 )),
1341 })
1342 .transpose()?;
1343 if write_authority.as_deref() == Some("worktree_write") && !options.worktree {
1344 return Err("task(): writeAuthority worktree_write requires worktree: true".to_string());
1345 }
1346 let role_kind = role.as_deref().and_then(task_role_kind);
1347 let type_kind = options.subagent_type.as_deref().and_then(task_role_kind);
1348 if let (Some(role_kind), Some(type_kind)) = (role_kind, type_kind)
1349 && role_kind != type_kind
1350 {
1351 return Err("task(): role and subagentType declare contradictory authorities".to_string());
1352 }
1353 let declared_kind = role_kind.or(type_kind);
1354 if matches!(declared_kind, Some(TaskRoleKind::ReadOnly))
1355 && write_authority
1356 .as_deref()
1357 .is_some_and(|authority| authority != "read_only")
1358 {
1359 return Err("task(): read-only roles cannot declare write-capable authority".to_string());
1360 }
1361 // A write-capable task with no declared scope is not refused here: the
1362 // spawn boundary (`validate_spawn_write_contract`) claims the workspace
1363 // root, or the task's deliverables, exactly as it does for a plain Agent
1364 // spawn, and the coordination ledger arbitrates contention with live peers.
1365 if let Some(attempts) = options.schema_repair_attempts
1366 && attempts > SCHEMA_REPAIR_MAX_ATTEMPTS
1367 {
1368 return Err(format!(
1369 "task(): schemaRepairAttempts is bounded to {SCHEMA_REPAIR_MAX_ATTEMPTS}; \
1370 repair is a bounded recovery, not a retry loop"
1371 ));
1372 }
1373 Ok(TaskRequest {
1374 description,
1375 subagent_type: options.subagent_type,
1376 role,
1377 profile,
1378 model: options.model,
1379 model_strength: options.model_strength,
1380 thinking: options.thinking,
1381 cwd,
1382 worktree: options.worktree,
1383 write_authority,
1384 write_roots: options.write_roots,
1385 exact_files: options.exact_files,
1386 coordination_contracts: options.coordination_contracts,
1387 dependencies: options.dependencies,
1388 acceptance: options.acceptance,
1389 allowed_tools: options.allowed_tools,
1390 // Host-imposed only: a script cannot set (or clear) a deny list.
1391 disallowed_tools: Vec::new(),
1392 max_depth: options.max_depth,
1393 token_budget: options.token_budget,
1394 max_steps: options.max_steps,
1395 wall_time_secs: options.wall_time_secs,
1396 response_schema: options.response_schema,
1397 schema_repair_attempts: options.schema_repair_attempts,
1398 label: options.label,
1399 phase: options.phase,
1400 })
1401 }
1402
1403 fn normalize_task_string_list(
1404 field: &str,
1405 values: Vec<String>,
1406 limit: usize,
1407 ) -> Result<Vec<String>, String> {
1408 if values.len() > limit {
1409 return Err(format!("task(): {field} accepts at most {limit} entries"));
1410 }
1411 let mut normalized = Vec::new();
1412 for value in values {
1413 let value = value.trim();
1414 if value.is_empty() || value.chars().count() > 512 {
1415 return Err(format!(
1416 "task(): {field} entries must be 1..=512 characters"
1417 ));
1418 }
1419 if !normalized.iter().any(|existing| existing == value) {
1420 normalized.push(value.to_string());
1421 }
1422 }
1423 Ok(normalized)
1424 }
1425
1426 /// Normalize a task working directory at both plan preflight and VM dispatch.
1427 /// The same bounded repo-relative policy applies to both entry points.
1428 pub fn normalize_task_cwd(value: &str) -> Result<String, String> {
1429 normalize_task_paths("cwd", vec![value.to_owned()], 1).map(|mut paths| paths.remove(0))
1430 }
1431
1432 /// [`normalize_task_cwd`] for a run that knows its workspace root. An
1433 /// absolute path that lies inside `workspace` is rewritten to the same
1434 /// bounded repo-relative form; an absolute path outside it, or one that
1435 /// escapes through `..`, is still rejected. Without a workspace every
1436 /// absolute path is rejected, exactly as before.
1437 ///
1438 /// The comparison is lexical and component-wise (no filesystem access), so
1439 /// `/ws-other` is never inside `/ws`. On Windows a drive letter or UNC share
1440 /// matches with or without the verbatim `\\?\` prefix, `/` and `\` are both
1441 /// separators, and components compare case-insensitively.
1442 pub fn normalize_task_cwd_in(
1443 value: &str,
1444 workspace: Option<&std::path::Path>,
1445 ) -> Result<String, String> {
1446 let raw = value.trim();
1447 let path = std::path::Path::new(raw);
1448 let Some(workspace) = workspace.filter(|_| path.is_absolute()) else {
1449 return normalize_task_cwd(value);
1450 };
1451 let outside = || {
1452 format!(
1453 "task(): cwd {raw:?} is outside the workspace {}; cwd entries must be bounded repo-relative paths or absolute paths inside the workspace",
1454 workspace.display()
1455 )
1456 };
1457 let mut remainder = path_keys(path).ok_or_else(outside)?.into_iter();
1458 for expected in path_keys(workspace).ok_or_else(outside)? {
1459 match remainder.next() {
1460 Some(actual) if actual.matches(&expected) => {}
1461 // `..` right after the workspace prefix is an escape, not a
1462 // different root: name it as traversal.
1463 Some(PathKey::Parent) => {
1464 return Err("task(): cwd paths cannot contain parent traversal".to_string());
1465 }
1466 _ => return Err(outside()),
1467 }
1468 }
1469 let mut segments = Vec::new();
1470 for key in remainder {
1471 match key {
1472 PathKey::Name { original, .. } => segments.push(original),
1473 PathKey::Parent => {
1474 return Err("task(): cwd paths cannot contain parent traversal".to_string());
1475 }
1476 PathKey::Root => return Err(outside()),
1477 }
1478 }
1479 if segments.is_empty() {
1480 return Ok(".".to_string());
1481 }
1482 normalize_task_cwd(&segments.join("/"))
1483 }
1484
1485 /// One lexical path component, keyed for comparison: a drive or UNC prefix
1486 /// with its verbatim marker dropped, and names case-folded on Windows.
1487 #[derive(Debug)]
1488 enum PathKey {
1489 Root,
1490 Parent,
1491 Name { key: String, original: String },
1492 }
1493
1494 impl PathKey {
1495 fn matches(&self, other: &Self) -> bool {
1496 match (self, other) {
1497 (Self::Root, Self::Root) | (Self::Parent, Self::Parent) => true,
1498 (Self::Name { key: left, .. }, Self::Name { key: right, .. }) => left == right,
1499 _ => false,
1500 }
1501 }
1502 }
1503
1504 fn path_keys(path: &std::path::Path) -> Option<Vec<PathKey>> {
1505 use std::path::{Component, Prefix};
1506 let fold = |value: &str| {
1507 if cfg!(windows) {
1508 value.to_lowercase()
1509 } else {
1510 value.to_string()
1511 }
1512 };
1513 let mut keys = Vec::new();
1514 for component in path.components() {
1515 match component {
1516 Component::Prefix(prefix) => {
1517 let key = match prefix.kind() {
1518 Prefix::Disk(drive) | Prefix::VerbatimDisk(drive) => {
1519 format!("{}:", drive.to_ascii_lowercase() as char)
1520 }
1521 Prefix::UNC(server, share) | Prefix::VerbatimUNC(server, share) => format!(
1522 "//{}/{}",
1523 server.to_str()?.to_lowercase(),
1524 share.to_str()?.to_lowercase()
1525 ),
1526 _ => prefix.as_os_str().to_str()?.to_lowercase(),
1527 };
1528 keys.push(PathKey::Name {
1529 original: key.clone(),
1530 key,
1531 });
1532 }
1533 Component::RootDir => keys.push(PathKey::Root),
1534 Component::CurDir => {}
1535 Component::ParentDir => keys.push(PathKey::Parent),
1536 Component::Normal(name) => {
1537 let original = name.to_str()?.to_string();
1538 keys.push(PathKey::Name {
1539 key: fold(&original),
1540 original,
1541 });
1542 }
1543 }
1544 }
1545 Some(keys)
1546 }
1547
1548 fn normalize_task_paths(
1549 field: &str,
1550 values: Vec<String>,
1551 limit: usize,
1552 ) -> Result<Vec<String>, String> {
1553 if values.len() > limit {
1554 return Err(format!("task(): {field} accepts at most {limit} entries"));
1555 }
1556 let mut normalized = Vec::new();
1557 for raw in values {
1558 let raw = raw.trim().replace('\\', "/");
1559 let windows_drive = raw.as_bytes().get(1) == Some(&b':')
1560 && raw.as_bytes().first().is_some_and(u8::is_ascii_alphabetic);
1561 if raw.is_empty()
1562 || raw.chars().count() > 512
1563 || raw.starts_with('/')
1564 || raw.starts_with("//")
1565 || windows_drive
1566 || raw.chars().any(|ch| matches!(ch, '\0' | '\r' | '\n'))
1567 {
1568 return Err(format!(
1569 "task(): {field} entries must be bounded repo-relative paths"
1570 ));
1571 }
1572 let mut segments = Vec::new();
1573 for segment in raw.split('/') {
1574 match segment {
1575 "" | "." => {}
1576 ".." => {
1577 return Err(format!(
1578 "task(): {field} paths cannot contain parent traversal"
1579 ));
1580 }
1581 value => segments.push(value),
1582 }
1583 }
1584 let path = if segments.is_empty() {
1585 ".".to_string()
1586 } else {
1587 segments.join("/")
1588 };
1589 if !normalized.contains(&path) {
1590 normalized.push(path);
1591 }
1592 }
1593 Ok(normalized)
1594 }
1595
1596 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
1597 enum TaskRoleKind {
1598 ReadOnly,
1599 General,
1600 Implementer,
1601 }
1602
1603 fn task_role_kind(value: &str) -> Option<TaskRoleKind> {
1604 match value.trim().to_ascii_lowercase().as_str() {
1605 "explore" | "explorer" | "scout" | "plan" | "planner" | "review" | "reviewer"
1606 | "verify" | "verifier" => Some(TaskRoleKind::ReadOnly),
1607 "general" | "worker" => Some(TaskRoleKind::General),
1608 "implement" | "implementer" | "builder" => Some(TaskRoleKind::Implementer),
1609 _ => None,
1610 }
1611 }
1612
1613 /// The JS prelude injected before every script: determinism bans, the
1614 /// `task`/`parallel`/`pipeline`/`log`/`phase` stdlib (design §7), and the
1615 /// `budget` global.
1616 fn prelude() -> String {
1617 PRELUDE_TEMPLATE.replace("__MAX_ITEMS__", &PARALLEL_MAX_ITEMS.to_string())
1618 }
1619
1620 const PRELUDE_TEMPLATE: &str = r#""use strict";
1621 (() => {
1622 const banned = (name) => () => {
1623 throw new Error(name + " is unavailable in Workflow scripts: runs must be deterministic for record/replay");
1624 };
1625 const BannedDate = function Date() {
1626 throw new Error("new Date()/Date() is unavailable in Workflow scripts: runs must be deterministic for record/replay");
1627 };
1628 BannedDate.now = banned("Date.now()");
1629 BannedDate.parse = banned("Date.parse()");
1630 BannedDate.UTC = banned("Date.UTC()");
1631 globalThis.Date = BannedDate;
1632 Math.random = banned("Math.random()");
1633
1634 // Capture temporary host bindings into this closure, then strip them from
1635 // globalThis so scripts only see the documented Workflow surface (#4129).
1636 const hostTask = __workflow_task;
1637 const hostLog = __workflow_log;
1638 const hostEverySlotFailed = __workflow_every_slot_failed;
1639 const hostSlotDropped = __workflow_slot_dropped;
1640 const hostPhase = __workflow_phase;
1641 const hostBudgetTotal = __workflow_budget_total;
1642 const hostBudgetSpent = __workflow_budget_spent;
1643 const hostBudgetRemaining = __workflow_budget_remaining;
1644
1645 const MAX_ITEMS = __MAX_ITEMS__;
1646 const taskErrorText = (err) => String(err && err.message !== undefined ? err.message : err);
1647
1648 // Typed slot errors (R9). Every error thrown by task() carries a
1649 // host-assigned `kind` copied off the task envelope; anything else that
1650 // reaches a slot was thrown by the script itself and reports as "script".
1651 //
1652 // The kind is read from the error object and never guessed from message
1653 // text. A substring classifier let a child's own words ("...budget
1654 // exhausted...", "...responseSchema...") forge a fatal classification and
1655 // abort a healthy run, and it could not tell a genuine subagent failure
1656 // apart from a plain `throw new Error(...)` in a stage.
1657 const HOST_KINDS = ["admission", "budget", "cancelled", "agent", "schema", "driver"];
1658 const SCRIPT_KIND = "script";
1659 const taskErrorKind = (err) =>
1660 err !== null && typeof err === "object" && HOST_KINDS.indexOf(err.kind) !== -1
1661 ? err.kind
1662 : SCRIPT_KIND;
1663 // Fatal kinds are never absorbed into a slot value: cancellation is the
1664 // run's own deadline, and a schema breach means the contract the caller
1665 // explicitly asked for was not met. `mode: "partial"` opts out for schema
1666 // (and only schema) by keeping it as a structured slot value instead.
1667 const isFatalTaskError = (err) => {
1668 const kind = taskErrorKind(err);
1669 return kind === "cancelled" || kind === "schema";
1670 };
1671
1672 // Stamp the resolved kind onto an error that is about to be rethrown, so a
1673 // script's own `catch (err) { err.kind }` reads the same vocabulary the
1674 // slot classifier used. Host errors already carry theirs; this only names
1675 // the script throws, which would otherwise surface as `undefined`.
1676 const stampKind = (err, kind) => {
1677 if (err !== null && typeof err === "object" && err.kind === undefined) {
1678 try {
1679 err.kind = kind;
1680 } catch (_) {
1681 // A frozen error keeps whatever it has; the log line still names it.
1682 }
1683 }
1684 return err;
1685 };
1686
1687 const SLOT_MODES = ["settled", "fail-fast", "partial"];
1688 // `settled` is the default and is exactly today's behavior: a non-fatal
1689 // slot failure resolves to `null` so an author need not handle every error.
1690 // An unrecognized mode throws rather than silently falling back — a typo
1691 // like `mode: "failfast"` used to read as `settled` and quietly keep
1692 // dropping slots the author believed were now fatal.
1693 const slotMode = (fn, opts) => {
1694 if (opts === null || typeof opts !== "object" || opts.mode === undefined) {
1695 return "settled";
1696 }
1697 if (SLOT_MODES.indexOf(opts.mode) === -1) {
1698 throw new Error(
1699 fn + "(): unknown mode " + JSON.stringify(opts.mode) +
1700 "; expected one of " + SLOT_MODES.join(", ")
1701 );
1702 }
1703 return opts.mode;
1704 };
1705
1706 // The failure ledger for one fan-out, attached to the resolved array as a
1707 // non-enumerable `errors` property. Non-enumerable and non-index, so the
1708 // array's contents, length, and JSON encoding are byte-identical to before:
1709 // `results.filter(Boolean)` still works, and a script that wants to know
1710 // WHY a slot is null can now ask instead of guessing.
1711 const attachSlotErrors = (results, errors) => {
1712 errors.sort((a, b) => a.index - b.index);
1713 Object.defineProperty(results, "errors", {
1714 value: Object.freeze(errors.map((entry) => Object.freeze(entry))),
1715 enumerable: false,
1716 configurable: false,
1717 writable: false,
1718 });
1719 return results;
1720 };
1721
1722 globalThis.task = async (opts) => {
1723 if (opts === null || typeof opts !== "object") {
1724 throw new TypeError("task(): expected an options object");
1725 }
1726 const envelope = JSON.parse(await hostTask(JSON.stringify(opts)));
1727 if (envelope.error !== undefined) {
1728 const err = new Error(envelope.error);
1729 // The host always names the kind; a missing one means an envelope this
1730 // prelude did not produce, which is not a typed task failure.
1731 err.kind = HOST_KINDS.indexOf(envelope.error_kind) !== -1
1732 ? envelope.error_kind
1733 : SCRIPT_KIND;
1734 throw err;
1735 }
1736 return envelope.value;
1737 };
1738
1739 globalThis.parallel = (thunks, opts) => {
1740 if (!Array.isArray(thunks)) {
1741 throw new TypeError("parallel(): expected an array of thunks");
1742 }
1743 if (thunks.length > MAX_ITEMS) {
1744 throw new Error("parallel(): max " + MAX_ITEMS + " items per call");
1745 }
1746 const mode = slotMode("parallel", opts);
1747 const failFast = mode === "fail-fast";
1748 const partial = mode === "partial";
1749 const errors = [];
1750 // Returns the slot value, or throws to reject the whole fan-out.
1751 const onSlotError = (index, err) => {
1752 const kind = taskErrorKind(err);
1753 const message = taskErrorText(err);
1754 stampKind(err, kind);
1755 // Cancellation is the run's deadline in every mode, partial included.
1756 if (kind === "cancelled") throw err;
1757 if (partial) {
1758 // Opt-in partial mode: every non-cancellation slot failure becomes a
1759 // structured value the script can branch on. It never masquerades as
1760 // a success -- `__taskError` is the whole point of the shape.
1761 errors.push({ index: index, kind: kind, message: message });
1762 hostLog(
1763 "parallel(): partial mode kept a failed slot as __taskError (kind=" +
1764 kind + ", slot " + index + "): " + message
1765 );
1766 return { __taskError: { index: index, kind: kind, message: message } };
1767 }
1768 if (isFatalTaskError(err)) throw err;
1769 if (failFast) {
1770 hostLog(
1771 "parallel(): fail-fast slot error (kind=" + kind + ", slot " + index + "): " + message
1772 );
1773 throw err;
1774 }
1775 errors.push({ index: index, kind: kind, message: message });
1776 hostLog(
1777 "parallel(): dropped a failed slot as null (kind=" + kind + ", slot " + index + "): " +
1778 message
1779 );
1780 hostSlotDropped("parallel", kind, index);
1781 return null;
1782 };
1783 const slots = thunks.map((thunk, index) => {
1784 try {
1785 return Promise.resolve(typeof thunk === "function" ? thunk() : thunk)
1786 .catch((err) => onSlotError(index, err));
1787 } catch (err) {
1788 try {
1789 return onSlotError(index, err);
1790 } catch (rethrown) {
1791 return Promise.reject(rethrown);
1792 }
1793 }
1794 });
1795 return Promise.all(slots).then((results) => {
1796 // A fan-out where nothing survived is a dead fan-out, not resilience.
1797 // The default stays ergonomic (the array still resolves) but the run
1798 // log says so in one line an operator can grep for, and the structured
1799 // event below lets the host status classifier refuse to call such a
1800 // run a plain success.
1801 if (results.length > 0 && errors.length === results.length) {
1802 hostLog(
1803 "parallel(): every slot failed (" + errors.length + " of " + results.length +
1804 "); no work survived this fan-out"
1805 );
1806 hostEverySlotFailed("parallel", errors.length, results.length);
1807 }
1808 return attachSlotErrors(results, errors);
1809 });
1810 };
1811
1812 globalThis.pipeline = (items, ...stages) => {
1813 if (!Array.isArray(items)) {
1814 throw new TypeError("pipeline(): expected an array of items");
1815 }
1816 if (items.length > MAX_ITEMS) {
1817 throw new Error("pipeline(): max " + MAX_ITEMS + " items per call");
1818 }
1819 // Options overload: pipeline(items, { stages: [...], mode: "fail-fast" }).
1820 let mode = "settled";
1821 if (
1822 stages.length === 1 &&
1823 stages[0] !== null &&
1824 typeof stages[0] === "object" &&
1825 Array.isArray(stages[0].stages)
1826 ) {
1827 mode = slotMode("pipeline", stages[0]);
1828 stages = stages[0].stages;
1829 }
1830 const failFast = mode === "fail-fast";
1831 const partial = mode === "partial";
1832 const errors = [];
1833 return Promise.all(items.map(async (item, index) => {
1834 let value = item;
1835 for (const stage of stages) {
1836 try {
1837 value = await stage(value, item, index);
1838 } catch (err) {
1839 const kind = taskErrorKind(err);
1840 const message = taskErrorText(err);
1841 stampKind(err, kind);
1842 if (kind === "cancelled") throw err;
1843 if (partial) {
1844 errors.push({ index: index, kind: kind, message: message });
1845 hostLog(
1846 "pipeline(): partial mode kept item " + index +
1847 " as __taskError (kind=" + kind + "): " + message
1848 );
1849 return { __taskError: { index: index, kind: kind, message: message } };
1850 }
1851 if (isFatalTaskError(err)) throw err;
1852 if (failFast) {
1853 hostLog(
1854 "pipeline(): fail-fast stage error on item " + index +
1855 " (kind=" + kind + "): " + message
1856 );
1857 throw err;
1858 }
1859 errors.push({ index: index, kind: kind, message: message });
1860 hostLog(
1861 "pipeline(): dropped item " + index + " as null (kind=" + kind + "): " + message
1862 );
1863 hostSlotDropped("pipeline", kind, index);
1864 return null;
1865 }
1866 }
1867 return value;
1868 })).then((results) => {
1869 if (results.length > 0 && errors.length === results.length) {
1870 hostLog(
1871 "pipeline(): every item failed (" + errors.length + " of " + results.length +
1872 "); no work survived this pipeline"
1873 );
1874 hostEverySlotFailed("pipeline", errors.length, results.length);
1875 }
1876 return attachSlotErrors(results, errors);
1877 });
1878 };
1879
1880 globalThis.log = (message) => {
1881 hostLog(typeof message === "string" ? message : (JSON.stringify(message) ?? String(message)));
1882 };
1883 globalThis.phase = (title) => {
1884 hostPhase(String(title));
1885 };
1886
1887 const total = hostBudgetTotal();
1888 globalThis.budget = Object.freeze({
1889 total: Number.isNaN(total) ? null : total,
1890 spent: () => hostBudgetSpent(),
1891 remaining: () => hostBudgetRemaining(),
1892 });
1893
1894 for (const name of [
1895 "__workflow_task",
1896 "__workflow_log",
1897 "__workflow_every_slot_failed",
1898 "__workflow_phase",
1899 "__workflow_budget_total",
1900 "__workflow_budget_spent",
1901 "__workflow_budget_remaining",
1902 ]) {
1903 try {
1904 delete globalThis[name];
1905 } catch (_) {
1906 // Non-configurable bindings stay; the inventory test will fail closed.
1907 }
1908 }
1909 })();
1910 "#;
1911
1911 lines RUST