Add the artifact, run diff and branch workspace platform records

`artifact.collected {execution, firing, attempt, path, blob, bytes, digest}`
records one file a stage left in its workspace, with the bytes in the blob
table; `run.diff {base_sha, head_sha, diff_summary, patch_blob}` records
the run branch against its base; `run.branch` names the workspace the
branch was created in. `RunProjection.artifacts` lists the collected
files, and `parse_blob_ref_encoded` reads Petri's `#json` marker so a
reader knows whether a blob is text or a JSON value.

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 2f4888c199
commit dc48197553
No known key found for this signature in database
6 changed files with 165 additions and 37 deletions

View file

@ -12,8 +12,10 @@ over both:
before and after the engine (`run.created`, `run.lifecycle`, `run.title`,
`run.parent`, `run.archived`, `run.superseded`, `run.notice`), who answered
a question (`interview.answered`), the branch and git identity a run works
under, a checkpoint commit, the pull request requests and outcomes, a
notification sent, a pairing. They are `PlatformRecord` values in
under, a checkpoint commit with its diff, a collected artifact
(`artifact.collected`), the run's diff (`run.diff`), the pull request
requests and outcomes, a notification sent, a pairing. They are
`PlatformRecord` values in
`fabro-store::platform_records`, stored in `platform_records` with a
per-run `seq`.

View file

