diff --git a/docs/internal/events-strategy.md b/docs/internal/events-strategy.md index 16610a2e2..4be756797 100644 --- a/docs/internal/events-strategy.md +++ b/docs/internal/events-strategy.md @@ -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`. diff --git a/lib/components/fabro-store/src/platform_records.rs b/lib/components/fabro-store/src/platform_records.rs index 123cfa1a4..096cb8f58 100644 --- a/lib/components/fabro-store/src/platform_records.rs +++ b/lib/components/fabro-store/src/platform_records.rs @@ -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, #[serde(default, skip_serializing_if = "Option::is_none")] pub base_sha: Option, + /// 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, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] @@ -425,6 +448,39 @@ pub struct CheckpointRecord { pub operation: Option, } +/// 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, +} + +/// 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, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub head_sha: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub diff_summary: Option, + /// The patch as a text blob; absent when the diff is empty. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub patch_blob: Option, +} + #[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, diff --git a/lib/foundation/fabro-types/src/blob_ref.rs b/lib/foundation/fabro-types/src/blob_ref.rs index f413cd6ff..d120bfdb8 100644 --- a/lib/foundation/fabro-types/src/blob_ref.rs +++ b/lib/foundation/fabro-types/src/blob_ref.rs @@ -14,6 +14,29 @@ pub fn parse_blob_ref(value: &str) -> Option { 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/#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 { 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"); diff --git a/lib/foundation/fabro-types/src/lib.rs b/lib/foundation/fabro-types/src/lib.rs index e182b742e..d1620e43d 100644 --- a/lib/foundation/fabro-types/src/lib.rs +++ b/lib/foundation/fabro-types/src/lib.rs @@ -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, }; diff --git a/lib/foundation/fabro-types/src/run_projection.rs b/lib/foundation/fabro-types/src/run_projection.rs index c13077bf3..fb0d06428 100644 --- a/lib/foundation/fabro-types/src/run_projection.rs +++ b/lib/foundation/fabro-types/src/run_projection.rs @@ -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, pub pending_interviews: BTreeMap, + /// 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, stages: HashMap, } +/// 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(), } } diff --git a/lib/packages/fabro-api-client/src/models/server-slate-db-settings.ts b/lib/packages/fabro-api-client/src/models/server-slate-db-settings.ts deleted file mode 100644 index 69c684618..000000000 --- a/lib/packages/fabro-api-client/src/models/server-slate-db-settings.ts +++ /dev/null @@ -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; -}