Read the tools, question details, script and condition off Petri's records

The Fabro halves of D1, D2 and D3. A stage's `agent_tools` is the union, by
name, of the `attractor.tools` payloads its native sessions record, with
`invoked` flipped by the envelope's `ToolCallStarted`; the payload carries
Petri's origin category, so Pebble's category is `subagent` for a sub-agent
tool and `other` for the rest. A pending question carries each option's
description and preview and the question's context, and its reference as
the review target when Fabro's validation admits it; the interview dock and
the Q&A renderer show the previews beside the descriptions, and the attach
prompt prints both under each choice. The web's command view reads the
script from the node's `meta.script`, the decision renderer the matched
condition from `meta.edges[edge].condition`, and `run events --pretty`
prints the condition on the transition line and one line per session
naming its tool count. The command view notes what the output capture did
not keep, from the final `step.finished` loss metrics.

The web fixtures are recaptured at the pin, so they carry the new facts.
VIEWS.md loses the two gap rows Petri filled and names the sources; the
README's list of what the fold leaves default shrinks to match.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 22:17:15 -04:00
parent 08b7a4fdd9
commit 572f89a6f2
No known key found for this signature in database
17 changed files with 6924 additions and 4788 deletions

View file

@ -191,7 +191,7 @@ describe("InterviewDock", () => {
expect(buttons.Revise).toBeDefined();
});
test("multiple choice renders option descriptions as display text", () => {
test("multiple choice renders option descriptions and previews as display text", () => {
const question = makeQuestion({
question_type: QuestionType.MULTIPLE_CHOICE,
options: [
@ -201,6 +201,7 @@ describe("InterviewDock", () => {
description: "Deploy the current patch",
preview: "<b>not rendered specially</b>",
},
{ key: "R", label: "[R] Revise" },
],
});
const tree = render(
@ -209,7 +210,13 @@ describe("InterviewDock", () => {
const text = textContent(tree.toJSON());
expect(text).toContain("Approve");
expect(text).toContain("Deploy the current patch");
expect(text).not.toContain("<b>not rendered specially</b>");
// The preview is shown as the text it is, never parsed as markup.
expect(text).toContain("<b>not rendered specially</b>");
const previews = tree.root.findAll(
(node) => node.props["data-testid"] === "interview-option-preview",
);
expect(previews).toHaveLength(1);
expect(shouldStackOptions(question.options ?? [])).toBe(true);
});
test("freeform question renders a textarea and disables send when empty", () => {

View file

@ -309,7 +309,9 @@ function ConfirmationBody({
export function shouldStackOptions(options: InterviewOption[]): boolean {
return options.some(
(option) =>
option.label.length > STACK_LABEL_LENGTH || Boolean(option.description),
option.label.length > STACK_LABEL_LENGTH ||
Boolean(option.description) ||
Boolean(option.preview),
);
}
@ -465,6 +467,14 @@ function OptionLabel({ option }: { option: InterviewOption }) {
{option.description}
</span>
)}
{option.preview && (
<span
data-testid="interview-option-preview"
className="mt-1 block whitespace-pre-wrap rounded bg-overlay-strong px-2 py-1 font-mono text-[11px]/4 font-normal text-fg-3"
>
{option.preview}
</span>
)}
</span>
);
}

View file

@ -200,6 +200,11 @@ function QuestionBlock({
{option.description}
</span>
)}
{option.preview && (
<span className="mt-1 block whitespace-pre-wrap rounded bg-overlay-strong px-2 py-1 font-mono text-[11px]/4 text-fg-3">
{option.preview}
</span>
)}
</span>
</li>
))}

View file

@ -1,9 +1,12 @@
import { describe, expect, test } from "bun:test";
import { loadPetriFixture } from "./petri-fixtures";
import type { RunStreamItem } from "@qltysh/fabro-api-client";
import {
agentEnvelopesOf,
commandOutcomeOf,
commandScriptOf,
debugRowsFromStream,
deriveRunPhasesFromStream,
extractPetriStageContext,
@ -12,6 +15,8 @@ import {
isTerminalLifecycleItem,
itemsForStage,
parallelOverviewFromProjection,
matchedCondition,
outputLossNote,
parsePetriInterviewPairs,
petriEventName,
petriStageLabel,
@ -201,9 +206,113 @@ describe("stage renderers", () => {
test("a command stage's outcome is read from its final step.finished", () => {
const say = itemsForStage(command.stream, "say@1");
expect(commandOutcomeOf(say).exitCode).toBe(0);
expect(commandOutcomeOf(say).outputLoss).toBeNull();
expect(extractPetriStageContext(say)).toBeNull();
});
test("an agent stage's projection lists the tools its session was offered", () => {
const names = (hello.projection.stages["greet@1"]?.agent_tools ?? []).map((tool) => tool.name);
expect(names).toContain("read_file");
expect(names).toContain("shell");
expect(names).toContain("request_user_input");
expect(hello.projection.stages["start@1"]?.agent_tools ?? []).toEqual([]);
});
test("a command stage's script rides on its node's meta", () => {
const say = itemsForStage(command.stream, "say@1");
expect(commandScriptOf(say)).toBe("echo hello from petri");
expect(commandScriptOf(itemsForStage(command.stream, "start@1"))).toBeNull();
});
test("the condition an edge matched is read from the node's edge table", () => {
const applied = (edge: number): RunStreamItem => ({
run_id: "run",
stream_seq: 9,
kind: "petri",
id: "9",
recorded_at: 1_789_706_579_000,
item: {
id: { log: "execution", execution: 0, seq: 9, index: 0 },
origin: "core",
context: { invocation: 0, execution: 0 },
subject: {
node: {
id: 2,
name: "build",
kind: "attractor/command",
meta: {
kind: "command",
edges: {
"0": { to: "ok", label: null, condition: "outcome=succeeded" },
"1": { to: "bad", label: null },
},
},
},
firing: 2,
visit: 1,
attempt: 1,
generation: 0,
branch: { role: "none" },
},
record: {
seq: 9,
body: { event: "route.applied", kind: "edge", firing: 2, group: 0, edge },
},
derived: { target: { name: edge === 0 ? "ok" : "bad" }, transition: "Continue", back: false },
},
});
expect(matchedCondition(applied(0))).toBe("outcome=succeeded");
expect(matchedCondition(applied(1))).toBeUndefined();
expect(findPetriEdgeForStage([applied(0)], "build@1")).toEqual({
fromNode: "build",
toNode: "ok",
reason: "condition",
condition: "outcome=succeeded",
isJump: false,
});
});
test("a command's output loss is read from its metrics and worded for the view", () => {
const finished = (custom: Record<string, unknown>): RunStreamItem => ({
run_id: "run",
stream_seq: 5,
kind: "petri",
id: "5",
recorded_at: 1_789_706_579_000,
item: {
id: { log: "execution", execution: 0, seq: 5, index: 0 },
origin: "external",
context: { invocation: 0, execution: 0 },
subject: { node: { id: 2, name: "say", kind: "attractor/command", meta: { kind: "command" } }, firing: 2, visit: 1, attempt: 1, generation: 0, branch: { role: "none" } },
record: {
seq: 5,
body: {
event: "step.finished",
firing: 2,
attempt: 1,
outcome: {
status: "success",
output: { stdout: "x", exit_status: 0 },
metrics: { duration_ms: 3, exit_code: 0, custom },
},
},
},
derived: { final: true, exhausted: false },
},
});
expect(commandOutcomeOf([finished({})]).outputLoss).toBeNull();
const cut = commandOutcomeOf([
finished({ "output.dropped_bytes": 2048, "output.truncated_lines": 1 }),
]).outputLoss;
expect(cut).toEqual({ droppedBytes: 2048, truncatedLines: 1, incomplete: false });
expect(outputLossNote(cut)).toBe("Output truncated: 2,048 bytes dropped, 1 line cut");
const silent = commandOutcomeOf([finished({ "output.incomplete": true })]).outputLoss;
expect(outputLossNote(silent)).toBe(
"Output may be incomplete: the capture ended on silence, so the tail may be missing",
);
expect(outputLossNote(null)).toBeNull();
});
test("an agent stage's Pebble envelopes are read with their variant and session", () => {
const envelopes = agentEnvelopesOf(itemsForStage(hello.stream, "greet@1"));
expect(envelopes.length).toBeGreaterThan(0);

View file

@ -499,12 +499,13 @@ export function findPetriEdgeForStage(
if (petriStageLabel(item) !== stageLabel) continue;
const target = getString(getObject(derived(item), "target"), "name");
if (!target) continue;
const kind = getString(petriBody(item), "kind") ?? "edge";
const body = petriBody(item);
const kind = getString(body, "kind") ?? "edge";
latest = {
fromNode: getString(subjectNode(item), "name") ?? stageLabel,
toNode: target,
reason: kind === "jump" ? "jump" : "condition",
condition: null,
condition: matchedCondition(item) ?? null,
isJump: kind === "jump",
};
}
@ -633,28 +634,56 @@ export function agentEnvelopesOf(items: PetriStream): PetriAgentEnvelope[] {
return out;
}
/** A command stage's script, from its `step.started` record, if recorded. */
/**
* The condition a `route.applied` item's edge matched, as written: the
* record's `edge` keys the subject node's `meta.edges`, whose entry carries
* the edge's `condition` when it has one (EVENTS.md "Source metadata").
*/
export function matchedCondition(item: RunStreamItem): string | undefined {
const edge = getNumber(petriBody(item), "edge");
if (edge === undefined) return undefined;
const edges = getObject(getObject(subjectNode(item), "meta"), "edges");
return getString(getObject(edges, String(edge)), "condition");
}
/**
* A command stage's script: the text the step runs rides on the node's
* `meta.script`, on every event of the stage (EVENTS.md "Source metadata").
*/
export function commandScriptOf(items: PetriStream): string | null {
for (const item of items) {
if (petriEventName(item) !== "step.started") continue;
const script =
getString(getObject(petriBody(item), "config"), "script") ??
getString(getObject(getObject(subjectNode(item), "meta"), "config"), "script");
const script = getString(getObject(subjectNode(item), "meta"), "script");
if (script) return script;
}
return null;
}
/**
* The exit code and duration of the stage's final `step.finished`: the
* command step's output carries `exit_status`, its metrics the duration.
* What a command's output capture did not keep, from the final
* `step.finished` metrics: `output.dropped_bytes` and
* `output.truncated_lines` count what the caps cut; `output.incomplete`
* says the capture ended on silence, so the tail may be missing by an
* amount nobody counted. Absent when the output is whole.
*/
export interface CommandOutputLoss {
droppedBytes: number;
truncatedLines: number;
incomplete: boolean;
}
/**
* The exit code, duration and output loss of the stage's final
* `step.finished`: the command step's output carries `exit_status`, its
* metrics the duration and the loss counters under `custom`.
*/
export function commandOutcomeOf(items: PetriStream): {
exitCode: number | null;
durationMs: number;
outputLoss: CommandOutputLoss | null;
} {
let exitCode: number | null = null;
let durationMs = 0;
let outputLoss: CommandOutputLoss | null = null;
for (const item of items) {
if (petriEventName(item) !== "step.finished") continue;
const outcome = getObject(petriBody(item), "outcome");
@ -663,8 +692,32 @@ export function commandOutcomeOf(items: PetriStream): {
exitCode =
getNumber(output, "exit_status") ?? getNumber(metrics, "exit_code") ?? exitCode;
durationMs = getNumber(metrics, "duration_ms") ?? durationMs;
const custom = getObject(metrics, "custom");
const droppedBytes = getNumber(custom, "output.dropped_bytes") ?? 0;
const truncatedLines = getNumber(custom, "output.truncated_lines") ?? 0;
const incomplete = getBool(custom, "output.incomplete") === true;
outputLoss =
droppedBytes > 0 || truncatedLines > 0 || incomplete
? { droppedBytes, truncatedLines, incomplete }
: null;
}
return { exitCode, durationMs };
return { exitCode, durationMs, outputLoss };
}
/** The one-line note the stage view shows beside output that is not whole. */
export function outputLossNote(loss: CommandOutputLoss | null): string | null {
if (!loss) return null;
const parts: string[] = [];
if (loss.droppedBytes > 0) {
parts.push(`${loss.droppedBytes.toLocaleString()} bytes dropped`);
}
if (loss.truncatedLines > 0) {
parts.push(`${loss.truncatedLines} ${loss.truncatedLines === 1 ? "line" : "lines"} cut`);
}
const counted = parts.length > 0 ? `Output truncated: ${parts.join(", ")}` : null;
if (!loss.incomplete) return counted;
const tail = "the capture ended on silence, so the tail may be missing";
return counted ? `${counted}; ${tail}` : `Output may be incomplete: ${tail}`;
}
// ── Stages from the projection ──────────────────────────────────────────

View file

@ -83,6 +83,7 @@ import {
agentEnvelopesOf,
commandOutcomeOf,
commandScriptOf,
outputLossNote,
debugRowSearchText,
debugRowsFromStream,
extractPetriStageContext,
@ -91,6 +92,7 @@ import {
parallelOverviewFromProjection,
parsePetriInterviewPairs,
reducerTranscriptFromProjection,
type CommandOutputLoss,
type DebugRow,
} from "../lib/petri-stream";
import {
@ -156,6 +158,8 @@ type TurnType =
exitCode: number | null;
durationMs: number;
outputBytes: number;
/** What the capture did not keep, or null when the output is whole. */
outputLoss: CommandOutputLoss | null;
};
type CommandTurn = Extract<TurnType, { kind: "command" }>;
@ -353,6 +357,7 @@ export function buildPetriStageActivity(
exitCode: outcome.exitCode,
durationMs: outcome.durationMs || (stage?.timing?.wall_time_ms ?? 0),
outputBytes: stage?.output_bytes ?? 0,
outputLoss: outcome.outputLoss,
});
}
return { turns, pendingTools: [] };
@ -1510,6 +1515,14 @@ function CommandLogs({
byteCount={turn.outputBytes}
enabled={!turn.running}
/>
{outputLossNote(turn.outputLoss) && (
<p
data-testid="command-output-loss"
className="text-xs text-amber"
>
{outputLossNote(turn.outputLoss)}
</p>
)}
</div>
);
}

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

View file

@ -472,6 +472,17 @@ async fn ask_attach_question(question: Question, styles: &'static Styles) -> Ans
opt.key,
opt.label,
);
// What choosing the option means, and a sample of what it
// would do, when the asking stage said.
for detail in [opt.description.as_deref(), opt.preview.as_deref()]
.into_iter()
.flatten()
.filter(|detail| !detail.trim().is_empty())
{
for line in detail.lines() {
eprintln!(" {}", styles.dim.apply_to(line));
}
}
}
if question.allow_freeform {
eprintln!(" Or type a free-text response");

View file

@ -161,6 +161,17 @@ impl<'a> PetriItem<'a> {
pub(crate) fn str_at(self, pointer: &str) -> Option<&'a str> {
self.value().pointer(pointer)?.as_str()
}
/// The condition the applied route matched, as written: the edge a
/// `route.applied` record names is a key into the subject node's
/// `meta.edges`, whose entry carries the edge's `condition` when it
/// has one.
pub(crate) fn matched_condition(self) -> Option<&'a str> {
let edge = self.body()?.get("edge")?.as_u64()?;
self.node()?
.pointer(&format!("/meta/edges/{edge}/condition"))?
.as_str()
}
}
/// Whether the item is the platform record of the run's terminal lifecycle
@ -393,15 +404,20 @@ pub(crate) fn format_pretty(
.and_then(Value::as_bool)
.unwrap_or(false);
let detail = if back { " (loop)" } else { "" };
// The condition the edge matched, when it has one.
let condition = view
.matched_condition()
.map_or_else(String::new, |condition| format!(" when {condition}"));
// Petri records the route after the next visit started, so the
// line names both ends of the edge.
Some(format!(
"{ts} {} {} {} {}{}",
"{ts} {} {} {} {}{}{}",
styles.dim.apply_to(view.node_name().unwrap_or("?")),
styles.dim.apply_to("\u{2192}"),
target,
styles.dim.apply_to(&transition),
styles.dim.apply_to(detail),
styles.dim.apply_to(&condition),
))
}
"fork.started" => {
@ -592,6 +608,19 @@ fn format_progress(ts: &str, view: PetriItem<'_>, label: &str, styles: &Styles)
let body = indented(styles, response, " ");
Some(format!("{header}\n{body}\n"))
}
"attractor.tools" => {
// The tools one native session was offered, once per session.
let count = custom
.get("tools")
.and_then(Value::as_array)
.map_or(0, Vec::len);
let noun = if count == 1 { "tool" } else { "tools" };
Some(format!(
"{ts} {} {}",
styles.dim.apply_to("\u{2699}"),
styles.dim.apply_to(format!("{count} {noun} available")),
))
}
"attractor.checkout" => {
let repository = custom
.get("repository")
@ -1106,6 +1135,59 @@ mod tests {
assert!(line.contains("answered by dev"), "{line}");
}
#[test]
fn a_route_line_names_the_condition_the_edge_matched() {
let styles = Styles::new(false);
let mut state = PrettyState::default();
let mut subject = subject("build", "command");
subject["node"]["meta"]["edges"] = json!({
"0": {"to": "ok", "label": null, "condition": "outcome=succeeded"},
"1": {"to": "bad", "label": null}
});
let applied = |edge: u64| {
petri(
7,
json!({
"origin": "core",
"context": {"invocation": 0, "execution": 0},
"subject": subject,
"record": {"seq": 9, "body": {"event": "route.applied", "kind": "edge",
"firing": 2, "group": 0, "edge": edge}},
"derived": {"target": {"name": if edge == 0 { "ok" } else { "bad" }},
"transition": "Continue", "back": false}
}),
)
};
let line = format_pretty(&applied(0), &styles, &mut state).expect("a route line");
assert!(
line.contains("build \u{2192} ok continue when outcome=succeeded"),
"{line}"
);
let line = format_pretty(&applied(1), &styles, &mut state).expect("a route line");
assert!(line.ends_with("build \u{2192} bad continue"), "{line}");
}
#[test]
fn a_sessions_tool_list_renders_as_its_count() {
let styles = Styles::new(false);
let mut state = PrettyState::default();
let listed = petri(
8,
json!({
"origin": "external",
"context": {"invocation": 0, "execution": 0},
"subject": subject("work", "agent"),
"record": {"seq": 10, "body": {"event": "step.progress.recorded", "firing": 2,
"ev": {"custom": {"kind": "attractor.tools", "session": "ses_1", "tools": [
{"name": "shell", "description": "Run a command", "source": {"kind": "native"}, "category": "builtin"},
{"name": "fabro_run_create", "description": "Create a run", "source": {"kind": "application"}, "category": "host"}
]}}}}
}),
);
let line = format_pretty(&listed, &styles, &mut state).expect("a tools line");
assert!(line.contains("2 tools available"), "{line}");
}
#[test]
fn the_engine_finish_decides_the_exit_code() {
let finished = petri(

View file

@ -122,14 +122,30 @@ prompt or agent stage carries the stub's text as its `response`.
`VIEWS.md` rows with no source yet, or whose source this crate does not read
yet, keep their default value in the projection: `Checkpoint`'s
engine-derived maps (`completed_nodes`, `node_retries`, `context_values`,
`node_outcomes`, `next_node_id`), `agent_tools`, `permission_level`,
`script_invocation` and `script_timing`, a stage's `notes`,
`StageCompletion` details for a `parsed.note`, the sandbox instance's
clone fields and workspace roots (Petri's checkout is a copy of the bound
repository, not a clone; the roots are the provider's, read live),
`Run.ask_fabro`, an interview option's `description` and
`preview`, the pull request `creation` state, and the run's notices,
notifications and pairings (recorded, not shown).
`node_outcomes`, `next_node_id`), `permission_level`,
`script_invocation` and `script_timing` (a command's script is on the
stream, as `subject.node.meta.script`, and the web's command view reads it
there), a stage's `notes`, `StageCompletion` details for a `parsed.note`,
the sandbox instance's clone fields and workspace roots (Petri's checkout
is a copy of the bound repository, not a clone; the roots are the
provider's, read live), `Run.ask_fabro`, the pull request `creation`
state, and the run's notices, notifications and pairings (recorded, not
shown).
Three facts the views once lacked a source for are read now. A stage's
`agent_tools` is the union, by name, of the `attractor.tools` payloads its
native sessions record (the node's own session, then each child session),
with `invoked` flipped by the envelope's `ToolCallStarted`; the payload
carries Petri's origin category, so Pebble's `category` is `subagent` for
a sub-agent tool and `other` for the rest. A pending question carries each
option's `description` and `preview` and the question's `context` as
`context_display`, and its `reference` as `review_target` when Fabro's
validation admits it. A decision's matched condition and a command's
script ride on the node's `meta` (`edges[edge].condition`, `script`), which
the CLI's `run events --pretty` and the web's stage renderers read off
the stream, beside the command's output loss counters
(`output.dropped_bytes`, `output.truncated_lines`, `output.incomplete`) on
its final `step.finished`.
### Retention

View file

@ -140,10 +140,11 @@ stages live in the child invocation and list under the fork (see Parallel).
| provider and model | `RunStage.provider_used`, `StageProjection.provider_used`, `model`, `permission_level` | `custom attractor.fallback.plan {requested, routes}` then envelope `SessionStarted {provider, model}`; `custom attractor.prompt {model}`; the node's config in the registered graph (`graph.registered`, blob by digest) for `reasoning_effort`, `speed`, `permission_level` and an ACP node's settings | attempt |
| prompt | `StageProjection.prompt`, the chat tab's `stage.prompt` | `custom attractor.prompt {prompt, sources}`; envelope `SessionStarted` and the first user message on the stream for an agent | attempt |
| response | `StageProjection.response`, `prompt.completed` | `custom attractor.prompt.completed {response, calls, repairs, usage, duration_ms}`; the final `step.finished` `outcome.output` for an agent | attempt |
| output, output bytes, streaming, termination | `StageProjection.output`, `output_bytes`, `live_streaming`, `termination`, `command.started` `script`, `command.completed` `exit_code`, the command log endpoint | `step.started`; `step.progress.recorded` `log {stream, line}` (the live log); `step.finished` `outcome.output`, `metrics.exit_code`, `metrics.duration_ms`; `timed_out` and `cancelled` statuses for `termination`; a `blob://` output through `get_blob`; the script from the node config | attempt |
| output, output bytes, streaming, termination | `StageProjection.output`, `output_bytes`, `live_streaming`, `termination`, `command.started` `script`, `command.completed` `exit_code`, the command log endpoint | `step.started`; `step.progress.recorded` `log {stream, line}` (the live log); `step.finished` `outcome.output`, `metrics.exit_code`, `metrics.duration_ms`; `timed_out` and `cancelled` statuses for `termination`; a `blob://` output through `get_blob`; the script is `subject.node.meta.script` on every event of the stage (the web's command view reads it off the stream) | attempt |
| output loss | the command view's "output truncated" note | the final `step.finished` `metrics.custom` `output.dropped_bytes`, `output.truncated_lines` (what the caps cut) and `output.incomplete` (the capture ended on silence); absent when the output is whole | attempt |
| script invocation and timing | `script_invocation`, `script_timing` | the node config; `metrics.duration_ms` | attempt |
| context updates, routing directive | `stage.completed` `context_updates`, `preferred_label`, `suggested_next_ids`, `jump_to_node` | `step.finished` `outcome.context_updates`; `routing.resolved` (per group the decision, overrides, jumps, blocks, the weighted draw; `derived.groups[].target`) | attempt |
| edge selected, loop restart | `edge.selected`, `loop.restart` | `route.applied` (`derived.target`, `transition`, `back`); a restart is `execution.finished {restart}` then `execution.declared {predecessor}` | stage, execution |
| edge selected, loop restart, the condition that matched | `edge.selected`, `loop.restart`, the decision renderer's `condition`, the `run events --pretty` transition line | `route.applied` (`derived.target`, `transition`, `back`); its `edge` keys `subject.node.meta.edges`, whose entry carries the edge's `condition` as written (absent on an unconditional edge); a restart is `execution.finished {restart}` then `execution.declared {predecessor}` | stage, execution |
| notes | `StageCompletion.notes` | `step.finished` `outcome` notes; `parsed.note {result_prepared, transition}` | attempt |
| files touched | `stage.completed` `files_touched` | Pebble's fold of envelope `ToolCallCompleted` (see Agent activity) | session |
| stage diff | `StageProjection.diff` | platform record `checkpoint {execution, firing, patch_blob}` | stage |
@ -178,9 +179,9 @@ interviews.
| Fabro fact | Fields | Source | Keyed on |
| --- | --- | --- | --- |
| pending | `pending_interviews[id] {question, started_at}`, `current_question`, `interview.started` | `step.progress.recorded` with `parsed.question` (`id`, `text`, `options[] {key, label}`, `default`, `freeform`, `sensitive`, `kind`, `reference {label, url, kind}`, `timeout_ms`); `wait.state.changed {awaiting_answer}`; pending until a closing row below | question |
| pending | `pending_interviews[id] {question, started_at}`, `current_question`, `interview.started` | `step.progress.recorded` with `parsed.question` (`id`, `text`, `options[] {key, label, description, preview}`, `default`, `freeform`, `sensitive`, `kind`, `reference {label, url, kind}`, `timeout_ms`, `context`); `wait.state.changed {awaiting_answer}`; pending until a closing row below | question |
| question fields | `InterviewQuestionRecord.id`, `text`, `stage`, `question_type`, `options`, `allow_freeform`, `timeout_seconds`, `review_target` | `parsed.question`: `kind` is `question_type`, `freeform` is `allow_freeform`, `reference` is `review_target`, `timeout_ms` is `timeout_seconds`; `stage` is the subject's label | question |
| option description and preview, context display | `InterviewOption.description`, `preview`, `context_display` | gap | question |
| option description and preview, context display | `InterviewOption.description`, `preview`, `context_display` | `parsed.question`: each option's `description` and `preview`, the question's `context` (a human gate reads them from its edges' `human.description` and `human.preview` and from the previous stage's response; a native agent's question carries Pebble's); `reference` is `review_target` when Fabro's validation admits it | question |
| answered | `interview.completed {answer, duration_ms}`, the `actor` | `control.requested` with `derived.answer` and `derived.deliverable = true`; a sensitive answer stays `{"$secret": "answer:<id>"}`; `wait.state.changed {running}` follows; duration is `control.requested` minus the question's `recorded_at`; the actor is platform record `interview.answered {question, principal}` | question |
| late answer | none today | `control.requested` with `derived.deliverable = false` | question |
| expired | `interview.timeout` | `parsed.question_expired {question, waited_ms, default}`; the gate's `step.finished` follows (success with the default, else class `retry_requested`) | question |
@ -237,7 +238,7 @@ keeps its shape.
| route and failover | `route`, `failovers[]`, `failover_stopped`, `prompt.failover` | `custom attractor.fallback.plan {requested, routes, notices}`; envelope `RouteFailover {from, to, attempt, usage, error, continuation}`, `RouteFailoverStopped {route, reason, error}`; `crates/petri/lib/tests/fallback_events.rs` is the rebuild | attempt, session |
| messages, tokens, cost | `messages`, `usage`, `agent.message {text, usage, tool_call_count}`, `prompts` | envelope `AssistantMessage {usage, tool_call_count, …}`; sum per session; the stage total is `pebble.usage` | session |
| tool calls | `tools{name: {calls, errors, open}}`, `agent.tool.started`, `agent.tool.completed`, `files_touched`, `last_file_touched`, `pending_writes` | envelope `ToolCallStarted {tool_name, tool_call_id, arguments}`, `ToolCallCompleted {tool_call_id, is_error, error_kind}`; Fabro's own run tools appear the same way (the `HostTools` capability) | tool call |
| tools available | `agent_tools` (`ToolSummary {name, description, source, category, invoked}`), `agent.tools.available` | gap for the list; `invoked` derives from `ToolCallStarted` | session |
| tools available | `agent_tools` (`ToolSummary {name, description, source, category, invoked}`), `agent.tools.available`, the `run events --pretty` tool count line | `custom attractor.tools {session, tools[] {name, description, source, category}}`, once per native session (the node's own, then each child session); the stage's list is the union by name; `source` is Pebble's as recorded; `category` is Pebble's class only for a `subagent` tool, `other` for the rest (the payload carries Petri's origin category, not Pebble's permission class); `invoked` derives from envelope `ToolCallStarted` | session |
| MCP servers | `mcp_servers{}`, `agent.mcp.*` | envelope `McpServerReady {server, tools, startup_ms}`, `McpServerFailed`, `McpServerDisconnected`; `custom attractor.mcp.unavailable {server, error}` | session |
| skills | `skills.available`, `skills.activated` | `custom attractor.skills` (directories and sources), `attractor.skills.warning`; envelope `SkillsDiscovered`, `SkillActivated` | session |
| sub-agents | `subagents[]`, `subagent_counts`, `descendants` | envelope `SubAgentSpawned {agent_id, depth, task}`, `SubAgentTurnStarted`, `SubAgentCompleted`, `SubAgentFailed`, `SubAgentClosed` under the parent session; the child's events under its own session with `parent_session_id`; `pebble.subagents` on `step.finished` | session |
@ -339,7 +340,7 @@ where it belongs to a stage. The proposed `record_json` fields follow.
| `live_inference_ms`, `live_tool_ms`, `tool_batch`, `inference`, `acp_started_at` | envelope brackets (Agent activity); `step.started` for ACP |
| `usage`, `usage_by_model`, `model` | `pebble.usage`, `prompt.usage`, `pebble.subagents.sessions`, envelope `AssistantMessage` per session |
| `permission_level` | the node config |
| `agent_tools` | gap |
| `agent_tools` | `custom attractor.tools` per session; `invoked` from envelope `ToolCallStarted` |
| `agent` | Pebble's fold over the stage's envelopes, unchanged |
| `state` | the Stages state rows |
@ -370,7 +371,7 @@ and served on the events stream, and no view row reads them.
| `cancel.requested`, `kill.requested` | Run summary: cancel reason; Questions: interrupted; Parallel: cancelled fork |
| `control.requested`, `derived.deliverable`, `derived.answer` | Questions: answered, late, interrupted, steer, interrupt; Agent activity: pair |
| `token.emitted` | not shown: the engine's token flow; `visit.started` carries the join's inputs |
| `route.applied` | Stages: edge selected; Run summary: current stage |
| `route.applied` | Stages: edge selected and the condition that matched; Run summary: current stage |
| `node.expanded` | Parallel: fork started (`for_each`) |
| `visit.started`, `visit.completed` | Stages: list, state, timing, retries; Run summary: current stage |
| `wait.state.changed` | Stages: state; Questions: pending; Run summary: blocked |
@ -381,6 +382,7 @@ and served on the events stream, and no view row reads them.
| `parsed.note` `hook`, `hook.activity`, `parsed.hook_activity` | Stages: hook decisions; Agent activity: hook agents |
| `custom attractor.prompt`, `attractor.prompt.completed` | Stages: prompt, response; Parallel: fan-in prompt |
| `custom attractor.thread` | Agent activity: sessions |
| `custom attractor.tools` | Agent activity: tools available |
| `custom attractor.fallback.plan` | Agent activity: route; Stages: provider and model; Run summary: models |
| `custom attractor.mcp.unavailable` | Agent activity: MCP servers |
| `custom attractor.skills`, `attractor.skills.warning` | Agent activity: skills |
@ -425,8 +427,6 @@ record where Fabro does.
| Fact | Views | Smallest source |
| --- | --- | --- |
| tools available to an agent | `agent_tools`, the insights sidebar's tool list | a `custom attractor.tools {node, firing, attempt, session, tools[] {name, description, source, category}}` from the native backend once per session, where it calls the `HostTools` builders; Pebble's `SessionStarted` carries only the provider and model |
| question option `description` and `preview`, `context_display` | the interview dock, the human Q&A renderer | optional fields on Petri's `QuestionOption` (`description`, `preview`) and `Question` (`context`), set by the human gate from the edge attributes Fabro's lowering already reads |
| who answered | `interview.completed` `actor`, Slack attribution | platform record `interview.answered {question, principal, channel}` written by Fabro's interviewer beside its `InterviewReply` |
| run branch and base sha | `StartRecord`, `run diff`, the commits picker | platform record `run.branch {run_branch, base_sha}` written when Fabro creates the run branch, at that checkpoint's stage position |
| Git identity | `git_identity` | platform record `git.identity {name, email, source}` |

View file

@ -36,12 +36,13 @@ use fabro_types::{
BlockedReason, CheckpointRecord as ViewCheckpoint, CodingAgentEvent, CodingEvent, Conclusion,
FailureCategory, FailureDetail, FailureReason, InterviewOption, InterviewQuestionRecord,
ModelRef, ModelUsage, ParallelBranchId, ParallelBranchResult, PendingInterviewRecord,
PullRequestCreation, PullRequestCreationStatus, PullRequestLink, RunApproval, RunApprovalState,
RunArtifact, RunControlAction, RunDiff, RunFailure, RunId, RunProjection, RunSandbox,
RunSandboxFailure, RunSandboxInstance, RunSandboxPlan, RunSandboxRuntime, RunStatus, RunTiming,
SandboxProviderKind, StageCompletion, StageHandler, StageId, StageInferenceProjection,
StageModelUsage, StageOutcome, StageProjection, StageState, StageTiming, StartRecord,
SuccessReason, first_event_seq, format_blob_ref, parse_blob_ref, timing, usage_rollup,
PullRequestCreation, PullRequestCreationStatus, PullRequestLink, ReviewTarget,
ReviewTargetKind, RunApproval, RunApprovalState, RunArtifact, RunControlAction, RunDiff,
RunFailure, RunId, RunProjection, RunSandbox, RunSandboxFailure, RunSandboxInstance,
RunSandboxPlan, RunSandboxRuntime, RunStatus, RunTiming, SandboxProviderKind, StageCompletion,
StageHandler, StageId, StageInferenceProjection, StageModelUsage, StageOutcome,
StageProjection, StageState, StageTiming, StartRecord, SuccessReason, ToolCategory, ToolSource,
ToolSummary, first_event_seq, format_blob_ref, parse_blob_ref, timing, usage_rollup,
};
use lithos_llm::catalog::{ModelId, ProviderId};
use lithos_llm::types::Usage;
@ -49,6 +50,7 @@ use petri_execution::events::{Derived, NodeRef, Parsed, RunEvent, Subject, ViewE
use petri_execution::{CoordinatorEvent, ExecutionId, InvocationId};
use petri_runtime::engine::{Admission, Event};
use petri_runtime::ir::{Metrics, SandboxInstance, Status, StepEvent};
use petri_runtime::steps::QuestionReference;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use tracing::debug;
@ -848,6 +850,27 @@ impl RunView {
}
}
}
// The tools a native session was offered, once per
// session (VIEWS.md "Agent activity", tools available):
// the stage's list is the union over its sessions, by
// name, in the order the sessions listed them.
"attractor.tools" => {
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
let tools = payload.get("tools").and_then(Value::as_array);
for tool in tools.into_iter().flatten() {
let Some(summary) = tool_summary(tool) else {
continue;
};
if !stage
.agent_tools
.iter()
.any(|known| known.name == summary.name)
{
stage.agent_tools.push(summary);
}
}
}
}
"attractor.parallel.branch.started" => {
let invocation = payload.get("invocation").and_then(Value::as_u64);
let index = payload
@ -916,16 +939,16 @@ impl RunView {
.map(|option| InterviewOption {
key: option.key.clone(),
label: option.label.clone(),
description: None,
preview: None,
description: option.description.clone(),
preview: option.preview.clone(),
})
.collect(),
allow_freeform: question.freeform,
timeout_seconds: question
.timeout_ms
.map(|timeout| timeout as f64 / 1000.0),
context_display: None,
review_target: None,
context_display: question.context.clone(),
review_target: question.reference.as_ref().and_then(review_target),
},
started_at: at,
});
@ -1004,6 +1027,16 @@ impl RunView {
if stage.completion.is_none() {
stage.usage = agent.usage.saturating_add(agent.descendant_usage());
}
// A tool the stage's list names was called, by any of its sessions.
if let CodingEvent::ToolCallStarted { tool_name, .. } = &envelope.event {
if let Some(tool) = stage
.agent_tools
.iter_mut()
.find(|tool| tool.name == *tool_name)
{
tool.invoked = true;
}
}
let is_root = envelope.parent_session_id.is_none();
#[expect(
clippy::wildcard_enum_match_arm,
@ -1589,6 +1622,48 @@ fn failure_message(status: &Status) -> Option<String> {
/// The finished attempt's metrics onto its stage: the timing and the usage
/// the backend reported.
/// One tool of an `attractor.tools` payload as the stage's list carries
/// it: the name and description as recorded, Pebble's `source` as it is,
/// and Pebble's behavioural category where Petri's says which (a
/// sub-agent tool); every other tool is `other`, because the payload
/// carries Petri's origin category (`builtin`, `mcp`, `host`, `question`),
/// not Pebble's permission class. `invoked` starts false and flips on the
/// session's `ToolCallStarted`.
fn tool_summary(tool: &Value) -> Option<ToolSummary> {
let name = tool.get("name").and_then(Value::as_str)?;
let description = tool
.get("description")
.and_then(Value::as_str)
.unwrap_or_default();
let source = tool
.get("source")
.cloned()
.and_then(|source| serde_json::from_value::<ToolSource>(source).ok())
.unwrap_or_default();
let category = match tool.get("category").and_then(Value::as_str) {
Some("subagent") => ToolCategory::Subagent,
_ => ToolCategory::Other,
};
Some(ToolSummary {
name: name.to_string(),
description: description.to_string(),
source,
category,
invoked: false,
})
}
/// The question's `reference` as Fabro's review target, when it is one
/// Fabro's validation admits (a `document`, or a reference without a kind,
/// with a label and an absolute HTTP URL within Fabro's limits).
fn review_target(reference: &QuestionReference) -> Option<ReviewTarget> {
let kind = match reference.kind.as_deref() {
Some("document") | None => ReviewTargetKind::Document,
Some(_) => return None,
};
ReviewTarget::new(reference.label.clone(), reference.url.clone(), kind).ok()
}
fn apply_metrics(stage: &mut StageProjection, metrics: &Metrics) {
let custom = &metrics.custom;
let inference = custom

View file

@ -23,9 +23,14 @@ use std::sync::Arc;
use std::time::Duration;
use fabro_petri::host_tools::recorded::{self, ExecutionId, InvocationId};
use fabro_petri::projection::{Item, RunView};
use fabro_petri::runtime::RuntimeSpec;
use fabro_store::platform_records::{PlatformRecord, RunCreatedRecord, StoredPlatformRecord};
use fabro_tool::fabro_client::ClientBackend;
use fabro_types::{BlobHash, RunId, WorkflowVersionId};
use fabro_types::{
BlobHash, RunId, StageId, ToolCategory, ToolSource, WorkflowVersionId,
test_support as types_support,
};
use fabro_workflow::run_tools::register_fabro_run_tools;
use fabro_workflow::services::FabroRunToolServices;
use httpmock::{Method, MockServer};
@ -33,12 +38,13 @@ use lithos_llm::types::Request;
use pebble_coding_agent::test_support::{
ScriptedCall, ScriptedProvider, scripted_client, text_response, tool_call_response,
};
use petri_execution::events::replay_run;
use petri_execution::host::{self, HostRun};
use petri_runtime::executor::Retention;
use petri_runtime::frontend::CompileInputs;
use petri_runtime::ir::RunStatus;
use petri_runtime::{RunOptions, Runtime};
use petri_store::{MemoryRunStore, RunKey, RunStore};
use petri_store::{Access, MemoryRunStore, RunKey, RunStore};
use serde_json::json;
use tokio::fs;
@ -258,6 +264,112 @@ async fn a_petri_stage_calls_a_run_tool_bound_to_the_run() {
assert_eq!(calls[0].payload["is_error"], true, "{:?}", calls[0].payload);
}
/// The stage's projection lists the tools its session was offered, host
/// tools included, as the `attractor.tools` payload records them: the
/// same names the model was advertised, each with its description and
/// Pebble's source, the sub-agent tools under Pebble's category, and the
/// one tool the model called marked invoked.
#[tokio::test]
async fn the_projection_lists_the_stages_tools_and_marks_the_one_called() {
if host_plugin().is_none() {
return;
}
let root = tempfile::tempdir().expect("a temp dir");
let workflow = install_bundle(root.path()).await;
let run_id = RunId::new();
let version_id: WorkflowVersionId = BlobHash::new(b"child workflow").into();
let server = MockServer::start_async().await;
server
.mock_async(|when, then| {
when.method(Method::POST).path("/api/v1/runs");
then.status(422).body("native admission rejection");
})
.await;
let (client, provider) = scripted_model(&version_id.to_string());
let store = Arc::new(MemoryRunStore::new());
let rt = runtime(
&root.path().join("run"),
&run_id.to_string(),
client,
Some(services(&server, run_id)),
&store,
);
run(&rt, &workflow).await;
// The run's events, folded as the projector folds them: the run's
// `run.created` record first, then Petri's events in record order.
let logs = store
.open(&RunKey::new(run_id.to_string()), Access::Read)
.await
.expect("the run opens");
let events = replay_run(&*logs).await.expect("the record replays");
let mut spec = types_support::test_run_spec();
spec.run_id = run_id;
let created = StoredPlatformRecord {
seq: 1,
recorded_at: 0,
record: PlatformRecord::RunCreated(RunCreatedRecord {
spec,
title: Some("Start a child run".to_string()),
parent_id: None,
retried_from: None,
web_url: None,
}),
position: None,
};
let mut view = RunView::new();
view.fold(&Item::Platform(&created), 1);
for (index, event) in events.iter().enumerate() {
view.fold(&Item::Petri(event), index as u64 + 2);
}
let projection = view.projection().expect("the run has a projection");
let stage = projection
.stage(&StageId::new("work", 1))
.expect("the agent stage is projected");
let mut listed: Vec<(String, String)> = stage
.agent_tools
.iter()
.map(|tool| (tool.name.clone(), tool.description.clone()))
.collect();
listed.sort();
let requests = provider.requests();
let mut offered = advertised(&requests[0]);
offered.sort();
assert_eq!(
listed, offered,
"the list is what the model was offered, by name and description"
);
let create = stage
.agent_tools
.iter()
.find(|tool| tool.name == "fabro_run_create")
.expect("the host tool is listed");
assert!(create.invoked, "the model called it");
assert_eq!(create.source, ToolSource::Application);
assert_eq!(create.category, ToolCategory::Other);
let shell = stage
.agent_tools
.iter()
.find(|tool| tool.name == "shell")
.expect("Pebble's own tool is listed");
assert!(!shell.invoked, "the model never called it");
assert_eq!(shell.source, ToolSource::Native);
assert!(
stage
.agent_tools
.iter()
.any(|tool| tool.category == ToolCategory::Subagent),
"Pebble's sub-agent tools keep their category: {:?}",
stage
.agent_tools
.iter()
.map(|tool| (&tool.name, tool.category))
.collect::<Vec<_>>()
);
}
/// Services bound to another run give the stage no run tools: the model
/// is not advertised them, and its call is refused as an unknown tool
/// rather than parenting a child run to the wrong run.

View file

@ -104,6 +104,29 @@ fn gate_workflow(markers: &Path, gate_attrs: &str) -> String {
)
}
/// A gate after a stage with a response, whose affirmative edge says what
/// choosing it means and shows a sample: the facts the interview dock
/// shows beside the choices.
fn described_gate_workflow(markers: &Path) -> String {
format!(
r#"digraph Gate {{
graph [goal="Ask with context"]
start [shape=Mdiamond]
exit [shape=Msquare]
plan [shape=parallelogram, output_schema="routing", script="echo '{{\"outcome\": \"succeeded\", \"context_updates\": {{\"last_stage\": \"plan\", \"response.plan\": \"Ship the fix in one commit.\"}}}}'"]
gate [shape=hexagon, label="Deploy?", timeout="1500ms", human.default_choice="no"]
yes [shape=parallelogram, script="touch {dir}/yes"]
no [shape=parallelogram, script="touch {dir}/no"]
start -> plan -> gate
gate -> yes [label="[Y] Yes", "human.description"="Merge and deploy to production", "human.preview"="deploy --prod"]
gate -> no [label="[N] No"]
yes -> exit
no -> exit
}}"#,
dir = markers.display()
)
}
fn host_plugin() -> Option<PathBuf> {
let found = env::var_os(HOST_PLUGIN_OVERRIDE)
.map(PathBuf::from)
@ -1156,12 +1179,16 @@ struct GateRun {
}
async fn gate_run(gate_attrs: &str) -> GateRun {
gate_run_of(|markers| gate_workflow(markers, gate_attrs)).await
}
async fn gate_run_of(workflow: impl FnOnce(&Path) -> String) -> GateRun {
let root = tempfile::tempdir().expect("a marker dir");
let markers = root.path().join("markers");
fs::create_dir_all(&markers)
.await
.expect("the marker dir creates");
let workflow = gate_workflow(&markers, gate_attrs);
let workflow = workflow(&markers);
let scenario = scenario(
"gate",
&[("workflow.fabro", &workflow), ("workflow.toml", SETTINGS)],
@ -1209,6 +1236,24 @@ impl GateRun {
self.projector.settle(self.scenario.run_id).await;
}
/// The stored projection once a question is pending in it.
async fn pending(&self) -> fabro_types::RunProjection {
let deadline = Instant::now() + Duration::from_secs(30);
loop {
let stored = projector::stored_projection(&self.scenario.pool, self.scenario.run_id)
.await
.expect("the stored projection reads");
if let Some(stored) = stored.filter(|stored| !stored.pending_interviews.is_empty()) {
return stored;
}
assert!(
Instant::now() < deadline,
"the question never showed as pending"
);
sleep(Duration::from_millis(10)).await;
}
}
async fn stored(&self) -> fabro_types::RunProjection {
projector::stored_projection(&self.scenario.pool, self.scenario.run_id)
.await
@ -1299,6 +1344,55 @@ async fn an_expired_question_is_pending_while_the_gate_waits_and_closes_on_the_e
assert_view_equals_rebuild(&gate.scenario.pool, gate.scenario.run_id).await;
}
/// A gate's choices carry what choosing them means and a sample of what
/// they would do, and the question carries the previous stage's response
/// as its context: the pending question in the projection shows all
/// three, as Petri's question record carries them, and leaves them absent
/// on a choice that has none.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_pending_question_carries_its_choice_descriptions_previews_and_context() {
if host_plugin().is_none() {
return;
}
let gate = Arc::new(gate_run_of(described_gate_workflow).await);
let running = {
let gate = Arc::clone(&gate);
tokio::spawn(async move { gate.run(Approval::Prompt).await })
};
let pending = gate.pending().await;
let (_, record) = pending
.pending_interviews
.iter()
.next()
.expect("one pending question");
let question = &record.question;
assert_eq!(question.stage, "gate@1");
assert_eq!(question.text, "Deploy?");
assert_eq!(
question.context_display.as_deref(),
Some("Ship the fix in one commit."),
"the context is the previous stage's response"
);
assert_eq!(question.options.len(), 2, "{:?}", question.options);
assert_eq!(question.options[0].key, "Y");
assert_eq!(
question.options[0].description.as_deref(),
Some("Merge and deploy to production")
);
assert_eq!(
question.options[0].preview.as_deref(),
Some("deploy --prod")
);
assert_eq!(question.options[1].key, "N");
assert_eq!(question.options[1].description, None);
assert_eq!(question.options[1].preview, None);
assert!(question.review_target.is_none());
running.await.expect("the run task ends");
assert!(gate.markers.join("no").exists(), "the default ran");
assert_view_equals_rebuild(&gate.scenario.pool, gate.scenario.run_id).await;
}
/// An auto-approved run answers its gate at once: the delivered answer
/// closes the question in the projection, the affirmative branch runs, and
/// the view rebuilds the same.