@ -15,9 +15,10 @@
//! today still append Fabro's legacy run events: for a Petri run the run
//! summary store derives the platform record from the legacy event through
//! [`platform_record_for`] and stores both in the event's transaction. The
//! `checkpoint`, `pull_request.created`, `notification.sent` and
//! `run.paired` kinds are defined here and written by the adapters that
//! perform those effects.
//! `run.branch`, `git.identity`, `checkpoint`, `artifact.collected`,
//! `run.diff`, `pull_request.created`, `notification.sent` and `run.paired`
//! kinds are defined here and written by the adapters that perform those
//! effects.
//!
//! Every record may carry an [`OperationKey`]: the identity of the external
//! effect it records (the execution, the Petri decision and the effect
@ -120,6 +121,12 @@ pub enum PlatformRecordKind {
#[serde(rename = "checkpoint")]
#[strum(serialize = "checkpoint")]
Checkpoint,
#[serde(rename = "artifact.collected")]
#[strum(serialize = "artifact.collected")]
ArtifactCollected,
#[serde(rename = "run.diff")]
#[strum(serialize = "run.diff")]
RunDiff,
#[serde(rename = "pull_request.requested")]
#[strum(serialize = "pull_request.requested")]
PullRequestRequested,
@ -183,6 +190,14 @@ pub enum PlatformRecord {
/// record, written after the commit succeeds.
#[serde(rename = "checkpoint")]
Checkpoint(CheckpointRecord),
/// A file a stage's attempt left in its workspace, collected under
/// `[run.artifacts] include` into the blob table.
#[serde(rename = "artifact.collected")]
ArtifactCollected(ArtifactCollectedRecord),
/// The run's whole diff, its run branch against its base commit, written
/// when the run finishes.
#[serde(rename = "run.diff")]
RunDiff(RunDiffRecord),
/// A pull request was asked for: the supervisor creates it.
#[serde(rename = "pull_request.requested")]
PullRequestRequested(PullRequestRequestedRecord),
@ -218,6 +233,8 @@ impl PlatformRecord {
Self::RunBranch(_) => PlatformRecordKind::RunBranch,
Self::GitIdentity(_) => PlatformRecordKind::GitIdentity,
Self::Checkpoint(_) => PlatformRecordKind::Checkpoint,
Self::ArtifactCollected(_) => PlatformRecordKind::ArtifactCollected,
Self::RunDiff(_) => PlatformRecordKind::RunDiff,
Self::PullRequestRequested(_) => PlatformRecordKind::PullRequestRequested,
Self::PullRequestCreated(_) => PlatformRecordKind::PullRequestCreated,
Self::PullRequestFailed(_) => PlatformRecordKind::PullRequestFailed,
@ -234,6 +251,7 @@ impl PlatformRecord {
pub fn operation(&self) -> Option<&OperationKey> {
match self {
Self::Checkpoint(record) => record.operation.as_ref(),
Self::ArtifactCollected(record) => record.operation.as_ref(),
Self::PullRequestCreated(record) => record.operation.as_ref(),
Self::NotificationSent(record) => record.operation.as_ref(),
Self::RunCreated(_)
@ -247,6 +265,7 @@ impl PlatformRecord {
| Self::InterviewAnswered(_)
| Self::RunBranch(_)
| Self::GitIdentity(_)
| Self::RunDiff(_)
| Self::PullRequestRequested(_)
| Self::PullRequestFailed(_)
| Self::PullRequestLinked(_)
@ -396,6 +415,10 @@ pub struct RunBranchRecord {
pub run_branch: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub base_sha: Option<String>,
/// The Petri workspace the branch was created in: where the run's diff
/// is measured.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub workspace: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
@ -425,6 +448,39 @@ pub struct CheckpointRecord {
pub operation: Option<OperationKey>,
}
/// One file collected from a stage's workspace after its attempt finished.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ArtifactCollectedRecord {
pub execution: u64,
pub firing: u64,
/// The attempt whose workspace the file was read from, 1-based.
pub attempt: u32,
/// The file's path relative to the workspace root.
pub path: String,
/// The blob that holds the file's bytes.
pub blob: BlobHash,
pub bytes: u64,
/// The SHA-256 of the bytes as lowercase hex: with `path`, the identity
/// a later capture of the same unchanged file is matched by.
pub digest: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub operation: Option<OperationKey>,
}
/// The run's diff: its run branch's head against its base commit.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RunDiffRecord {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub base_sha: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub head_sha: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub diff_summary: Option<DiffSummary>,
/// The patch as a text blob; absent when the diff is empty.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub patch_blob: Option<BlobHash>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PullRequestCreatedRecord {
pub number: u64,
@ -735,6 +791,7 @@ mod tests {
PlatformRecordKind::RunBranch => PlatformRecord::RunBranch(RunBranchRecord {
run_branch: Some("fabro/run-1".to_string()),
base_sha: Some("abc".to_string()),
workspace: Some("invocation-0-scope-0".to_string()),
}),
PlatformRecordKind::GitIdentity => PlatformRecord::GitIdentity(GitIdentityRecord {
identity: GitIdentity {
@ -764,6 +821,35 @@ mod tests {
effect: "commit".to_string(),
}),
}),
PlatformRecordKind::ArtifactCollected => {
PlatformRecord::ArtifactCollected(ArtifactCollectedRecord {
execution: 0,
firing: 3,
attempt: 1,
path: "assets/report.txt".to_string(),
blob: BlobHash::new(b"report"),
bytes: 6,
digest: BlobHash::new(b"report").to_string(),
operation: Some(OperationKey {
execution: 0,
decision: DecisionRef::AttemptStart {
firing: 3,
attempt: 1,
},
effect: "artifact".to_string(),
}),
})
}
PlatformRecordKind::RunDiff => PlatformRecord::RunDiff(RunDiffRecord {
base_sha: Some("abc".to_string()),
head_sha: Some("def".to_string()),
diff_summary: Some(DiffSummary {
files_changed: 1,
additions: 2,
deletions: 0,
}),
patch_blob: Some(BlobHash::new(b"patch")),
}),
PlatformRecordKind::PullRequestCreated => {
PlatformRecord::PullRequestCreated(PullRequestCreatedRecord {
number: 7,

View file

@ -14,6 +14,29 @@ pub fn parse_blob_ref(value: &str) -> Option<BlobHash> {
value.strip_prefix(BLOB_REF_PREFIX)?.parse().ok()
}
/// How the bytes behind a blob reference decode back into a value.
///
/// Petri stores a large string as its own bytes and marks a structured value
/// with a `#json` suffix on the reference (`blob://sha256/<hex>#json`), so a
/// reader knows whether to parse the bytes or take them as text.
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum BlobRefEncoding {
/// The bytes are the text of a string value.
Text,
/// The bytes are compact JSON of a structured value.
Json,
}
/// A blob reference with its encoding: a plain reference is text, one with
/// the `#json` suffix is JSON.
#[must_use]
pub fn parse_blob_ref_encoded(value: &str) -> Option<(BlobHash, BlobRefEncoding)> {
match value.strip_suffix("#json") {
Some(body) => parse_blob_ref(body).map(|hash| (hash, BlobRefEncoding::Json)),
None => parse_blob_ref(value).map(|hash| (hash, BlobRefEncoding::Text)),
}
}
#[must_use]
pub fn parse_managed_blob_file_ref(value: &str) -> Option<BlobHash> {
let path = value.strip_prefix("file://")?;
@ -45,7 +68,10 @@ fn has_path_suffix(path: &str, suffix: &[&str]) -> bool {
#[cfg(test)]
mod tests {
use super::{format_blob_ref, parse_blob_ref, parse_managed_blob_file_ref};
use super::{
BlobRefEncoding, format_blob_ref, parse_blob_ref, parse_blob_ref_encoded,
parse_managed_blob_file_ref,
};
use crate::BlobHash;
#[test]
@ -56,6 +82,21 @@ mod tests {
assert_eq!(parse_blob_ref(&formatted), Some(blob_hash));
}
#[test]
fn a_json_suffix_names_the_encoding() {
let blob_hash = BlobHash::new(b"text");
let formatted = format_blob_ref(&blob_hash);
assert_eq!(
parse_blob_ref_encoded(&formatted),
Some((blob_hash, BlobRefEncoding::Text))
);
assert_eq!(
parse_blob_ref_encoded(&format!("{formatted}#json")),
Some((blob_hash, BlobRefEncoding::Json))
);
assert_eq!(parse_blob_ref_encoded("not a reference"), None);
}
#[test]
fn managed_local_blob_file_ref_is_recognized() {
let blob_hash = BlobHash::new(b"hello");

View file

@ -72,7 +72,10 @@ pub use agent_props::{
pub use artifact::ArtifactUpload;
pub use auth::{IdpIdentity, IdpIdentityError};
pub use blob_hash::BlobHash;
pub use blob_ref::{format_blob_ref, parse_blob_ref, parse_managed_blob_file_ref};
pub use blob_ref::{
BlobRefEncoding, format_blob_ref, parse_blob_ref, parse_blob_ref_encoded,
parse_managed_blob_file_ref,
};
pub use catalog_api::{Model, ModelControls, ModelCosts, ModelFeatures, ModelLimits, Provider};
pub use checkpoint::Checkpoint;
pub use command_output::{CommandOutputStream, CommandTermination};
@ -145,7 +148,7 @@ pub use run_intent::{
TargetValidationError, ValidatedGitRunTarget, ValidatedRunTarget,
};
pub use run_projection::{
CheckpointRecord, PendingInterviewRecord, RunProjection, StageContextWindow,
CheckpointRecord, PendingInterviewRecord, RunArtifact, RunProjection, StageContextWindow,
StageContextWindowUnavailableReason, StageInferenceProjection, StageModelUsage,
StageProjection, StageToolBatchProjection, first_event_seq,
};

View file

@ -13,10 +13,11 @@ use strum::{Display, EnumString, IntoStaticStr};
use crate::agent_props::{AgentSessionActivatedProps, StagePromptProps};
use crate::{
AgentBackend, Checkpoint, Conclusion, GitIdentity, InterviewQuestionRecord, InvalidTransition,
ModelRef, ModelUsage, ParallelBranchId, PullRequestCreation, PullRequestLink, RunApproval,
RunControlAction, RunDiff, RunId, RunSandbox, RunSpec, RunStatus, RunTiming, StageCompletion,
StageHandler, StageId, StageState, StageTiming, StartRecord, timing,
AgentBackend, BlobHash, Checkpoint, Conclusion, GitIdentity, InterviewQuestionRecord,
InvalidTransition, ModelRef, ModelUsage, ParallelBranchId, PullRequestCreation,
PullRequestLink, RunApproval, RunControlAction, RunDiff, RunId, RunSandbox, RunSpec, RunStatus,
RunTiming, StageCompletion, StageHandler, StageId, StageState, StageTiming, StartRecord,
timing,
};
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
@ -51,9 +52,28 @@ pub struct RunProjection {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub git_identity: Option<GitIdentity>,
pub pending_interviews: BTreeMap<String, PendingInterviewRecord>,
/// The files collected from the run's workspaces under
/// `[run.artifacts] include`, one entry per capture, in the order they
/// were recorded. The bytes are in the blob table under `blob`.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub artifacts: Vec<RunArtifact>,
stages: HashMap<StageId, StageProjection>,
}
/// One file a stage's attempt left in its workspace and the run collected.
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct RunArtifact {
/// The stage that produced the file, as the projection labels it.
pub stage_id: StageId,
/// The attempt of the stage, 1-based, as the artifact listing's `retry`.
pub retry: u32,
/// The file's path relative to the workspace root.
pub relative_path: String,
pub size: u64,
/// The blob that holds the file's bytes.
pub blob: BlobHash,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
pub struct PendingInterviewRecord {
pub question: InterviewQuestionRecord,
@ -694,6 +714,7 @@ impl RunProjection {
retried_from: None,
git_identity: None,
pending_interviews: BTreeMap::new(),
artifacts: Vec::new(),
stages: HashMap::new(),
}
}

View file

@ -1,25 +0,0 @@
/* tslint:disable */
/* eslint-disable */
/**
* Fabro Run API
* HTTP API for managing Fabro workflow run executions.
*
* The version of the OpenAPI document: 0.2.0
*
*
* NOTE: This class is auto generated by OpenAPI Generator (https://openapi-generator.tech).
* https://openapi-generator.tech
* Do not edit the class manually.
*/
// May contain unused imports in some cases
// @ts-ignore
import type { ObjectStoreSettings } from './object-store-settings';
export interface ServerSlateDbSettings {
'prefix': string;
'store': ObjectStoreSettings;
'flush_interval': string;
'disk_cache': boolean;
}