task_runs/types.rs
1//! Core types for the TaskRun system.
2//!
3//! These mirror the data model in `.yah/docs/working/yah-task-runs.md`.
4//! Intentionally kept free of I/O — the store layer owns persistence.
5//!
6//! The types `TaskRunId`, `Level`, `EventSource`, `ChunkRef`, `Event`,
7//! `Diagnostic`, `Initiator`, and `RunStatus` live in
8//! `crates/yah/observation/` and are re-exported here for backward
9//! compatibility.
10
11use serde::{Deserialize, Serialize};
12use std::path::PathBuf;
13
14// Re-export the hoisted observation types so all existing callers continue to
15// work via `use task_runs::{TaskRunId, Event, ...}`.
16pub use observation::{
17 ChunkRef, Diagnostic, Event, EventScope, EventSource, ForgeId, Initiator, Level, RunStatus,
18 TaskRunId, RESERVED_FIELD_PATHS,
19};
20
21// ─── BeholderStatus ───────────────────────────────────────────────────────────
22
23/// Opaque string surfaced on `TaskRunMeta.beholder_status`.
24///
25/// Examples: `"attached:cargo@1.78"`, `"none:auto"`, `"declined:cargo
26/// reason=\"explicit --message-format=human\""`, `"unknown_format"`,
27/// `"attached:cargo@1.78 rewrite=\"--message-format=json-render-diagnostics\""`.
28#[derive(Debug, Clone, Serialize, Deserialize)]
29pub struct BeholderStatus {
30 pub text: String,
31 /// Args added to argv by a `Rewriter` beholder. `None` for `Parser` mode or
32 /// when no beholder attached.
33 #[serde(skip_serializing_if = "Option::is_none")]
34 pub rewrite_added: Option<Vec<String>>,
35}
36
37impl BeholderStatus {
38 fn make(text: String) -> Self {
39 Self { text, rewrite_added: None }
40 }
41
42 pub fn none_auto() -> Self {
43 Self::make("none:auto".to_string())
44 }
45 /// Bytes-only because the caller explicitly set `BeholderSelect::None`.
46 pub fn none_explicit() -> Self {
47 Self::make("none:explicit".to_string())
48 }
49 pub fn attached(name: &str, version: &str) -> Self {
50 Self::make(format!("attached:{name}@{version}"))
51 }
52 pub fn declined(name: &str, reason: &str) -> Self {
53 Self::make(format!("declined:{name} reason=\"{reason}\""))
54 }
55 /// `BeholderSelect::Force` matched; `matches` predicate agreed.
56 pub fn forced(name: &str, version: &str) -> Self {
57 Self::make(format!("forced:{name}@{version}"))
58 }
59 /// `BeholderSelect::Force` matched; beholder's `matches` would have declined.
60 pub fn forced_against_flags(name: &str, version: &str) -> Self {
61 Self::make(format!("forced-against-flags:{name}@{version}"))
62 }
63 /// `BeholderSelect::Force` matched on a run whose output is consumed
64 /// verbatim, where a Rewriter beholder would normally decline rather than
65 /// change what the caller reads.
66 pub fn forced_against_verbatim(name: &str, version: &str) -> Self {
67 Self::make(format!("forced-against-verbatim:{name}@{version}"))
68 }
69 pub fn unknown_format() -> Self {
70 Self::make("unknown_format".to_string())
71 }
72 /// Beholder detected that the tool's output format is unrecognized (schema
73 /// drift). Records which beholder made the call and why.
74 pub fn unknown_format_with_reason(name: &str, reason: &str) -> Self {
75 Self::make(format!("unknown_format:{name} reason=\"{reason}\""))
76 }
77
78 /// Append rewrite info to this status if `added` is non-empty.
79 ///
80 /// Called after a `Rewriter` beholder adjusts argv so agents can see exactly
81 /// what was injected into the command line.
82 pub fn with_rewrite(mut self, added: Vec<String>) -> Self {
83 if !added.is_empty() {
84 let repr = added.join(" ");
85 self.text = format!("{} rewrite=\"{repr}\"", self.text);
86 self.rewrite_added = Some(added);
87 }
88 self
89 }
90}
91
92// ─── TaskRunMeta ──────────────────────────────────────────────────────────────
93
94#[derive(Debug, Clone, Serialize, Deserialize)]
95pub struct TaskRunMeta {
96 pub id: TaskRunId,
97 pub command: String,
98 pub cwd: PathBuf,
99 pub env: Vec<(String, String)>,
100 pub started_at: u64,
101 pub status: RunStatus,
102 pub label: Option<String>,
103 pub initiator: Initiator,
104 pub beholder_status: Option<BeholderStatus>,
105 /// If true, the GC sweep will not drop this run's output during warm rolloff.
106 /// Pinned runs are exempt until explicitly unpinned or archived.
107 #[serde(default)]
108 pub pinned: bool,
109 /// What surface spawned this run. `None` (the default) is an ordinary
110 /// `task.run` job; `Some("terminal")` marks an interactive terminal
111 /// session (SSH / local PTY / camp shell). A generic provenance tag, not
112 /// a UI concept — it lets a consumer (e.g. the desktop terminal rail) list
113 /// just its own runs from `task.list` without scooping up unrelated jobs.
114 #[serde(default)]
115 pub origin: Option<String>,
116 /// PID of the process whose [`crate::TaskDriver`] spawned this run — the
117 /// run's **owner**, not the child.
118 ///
119 /// Recorded so a driver starting up in one process can tell a genuinely
120 /// abandoned run from one a live peer process is still driving. Without
121 /// it, "leftover `Running` rows are stale" is only true when exactly one
122 /// process ever writes the store, and the moment a second one attaches it
123 /// tombstones the first one's live runs. See
124 /// [`crate::driver::StaleRunPolicy`].
125 ///
126 /// `None` for rows written before the column existed, which the policy
127 /// reads as "owner unknown" and treats conservatively (tombstone).
128 #[serde(default)]
129 pub host_pid: Option<u32>,
130}
131
132// ─── Stream ───────────────────────────────────────────────────────────────────
133
134#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
135#[serde(rename_all = "lowercase")]
136pub enum Stream {
137 Stdout,
138 Stderr,
139 Synth,
140}
141
142impl Stream {
143 pub fn as_str(self) -> &'static str {
144 match self {
145 Stream::Stdout => "stdout",
146 Stream::Stderr => "stderr",
147 Stream::Synth => "synth",
148 }
149 }
150}
151
152impl std::str::FromStr for Stream {
153 type Err = String;
154 fn from_str(s: &str) -> Result<Self, Self::Err> {
155 match s {
156 "stdout" => Ok(Stream::Stdout),
157 "stderr" => Ok(Stream::Stderr),
158 "synth" => Ok(Stream::Synth),
159 other => Err(format!("unknown stream: {other}")),
160 }
161 }
162}
163
164// ─── OutputChunk ──────────────────────────────────────────────────────────────
165
166#[derive(Debug, Clone, Serialize, Deserialize)]
167pub struct OutputChunk {
168 pub run_id: TaskRunId,
169 pub seq: u32,
170 pub offset_ms: u32,
171 pub stream: Stream,
172 pub bytes: Vec<u8>,
173}
174
175// ─── SeqRange ─────────────────────────────────────────────────────────────────
176
177/// Inclusive range over chunk `seq` numbers within a single run.
178#[derive(Debug, Clone, Serialize, Deserialize)]
179pub struct SeqRange {
180 pub lo: u32,
181 pub hi: u32,
182}
183
184// ─── Triage (Tier 1.75 — pruner output) ───────────────────────────────────────
185
186/// Pruner output for a run: a list of verbatim chunk ranges + a human-facing
187/// synopsis.
188///
189/// **Agents must read `keep`/`primary` ranges and resolve them to bytes via
190/// `task.lines`. The `synopsis` is for human display only — it can paraphrase
191/// or hallucinate. Never parse it programmatically.**
192#[derive(Debug, Clone, Serialize, Deserialize)]
193pub struct Triage {
194 pub run_id: TaskRunId,
195 pub synopsis: String,
196 pub keep: Vec<KeepRange>,
197 pub primary: SeqRange,
198 pub model: String,
199 pub prompt_version: u32,
200 pub cached_at: u64,
201 pub partial: bool,
202}
203
204#[derive(Debug, Clone, Serialize, Deserialize)]
205pub struct KeepRange {
206 pub range: SeqRange,
207 pub reason: String,
208}