diff --git a/lib/components/fabro-petri/src/projection.rs b/lib/components/fabro-petri/src/projection.rs index bf910fdd6..420695be3 100644 --- a/lib/components/fabro-petri/src/projection.rs +++ b/lib/components/fabro-petri/src/projection.rs @@ -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, #[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, #[serde(default)] pub health: RecordHealth, /// Firings (`":"`) 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.` the step wrote // into the run context, as the prompt step writes it. diff --git a/lib/components/fabro-petri/tests/projection.rs b/lib/components/fabro-petri/tests/projection.rs index a8c5f2e78..2e6ab7a71 100644 --- a/lib/components/fabro-petri/tests/projection.rs +++ b/lib/components/fabro-petri/tests/projection.rs @@ -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, + blobs: Arc, 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 { 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 { 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() {