#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
#[non_exhaustive]
pub enum DriverLiveness {
Driving,
DriverDead,
Parked,
Undriven,
}
use std::path::{Path, PathBuf};
use crate::event::{Envelope, Source};
use crate::filter::EventFilter;
use crate::graph::{self, Landing, NodeStatus};
use crate::journal::PipelineKind;
use crate::ledger::{self, LaunchRecord, LockRecord, RunPaths};
use crate::projection::{self, MemberLabel, Refusal, RunState};
use crate::sys;
pub use crate::ledger::Skipped;
pub const DEFAULT_PARKED_AFTER_SECONDS: u64 = 1_800;
pub const PARKED_AFTER_ENV: &str = "ONEPIPELINE_PARKED_AFTER_SECONDS";
const ENDED_BY_THE_STOP: &str = "worker ended when the run was stopped";
const OUTLIVED_THE_STOP: &str = "worker may still be running: the stop could not reach it";
fn became_of_the_worker(state: &crate::projection::RunState) -> &'static str {
match state.stop {
crate::projection::StopState::WorkersUndetermined => OUTLIVED_THE_STOP,
_ => ENDED_BY_THE_STOP,
}
}
impl DriverLiveness {
pub fn as_str(self) -> &'static str {
match self {
Self::Driving => "ACTIVE",
Self::DriverDead => "DRIVER DEAD",
Self::Parked => "PARKED",
Self::Undriven => "UNDRIVEN",
}
}
pub fn is_undriven(self) -> bool {
matches!(self, Self::DriverDead | Self::Parked)
}
}
pub fn parked_after_seconds() -> u64 {
std::env::var(PARKED_AFTER_ENV)
.ok()
.and_then(|value| value.parse().ok())
.filter(|seconds| *seconds > 0)
.unwrap_or(DEFAULT_PARKED_AFTER_SECONDS)
}
pub fn decision_outstanding(state: &RunState, paths: &RunPaths) -> bool {
state.awaiting_human_action() || blocking_surface(paths)
}
fn blocking_surface(paths: &RunPaths) -> bool {
let queue = crate::channel::ChannelState::new(paths).queue();
queue
.waiting
.iter()
.chain(queue.pending.iter())
.any(|surface| surface.blocking)
}
pub fn liveness(launch: &LaunchRecord, state: &RunState, paths: &RunPaths) -> DriverLiveness {
if state.stop_recorded() {
return DriverLiveness::DriverDead;
}
let ours = launch.host == sys::hostname();
if ours && !sys::process_may_be_live(launch.pid) {
return DriverLiveness::DriverDead;
}
let quiet_for = state
.last_write_at
.map(|last| sys::now_millis().saturating_sub(last) / 1_000);
match quiet_for {
Some(seconds)
if seconds > parked_after_seconds() && !decision_outstanding(state, paths) =>
{
DriverLiveness::Parked
}
_ => DriverLiveness::Driving,
}
}
#[derive(Debug)]
pub struct RunView {
pub paths: RunPaths,
pub launch: LaunchRecord,
pub events: Vec<Envelope>,
pub state: RunState,
}
impl RunView {
pub fn open(paths: &RunPaths) -> crate::Result<Self> {
if !paths.exists() {
return Err(crate::Error::NoSuchRun {
run: paths.run.clone(),
root: paths.dir.parent().unwrap_or(Path::new(".")).to_path_buf(),
});
}
let launch: LaunchRecord = ledger::read_json(&paths.launch())?;
let mut events = crate::journal::read(&paths.journal());
crate::journal::merge_order(&mut events);
let mut state = projection::fold(&events);
state.cross_dag = crate::crossdag::resolve_quietly(
&paths
.dir
.parent()
.map_or_else(ledger::runs_root, Path::to_path_buf),
&state.graph,
);
Ok(Self {
paths: paths.clone(),
launch,
events,
state,
})
}
pub fn liveness(&self) -> DriverLiveness {
liveness(&self.launch, &self.state, &self.paths)
}
pub fn unread_surfaces(&self) -> (usize, Option<u64>) {
let queue = crate::channel::ChannelState::new(&self.paths).queue();
let oldest = queue
.waiting
.iter()
.map(|surface| sys::now_millis().saturating_sub(surface.queued_at) / 1_000)
.max();
(queue.waiting.len(), oldest)
}
pub fn summary(&self) -> String {
let statuses = self.state.statuses();
let done = statuses
.values()
.filter(|status| **status == NodeStatus::Done)
.count();
let unlanded = match unlanded_nodes(self).len() {
0 => String::new(),
count => format!(", {count} not landed as of settlement"),
};
format!("{done}/{} done{unlanded}", statuses.len())
}
}
#[derive(Debug)]
pub struct Survey {
pub root: PathBuf,
pub views: Vec<RunView>,
pub skipped: Vec<Skipped>,
}
impl Survey {
pub fn of(root: &Path) -> Self {
let index = ledger::all_runs(root);
let mut views = Vec::new();
let mut skipped = index.skipped;
for paths in index.runs {
match RunView::open(&paths) {
Ok(view) => views.push(view),
Err(error) => skipped.push(Skipped {
path: paths.dir,
reason: error.to_string(),
}),
}
}
skipped.sort_by(|a, b| a.path.cmp(&b.path));
Self {
root: root.to_path_buf(),
views,
skipped,
}
}
pub fn of_one(view: RunView) -> Self {
let root = view
.paths
.dir
.parent()
.map_or_else(ledger::runs_root, Path::to_path_buf);
Self {
root,
views: vec![view],
skipped: Vec::new(),
}
}
}
fn skipped_lines(skipped: &[Skipped]) -> String {
if skipped.is_empty() {
return String::new();
}
let mut out = format!("{} run root(s) skipped:\n", skipped.len());
for root in skipped {
out.push_str(&format!(
" {}: {}\n",
one_line(&root.path.display().to_string()),
one_line(&root.reason)
));
}
out
}
fn nothing_to_report(survey: &Survey) -> String {
let mut out = if survey.views.is_empty() && !survey.skipped.is_empty() {
format!(
"no run under {} could be read\n",
one_line(&survey.root.display().to_string())
)
} else {
"no runs recorded\n".to_string()
};
out.push_str(&skipped_lines(&survey.skipped));
out
}
pub fn liveness_word(view: &RunView) -> &'static str {
let statuses = view.state.statuses();
if !statuses.is_empty() && graph::state_of(&statuses) == graph::GraphState::Complete {
return "SETTLED";
}
view.liveness().as_str()
}
pub fn runs(root: &Path, mine_only: bool, session: &str) -> String {
let survey = Survey::of(root);
let mut out = String::new();
for view in &survey.views {
let owned = view.launch.owned_by(session);
if mine_only && !owned {
continue;
}
let marker = if owned { '*' } else { ' ' };
out.push_str(&format!(
"{marker} {:<24} {:<24} {} {}\n",
view.paths.run,
view.launch.owner_label(session),
view.summary(),
liveness_word(view)
));
if view.liveness().is_undriven() {
out.push_str(&format!(
" {} — its ledger is intact; attach a fresh driver with: \
onepipeline adopt {}\n",
view.liveness().as_str(),
view.paths.run
));
continue;
}
if let (count, Some(stale)) = view.unread_surfaces() {
if count > 0 {
out.push_str(&format!(
" {count} planner update(s) waiting, unread for {}; \
read them with: onepipeline next {}\n",
crate::telemetry::duration(stale * 1_000),
view.paths.run
));
}
}
}
if out.is_empty() {
return nothing_to_report(&survey);
}
out.push_str(&skipped_lines(&survey.skipped));
out
}
pub fn status(survey: &Survey) -> String {
let mut out = String::new();
for view in &survey.views {
out.push_str(&format!(
"{} {} {}\n",
view.paths.run,
liveness_word(view),
view.summary()
));
if view.liveness().is_undriven() {
out.push_str(&format!(
" {}: nothing is driving this run; adopt it or stop it\n",
view.liveness().as_str()
));
}
if let Some(pending) = crate::channel::ChannelState::new(&view.paths).pending() {
out.push_str(&format!(
" waiting for planner {}: {} — {}\n",
if pending.blocking {
"decision"
} else {
"reply"
},
pending.kind,
pending.message
));
}
let (unread, stale) = view.unread_surfaces();
if unread > 0 {
out.push_str(&format!(
" {unread} planner update(s) waiting, unread for {}\n",
crate::telemetry::duration(stale.unwrap_or(0) * 1_000)
));
}
let statuses = view.state.statuses();
for (id, node_status) in &statuses {
if *node_status != NodeStatus::Running {
continue;
}
let age = view
.state
.dispatched_at
.get(id)
.map(|at| sys::now_millis().saturating_sub(*at));
let age = crate::telemetry::duration(age.unwrap_or(0));
if view.state.stop_recorded() {
let became = became_of_the_worker(&view.state);
out.push_str(&format!(" {id}: {became}, {age} in\n"));
continue;
}
out.push_str(&format!(" {id}: running for {age}"));
match view.state.activity.get(id) {
None => out.push_str(&format!(" — {}", DriverLiveness::Undriven.as_str())),
Some(activity) => out.push_str(&format!(" — {}", working(activity))),
}
out.push('\n');
}
for (id, node_status) in &statuses {
if *node_status != NodeStatus::Parked {
continue;
}
let Some(pending) = cancelling_for(&view.state, id) else {
continue;
};
out.push_str(&format!(
" {id}: cancelling — asked to stop {pending} ago and its dispatch has not \
settled; it still holds the node's workspace, so wait for it rather than \
requeueing the node\n"
));
}
for (id, node_status) in &statuses {
if *node_status != NodeStatus::Ready {
continue;
}
out.push_str(&format!(
" {id}: ready — {}\n",
waiting_on(&view.state, id)
));
}
for (id, node_status) in &statuses {
if *node_status != NodeStatus::Failed {
continue;
}
for refusal in refusals_of(&view.state, id) {
out.push_str(&format!(" {id}: failed — {}\n", refusal_phrase(refusal)));
}
}
let unlanded = unlanded_nodes(view);
if !unlanded.is_empty() {
out.push_str(&format!(
" {} node(s) settled without landing: {} — as each settled, not as of now; \
`results {}` names the change to open\n",
unlanded.len(),
unlanded.join(", "),
view.paths.run
));
}
if let Some(health) = crate::agentgraph::health() {
out.push_str(&format!(" providers: {health}\n"));
}
}
if out.is_empty() {
return nothing_to_report(survey);
}
out.push_str(&skipped_lines(&survey.skipped));
out
}
fn cancelling_for(state: &RunState, id: &str) -> Option<String> {
let since = state.recorded.get(id)?.cancelling_since()?;
Some(crate::telemetry::duration(
sys::now_millis().saturating_sub(since),
))
}
fn refusals_of<'a>(state: &'a RunState, node: &str) -> &'a [Refusal] {
state.refusals.get(node).map_or(&[], Vec::as_slice)
}
fn refusal_phrase(refusal: &Refusal) -> String {
let role = refusal
.advanced
.role
.and_then(|role| serde_json::to_value(role).ok());
let side = match (
role.as_ref().and_then(serde_json::Value::as_str),
&refusal.member,
) {
(Some(role), _) => format!("the {role} side"),
(None, MemberLabel::Named(member)) => format!("member '{member}'"),
(None, MemberLabel::Unstamped) => "a side the record does not name".to_string(),
(None, MemberLabel::Unreadable) => "a side this build cannot read".to_string(),
};
let reason = if refusal.advanced.reason.is_empty() {
"for a reason the record does not carry".to_string()
} else {
format!("({})", refusal.advanced.reason)
};
let again = if refusal.records.get() > 1 {
format!(", recorded {} times", refusal.records)
} else {
String::new()
};
one_line(&format!(
"{side}: identity '{}' refused {reason}{again}",
refusal.advanced.identity
))
}
fn unlanded_nodes(view: &RunView) -> Vec<String> {
view.state
.landings
.iter()
.filter(|(_, landing)| **landing == Landing::Unlanded)
.map(|(node, _)| node.clone())
.collect()
}
fn waiting_on(state: &RunState, id: &str) -> String {
const QUEUED: &str = "queued for dispatch";
let Some(repo) = state.graph.get(id).and_then(|node| node.repo.as_deref()) else {
return QUEUED.to_string();
};
let holders = match crate::vcs::holders_of(repo) {
Ok(holders) => holders,
Err(why) => {
return format!(
"{QUEUED}, and this host cannot say whether the '{repo}' workspace is \
free: {}",
one_line(&why)
)
}
};
let held: Vec<String> = holders
.into_iter()
.filter(|holder| {
holder.state == onevcs::Lifecycle::Open && holder.liveness == onevcs::Liveness::Live
})
.map(|holder| {
format!(
"session '{}' (owner_pid {})",
holder.token.0, holder.owner_pid
)
})
.collect();
if held.is_empty() {
return QUEUED.to_string();
}
format!(
"waiting for the '{repo}' workspace, held by {}",
held.join(", ")
)
}
fn working(activity: &crate::projection::NodeActivity) -> String {
let ago = |at: u64| crate::telemetry::duration(sys::now_millis().saturating_sub(at));
let alive = activity
.last_heartbeat_at
.filter(|beat| activity.progress.is_none_or(|done| *beat > done.last_at()))
.map(|beat| format!("; alive {} ago", ago(beat)))
.unwrap_or_default();
let Some(progress) = activity.progress else {
return format!("nothing recorded yet{alive}");
};
let counted = format!(
"{} event(s), {} ago{alive}",
progress.events(),
ago(progress.last_at())
);
match &activity.doing {
Some(doing) => format!("now {doing} ({counted})"),
None => counted,
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum Proof {
Held,
Stale(String),
Unproven(String),
}
fn dispatch_proof(view: &RunView) -> Proof {
if view.state.stop_recorded() {
return Proof::Stale("the run was stopped".to_string());
}
let path = view.paths.lock();
match std::fs::metadata(&path) {
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
return Proof::Stale(
"nothing holds the run's ownership lock, so no driver is running it".to_string(),
)
}
Err(error) => {
return Proof::Unproven(format!(
"the run's ownership lock cannot be described: {error}"
))
}
Ok(about) if !about.is_file() => {
return Proof::Unproven(
"the run's ownership lock is not a file, so nothing here holds it".to_string(),
)
}
Ok(_) => {}
}
let Some(held) = ledger::read_json_opt::<LockRecord>(&path) else {
return Proof::Unproven("the run's ownership lock cannot be read".to_string());
};
if held.host != sys::hostname() {
return Proof::Unproven(format!(
"its driver holds the lock on {}, and a pid means nothing across machines",
held.host
));
}
if !sys::process_may_be_live(held.pid) {
return Proof::Stale(format!("its driver (pid {}) is gone", held.pid));
}
if held.started.is_empty() {
return Proof::Unproven(format!(
"the run's lock carries no start token for pid {}, so nothing says it is still \
the process that took it",
held.pid
));
}
match sys::process_start_token(held.pid) {
None => Proof::Unproven(format!(
"this host will not say when pid {} started",
held.pid
)),
Some(token) if token.matches(&held.started) => Proof::Held,
Some(_) => Proof::Stale(format!(
"pid {} is a different process from the one that took the run's lock",
held.pid
)),
}
}
pub fn host(survey: &Survey) -> String {
let mut out = format!("host {}\n", sys::hostname());
out.push_str(&format!(
" reading {}\n",
one_line(&survey.root.display().to_string())
));
let mut rendered = false;
let mut ignored: Vec<String> = Vec::new();
for view in &survey.views {
let proof = dispatch_proof(view);
for (id, status) in &view.state.statuses() {
if *status != NodeStatus::Running {
continue;
}
if let Proof::Stale(why) = &proof {
ignored.push(one_line(&format!("{}/{id}: {why}", view.paths.run)));
continue;
}
let age = view
.state
.dispatched_at
.get(id)
.map(|at| sys::now_millis().saturating_sub(*at))
.unwrap_or(0);
rendered = true;
out.push_str(&format!(
" {:<24} {:<20} {:<16} {}",
view.paths.run,
id,
view.launch.launcher,
crate::telemetry::duration(age)
));
if let Proof::Unproven(why) = &proof {
out.push_str(&format!(" UNPROVEN: {}", one_line(why)));
}
out.push('\n');
}
}
if !rendered {
out.push_str(" no live dispatches\n");
}
if !ignored.is_empty() {
out.push_str(&format!(
" {} stale registry entr{} ignored: {}\n",
ignored.len(),
if ignored.len() == 1 { "y" } else { "ies" },
ignored.join("; ")
));
}
out.push_str(&skipped_lines(&survey.skipped));
out
}
pub fn shaped<'a>(view: &'a RunView, filter: &EventFilter) -> Vec<&'a Envelope> {
view.events
.iter()
.filter(|event| filter.matches(event))
.collect()
}
pub fn monitor(view: &RunView, filter: &EventFilter) -> String {
let mut out = String::from(
"Concise graph events; ask the producing library for full detail by stream id.\n",
);
for event in shaped(view, filter) {
let id = match event.source {
Source::Pipeline => format!("graph:{}", event.labels.node.as_deref().unwrap_or("-")),
Source::Agentgraph => format!("agent:{}", event.stream),
Source::Vcs => format!("vcs:{}", event.stream),
};
out.push_str(&format!("{} {:<28} {}\n", event.ts, id, summarize(event)));
}
out.push_str(&format!(
"-- {} {} {} {}\n",
view.paths.run,
view.summary(),
liveness_word(view),
graph::state_of(&view.state.statuses()).as_str()
));
out
}
fn summarize(event: &Envelope) -> String {
const CAP: usize = 96;
let mut detail = event.kind.0.clone();
for key in ["status", "landing", "outcome", "state", "message", "reason"] {
if let Some(value) = event.payload.get(key).and_then(|v| v.as_str()) {
detail.push_str(&format!(" {value}"));
}
}
let stripped: String = detail
.chars()
.map(|c| if c.is_control() { ' ' } else { c })
.collect();
if stripped.chars().count() <= CAP {
return stripped;
}
stripped.chars().take(CAP).collect()
}
fn landed_phrase(landing: Landing, settled_at: Option<u64>) -> String {
let ago = match settled_at {
Some(at) => format!(
" {} ago",
crate::telemetry::duration(sys::now_millis().saturating_sub(at))
),
None => String::new(),
};
match landing {
Landing::Landed => "landed on its base".to_string(),
Landing::Unlanded => format!(
"NOT landed: the change had not reached its base when this settled{ago}, and \
nothing has re-read it since — open the change for where it is now"
),
}
}
pub fn results(view: &RunView) -> String {
let mut out = format!(
"{} {}\n",
view.paths.run,
graph::state_of(&view.state.statuses()).as_str()
);
let statuses = view.state.statuses();
for node in view.state.graph.iter() {
let status = statuses
.get(&node.id)
.copied()
.unwrap_or(NodeStatus::Pending);
out.push_str(&format!(" {:<24} {}", node.id, status.as_str()));
if let Some(outcome) = view.state.outcomes.get(&node.id) {
out.push_str(&format!(" ({outcome})"));
}
if let Some(pending) = cancelling_for(&view.state, &node.id) {
out.push_str(&format!(" — cancelling, asked to stop {pending} ago"));
}
if let Some(landing) = view.state.landings.get(&node.id) {
let settled_at = view.state.settled_at.get(&node.id).copied();
out.push_str(&format!(" — {}", landed_phrase(*landing, settled_at)));
}
if status == NodeStatus::Running && view.state.stop_recorded() {
out.push_str(&format!(" — {}", became_of_the_worker(&view.state)));
}
let branch = view
.state
.branches
.get(&node.id)
.or(node.branch.as_ref())
.cloned();
if let (NodeStatus::Parked | NodeStatus::Failed | NodeStatus::Cancelled, Some(branch)) =
(status, &branch)
{
out.push_str(&format!(" — preserved on {branch}"));
}
if let Some(url) = view.state.change_urls.get(&node.id) {
out.push_str(&format!(" — {url}"));
}
out.push('\n');
if let Some(detail) = view
.events
.iter()
.rev()
.find(|event| {
event.kind.0 == PipelineKind::NodeSettled.as_str()
&& event.labels.node.as_deref() == Some(node.id.as_str())
})
.and_then(|event| event.payload.get("detail"))
.and_then(|detail| detail.as_str())
{
out.push_str(&format!(" detail: {}\n", one_line(detail)));
}
if status == NodeStatus::Failed {
for refusal in refusals_of(&view.state, &node.id) {
out.push_str(&format!(" provider: {}\n", refusal_phrase(refusal)));
}
}
if status == NodeStatus::Waiting {
if let Some(task) = &node.task {
out.push_str(&format!(" action: {task}\n"));
}
let unblocks = graph::unblocks(&view.state.graph, &node.id);
if !unblocks.is_empty() {
out.push_str(&format!(" unblocks: {}\n", unblocks.join(", ")));
}
}
}
out
}
pub fn transcript(view: &RunView, only: Option<&str>) -> String {
let mut out = String::new();
let settlements = crate::report::evidence(&view.paths, &view.events);
for node in nodes_with_agent_records(view, only) {
out.push_str(&format!("{} {}\n", view.paths.run, one_line(&node)));
for event in view
.events
.iter()
.filter(|event| event.source == Source::Agentgraph)
.filter(|event| event.labels.node.as_deref() == Some(node.as_str()))
{
let field = |key: &str| {
event
.payload
.get(key)
.and_then(|value| value.as_str())
.unwrap_or_default()
};
match event.kind.0.as_str() {
"turn-started" => out.push_str(&format!(
" turn {}\n",
event
.payload
.get("turn")
.map_or_else(|| "-".to_string(), ToString::to_string)
)),
"turn-activity" => out.push_str(&format!(
" {} {} {}\n",
one_line(field("kind")),
one_line(field("name")),
one_line(field("detail"))
)),
_ => {}
}
}
for settled in settlements
.iter()
.filter(|settled| settled.node.as_deref() == Some(node.as_str()))
{
out.push_str(&format!(
" report {} {}\n",
one_line(settled.member.as_deref().unwrap_or("-")),
one_line(&settled.named.display().to_string())
));
let Some(document) = crate::report::read(&settled.kept) else {
out.push_str(
" not retained by this run, so it is not read: only this run's own \
copy of a report is ever opened\n",
);
continue;
};
let turns = crate::report::turns(&document);
if turns.is_empty() {
out.push_str(" it carries no transcript\n");
}
for turn in turns {
out.push_str(&format!(" {}\n", one_line(&turn.role)));
for line in turn.text.lines() {
out.push_str(&format!(" {}\n", one_line(line)));
}
for tool in turn.tools {
out.push_str(&format!(
" {} {} {}\n",
one_line(&tool.kind),
one_line(&tool.name),
one_line(&tool.detail)
));
}
}
}
}
if out.is_empty() {
out.push_str("no dispatch has recorded a transcript\n");
}
out
}
pub(crate) fn nodes_with_agent_records(view: &RunView, only: Option<&str>) -> Vec<String> {
let mut nodes: Vec<String> = view
.events
.iter()
.filter(|event| event.source == Source::Agentgraph)
.filter_map(|event| event.labels.node.clone())
.filter(|node| only.is_none_or(|wanted| wanted == node))
.collect();
nodes.sort_unstable();
nodes.dedup();
nodes
}
fn one_line(text: &str) -> String {
text.chars()
.map(|c| if c.is_control() { ' ' } else { c })
.collect()
}
pub fn goals(survey: &Survey) -> String {
let mut out = String::new();
for view in &survey.views {
let goal = view
.state
.plan
.as_ref()
.and_then(|plan| plan.goal.as_ref())
.map(|goal| goal.text.clone())
.unwrap_or_else(|| crate::plan::NO_GOAL.to_string());
out.push_str(&format!(
"{} {}\n {}\n {}\n",
view.paths.run,
liveness_word(view),
goal,
view.summary()
));
let mut repos: Vec<&str> = view
.state
.graph
.iter()
.filter_map(|node| node.repo.as_deref())
.collect();
repos.sort_unstable();
repos.dedup();
if !repos.is_empty() {
out.push_str(&format!(" identities: {}\n", repos.join(", ")));
}
}
if out.is_empty() {
return nothing_to_report(survey);
}
out.push_str(&skipped_lines(&survey.skipped));
out
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event::{EventKind, Labels, ENVELOPE_VERSION};
use crate::filter::Filters;
use crate::plan::{Node, Plan, PLAN_SCHEMA_VERSION};
use serde_json::json;
use std::path::PathBuf;
fn scratch(name: &str) -> PathBuf {
let dir = std::env::temp_dir().join(format!("onepipeline-views-{name}-{}", sys::pid()));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).expect("a scratch root");
dir
}
fn plan() -> Plan {
Plan {
schema_version: PLAN_SCHEMA_VERSION,
goal: Some(crate::plan::Goal {
text: "close the coverage gap".into(),
}),
name: Some("demo".into()),
concurrency: 4,
tasks: vec![Node {
id: "build".into(),
persona: Some("engineer".into()),
task: Some("## What\ndo it".into()),
..Node::default()
}],
}
}
fn launch(pid: u32) -> LaunchRecord {
LaunchRecord {
run_id: "demo".into(),
plan: PathBuf::from("plan.json"),
dir: PathBuf::from("/tmp/launch"),
graph: "graphs/dag-scope.yaml".into(),
graph_run: String::new(),
node_graph: String::new(),
pr_author_graph: String::new(),
launcher: "claude-code".into(),
session: "session-a".into(),
pid,
host: sys::hostname(),
started: sys::process_start_token(pid)
.map(|token| token.recorded().to_string())
.unwrap_or_default(),
started_at: sys::now_rfc3339(),
heartbeat_interval: 1_800,
dag_sets: Vec::new(),
node_sets: Vec::new(),
adoptions: 0,
filters: Filters::default(),
}
}
fn write_run(root: &Path, run: &str, pid: u32, events: &[Envelope]) -> RunPaths {
let paths = RunPaths::under(root, run);
paths.create().expect("the run directory");
let mut record = launch(pid);
record.run_id = run.to_string();
ledger::write_json(&paths.launch(), &record).expect("a launch record");
for event in events {
ledger::append_line(
&paths.journal(),
&serde_json::to_string(event).expect("an event"),
)
.expect("appended");
}
paths
}
fn hold_lock(paths: &RunPaths) {
ledger::write_json(
&paths.lock(),
&LockRecord {
pid: sys::pid(),
host: sys::hostname(),
acquired_at: sys::now_rfc3339(),
verb: "drive".into(),
started: sys::process_start_token(sys::pid())
.map(|token| token.recorded().to_string())
.unwrap_or_default(),
},
)
.expect("a held lock");
}
fn event(
kind: crate::journal::PipelineKind,
node: Option<&str>,
fields: &[(&str, serde_json::Value)],
) -> Envelope {
relayed(
EventKind(kind.as_str().into()),
Source::Pipeline,
node,
fields,
)
}
fn relayed(
kind: EventKind,
source: Source,
node: Option<&str>,
fields: &[(&str, serde_json::Value)],
) -> Envelope {
Envelope {
v: ENVELOPE_VERSION,
ts: sys::now_rfc3339(),
stream: "s".into(),
seq: 0,
source,
kind,
labels: Labels {
run_id: Some("demo".into()),
round: Some(1),
node: node.map(str::to_string),
..Labels::default()
},
payload: crate::journal::payload(fields),
artifacts: Vec::new(),
}
}
fn dead_pid() -> u32 {
sys::reaped_pid()
}
#[test]
fn a_driver_this_host_can_prove_is_gone_reads_as_driver_dead() {
let root = scratch("dead");
write_run(
&root,
"demo",
dead_pid(),
&[event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
)],
);
let view = RunView::open(&RunPaths::under(&root, "demo")).expect("the run reads");
assert_eq!(view.liveness(), DriverLiveness::DriverDead);
assert!(view.liveness().is_undriven());
assert!(runs(&root, false, "session-a").contains("DRIVER DEAD"));
std::fs::remove_dir_all(&root).ok();
}
fn quiet_run(root: &Path, run: &str) -> RunPaths {
let mut stale = event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
);
stale.ts = "2020-01-01T00:00:00Z".into();
write_run(root, run, sys::pid(), &[stale])
}
#[test]
fn a_live_driver_that_has_gone_quiet_with_nothing_outstanding_reads_as_parked() {
let root = scratch("quiet-parked");
let paths = quiet_run(&root, "demo");
let view = RunView::open(&paths).expect("the run reads");
assert_eq!(view.liveness(), DriverLiveness::Parked);
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_live_driver_quiet_behind_a_blocking_surface_reads_as_active() {
let root = scratch("quiet-blocking");
let paths = quiet_run(&root, "demo");
crate::channel::ChannelState::new(&paths)
.push(crate::channel::Surface {
id: 0,
kind: "blocker".into(),
message: "Node build needs a decision; proceed?".into(),
source: "monitor".into(),
blocking: true,
queued_at: sys::now_millis(),
workstream: Some("build".into()),
})
.expect("the surface queues");
let view = RunView::open(&paths).expect("the run reads");
assert_eq!(view.liveness(), DriverLiveness::Driving);
assert!(!view.liveness().is_undriven());
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_live_driver_that_is_writing_reads_as_active() {
let root = scratch("live");
write_run(
&root,
"demo",
sys::pid(),
&[event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
)],
);
let view = RunView::open(&RunPaths::under(&root, "demo")).expect("the run reads");
assert_eq!(view.liveness(), DriverLiveness::Driving);
assert!(!view.liveness().is_undriven());
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_pid_recorded_on_another_host_never_reads_as_dead() {
let root = scratch("elsewhere");
let paths = RunPaths::under(&root, "demo");
paths.create().expect("the run directory");
let mut record = launch(dead_pid());
record.host = "some-other-host".into();
ledger::write_json(&paths.launch(), &record).expect("a launch record");
ledger::append_line(
&paths.journal(),
&serde_json::to_string(&event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
))
.expect("an event"),
)
.expect("appended");
let view = RunView::open(&paths).expect("the run reads");
assert_eq!(
view.liveness(),
DriverLiveness::Driving,
"a pid means nothing across machines"
);
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_run_nobody_recorded_is_no_such_run() {
let root = scratch("missing");
let error = RunView::open(&RunPaths::under(&root, "nowhere")).unwrap_err();
assert!(matches!(error, crate::Error::NoSuchRun { .. }));
assert!(runs(&root, false, "session-a").contains("no runs recorded"));
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn only_the_reader_sees_mine_and_a_foreign_run_is_labelled_by_digest() {
let root = scratch("owner");
write_run(
&root,
"demo",
sys::pid(),
&[event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
)],
);
let listing = runs(&root, false, "session-a");
assert!(listing.contains("[mine]"), "{listing}");
let foreign = runs(&root, false, "session-b");
assert!(!foreign.contains("[mine]"), "{foreign}");
assert!(
!foreign.contains("session-a"),
"{foreign} leaks the session id"
);
assert!(runs(&root, true, "session-b").contains("no runs recorded"));
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_run_root_this_build_refuses_is_named_rather_than_dropped() {
let root = scratch("skipped");
write_run(
&root,
"readable",
sys::pid(),
&[event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
)],
);
std::fs::create_dir_all(root.join("no-launch")).expect("a directory with no launch");
let typo = RunPaths::under(&root, "typo");
typo.create().expect("the run directory");
std::fs::write(typo.launch(), json!({"oops": true}).to_string()).expect("a launch record");
let survey = Survey::of(&root);
assert_eq!(survey.views.len(), 1, "{:?}", survey.skipped);
assert_eq!(survey.skipped.len(), 2, "{:?}", survey.skipped);
for rendered in [
runs(&root, false, "session-a"),
status(&survey),
goals(&survey),
] {
assert!(rendered.contains("readable"), "{rendered}");
}
for rendered in [
runs(&root, false, "session-a"),
status(&survey),
goals(&survey),
host(&survey),
] {
assert!(rendered.contains("2 run root(s) skipped"), "{rendered}");
assert!(rendered.contains("no-launch"), "{rendered}");
assert!(rendered.contains("launch.json"), "{rendered}");
assert!(rendered.contains("oops"), "{rendered}");
}
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_root_whose_every_run_is_refused_does_not_read_as_no_runs_recorded() {
let root = scratch("all-refused");
std::fs::create_dir_all(root.join("no-launch")).expect("a directory with no launch");
let survey = Survey::of(&root);
assert!(survey.views.is_empty());
for rendered in [
runs(&root, false, "session-a"),
status(&survey),
goals(&survey),
] {
assert!(
!rendered.contains("no runs recorded"),
"a rejected root reported as an absence: {rendered}"
);
assert!(rendered.contains("no run under"), "{rendered}");
assert!(rendered.contains("1 run root(s) skipped"), "{rendered}");
assert!(rendered.contains("no-launch"), "{rendered}");
}
let empty = scratch("all-refused-empty");
assert_eq!(runs(&empty, false, "session-a"), "no runs recorded\n");
std::fs::remove_dir_all(&root).ok();
std::fs::remove_dir_all(&empty).ok();
}
#[test]
fn a_host_row_whose_driver_is_gone_is_counted_rather_than_rendered_live() {
let root = scratch("host-stale");
let paths = write_run(
&root,
"ghosted",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
event(
crate::journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
],
);
ledger::write_json(
&paths.lock(),
&LockRecord {
pid: dead_pid(),
host: sys::hostname(),
acquired_at: sys::now_rfc3339(),
verb: "drive".into(),
started: "a token from the process that died".into(),
},
)
.expect("a stale lock");
let rendered = host(&Survey::of(&root));
assert!(
!rendered.contains("ghosted "),
"a dispatch nothing is driving was rendered as a live row: {rendered}"
);
assert!(rendered.contains("no live dispatches"), "{rendered}");
assert!(
rendered.contains("1 stale registry entry ignored"),
"{rendered}"
);
assert!(rendered.contains("ghosted/build"), "{rendered}");
assert!(rendered.contains("is gone"), "{rendered}");
assert!(rendered.contains(&root.display().to_string()), "{rendered}");
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_host_row_backed_by_a_held_lock_renders_as_a_live_dispatch() {
let root = scratch("host-live");
let paths = write_run(
&root,
"driven",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
event(
crate::journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
],
);
hold_lock(&paths);
let rendered = host(&Survey::of(&root));
assert!(rendered.contains("driven"), "{rendered}");
assert!(rendered.contains("build"), "{rendered}");
assert!(!rendered.contains("no live dispatches"), "{rendered}");
assert!(!rendered.contains("stale registry"), "{rendered}");
assert!(!rendered.contains("UNPROVEN"), "{rendered}");
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_host_row_this_host_cannot_prove_either_way_says_so_rather_than_reading_live() {
let root = scratch("host-unproven");
let paths = write_run(
&root,
"elsewhere",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
event(
crate::journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
],
);
ledger::write_json(
&paths.lock(),
&LockRecord {
pid: sys::pid(),
host: "some-other-host".into(),
acquired_at: sys::now_rfc3339(),
verb: "drive".into(),
started: String::new(),
},
)
.expect("a lock taken elsewhere");
let rendered = host(&Survey::of(&root));
assert!(rendered.contains("elsewhere"), "{rendered}");
assert!(rendered.contains("UNPROVEN"), "{rendered}");
assert!(rendered.contains("some-other-host"), "{rendered}");
assert!(!rendered.contains("stale registry"), "{rendered}");
ledger::write_json(
&paths.lock(),
&LockRecord {
pid: sys::pid(),
host: sys::hostname(),
acquired_at: sys::now_rfc3339(),
verb: "drive".into(),
started: String::new(),
},
)
.expect("a lock from a build that predates the stamp");
let rendered = host(&Survey::of(&root));
assert!(rendered.contains("UNPROVEN"), "{rendered}");
assert!(rendered.contains("no start token"), "{rendered}");
ledger::write_json(
&paths.lock(),
&LockRecord {
pid: sys::pid(),
host: sys::hostname(),
acquired_at: sys::now_rfc3339(),
verb: "drive".into(),
started: "the process that took it, which was not this one".into(),
},
)
.expect("a lock a reused pid now answers for");
let rendered = host(&Survey::of(&root));
assert!(
rendered.contains("1 stale registry entry ignored"),
"{rendered}"
);
assert!(rendered.contains("different process"), "{rendered}");
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_failed_node_names_the_side_and_the_identity_that_refused() {
let root = scratch("refusal");
let refused = |role: Option<&str>, identity: &str, reason: &str| {
let mut fields = vec![("identity", json!(identity)), ("reason", json!(reason))];
if let Some(role) = role {
fields.push(("role", json!(role)));
}
let mut envelope = relayed(
EventKind("fallback-advanced".into()),
Source::Agentgraph,
Some("build"),
&fields,
);
envelope.stream = "oneagentgraph-1".into();
envelope
.labels
.extra
.insert("member".into(), "worker".into());
envelope
};
write_run(
&root,
"refused",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
event(
crate::journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
refused(Some("agent"), "claude-code", "quota"),
refused(Some("judge"), "codex", "rate_limit"),
refused(Some("judge"), "codex", "rate_limit"),
event(
crate::journal::PipelineKind::NodeSettled,
Some("build"),
&[
("status", json!("failed")),
("outcome", json!("task-failed")),
],
),
],
);
let survey = Survey::of(&root);
let rendered = results(&survey.views[0]);
assert!(rendered.contains("the agent side"), "{rendered}");
assert!(rendered.contains("claude-code"), "{rendered}");
assert!(rendered.contains("(quota)"), "{rendered}");
assert!(rendered.contains("the judge side"), "{rendered}");
assert!(rendered.contains("codex"), "{rendered}");
assert!(rendered.contains("recorded 2 times"), "{rendered}");
let rendered = status(&survey);
assert!(rendered.contains("build: failed —"), "{rendered}");
assert!(rendered.contains("the judge side"), "{rendered}");
assert!(rendered.contains("codex"), "{rendered}");
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn an_unattributed_refusal_is_never_given_a_side_it_did_not_carry() {
let advanced = |reason: &str| oneagentgraph::event::FallbackAdvanced {
identity: "codex".into(),
reason: reason.into(),
role: None,
turn: None,
};
let single = Refusal {
advanced: advanced("auth"),
member: MemberLabel::Named("worker".into()),
records: std::num::NonZeroU64::MIN,
};
assert_eq!(
refusal_phrase(&single),
"member 'worker': identity 'codex' refused (auth)"
);
let bare = Refusal {
advanced: advanced(""),
member: MemberLabel::Unstamped,
records: std::num::NonZeroU64::MIN,
};
let phrase = refusal_phrase(&bare);
assert!(
phrase.contains("a side the record does not name"),
"{phrase}"
);
assert!(
phrase.contains("for a reason the record does not carry"),
"{phrase}"
);
let mut nameless = relayed(
EventKind("fallback-advanced".into()),
Source::Agentgraph,
Some("build"),
&[("reason", json!("quota"))],
);
nameless.stream = "oneagentgraph-1".into();
assert!(projection::fold(&[nameless]).refusals.is_empty());
}
#[test]
fn every_view_renders_from_the_merged_stream() {
let root = scratch("render");
let mut agent = relayed(
EventKind("turn-finished".into()),
Source::Agentgraph,
Some("build"),
&[("message", json!("ran the gate"))],
);
agent.stream = "oneagentgraph-1".into();
let mut vcs = relayed(
EventKind("session-opened".into()),
Source::Vcs,
Some("build"),
&[("branch", json!("feature"))],
);
vcs.stream = "onevcs-tok".into();
write_run(
&root,
"demo",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
event(crate::journal::PipelineKind::NodeReady, Some("build"), &[]),
event(
crate::journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
agent,
vcs,
],
);
let view = RunView::open(&RunPaths::under(&root, "demo")).expect("the run reads");
let stream = monitor(&view, &EventFilter::default());
assert!(stream.starts_with("Concise graph events;"), "{stream}");
assert!(stream.contains("agent:oneagentgraph-1"), "{stream}");
assert!(stream.contains("vcs:onevcs-tok"), "{stream}");
assert!(stream.contains("graph:build"), "{stream}");
assert!(stream.contains("-- demo 0/1 done"), "{stream}");
assert!(
!stream.contains("round"),
"a round reached a view: {stream}"
);
hold_lock(&RunPaths::under(&root, "demo"));
let survey = Survey::of(&root);
assert!(status(&survey).contains("build: running"));
assert!(host(&survey).contains("build"));
assert!(goals(&survey).contains("close the coverage gap"));
assert!(results(&survey.views[0]).contains("build"));
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_node_the_ledger_calls_running_that_nothing_drives_is_undriven() {
let root = scratch("undriven");
write_run(
&root,
"demo",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
event(
crate::journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
],
);
let rendered = status(&Survey::of(&root));
assert!(rendered.contains("UNDRIVEN"), "{rendered}");
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_live_dispatch_reports_what_it_is_doing_now_with_a_count_and_an_age() {
let root = scratch("activity");
let mut turn = relayed(
EventKind("turn-activity".into()),
Source::Agentgraph,
Some("build"),
&[
("kind", json!("tool_call")),
("name", json!("Bash")),
("detail", json!("cargo llvm-cov --workspace")),
],
);
turn.stream = "oneagentgraph-1".into();
write_run(
&root,
"demo",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
event(
crate::journal::PipelineKind::NodeDispatched,
Some("build"),
&[],
),
turn,
],
);
let rendered = status(&Survey::of(&root));
assert!(
rendered.contains("now Bash cargo llvm-cov --workspace"),
"{rendered}"
);
assert!(rendered.contains("1 event(s)"), "{rendered}");
assert!(rendered.contains("ago"), "{rendered}");
assert!(
!rendered.contains(DriverLiveness::Undriven.as_str()),
"a dispatch that is recording was reported as driving nothing: {rendered}"
);
std::fs::remove_dir_all(&root).ok();
}
fn recorded(events: u64, last_at: u64) -> Option<crate::projection::Progress> {
(0..events).fold(None, |progress, _| match progress {
None => crate::projection::Progress::first(Some(last_at)),
Some(progress) => Some(progress.and(Some(last_at))),
})
}
#[test]
fn a_dispatch_that_has_named_no_tool_reports_its_count_rather_than_a_guess() {
let rendered = working(&crate::projection::NodeActivity {
doing: None,
progress: recorded(3, sys::now_millis()),
last_heartbeat_at: None,
});
assert_eq!(rendered, "3 event(s), 0s ago");
assert!(!rendered.contains("now"), "{rendered}");
}
#[test]
fn a_heartbeat_is_reported_beside_the_age_of_the_work_rather_than_as_work() {
let now = sys::now_millis();
let rendered = working(&crate::projection::NodeActivity {
doing: Some("Bash red-green.sh".into()),
progress: recorded(4, now - 600_000),
last_heartbeat_at: Some(now),
});
assert!(
rendered.contains("4 event(s), 10m00s ago"),
"the age of the work was taken from the heartbeat: {rendered}"
);
assert!(
rendered.contains("alive 0s ago"),
"a dispatch that is alive and doing nothing is not reported as alive: {rendered}"
);
}
#[test]
fn a_dispatch_that_has_only_heartbeated_reports_no_work_and_still_reads_as_alive() {
let rendered = working(&crate::projection::NodeActivity {
doing: None,
progress: None,
last_heartbeat_at: Some(sys::now_millis()),
});
assert_eq!(rendered, "nothing recorded yet; alive 0s ago");
}
#[test]
fn a_transcript_renders_the_turns_tools_and_the_report_it_settled_with() {
let root = scratch("transcript");
let paths = RunPaths::under(&root, "demo");
let stored = paths.report_for("s", 0);
std::fs::create_dir_all(paths.reports_dir()).expect("the run's report storage");
std::fs::write(
&stored,
json!({
"schema_version": 7,
"transcript": {"messages": [
{"role": "assistant", "content": "Ran the gate.\nIt passed.", "events": [
{"kind": "tool_call", "name": "bash", "input": {"command": "just check"}},
]},
]},
})
.to_string(),
)
.expect("a stored report");
let mut started = relayed(
EventKind("turn-started".into()),
Source::Agentgraph,
Some("build"),
&[("turn", json!(1))],
);
started.stream = "oneagentgraph-1".into();
let mut activity = relayed(
EventKind("turn-activity".into()),
Source::Agentgraph,
Some("build"),
&[
("kind", json!("tool_call")),
("name", json!("bash")),
("detail", json!("just check")),
],
);
activity.stream = "oneagentgraph-1".into();
let settled = relayed(
EventKind(crate::report::MEMBER_SETTLED.into()),
Source::Agentgraph,
Some("build"),
&[(crate::report::REPORT_PATH, json!("/elsewhere/report.json"))],
);
write_run(
&root,
"demo",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
started,
activity,
settled,
],
);
let view = RunView::open(&RunPaths::under(&root, "demo")).expect("the run reads");
let rendered = transcript(&view, None);
assert!(rendered.contains("demo build"), "{rendered}");
assert!(rendered.contains("turn 1"), "{rendered}");
assert!(
rendered.contains("tool_call bash just check"),
"{rendered}"
);
assert!(rendered.contains("assistant"), "{rendered}");
assert!(rendered.contains("Ran the gate."), "{rendered}");
assert!(rendered.contains("It passed."), "{rendered}");
assert!(transcript(&view, Some("elsewhere")).contains("no dispatch"));
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_report_this_run_did_not_keep_is_named_as_unretained_and_never_opened() {
let root = scratch("transcript-unread");
let settled = relayed(
EventKind(crate::report::MEMBER_SETTLED.into()),
Source::Agentgraph,
Some("build"),
&[(
crate::report::REPORT_PATH,
json!("/nowhere/onepipeline/report.json"),
)],
);
write_run(
&root,
"demo",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
settled,
],
);
let view = RunView::open(&RunPaths::under(&root, "demo")).expect("the run reads");
let rendered = transcript(&view, None);
assert!(rendered.contains("not retained by this run"), "{rendered}");
assert!(
rendered.contains("/nowhere/onepipeline/report.json"),
"the path that was not read is not named: {rendered}"
);
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_report_without_a_transcript_says_so() {
let root = scratch("transcript-none");
let paths = RunPaths::under(&root, "demo");
std::fs::create_dir_all(paths.reports_dir()).expect("the run's report storage");
std::fs::write(
paths.report_for("s", 0),
json!({"usage": {"input_tokens": 1}}).to_string(),
)
.expect("a stored report");
let settled = relayed(
EventKind(crate::report::MEMBER_SETTLED.into()),
Source::Agentgraph,
Some("build"),
&[(crate::report::REPORT_PATH, json!("/elsewhere/report.json"))],
);
write_run(
&root,
"demo",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
),
settled,
],
);
let view = RunView::open(&RunPaths::under(&root, "demo")).expect("the run reads");
assert!(transcript(&view, None).contains("carries no transcript"));
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_run_that_dispatched_nothing_has_no_transcript_to_render() {
let root = scratch("transcript-empty");
write_run(
&root,
"demo",
sys::pid(),
&[event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(plan()))],
)],
);
let view = RunView::open(&RunPaths::under(&root, "demo")).expect("the run reads");
assert_eq!(
transcript(&view, None),
"no dispatch has recorded a transcript\n"
);
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_waiting_human_reports_its_action_and_what_it_unblocks() {
let root = scratch("waiting");
let mut waiting_plan = plan();
waiting_plan.tasks = vec![
Node {
id: "approve".into(),
kind: crate::plan::NodeKind::Human,
task: Some("approve the release".into()),
..Node::default()
},
Node {
id: "ship".into(),
persona: Some("engineer".into()),
task: Some("## What\nship".into()),
deps: vec!["approve".into()],
..Node::default()
},
];
write_run(
&root,
"demo",
sys::pid(),
&[
event(
crate::journal::PipelineKind::RunStarted,
None,
&[("plan", json!(waiting_plan))],
),
event(
crate::journal::PipelineKind::NodeSettled,
Some("approve"),
&[("status", json!("waiting"))],
),
],
);
let view = RunView::open(&RunPaths::under(&root, "demo")).expect("the run reads");
let rendered = results(&view);
assert!(rendered.contains("approve the release"), "{rendered}");
assert!(rendered.contains("unblocks: ship"), "{rendered}");
assert!(
rendered.contains("ship") && rendered.contains("blocked"),
"{rendered}"
);
std::fs::remove_dir_all(&root).ok();
}
#[test]
fn a_summary_line_is_capped_and_control_stripped() {
let long = "x".repeat(500);
let stripped = summarize(&relayed(
EventKind("kind".into()),
Source::Agentgraph,
None,
&[("message", json!(format!("a\nb{long}")))],
));
assert!(!stripped.contains('\n'), "{stripped}");
assert_eq!(stripped.chars().count(), 96);
}
#[test]
fn the_parked_threshold_is_read_from_the_environment_or_defaults() {
assert!(parked_after_seconds() > 0);
}
#[test]
fn every_liveness_verdict_has_the_word_the_contract_fixes() {
assert_eq!(DriverLiveness::Driving.as_str(), "ACTIVE");
assert_eq!(DriverLiveness::DriverDead.as_str(), "DRIVER DEAD");
assert_eq!(DriverLiveness::Parked.as_str(), "PARKED");
assert_eq!(DriverLiveness::Undriven.as_str(), "UNDRIVEN");
}
}