Skip to main content

codewhale_protocol/
fleet.rs

1//! Agent Fleet control-plane protocol types.
2//!
3//! These types define the durable, serializable contract between the fleet
4//! manager, workers, CLI/TUI surfaces, and the Runtime API. They are
5//! intentionally additive: existing runtime-event consumers ignore unknown
6//! fields and are unaffected by fleet extensions.
7//!
8//! See:
9//! - <https://github.com/Hmbown/CodeWhale/issues/3154> (Agent Fleet control plane)
10//! - <https://github.com/Hmbown/CodeWhale/issues/3096> (Runtime API sub-agent direction)
11
12use std::collections::BTreeMap;
13use std::path::PathBuf;
14
15use serde::{Deserialize, Deserializer, Serialize, Serializer, de};
16use serde_json::Value;
17
18use super::Status;
19
20pub const FLEET_PROTOCOL_VERSION: &str = "0.1.0";
21
22/// Globally unique identifier for a fleet run.
23#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq, Hash)]
24pub struct FleetRunId(pub String);
25
26impl From<String> for FleetRunId {
27    fn from(value: String) -> Self {
28        Self(value)
29    }
30}
31
32impl From<&str> for FleetRunId {
33    fn from(value: &str) -> Self {
34        Self(value.to_string())
35    }
36}
37
38/// Top-level fleet run handle.
39#[derive(Debug, Clone, Serialize, Deserialize)]
40pub struct FleetRun {
41    pub id: FleetRunId,
42    pub name: String,
43    pub status: FleetRunStatus,
44    /// Explicit execution target selected by the managed client.
45    ///
46    /// Older CLI-created runs predate target selection and therefore omit
47    /// this field. Runtime API creation always persists it and currently
48    /// accepts only [`FleetRuntimeTarget::ThisComputer`].
49    #[serde(default, skip_serializing_if = "Option::is_none")]
50    pub target: Option<FleetRuntimeTarget>,
51    /// Named Workflow descriptor that owns this Fleet run.
52    ///
53    /// The durable task specs below remain the executable source of truth;
54    /// this descriptor keeps the product identity and scheduling policy
55    /// inspectable without smuggling them through labels.
56    #[serde(default, skip_serializing_if = "Option::is_none")]
57    pub workflow: Option<FleetWorkflowDescriptor>,
58    /// Canonical named roles declared for the run.
59    #[serde(default, skip_serializing_if = "Vec::is_empty")]
60    pub roles: Vec<String>,
61    /// Maximum number of workers the manager may drive concurrently.
62    ///
63    /// Older ledgers omit this field; callers fall back to the persisted
64    /// worker roster when resuming those runs.
65    #[serde(default, skip_serializing_if = "Option::is_none")]
66    pub max_workers: Option<usize>,
67    #[serde(default)]
68    pub task_specs: Vec<FleetTaskSpec>,
69    #[serde(default)]
70    pub worker_specs: Vec<FleetWorkerSpec>,
71    #[serde(default)]
72    pub labels: BTreeMap<String, String>,
73    #[serde(skip_serializing_if = "Option::is_none")]
74    pub security_policy: Option<FleetSecurityPolicy>,
75    pub created_at: String,
76    #[serde(skip_serializing_if = "Option::is_none")]
77    pub updated_at: Option<String>,
78    #[serde(skip_serializing_if = "Option::is_none")]
79    pub completed_at: Option<String>,
80}
81
82/// Product-level Runtime target for a managed Fleet run.
83///
84/// The enum intentionally names unsupported targets as contract values so a
85/// client receives a precise capability refusal instead of silently falling
86/// back to local execution.
87#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
88#[serde(rename_all = "snake_case")]
89pub enum FleetRuntimeTarget {
90    ThisComputer,
91    AnotherComputer,
92    Cloud,
93}
94
95/// Scheduling shape currently executable by the durable Fleet manager.
96///
97/// Fleet tasks are independent queue entries today, so only parallel
98/// workflows are advertised. Sequence/pipeline support must not be accepted
99/// until dependencies are durable in the Fleet ledger.
100#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
101#[serde(rename_all = "snake_case")]
102pub enum FleetWorkflowKind {
103    Parallel,
104}
105
106/// Durable identity for the Workflow that coordinates a Fleet run.
107#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
108pub struct FleetWorkflowDescriptor {
109    pub id: String,
110    pub kind: FleetWorkflowKind,
111}
112
113/// One privacy-bounded durable event exposed to managed Fleet clients.
114///
115/// `cursor` is an opaque stable digest of the underlying ledger transition.
116/// Clients persist it and send it back on reconnect; they must not parse it.
117/// Worker-local sequence numbers remain available separately because they are
118/// monotonic only within one `(worker, task)` lifecycle, not across a run.
119#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
120pub struct FleetRuntimeEvent {
121    pub cursor: String,
122    pub event: String,
123    pub run_id: FleetRunId,
124    #[serde(default, skip_serializing_if = "Option::is_none")]
125    pub worker_id: Option<String>,
126    #[serde(default, skip_serializing_if = "Option::is_none")]
127    pub task_id: Option<String>,
128    #[serde(default, skip_serializing_if = "Option::is_none")]
129    pub timestamp: Option<String>,
130    #[serde(default, skip_serializing_if = "Option::is_none")]
131    pub worker_seq: Option<u64>,
132    #[serde(default)]
133    pub payload: Value,
134}
135
136/// Bounded durable replay page.
137#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
138pub struct FleetEventReplay {
139    pub run_id: FleetRunId,
140    pub events: Vec<FleetRuntimeEvent>,
141    #[serde(default)]
142    pub has_more: bool,
143    /// True when a no-cursor request returned only the newest bounded tail.
144    #[serde(default)]
145    pub history_truncated: bool,
146    #[serde(default, skip_serializing_if = "Option::is_none")]
147    pub next_cursor: Option<String>,
148}
149
150/// Lifecycle status for an entire fleet run.
151#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
152#[serde(rename_all = "snake_case")]
153pub enum FleetRunStatus {
154    Pending,
155    Queued,
156    Running,
157    Paused,
158    Completed,
159    Failed,
160    Cancelled,
161}
162
163impl Status for FleetRunStatus {
164    fn is_terminal(&self) -> bool {
165        matches!(self, Self::Completed | Self::Failed | Self::Cancelled)
166    }
167    fn is_active(&self) -> bool {
168        matches!(self, Self::Pending | Self::Queued | Self::Running)
169    }
170    fn is_paused(&self) -> bool {
171        matches!(self, Self::Paused)
172    }
173}
174
175/// Specification of a single unit of work within a run.
176#[derive(Debug, Clone, Serialize, Deserialize)]
177pub struct FleetTaskSpec {
178    pub id: String,
179    pub name: String,
180    #[serde(skip_serializing_if = "Option::is_none")]
181    pub description: Option<String>,
182    #[serde(skip_serializing_if = "Option::is_none")]
183    pub objective: Option<String>,
184    pub instructions: String,
185    #[serde(skip_serializing_if = "Option::is_none")]
186    pub worker: Option<FleetTaskWorkerProfile>,
187    #[serde(skip_serializing_if = "Option::is_none")]
188    pub workspace: Option<FleetWorkspaceRequirements>,
189    #[serde(default)]
190    #[serde(skip_serializing_if = "Vec::is_empty")]
191    pub input_files: Vec<PathBuf>,
192    #[serde(default)]
193    #[serde(skip_serializing_if = "Vec::is_empty")]
194    pub context: Vec<String>,
195    #[serde(skip_serializing_if = "Option::is_none")]
196    pub budget: Option<FleetTaskBudget>,
197    #[serde(default)]
198    #[serde(skip_serializing_if = "Vec::is_empty")]
199    pub tags: Vec<String>,
200    #[serde(default)]
201    pub expected_artifacts: Vec<FleetArtifactKind>,
202    #[serde(skip_serializing_if = "Option::is_none")]
203    pub scorer: Option<FleetScorerSpec>,
204    #[serde(skip_serializing_if = "Option::is_none")]
205    pub retry_policy: Option<FleetRetryPolicy>,
206    #[serde(skip_serializing_if = "Option::is_none")]
207    pub alert_policy: Option<FleetAlertPolicy>,
208    #[serde(default)]
209    pub timeout_seconds: Option<u64>,
210    #[serde(default)]
211    pub metadata: BTreeMap<String, Value>,
212}
213
214/// Worker role and tool expectations for a task.
215#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq, Eq)]
216pub struct FleetTaskWorkerProfile {
217    /// Named agent profile/persona posture to layer onto this worker.
218    ///
219    /// `profile` is accepted as a shorter authoring alias. This is an intent
220    /// reference only; profile loading and permission narrowing happen in the
221    /// Fleet runtime layer.
222    #[serde(default, alias = "profile", skip_serializing_if = "Option::is_none")]
223    pub agent_profile: Option<String>,
224    #[serde(skip_serializing_if = "Option::is_none")]
225    pub role: Option<String>,
226    /// Fleet loadout intent such as `auto`, `fast`, or `review`.
227    ///
228    /// This is not a concrete provider/model selection; route resolution owns
229    /// the executable provider/model/wire-model decision.
230    #[serde(default, skip_serializing_if = "Option::is_none")]
231    pub loadout: Option<String>,
232    /// Fleet model class hint such as `strong`, `balanced`, or `fast`.
233    #[serde(default, skip_serializing_if = "Option::is_none")]
234    pub model_class: Option<String>,
235    /// Optional explicit model id for this worker.
236    ///
237    /// Task-level model overrides are visible authoring data and take
238    /// precedence over the referenced agent profile's model hint. Provider and
239    /// wire-model validation still belong to route resolution.
240    #[serde(default, skip_serializing_if = "Option::is_none")]
241    pub model: Option<String>,
242    #[serde(skip_serializing_if = "Option::is_none")]
243    pub tool_profile: Option<String>,
244    #[serde(default)]
245    #[serde(skip_serializing_if = "Vec::is_empty")]
246    pub tools: Vec<String>,
247    #[serde(default)]
248    #[serde(skip_serializing_if = "Vec::is_empty")]
249    pub capabilities: Vec<String>,
250}
251
252/// Workspace and environment constraints needed before a task starts.
253#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq, Eq)]
254pub struct FleetWorkspaceRequirements {
255    #[serde(skip_serializing_if = "Option::is_none")]
256    pub root: Option<PathBuf>,
257    #[serde(default)]
258    #[serde(skip_serializing_if = "Vec::is_empty")]
259    pub required_files: Vec<PathBuf>,
260    #[serde(default)]
261    #[serde(skip_serializing_if = "Vec::is_empty")]
262    pub writable_paths: Vec<PathBuf>,
263    #[serde(skip_serializing_if = "Option::is_none")]
264    pub environment: Option<FleetEnvironmentRequirements>,
265}
266
267/// Environment variables a task requires or may pass through to workers.
268#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq, Eq)]
269pub struct FleetEnvironmentRequirements {
270    #[serde(default)]
271    #[serde(skip_serializing_if = "Vec::is_empty")]
272    pub required: Vec<String>,
273    #[serde(default)]
274    #[serde(skip_serializing_if = "Vec::is_empty")]
275    pub allowlist: Vec<String>,
276}
277
278/// Budget limits for a task.
279#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq, Eq)]
280pub struct FleetTaskBudget {
281    #[serde(skip_serializing_if = "Option::is_none")]
282    pub max_tokens: Option<u64>,
283    #[serde(skip_serializing_if = "Option::is_none")]
284    pub max_tool_calls: Option<u32>,
285    #[serde(skip_serializing_if = "Option::is_none")]
286    pub max_seconds: Option<u64>,
287}
288
289/// Reference to an artifact produced or consumed by a task.
290#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
291pub struct FleetArtifactRef {
292    pub kind: FleetArtifactKind,
293    pub path: PathBuf,
294    #[serde(skip_serializing_if = "Option::is_none")]
295    pub checksum: Option<String>,
296    #[serde(skip_serializing_if = "Option::is_none")]
297    pub mime_type: Option<String>,
298    #[serde(default)]
299    pub size_bytes: Option<u64>,
300}
301
302/// Kind of artifact a task may produce or consume.
303#[derive(Debug, Clone, PartialEq, Eq)]
304pub enum FleetArtifactKind {
305    Log,
306    Patch,
307    TestResult,
308    Report,
309    Checkpoint,
310    Receipt,
311    Other(String),
312}
313
314impl FleetArtifactKind {
315    fn as_wire_str(&self) -> &str {
316        match self {
317            Self::Log => "log",
318            Self::Patch => "patch",
319            Self::TestResult => "test_result",
320            Self::Report => "report",
321            Self::Checkpoint => "checkpoint",
322            Self::Receipt => "receipt",
323            Self::Other(kind) => kind.as_str(),
324        }
325    }
326
327    fn from_wire_str(value: &str) -> Self {
328        match value {
329            "log" => Self::Log,
330            "patch" => Self::Patch,
331            "test_result" => Self::TestResult,
332            "report" => Self::Report,
333            "checkpoint" => Self::Checkpoint,
334            "receipt" => Self::Receipt,
335            other => Self::Other(other.to_string()),
336        }
337    }
338}
339
340impl Serialize for FleetArtifactKind {
341    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
342    where
343        S: Serializer,
344    {
345        serializer.serialize_str(self.as_wire_str())
346    }
347}
348
349impl<'de> Deserialize<'de> for FleetArtifactKind {
350    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
351    where
352        D: Deserializer<'de>,
353    {
354        let value = String::deserialize(deserializer)?;
355        Ok(Self::from_wire_str(&value))
356    }
357}
358
359/// Scoring rule used to verify a task result.
360#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
361#[serde(tag = "kind", rename_all = "snake_case")]
362pub enum FleetScorerSpec {
363    ExitCode,
364    FileExists {
365        path: PathBuf,
366    },
367    RegexMatch {
368        path: PathBuf,
369        pattern: String,
370    },
371    JsonPath {
372        path: PathBuf,
373        expression: String,
374    },
375    Command {
376        command: String,
377        #[serde(default)]
378        args: Vec<String>,
379    },
380    CodeWhaleVerifierPrompt {
381        prompt: String,
382    },
383    Manual,
384}
385
386/// Worker specification.
387#[derive(Debug, Clone, Serialize, Deserialize)]
388pub struct FleetWorkerSpec {
389    pub id: String,
390    pub name: String,
391    pub host: FleetHostSpec,
392    #[serde(default)]
393    #[serde(skip_serializing_if = "Option::is_none")]
394    pub trust_level: Option<FleetTrustLevel>,
395    #[serde(default)]
396    pub labels: BTreeMap<String, String>,
397    #[serde(default)]
398    pub capabilities: Vec<String>,
399    #[serde(skip_serializing_if = "Option::is_none")]
400    pub max_concurrent_tasks: Option<usize>,
401}
402
403/// Host on which a worker runs.
404#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
405#[serde(tag = "kind", rename_all = "snake_case")]
406pub enum FleetHostSpec {
407    Local,
408    Ssh {
409        host: String,
410        #[serde(skip_serializing_if = "Option::is_none")]
411        port: Option<u16>,
412        #[serde(skip_serializing_if = "Option::is_none")]
413        user: Option<String>,
414        #[serde(skip_serializing_if = "Option::is_none")]
415        identity: Option<PathBuf>,
416        /// Known hosts file for host-key verification.
417        #[serde(skip_serializing_if = "Option::is_none")]
418        known_hosts: Option<PathBuf>,
419        /// Expected host key fingerprint (SHA256:...) for key pinning.
420        /// When set, the connection is only trusted if the server's
421        /// host key matches this fingerprint exactly.
422        #[serde(skip_serializing_if = "Option::is_none")]
423        host_key_fingerprint: Option<String>,
424        #[serde(skip_serializing_if = "Option::is_none")]
425        working_directory: Option<PathBuf>,
426        #[serde(default)]
427        #[serde(skip_serializing_if = "Vec::is_empty")]
428        env_allowlist: Vec<String>,
429        #[serde(skip_serializing_if = "Option::is_none")]
430        codewhale_binary: Option<String>,
431    },
432    #[serde(alias = "container")]
433    #[serde(alias = "Container")]
434    Docker {
435        image: String,
436        #[serde(default)]
437        args: Vec<String>,
438    },
439}
440
441// ── Security and trust types ────────────────────────────────────────────────
442
443/// Trust classification assigned to a worker host.
444///
445/// The trust level determines what a worker is allowed to do and what
446/// secrets it may access. The default for new workers is [`FleetTrustLevel::Sandbox`];
447/// operators must explicitly raise trust for SSH or container workers.
448#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord, Default)]
449#[serde(rename_all = "snake_case")]
450pub enum FleetTrustLevel {
451    /// Fully isolated: no network, no secrets, no writes outside `.codewhale/fleet/`.
452    /// Suitable for untrusted code review, community PR checks, or third-party tool runs.
453    #[default]
454    Sandbox = 0,
455    /// Local-only worker with access to the workspace and configured secrets.
456    /// Default for local workers. May read repo files but writes are gated.
457    Local = 1,
458    /// Worker on a known remote host with verified identity and a bounded
459    /// set of explicitly granted capabilities. Requires SSH host-key
460    /// verification or equivalent attestation.
461    #[serde(alias = "remote-verified", alias = "remoteVerified")]
462    RemoteVerified = 2,
463    /// Fully trusted worker (e.g. operator's own machine, CI runner).
464    /// Has access to all configured secrets and may perform any action the
465    /// operator can. Reserved for dogfood smoke and operator-owned machines.
466    Operator = 3,
467}
468
469impl FleetTrustLevel {
470    /// Whether this trust level is allowed to access provider secrets.
471    #[must_use]
472    pub fn may_access_secrets(&self) -> bool {
473        matches!(self, Self::Operator | Self::RemoteVerified | Self::Local)
474    }
475
476    /// Whether this trust level is allowed to write outside `.codewhale/fleet/`.
477    #[must_use]
478    pub fn may_write_workspace(&self) -> bool {
479        matches!(self, Self::Operator | Self::Local)
480    }
481
482    /// Whether this trust level is allowed network access.
483    #[must_use]
484    pub fn may_access_network(&self) -> bool {
485        matches!(self, Self::Operator | Self::RemoteVerified | Self::Local)
486    }
487}
488
489/// Security policy applied to a fleet run.
490///
491/// A policy defines the default trust level for workers, which secrets
492/// may be resolved, and what capabilities are granted. When a run has no
493/// explicit policy, workers inherit conservative defaults.
494#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
495pub struct FleetSecurityPolicy {
496    /// Default trust level for workers that don't declare one explicitly.
497    #[serde(default)]
498    pub default_trust_level: FleetTrustLevel,
499    /// Secret refs that workers may resolve. An empty list means no secrets
500    /// are available. Each entry is a key name, not a value.
501    #[serde(default)]
502    #[serde(skip_serializing_if = "Vec::is_empty")]
503    pub allowed_secrets: Vec<FleetSecretRef>,
504    /// Capability grants for workers in this run.
505    #[serde(default)]
506    #[serde(skip_serializing_if = "Vec::is_empty")]
507    pub capability_grants: Vec<FleetCapabilityGrant>,
508    /// Maximum trust level any worker in this run may have, even if the
509    /// worker spec requests higher. Defaults to Operator (no ceiling).
510    #[serde(default = "default_max_trust_level")]
511    pub max_trust_level: FleetTrustLevel,
512    /// Require identity verification for remote workers. When true, SSH
513    /// workers must pass host-key verification before being trusted at
514    /// RemoteVerified level; unverified remotes stay at Sandbox.
515    #[serde(default)]
516    pub require_identity_verification: bool,
517    /// Allow conservative parallel execution of read-only tools (#2983).
518    /// When true, workers may batch independent read-only tool calls
519    /// (reads, searches, greps) into concurrent turns. Disabled by default
520    /// to avoid overwhelming providers or hitting rate limits.
521    #[serde(default)]
522    pub allow_parallel_reads: bool,
523}
524
525fn default_max_trust_level() -> FleetTrustLevel {
526    FleetTrustLevel::Operator
527}
528
529impl Default for FleetSecurityPolicy {
530    fn default() -> Self {
531        Self {
532            default_trust_level: FleetTrustLevel::Sandbox,
533            allowed_secrets: Vec::new(),
534            capability_grants: Vec::new(),
535            max_trust_level: FleetTrustLevel::Operator,
536            require_identity_verification: false,
537            allow_parallel_reads: false,
538        }
539    }
540}
541
542/// A reference to a secret that should be resolved at runtime, never
543/// serialized as a plaintext value.
544///
545/// Secret refs appear in task specs, alert configs, and worker definitions.
546/// The actual secret value is resolved by the fleet manager from the
547/// secrets backend (OS keyring, environment, or file store) just before
548/// the worker starts.
549#[derive(Debug, Clone, Serialize, PartialEq, Eq, Hash, PartialOrd, Ord)]
550pub struct FleetSecretRef {
551    /// The secret key name (e.g. `"CODEWHALE_API_KEY"`, `"GH_TOKEN"`).
552    pub key: String,
553    /// Optional source hint for resolution order.
554    /// - `"env"` — resolve from environment variable
555    /// - `"keyring"` — resolve from OS keyring
556    /// - `"file"` — resolve from `~/.codewhale/secrets/`
557    /// - absent / null — try all sources in default order
558    #[serde(skip_serializing_if = "Option::is_none")]
559    pub source: Option<String>,
560}
561
562impl FleetSecretRef {
563    /// Create a secret ref from a key name with default resolution.
564    #[must_use]
565    pub fn new(key: impl Into<String>) -> Self {
566        Self {
567            key: key.into(),
568            source: None,
569        }
570    }
571
572    /// Create a secret ref with an explicit source.
573    #[must_use]
574    pub fn with_source(key: impl Into<String>, source: impl Into<String>) -> Self {
575        Self {
576            key: key.into(),
577            source: Some(source.into()),
578        }
579    }
580
581    /// Redacted display form for logging. Shows the key name and source
582    /// but never the resolved value.
583    #[must_use]
584    pub fn redacted(&self) -> String {
585        match &self.source {
586            Some(src) => format!("<secret:{}.{}>", src, self.key),
587            None => format!("<secret:{}>", self.key),
588        }
589    }
590}
591
592impl std::fmt::Display for FleetSecretRef {
593    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
594        write!(f, "{}", self.redacted())
595    }
596}
597
598impl From<&str> for FleetSecretRef {
599    fn from(key: &str) -> Self {
600        Self::new(key)
601    }
602}
603
604impl From<String> for FleetSecretRef {
605    fn from(key: String) -> Self {
606        Self::new(key)
607    }
608}
609
610impl<'de> Deserialize<'de> for FleetSecretRef {
611    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
612    where
613        D: Deserializer<'de>,
614    {
615        #[derive(Deserialize)]
616        #[serde(untagged)]
617        enum SecretRefWire {
618            Key(String),
619            Structured {
620                key: String,
621                #[serde(default)]
622                source: Option<String>,
623            },
624        }
625
626        match SecretRefWire::deserialize(deserializer)? {
627            SecretRefWire::Key(key) if !key.trim().is_empty() => Ok(FleetSecretRef::new(key)),
628            SecretRefWire::Key(_) => Err(de::Error::custom("secret ref key cannot be empty")),
629            SecretRefWire::Structured { key, source } if !key.trim().is_empty() => {
630                Ok(FleetSecretRef { key, source })
631            }
632            SecretRefWire::Structured { .. } => {
633                Err(de::Error::custom("secret ref key cannot be empty"))
634            }
635        }
636    }
637}
638
639/// How a worker authenticates to the fleet manager.
640#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
641#[serde(tag = "method", rename_all = "snake_case")]
642pub enum FleetWorkerAuth {
643    /// No authentication (local workers share the same uid).
644    None,
645    /// SSH key-based authentication with host-key verification.
646    SshKey {
647        /// Path to the SSH identity file (may be a FleetSecretRef in JSON
648        /// as `{"key": "...", "source": "file"}`).
649        identity: PathBuf,
650        /// Known hosts file for host-key verification.
651        #[serde(skip_serializing_if = "Option::is_none")]
652        known_hosts: Option<PathBuf>,
653        /// Expected host key fingerprint for pinning.
654        #[serde(skip_serializing_if = "Option::is_none")]
655        host_key_fingerprint: Option<String>,
656        /// SSH user for the connection.
657        #[serde(skip_serializing_if = "Option::is_none")]
658        user: Option<String>,
659    },
660    /// Token-based authentication for remote workers behind a fleet proxy.
661    Token {
662        /// Reference to the token secret.
663        token_ref: FleetSecretRef,
664    },
665    /// mTLS certificate-based authentication.
666    Mtls {
667        /// Path to the client certificate.
668        cert_path: PathBuf,
669        /// Reference to the private key secret.
670        key_ref: FleetSecretRef,
671    },
672}
673
674/// A capability grant that explicitly authorizes a worker to perform
675/// a specific class of action.
676///
677/// By default, new workers get no grants (least privilege). Grants are
678/// additive: a worker's effective capabilities are the union of its
679/// trust-level defaults plus any explicit grants.
680#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
681pub struct FleetCapabilityGrant {
682    /// The capability being granted (e.g. `"network"`, `"git-push"`,
683    /// `"provider-secrets"`, `"release"`).
684    pub capability: String,
685    /// Optional scope limiting the grant (e.g. `"github.com"` for network,
686    /// `"crates/tui/**"` for file writes).
687    #[serde(skip_serializing_if = "Option::is_none")]
688    pub scope: Option<String>,
689    /// Optional justification for the grant (audit trail).
690    #[serde(skip_serializing_if = "Option::is_none")]
691    pub reason: Option<String>,
692}
693
694/// Runtime status of a worker.
695#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
696#[serde(rename_all = "snake_case")]
697pub enum FleetWorkerStatus {
698    Unknown,
699    Online,
700    Busy,
701    Offline,
702    Unhealthy,
703    Draining,
704    Retired,
705}
706
707impl Status for FleetWorkerStatus {
708    fn is_terminal(&self) -> bool {
709        matches!(self, Self::Retired)
710    }
711    fn is_active(&self) -> bool {
712        matches!(self, Self::Online | Self::Busy)
713    }
714    fn is_paused(&self) -> bool {
715        false
716    }
717}
718
719/// Durable inbox entry: a task waiting to be leased to a worker.
720#[derive(Debug, Clone, Serialize, Deserialize)]
721pub struct FleetInboxEntry {
722    pub run_id: FleetRunId,
723    pub task_id: String,
724    pub priority: i32,
725    pub enqueued_at: String,
726    #[serde(default)]
727    pub lease_deadline: Option<String>,
728    #[serde(default)]
729    pub attempts: u32,
730}
731
732/// Worker event envelope.
733#[derive(Debug, Clone, Serialize, Deserialize)]
734pub struct FleetWorkerEvent {
735    pub seq: u64,
736    pub run_id: FleetRunId,
737    pub worker_id: String,
738    pub task_id: String,
739    pub timestamp: String,
740    #[serde(flatten)]
741    pub payload: FleetWorkerEventPayload,
742    #[serde(default)]
743    #[serde(skip_serializing_if = "BTreeMap::is_empty")]
744    pub extra: BTreeMap<String, Value>,
745}
746
747/// Union of all worker event payloads.
748#[derive(Debug, Clone, Serialize, Deserialize)]
749#[serde(tag = "state", rename_all = "snake_case")]
750pub enum FleetWorkerEventPayload {
751    Queued,
752    Leased {
753        #[serde(skip_serializing_if = "Option::is_none")]
754        lease_expires_at: Option<String>,
755    },
756    Starting,
757    Running,
758    ModelWait {
759        #[serde(skip_serializing_if = "Option::is_none")]
760        model: Option<String>,
761    },
762    RunningTool {
763        tool: String,
764        #[serde(skip_serializing_if = "Option::is_none")]
765        call_id: Option<String>,
766    },
767    /// Typed receipt emitted by a Workflow running inside this worker.
768    WorkflowEvent {
769        /// Inner Workflow run id. Named distinctly from the outer Fleet run id
770        /// because payloads are flattened into `FleetWorkerEvent`.
771        workflow_run_id: String,
772        event: Value,
773    },
774    Heartbeat {
775        #[serde(default)]
776        #[serde(skip_serializing_if = "Option::is_none")]
777        cpu_percent: Option<f32>,
778        #[serde(default)]
779        #[serde(skip_serializing_if = "Option::is_none")]
780        memory_mb: Option<u64>,
781    },
782    Artifact(FleetArtifactRef),
783    Completed {
784        #[serde(default)]
785        #[serde(skip_serializing_if = "Option::is_none")]
786        exit_code: Option<i32>,
787        #[serde(skip_serializing_if = "Option::is_none")]
788        summary: Option<String>,
789    },
790    Failed {
791        reason: String,
792        #[serde(default)]
793        recoverable: bool,
794    },
795    Cancelled {
796        #[serde(skip_serializing_if = "Option::is_none")]
797        cancelled_by: Option<String>,
798    },
799    Interrupted {
800        #[serde(skip_serializing_if = "Option::is_none")]
801        signal: Option<String>,
802    },
803    Stale {
804        #[serde(skip_serializing_if = "Option::is_none")]
805        last_heartbeat_at: Option<String>,
806    },
807    Restarted {
808        #[serde(default)]
809        restart_count: u32,
810    },
811    Escalated {
812        channel: String,
813        #[serde(skip_serializing_if = "Option::is_none")]
814        alert_id: Option<String>,
815    },
816}
817
818/// Retry policy for a task or worker.
819#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
820pub struct FleetRetryPolicy {
821    #[serde(default = "default_retry_max_attempts")]
822    pub max_attempts: u32,
823    #[serde(default = "default_retry_initial_backoff_seconds")]
824    pub initial_backoff_seconds: u64,
825    #[serde(default = "default_retry_max_backoff_seconds")]
826    pub max_backoff_seconds: u64,
827    #[serde(default = "default_retry_backoff_multiplier")]
828    pub backoff_multiplier: u32,
829}
830
831impl Default for FleetRetryPolicy {
832    fn default() -> Self {
833        Self {
834            max_attempts: 3,
835            initial_backoff_seconds: 5,
836            max_backoff_seconds: 300,
837            backoff_multiplier: 2,
838        }
839    }
840}
841
842fn default_retry_max_attempts() -> u32 {
843    FleetRetryPolicy::default().max_attempts
844}
845
846fn default_retry_initial_backoff_seconds() -> u64 {
847    FleetRetryPolicy::default().initial_backoff_seconds
848}
849
850fn default_retry_max_backoff_seconds() -> u64 {
851    FleetRetryPolicy::default().max_backoff_seconds
852}
853
854fn default_retry_backoff_multiplier() -> u32 {
855    FleetRetryPolicy::default().backoff_multiplier
856}
857
858/// Alert/escalation policy attached to a task or run.
859#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
860pub struct FleetAlertPolicy {
861    #[serde(default)]
862    #[serde(skip_serializing_if = "Vec::is_empty")]
863    pub events: Vec<FleetAlertEventClass>,
864    #[serde(default)]
865    pub channels: Vec<FleetAlertChannel>,
866    #[serde(default)]
867    pub after_attempts: Option<u32>,
868    #[serde(default)]
869    pub after_minutes_stale: Option<u64>,
870}
871
872#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord, Hash)]
873#[serde(rename_all = "snake_case")]
874pub enum FleetAlertEventClass {
875    Stale,
876    RestartExhausted,
877    NeedsHuman,
878    BudgetExceeded,
879    VerifierFailed,
880    RunCompleted,
881}
882
883#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
884#[serde(tag = "kind", rename_all = "snake_case")]
885pub enum FleetAlertChannel {
886    Slack {
887        /// Webhook URL, resolved from a secret ref or inline.
888        #[serde(flatten)]
889        webhook: FleetAlertEndpoint,
890    },
891    Webhook {
892        #[serde(flatten)]
893        endpoint: FleetAlertEndpoint,
894    },
895    #[serde(alias = "pager_duty")]
896    #[serde(alias = "pagerduty")]
897    PagerDuty {
898        routing_key: String,
899        severity: String,
900    },
901}
902
903/// An alert channel endpoint, supporting both inline URLs and secret refs.
904///
905/// For Slack and generic webhook channels, the URL may be provided directly
906/// or as a secret reference resolved at send time. When both `url` and
907/// `url_ref` are present, `url_ref` takes precedence after resolution.
908#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
909pub struct FleetAlertEndpoint {
910    /// Inline URL (plaintext; only for non-sensitive endpoints).
911    #[serde(
912        alias = "webhook_url",
913        alias = "endpoint_url",
914        skip_serializing_if = "Option::is_none"
915    )]
916    pub url: Option<String>,
917    /// Reference to a secret containing the webhook URL.
918    #[serde(
919        alias = "webhook_url_ref",
920        alias = "webhook_ref",
921        alias = "url_secret_ref",
922        skip_serializing_if = "Option::is_none"
923    )]
924    pub url_ref: Option<FleetSecretRef>,
925    /// Optional HMAC secret for webhook payload signing, as a secret ref.
926    #[serde(
927        alias = "secret",
928        alias = "webhook_secret",
929        alias = "signing_secret",
930        skip_serializing_if = "Option::is_none"
931    )]
932    pub secret_ref: Option<FleetSecretRef>,
933}
934
935impl FleetAlertEndpoint {
936    /// Create an inline URL endpoint (for non-sensitive use).
937    #[must_use]
938    pub fn inline(url: impl Into<String>) -> Self {
939        Self {
940            url: Some(url.into()),
941            url_ref: None,
942            secret_ref: None,
943        }
944    }
945
946    /// Create a secret-backed URL endpoint.
947    #[must_use]
948    pub fn from_secret(url_ref: FleetSecretRef) -> Self {
949        Self {
950            url: None,
951            url_ref: Some(url_ref),
952            secret_ref: None,
953        }
954    }
955
956    /// Redacted display form for logging.
957    #[must_use]
958    pub fn redacted(&self) -> String {
959        self.url_ref
960            .as_ref()
961            .map_or_else(|| "<inline-url>".to_string(), |r| r.redacted())
962    }
963}
964
965/// Resolved-route detail persisted on a [`FleetReceipt`] (#3154).
966///
967/// This is an additive, *plain-strings* snapshot of the route a fleet worker
968/// resolved to. It deliberately does NOT depend on any `codewhale-config` route
969/// type so the protocol crate stays free of the route model.
970///
971/// CRITICAL no-secrets invariant: this struct carries ONLY non-sensitive route
972/// shape — provider id/kind, model ids, wire protocol, role/loadout/model-class
973/// intent, reasoning tier when known, and deterministic intent sources. It
974/// must NEVER hold a credential, API key, bearer token, or a base URL that
975/// embeds credentials. There is intentionally no field that could carry a
976/// secret.
977#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
978pub struct FleetResolvedRoute {
979    /// Resolved provider canonical id (e.g. `"deepseek"`).
980    pub provider_id: String,
981    /// Exact configured provider-table id when the worker used one.
982    ///
983    /// This is intentionally additive to `provider_id`: literal
984    /// `[providers.custom]` resolves to `Some("custom")`, while the legacy
985    /// idless root custom route resolves to `None`. Keeping the distinction
986    /// prevents a receipt from silently collapsing two different credential
987    /// and endpoint authorities into the same generic `custom` label.
988    #[serde(default, skip_serializing_if = "Option::is_none")]
989    pub provider_exact_id: Option<String>,
990    /// Resolved provider kind (e.g. `"deepseek"`).
991    pub provider_kind: String,
992    /// Canonical, provider-agnostic model identity, when known.
993    #[serde(default, skip_serializing_if = "Option::is_none")]
994    pub canonical_model: Option<String>,
995    /// Provider-owned wire model id placed on the request.
996    pub wire_model_id: String,
997    /// Selected wire protocol (e.g. `"chat_completions"`).
998    pub protocol: String,
999    /// Effective Fleet role intent, when one applied.
1000    #[serde(default, skip_serializing_if = "Option::is_none")]
1001    pub role: Option<String>,
1002    /// Effective Fleet loadout intent, when one applied.
1003    #[serde(default, skip_serializing_if = "Option::is_none")]
1004    pub loadout: Option<String>,
1005    /// Original task-level model-class intent, when authored separately from
1006    /// `loadout`. Profile `model_class_hint` is normalized into `loadout`.
1007    #[serde(default, skip_serializing_if = "Option::is_none")]
1008    pub model_class: Option<String>,
1009    /// Runtime model-route seam used by sub-agent routing (`inherit`, `faster`,
1010    /// `auto`, or `fixed`).
1011    #[serde(default, skip_serializing_if = "Option::is_none")]
1012    pub model_route: Option<String>,
1013    /// Concrete reasoning tier, when it is known by the route resolver path.
1014    #[serde(default, skip_serializing_if = "Option::is_none")]
1015    pub reasoning_effort: Option<String>,
1016    /// Deterministic source for the effective role intent.
1017    #[serde(default, skip_serializing_if = "Option::is_none")]
1018    pub role_source: Option<String>,
1019    /// Deterministic source for the effective loadout intent.
1020    #[serde(default, skip_serializing_if = "Option::is_none")]
1021    pub loadout_source: Option<String>,
1022    /// Deterministic source for the model-class hint, when present.
1023    #[serde(default, skip_serializing_if = "Option::is_none")]
1024    pub model_class_source: Option<String>,
1025    /// Deterministic source for the model selector used by the resolver.
1026    #[serde(default, skip_serializing_if = "Option::is_none")]
1027    pub model_source: Option<String>,
1028    /// How the route was produced (e.g. `"resolver"`).
1029    pub source: String,
1030}
1031
1032/// Effective worker authority persisted on a [`FleetReceipt`] (#3211).
1033///
1034/// This is a non-secret snapshot of the already-computed runtime profile. It
1035/// records what the worker was allowed to do; it does not grant permissions and
1036/// does not carry credentials, sandbox paths, or provider endpoints.
1037#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1038pub struct FleetEffectivePermissions {
1039    /// Whether the worker profile may modify workspace files.
1040    pub write: bool,
1041    /// Whether the worker profile may use network-capable tools.
1042    pub network: bool,
1043    /// Shell posture (`none`, `read_only`, or `full`).
1044    pub shell: String,
1045    /// Tool-surface posture (`inherit` or `explicit`).
1046    pub tool_scope: String,
1047    /// Explicit tool names when `tool_scope` is `explicit`.
1048    #[serde(default, skip_serializing_if = "Vec::is_empty")]
1049    pub tools: Vec<String>,
1050    /// Whether the worker is intended to run detached/background.
1051    pub background: bool,
1052    /// Remaining nested-delegation budget after parent intersection/hardening.
1053    pub max_spawn_depth: u32,
1054    /// Roster profile id that contributed to this worker, when any.
1055    #[serde(default, skip_serializing_if = "Option::is_none")]
1056    pub profile_id: Option<String>,
1057    /// Roster layer for `profile_id` (`built_in`, `config`, or `workspace`).
1058    #[serde(default, skip_serializing_if = "Option::is_none")]
1059    pub profile_origin: Option<String>,
1060    /// How this snapshot was produced (e.g. `"worker_runtime_profile"`).
1061    pub source: String,
1062}
1063
1064/// Receipt produced when a task completes verification.
1065#[derive(Debug, Clone, Serialize, Deserialize)]
1066pub struct FleetReceipt {
1067    pub run_id: FleetRunId,
1068    pub task_id: String,
1069    pub worker_id: String,
1070    /// Durable lease generation that produced this receipt.
1071    ///
1072    /// Optional for backward compatibility with receipts written before Fleet
1073    /// attempts were fenced explicitly.
1074    #[serde(default, skip_serializing_if = "Option::is_none")]
1075    pub attempt: Option<u32>,
1076    /// Sequence of the terminal worker event finalized with this receipt.
1077    ///
1078    /// Optional so older ledger records remain replayable.
1079    #[serde(default, skip_serializing_if = "Option::is_none")]
1080    pub terminal_seq: Option<u64>,
1081    pub completed_at: String,
1082    pub result: FleetTaskResult,
1083    #[serde(skip_serializing_if = "Option::is_none")]
1084    pub failure_kind: Option<FleetTaskFailureKind>,
1085    #[serde(default)]
1086    pub artifacts: Vec<FleetArtifactRef>,
1087    #[serde(default)]
1088    pub score: Option<FleetScore>,
1089    /// Resolved-route snapshot for this task (#3154).
1090    ///
1091    /// `#[serde(default)]` keeps older ledgers (written before this field
1092    /// existed) deserializable.
1093    #[serde(default, skip_serializing_if = "Option::is_none")]
1094    pub resolved_route: Option<FleetResolvedRoute>,
1095    /// Effective worker authority for this task (#3211).
1096    #[serde(default, skip_serializing_if = "Option::is_none")]
1097    pub effective_permissions: Option<FleetEffectivePermissions>,
1098}
1099
1100#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1101#[serde(rename_all = "snake_case")]
1102pub enum FleetTaskResult {
1103    Pass,
1104    Partial,
1105    Fail,
1106    Skip,
1107    Timeout,
1108}
1109
1110/// Source category for a failed task receipt.
1111#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
1112#[serde(rename_all = "snake_case")]
1113pub enum FleetTaskFailureKind {
1114    Transport,
1115    Task,
1116    Verifier,
1117}
1118
1119#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
1120pub struct FleetScore {
1121    pub value: f64,
1122    #[serde(skip_serializing_if = "Option::is_none")]
1123    pub max: Option<f64>,
1124    #[serde(skip_serializing_if = "Option::is_none")]
1125    pub notes: Option<String>,
1126}
1127
1128#[cfg(test)]
1129mod tests {
1130    use super::*;
1131
1132    #[test]
1133    fn fleet_run_round_trip() {
1134        let run = FleetRun {
1135            id: FleetRunId::from("run-001"),
1136            name: "dogfood smoke".to_string(),
1137            status: FleetRunStatus::Running,
1138            target: Some(FleetRuntimeTarget::ThisComputer),
1139            workflow: Some(FleetWorkflowDescriptor {
1140                id: "release-checks".to_string(),
1141                kind: FleetWorkflowKind::Parallel,
1142            }),
1143            roles: vec!["release-checker".to_string()],
1144            max_workers: Some(1),
1145            task_specs: vec![FleetTaskSpec {
1146                id: "task-1".to_string(),
1147                name: "lint".to_string(),
1148                description: None,
1149                objective: Some("Keep the workspace lint-clean".to_string()),
1150                instructions: "run cargo clippy".to_string(),
1151                worker: Some(FleetTaskWorkerProfile {
1152                    agent_profile: None,
1153                    role: Some("release-checker".to_string()),
1154                    loadout: None,
1155                    model_class: None,
1156                    model: None,
1157                    tool_profile: Some("read-only".to_string()),
1158                    tools: vec!["cargo".to_string()],
1159                    capabilities: vec!["rust".to_string()],
1160                }),
1161                workspace: Some(FleetWorkspaceRequirements {
1162                    root: Some(PathBuf::from(".")),
1163                    required_files: vec![PathBuf::from("Cargo.toml")],
1164                    writable_paths: vec![],
1165                    environment: Some(FleetEnvironmentRequirements {
1166                        required: vec!["PATH".to_string()],
1167                        allowlist: vec!["RUST_LOG".to_string()],
1168                    }),
1169                }),
1170                input_files: vec![PathBuf::from("crates/tui/src/main.rs")],
1171                context: vec!["release gate".to_string()],
1172                budget: Some(FleetTaskBudget {
1173                    max_tokens: Some(8000),
1174                    max_tool_calls: Some(20),
1175                    max_seconds: Some(300),
1176                }),
1177                tags: vec!["release".to_string()],
1178                expected_artifacts: vec![FleetArtifactKind::Log],
1179                scorer: Some(FleetScorerSpec::ExitCode),
1180                retry_policy: Some(FleetRetryPolicy::default()),
1181                alert_policy: None,
1182                timeout_seconds: Some(300),
1183                metadata: BTreeMap::new(),
1184            }],
1185            worker_specs: vec![],
1186            labels: BTreeMap::new(),
1187            security_policy: None,
1188            created_at: "2026-06-12T17:00:00Z".to_string(),
1189            updated_at: None,
1190            completed_at: None,
1191        };
1192        let json = serde_json::to_string(&run).unwrap();
1193        let back: FleetRun = serde_json::from_str(&json).unwrap();
1194        assert_eq!(back.id, run.id);
1195        assert_eq!(back.status, FleetRunStatus::Running);
1196        assert_eq!(back.target, Some(FleetRuntimeTarget::ThisComputer));
1197        assert_eq!(back.roles, vec!["release-checker"]);
1198        assert_eq!(
1199            back.workflow.as_ref().map(|workflow| workflow.id.as_str()),
1200            Some("release-checks")
1201        );
1202        assert_eq!(back.task_specs.len(), 1);
1203        assert_eq!(
1204            back.task_specs[0].worker.as_ref().unwrap().role.as_deref(),
1205            Some("release-checker")
1206        );
1207        assert_eq!(
1208            back.task_specs[0]
1209                .workspace
1210                .as_ref()
1211                .unwrap()
1212                .required_files,
1213            vec![PathBuf::from("Cargo.toml")]
1214        );
1215    }
1216
1217    #[test]
1218    fn worker_profile_carries_agent_profile_and_loadout_intent() {
1219        let json = r#"{
1220            "profile": "adversarial_reviewer",
1221            "role": "reviewer",
1222            "loadout": "auto",
1223            "model_class": "balanced",
1224            "model": "deepseek-v4-pro",
1225            "tool_profile": "read-only",
1226            "tools": ["read_file"],
1227            "capabilities": ["rust"]
1228        }"#;
1229
1230        let profile: FleetTaskWorkerProfile = serde_json::from_str(json).unwrap();
1231
1232        assert_eq!(
1233            profile.agent_profile.as_deref(),
1234            Some("adversarial_reviewer")
1235        );
1236        assert_eq!(profile.role.as_deref(), Some("reviewer"));
1237        assert_eq!(profile.loadout.as_deref(), Some("auto"));
1238        assert_eq!(profile.model_class.as_deref(), Some("balanced"));
1239        assert_eq!(profile.model.as_deref(), Some("deepseek-v4-pro"));
1240        assert_eq!(profile.tool_profile.as_deref(), Some("read-only"));
1241
1242        let serialized = serde_json::to_value(&profile).unwrap();
1243        assert_eq!(serialized["agent_profile"], "adversarial_reviewer");
1244        assert_eq!(serialized["model"], "deepseek-v4-pro");
1245        assert!(serialized.get("profile").is_none());
1246    }
1247
1248    #[test]
1249    fn worker_event_lifecycle_round_trip() {
1250        let events = vec![
1251            FleetWorkerEvent {
1252                seq: 1,
1253                run_id: FleetRunId::from("run-002"),
1254                worker_id: "worker-a".to_string(),
1255                task_id: "task-1".to_string(),
1256                timestamp: "2026-06-12T17:01:00Z".to_string(),
1257                payload: FleetWorkerEventPayload::Queued,
1258                extra: BTreeMap::new(),
1259            },
1260            FleetWorkerEvent {
1261                seq: 2,
1262                run_id: FleetRunId::from("run-002"),
1263                worker_id: "worker-a".to_string(),
1264                task_id: "task-1".to_string(),
1265                timestamp: "2026-06-12T17:01:05Z".to_string(),
1266                payload: FleetWorkerEventPayload::RunningTool {
1267                    tool: "bash".to_string(),
1268                    call_id: Some("call-1".to_string()),
1269                },
1270                extra: BTreeMap::new(),
1271            },
1272            FleetWorkerEvent {
1273                seq: 3,
1274                run_id: FleetRunId::from("run-002"),
1275                worker_id: "worker-a".to_string(),
1276                task_id: "task-1".to_string(),
1277                timestamp: "2026-06-12T17:02:00Z".to_string(),
1278                payload: FleetWorkerEventPayload::Completed {
1279                    exit_code: Some(0),
1280                    summary: Some("ok".to_string()),
1281                },
1282                extra: BTreeMap::new(),
1283            },
1284        ];
1285        let json = serde_json::to_string(&events).unwrap();
1286        let back: Vec<FleetWorkerEvent> = serde_json::from_str(&json).unwrap();
1287        assert_eq!(back.len(), 3);
1288        assert!(matches!(back[0].payload, FleetWorkerEventPayload::Queued));
1289        assert!(matches!(
1290            back[2].payload,
1291            FleetWorkerEventPayload::Completed { .. }
1292        ));
1293    }
1294
1295    #[test]
1296    fn workflow_receipt_round_trip_keeps_outer_and_inner_run_ids_distinct() {
1297        let event = FleetWorkerEvent {
1298            seq: 3,
1299            run_id: FleetRunId::from("fleet-run-1"),
1300            worker_id: "worker-a".to_string(),
1301            task_id: "task-1".to_string(),
1302            timestamp: "2026-07-10T00:00:00Z".to_string(),
1303            payload: FleetWorkerEventPayload::WorkflowEvent {
1304                workflow_run_id: "workflow_1".to_string(),
1305                event: serde_json::json!({"type": "task_completed"}),
1306            },
1307            extra: BTreeMap::new(),
1308        };
1309        let value = serde_json::to_value(&event).unwrap();
1310        assert_eq!(value["run_id"], "fleet-run-1");
1311        assert_eq!(value["workflow_run_id"], "workflow_1");
1312        let back: FleetWorkerEvent = serde_json::from_value(value).unwrap();
1313        assert!(matches!(
1314            back.payload,
1315            FleetWorkerEventPayload::WorkflowEvent {
1316                workflow_run_id,
1317                ref event,
1318            } if workflow_run_id == "workflow_1" && event["type"] == "task_completed"
1319        ));
1320    }
1321
1322    #[test]
1323    fn alert_policy_round_trip() {
1324        let policy = FleetAlertPolicy {
1325            events: vec![FleetAlertEventClass::Stale],
1326            channels: vec![FleetAlertChannel::Slack {
1327                webhook: FleetAlertEndpoint::inline("https://hooks.slack.com/test"),
1328            }],
1329            after_attempts: Some(2),
1330            after_minutes_stale: Some(10),
1331        };
1332        let json = serde_json::to_string(&policy).unwrap();
1333        assert!(json.contains("\"events\":[\"stale\"]"));
1334        assert!(json.contains("\"kind\":\"slack\""));
1335        let back: FleetAlertPolicy = serde_json::from_str(&json).unwrap();
1336        assert_eq!(back.events, vec![FleetAlertEventClass::Stale]);
1337        assert_eq!(back.after_attempts, Some(2));
1338    }
1339
1340    #[test]
1341    fn artifact_other_kind_round_trip() {
1342        let artifact = FleetArtifactRef {
1343            kind: FleetArtifactKind::Other("coverage.xml".to_string()),
1344            path: PathBuf::from("/tmp/coverage.xml"),
1345            checksum: Some("sha256:abc".to_string()),
1346            mime_type: Some("application/xml".to_string()),
1347            size_bytes: Some(1024),
1348        };
1349        let json = serde_json::to_string(&artifact).unwrap();
1350        let back: FleetArtifactRef = serde_json::from_str(&json).unwrap();
1351        assert_eq!(back.kind, artifact.kind);
1352        assert_eq!(back.size_bytes, Some(1024));
1353    }
1354
1355    #[test]
1356    fn ssh_host_spec_accepts_minimal_legacy_json() {
1357        let json = r#"{"kind":"ssh","host":"builder.example.test"}"#;
1358        let host: FleetHostSpec = serde_json::from_str(json).unwrap();
1359
1360        match host {
1361            FleetHostSpec::Ssh {
1362                host,
1363                port,
1364                user,
1365                identity,
1366                known_hosts,
1367                host_key_fingerprint,
1368                working_directory,
1369                env_allowlist,
1370                codewhale_binary,
1371            } => {
1372                assert_eq!(host, "builder.example.test");
1373                assert_eq!(port, None);
1374                assert_eq!(user, None);
1375                assert_eq!(identity, None);
1376                assert_eq!(known_hosts, None);
1377                assert_eq!(host_key_fingerprint, None);
1378                assert_eq!(working_directory, None);
1379                assert!(env_allowlist.is_empty());
1380                assert_eq!(codewhale_binary, None);
1381            }
1382            other => panic!("expected ssh host spec, got {other:?}"),
1383        }
1384    }
1385
1386    #[test]
1387    fn artifact_kind_uses_flat_string_json() {
1388        let known = serde_json::to_string(&FleetArtifactKind::TestResult).unwrap();
1389        assert_eq!(known, "\"test_result\"");
1390
1391        let custom =
1392            serde_json::to_string(&FleetArtifactKind::Other("coverage.xml".to_string())).unwrap();
1393        assert_eq!(custom, "\"coverage.xml\"");
1394
1395        let parsed: FleetArtifactKind = serde_json::from_str("\"coverage.xml\"").unwrap();
1396        assert_eq!(parsed, FleetArtifactKind::Other("coverage.xml".to_string()));
1397    }
1398
1399    #[test]
1400    fn retry_policy_missing_fields_use_nonzero_defaults() {
1401        let policy: FleetRetryPolicy = serde_json::from_value(serde_json::json!({})).unwrap();
1402        assert_eq!(policy, FleetRetryPolicy::default());
1403
1404        let policy: FleetRetryPolicy =
1405            serde_json::from_value(serde_json::json!({"max_attempts": 5})).unwrap();
1406        assert_eq!(policy.max_attempts, 5);
1407        assert_eq!(
1408            policy.initial_backoff_seconds,
1409            FleetRetryPolicy::default().initial_backoff_seconds
1410        );
1411        assert_eq!(
1412            policy.max_backoff_seconds,
1413            FleetRetryPolicy::default().max_backoff_seconds
1414        );
1415        assert_eq!(
1416            policy.backoff_multiplier,
1417            FleetRetryPolicy::default().backoff_multiplier
1418        );
1419    }
1420
1421    #[test]
1422    fn sparse_worker_events_omit_absent_optional_fields() {
1423        let heartbeat = FleetWorkerEventPayload::Heartbeat {
1424            cpu_percent: None,
1425            memory_mb: None,
1426        };
1427        let heartbeat_json = serde_json::to_value(&heartbeat).unwrap();
1428        assert_eq!(heartbeat_json, serde_json::json!({"state": "heartbeat"}));
1429
1430        let completed = FleetWorkerEventPayload::Completed {
1431            exit_code: None,
1432            summary: None,
1433        };
1434        let completed_json = serde_json::to_value(&completed).unwrap();
1435        assert_eq!(completed_json, serde_json::json!({"state": "completed"}));
1436    }
1437
1438    #[test]
1439    fn receipt_round_trip() {
1440        let receipt = FleetReceipt {
1441            run_id: FleetRunId::from("run-003"),
1442            task_id: "task-1".to_string(),
1443            worker_id: "worker-b".to_string(),
1444            attempt: Some(2),
1445            terminal_seq: Some(7),
1446            completed_at: "2026-06-12T17:03:00Z".to_string(),
1447            result: FleetTaskResult::Pass,
1448            failure_kind: None,
1449            artifacts: vec![],
1450            score: Some(FleetScore {
1451                value: 0.95,
1452                max: Some(1.0),
1453                notes: None,
1454            }),
1455            resolved_route: None,
1456            effective_permissions: None,
1457        };
1458        let json = serde_json::to_string(&receipt).unwrap();
1459        let back: FleetReceipt = serde_json::from_str(&json).unwrap();
1460        assert_eq!(back.result, FleetTaskResult::Pass);
1461        assert_eq!(back.score.as_ref().unwrap().value, 0.95);
1462        assert_eq!(back.attempt, Some(2));
1463        assert_eq!(back.terminal_seq, Some(7));
1464    }
1465
1466    #[test]
1467    fn partial_receipt_records_failure_source_when_needed() {
1468        let receipt = FleetReceipt {
1469            run_id: FleetRunId::from("run-004"),
1470            task_id: "task-2".to_string(),
1471            worker_id: "worker-c".to_string(),
1472            attempt: None,
1473            terminal_seq: None,
1474            completed_at: "2026-06-12T17:04:00Z".to_string(),
1475            result: FleetTaskResult::Partial,
1476            failure_kind: Some(FleetTaskFailureKind::Verifier),
1477            artifacts: vec![],
1478            score: Some(FleetScore {
1479                value: 0.5,
1480                max: Some(1.0),
1481                notes: Some("manual verification required".to_string()),
1482            }),
1483            resolved_route: None,
1484            effective_permissions: None,
1485        };
1486
1487        let json = serde_json::to_string(&receipt).unwrap();
1488        assert!(json.contains("\"result\":\"partial\""));
1489        assert!(json.contains("\"failure_kind\":\"verifier\""));
1490        let back: FleetReceipt = serde_json::from_str(&json).unwrap();
1491        assert_eq!(back.result, FleetTaskResult::Partial);
1492        assert_eq!(back.failure_kind, Some(FleetTaskFailureKind::Verifier));
1493    }
1494
1495    #[test]
1496    fn ssh_host_spec_with_key_pinning_round_trip() {
1497        let spec = FleetHostSpec::Ssh {
1498            host: "builder.trusted.example.com".to_string(),
1499            port: Some(22),
1500            user: Some("codewhale".to_string()),
1501            identity: Some(PathBuf::from("~/.ssh/codewhale_fleet")),
1502            known_hosts: Some(PathBuf::from("~/.ssh/known_hosts")),
1503            host_key_fingerprint: Some("SHA256:aLGqZo1M6c...".to_string()),
1504            working_directory: Some(PathBuf::from("/srv/codewhale/work")),
1505            env_allowlist: vec!["CODEWHALE_PROFILE".to_string()],
1506            codewhale_binary: Some("/usr/local/bin/codewhale".to_string()),
1507        };
1508        let json = serde_json::to_string_pretty(&spec).unwrap();
1509        assert!(json.contains("\"known_hosts\""));
1510        assert!(json.contains("\"host_key_fingerprint\""));
1511        assert!(json.contains("SHA256:aLGqZo1M6c..."));
1512
1513        let back: FleetHostSpec = serde_json::from_str(&json).unwrap();
1514        match back {
1515            FleetHostSpec::Ssh {
1516                host,
1517                known_hosts,
1518                host_key_fingerprint,
1519                ..
1520            } => {
1521                assert_eq!(host, "builder.trusted.example.com");
1522                assert_eq!(known_hosts, Some(PathBuf::from("~/.ssh/known_hosts")));
1523                assert_eq!(
1524                    host_key_fingerprint,
1525                    Some("SHA256:aLGqZo1M6c...".to_string())
1526                );
1527            }
1528            other => panic!("expected ssh host spec, got {other:?}"),
1529        }
1530    }
1531
1532    #[test]
1533    fn secret_ref_redacted_never_exposes_value() {
1534        let ref_ = FleetSecretRef::new("DEEPSEEK_API_KEY");
1535        let redacted = ref_.redacted();
1536        assert!(redacted.contains("DEEPSEEK_API_KEY"));
1537        assert!(!redacted.contains("sk-"));
1538        assert!(redacted.contains("<secret:"));
1539
1540        let ref_ = FleetSecretRef::with_source("GH_TOKEN", "env");
1541        let redacted = ref_.redacted();
1542        assert!(redacted.contains("env.GH_TOKEN"));
1543        assert!(!redacted.contains("ghp_"));
1544    }
1545
1546    #[test]
1547    fn alert_endpoint_from_secret_round_trip() {
1548        let endpoint = FleetAlertEndpoint::from_secret(FleetSecretRef::new("SLACK_WEBHOOK"));
1549        let json = serde_json::to_string(&endpoint).unwrap();
1550        assert!(json.contains("SLACK_WEBHOOK"));
1551        assert!(!json.contains("hooks.slack.com"));
1552
1553        let back: FleetAlertEndpoint = serde_json::from_str(&json).unwrap();
1554        assert_eq!(back.url_ref.as_ref().unwrap().key, "SLACK_WEBHOOK");
1555        assert_eq!(back.url, None);
1556    }
1557
1558    #[test]
1559    fn secret_ref_accepts_legacy_string_wire_shape() {
1560        let ref_: FleetSecretRef = serde_json::from_str(r#""CODEWHALE_FLEET_TOKEN""#).unwrap();
1561        assert_eq!(ref_, FleetSecretRef::new("CODEWHALE_FLEET_TOKEN"));
1562
1563        let ref_: FleetSecretRef =
1564            serde_json::from_str(r#"{"key":"GH_TOKEN","source":"env"}"#).unwrap();
1565        assert_eq!(ref_, FleetSecretRef::with_source("GH_TOKEN", "env"));
1566    }
1567
1568    #[test]
1569    fn trust_level_accepts_hyphenated_remote_verified() {
1570        let trust: FleetTrustLevel = serde_json::from_str(r#""remote-verified""#).unwrap();
1571        assert_eq!(trust, FleetTrustLevel::RemoteVerified);
1572
1573        let canonical = serde_json::to_string(&trust).unwrap();
1574        assert_eq!(canonical, r#""remote_verified""#);
1575    }
1576
1577    #[test]
1578    fn alert_channel_accepts_legacy_webhook_fields() {
1579        let channel: FleetAlertChannel = serde_json::from_str(
1580            r#"{
1581                "kind": "slack",
1582                "webhook_url": "https://hooks.slack.com/test",
1583                "secret": "SLACK_SIGNING_SECRET"
1584            }"#,
1585        )
1586        .unwrap();
1587
1588        match channel {
1589            FleetAlertChannel::Slack { webhook } => {
1590                assert_eq!(webhook.url.as_deref(), Some("https://hooks.slack.com/test"));
1591                assert_eq!(
1592                    webhook.secret_ref,
1593                    Some(FleetSecretRef::new("SLACK_SIGNING_SECRET"))
1594                );
1595            }
1596            other => panic!("expected slack channel, got {other:?}"),
1597        }
1598    }
1599
1600    #[test]
1601    fn security_policy_defaults_are_conservative() {
1602        let policy = FleetSecurityPolicy::default();
1603        assert_eq!(policy.default_trust_level, FleetTrustLevel::Sandbox);
1604        assert!(policy.allowed_secrets.is_empty());
1605        assert!(policy.capability_grants.is_empty());
1606        assert_eq!(policy.max_trust_level, FleetTrustLevel::Operator);
1607        assert!(!policy.require_identity_verification);
1608    }
1609
1610    #[test]
1611    fn trust_level_ordinal_reflects_privilege() {
1612        assert!(FleetTrustLevel::Operator > FleetTrustLevel::RemoteVerified);
1613        assert!(FleetTrustLevel::RemoteVerified > FleetTrustLevel::Local);
1614        assert!(FleetTrustLevel::Local > FleetTrustLevel::Sandbox);
1615
1616        assert!(FleetTrustLevel::Operator.may_access_secrets());
1617        assert!(!FleetTrustLevel::Sandbox.may_access_secrets());
1618        assert!(!FleetTrustLevel::Sandbox.may_write_workspace());
1619        assert!(FleetTrustLevel::Operator.may_write_workspace());
1620    }
1621
1622    fn sample_receipt_with_route() -> FleetReceipt {
1623        FleetReceipt {
1624            run_id: FleetRunId::from("run-route"),
1625            task_id: "task-route".to_string(),
1626            worker_id: "worker-route".to_string(),
1627            attempt: Some(1),
1628            terminal_seq: Some(4),
1629            completed_at: "2026-06-23T00:00:00Z".to_string(),
1630            result: FleetTaskResult::Pass,
1631            failure_kind: None,
1632            artifacts: vec![],
1633            score: None,
1634            resolved_route: Some(FleetResolvedRoute {
1635                provider_id: "deepseek".to_string(),
1636                provider_exact_id: None,
1637                provider_kind: "deepseek".to_string(),
1638                canonical_model: Some("deepseek-v4-pro".to_string()),
1639                wire_model_id: "deepseek-v4-pro".to_string(),
1640                protocol: "chat_completions".to_string(),
1641                role: Some("builder".to_string()),
1642                loadout: Some("auto".to_string()),
1643                model_class: Some("balanced".to_string()),
1644                model_route: Some("auto".to_string()),
1645                reasoning_effort: Some("high".to_string()),
1646                role_source: Some("task.role".to_string()),
1647                loadout_source: Some("task.loadout".to_string()),
1648                model_class_source: Some("task.model_class".to_string()),
1649                model_source: Some("task.model".to_string()),
1650                source: "resolver".to_string(),
1651            }),
1652            effective_permissions: Some(FleetEffectivePermissions {
1653                write: true,
1654                network: true,
1655                shell: "full".to_string(),
1656                tool_scope: "explicit".to_string(),
1657                tools: vec!["read_file".to_string(), "apply_patch".to_string()],
1658                background: true,
1659                max_spawn_depth: 2,
1660                profile_id: Some("builder".to_string()),
1661                profile_origin: Some("built_in".to_string()),
1662                source: "worker_runtime_profile".to_string(),
1663            }),
1664        }
1665    }
1666
1667    #[test]
1668    fn fleet_resolved_route_round_trips() {
1669        let receipt = sample_receipt_with_route();
1670        let json = serde_json::to_string(&receipt).unwrap();
1671        let back: FleetReceipt = serde_json::from_str(&json).unwrap();
1672        assert_eq!(back.resolved_route, receipt.resolved_route);
1673        assert_eq!(back.effective_permissions, receipt.effective_permissions);
1674        let route = back.resolved_route.unwrap();
1675        assert_eq!(route.provider_id, "deepseek");
1676        assert_eq!(route.wire_model_id, "deepseek-v4-pro");
1677        assert_eq!(route.protocol, "chat_completions");
1678        assert_eq!(route.role.as_deref(), Some("builder"));
1679        assert_eq!(route.loadout.as_deref(), Some("auto"));
1680        assert_eq!(route.model_class.as_deref(), Some("balanced"));
1681        assert_eq!(route.model_route.as_deref(), Some("auto"));
1682        assert_eq!(route.reasoning_effort.as_deref(), Some("high"));
1683        assert_eq!(route.role_source.as_deref(), Some("task.role"));
1684        assert_eq!(route.loadout_source.as_deref(), Some("task.loadout"));
1685        assert_eq!(
1686            route.model_class_source.as_deref(),
1687            Some("task.model_class")
1688        );
1689        assert_eq!(route.model_source.as_deref(), Some("task.model"));
1690        assert_eq!(route.source, "resolver");
1691
1692        let permissions = back
1693            .effective_permissions
1694            .expect("effective permissions should round-trip");
1695        assert!(permissions.write);
1696        assert!(permissions.network);
1697        assert_eq!(permissions.shell, "full");
1698        assert_eq!(permissions.tool_scope, "explicit");
1699        assert_eq!(
1700            permissions.tools,
1701            vec!["read_file".to_string(), "apply_patch".to_string()]
1702        );
1703        assert!(permissions.background);
1704        assert_eq!(permissions.max_spawn_depth, 2);
1705        assert_eq!(permissions.profile_id.as_deref(), Some("builder"));
1706        assert_eq!(permissions.profile_origin.as_deref(), Some("built_in"));
1707        assert_eq!(permissions.source, "worker_runtime_profile");
1708    }
1709
1710    #[test]
1711    fn fleet_receipt_without_resolved_route_still_deserializes() {
1712        // An old ledger receipt JSON written before #3154 has no
1713        // `resolved_route` key; `#[serde(default)]` must keep it readable.
1714        let legacy = r#"{
1715            "run_id": "run-legacy",
1716            "task_id": "task-legacy",
1717            "worker_id": "worker-legacy",
1718            "completed_at": "2026-06-01T00:00:00Z",
1719            "result": "pass",
1720            "artifacts": [],
1721            "score": null
1722        }"#;
1723        let receipt: FleetReceipt = serde_json::from_str(legacy).unwrap();
1724        assert_eq!(receipt.task_id, "task-legacy");
1725        assert!(receipt.resolved_route.is_none());
1726        assert!(receipt.attempt.is_none());
1727        assert!(receipt.terminal_seq.is_none());
1728    }
1729
1730    #[test]
1731    fn fleet_resolved_route_legacy_shape_still_deserializes() {
1732        let legacy = r#"{
1733            "run_id": "run-route",
1734            "task_id": "task-route",
1735            "worker_id": "worker-route",
1736            "completed_at": "2026-06-23T00:00:00Z",
1737            "result": "pass",
1738            "artifacts": [],
1739            "score": null,
1740            "resolved_route": {
1741                "provider_id": "deepseek",
1742                "provider_kind": "deepseek",
1743                "canonical_model": "deepseek-v4-pro",
1744                "wire_model_id": "deepseek-v4-pro",
1745                "protocol": "chat_completions",
1746                "role": "builder",
1747                "loadout": "fast",
1748                "source": "resolver"
1749            }
1750        }"#;
1751
1752        let receipt: FleetReceipt = serde_json::from_str(legacy).unwrap();
1753        let route = receipt.resolved_route.expect("legacy route should parse");
1754        assert_eq!(route.source, "resolver");
1755        assert_eq!(route.role.as_deref(), Some("builder"));
1756        assert_eq!(route.loadout.as_deref(), Some("fast"));
1757        assert_eq!(route.model_class, None);
1758        assert_eq!(route.model_route, None);
1759        assert_eq!(route.reasoning_effort, None);
1760        assert_eq!(route.role_source, None);
1761        assert_eq!(route.loadout_source, None);
1762        assert_eq!(route.model_class_source, None);
1763        assert_eq!(route.model_source, None);
1764    }
1765
1766    #[test]
1767    fn fleet_resolved_route_serialization_carries_no_secrets() {
1768        let receipt = sample_receipt_with_route();
1769        // Scan the serialized resolved-route object: this is the field whose
1770        // no-secrets invariant we are asserting. Scoping to the route value
1771        // avoids false positives from unrelated envelope ids (e.g. a task id
1772        // such as "task-foo" innocently contains the substring "sk-").
1773        let route_json = serde_json::to_string(receipt.resolved_route.as_ref().unwrap()).unwrap();
1774        assert_no_secret_markers(&route_json);
1775        // The envelope as a whole must also stay credential-free.
1776        let receipt_json = serde_json::to_string(&receipt).unwrap();
1777        for needle in SECRET_KEY_MARKERS {
1778            assert!(
1779                !receipt_json.to_ascii_lowercase().contains(needle),
1780                "receipt JSON must not contain secret-key marker {needle:?}: {receipt_json}"
1781            );
1782        }
1783    }
1784
1785    /// Substrings that indicate a leaked credential field/value. These are
1786    /// deliberately specific so legitimate ids/model names do not trip them.
1787    const SECRET_KEY_MARKERS: &[&str] = &[
1788        "api_key",
1789        "apikey",
1790        "api-key",
1791        "authorization",
1792        "bearer ",
1793        "auth_token",
1794        "auth-token",
1795        "password",
1796        "credential",
1797        "sk-ant-",
1798        "sk-proj-",
1799        "sk-or-",
1800        "secret",
1801    ];
1802
1803    fn assert_no_secret_markers(json: &str) {
1804        let haystack = json.to_ascii_lowercase();
1805        for needle in SECRET_KEY_MARKERS {
1806            assert!(
1807                !haystack.contains(needle),
1808                "resolved-route JSON must not contain secret marker {needle:?}: {json}"
1809            );
1810        }
1811    }
1812}