返回 CodeWhale
vm_tests.rs
根目录 / crates / workflow-js / tests / vm_tests.rs
1 //! End-to-end tests for the Workflow JS runtime against a fake driver.
2
3 use std::sync::Arc;
4 use std::time::Duration;
5
6 use codewhale_workflow_js::testing::{FakeDriver, FakeReply};
7 use codewhale_workflow_js::{
8 ProgressEvent, WORKFLOW_LIFETIME_CAP, WorkflowJsError, WorkflowRunCancel, WorkflowVm,
9 };
10 use serde_json::json;
11
12 async fn run(
13 driver: &Arc<FakeDriver>,
14 source: &str,
15 args: serde_json::Value,
16 ) -> Result<serde_json::Value, WorkflowJsError> {
17 WorkflowVm::new()
18 .run_script(
19 source,
20 args,
21 driver.clone() as Arc<dyn codewhale_workflow_js::WorkflowDriver>,
22 )
23 .await
24 }
25
26 fn script_message(result: Result<serde_json::Value, WorkflowJsError>) -> String {
27 match result {
28 Err(WorkflowJsError::Script(message)) => message,
29 other => panic!("expected script error, got {other:?}"),
30 }
31 }
32
33 #[tokio::test]
34 async fn plain_return_value_round_trips() {
35 let driver = Arc::new(FakeDriver::new());
36 let value = run(&driver, "return 1 + 1;", json!(null)).await.unwrap();
37 assert_eq!(value, json!(2));
38 }
39
40 #[tokio::test]
41 async fn undefined_return_becomes_null() {
42 let driver = Arc::new(FakeDriver::new());
43 let value = run(&driver, "const x = 1;", json!(null)).await.unwrap();
44 assert_eq!(value, json!(null));
45 }
46
47 #[tokio::test]
48 async fn args_global_is_the_invocation_input() {
49 let driver = Arc::new(FakeDriver::new());
50 let value = run(
51 &driver,
52 "return { sum: args.x + 1, tag: args.tags[0] };",
53 json!({"x": 41, "tags": ["release"]}),
54 )
55 .await
56 .unwrap();
57 assert_eq!(value, json!({"sum": 42, "tag": "release"}));
58 }
59
60 #[tokio::test]
61 async fn checked_in_best_of_n_search_recipe_runs_with_structured_receipts() {
62 let driver = Arc::new(FakeDriver::new());
63 for index in 1..=2 {
64 driver.on(
65 // Rules match the driver-visible `TaskRequest.description`, which
66 // is the full instruction text (the VM's `prompt` alias wins over
67 // a short label). Match the unique per-candidate suffix line.
68 &format!("candidate_id=cand_{index:03} of 2."),
69 FakeReply::Complete(
70 json!({
71 "candidate_id": format!("cand_{index:03}"),
72 "hypothesis": "bounded fixture",
73 "modified_paths": ["src/lib.rs"],
74 "commands_run": ["cargo test --locked"],
75 "self_verdict": "pass",
76 "known_risks": [],
77 "artifact_refs": [format!("patch:cand_{index:03}")]
78 })
79 .to_string(),
80 ),
81 );
82 }
83 driver.on(
84 "read-only tournament judge",
85 FakeReply::Complete(
86 json!({
87 "winner_id": "cand_001",
88 "ranking": ["cand_001", "cand_002"],
89 "verification_required": true,
90 "reasons": ["fixture score"]
91 })
92 .to_string(),
93 ),
94 );
95
96 let value = run(
97 &driver,
98 include_str!("../../../workflows/operate_best_of_n.workflow.js"),
99 json!({
100 "brief": "Implement the fixture",
101 "strategy": "search",
102 "n": 2,
103 "writeRoots": ["src"],
104 "model": "deepseek-v4-flash",
105 "thinking": "max"
106 }),
107 )
108 .await
109 .expect("checked-in search recipe should execute");
110
111 assert_eq!(value["scenario"], "operate-search");
112 assert_eq!(value["review"]["winner_id"], "cand_001");
113 assert_eq!(driver.spawn_count(), 3);
114 let requests = driver.requests();
115 assert_eq!(requests[0].model.as_deref(), Some("deepseek-v4-flash"));
116 assert_eq!(requests[0].thinking.as_deref(), Some("max"));
117 assert_eq!(requests[0].write_roots, ["src"]);
118 assert_eq!(requests[2].write_authority.as_deref(), Some("read_only"));
119 // Regression: the driver-visible description is the full instruction text,
120 // so reply rules must target text that actually reaches the driver. If a
121 // future recipe reintroduces a separate short `description` next to a long
122 // `prompt`, these needles stop matching, the FakeDriver falls back to its
123 // non-JSON "done:..." reply, and the structured receipts fail loudly.
124 assert!(
125 requests[0]
126 .description
127 .starts_with("You are one independent candidate")
128 );
129 assert!(
130 requests[0]
131 .description
132 .contains("CANDIDATE-SPECIFIC INSTRUCTION: candidate_id=cand_001 of 2.")
133 );
134 assert!(
135 requests[2]
136 .description
137 .starts_with("You are the read-only tournament judge")
138 );
139 }
140
141 #[tokio::test]
142 async fn task_prompt_wins_over_description_as_driver_visible_text() {
143 let driver = Arc::new(FakeDriver::new());
144 let value = run(
145 &driver,
146 r#"
147 return await task({
148 description: "short progress label",
149 prompt: "the real instruction",
150 });
151 "#,
152 json!(null),
153 )
154 .await
155 .unwrap();
156
157 // No rules were registered, so the FakeDriver fallback echoes the
158 // driver-visible description. The reply text proves the driver received
159 // the prompt, not the short label.
160 assert_eq!(value, json!("done:the real instruction"));
161 let requests = driver.requests();
162 assert_eq!(requests.len(), 1);
163 assert_eq!(requests[0].description, "the real instruction");
164 assert_ne!(requests[0].description, "short progress label");
165 }
166
167 #[tokio::test]
168 async fn task_round_trip_carries_all_options_and_normalizes_profile() {
169 let driver = Arc::new(FakeDriver::new());
170 let value = run(
171 &driver,
172 r#"
173 return await task({
174 description: "implement the bounded change",
175 subagentType: "implementer",
176 profile: " ALpha-1 ",
177 model: "deepseek-chat",
178 modelStrength: "faster",
179 thinking: "low",
180 cwd: "repo-a",
181 worktree: true,
182 writeAuthority: "worktree_write",
183 writeRoots: ["crates/tui/src"],
184 exactFiles: ["Cargo.toml"],
185 coordinationContracts: ["public-api"],
186 dependencies: ["issue-4619"],
187 acceptance: ["locked tests pass"],
188 allowedTools: ["read", "grep"],
189 maxDepth: 2,
190 tokenBudget: 5000,
191 maxSteps: 4,
192 wallTimeSecs: 90,
193 label: "L1",
194 phase: "P1",
195 });
196 "#,
197 json!(null),
198 )
199 .await
200 .unwrap();
201 assert_eq!(value, json!("done:implement the bounded change"));
202
203 let requests = driver.requests();
204 assert_eq!(requests.len(), 1);
205 let request = &requests[0];
206 assert_eq!(request.description, "implement the bounded change");
207 assert_eq!(request.subagent_type.as_deref(), Some("implementer"));
208 assert_eq!(request.profile.as_deref(), Some("alpha-1"));
209 assert_eq!(request.model.as_deref(), Some("deepseek-chat"));
210 assert_eq!(request.model_strength.as_deref(), Some("faster"));
211 assert_eq!(request.thinking.as_deref(), Some("low"));
212 assert_eq!(request.cwd.as_deref(), Some("repo-a"));
213 assert!(request.worktree);
214 assert_eq!(request.write_authority.as_deref(), Some("worktree_write"));
215 assert_eq!(request.write_roots, ["crates/tui/src"]);
216 assert_eq!(request.exact_files, ["Cargo.toml"]);
217 assert_eq!(request.coordination_contracts, ["public-api"]);
218 assert_eq!(request.dependencies, ["issue-4619"]);
219 assert_eq!(request.acceptance, ["locked tests pass"]);
220 assert_eq!(
221 request.allowed_tools.as_deref(),
222 Some(["read".to_string(), "grep".to_string()].as_slice())
223 );
224 assert_eq!(request.max_depth, Some(2));
225 assert_eq!(request.token_budget, Some(5000));
226 assert_eq!(request.max_steps, Some(4));
227 assert_eq!(request.wall_time_secs, Some(90));
228 assert_eq!(request.response_schema, None);
229 assert_eq!(request.label.as_deref(), Some("L1"));
230 assert_eq!(request.phase.as_deref(), Some("P1"));
231 }
232
233 #[tokio::test]
234 async fn scopeless_write_tasks_dispatch_and_leave_scope_to_the_spawn_boundary() {
235 // A Workflow script starts a writing child the same way the Agent tool
236 // does: no declared scope, so the spawn boundary claims the workspace
237 // root. The VM must not refuse it before dispatch.
238 let driver = Arc::new(FakeDriver::new());
239 run(
240 &driver,
241 r#"
242 await task({ prompt: "implement it", type: "implementer" });
243 await task({ prompt: "build it", type: "builder" });
244 await task({ prompt: "do it", type: "general" });
245 await task({ prompt: "lead it", profile: "release-lead" });
246 return await task({
247 prompt: "edit without a claim",
248 writeAuthority: "workspace_write",
249 });
250 "#,
251 json!(null),
252 )
253 .await
254 .expect("scopeless write tasks dispatch");
255 let requests = driver.requests();
256 assert_eq!(requests.len(), 5);
257 for request in &requests {
258 assert!(request.write_roots.is_empty());
259 assert!(request.exact_files.is_empty());
260 assert!(request.coordination_contracts.is_empty());
261 }
262 assert_eq!(requests[0].subagent_type.as_deref(), Some("implementer"));
263 assert_eq!(
264 requests[4].write_authority.as_deref(),
265 Some("workspace_write")
266 );
267
268 // A read-only role still cannot claim write authority.
269 let driver = Arc::new(FakeDriver::new());
270 let error = run(
271 &driver,
272 r#"return await task({ prompt: "x", type: "reviewer", writeAuthority: "workspace_write" });"#,
273 json!(null),
274 )
275 .await
276 .expect_err("read-only role with write authority")
277 .to_string();
278 assert!(
279 error.contains("read-only roles cannot declare write-capable authority"),
280 "{error}"
281 );
282 assert!(driver.requests().is_empty());
283 }
284
285 #[tokio::test]
286 async fn task_coordination_lists_deduplicate_with_hard_count_bounds() {
287 let driver = Arc::new(FakeDriver::new());
288 run(
289 &driver,
290 r#"
291 return await task({
292 prompt: "bounded edit",
293 type: "implementer",
294 writeAuthority: "workspace_write",
295 exactFiles: ["src/a.rs", "src/a.rs"],
296 dependencies: ["A", "A"],
297 acceptance: ["tests pass", "tests pass"],
298 });
299 "#,
300 json!(null),
301 )
302 .await
303 .expect("bounded unique coordination values");
304 let request = driver.requests().pop().expect("request");
305 assert_eq!(request.exact_files, ["src/a.rs"]);
306 assert_eq!(request.dependencies, ["A"]);
307 assert_eq!(request.acceptance, ["tests pass"]);
308 }
309
310 #[tokio::test]
311 async fn task_write_paths_normalize_and_reject_escape_spellings() {
312 let driver = Arc::new(FakeDriver::new());
313 run(
314 &driver,
315 r#"return await task({
316 prompt: "bounded edit",
317 type: "implementer",
318 writeRoots: ["./src//", "src"],
319 exactFiles: ["src\\lib.rs"]
320 });"#,
321 json!(null),
322 )
323 .await
324 .expect("normalized repo-relative paths");
325 let request = driver.requests().pop().expect("request");
326 assert_eq!(request.write_roots, ["src"]);
327 assert_eq!(request.exact_files, ["src/lib.rs"]);
328
329 for path in [
330 "../outside",
331 "/tmp/outside",
332 "C:\\outside",
333 "src/../../outside",
334 ] {
335 let driver = Arc::new(FakeDriver::new());
336 let source = format!(
337 "return await task({{ prompt: 'escape', type: 'implementer', writeRoots: [{}] }});",
338 serde_json::to_string(path).expect("path json")
339 );
340 let message = script_message(run(&driver, &source, json!(null)).await);
341 assert!(
342 message.contains("repo-relative") || message.contains("traversal"),
343 "{path}: {message}"
344 );
345 assert!(driver.requests().is_empty());
346 }
347 }
348
349 #[tokio::test]
350 async fn task_read_only_roles_reject_write_escalation() {
351 // Scopeless write-capable tasks dispatch (see
352 // scopeless_write_tasks_dispatch_and_leave_scope_to_the_spawn_boundary);
353 // a read-only role still cannot claim write authority, and contradictory
354 // identities are still refused.
355 for source in [
356 r#"return await task({prompt: "wrong authority", type: "reviewer", writeAuthority: "workspace_write", writeRoots: ["src"]});"#,
357 r#"return await task({prompt: "wrong authority", type: "scout", writeAuthority: "workspace_write", writeRoots: ["src"]});"#,
358 r#"return await task({prompt: "role conflict", type: "implementer", role: "reviewer", writeRoots: ["src"]});"#,
359 ] {
360 let driver = Arc::new(FakeDriver::new());
361 let message = script_message(run(&driver, source, json!(null)).await);
362 assert!(
363 message.contains("require")
364 || message.contains("cannot")
365 || message.contains("contradictory"),
366 "{message}"
367 );
368 assert!(driver.requests().is_empty());
369 }
370 }
371
372 #[tokio::test]
373 async fn task_implementer_identity_can_be_narrowed_to_read_only_authority() {
374 let driver = Arc::new(FakeDriver::new());
375 let value = run(
376 &driver,
377 r#"return await task({prompt: "verification-only plan", type: "implementer", writeAuthority: "read_only"});"#,
378 json!(null),
379 )
380 .await
381 .expect("read-only authority must safely narrow an implementer identity");
382 assert_eq!(value, json!("done:verification-only plan"));
383 let request = driver.requests().pop().expect("request");
384 assert_eq!(request.subagent_type.as_deref(), Some("implementer"));
385 assert_eq!(request.write_authority.as_deref(), Some("read_only"));
386 assert!(request.write_roots.is_empty());
387 }
388
389 #[tokio::test]
390 async fn task_accepts_prompt_and_type_aliases() {
391 let driver = Arc::new(FakeDriver::new());
392 run(
393 &driver,
394 r#"return await task({ prompt: "aliased", type: "verifier" });"#,
395 json!(null),
396 )
397 .await
398 .unwrap();
399 let request = &driver.requests()[0];
400 assert_eq!(request.description, "aliased");
401 assert_eq!(request.subagent_type.as_deref(), Some("verifier"));
402 }
403
404 #[tokio::test]
405 async fn task_title_alias_routes_to_description() {
406 let driver = Arc::new(FakeDriver::new());
407 run(
408 &driver,
409 r#"return await task({ title: "inspect the release candidate", type: "verifier" });"#,
410 json!(null),
411 )
412 .await
413 .expect("title is accepted as the task description");
414
415 let request = &driver.requests()[0];
416 assert_eq!(request.description, "inspect the release candidate");
417 assert_eq!(request.subagent_type.as_deref(), Some("verifier"));
418 }
419
420 #[tokio::test]
421 async fn task_prompt_takes_precedence_over_short_description() {
422 let driver = Arc::new(FakeDriver::new());
423 run(
424 &driver,
425 r#"return await task({
426 description: "Short progress summary",
427 prompt: "Detailed child instructions",
428 label: "fixture-compatible"
429 });"#,
430 json!(null),
431 )
432 .await
433 .unwrap();
434 let request = &driver.requests()[0];
435 assert_eq!(request.description, "Detailed child instructions");
436 assert_eq!(request.label.as_deref(), Some("fixture-compatible"));
437 }
438
439 #[tokio::test]
440 async fn task_rejects_invalid_profile_tokens() {
441 for bad in ["two words", "a=b", "a\"b", "a`b", " "] {
442 let driver = Arc::new(FakeDriver::new());
443 let source = format!(
444 "return await task({{ description: \"x\", profile: {} }});",
445 serde_json::Value::String(bad.to_string())
446 );
447 let message = script_message(run(&driver, &source, json!(null)).await);
448 assert!(message.contains("profile"), "profile {bad:?}: {message}");
449 assert_eq!(driver.spawn_count(), 0, "invalid profile must not spawn");
450 }
451 }
452
453 #[tokio::test]
454 async fn task_requires_a_description() {
455 let driver = Arc::new(FakeDriver::new());
456 let message = script_message(run(&driver, "return await task({});", json!(null)).await);
457 assert!(message.contains("description"), "{message}");
458 assert_eq!(driver.spawn_count(), 0);
459 }
460
461 #[tokio::test]
462 async fn task_rejects_unknown_option_names() {
463 let driver = Arc::new(FakeDriver::new());
464 let message = script_message(
465 run(
466 &driver,
467 r#"return await task({ description: "x", responseschema: {} });"#,
468 json!(null),
469 )
470 .await,
471 );
472 assert!(message.contains("invalid options"), "{message}");
473 assert_eq!(driver.spawn_count(), 0);
474 }
475
476 #[tokio::test]
477 async fn driver_rejection_is_catchable_in_script() {
478 let driver = Arc::new(FakeDriver::new());
479 driver.on("bad", FakeReply::Reject("admission cap".to_string()));
480 let value = run(
481 &driver,
482 r#"
483 try {
484 await task({ description: "bad idea" });
485 return "no-throw";
486 } catch (err) {
487 return String(err);
488 }
489 "#,
490 json!(null),
491 )
492 .await
493 .unwrap();
494 let text = value.as_str().unwrap();
495 assert!(text.contains("admission cap"), "{text}");
496 }
497
498 #[tokio::test]
499 async fn parallel_fan_out_maps_one_failure_to_null_slot() {
500 let driver = Arc::new(FakeDriver::new());
501 driver.on("beta", FakeReply::Fail("boom".to_string()));
502 let value = run(
503 &driver,
504 r#"
505 return await parallel([
506 () => task({ description: "alpha" }),
507 () => task({ description: "beta" }),
508 () => task({ description: "gamma" }),
509 ]);
510 "#,
511 json!(null),
512 )
513 .await
514 .unwrap();
515 assert_eq!(value, json!(["done:alpha", null, "done:gamma"]));
516 assert_eq!(driver.spawn_count(), 3);
517 }
518
519 #[tokio::test]
520 async fn parallel_logs_a_breadcrumb_when_a_slot_is_dropped_to_null() {
521 // #dogfood 0.8.67: a fan-out slot that fails for a non-schema reason still
522 // resolves to null (documented resilience), but must leave a breadcrumb in
523 // the run log so an operator can see why a slot came back null / nothing
524 // spawned — instead of a silent "completed" with no explanation.
525 let driver = Arc::new(FakeDriver::new());
526 driver.on("beta", FakeReply::Fail("boom".to_string()));
527 let value = run(
528 &driver,
529 r#"
530 return await parallel([
531 () => task({ description: "alpha" }),
532 () => task({ description: "beta" }),
533 ]);
534 "#,
535 json!(null),
536 )
537 .await
538 .unwrap();
539 assert_eq!(value, json!(["done:alpha", null]));
540 assert!(
541 driver.events().iter().any(|event| matches!(
542 event,
543 ProgressEvent::Log { message } if message.contains("dropped a failed slot")
544 )),
545 "a dropped parallel slot should leave a breadcrumb in the run log"
546 );
547 }
548
549 #[tokio::test]
550 async fn parallel_surfaces_response_schema_errors_instead_of_null() {
551 let driver = Arc::new(FakeDriver::new());
552 driver.on(
553 "bad schema",
554 FakeReply::Complete(r#"{"refuted":"yes"}"#.to_string()),
555 );
556
557 let message = script_message(
558 run(
559 &driver,
560 r#"
561 return await parallel([
562 () => task({
563 description: "bad schema",
564 responseSchema: {
565 type: "object",
566 properties: { refuted: { type: "boolean" } },
567 required: ["refuted"],
568 },
569 }),
570 ]);
571 "#,
572 json!(null),
573 )
574 .await,
575 );
576
577 // The default bounded repair (#5583) re-asks once — the fake's rule
578 // matches the repair too, so it fails identically and the run still
579 // fails loud instead of degrading to a null slot.
580 assert!(message.contains("responseSchema validation"), "{message}");
581 assert_eq!(
582 driver.spawn_count(),
583 2,
584 "default repair re-asks exactly once"
585 );
586 assert!(
587 driver.events().iter().any(|event| matches!(
588 event,
589 ProgressEvent::TaskSchemaRepairAttempted { attempt: 1, raw, .. }
590 if raw.contains("yes")
591 )),
592 "the failed first attempt should be receipted before the repair"
593 );
594 assert!(
595 driver.events().iter().any(|event| matches!(
596 event,
597 ProgressEvent::TaskSchemaValidationFailed { message, attempt: 2, .. }
598 if message.contains("responseSchema validation")
599 )),
600 "schema validation error should be emitted as workflow progress"
601 );
602 }
603
604 #[tokio::test]
605 async fn parallel_partial_mode_keeps_schema_failures_as_structured_slots() {
606 let driver = Arc::new(FakeDriver::new());
607 // Repair is disabled per-task so each slot fails terminally on its own
608 // reply; the mixed fan-out then exercises partial mode directly.
609 driver.on(
610 "good slot",
611 FakeReply::Complete(r#"{"refuted": true}"#.to_string()),
612 );
613 driver.on(
614 "bad slot",
615 FakeReply::Complete("not json at all".to_string()),
616 );
617 driver.on("dead slot", FakeReply::Fail("boom".to_string()));
618
619 let value = run(
620 &driver,
621 r#"
622 const results = await parallel([
623 () => task({
624 description: "good slot",
625 responseSchema: { "type": "object" },
626 }),
627 () => task({
628 description: "bad slot",
629 schemaRepairAttempts: 0,
630 responseSchema: { "type": "object" },
631 }),
632 () => task({ description: "dead slot" }),
633 ], { mode: "partial" });
634 return results.map((slot) =>
635 slot && typeof slot === "object" && slot.__taskError !== undefined
636 ? "error:" + slot.__taskError.kind
637 : slot === null
638 ? "null"
639 : "value:" + JSON.stringify(slot)
640 );
641 "#,
642 json!(null),
643 )
644 .await
645 .expect("partial mode completes the fan-out");
646
647 assert_eq!(
648 value,
649 json!([
650 "value:{\"refuted\":true}",
651 // The JS-level kind is the fatal "schema"; the finer decode kind
652 // (json_parse) lives on the receipt events, asserted below.
653 "error:schema",
654 // R9 behavior change: partial mode used to drop a dead subagent
655 // to `null` — indistinguishable from a slot that legitimately
656 // returned nothing. It is now a typed, inspectable failure.
657 "error:agent"
658 ])
659 );
660 // Every failed slot still leaves its terminal receipt.
661 assert!(
662 driver.events().iter().any(|event| matches!(
663 event,
664 ProgressEvent::TaskSchemaValidationFailed { kind, .. } if kind == "json_parse"
665 )),
666 "partial mode must not swallow the schema-failure receipt"
667 );
668 }
669
670 /// Audit R06-10: a cancel that lands while `task()` is still waiting on
671 /// admission ends the run instead of leaving it parked on the driver.
672 #[tokio::test]
673 async fn cancel_during_task_admission_ends_the_run() {
674 let driver = Arc::new(FakeDriver::new());
675 driver.on("gate", FakeReply::HoldAdmission);
676 let cancel = WorkflowRunCancel::new();
677 let run_cancel = cancel.clone();
678 let run_driver = driver.clone();
679 let handle = tokio::spawn(async move {
680 WorkflowVm::new()
681 .run_script_with_cancel(
682 r#"await task({ description: "gate" });"#,
683 json!(null),
684 run_driver as Arc<dyn codewhale_workflow_js::WorkflowDriver>,
685 run_cancel,
686 )
687 .await
688 });
689
690 tokio::time::timeout(Duration::from_secs(2), async {
691 while driver.spawn_count() == 0 {
692 tokio::task::yield_now().await;
693 }
694 })
695 .await
696 .expect("task should reach admission");
697 cancel.cancel();
698
699 let result = tokio::time::timeout(Duration::from_secs(5), handle)
700 .await
701 .expect("the run must end once cancelled")
702 .expect("VM task should join");
703 assert!(
704 matches!(result, Err(WorkflowJsError::Cancelled)),
705 "{result:?}"
706 );
707 }
708
709 #[tokio::test]
710 async fn parallel_partial_mode_still_fails_the_run_on_cancellation() {
711 let driver = Arc::new(FakeDriver::new());
712 driver.on("hang", FakeReply::Never);
713 let cancel = WorkflowRunCancel::new();
714 let run_cancel = cancel.clone();
715 let run_driver = driver.clone();
716 let handle = tokio::spawn(async move {
717 WorkflowVm::new()
718 .run_script_with_cancel(
719 r#"
720 await parallel([
721 () => task({ description: "hang", responseSchema: { "type": "object" } }),
722 ], { mode: "partial" });
723 "#,
724 json!(null),
725 run_driver as Arc<dyn codewhale_workflow_js::WorkflowDriver>,
726 run_cancel,
727 )
728 .await
729 });
730
731 tokio::time::timeout(Duration::from_secs(2), async {
732 while driver.spawn_count() == 0 {
733 tokio::task::yield_now().await;
734 }
735 })
736 .await
737 .expect("task should start");
738 cancel.cancel();
739
740 let result = handle.await.expect("VM task should join");
741 assert!(
742 matches!(result, Err(WorkflowJsError::Cancelled)),
743 "partial mode must not downgrade cancellation into a slot value: {result:?}"
744 );
745 }
746
747 #[tokio::test]
748 async fn pipeline_surfaces_response_schema_errors_instead_of_null() {
749 let driver = Arc::new(FakeDriver::new());
750 driver.on(
751 "bad schema",
752 FakeReply::Complete("not json at all".to_string()),
753 );
754
755 let message = script_message(
756 run(
757 &driver,
758 r#"
759 return await pipeline(
760 ["bad schema"],
761 (description) => task({
762 description,
763 schemaRepairAttempts: 0,
764 responseSchema: {
765 type: "object",
766 properties: { refuted: { type: "boolean" } },
767 required: ["refuted"],
768 },
769 }),
770 );
771 "#,
772 json!(null),
773 )
774 .await,
775 );
776
777 // Repair disabled: the first decode failure is terminal.
778 assert!(message.contains("not valid JSON"), "{message}");
779 assert_eq!(driver.spawn_count(), 1);
780 assert!(
781 driver.events().iter().any(|event| matches!(
782 event,
783 ProgressEvent::TaskSchemaValidationFailed { kind, attempt: 1, .. }
784 if kind == "json_parse"
785 )),
786 "a disabled repair must fail terminally on attempt 1 with the parse kind"
787 );
788 }
789
790 #[tokio::test]
791 async fn prose_wrapped_json_repairs_in_one_attempt() {
792 let driver = Arc::new(FakeDriver::new());
793 // First match wins: the repair spawn's description carries the
794 // "[schema repair 2]" marker, the first attempt's does not.
795 driver.on(
796 "[schema repair",
797 FakeReply::Complete(r#"{"refuted": true}"#.to_string()),
798 );
799 driver.on(
800 "score the claim",
801 FakeReply::Complete(
802 "Sure! Happy to help. Here is my verdict:\n\
803 ```json\n{\"refuted\": true}\n```\n\
804 Let me know if you need anything else."
805 .to_string(),
806 ),
807 );
808
809 let value = run(
810 &driver,
811 r#"
812 return await task({
813 description: "score the claim",
814 responseSchema: {
815 type: "object",
816 properties: { refuted: { type: "boolean" } },
817 required: ["refuted"],
818 },
819 });
820 "#,
821 json!(null),
822 )
823 .await
824 .expect("prose-wrapped JSON should repair in one attempt");
825
826 assert_eq!(value, json!({ "refuted": true }));
827 assert_eq!(driver.spawn_count(), 2);
828 let requests = driver.requests();
829 assert_eq!(requests[0].response_schema, requests[1].response_schema);
830 assert!(
831 requests[1].description.starts_with("[schema repair 2]"),
832 "the repair spawn must identify itself: {}",
833 requests[1].description
834 );
835 assert!(
836 requests[1].description.contains("score the claim"),
837 "the repair prompt must embed the original task"
838 );
839 assert!(
840 driver.events().iter().any(|event| matches!(
841 event,
842 ProgressEvent::TaskSchemaRepairAttempted {
843 kind, attempt: 1, raw, raw_truncated: false, ..
844 } if kind == "json_parse" && raw.contains("Happy to help")
845 )),
846 "the prose failure should be receipted with the parse kind"
847 );
848 assert!(
849 !driver
850 .events()
851 .iter()
852 .any(|event| matches!(event, ProgressEvent::TaskSchemaValidationFailed { .. })),
853 "a successful repair must not leave a terminal schema-failure receipt"
854 );
855 }
856
857 #[tokio::test]
858 async fn schema_violation_receipt_names_the_validation_kind() {
859 let driver = Arc::new(FakeDriver::new());
860 driver.on(
861 "[schema repair",
862 FakeReply::Complete(r#"{"refuted": false}"#.to_string()),
863 );
864 driver.on(
865 "check the gate",
866 FakeReply::Complete(r#"{"refuted":"no"}"#.to_string()),
867 );
868
869 let value = run(
870 &driver,
871 r#"
872 return await task({
873 description: "check the gate",
874 responseSchema: {
875 type: "object",
876 properties: { refuted: { type: "boolean" } },
877 required: ["refuted"],
878 },
879 });
880 "#,
881 json!(null),
882 )
883 .await
884 .expect("valid JSON of the wrong shape should repair");
885
886 assert_eq!(value, json!({ "refuted": false }));
887 assert!(
888 driver.events().iter().any(|event| matches!(
889 event,
890 ProgressEvent::TaskSchemaRepairAttempted { kind, message, .. }
891 if kind == "schema_validation"
892 && message.contains("responseSchema validation")
893 )),
894 "a parsed-but-invalid reply must receipt as schema_validation, not json_parse"
895 );
896 }
897
898 #[tokio::test]
899 async fn schema_repair_attempts_is_bounded_at_the_parse_gate() {
900 let driver = Arc::new(FakeDriver::new());
901 let message = script_message(
902 run(
903 &driver,
904 r#"
905 return await task({
906 description: "bound me",
907 schemaRepairAttempts: 4,
908 responseSchema: { "type": "object" },
909 });
910 "#,
911 json!(null),
912 )
913 .await,
914 );
915 assert!(message.contains("bounded to 3"), "{message}");
916 assert_eq!(
917 driver.spawn_count(),
918 0,
919 "no child may spawn for a bad option"
920 );
921 }
922
923 #[tokio::test]
924 async fn repair_is_refused_when_the_shared_budget_is_exhausted() {
925 let driver = Arc::new(FakeDriver::new());
926 // Attempt 1 is admitted with an empty pool and debits it fully at spawn.
927 driver.set_budget(Some(100), 100);
928 driver.on(
929 "spend it all",
930 FakeReply::Complete("sure thing, no JSON here".to_string()),
931 );
932
933 let message = script_message(
934 run(
935 &driver,
936 r#"
937 return await task({
938 description: "spend it all",
939 responseSchema: { "type": "object" },
940 });
941 "#,
942 json!(null),
943 )
944 .await,
945 );
946
947 assert!(
948 message.contains("repair skipped: budget exhausted"),
949 "{message}"
950 );
951 assert_eq!(
952 driver.spawn_count(),
953 1,
954 "the repair must not spawn on an empty pool"
955 );
956 assert!(
957 driver.events().iter().any(|event| matches!(
958 event,
959 ProgressEvent::TaskSchemaValidationFailed { attempt: 1, message, .. }
960 if message.contains("repair skipped: budget exhausted")
961 )),
962 "the refused repair must stay a schema failure with the reason named"
963 );
964 }
965
966 #[tokio::test]
967 async fn repair_is_refused_when_the_shared_wall_clock_is_spent() {
968 let driver = Arc::new(FakeDriver::new());
969 driver.on_with_delay(
970 "slow prose",
971 FakeReply::Complete("eventually, still not json".to_string()),
972 Duration::from_millis(1_100),
973 );
974
975 let message = script_message(
976 run(
977 &driver,
978 r#"
979 return await task({
980 description: "slow prose",
981 wallTimeSecs: 1,
982 responseSchema: { "type": "object" },
983 });
984 "#,
985 json!(null),
986 )
987 .await,
988 );
989
990 assert!(
991 message.contains("repair skipped: no wall-time left from wallTimeSecs"),
992 "{message}"
993 );
994 assert_eq!(
995 driver.spawn_count(),
996 1,
997 "the repair inherits the spent clock, not a fresh one"
998 );
999 }
1000
1001 #[tokio::test]
1002 async fn cancellation_during_repair_terminates_cleanly() {
1003 let driver = Arc::new(FakeDriver::new());
1004 driver.on("[schema repair", FakeReply::Never);
1005 driver.on(
1006 "hang the repair",
1007 FakeReply::Complete("prose, no json".to_string()),
1008 );
1009 let cancel = WorkflowRunCancel::new();
1010 let run_cancel = cancel.clone();
1011 let run_driver = driver.clone();
1012 let handle = tokio::spawn(async move {
1013 WorkflowVm::new()
1014 .run_script_with_cancel(
1015 r#"
1016 return await task({
1017 description: "hang the repair",
1018 responseSchema: { "type": "object" },
1019 });
1020 "#,
1021 json!(null),
1022 run_driver as Arc<dyn codewhale_workflow_js::WorkflowDriver>,
1023 run_cancel,
1024 )
1025 .await
1026 });
1027
1028 tokio::time::timeout(Duration::from_secs(2), async {
1029 while driver.spawn_count() < 2 {
1030 tokio::task::yield_now().await;
1031 }
1032 })
1033 .await
1034 .expect("repair should start");
1035 cancel.cancel();
1036
1037 let result = handle.await.expect("VM task should join");
1038 assert!(
1039 matches!(result, Err(WorkflowJsError::Cancelled)),
1040 "{result:?}"
1041 );
1042 assert!(
1043 !driver
1044 .events()
1045 .iter()
1046 .any(|event| matches!(event, ProgressEvent::TaskSchemaValidationFailed { .. })),
1047 "cancellation must not be rewritten into a schema failure"
1048 );
1049 }
1050
1051 #[tokio::test]
1052 async fn parallel_fail_fast_rejects_with_the_typed_slot_error() {
1053 let driver = Arc::new(FakeDriver::new());
1054 driver.on("beta", FakeReply::Fail("boom".to_string()));
1055 let value = run(
1056 &driver,
1057 r#"
1058 try {
1059 await parallel([
1060 () => task({ description: "alpha" }),
1061 () => task({ description: "beta" }),
1062 ], { mode: "fail-fast" });
1063 return "no-error";
1064 } catch (err) {
1065 return (err && err.kind) + ":" + (err && err.message);
1066 }
1067 "#,
1068 json!(null),
1069 )
1070 .await
1071 .unwrap();
1072 let text = value.as_str().unwrap();
1073 assert!(
1074 // R9: a child that ran and failed is `agent`, distinct from the
1075 // `script` kind a plain `throw` in a thunk produces.
1076 text.starts_with("agent:") && text.contains("boom"),
1077 "fail-fast must reject with the typed slot error: {text}"
1078 );
1079 assert!(
1080 driver.events().iter().any(|event| matches!(
1081 event,
1082 ProgressEvent::Log { message } if message.contains("fail-fast slot error")
1083 )),
1084 "fail-fast must leave a breadcrumb with the slot error"
1085 );
1086 }
1087
1088 #[tokio::test]
1089 async fn task_errors_carry_typed_kinds() {
1090 let driver = Arc::new(FakeDriver::new());
1091 driver.on("budget", FakeReply::BudgetExhausted("limit 10".to_string()));
1092 driver.on("cancelled", FakeReply::Cancelled);
1093 driver.on("admission", FakeReply::Reject("admission cap".to_string()));
1094 let value = run(
1095 &driver,
1096 r#"
1097 const kinds = {};
1098 for (const description of ["budget", "cancelled", "admission"]) {
1099 try {
1100 await task({ description });
1101 kinds[description] = "none";
1102 } catch (err) {
1103 kinds[description] = err && err.kind;
1104 }
1105 }
1106 return kinds;
1107 "#,
1108 json!(null),
1109 )
1110 .await
1111 .unwrap();
1112 assert_eq!(
1113 value,
1114 json!({"budget": "budget", "cancelled": "cancelled", "admission": "admission"})
1115 );
1116 }
1117
1118 #[tokio::test]
1119 async fn pipeline_fail_fast_rejects_instead_of_nulling_the_item() {
1120 let value = run(
1121 &Arc::new(FakeDriver::new()),
1122 r#"
1123 const stage = async (value) => {
1124 if (value === 1) throw new Error("stage boom");
1125 return value * 2;
1126 };
1127 try {
1128 await pipeline([1, 2], { stages: [stage], mode: "fail-fast" });
1129 return "no-error";
1130 } catch (err) {
1131 return (err && err.kind) + ":" + (err && err.message);
1132 }
1133 "#,
1134 json!(null),
1135 )
1136 .await
1137 .unwrap();
1138 assert_eq!(value, json!("script:stage boom"));
1139 }
1140
1141 #[tokio::test]
1142 async fn parallel_enforces_the_1000_item_cap_without_spawning() {
1143 let driver = Arc::new(FakeDriver::new());
1144 let value = run(
1145 &driver,
1146 r#"
1147 const thunks = new Array(1001).fill(() => task({ description: "x" }));
1148 try {
1149 await parallel(thunks);
1150 return "no-throw";
1151 } catch (err) {
1152 return String(err);
1153 }
1154 "#,
1155 json!(null),
1156 )
1157 .await
1158 .unwrap();
1159 let text = value.as_str().unwrap();
1160 assert!(text.contains("max 1000"), "{text}");
1161 assert_eq!(driver.spawn_count(), 0, "cap must reject before any spawn");
1162 }
1163
1164 #[tokio::test]
1165 async fn parallel_accepts_exactly_1000_items() {
1166 let driver = Arc::new(FakeDriver::new());
1167 let value = run(
1168 &driver,
1169 r#"
1170 const thunks = new Array(1000).fill(() => Promise.resolve(1));
1171 const results = await parallel(thunks);
1172 return results.length;
1173 "#,
1174 json!(null),
1175 )
1176 .await
1177 .unwrap();
1178 assert_eq!(value, json!(1000));
1179 }
1180
1181 #[tokio::test]
1182 async fn pipeline_has_no_barrier_between_stages() {
1183 let driver = Arc::new(FakeDriver::new());
1184 // Item A crawls through stage 1; item B sprints through both stages.
1185 driver.on_with_delay(
1186 "s1:A",
1187 FakeReply::Complete("A1".to_string()),
1188 Duration::from_millis(300),
1189 );
1190 driver.on_with_delay(
1191 "s1:B",
1192 FakeReply::Complete("B1".to_string()),
1193 Duration::from_millis(20),
1194 );
1195 driver.on_with_delay(
1196 "s2:B1",
1197 FakeReply::Complete("B2".to_string()),
1198 Duration::from_millis(20),
1199 );
1200 driver.on("s2:A1", FakeReply::Complete("A2".to_string()));
1201
1202 let value = run(
1203 &driver,
1204 r#"
1205 return await pipeline(
1206 ["A", "B"],
1207 (v) => task({ description: "s1:" + v }),
1208 (v) => task({ description: "s2:" + v }),
1209 );
1210 "#,
1211 json!(null),
1212 )
1213 .await
1214 .unwrap();
1215 assert_eq!(value, json!(["A2", "B2"]));
1216
1217 // B's stage 2 must have been requested while A was still in stage 1 —
1218 // per-item chains, no stage barrier.
1219 let descriptions = driver.request_descriptions();
1220 assert_eq!(descriptions[..2], ["s1:A".to_string(), "s1:B".to_string()]);
1221 assert_eq!(
1222 descriptions[2], "s2:B1",
1223 "expected B to reach stage 2 while A was still in stage 1: {descriptions:?}"
1224 );
1225 assert_eq!(descriptions[3], "s2:A1");
1226 }
1227
1228 #[tokio::test]
1229 async fn pipeline_stage_error_drops_only_that_item() {
1230 let driver = Arc::new(FakeDriver::new());
1231 driver.on("s1:B", FakeReply::Fail("boom".to_string()));
1232 let value = run(
1233 &driver,
1234 r#"
1235 return await pipeline(
1236 ["A", "B"],
1237 (v) => task({ description: "s1:" + v }),
1238 (v) => v + "+2",
1239 );
1240 "#,
1241 json!(null),
1242 )
1243 .await
1244 .unwrap();
1245 assert_eq!(value, json!(["done:s1:A+2", null]));
1246 }
1247
1248 #[tokio::test]
1249 async fn task_throws_once_budget_spent_reaches_total() {
1250 let driver = Arc::new(FakeDriver::new());
1251 driver.set_budget(Some(100), 60);
1252 let value = run(
1253 &driver,
1254 r#"
1255 let completed = 0;
1256 try {
1257 while (true) {
1258 await task({ description: "chunk " + completed });
1259 completed++;
1260 }
1261 } catch (err) {
1262 return { completed, message: String(err) };
1263 }
1264 "#,
1265 json!(null),
1266 )
1267 .await
1268 .unwrap();
1269 assert_eq!(value["completed"], json!(2));
1270 let message = value["message"].as_str().unwrap();
1271 assert!(message.contains("budget exhausted"), "{message}");
1272 assert_eq!(driver.spawn_count(), 2);
1273 }
1274
1275 #[tokio::test]
1276 async fn budget_globals_reflect_live_driver_snapshots() {
1277 let driver = Arc::new(FakeDriver::new());
1278 driver.set_budget(Some(1000), 100);
1279 let value = run(
1280 &driver,
1281 r#"
1282 const before = budget.remaining();
1283 await task({ description: "one" });
1284 return {
1285 total: budget.total,
1286 before,
1287 spent: budget.spent(),
1288 after: budget.remaining(),
1289 };
1290 "#,
1291 json!(null),
1292 )
1293 .await
1294 .unwrap();
1295 assert_eq!(
1296 value,
1297 json!({"total": 1000, "before": 1000, "spent": 100, "after": 900})
1298 );
1299 }
1300
1301 #[tokio::test]
1302 async fn unbounded_budget_reads_as_null_total_and_infinite_remaining() {
1303 let driver = Arc::new(FakeDriver::new());
1304 let value = run(
1305 &driver,
1306 "return budget.total === null && budget.remaining() === Infinity;",
1307 json!(null),
1308 )
1309 .await
1310 .unwrap();
1311 assert_eq!(value, json!(true));
1312 }
1313
1314 #[tokio::test]
1315 async fn lifetime_cap_throws_on_spawn_attempt_1001() {
1316 let driver = Arc::new(FakeDriver::new());
1317 let value = run(
1318 &driver,
1319 r#"
1320 let completed = 0;
1321 try {
1322 for (let i = 0; i < 1001; i++) {
1323 await task({ description: "t" + i });
1324 completed++;
1325 }
1326 return "no-throw";
1327 } catch (err) {
1328 return { completed, message: String(err) };
1329 }
1330 "#,
1331 json!(null),
1332 )
1333 .await
1334 .unwrap();
1335 assert_eq!(value["completed"], json!(WORKFLOW_LIFETIME_CAP));
1336 let message = value["message"].as_str().unwrap();
1337 assert!(message.contains("lifetime agent cap (1000)"), "{message}");
1338 assert_eq!(driver.spawn_count(), WORKFLOW_LIFETIME_CAP as usize);
1339 }
1340
1341 #[tokio::test]
1342 async fn response_schema_returns_the_parsed_validated_object() {
1343 let driver = Arc::new(FakeDriver::new());
1344 driver.on(
1345 "check",
1346 FakeReply::Complete(r#"{"refuted": true, "confidence": 0.9}"#.to_string()),
1347 );
1348 let value = run(
1349 &driver,
1350 r#"
1351 const verdict = await task({
1352 description: "check the claim",
1353 responseSchema: {
1354 type: "object",
1355 properties: { refuted: { type: "boolean" } },
1356 required: ["refuted"],
1357 },
1358 });
1359 return verdict.refuted === true ? "refuted" : "upheld";
1360 "#,
1361 json!(null),
1362 )
1363 .await
1364 .unwrap();
1365 assert_eq!(value, json!("refuted"));
1366 assert!(driver.requests()[0].response_schema.is_some());
1367 }
1368
1369 #[tokio::test]
1370 async fn response_schema_rejects_non_json_replies() {
1371 let driver = Arc::new(FakeDriver::new());
1372 driver.on(
1373 "check",
1374 FakeReply::Complete("definitely not json".to_string()),
1375 );
1376 let message = script_message(
1377 run(
1378 &driver,
1379 r#"
1380 return await task({
1381 description: "check",
1382 responseSchema: { type: "object" },
1383 });
1384 "#,
1385 json!(null),
1386 )
1387 .await,
1388 );
1389 assert!(message.contains("not valid JSON"), "{message}");
1390 }
1391
1392 #[tokio::test]
1393 async fn response_schema_rejects_schema_violations() {
1394 let driver = Arc::new(FakeDriver::new());
1395 driver.on(
1396 "check",
1397 FakeReply::Complete(r#"{"refuted": "yes"}"#.to_string()),
1398 );
1399 let message = script_message(
1400 run(
1401 &driver,
1402 r#"
1403 return await task({
1404 description: "check",
1405 responseSchema: {
1406 type: "object",
1407 properties: { refuted: { type: "boolean" } },
1408 required: ["refuted"],
1409 },
1410 });
1411 "#,
1412 json!(null),
1413 )
1414 .await,
1415 );
1416 assert!(message.contains("responseSchema validation"), "{message}");
1417 }
1418
1419 #[tokio::test]
1420 async fn determinism_ban_date_now() {
1421 let driver = Arc::new(FakeDriver::new());
1422 let message = script_message(run(&driver, "return Date.now();", json!(null)).await);
1423 assert!(message.contains("Date.now()"), "{message}");
1424 }
1425
1426 #[tokio::test]
1427 async fn determinism_ban_math_random() {
1428 let driver = Arc::new(FakeDriver::new());
1429 let message = script_message(run(&driver, "return Math.random();", json!(null)).await);
1430 assert!(message.contains("Math.random()"), "{message}");
1431 }
1432
1433 #[tokio::test]
1434 async fn determinism_ban_new_date() {
1435 let driver = Arc::new(FakeDriver::new());
1436 let message = script_message(run(&driver, "return new Date();", json!(null)).await);
1437 assert!(message.contains("unavailable"), "{message}");
1438 }
1439
1440 /// Explicit product surface for the sandboxed Workflow VM (#4129).
1441 ///
1442 /// Only these Workflow-owned calls may exist on `globalThis` beyond standard
1443 /// ECMAScript intrinsics. If a new host global is intentionally added, update
1444 /// this list in the same PR — the fail-closed inventory test below will break
1445 /// until the allowlist is extended deliberately.
1446 const WORKFLOW_ALLOWED_GLOBALS: &[&str] = &[
1447 "task", "parallel", "pipeline", "phase", "log", "budget", "args",
1448 ];
1449
1450 /// Host / Node / Deno / browser surfaces that must never leak into the VM.
1451 ///
1452 /// Standard ECMAScript intrinsics (`Object`, `Function`, `eval`, `Promise`, …)
1453 /// remain available; this list is only host escape hatches.
1454 const SANDBOX_BANNED_GLOBALS: &[&str] = &[
1455 "process",
1456 "require",
1457 "module",
1458 "exports",
1459 "__dirname",
1460 "__filename",
1461 "Buffer",
1462 "fs",
1463 "child_process",
1464 "os",
1465 "path",
1466 "net",
1467 "http",
1468 "https",
1469 "fetch",
1470 "XMLHttpRequest",
1471 "WebSocket",
1472 "Deno",
1473 "Bun",
1474 "Worker",
1475 ];
1476
1477 #[tokio::test]
1478 async fn sandbox_exposes_only_the_documented_workflow_calls() {
1479 let driver = Arc::new(FakeDriver::new());
1480 let value = run(
1481 &driver,
1482 r#"
1483 return {
1484 task: typeof task,
1485 parallel: typeof parallel,
1486 pipeline: typeof pipeline,
1487 phase: typeof phase,
1488 log: typeof log,
1489 budget: typeof budget,
1490 args: typeof args,
1491 };
1492 "#,
1493 json!({"ok": true}),
1494 )
1495 .await
1496 .unwrap();
1497 assert_eq!(
1498 value,
1499 json!({
1500 "task": "function",
1501 "parallel": "function",
1502 "pipeline": "function",
1503 "phase": "function",
1504 "log": "function",
1505 "budget": "object",
1506 "args": "object",
1507 })
1508 );
1509 // Keep the constant and the live typeof probe in lockstep.
1510 assert_eq!(
1511 WORKFLOW_ALLOWED_GLOBALS,
1512 &[
1513 "task", "parallel", "pipeline", "phase", "log", "budget", "args"
1514 ]
1515 );
1516 }
1517
1518 #[tokio::test]
1519 async fn sandbox_blocks_host_filesystem_shell_network_and_env_surfaces() {
1520 // Each probe must either throw / reject or resolve to a clearly absent
1521 // binding. We never allow a successful host escape.
1522 let probes: &[(&str, &str)] = &[
1523 (
1524 "process.env",
1525 r#"
1526 if (typeof process !== "undefined") {
1527 return process.env;
1528 }
1529 throw new Error("process is unavailable");
1530 "#,
1531 ),
1532 (
1533 "require('fs')",
1534 r#"
1535 if (typeof require === "function") {
1536 return require("fs");
1537 }
1538 throw new Error("require is unavailable");
1539 "#,
1540 ),
1541 (
1542 "import",
1543 r#"
1544 // Dynamic import is a module-loader surface; the VM has no loader.
1545 return await import("fs");
1546 "#,
1547 ),
1548 (
1549 "fetch",
1550 r#"
1551 if (typeof fetch === "function") {
1552 return await fetch("https://example.invalid/");
1553 }
1554 throw new Error("fetch is unavailable");
1555 "#,
1556 ),
1557 (
1558 "child_process",
1559 r#"
1560 if (typeof require === "function") {
1561 return require("child_process");
1562 }
1563 if (typeof child_process !== "undefined") {
1564 return child_process;
1565 }
1566 throw new Error("child_process is unavailable");
1567 "#,
1568 ),
1569 (
1570 "Deno.env",
1571 r#"
1572 if (typeof Deno !== "undefined") {
1573 return Deno.env.toObject();
1574 }
1575 throw new Error("Deno is unavailable");
1576 "#,
1577 ),
1578 ];
1579
1580 for (label, source) in probes {
1581 let driver = Arc::new(FakeDriver::new());
1582 let result = run(&driver, source, json!(null)).await;
1583 assert!(
1584 result.is_err(),
1585 "sandbox probe `{label}` must fail closed, got {result:?}"
1586 );
1587 // No driver side-effect is expected from a sandbox probe.
1588 assert_eq!(
1589 driver.spawn_count(),
1590 0,
1591 "probe `{label}` must not spawn tasks"
1592 );
1593 }
1594 }
1595
1596 #[tokio::test]
1597 async fn sandbox_global_inventory_fails_closed_on_new_host_leaks() {
1598 let driver = Arc::new(FakeDriver::new());
1599 let value = run(
1600 &driver,
1601 r#"
1602 // Own enumerable + non-enumerable names on the global object.
1603 // Anything beyond standard ECMAScript + the Workflow allowlist is a
1604 // regression that must break this test so new leaks cannot land quietly.
1605 const names = Reflect.ownKeys(globalThis)
1606 .map((k) => String(k))
1607 .sort();
1608 return names;
1609 "#,
1610 json!(null),
1611 )
1612 .await
1613 .unwrap();
1614 let names: Vec<String> = serde_json::from_value(value).expect("name list is a JSON array");
1615
1616 // Fail closed: none of the banned host surfaces may appear.
1617 for banned in SANDBOX_BANNED_GLOBALS {
1618 assert!(
1619 !names.iter().any(|n| n == *banned),
1620 "banned global `{banned}` leaked into the Workflow VM: {names:?}"
1621 );
1622 }
1623
1624 // Every Workflow-owned call must still be present.
1625 for allowed in WORKFLOW_ALLOWED_GLOBALS {
1626 assert!(
1627 names.iter().any(|n| n == *allowed),
1628 "expected Workflow global `{allowed}` missing from inventory: {names:?}"
1629 );
1630 }
1631
1632 // Internal host helpers must not be script-visible.
1633 for internal in [
1634 "__workflow_task",
1635 "__workflow_log",
1636 "__workflow_every_slot_failed",
1637 "__workflow_phase",
1638 "__workflow_budget_total",
1639 "__workflow_budget_spent",
1640 "__workflow_budget_remaining",
1641 ] {
1642 assert!(
1643 !names.iter().any(|n| n == internal),
1644 "internal host binding `{internal}` must stay hidden: {names:?}"
1645 );
1646 }
1647 }
1648
1649 #[tokio::test]
1650 async fn sandbox_rejects_commonjs_module_loader_and_eval_style_constructors() {
1651 let driver = Arc::new(FakeDriver::new());
1652 // `eval` / `Function` are standard ES, but if they are present they must
1653 // still be unable to reach host modules. The banned-global inventory above
1654 // already fails closed if Node-style loaders appear; this probe documents
1655 // the intended product message for module load attempts.
1656 let message = script_message(
1657 run(
1658 &driver,
1659 r#"
1660 if (typeof require === "function") {
1661 return require("node:fs");
1662 }
1663 throw new Error("require is unavailable");
1664 "#,
1665 json!(null),
1666 )
1667 .await,
1668 );
1669 assert!(
1670 message.contains("unavailable") || message.contains("require"),
1671 "{message}"
1672 );
1673 }
1674
1675 #[tokio::test]
1676 async fn dropping_the_run_future_cancels_outstanding_tasks() {
1677 let driver = Arc::new(FakeDriver::new());
1678 driver.on("hang", FakeReply::Never);
1679 let vm = WorkflowVm::new();
1680 {
1681 let fut = vm.run_script(
1682 "await task({ description: 'hang forever' }); return 'unreachable';",
1683 json!(null),
1684 driver.clone() as Arc<dyn codewhale_workflow_js::WorkflowDriver>,
1685 );
1686 let outcome = tokio::time::timeout(Duration::from_millis(400), fut).await;
1687 assert!(outcome.is_err(), "run should still be pending at timeout");
1688 // The timed-out future is dropped here.
1689 }
1690 assert!(
1691 driver.cancel_all_calls() >= 1,
1692 "dropping the run future must cancel outstanding driver tasks"
1693 );
1694 assert_eq!(driver.spawn_count(), 1);
1695 }
1696
1697 #[tokio::test]
1698 async fn parallel_does_not_continue_after_external_run_cancellation() {
1699 let driver = Arc::new(FakeDriver::new());
1700 driver.on("hang", FakeReply::Never);
1701 let cancel = WorkflowRunCancel::new();
1702 let run_cancel = cancel.clone();
1703 let run_driver = driver.clone();
1704 let handle = tokio::spawn(async move {
1705 WorkflowVm::new()
1706 .run_script_with_cancel(
1707 r#"
1708 await parallel([() => task({ description: "hang" })]);
1709 phase("unreachable after cancellation");
1710 return "wrong";
1711 "#,
1712 json!(null),
1713 run_driver as Arc<dyn codewhale_workflow_js::WorkflowDriver>,
1714 run_cancel,
1715 )
1716 .await
1717 });
1718
1719 tokio::time::timeout(Duration::from_secs(2), async {
1720 while driver.spawn_count() == 0 {
1721 tokio::task::yield_now().await;
1722 }
1723 })
1724 .await
1725 .expect("task should start");
1726 cancel.cancel();
1727
1728 let result = handle.await.expect("VM task should join");
1729 assert!(
1730 matches!(result, Err(WorkflowJsError::Cancelled)),
1731 "{result:?}"
1732 );
1733 assert!(
1734 !driver.events().iter().any(|event| matches!(
1735 event,
1736 ProgressEvent::Phase { title } if title == "unreachable after cancellation"
1737 )),
1738 "parallel() must not downgrade run cancellation into a null slot"
1739 );
1740 }
1741
1742 #[tokio::test]
1743 async fn script_error_rejects_cleanly_and_cancels_children() {
1744 let driver = Arc::new(FakeDriver::new());
1745 let result = run(
1746 &driver,
1747 r#"await task({ description: "quick" }); throw new Error("boom");"#,
1748 json!(null),
1749 )
1750 .await;
1751 let message = script_message(result);
1752 assert!(message.contains("boom"), "{message}");
1753 assert!(
1754 driver.cancel_all_calls() >= 1,
1755 "a failed run must cancel its cascade"
1756 );
1757 }
1758
1759 #[tokio::test]
1760 async fn log_and_phase_events_reach_the_driver_in_order() {
1761 let driver = Arc::new(FakeDriver::new());
1762 run(
1763 &driver,
1764 r#"
1765 phase("scan");
1766 log("a");
1767 log({ found: 2 });
1768 phase("verify");
1769 log("b");
1770 return null;
1771 "#,
1772 json!(null),
1773 )
1774 .await
1775 .unwrap();
1776 assert_eq!(
1777 driver.events(),
1778 vec![
1779 ProgressEvent::Phase {
1780 title: "scan".to_string()
1781 },
1782 ProgressEvent::Log {
1783 message: "a".to_string()
1784 },
1785 ProgressEvent::Log {
1786 message: r#"{"found":2}"#.to_string()
1787 },
1788 ProgressEvent::Phase {
1789 title: "verify".to_string()
1790 },
1791 ProgressEvent::Log {
1792 message: "b".to_string()
1793 },
1794 ]
1795 );
1796 }
1797
1798 #[tokio::test]
1799 async fn promise_all_of_tasks_resolves_concurrently() {
1800 let driver = Arc::new(FakeDriver::new());
1801 driver.on_with_delay(
1802 "left",
1803 FakeReply::Complete("L".to_string()),
1804 Duration::from_millis(50),
1805 );
1806 driver.on_with_delay(
1807 "right",
1808 FakeReply::Complete("R".to_string()),
1809 Duration::from_millis(50),
1810 );
1811 let started = std::time::Instant::now();
1812 let value = run(
1813 &driver,
1814 r#"
1815 const [a, b] = await Promise.all([
1816 task({ description: "left" }),
1817 task({ description: "right" }),
1818 ]);
1819 return a + "/" + b;
1820 "#,
1821 json!(null),
1822 )
1823 .await
1824 .unwrap();
1825 assert_eq!(value, json!("L/R"));
1826 // Two 50ms tasks awaited concurrently should not take ~100ms serially.
1827 // Generous bound to stay green on slow CI.
1828 assert!(
1829 started.elapsed() < Duration::from_millis(3000),
1830 "took {:?}",
1831 started.elapsed()
1832 );
1833 assert_eq!(driver.spawn_count(), 2);
1834 }
1835
1836 #[tokio::test]
1837 async fn export_default_async_function_runs_with_args() {
1838 let driver = Arc::new(FakeDriver::new());
1839 let source = r#"
1840 export default async function (args) {
1841 return { doubled: args.n * 2 };
1842 }
1843 "#;
1844 let value = run(&driver, source, json!({ "n": 21 })).await.unwrap();
1845 assert_eq!(value, json!({ "doubled": 42 }));
1846 }
1847
1848 #[tokio::test]
1849 async fn export_default_function_result_becomes_run_result() {
1850 let driver = Arc::new(FakeDriver::new());
1851 let source = r#"
1852 function helper() {
1853 return "from-helper";
1854 }
1855 export default function () {
1856 return helper();
1857 }
1858 "#;
1859 let value = run(&driver, source, json!(null)).await.unwrap();
1860 assert_eq!(value, json!("from-helper"));
1861 }
1862
1863 #[tokio::test]
1864 async fn export_default_non_function_value_is_returned() {
1865 let driver = Arc::new(FakeDriver::new());
1866 let value = run(&driver, "export default 7;", json!(null))
1867 .await
1868 .unwrap();
1869 assert_eq!(value, json!(7));
1870 }
1871
1872 #[tokio::test]
1873 async fn plain_scripts_are_untouched_by_export_desugaring() {
1874 let driver = Arc::new(FakeDriver::new());
1875 // A string literal mentioning `export default` must not trigger the
1876 // module desugaring path.
1877 let value = run(
1878 &driver,
1879 "const note = \"export default docs\";\nreturn note.length;",
1880 json!(null),
1881 )
1882 .await
1883 .unwrap();
1884 assert_eq!(value, json!(19));
1885 }
1886
1887 #[tokio::test]
1888 async fn export_default_examples_inside_multiline_text_are_not_desugared() {
1889 let driver = Arc::new(FakeDriver::new());
1890 let value = run(
1891 &driver,
1892 r#"
1893 const template = `
1894 export default async function (args) {
1895 return args;
1896 }
1897 `;
1898 /*
1899 export default function () {
1900 return "comment example";
1901 }
1902 */
1903 return template.includes("export default async function");
1904 "#,
1905 json!(null),
1906 )
1907 .await
1908 .unwrap();
1909 assert_eq!(value, json!(true));
1910 }
1911
1912 #[tokio::test]
1913 async fn task_accepts_agent_tool_spellings() {
1914 // The `agent` tool and `task()` are written by the same authors; a schema
1915 // that runs on one surface must not be an unknown-field error on the
1916 // other. snake_case spellings and `workspace_policy` are aliases.
1917 let driver = Arc::new(FakeDriver::new());
1918 let value = run(
1919 &driver,
1920 r#"
1921 return await task({
1922 prompt: "cross-surface schema",
1923 subagent_type: "implementer",
1924 workspace_policy: "worktree",
1925 write_authority: "worktree_write",
1926 write_roots: ["crates/tui/src"],
1927 token_budget: 5000,
1928 max_steps: 4,
1929 });
1930 "#,
1931 json!(null),
1932 )
1933 .await
1934 .unwrap();
1935 assert_eq!(value, json!("done:cross-surface schema"));
1936 let requests = driver.requests();
1937 assert_eq!(requests.len(), 1);
1938 assert!(
1939 requests[0].worktree,
1940 "workspace_policy worktree maps to worktree isolation"
1941 );
1942 assert_eq!(
1943 requests[0].write_authority.as_deref(),
1944 Some("worktree_write")
1945 );
1946 assert_eq!(requests[0].token_budget, Some(5000));
1947
1948 // "shared" is accepted and stays non-worktree; contradictions and unknown
1949 // values still fail loudly.
1950 let error = run(
1951 &driver,
1952 r#"return await task({ prompt: "x", workspacePolicy: "shared", worktree: true });"#,
1953 json!(null),
1954 )
1955 .await
1956 .unwrap_err();
1957 assert!(script_message(Err(error)).contains("conflicts with worktree"));
1958 let error = run(
1959 &driver,
1960 r#"return await task({ prompt: "x", workspacePolicy: "solo" });"#,
1961 json!(null),
1962 )
1963 .await
1964 .unwrap_err();
1965 assert!(script_message(Err(error)).contains("must be shared or worktree"));
1966 }
1967
1968 #[tokio::test]
1969 async fn vm_rejected_task_options_notify_the_driver() {
1970 // A task() whose options fail VM validation throws before spawn_task, and
1971 // inside parallel() that throw collapses to a null slot. The driver must
1972 // still receive a TaskRejected event so the run record can refuse to call
1973 // the run a plain success (morning-report issue #2).
1974 let driver = Arc::new(FakeDriver::new());
1975 let value = run(
1976 &driver,
1977 r#"
1978 return await parallel([
1979 () => task({ prompt: "bad slot", label: "L-bad", phase: "P1", cwd: "/absolute/path" }),
1980 ]);
1981 "#,
1982 json!(null),
1983 )
1984 .await
1985 .unwrap();
1986 assert_eq!(value, json!([null]));
1987 assert!(
1988 driver.requests().is_empty(),
1989 "no dispatch reached the driver"
1990 );
1991 let rejected: Vec<_> = driver
1992 .events()
1993 .into_iter()
1994 .filter_map(|event| match event {
1995 ProgressEvent::TaskRejected {
1996 label,
1997 phase,
1998 message,
1999 } => Some((label, phase, message)),
2000 _ => None,
2001 })
2002 .collect();
2003 assert_eq!(rejected.len(), 1, "one rejection event per refused slot");
2004 let (label, phase, message) = &rejected[0];
2005 assert_eq!(label.as_deref(), Some("L-bad"));
2006 assert_eq!(phase.as_deref(), Some("P1"));
2007 assert!(message.contains("bounded repo-relative paths"), "{message}");
2008 }
2009
2010 // ---------------------------------------------------------------------------
2011 // R9: typed slot errors, inspectable settled failures, explicit modes.
2012 // ---------------------------------------------------------------------------
2013
2014 /// Every way a `task()` can die gets its own kind, assigned by the host where
2015 /// the failure happened. Before R9 all six collapsed into two buckets
2016 /// ("budget"/"cancelled" if the message happened to say so, "task" otherwise),
2017 /// so a dead subagent and a typo'd script throw were the same thing.
2018 #[tokio::test]
2019 async fn every_task_failure_mode_carries_its_own_kind() {
2020 let driver = Arc::new(FakeDriver::new());
2021 driver.on("agent case", FakeReply::Fail("boom".to_string()));
2022 driver.on(
2023 "budget case",
2024 FakeReply::BudgetExhausted("limit 10".to_string()),
2025 );
2026 driver.on("cancelled case", FakeReply::Cancelled);
2027 driver.on(
2028 "admission case",
2029 FakeReply::Reject("admission cap".to_string()),
2030 );
2031 driver.on(
2032 "driver case",
2033 FakeReply::Unavailable("driver gone".to_string()),
2034 );
2035 driver.on("dropped case", FakeReply::DropCompletion);
2036 driver.on(
2037 "schema case",
2038 FakeReply::Complete("not json at all".to_string()),
2039 );
2040
2041 let value = run(
2042 &driver,
2043 r#"
2044 const kinds = {};
2045 const probe = async (name, opts) => {
2046 try {
2047 await task(opts);
2048 kinds[name] = "none";
2049 } catch (err) {
2050 kinds[name] = err && err.kind;
2051 }
2052 };
2053 await probe("agent", { description: "agent case" });
2054 await probe("budget", { description: "budget case" });
2055 await probe("cancelled", { description: "cancelled case" });
2056 await probe("admission", { description: "admission case" });
2057 await probe("driver", { description: "driver case" });
2058 await probe("dropped", { description: "dropped case" });
2059 await probe("schema", {
2060 description: "schema case",
2061 schemaRepairAttempts: 0,
2062 responseSchema: { type: "object" },
2063 });
2064 // A malformed options object never reaches a child either.
2065 await probe("bad-options", { description: "x", nosuchoption: 1 });
2066 try {
2067 await task("not an object");
2068 kinds["not-an-object"] = "none";
2069 } catch (err) {
2070 kinds["not-an-object"] = String(err && err.kind);
2071 }
2072 return kinds;
2073 "#,
2074 json!(null),
2075 )
2076 .await
2077 .unwrap();
2078
2079 assert_eq!(
2080 value,
2081 json!({
2082 "agent": "agent",
2083 "budget": "budget",
2084 "cancelled": "cancelled",
2085 "admission": "admission",
2086 "driver": "driver",
2087 "dropped": "driver",
2088 "schema": "schema",
2089 "bad-options": "admission",
2090 // A TypeError raised by the prelude's own argument check never
2091 // came from the host, so it is a script error, not a task kind.
2092 "not-an-object": "undefined",
2093 })
2094 );
2095 }
2096
2097 /// The classifier reads `Error.kind`, never the message text. A child is free
2098 /// to say "budget exhausted" or "responseSchema" in its own failure prose;
2099 /// under the old substring classifier that forged a fatal kind and aborted an
2100 /// otherwise healthy fan-out.
2101 #[tokio::test]
2102 async fn slot_kinds_cannot_be_forged_from_child_failure_text() {
2103 let driver = Arc::new(FakeDriver::new());
2104 driver.on(
2105 "liar",
2106 FakeReply::Fail(
2107 "the reviewer said the run cancelled because budget exhausted and responseSchema \
2108 validation failed"
2109 .to_string(),
2110 ),
2111 );
2112
2113 let value = run(
2114 &driver,
2115 r#"
2116 const results = await parallel([
2117 () => task({ description: "liar" }),
2118 () => task({ description: "honest" }),
2119 ]);
2120 return {
2121 slots: results,
2122 kinds: results.errors.map((entry) => entry.kind),
2123 };
2124 "#,
2125 json!(null),
2126 )
2127 .await
2128 .expect("a child's prose must not cancel the run");
2129
2130 assert_eq!(
2131 value,
2132 json!({
2133 "slots": [null, "done:honest"],
2134 "kinds": ["agent"],
2135 })
2136 );
2137 }
2138
2139 /// The settled default is unchanged on the wire — same slots, same length,
2140 /// same JSON — but the run can now ask why a slot is null instead of guessing.
2141 #[tokio::test]
2142 async fn settled_parallel_keeps_null_slots_and_attaches_an_inspectable_ledger() {
2143 let driver = Arc::new(FakeDriver::new());
2144 driver.on("beta", FakeReply::Fail("boom".to_string()));
2145 driver.on(
2146 "delta",
2147 FakeReply::BudgetExhausted("pool drained".to_string()),
2148 );
2149
2150 let value = run(
2151 &driver,
2152 r#"
2153 const results = await parallel([
2154 () => task({ description: "alpha" }),
2155 () => task({ description: "beta" }),
2156 () => task({ description: "gamma" }),
2157 () => task({ description: "delta" }),
2158 ]);
2159 return {
2160 slots: results,
2161 length: results.length,
2162 // Non-enumerable: the array still serializes as a plain array.
2163 encoded: JSON.stringify(results),
2164 errors: results.errors.map((entry) => ({
2165 index: entry.index,
2166 kind: entry.kind,
2167 says: entry.message.indexOf("boom") !== -1
2168 || entry.message.indexOf("pool drained") !== -1,
2169 })),
2170 };
2171 "#,
2172 json!(null),
2173 )
2174 .await
2175 .unwrap();
2176
2177 assert_eq!(
2178 value,
2179 json!({
2180 "slots": ["done:alpha", null, "done:gamma", null],
2181 "length": 4,
2182 "encoded": "[\"done:alpha\",null,\"done:gamma\",null]",
2183 "errors": [
2184 {"index": 1, "kind": "agent", "says": true},
2185 {"index": 3, "kind": "budget", "says": true},
2186 ],
2187 })
2188 );
2189 }
2190
2191 /// A clean fan-out still gets the ledger, empty and frozen — a script can read
2192 /// `results.errors.length` unconditionally.
2193 #[tokio::test]
2194 async fn a_clean_fan_out_still_reports_an_empty_frozen_error_ledger() {
2195 let value = run(
2196 &Arc::new(FakeDriver::new()),
2197 r#"
2198 const results = await parallel([() => task({ description: "alpha" })]);
2199 let mutated = false;
2200 try {
2201 results.errors = ["forged"];
2202 mutated = true;
2203 } catch (_) {
2204 mutated = false;
2205 }
2206 return {
2207 count: results.errors.length,
2208 frozen: Object.isFrozen(results.errors),
2209 mutated: mutated,
2210 stillEmpty: results.errors.length,
2211 };
2212 "#,
2213 json!(null),
2214 )
2215 .await
2216 .unwrap();
2217 assert_eq!(
2218 value,
2219 json!({"count": 0, "frozen": true, "mutated": false, "stillEmpty": 0})
2220 );
2221 }
2222
2223 /// `settled` is the spelling of today's default; naming it explicitly changes
2224 /// nothing.
2225 #[tokio::test]
2226 async fn explicit_settled_mode_matches_the_default() {
2227 let driver = Arc::new(FakeDriver::new());
2228 driver.on("beta", FakeReply::Fail("boom".to_string()));
2229 let value = run(
2230 &driver,
2231 r#"
2232 const thunks = () => [
2233 () => task({ description: "alpha" }),
2234 () => task({ description: "beta" }),
2235 ];
2236 const implicit = await parallel(thunks());
2237 const explicit = await parallel(thunks(), { mode: "settled" });
2238 return {
2239 implicit: implicit,
2240 explicit: explicit,
2241 sameKinds: JSON.stringify(implicit.errors.map((e) => e.kind))
2242 === JSON.stringify(explicit.errors.map((e) => e.kind)),
2243 };
2244 "#,
2245 json!(null),
2246 )
2247 .await
2248 .unwrap();
2249 assert_eq!(
2250 value,
2251 json!({
2252 "implicit": ["done:alpha", null],
2253 "explicit": ["done:alpha", null],
2254 "sameKinds": true,
2255 })
2256 );
2257 }
2258
2259 /// A typo'd mode used to read as `settled`: the author believed slots were
2260 /// now fatal while they kept silently dropping. It throws instead.
2261 #[tokio::test]
2262 async fn an_unknown_mode_is_refused_rather_than_silently_settled() {
2263 let driver = Arc::new(FakeDriver::new());
2264 let value = run(
2265 &driver,
2266 r#"
2267 const errs = [];
2268 for (const mode of ["failfast", "all-settled", 7]) {
2269 try {
2270 await parallel([() => task({ description: "alpha" })], { mode });
2271 errs.push("no-error");
2272 } catch (err) {
2273 errs.push(err.message);
2274 }
2275 }
2276 try {
2277 await pipeline([1], { stages: [(v) => v], mode: "failfast" });
2278 errs.push("no-error");
2279 } catch (err) {
2280 errs.push(err.message);
2281 }
2282 return errs;
2283 "#,
2284 json!(null),
2285 )
2286 .await
2287 .unwrap();
2288 let messages = value.as_array().unwrap();
2289 assert_eq!(messages.len(), 4, "{value}");
2290 for message in messages {
2291 let text = message.as_str().unwrap();
2292 assert!(
2293 text.contains("unknown mode") && text.contains("settled, fail-fast, partial"),
2294 "{text}"
2295 );
2296 }
2297 assert_eq!(
2298 driver.spawn_count(),
2299 0,
2300 "a refused mode must not spawn anything"
2301 );
2302 }
2303
2304 /// Partial mode is the "inspect every outcome" contract: no failure is erased,
2305 /// and none of them can be mistaken for a value.
2306 #[tokio::test]
2307 async fn partial_mode_types_every_non_cancellation_failure() {
2308 let driver = Arc::new(FakeDriver::new());
2309 driver.on("dead", FakeReply::Fail("boom".to_string()));
2310 driver.on("broke", FakeReply::BudgetExhausted("drained".to_string()));
2311 driver.on("refused", FakeReply::Reject("admission cap".to_string()));
2312
2313 let value = run(
2314 &driver,
2315 r#"
2316 const results = await parallel([
2317 () => task({ description: "alive" }),
2318 () => task({ description: "dead" }),
2319 () => task({ description: "broke" }),
2320 () => task({ description: "refused" }),
2321 ], { mode: "partial" });
2322 return {
2323 shapes: results.map((slot) =>
2324 slot && typeof slot === "object" && slot.__taskError
2325 ? slot.__taskError.kind + "@" + slot.__taskError.index
2326 : String(slot)
2327 ),
2328 ledger: results.errors.map((entry) => entry.kind),
2329 noNulls: results.every((slot) => slot !== null),
2330 };
2331 "#,
2332 json!(null),
2333 )
2334 .await
2335 .expect("partial mode completes the fan-out");
2336
2337 assert_eq!(
2338 value,
2339 json!({
2340 "shapes": ["done:alive", "agent@1", "budget@2", "admission@3"],
2341 "ledger": ["agent", "budget", "admission"],
2342 "noNulls": true,
2343 })
2344 );
2345 }
2346
2347 /// `pipeline` speaks the same three modes and keeps the same ledger.
2348 #[tokio::test]
2349 async fn pipeline_supports_settled_fail_fast_and_partial_with_a_ledger() {
2350 let driver = Arc::new(FakeDriver::new());
2351 driver.on("bad-1", FakeReply::Fail("boom".to_string()));
2352
2353 let value = run(
2354 &driver,
2355 r#"
2356 const stage = (item) => task({ description: item });
2357 const settled = await pipeline(["ok-0", "bad-1", "ok-2"], stage);
2358 const partial = await pipeline(["ok-0", "bad-1"], {
2359 stages: [stage],
2360 mode: "partial",
2361 });
2362 let failFast = "no-error";
2363 try {
2364 await pipeline(["ok-0", "bad-1"], { stages: [stage], mode: "fail-fast" });
2365 } catch (err) {
2366 failFast = err.kind + ":" + (err.message.indexOf("boom") !== -1);
2367 }
2368 return {
2369 settled: settled,
2370 settledLedger: settled.errors.map((e) => e.index + ":" + e.kind),
2371 partial: partial.map((slot) =>
2372 slot && typeof slot === "object" && slot.__taskError
2373 ? slot.__taskError.kind
2374 : slot
2375 ),
2376 failFast: failFast,
2377 };
2378 "#,
2379 json!(null),
2380 )
2381 .await
2382 .unwrap();
2383
2384 assert_eq!(
2385 value,
2386 json!({
2387 "settled": ["done:ok-0", null, "done:ok-2"],
2388 "settledLedger": ["1:agent"],
2389 "partial": ["done:ok-0", "agent"],
2390 "failFast": "agent:true",
2391 })
2392 );
2393 }
2394
2395 /// A fan-out where nothing survived is a dead fan-out. The default still
2396 /// resolves (existing scripts keep working) but the run log says so in a line
2397 /// the host status classifier and an operator can both find.
2398 #[tokio::test]
2399 async fn a_fan_out_where_every_slot_failed_says_so_in_the_run_log() {
2400 let driver = Arc::new(FakeDriver::new());
2401 driver.on("doomed", FakeReply::Fail("boom".to_string()));
2402 let value = run(
2403 &driver,
2404 r#"
2405 const results = await parallel([
2406 () => task({ description: "doomed a" }),
2407 () => task({ description: "doomed b" }),
2408 ]);
2409 return { slots: results, failed: results.errors.length };
2410 "#,
2411 json!(null),
2412 )
2413 .await
2414 .unwrap();
2415 assert_eq!(value, json!({"slots": [null, null], "failed": 2}));
2416 assert!(
2417 driver.events().iter().any(|event| matches!(
2418 event,
2419 ProgressEvent::Log { message }
2420 if message.contains("every slot failed (2 of 2)")
2421 && message.contains("no work survived")
2422 )),
2423 "a fully-failed fan-out must be named in the run log: {:?}",
2424 driver.events()
2425 );
2426 }
2427
2428 /// Cancellation stays fatal in every mode — it is the run's deadline, not a
2429 /// per-slot outcome — and partial mode does not get to keep it as a value.
2430 #[tokio::test]
2431 async fn cancellation_is_fatal_in_partial_pipeline_mode() {
2432 let driver = Arc::new(FakeDriver::new());
2433 driver.on("hang", FakeReply::Never);
2434 let cancel = WorkflowRunCancel::new();
2435 let run_cancel = cancel.clone();
2436 let run_driver = driver.clone();
2437 let handle = tokio::spawn(async move {
2438 WorkflowVm::new()
2439 .run_script_with_cancel(
2440 r#"
2441 await pipeline(["hang"], {
2442 stages: [(item) => task({ description: item })],
2443 mode: "partial",
2444 });
2445 "#,
2446 json!(null),
2447 run_driver as Arc<dyn codewhale_workflow_js::WorkflowDriver>,
2448 run_cancel,
2449 )
2450 .await
2451 });
2452
2453 tokio::time::timeout(Duration::from_secs(2), async {
2454 while driver.spawn_count() == 0 {
2455 tokio::task::yield_now().await;
2456 }
2457 })
2458 .await
2459 .expect("task should start");
2460 cancel.cancel();
2461
2462 let result = handle.await.expect("VM task should join");
2463 assert!(
2464 matches!(result, Err(WorkflowJsError::Cancelled)),
2465 "pipeline partial mode must not downgrade cancellation into a slot value: {result:?}"
2466 );
2467 }
2468
2469 /// The dropped-slot breadcrumb now names the kind, so the run log alone
2470 /// distinguishes "the agent failed" from "we never got to spawn it".
2471 #[tokio::test]
2472 async fn the_dropped_slot_breadcrumb_names_the_kind_and_the_slot() {
2473 let driver = Arc::new(FakeDriver::new());
2474 driver.on("beta", FakeReply::Fail("boom".to_string()));
2475 run(
2476 &driver,
2477 r#"
2478 return await parallel([
2479 () => task({ description: "alpha" }),
2480 () => task({ description: "beta" }),
2481 ]);
2482 "#,
2483 json!(null),
2484 )
2485 .await
2486 .unwrap();
2487 assert!(
2488 driver.events().iter().any(|event| matches!(
2489 event,
2490 ProgressEvent::Log { message }
2491 if message.contains("dropped a failed slot as null")
2492 && message.contains("kind=agent")
2493 && message.contains("slot 1")
2494 )),
2495 "the breadcrumb must name the kind and the slot: {:?}",
2496 driver.events()
2497 );
2498 }
2499
2500 struct EchoInvoker;
2501
2502 #[async_trait::async_trait]
2503 impl codewhale_workflow_js::ToolInvoker for EchoInvoker {
2504 async fn invoke(
2505 &self,
2506 request: codewhale_workflow_js::ToolCallRequest,
2507 ) -> Result<codewhale_workflow_js::ToolCallResponse, codewhale_workflow_js::DriverError> {
2508 use codewhale_workflow_js::{DriverError, ToolCallResponse};
2509 if request.tool == "boom" {
2510 return Ok(ToolCallResponse {
2511 ok: false,
2512 result: json!("kaput"),
2513 });
2514 }
2515 if request.tool == "deny" {
2516 return Err(DriverError::Rejected("nope".to_string()));
2517 }
2518 Ok(ToolCallResponse {
2519 ok: true,
2520 result: json!({ "echo": request.input }),
2521 })
2522 }
2523 }
2524
2525 async fn run_tools(source: &str) -> Result<serde_json::Value, WorkflowJsError> {
2526 let driver = Arc::new(FakeDriver::new());
2527 WorkflowVm::new()
2528 .run_tools_script(
2529 source,
2530 json!(null),
2531 driver.clone() as Arc<dyn codewhale_workflow_js::WorkflowDriver>,
2532 Arc::new(EchoInvoker) as Arc<dyn codewhale_workflow_js::ToolInvoker>,
2533 WorkflowRunCancel::new(),
2534 )
2535 .await
2536 }
2537
2538 #[tokio::test]
2539 async fn tools_surface_is_absent_without_invoker() {
2540 let driver = Arc::new(FakeDriver::new());
2541 let value = run(&driver, "return typeof tools;", json!(null))
2542 .await
2543 .unwrap();
2544 assert_eq!(value, json!("undefined"));
2545 }
2546
2547 #[tokio::test]
2548 async fn tools_call_round_trips() {
2549 let value =
2550 run_tools(r#"const r = await tools.call("read", { path: "x" }); return r.echo.path;"#)
2551 .await
2552 .unwrap();
2553 assert_eq!(value, json!("x"));
2554 }
2555
2556 #[tokio::test]
2557 async fn native_callbacks_stay_on_the_originating_host_runtime() {
2558 use codewhale_workflow_js::{
2559 BudgetSnapshot, DriverError, SpawnedTask, TaskRequest, ToolCallRequest, ToolCallResponse,
2560 ToolInvoker, WorkflowDriver,
2561 };
2562 use std::sync::atomic::{AtomicUsize, Ordering};
2563
2564 struct OriginCallbacks {
2565 origin: std::thread::ThreadId,
2566 calls: AtomicUsize,
2567 driver: FakeDriver,
2568 }
2569
2570 #[async_trait::async_trait]
2571 impl WorkflowDriver for OriginCallbacks {
2572 async fn spawn_task(&self, request: TaskRequest) -> Result<SpawnedTask, DriverError> {
2573 assert_eq!(std::thread::current().id(), self.origin);
2574 self.calls.fetch_add(1, Ordering::SeqCst);
2575 self.driver.spawn_task(request).await
2576 }
2577
2578 fn cancel_all(&self) {
2579 self.driver.cancel_all();
2580 }
2581
2582 fn budget(&self) -> BudgetSnapshot {
2583 self.driver.budget()
2584 }
2585
2586 fn progress(&self, event: ProgressEvent) {
2587 self.driver.progress(event);
2588 }
2589 }
2590
2591 #[async_trait::async_trait]
2592 impl ToolInvoker for OriginCallbacks {
2593 async fn invoke(&self, request: ToolCallRequest) -> Result<ToolCallResponse, DriverError> {
2594 assert_eq!(std::thread::current().id(), self.origin);
2595 self.calls.fetch_add(1, Ordering::SeqCst);
2596 EchoInvoker.invoke(request).await
2597 }
2598 }
2599
2600 let callbacks = Arc::new(OriginCallbacks {
2601 origin: std::thread::current().id(),
2602 calls: AtomicUsize::new(0),
2603 driver: FakeDriver::new(),
2604 });
2605 let value = WorkflowVm::new()
2606 .run_tools_script(
2607 r#"const child = await task({ description: "origin-thread" });
2608 const tool = await tools.call("read", { path: "original" });
2609 return { child, path: tool.echo.path };"#,
2610 json!(null),
2611 callbacks.clone(),
2612 callbacks.clone(),
2613 WorkflowRunCancel::new(),
2614 )
2615 .await
2616 .unwrap();
2617 assert_eq!(
2618 value,
2619 json!({"child": "done:origin-thread", "path": "original"})
2620 );
2621 assert_eq!(callbacks.calls.load(Ordering::SeqCst), 2);
2622 assert_eq!(callbacks.driver.spawn_count(), 1);
2623 }
2624
2625 #[tokio::test]
2626 async fn tools_call_refusal_throws_admission_kind() {
2627 let value = run_tools(
2628 r#"try { await tools.call("deny", {}); return "no-throw"; } catch (e) { return e.kind; }"#,
2629 )
2630 .await
2631 .unwrap();
2632 assert_eq!(value, json!("admission"));
2633 }
2634
2635 #[tokio::test]
2636 async fn tools_call_failure_throws_agent_kind() {
2637 let value = run_tools(
2638 r#"try { await tools.call("boom", {}); return "no-throw"; } catch (e) { return e.kind; }"#,
2639 )
2640 .await
2641 .unwrap();
2642 assert_eq!(value, json!("agent"));
2643 }
2644
2645 #[tokio::test]
2646 async fn tools_call_cap_rejects_runaway_loops() {
2647 let err = run_tools(
2648 r#"for (let i = 0; i < 55; i++) { await tools.call("read", {}); } return "never";"#,
2649 )
2650 .await;
2651 let message = script_message(err);
2652 assert!(
2653 message.contains("per-run tool-call cap"),
2654 "unexpected: {message}"
2655 );
2656 }
2657
2658 /// A platform-native absolute path for `unix_path` (`/a/b`): unchanged on
2659 /// Unix, `C:\a\b` on Windows, where `/a/b` has no drive and is not absolute.
2660 fn native_abs(unix_path: &str) -> String {
2661 if cfg!(windows) {
2662 format!("C:{}", unix_path.replace('/', "\\"))
2663 } else {
2664 unix_path.to_string()
2665 }
2666 }
2667
2668 /// An absolute `cwd` inside the run's workspace normalizes to the
2669 /// repo-relative form; outside it, through `..`, or with no known workspace
2670 /// it is still refused (the trust boundary is unchanged).
2671 #[test]
2672 fn absolute_cwd_inside_the_workspace_normalizes_and_outside_is_refused() {
2673 use codewhale_workflow_js::{normalize_task_cwd, normalize_task_cwd_in};
2674 use std::path::PathBuf;
2675 let workspace = PathBuf::from(native_abs("/Volumes/VIXinSSD/CW"));
2676 let workspace = Some(workspace.as_path());
2677 assert_eq!(
2678 normalize_task_cwd_in(&native_abs("/Volumes/VIXinSSD/CW/codewhale"), workspace).unwrap(),
2679 "codewhale"
2680 );
2681 assert_eq!(
2682 normalize_task_cwd_in(&native_abs("/Volumes/VIXinSSD/CW/./crates/tui"), workspace).unwrap(),
2683 "crates/tui"
2684 );
2685 assert_eq!(
2686 normalize_task_cwd_in(&native_abs("/Volumes/VIXinSSD/CW/"), workspace).unwrap(),
2687 "."
2688 );
2689 assert_eq!(
2690 normalize_task_cwd_in("crates/tui", workspace).unwrap(),
2691 normalize_task_cwd("crates/tui").unwrap()
2692 );
2693 for outside in [
2694 native_abs("/Volumes/VIXinSSD/CW-other/codewhale"),
2695 native_abs("/Volumes/VIXinSSD"),
2696 native_abs("/etc"),
2697 native_abs("/Volumes/VIXinSSD/CW/../secrets"),
2698 native_abs("/Volumes/VIXinSSD/CW/codewhale/../../secrets"),
2699 native_abs("/Volumes/VIXinSSD/X/../CW/codewhale"),
2700 ] {
2701 let error = normalize_task_cwd_in(&outside, workspace).unwrap_err();
2702 assert!(
2703 error.contains("outside the workspace") || error.contains("parent traversal"),
2704 "{outside}: {error}"
2705 );
2706 }
2707 let error =
2708 normalize_task_cwd_in(&native_abs("/Volumes/VIXinSSD/CW/codewhale"), None).unwrap_err();
2709 assert!(error.contains("bounded repo-relative paths"), "{error}");
2710 }
2711
2712 /// Windows spellings of the same workspace path: a verbatim `\\?\` prefix on
2713 /// either side, forward slashes, and a different case all match; another
2714 /// drive or a sibling directory does not.
2715 #[cfg(windows)]
2716 #[test]
2717 fn absolute_cwd_matches_windows_spellings_of_the_workspace() {
2718 use codewhale_workflow_js::normalize_task_cwd_in;
2719 use std::path::Path;
2720 let plain = Path::new(r"C:\Users\dev\CW");
2721 let verbatim = Path::new(r"\\?\C:\Users\dev\CW");
2722 for workspace in [plain, verbatim] {
2723 for inside in [
2724 r"C:\Users\dev\CW\codewhale",
2725 "C:/Users/dev/CW/codewhale",
2726 r"c:\users\DEV\cw\codewhale",
2727 r"\\?\C:\Users\dev\CW\codewhale",
2728 ] {
2729 assert_eq!(
2730 normalize_task_cwd_in(inside, Some(workspace)).unwrap(),
2731 "codewhale",
2732 "{inside} in {}",
2733 workspace.display()
2734 );
2735 }
2736 for outside in [
2737 r"D:\Users\dev\CW\codewhale",
2738 r"C:\Users\dev\CW-other",
2739 r"C:\Users\dev\CW\..\secrets",
2740 r"\\server\share\Users\dev\CW",
2741 ] {
2742 let error = normalize_task_cwd_in(outside, Some(workspace)).unwrap_err();
2743 assert!(
2744 error.contains("outside the workspace") || error.contains("parent traversal"),
2745 "{outside}: {error}"
2746 );
2747 }
2748 }
2749 // A UNC workspace matches its own share, case-insensitively.
2750 let unc = Path::new(r"\\Server\Share\CW");
2751 assert_eq!(
2752 normalize_task_cwd_in(r"\\server\share\cw\codewhale", Some(unc)).unwrap(),
2753 "codewhale"
2754 );
2755 }
2756
2756 lines RUST