返回 CodeWhale
task_spec.rs
根目录 / crates / tui / src / fleet / task_spec.rs
1 //! Typed task-spec loading, artifact refs, deterministic scorers, and receipts.
2
3 #![allow(dead_code)]
4
5 use std::collections::{BTreeMap, BTreeSet};
6 use std::path::{Path, PathBuf};
7
8 use anyhow::{Context, Result, bail};
9 use chrono::{SecondsFormat, Utc};
10 use codewhale_protocol::fleet::*;
11 use regex::Regex;
12 use serde::{Deserialize, Serialize};
13 use serde_json::{Value, json};
14
15 use super::ledger::FleetLedger;
16
17 const MAX_SCORER_READ_BYTES: u64 = 1_000_000;
18 const MAX_FLEET_ID_BYTES: usize = 128;
19 const MAX_FLEET_NAME_BYTES: usize = 256;
20
21 #[derive(Debug, Clone, Serialize, Deserialize)]
22 pub struct FleetTaskSpecDocument {
23 #[serde(default)]
24 pub name: Option<String>,
25 #[serde(default)]
26 pub labels: BTreeMap<String, String>,
27 #[serde(default)]
28 #[serde(skip_serializing_if = "Option::is_none")]
29 /// Legacy replay/input compatibility only. New run validation rejects it;
30 /// Runtime policy owns execution authority.
31 pub security_policy: Option<FleetSecurityPolicy>,
32 #[serde(default, alias = "worker_specs")]
33 pub workers: Vec<FleetWorkerSpec>,
34 #[serde(default)]
35 pub tasks: Vec<FleetTaskSpec>,
36 /// Optional run-wide usage ceiling (R6, #5567), e.g.
37 /// `usage_ceiling = { max_total_tokens = 2_000_000 }`.
38 #[serde(default)]
39 #[serde(skip_serializing_if = "Option::is_none")]
40 pub usage_ceiling: Option<codewhale_protocol::fleet::FleetUsageCeiling>,
41 }
42
43 /// A parsed spec file in one of its three accepted shapes. The shape is
44 /// chosen from the file's structure first ([`FleetTaskSpecShape::detect`]) and
45 /// only then deserialized into the matching type, so a malformed spec reports
46 /// the real field error instead of serde's opaque "did not match any variant
47 /// of untagged enum".
48 #[derive(Debug, Clone)]
49 enum FleetTaskSpecFile {
50 Document(FleetTaskSpecDocument),
51 Tasks(Vec<FleetTaskSpec>),
52 Single(Box<FleetTaskSpec>),
53 }
54
55 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
56 enum FleetTaskSpecShape {
57 /// `{ name?, labels?, workers?, tasks = [...] }`
58 Document,
59 /// A bare JSON array of task objects.
60 Tasks,
61 /// A single task object (`{ id, name, instructions, ... }`).
62 Single,
63 }
64
65 #[derive(Debug, Clone, Copy, PartialEq, Eq)]
66 enum FleetTaskSpecFormat {
67 Json,
68 Toml,
69 }
70
71 impl FleetTaskSpecFormat {
72 fn label(self) -> &'static str {
73 match self {
74 Self::Json => "JSON",
75 Self::Toml => "TOML",
76 }
77 }
78 }
79
80 /// Top-level keys that only a spec document carries.
81 const DOCUMENT_KEYS: &[&str] = &["tasks", "workers", "worker_specs"];
82 /// Top-level keys that mark a bare single-task file.
83 const SINGLE_TASK_KEYS: &[&str] = &["id", "instructions"];
84
85 impl FleetTaskSpecShape {
86 fn detect(value: &Value) -> Result<Self> {
87 match value {
88 Value::Array(_) => Ok(Self::Tasks),
89 Value::Object(map) => {
90 if DOCUMENT_KEYS.iter().any(|key| map.contains_key(*key)) {
91 Ok(Self::Document)
92 } else if SINGLE_TASK_KEYS.iter().any(|key| map.contains_key(*key)) {
93 Ok(Self::Single)
94 } else {
95 // Name/labels-only (or empty) objects are documents; the
96 // validator then reports the missing `tasks`.
97 Ok(Self::Document)
98 }
99 }
100 other => bail!(
101 "a fleet task spec must be a document object with `tasks`, an array of task objects, or a single task object; found {}",
102 json_kind(other)
103 ),
104 }
105 }
106
107 fn label(self) -> &'static str {
108 match self {
109 Self::Document => "spec document",
110 Self::Tasks => "task array",
111 Self::Single => "single task",
112 }
113 }
114 }
115
116 fn json_kind(value: &Value) -> &'static str {
117 match value {
118 Value::Null => "null",
119 Value::Bool(_) => "a boolean",
120 Value::Number(_) => "a number",
121 Value::String(_) => "a string",
122 Value::Array(_) => "an array",
123 Value::Object(_) => "an object",
124 }
125 }
126
127 fn parse_task_spec_file(raw: &str, format: FleetTaskSpecFormat) -> Result<FleetTaskSpecFile> {
128 // Pass 1: syntax only, to learn the shape.
129 let value = match format {
130 FleetTaskSpecFormat::Json => serde_json::from_str::<Value>(raw)
131 .map_err(|err| anyhow::anyhow!("invalid JSON: {err}"))?,
132 FleetTaskSpecFormat::Toml => {
133 let table = toml::from_str::<toml::Table>(raw)
134 .map_err(|err| anyhow::anyhow!("invalid TOML: {err}"))?;
135 serde_json::to_value(table).context("converting TOML fleet task spec")?
136 }
137 };
138 let shape = FleetTaskSpecShape::detect(&value)?;
139
140 // Pass 2: typed deserialize of exactly that shape, from the raw text so
141 // the error keeps its line/column.
142 fn typed<T: serde::de::DeserializeOwned>(
143 raw: &str,
144 format: FleetTaskSpecFormat,
145 ) -> std::result::Result<T, String> {
146 match format {
147 FleetTaskSpecFormat::Json => serde_json::from_str::<T>(raw).map_err(|e| e.to_string()),
148 FleetTaskSpecFormat::Toml => toml::from_str::<T>(raw).map_err(|e| e.to_string()),
149 }
150 }
151 let parsed = match shape {
152 FleetTaskSpecShape::Document => {
153 typed::<FleetTaskSpecDocument>(raw, format).map(FleetTaskSpecFile::Document)
154 }
155 FleetTaskSpecShape::Tasks => {
156 typed::<Vec<FleetTaskSpec>>(raw, format).map(FleetTaskSpecFile::Tasks)
157 }
158 FleetTaskSpecShape::Single => typed::<FleetTaskSpec>(raw, format)
159 .map(|task| FleetTaskSpecFile::Single(Box::new(task))),
160 };
161 parsed.map_err(|err| {
162 let location = locate_spec_error(shape, &value)
163 .map(|loc| format!(" at {loc}"))
164 .unwrap_or_default();
165 anyhow::anyhow!(
166 "{} {}{location}: {}",
167 format.label(),
168 shape.label(),
169 err.trim()
170 )
171 })
172 }
173
174 /// Name the first task (or worker) entry that fails to deserialize on its
175 /// own, e.g. `tasks[1] (id "review")`, so a long spec's error points at the
176 /// entry and not only at a line number.
177 fn locate_spec_error(shape: FleetTaskSpecShape, value: &Value) -> Option<String> {
178 fn first_bad<T: serde::de::DeserializeOwned>(prefix: &str, items: &[Value]) -> Option<String> {
179 items.iter().enumerate().find_map(|(index, item)| {
180 serde_json::from_value::<T>(item.clone()).err().map(|_| {
181 match item.get("id").and_then(Value::as_str) {
182 Some(id) => format!("{prefix}[{index}] (id {id:?})"),
183 None => format!("{prefix}[{index}]"),
184 }
185 })
186 })
187 }
188 match shape {
189 FleetTaskSpecShape::Document => {
190 if let Some(tasks) = value.get("tasks").and_then(Value::as_array)
191 && let Some(loc) = first_bad::<FleetTaskSpec>("tasks", tasks)
192 {
193 return Some(loc);
194 }
195 let workers = value
196 .get("workers")
197 .or_else(|| value.get("worker_specs"))
198 .and_then(Value::as_array)?;
199 first_bad::<FleetWorkerSpec>("workers", workers)
200 }
201 FleetTaskSpecShape::Tasks => first_bad::<FleetTaskSpec>("", value.as_array()?),
202 FleetTaskSpecShape::Single => None,
203 }
204 }
205
206 impl FleetTaskSpecFile {
207 fn into_document(self, fallback_name: String) -> FleetTaskSpecDocument {
208 match self {
209 Self::Document(mut doc) => {
210 if doc.name.as_deref().is_none_or(str::is_empty) {
211 doc.name = Some(fallback_name);
212 }
213 doc
214 }
215 Self::Tasks(tasks) => FleetTaskSpecDocument {
216 name: Some(fallback_name),
217 labels: BTreeMap::new(),
218 security_policy: None,
219 workers: Vec::new(),
220 tasks,
221 usage_ceiling: None,
222 },
223 Self::Single(task) => FleetTaskSpecDocument {
224 name: Some(fallback_name),
225 labels: BTreeMap::new(),
226 security_policy: None,
227 workers: Vec::new(),
228 tasks: vec![*task],
229 usage_ceiling: None,
230 },
231 }
232 }
233 }
234
235 /// The worker's visible final answer as carried by the terminal exec
236 /// `metadata` receipt: `excerpt` is already bounded and secret-redacted by the
237 /// emitter (`visible_final_answer_excerpt`), `chars` is the real
238 /// pre-truncation length (`visible_final_answer_chars`).
239 #[derive(Debug, Clone, PartialEq, Eq)]
240 pub struct FleetWorkerFinalAnswer {
241 pub excerpt: String,
242 pub chars: usize,
243 }
244
245 impl FleetWorkerFinalAnswer {
246 /// The receipt note for a task whose only deliverable is its answer text.
247 /// The excerpt is used verbatim: it was bounded and redacted once at the
248 /// emitter, and the ledger redacts receipt notes again on write.
249 pub fn receipt_note(&self) -> String {
250 format!(
251 "worker produced {} characters of deliverable: {}",
252 self.chars, self.excerpt
253 )
254 }
255 }
256
257 #[derive(Debug, Clone)]
258 pub struct FleetTaskVerificationInput {
259 pub run_id: FleetRunId,
260 pub task_id: String,
261 pub worker_id: String,
262 /// Durable lease generation whose result is being verified.
263 pub attempt: u32,
264 pub exit_code: Option<i32>,
265 pub artifacts: Vec<FleetArtifactRef>,
266 /// The worker's visible final answer, as reported by its terminal exec
267 /// receipt. Report/summary tasks with no scorer and no file artifact
268 /// surface this as their deliverable instead of "no verifiable output".
269 pub final_answer: Option<FleetWorkerFinalAnswer>,
270 /// Saved exec session id holding the worker's full transcript, when the
271 /// worker persisted one on completion.
272 pub saved_session_id: Option<String>,
273 /// Resolved-route snapshot to persist on the receipt (#3154).
274 pub resolved_route: Option<FleetResolvedRoute>,
275 /// Effective worker authority snapshot to persist on the receipt (#3211).
276 pub effective_permissions: Option<FleetEffectivePermissions>,
277 }
278
279 #[derive(Debug, Clone)]
280 pub struct FleetTaskVerification {
281 pub result: FleetTaskResult,
282 pub failure_kind: Option<FleetTaskFailureKind>,
283 pub score: FleetScore,
284 pub evidence: Vec<String>,
285 }
286
287 pub fn load_task_spec_document(path: &Path) -> Result<FleetTaskSpecDocument> {
288 let raw = std::fs::read_to_string(path)
289 .with_context(|| format!("reading fleet task spec {}", path.display()))?;
290 let fallback_name = path
291 .file_stem()
292 .and_then(|s| s.to_str())
293 .filter(|s| !s.is_empty())
294 .unwrap_or("fleet-run")
295 .to_string();
296 let format = match path.extension().and_then(|s| s.to_str()) {
297 Some("toml") => FleetTaskSpecFormat::Toml,
298 _ => FleetTaskSpecFormat::Json,
299 };
300 let parsed = parse_task_spec_file(&raw, format).with_context(|| {
301 format!(
302 "parsing {} fleet task spec {}",
303 format.label(),
304 path.display()
305 )
306 })?;
307 let doc = parsed.into_document(fallback_name);
308 validate_task_spec_document(&doc)?;
309 Ok(doc)
310 }
311
312 pub fn validate_task_spec_document(doc: &FleetTaskSpecDocument) -> Result<()> {
313 if doc.security_policy.is_some() {
314 bail!(
315 "fleet task spec security_policy is a legacy compatibility field, not executable Fleet identity; configure trust, secrets, approvals, sandboxing, and tool authority through Runtime policy"
316 );
317 }
318 if doc.tasks.is_empty() {
319 bail!("fleet task spec must include at least one task");
320 }
321 let mut ids = BTreeSet::new();
322 for task in &doc.tasks {
323 validate_fleet_identity("task id", &task.id)?;
324 if !ids.insert(task.id.clone()) {
325 bail!("duplicate fleet task id {}", task.id);
326 }
327 validate_fleet_name(&format!("task {} name", task.id), &task.name)?;
328 if task.instructions.trim().is_empty() {
329 bail!("fleet task {} instructions cannot be empty", task.id);
330 }
331 if let Some(objective) = &task.objective
332 && objective.trim().is_empty()
333 {
334 bail!("fleet task {} objective cannot be empty", task.id);
335 }
336 validate_worker_profile(&task.id, task.worker.as_ref())?;
337 if task
338 .metadata
339 .contains_key(super::worker_runtime::FROZEN_FLEET_MEMBER_METADATA_KEY)
340 {
341 bail!(
342 "fleet task {} metadata key {} is reserved for the durable Runtime selection receipt",
343 task.id,
344 super::worker_runtime::FROZEN_FLEET_MEMBER_METADATA_KEY
345 );
346 }
347 validate_tags(&task.id, &task.tags)?;
348 validate_workspace_requirements(task)?;
349 }
350 let mut worker_ids = BTreeSet::new();
351 for worker in &doc.workers {
352 validate_fleet_identity("worker id", &worker.id)?;
353 if !worker_ids.insert(worker.id.clone()) {
354 bail!("duplicate fleet worker id {}", worker.id);
355 }
356 validate_fleet_name(&format!("worker {} name", worker.id), &worker.name)?;
357 if worker.trust_level.is_some() {
358 bail!(
359 "fleet worker {} trust_level is a legacy compatibility field, not Fleet identity; configure execution authority through Runtime policy",
360 worker.id
361 );
362 }
363 }
364 Ok(())
365 }
366
367 fn validate_fleet_identity(field: &str, value: &str) -> Result<()> {
368 if value.is_empty() {
369 bail!("fleet {field} cannot be empty");
370 }
371 if value.len() > MAX_FLEET_ID_BYTES || !value.chars().all(is_worker_token_char) {
372 bail!(
373 "fleet {field} must be a simple ASCII token no longer than {MAX_FLEET_ID_BYTES} bytes"
374 );
375 }
376 Ok(())
377 }
378
379 fn validate_fleet_name(field: &str, value: &str) -> Result<()> {
380 if value.trim().is_empty() {
381 bail!("fleet {field} cannot be empty");
382 }
383 if value.len() > MAX_FLEET_NAME_BYTES || value.chars().any(char::is_control) {
384 bail!(
385 "fleet {field} must be one printable line no longer than {MAX_FLEET_NAME_BYTES} bytes"
386 );
387 }
388 Ok(())
389 }
390
391 fn validate_worker_profile(task_id: &str, worker: Option<&FleetTaskWorkerProfile>) -> Result<()> {
392 let Some(worker) = worker else {
393 return Ok(());
394 };
395 validate_worker_selector(
396 task_id,
397 "worker.agent_profile",
398 worker.agent_profile.as_deref(),
399 )?;
400 validate_worker_token(task_id, "worker.loadout", worker.loadout.as_deref())?;
401 validate_worker_token(task_id, "worker.model_class", worker.model_class.as_deref())?;
402 validate_worker_model(task_id, worker.model.as_deref())?;
403 Ok(())
404 }
405
406 fn validate_worker_selector(task_id: &str, field: &str, value: Option<&str>) -> Result<()> {
407 let Some(value) = value else {
408 return Ok(());
409 };
410 let trimmed = value.trim();
411 if trimmed.is_empty() {
412 bail!("fleet task {task_id} {field} cannot be empty");
413 }
414 if trimmed != value || value.len() > MAX_FLEET_NAME_BYTES || value.chars().any(char::is_control)
415 {
416 bail!(
417 "fleet task {task_id} {field} must be one printable selector no longer than {MAX_FLEET_NAME_BYTES} bytes"
418 );
419 }
420 Ok(())
421 }
422
423 fn validate_worker_token(task_id: &str, field: &str, value: Option<&str>) -> Result<()> {
424 let Some(value) = value else {
425 return Ok(());
426 };
427 let trimmed = value.trim();
428 if trimmed.is_empty() {
429 bail!("fleet task {task_id} {field} cannot be empty");
430 }
431 if trimmed != value || !trimmed.chars().all(is_worker_token_char) {
432 bail!(
433 "fleet task {task_id} {field} must be a simple token, not a path or provider/model id"
434 );
435 }
436 Ok(())
437 }
438
439 fn is_worker_token_char(ch: char) -> bool {
440 ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.')
441 }
442
443 fn validate_worker_model(task_id: &str, value: Option<&str>) -> Result<()> {
444 let Some(value) = value else {
445 return Ok(());
446 };
447 let trimmed = value.trim();
448 if trimmed.is_empty() {
449 bail!("fleet task {task_id} worker.model cannot be empty");
450 }
451 if trimmed != value
452 || !trimmed
453 .chars()
454 .all(|ch| ch.is_ascii_graphic() && !matches!(ch, '=' | '\'' | '"'))
455 {
456 bail!(
457 "fleet task {task_id} worker.model must be a visible model id without whitespace or secrets"
458 );
459 }
460 Ok(())
461 }
462
463 #[allow(clippy::too_many_arguments)]
464 pub fn write_fleet_artifact_ref(
465 workspace: &Path,
466 run_id: &FleetRunId,
467 task_id: &str,
468 worker_id: &str,
469 kind: FleetArtifactKind,
470 filename: &str,
471 contents: &[u8],
472 mime_type: Option<&str>,
473 ) -> Result<FleetArtifactRef> {
474 let rel_path = PathBuf::from(".codewhale")
475 .join("fleet")
476 .join(safe_path_segment(&run_id.0))
477 .join(safe_path_segment(task_id))
478 .join(safe_path_segment(worker_id))
479 .join(safe_path_segment(filename));
480 super::artifacts::write(workspace, &rel_path, contents)?;
481 Ok(FleetArtifactRef {
482 kind,
483 path: rel_path,
484 checksum: Some(format!("sha256:{}", crate::hashing::sha256_hex(contents))),
485 mime_type: mime_type.map(str::to_string),
486 size_bytes: Some(contents.len() as u64),
487 })
488 }
489
490 pub fn verify_task_result(
491 workspace: &Path,
492 task: &FleetTaskSpec,
493 input: &FleetTaskVerificationInput,
494 ) -> FleetTaskVerification {
495 match &task.scorer {
496 Some(FleetScorerSpec::ExitCode) => verify_exit_code(input.exit_code),
497 Some(FleetScorerSpec::FileExists { path }) => verify_file_exists(workspace, path),
498 Some(FleetScorerSpec::RegexMatch { path, pattern }) => {
499 verify_regex_match(workspace, path, pattern)
500 }
501 Some(FleetScorerSpec::JsonPath { path, expression }) => {
502 verify_json_path(workspace, path, expression)
503 }
504 Some(FleetScorerSpec::Command { command, .. }) => partial(
505 format!("external scorer command configured: {command}"),
506 "run the configured scorer command to finalize this receipt",
507 ),
508 Some(FleetScorerSpec::CodeWhaleVerifierPrompt { .. }) => partial(
509 "Codewhale verifier prompt configured",
510 "run a verifier prompt pass to finalize this receipt",
511 ),
512 Some(FleetScorerSpec::Manual) => partial(
513 "manual scorer configured",
514 "manual verification is required to finalize this receipt",
515 ),
516 None if !has_verifiable_artifact(input) => match input
517 .final_answer
518 .as_ref()
519 .filter(|answer| !answer.excerpt.trim().is_empty())
520 {
521 Some(answer) => partial(
522 "no scorer configured; worker produced a summary deliverable",
523 answer.receipt_note(),
524 ),
525 None => partial(
526 "no scorer configured and no verifiable artifacts recorded",
527 "worker exited successfully but produced no verifiable output",
528 ),
529 },
530 None => partial(
531 "no scorer configured",
532 "task has artifacts but no deterministic scorer",
533 ),
534 }
535 }
536
537 pub fn prepare_verification_receipt(
538 workspace: &Path,
539 input: &FleetTaskVerificationInput,
540 verification: FleetTaskVerification,
541 ) -> Result<FleetReceipt> {
542 let evidence = json!({
543 "run_id": input.run_id.0.clone(),
544 "task_id": input.task_id.clone(),
545 "worker_id": input.worker_id.clone(),
546 "attempt": input.attempt,
547 "result": verification.result.clone(),
548 "failure_kind": verification.failure_kind.clone(),
549 "score": verification.score.clone(),
550 "evidence": verification.evidence.clone(),
551 "artifacts": input.artifacts.clone(),
552 });
553 let bytes =
554 serde_json::to_vec_pretty(&evidence).context("serializing fleet receipt evidence")?;
555 // Content-address the evidence as well as namespacing it by attempt. A
556 // stale verifier may finish after a retry has started; it is allowed to
557 // leave an orphaned evidence file, but it must never overwrite the file a
558 // winning attempt's durable receipt references.
559 let evidence_hash = crate::hashing::sha256_hex(&bytes);
560 let filename = format!(
561 "verification-receipt-attempt-{:010}-{}.json",
562 input.attempt, evidence_hash
563 );
564 let receipt_artifact = write_fleet_artifact_ref(
565 workspace,
566 &input.run_id,
567 &input.task_id,
568 &input.worker_id,
569 FleetArtifactKind::Receipt,
570 &filename,
571 &bytes,
572 Some("application/json"),
573 )?;
574 let mut artifacts = input.artifacts.clone();
575 artifacts.push(receipt_artifact);
576 let receipt = FleetReceipt {
577 run_id: input.run_id.clone(),
578 task_id: input.task_id.clone(),
579 worker_id: input.worker_id.clone(),
580 attempt: Some(input.attempt),
581 terminal_seq: None,
582 completed_at: timestamp(),
583 result: verification.result,
584 failure_kind: verification.failure_kind,
585 artifacts,
586 score: Some(verification.score),
587 resolved_route: input.resolved_route.clone(),
588 saved_session_id: input.saved_session_id.clone(),
589 effective_permissions: input.effective_permissions.clone(),
590 };
591 Ok(receipt)
592 }
593
594 pub fn record_verification_receipt(
595 ledger: &FleetLedger,
596 workspace: &Path,
597 input: &FleetTaskVerificationInput,
598 verification: FleetTaskVerification,
599 ) -> Result<FleetReceipt> {
600 let receipt = prepare_verification_receipt(workspace, input, verification)?;
601 ledger.record_receipt(receipt.clone())?;
602 Ok(receipt)
603 }
604
605 fn validate_tags(task_id: &str, tags: &[String]) -> Result<()> {
606 let mut seen = BTreeSet::new();
607 for tag in tags {
608 if tag.trim().is_empty() {
609 bail!("fleet task {task_id} tag cannot be empty");
610 }
611 if !seen.insert(tag) {
612 bail!("fleet task {task_id} has duplicate tag {tag}");
613 }
614 }
615 Ok(())
616 }
617
618 fn validate_workspace_requirements(task: &FleetTaskSpec) -> Result<()> {
619 let Some(workspace) = &task.workspace else {
620 return Ok(());
621 };
622 let env = workspace.environment.as_ref();
623 for name in env
624 .into_iter()
625 .flat_map(|env| env.required.iter().chain(env.allowlist.iter()))
626 {
627 if name.trim().is_empty() {
628 bail!(
629 "fleet task {} environment variable name cannot be empty",
630 task.id
631 );
632 }
633 }
634 Ok(())
635 }
636
637 fn verify_exit_code(exit_code: Option<i32>) -> FleetTaskVerification {
638 match exit_code {
639 Some(0) => pass("exit_code=0"),
640 Some(code) => fail(
641 FleetTaskFailureKind::Task,
642 0.0,
643 format!("exit_code={code}"),
644 "worker task exited unsuccessfully",
645 ),
646 None => fail(
647 FleetTaskFailureKind::Transport,
648 0.0,
649 "missing exit code",
650 "worker transport did not report a process result",
651 ),
652 }
653 }
654
655 fn verify_file_exists(workspace: &Path, path: &Path) -> FleetTaskVerification {
656 let abs_path = resolve_workspace_path(workspace, path);
657 if abs_path.is_file() {
658 pass(format!("file exists: {}", path.display()))
659 } else {
660 fail(
661 FleetTaskFailureKind::Task,
662 0.0,
663 format!("missing file: {}", path.display()),
664 "expected artifact file was not produced",
665 )
666 }
667 }
668
669 fn verify_regex_match(workspace: &Path, path: &Path, pattern: &str) -> FleetTaskVerification {
670 let regex = match Regex::new(pattern) {
671 Ok(regex) => regex,
672 Err(err) => {
673 return fail(
674 FleetTaskFailureKind::Verifier,
675 0.0,
676 format!("invalid regex: {err}"),
677 "regex scorer could not be compiled",
678 );
679 }
680 };
681 let contents = match read_bounded_to_string(workspace, path) {
682 Ok(contents) => contents,
683 Err(err) => {
684 return fail(
685 err.failure_kind,
686 0.0,
687 err.evidence,
688 "regex scorer could not read bounded evidence",
689 );
690 }
691 };
692 if regex.is_match(&contents) {
693 pass(format!("regex matched {}: {pattern}", path.display()))
694 } else {
695 fail(
696 FleetTaskFailureKind::Task,
697 0.0,
698 format!("regex did not match {}: {pattern}", path.display()),
699 "worker output did not satisfy the regex scorer",
700 )
701 }
702 }
703
704 fn verify_json_path(workspace: &Path, path: &Path, expression: &str) -> FleetTaskVerification {
705 let Some(segments) = json_path_segments(expression) else {
706 return fail(
707 FleetTaskFailureKind::Verifier,
708 0.0,
709 format!("unsupported JSON path expression: {expression}"),
710 "json_path scorer supports $.field or .field paths",
711 );
712 };
713 let contents = match read_bounded_to_string(workspace, path) {
714 Ok(contents) => contents,
715 Err(err) => {
716 return fail(
717 err.failure_kind,
718 0.0,
719 err.evidence,
720 "json_path scorer could not read bounded evidence",
721 );
722 }
723 };
724 let value: Value = match serde_json::from_str(&contents) {
725 Ok(value) => value,
726 Err(err) => {
727 return fail(
728 FleetTaskFailureKind::Task,
729 0.0,
730 format!("invalid JSON in {}: {err}", path.display()),
731 "worker artifact was not valid JSON",
732 );
733 }
734 };
735 match json_path_lookup(&value, &segments) {
736 Some(found) if json_truthy(found) => pass(format!(
737 "json_path matched {}: {expression}",
738 path.display()
739 )),
740 _ => fail(
741 FleetTaskFailureKind::Task,
742 0.0,
743 format!(
744 "json_path missing or false in {}: {expression}",
745 path.display()
746 ),
747 "worker JSON artifact did not satisfy the scorer",
748 ),
749 }
750 }
751
752 fn pass(evidence: impl Into<String>) -> FleetTaskVerification {
753 let evidence = evidence.into();
754 FleetTaskVerification {
755 result: FleetTaskResult::Pass,
756 failure_kind: None,
757 score: FleetScore {
758 value: 1.0,
759 max: Some(1.0),
760 notes: Some(evidence.clone()),
761 },
762 evidence: vec![evidence],
763 }
764 }
765
766 fn partial(evidence: impl Into<String>, notes: impl Into<String>) -> FleetTaskVerification {
767 let evidence = evidence.into();
768 let notes = notes.into();
769 FleetTaskVerification {
770 result: FleetTaskResult::Partial,
771 failure_kind: None,
772 score: FleetScore {
773 value: 0.5,
774 max: Some(1.0),
775 notes: Some(notes),
776 },
777 evidence: vec![evidence],
778 }
779 }
780
781 fn fail(
782 failure_kind: FleetTaskFailureKind,
783 value: f64,
784 evidence: impl Into<String>,
785 notes: impl Into<String>,
786 ) -> FleetTaskVerification {
787 let evidence = evidence.into();
788 FleetTaskVerification {
789 result: FleetTaskResult::Fail,
790 failure_kind: Some(failure_kind),
791 score: FleetScore {
792 value,
793 max: Some(1.0),
794 notes: Some(notes.into()),
795 },
796 evidence: vec![evidence],
797 }
798 }
799
800 fn has_verifiable_artifact(input: &FleetTaskVerificationInput) -> bool {
801 input.artifacts.iter().any(|artifact| {
802 !matches!(
803 artifact.kind,
804 FleetArtifactKind::Log | FleetArtifactKind::Receipt
805 )
806 })
807 }
808
809 #[derive(Debug)]
810 struct EvidenceReadError {
811 failure_kind: FleetTaskFailureKind,
812 evidence: String,
813 }
814
815 fn read_bounded_to_string(
816 workspace: &Path,
817 path: &Path,
818 ) -> std::result::Result<String, EvidenceReadError> {
819 let abs_path = resolve_workspace_path(workspace, path);
820 let metadata = std::fs::metadata(&abs_path).map_err(|err| EvidenceReadError {
821 failure_kind: if err.kind() == std::io::ErrorKind::NotFound {
822 FleetTaskFailureKind::Task
823 } else {
824 FleetTaskFailureKind::Verifier
825 },
826 evidence: format!("cannot read {}: {err}", path.display()),
827 })?;
828 if metadata.len() > MAX_SCORER_READ_BYTES {
829 return Err(EvidenceReadError {
830 failure_kind: FleetTaskFailureKind::Verifier,
831 evidence: format!(
832 "refusing to read oversized evidence {}: {} bytes",
833 path.display(),
834 metadata.len()
835 ),
836 });
837 }
838 std::fs::read_to_string(&abs_path).map_err(|err| EvidenceReadError {
839 failure_kind: FleetTaskFailureKind::Verifier,
840 evidence: format!("cannot decode {} as UTF-8: {err}", path.display()),
841 })
842 }
843
844 fn resolve_workspace_path(workspace: &Path, path: &Path) -> PathBuf {
845 if path.is_absolute() {
846 path.to_path_buf()
847 } else {
848 workspace.join(path)
849 }
850 }
851
852 fn json_path_segments(expression: &str) -> Option<Vec<&str>> {
853 let trimmed = expression.trim();
854 let path = trimmed
855 .strip_prefix("$.")
856 .or_else(|| trimmed.strip_prefix('.'))?;
857 if path.is_empty() {
858 return None;
859 }
860 let segments: Vec<_> = path.split('.').collect();
861 if segments.iter().any(|segment| segment.is_empty()) {
862 return None;
863 }
864 Some(segments)
865 }
866
867 fn json_path_lookup<'a>(value: &'a Value, segments: &[&str]) -> Option<&'a Value> {
868 let mut current = value;
869 for segment in segments {
870 current = current.as_object()?.get(*segment)?;
871 }
872 Some(current)
873 }
874
875 fn json_truthy(value: &Value) -> bool {
876 !matches!(value, Value::Null | Value::Bool(false))
877 }
878
879 fn timestamp() -> String {
880 Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true)
881 }
882
883 fn safe_path_segment(value: &str) -> String {
884 value
885 .chars()
886 .map(|ch| {
887 if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') {
888 ch
889 } else {
890 '_'
891 }
892 })
893 .collect()
894 }
895
896 #[cfg(test)]
897 mod tests {
898 use super::*;
899 use serde_json::json;
900 use tempfile::TempDir;
901
902 fn task(id: &str, scorer: Option<FleetScorerSpec>) -> FleetTaskSpec {
903 FleetTaskSpec {
904 id: id.to_string(),
905 name: id.to_string(),
906 description: None,
907 objective: Some(format!("Verify {id}")),
908 instructions: format!("do {id}"),
909 worker: Some(FleetTaskWorkerProfile {
910 agent_profile: None,
911 role: Some("reviewer".to_string()),
912 loadout: None,
913 model_class: None,
914 model: None,
915 tool_profile: Some("read-only".to_string()),
916 tools: vec!["git".to_string()],
917 capabilities: vec!["rust".to_string()],
918 }),
919 workspace: Some(FleetWorkspaceRequirements {
920 root: Some(PathBuf::from(".")),
921 required_files: vec![PathBuf::from("Cargo.toml")],
922 writable_paths: vec![PathBuf::from(".codewhale/fleet")],
923 environment: Some(FleetEnvironmentRequirements {
924 required: vec!["PATH".to_string()],
925 allowlist: vec!["RUST_LOG".to_string()],
926 }),
927 }),
928 input_files: vec![PathBuf::from("Cargo.toml")],
929 context: vec!["fleet verifier test".to_string()],
930 budget: Some(FleetTaskBudget {
931 max_tokens: Some(4000),
932 max_steps: None,
933 max_tool_calls: Some(12),
934 max_seconds: Some(120),
935 }),
936 expected_artifacts: vec![FleetArtifactKind::Log, FleetArtifactKind::Receipt],
937 scorer,
938 retry_policy: Some(FleetRetryPolicy::default()),
939 alert_policy: None,
940 timeout_seconds: Some(120),
941 tags: vec!["review".to_string()],
942 metadata: BTreeMap::new(),
943 }
944 }
945
946 #[test]
947 fn fleet_task_spec_document_parses_multi_task_verified_shape() {
948 let tmp = TempDir::new().unwrap();
949 let path = tmp.path().join("fleet-tasks.json");
950 let doc = json!({
951 "name": "release triage",
952 "labels": {"milestone": "v0.8.60"},
953 "tasks": [
954 task("release-notes", Some(FleetScorerSpec::ExitCode)),
955 task("risk-review", Some(FleetScorerSpec::Manual))
956 ]
957 });
958 std::fs::write(&path, serde_json::to_string_pretty(&doc).unwrap()).unwrap();
959
960 let parsed = load_task_spec_document(&path).unwrap();
961
962 assert_eq!(parsed.name.as_deref(), Some("release triage"));
963 assert_eq!(parsed.tasks.len(), 2);
964 assert_eq!(
965 parsed.tasks[0].objective.as_deref(),
966 Some("Verify release-notes")
967 );
968 assert_eq!(
969 parsed.tasks[0].worker.as_ref().unwrap().role.as_deref(),
970 Some("reviewer")
971 );
972 assert_eq!(parsed.tasks[1].tags, vec!["review"]);
973 }
974
975 #[test]
976 fn fleet_task_spec_document_parses_worker_profile_loadout_intent() {
977 let tmp = TempDir::new().unwrap();
978 let path = tmp.path().join("fleet-profile-task.json");
979 let doc = json!({
980 "name": "profile loadout smoke",
981 "tasks": [{
982 "id": "review",
983 "name": "review",
984 "instructions": "review the patch",
985 "worker": {
986 "profile": "DeepSeek V4 Flash",
987 "role": "reviewer",
988 "loadout": "auto",
989 "model_class": "balanced",
990 "model": "deepseek-v4-pro",
991 "tool_profile": "read-only",
992 "tools": ["read_file", "grep_files"],
993 "capabilities": ["rust"]
994 }
995 }]
996 });
997 std::fs::write(&path, serde_json::to_string_pretty(&doc).unwrap()).unwrap();
998
999 let parsed = load_task_spec_document(&path).unwrap();
1000 let worker = parsed.tasks[0].worker.as_ref().unwrap();
1001
1002 assert_eq!(worker.agent_profile.as_deref(), Some("DeepSeek V4 Flash"));
1003 assert_eq!(worker.role.as_deref(), Some("reviewer"));
1004 assert_eq!(worker.loadout.as_deref(), Some("auto"));
1005 assert_eq!(worker.model_class.as_deref(), Some("balanced"));
1006 assert_eq!(worker.model.as_deref(), Some("deepseek-v4-pro"));
1007 assert_eq!(worker.tool_profile.as_deref(), Some("read-only"));
1008 }
1009
1010 #[test]
1011 fn fleet_task_spec_rejects_unsafe_worker_profile_intent_tokens() {
1012 let tmp = TempDir::new().unwrap();
1013 let path = tmp.path().join("unsafe-profile-task.json");
1014 let doc = json!({
1015 "tasks": [{
1016 "id": "review",
1017 "name": "review",
1018 "instructions": "review the patch",
1019 "worker": {
1020 "profile": "reviewer\n../../secrets",
1021 "loadout": "openrouter/deepseek",
1022 "model_class": "",
1023 "model": "deepseek/deepseek-v4-pro"
1024 }
1025 }]
1026 });
1027 std::fs::write(&path, serde_json::to_string_pretty(&doc).unwrap()).unwrap();
1028
1029 let err = load_task_spec_document(&path).unwrap_err().to_string();
1030
1031 assert!(
1032 err.contains("worker.agent_profile must be one printable selector"),
1033 "unexpected error: {err}"
1034 );
1035 }
1036
1037 #[test]
1038 fn fleet_task_spec_rejects_secret_like_worker_model() {
1039 let tmp = TempDir::new().unwrap();
1040 let path = tmp.path().join("unsafe-worker-model.json");
1041 let doc = json!({
1042 "tasks": [{
1043 "id": "review",
1044 "name": "review",
1045 "instructions": "review the patch",
1046 "worker": {
1047 "model": "deepseek-v4-pro api_key=secret"
1048 }
1049 }]
1050 });
1051 std::fs::write(&path, serde_json::to_string_pretty(&doc).unwrap()).unwrap();
1052
1053 let err = load_task_spec_document(&path).unwrap_err().to_string();
1054
1055 assert!(
1056 err.contains("worker.model must be a visible model id"),
1057 "unexpected error: {err}"
1058 );
1059 }
1060
1061 #[test]
1062 fn fleet_task_spec_rejects_unbounded_or_multiline_task_and_worker_identities() {
1063 let tmp = TempDir::new().unwrap();
1064 let path = tmp.path().join("unsafe-identities.json");
1065 let doc = json!({
1066 "workers": [{
1067 "id": "worker\r\nforged",
1068 "name": "forged worker",
1069 "host": {"kind": "local"}
1070 }],
1071 "tasks": [{
1072 "id": "review",
1073 "name": "review",
1074 "instructions": "review the patch"
1075 }]
1076 });
1077 std::fs::write(&path, serde_json::to_string_pretty(&doc).unwrap()).unwrap();
1078
1079 let err = load_task_spec_document(&path).unwrap_err().to_string();
1080 assert!(
1081 err.contains("worker id must be a simple ASCII token"),
1082 "unexpected error: {err}"
1083 );
1084
1085 let mut doc = task("review", None);
1086 doc.id = "a".repeat(MAX_FLEET_ID_BYTES + 1);
1087 let err = validate_task_spec_document(&FleetTaskSpecDocument {
1088 name: None,
1089 labels: BTreeMap::new(),
1090 security_policy: None,
1091 workers: Vec::new(),
1092 tasks: vec![doc],
1093 usage_ceiling: None,
1094 })
1095 .unwrap_err()
1096 .to_string();
1097 assert!(
1098 err.contains("task id must be a simple ASCII token"),
1099 "unexpected error: {err}"
1100 );
1101 }
1102
1103 #[test]
1104 fn fleet_task_spec_rejects_runtime_authority_fields_for_new_runs() {
1105 let mut security_doc = FleetTaskSpecDocument {
1106 name: None,
1107 labels: BTreeMap::new(),
1108 security_policy: Some(FleetSecurityPolicy::default()),
1109 workers: Vec::new(),
1110 tasks: vec![task("review", None)],
1111 usage_ceiling: None,
1112 };
1113 let error = validate_task_spec_document(&security_doc)
1114 .unwrap_err()
1115 .to_string();
1116 assert!(
1117 error.contains("security_policy is a legacy compatibility field"),
1118 "unexpected error: {error}"
1119 );
1120
1121 security_doc.security_policy = None;
1122 security_doc.workers.push(FleetWorkerSpec {
1123 id: "worker-1".to_string(),
1124 name: "Worker 1".to_string(),
1125 host: FleetHostSpec::Local,
1126 trust_level: Some(FleetTrustLevel::Local),
1127 labels: BTreeMap::new(),
1128 capabilities: vec!["local".to_string()],
1129 max_concurrent_tasks: Some(1),
1130 });
1131 let error = validate_task_spec_document(&security_doc)
1132 .unwrap_err()
1133 .to_string();
1134 assert!(
1135 error.contains("trust_level is a legacy compatibility field"),
1136 "unexpected error: {error}"
1137 );
1138 }
1139
1140 #[cfg(unix)]
1141 #[test]
1142 fn fleet_artifact_publication_rejects_symlinked_workspace_paths() {
1143 use std::os::unix::fs::symlink;
1144 let workspace = TempDir::new().unwrap();
1145 let outside = TempDir::new().unwrap();
1146 std::fs::create_dir_all(workspace.path().join(".codewhale")).unwrap();
1147 symlink(outside.path(), workspace.path().join(".codewhale/fleet")).unwrap();
1148 let result = write_fleet_artifact_ref(
1149 workspace.path(),
1150 &FleetRunId::from("run-1"),
1151 "task-a",
1152 "worker-1",
1153 FleetArtifactKind::Receipt,
1154 "receipt.json",
1155 b"synthetic receipt",
1156 Some("application/json"),
1157 );
1158 assert!(
1159 result.is_err(),
1160 "publication must reject a symlinked parent"
1161 );
1162 assert!(!outside.path().join("run-1").exists());
1163 }
1164
1165 #[test]
1166 fn fleet_task_spec_artifact_refs_are_bounded_paths() {
1167 let tmp = TempDir::new().unwrap();
1168 let artifact = write_fleet_artifact_ref(
1169 tmp.path(),
1170 &FleetRunId::from("run-1"),
1171 "task-a",
1172 "worker-1",
1173 FleetArtifactKind::Log,
1174 "worker.log",
1175 b"this is artifact content",
1176 Some("text/plain"),
1177 )
1178 .unwrap();
1179
1180 let json = serde_json::to_string(&artifact).unwrap();
1181 assert!(!json.contains("this is artifact content"));
1182 assert!(json.contains("worker.log"));
1183 assert_eq!(artifact.size_bytes, Some(24));
1184 assert!(artifact.checksum.as_deref().unwrap().starts_with("sha256:"));
1185 assert!(tmp.path().join(&artifact.path).exists());
1186 }
1187
1188 #[test]
1189 fn fleet_task_spec_scorers_record_pass_fail_partial_evidence() {
1190 let tmp = TempDir::new().unwrap();
1191 std::fs::write(tmp.path().join("result.txt"), "status=ok\n").unwrap();
1192 std::fs::write(tmp.path().join("result.json"), r#"{"status":"ok"}"#).unwrap();
1193 let input = FleetTaskVerificationInput {
1194 run_id: FleetRunId::from("run-1"),
1195 task_id: "task-a".to_string(),
1196 worker_id: "worker-1".to_string(),
1197 attempt: 1,
1198 exit_code: Some(0),
1199 artifacts: vec![],
1200 final_answer: None,
1201 saved_session_id: None,
1202 resolved_route: None,
1203 effective_permissions: None,
1204 };
1205
1206 let pass = verify_task_result(
1207 tmp.path(),
1208 &task("exit", Some(FleetScorerSpec::ExitCode)),
1209 &input,
1210 );
1211 assert_eq!(pass.result, FleetTaskResult::Pass);
1212 assert_eq!(pass.failure_kind, None);
1213
1214 let regex = verify_task_result(
1215 tmp.path(),
1216 &task(
1217 "regex",
1218 Some(FleetScorerSpec::RegexMatch {
1219 path: PathBuf::from("result.txt"),
1220 pattern: "status=ok".to_string(),
1221 }),
1222 ),
1223 &input,
1224 );
1225 assert_eq!(regex.result, FleetTaskResult::Pass);
1226
1227 let json_path = verify_task_result(
1228 tmp.path(),
1229 &task(
1230 "json",
1231 Some(FleetScorerSpec::JsonPath {
1232 path: PathBuf::from("result.json"),
1233 expression: "$.status".to_string(),
1234 }),
1235 ),
1236 &input,
1237 );
1238 assert_eq!(json_path.result, FleetTaskResult::Pass);
1239
1240 let manual = verify_task_result(
1241 tmp.path(),
1242 &task("manual", Some(FleetScorerSpec::Manual)),
1243 &input,
1244 );
1245 assert_eq!(manual.result, FleetTaskResult::Partial);
1246
1247 let no_scorer_empty = verify_task_result(tmp.path(), &task("unscored", None), &input);
1248 assert_eq!(no_scorer_empty.result, FleetTaskResult::Partial);
1249 assert!(
1250 no_scorer_empty
1251 .score
1252 .notes
1253 .as_deref()
1254 .unwrap_or_default()
1255 .contains("no verifiable output")
1256 );
1257
1258 let failed = verify_task_result(
1259 tmp.path(),
1260 &task(
1261 "missing",
1262 Some(FleetScorerSpec::FileExists {
1263 path: PathBuf::from("missing.txt"),
1264 }),
1265 ),
1266 &input,
1267 );
1268 assert_eq!(failed.result, FleetTaskResult::Fail);
1269 assert_eq!(failed.failure_kind, Some(FleetTaskFailureKind::Task));
1270
1271 let verifier_failed = verify_task_result(
1272 tmp.path(),
1273 &task(
1274 "bad-regex",
1275 Some(FleetScorerSpec::RegexMatch {
1276 path: PathBuf::from("result.txt"),
1277 pattern: "[".to_string(),
1278 }),
1279 ),
1280 &input,
1281 );
1282 assert_eq!(verifier_failed.result, FleetTaskResult::Fail);
1283 assert_eq!(
1284 verifier_failed.failure_kind,
1285 Some(FleetTaskFailureKind::Verifier)
1286 );
1287 }
1288
1289 #[test]
1290 fn unscored_worker_surfaces_summary_deliverable_instead_of_no_output() {
1291 let tmp = TempDir::new().unwrap();
1292 let input = FleetTaskVerificationInput {
1293 run_id: FleetRunId::from("run-1"),
1294 task_id: "task-a".to_string(),
1295 worker_id: "worker-1".to_string(),
1296 attempt: 1,
1297 exit_code: Some(0),
1298 artifacts: vec![],
1299 final_answer: Some(FleetWorkerFinalAnswer {
1300 excerpt: "The Changelog review is complete".to_string(),
1301 chars: 32,
1302 }),
1303 saved_session_id: None,
1304 resolved_route: None,
1305 effective_permissions: None,
1306 };
1307 let verification = verify_task_result(tmp.path(), &task("unscored", None), &input);
1308 assert_eq!(verification.result, FleetTaskResult::Partial);
1309 let notes = verification
1310 .score
1311 .notes
1312 .as_deref()
1313 .unwrap_or_default()
1314 .to_string();
1315 assert!(
1316 notes.contains("worker produced 32 characters of deliverable"),
1317 "unexpected notes: {notes}"
1318 );
1319 assert!(notes.contains("Changelog review is complete"));
1320 assert!(!notes.contains("no verifiable output"));
1321 }
1322
1323 #[test]
1324 fn fleet_task_spec_receipt_records_artifacts_scores_and_failure_kind() {
1325 let tmp = TempDir::new().unwrap();
1326 let ledger = FleetLedger::open(tmp.path()).unwrap();
1327 let log = write_fleet_artifact_ref(
1328 tmp.path(),
1329 &FleetRunId::from("run-1"),
1330 "task-a",
1331 "worker-1",
1332 FleetArtifactKind::Log,
1333 "worker.log",
1334 b"exit_code=1",
1335 Some("text/plain"),
1336 )
1337 .unwrap();
1338 let input = FleetTaskVerificationInput {
1339 run_id: FleetRunId::from("run-1"),
1340 task_id: "task-a".to_string(),
1341 worker_id: "worker-1".to_string(),
1342 attempt: 3,
1343 exit_code: Some(1),
1344 artifacts: vec![log],
1345 final_answer: None,
1346 saved_session_id: None,
1347 resolved_route: None,
1348 effective_permissions: Some(FleetEffectivePermissions {
1349 write: false,
1350 network: false,
1351 shell: "read_only".to_string(),
1352 tool_scope: "explicit".to_string(),
1353 tools: vec!["read_file".to_string()],
1354 background: true,
1355 max_spawn_depth: 0,
1356 profile_id: None,
1357 profile_origin: None,
1358 source: "worker_runtime_profile".to_string(),
1359 }),
1360 };
1361 let verification = verify_task_result(
1362 tmp.path(),
1363 &task("task-a", Some(FleetScorerSpec::ExitCode)),
1364 &input,
1365 );
1366
1367 let receipt =
1368 record_verification_receipt(&ledger, tmp.path(), &input, verification).unwrap();
1369
1370 assert_eq!(receipt.result, FleetTaskResult::Fail);
1371 assert_eq!(receipt.failure_kind, Some(FleetTaskFailureKind::Task));
1372 assert_eq!(receipt.attempt, Some(3));
1373 assert_eq!(receipt.terminal_seq, None);
1374 assert_eq!(receipt.effective_permissions, input.effective_permissions);
1375 assert_eq!(receipt.artifacts.len(), 2);
1376 assert!(matches!(
1377 receipt.artifacts.last().unwrap().kind,
1378 FleetArtifactKind::Receipt
1379 ));
1380 assert!(
1381 receipt
1382 .artifacts
1383 .last()
1384 .unwrap()
1385 .path
1386 .to_string_lossy()
1387 .contains("verification-receipt-attempt-0000000003-")
1388 );
1389 let state = ledger.rebuild_state().unwrap();
1390 assert_eq!(
1391 state.receipts["run-1:task-a"].failure_kind,
1392 Some(FleetTaskFailureKind::Task)
1393 );
1394 }
1395
1396 #[test]
1397 fn verification_evidence_is_attempt_and_content_addressed() {
1398 let tmp = TempDir::new().unwrap();
1399 let mut input = FleetTaskVerificationInput {
1400 run_id: FleetRunId::from("run-1"),
1401 task_id: "task-a".to_string(),
1402 worker_id: "worker-1".to_string(),
1403 attempt: 1,
1404 exit_code: Some(1),
1405 artifacts: Vec::new(),
1406 final_answer: None,
1407 saved_session_id: None,
1408 resolved_route: None,
1409 effective_permissions: None,
1410 };
1411 let scorer = task("task-a", Some(FleetScorerSpec::ExitCode));
1412 let stale_verification = verify_task_result(tmp.path(), &scorer, &input);
1413 let stale = prepare_verification_receipt(tmp.path(), &input, stale_verification).unwrap();
1414
1415 input.attempt = 2;
1416 input.exit_code = Some(0);
1417 let winning_verification = verify_task_result(tmp.path(), &scorer, &input);
1418 let winning =
1419 prepare_verification_receipt(tmp.path(), &input, winning_verification).unwrap();
1420
1421 let stale_path = &stale.artifacts.last().unwrap().path;
1422 let winning_path = &winning.artifacts.last().unwrap().path;
1423 assert_ne!(stale_path, winning_path);
1424 assert!(stale_path.to_string_lossy().contains("attempt-0000000001-"));
1425 assert!(
1426 winning_path
1427 .to_string_lossy()
1428 .contains("attempt-0000000002-")
1429 );
1430 assert!(tmp.path().join(stale_path).is_file());
1431 assert!(tmp.path().join(winning_path).is_file());
1432 assert_eq!(stale.result, FleetTaskResult::Fail);
1433 assert_eq!(winning.result, FleetTaskResult::Pass);
1434 }
1435
1436 fn load_error(file_name: &str, body: &str) -> String {
1437 let tmp = TempDir::new().unwrap();
1438 let path = tmp.path().join(file_name);
1439 std::fs::write(&path, body).unwrap();
1440 let err = load_task_spec_document(&path).expect_err("spec should be rejected");
1441 format!("{err:#}")
1442 }
1443
1444 #[test]
1445 fn fleet_task_spec_document_shape_error_names_missing_field_and_task() {
1446 let err = load_error(
1447 "doc.json",
1448 r#"{"name": "n", "tasks": [
1449 {"id": "ok", "name": "ok", "instructions": "do it"},
1450 {"id": "review", "name": "review"}
1451 ]}"#,
1452 );
1453 assert!(!err.contains("untagged enum"), "{err}");
1454 assert!(err.contains("JSON spec document"), "{err}");
1455 assert!(err.contains("missing field `instructions`"), "{err}");
1456 assert!(err.contains(r#"tasks[1] (id "review")"#), "{err}");
1457 }
1458
1459 #[test]
1460 fn fleet_task_spec_task_array_shape_error_names_missing_field() {
1461 let err = load_error("tasks.json", r#"[{"id": "a", "instructions": "do it"}]"#);
1462 assert!(!err.contains("untagged enum"), "{err}");
1463 assert!(err.contains("JSON task array"), "{err}");
1464 assert!(err.contains("missing field `name`"), "{err}");
1465 assert!(err.contains(r#"[0] (id "a")"#), "{err}");
1466 }
1467
1468 #[test]
1469 fn fleet_task_spec_single_task_shape_error_names_missing_field() {
1470 let err = load_error("one.json", r#"{"id": "a", "name": "a"}"#);
1471 assert!(!err.contains("untagged enum"), "{err}");
1472 assert!(err.contains("JSON single task"), "{err}");
1473 assert!(err.contains("missing field `instructions`"), "{err}");
1474 }
1475
1476 #[test]
1477 fn fleet_task_spec_toml_document_error_names_missing_field() {
1478 let err = load_error(
1479 "doc.toml",
1480 "name = \"n\"\n\n[[tasks]]\nid = \"a\"\ninstructions = \"do it\"\n",
1481 );
1482 assert!(!err.contains("untagged enum"), "{err}");
1483 assert!(err.contains("TOML spec document"), "{err}");
1484 assert!(err.contains("missing field `name`"), "{err}");
1485 assert!(err.contains(r#"tasks[0] (id "a")"#), "{err}");
1486 }
1487
1488 #[test]
1489 fn fleet_task_spec_rejects_scalar_top_level_with_shape_hint() {
1490 let err = load_error("scalar.json", "\"just a string\"");
1491 assert!(err.contains("found a string"), "{err}");
1492 assert!(err.contains("array of task objects"), "{err}");
1493 }
1494
1495 #[test]
1496 fn fleet_task_spec_single_and_array_shapes_load_with_fallback_name() {
1497 let tmp = TempDir::new().unwrap();
1498 let single = tmp.path().join("solo.json");
1499 std::fs::write(
1500 &single,
1501 r#"{"id": "a", "name": "a", "instructions": "do it"}"#,
1502 )
1503 .unwrap();
1504 let doc = load_task_spec_document(&single).unwrap();
1505 assert_eq!(doc.name.as_deref(), Some("solo"));
1506 assert_eq!(doc.tasks.len(), 1);
1507
1508 let array = tmp.path().join("pair.json");
1509 std::fs::write(
1510 &array,
1511 r#"[{"id": "a", "name": "a", "instructions": "x"},
1512 {"id": "b", "name": "b", "instructions": "y"}]"#,
1513 )
1514 .unwrap();
1515 let doc = load_task_spec_document(&array).unwrap();
1516 assert_eq!(doc.name.as_deref(), Some("pair"));
1517 assert_eq!(doc.tasks.len(), 2);
1518 }
1519
1520 #[test]
1521 fn fleet_task_spec_toml_single_task_loads_with_fallback_name() {
1522 let tmp = TempDir::new().unwrap();
1523 let path = tmp.path().join("solo.toml");
1524 std::fs::write(
1525 &path,
1526 "id = \"a\"\nname = \"a\"\ninstructions = \"do it\"\n",
1527 )
1528 .unwrap();
1529 let doc = load_task_spec_document(&path).unwrap();
1530 assert_eq!(doc.name.as_deref(), Some("solo"));
1531 assert_eq!(doc.tasks.len(), 1);
1532 }
1533
1534 fn repo_doc(relative: &str) -> PathBuf {
1535 Path::new(env!("CARGO_MANIFEST_DIR"))
1536 .join("../..")
1537 .join(relative)
1538 }
1539
1540 #[test]
1541 fn fleet_dogfood_example_spec_parses_and_validates() {
1542 let doc = load_task_spec_document(&repo_doc("docs/examples/fleet-dogfood.toml"))
1543 .expect("docs/examples/fleet-dogfood.toml must stay a valid fleet spec");
1544 assert_eq!(doc.name.as_deref(), Some("dogfood smoke"));
1545 let ids: Vec<_> = doc.tasks.iter().map(|task| task.id.as_str()).collect();
1546 assert_eq!(ids, ["cargo-check", "protocol-review"]);
1547 }
1548
1549 #[test]
1550 fn fleet_workflow_tutorial_json_spec_parses_and_validates() {
1551 let tutorial = std::fs::read_to_string(repo_doc("docs/FLEET_WORKFLOW_TUTORIAL.md"))
1552 .expect("read fleet tutorial");
1553 let start = tutorial
1554 .find("```json\n")
1555 .expect("tutorial should carry a JSON task spec")
1556 + "```json\n".len();
1557 let end = start + tutorial[start..].find("```").expect("closed JSON fence");
1558 let tmp = TempDir::new().unwrap();
1559 let path = tmp.path().join("tasks.json");
1560 std::fs::write(&path, &tutorial[start..end]).unwrap();
1561 let doc = load_task_spec_document(&path)
1562 .expect("the tutorial's tasks.json must stay a valid fleet spec");
1563 assert_eq!(doc.name.as_deref(), Some("docs readiness check"));
1564 assert_eq!(doc.tasks.len(), 2);
1565 }
1566 }
1567
1567 lines RUST