Fold artifacts, diffs, blob references and dry-run responses into the projection

The projection lists every `artifact.collected` record under the stage's
label, carries a checkpoint's patch blob as its `blob://` reference on the
stage and the checkpoint, and takes `run.diff` as the conclusion's diff,
whichever side of the run's finish it arrives on. A command's `stdout`
from `step.finished` becomes the stage's output, so an offloaded output
shows as its reference rather than the live log's bytes, and a simulated
prompt or agent stage carries the stub's text as its response.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 16:00:11 -04:00
parent 67ba595b01
commit 2407ce3d8d
No known key found for this signature in database
2 changed files with 140 additions and 10 deletions

View file

@ -37,10 +37,11 @@ use fabro_types::{
FailureCategory, FailureDetail, FailureReason, InterviewOption, InterviewQuestionRecord,
ModelRef, ModelUsage, ParallelBranchId, ParallelBranchResult, PendingInterviewRecord,
PullRequestCreation, PullRequestCreationStatus, PullRequestLink, RunApproval, RunApprovalState,
RunControlAction, RunDiff, RunFailure, RunId, RunProjection, RunSandbox, RunSandboxPlan,
RunStatus, RunTiming, SandboxProviderKind, StageCompletion, StageHandler, StageId,
StageInferenceProjection, StageModelUsage, StageOutcome, StageProjection, StageState,
StageTiming, StartRecord, SuccessReason, first_event_seq, timing, usage_rollup,
RunArtifact, RunControlAction, RunDiff, RunFailure, RunId, RunProjection, RunSandbox,
RunSandboxPlan, 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,
};
use lithos_llm::catalog::{ModelId, ProviderId};
use lithos_llm::types::Usage;
@ -129,6 +130,10 @@ pub struct FoldState {
pub base_sha: Option<String>,
#[serde(default)]
pub checkpoints: u32,
/// The run's diff as its `run.diff` record gave it, whichever side of
/// the run's finish it arrived on.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub run_diff: Option<RunDiff>,
#[serde(default)]
pub health: RecordHealth,
/// Firings (`"<execution>:<firing>"`) whose attempt has recorded a
@ -236,20 +241,62 @@ impl RunView {
.stages
.get(&stage_key(record.execution, record.firing));
let current_node = stage.map_or_else(String::new, |stage| stage.node_name.clone());
let stage_id = stage
.filter(|stage| stage.shown)
.map(|stage| stage.stage_id.clone());
let checkpoint = fabro_types::Checkpoint {
timestamp: at,
current_node: current_node.clone(),
git_commit_sha: record.git_commit_sha.clone(),
};
// The patch stays in the blob table; the view carries its
// reference for a reader to resolve.
let patch = record.patch_blob.as_ref().map(format_blob_ref);
if let Some(stage) = stage_id.and_then(|stage_id| projection.stage_mut(&stage_id)) {
if patch.is_some() {
stage.diff.clone_from(&patch);
}
}
projection.checkpoints.push(ViewCheckpoint {
seq: u32::try_from(stream_seq).unwrap_or(u32::MAX),
checkpoint,
diff: RunDiff {
patch: None,
patch,
summary: record.diff_summary,
},
});
}
PlatformRecord::ArtifactCollected(record) => {
let stage = self
.state
.stages
.get(&stage_key(record.execution, record.firing));
let Some(stage_id) = stage.map(|stage| stage.stage_id.clone()) else {
debug!(
seq = stored.seq,
path = record.path,
"artifact record for an unknown firing; not folded"
);
return;
};
projection.artifacts.push(RunArtifact {
stage_id,
retry: record.attempt,
relative_path: record.path.clone(),
size: record.bytes,
blob: record.blob,
});
}
PlatformRecord::RunDiff(record) => {
let diff = RunDiff {
patch: record.patch_blob.as_ref().map(format_blob_ref),
summary: record.diff_summary,
};
if let Some(conclusion) = projection.conclusion.as_mut() {
conclusion.diff = diff.clone();
}
self.state.run_diff = Some(diff);
}
PlatformRecord::PullRequestRequested(record) => {
projection.pull_request_creation = Some(PullRequestCreation {
id: record.creation_id,
@ -491,8 +538,11 @@ impl RunView {
stages,
usage: rollup.usage_if_present(),
total_retries,
diff: last_checkpoint
.map(|checkpoint| checkpoint.diff.clone())
diff: self
.state
.run_diff
.clone()
.or_else(|| last_checkpoint.map(|checkpoint| checkpoint.diff.clone()))
.unwrap_or_default(),
});
}
@ -546,9 +596,35 @@ impl RunView {
.as_ref()
.map(|subject| subject.node.name.to_string());
if let Some(stage) = self.stage_of(execution, event.subject.as_ref()) {
if let Some(output) = outcome.output.as_str() {
// The step's output: a string, or a command's `stdout`,
// either of which is a `blob://` reference when the
// step offloaded it. The reference stays as it is; the
// bytes it names are the live log's.
let output = outcome
.output
.as_str()
.or_else(|| outcome.output.get("stdout").and_then(Value::as_str));
if let Some(output) = output {
if parse_blob_ref(output).is_none() {
stage.output_bytes = Some(output.len() as u64);
}
stage.output = Some(output.to_string());
stage.output_bytes = Some(output.len() as u64);
}
// A simulated step (a dry run) answers with its text.
let simulated = outcome
.output
.get("simulated")
.and_then(Value::as_bool)
.unwrap_or(false);
if simulated
&& matches!(
stage.handler,
Some(StageHandler::Prompt | StageHandler::Agent)
)
{
if let Some(text) = outcome.output.get("text").and_then(Value::as_str) {
stage.response = Some(text.to_string());
}
}
// An agent's answer: the `response.<node>` the step wrote
// into the run context, as the prompt step writes it.

View file

@ -26,6 +26,7 @@ use std::time::{Duration, Instant};
use fabro_db::DbPool;
use fabro_interview::ControlInterviewer;
use fabro_petri::SqliteRunStore;
use fabro_petri::blobs::{Blobs, RunBlobs};
use fabro_petri::check::Launch;
use fabro_petri::engine::{self, RunStatus as EngineRunStatus};
use fabro_petri::interview::{Approval, FabroInterviewer};
@ -34,7 +35,7 @@ use fabro_petri::runtime::RuntimeSpec;
use fabro_store::platform_records::{
PlatformRecord, PlatformRecordStore, RunCreatedRecord, RunLifecycleKind, RunLifecycleRecord,
};
use fabro_store::test_support;
use fabro_store::{BlobStore, test_support};
use fabro_types::{
BlobHash, PetriAdmission, PetriGraphRef, RunId, RunStatus, StageHandler, StageId, StageState,
test_support as types_support,
@ -204,6 +205,7 @@ async fn create_run(pool: &DbPool, run_id: RunId, goal: &str) {
/// Run `workflow` to completion on the real registry over `store`.
async fn run_workflow(
store: Arc<dyn RunStore>,
blobs: Arc<dyn Blobs>,
run_dir: &Path,
run_id: RunId,
workflow: &Path,
@ -215,6 +217,9 @@ async fn run_workflow(
} else {
petri_attractor_steps::register(runtime)
};
// A large stage value goes to Fabro's blob table, as it does under the
// worker and the server.
let runtime = runtime.capability(RunBlobs::output_store(blobs));
let rt = runtime.store(store).options(run_options(run_dir, run_id));
let lowered = rt
.check(workflow, None, None, &CompileInputs::new())
@ -288,6 +293,27 @@ async fn command_scenario() -> Scenario {
.await
}
/// A command whose output is above Petri's offload threshold.
const LARGE_OUTPUT_WORKFLOW: &str = r#"digraph Large {
graph [goal="Print a lot"]
start [shape=Mdiamond]
exit [shape=Msquare]
big [shape=parallelogram, script="yes xxxxxxxxxxxxxxxx | head -n 8000"]
start -> big -> exit
}"#;
async fn large_output_scenario() -> Scenario {
scenario(
"large",
&[
("workflow.fabro", LARGE_OUTPUT_WORKFLOW),
("workflow.toml", SETTINGS),
],
false,
)
.await
}
async fn parallel_scenario() -> Scenario {
scenario(
"parallel",
@ -308,6 +334,7 @@ async fn run_live(scenario: &Scenario) -> Arc<Projector> {
let store = projector.observe_store(Arc::new(SqliteRunStore::new(scenario.pool.clone())));
run_workflow(
store,
Arc::new(BlobStore::new(scenario.pool.clone())),
&scenario.run_dir,
scenario.run_id,
&scenario.workflow,
@ -323,6 +350,7 @@ async fn run_live(scenario: &Scenario) -> Arc<Projector> {
async fn run_unobserved(scenario: &Scenario) {
run_workflow(
Arc::new(SqliteRunStore::new(scenario.pool.clone())),
Arc::new(BlobStore::new(scenario.pool.clone())),
&scenario.run_dir,
scenario.run_id,
&scenario.workflow,
@ -438,6 +466,32 @@ async fn the_hello_bundle_projects_live_as_it_rebuilds() {
);
}
/// A command's offloaded output reaches the view as its `blob://`
/// reference, never as the bytes the live log accumulated.
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_large_output_projects_as_its_blob_reference() {
if host_plugin().is_none() {
return;
}
let scenario = large_output_scenario().await;
run_live(&scenario).await;
assert_view_equals_rebuild(&scenario.pool, scenario.run_id).await;
let stored = projector::stored_projection(&scenario.pool, scenario.run_id)
.await
.expect("reads")
.expect("stored");
let big = stored
.stage(&StageId::new("big", 1))
.expect("the command stage is shown");
let output = big.output.as_deref().expect("the stage has an output");
assert!(
fabro_types::parse_blob_ref(output).is_some(),
"the output is a blob reference: {} bytes, {}",
output.len(),
&output[..output.len().min(80)]
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_command_workflow_projects_live_as_it_rebuilds() {
if host_plugin().is_none() {