refactor(run): simplify CAS-backed run definitions

Drop compatibility versioning from run-definition blobs, remove the
read-after-write polling added around CAS access, and tighten tests to
assert workflow_bundle.json is never written.
This commit is contained in:
Bryan Helmkamp 2026-04-07 23:10:00 -04:00
parent 6e3cd5fc12
commit 4a417b6013
No known key found for this signature in database
5 changed files with 25 additions and 63 deletions

View file

@ -6533,7 +6533,10 @@ mod tests {
.expect("accepted definition blob should exist");
let accepted_definition: serde_json::Value =
serde_json::from_slice(&accepted_definition_bytes).unwrap();
assert_eq!(accepted_definition["version"], 1);
assert!(
accepted_definition.get("version").is_none(),
"accepted run definition should not carry compatibility versioning"
);
assert_eq!(accepted_definition["workflow_path"], "workflow.fabro");
assert!(accepted_definition["workflows"]["workflow.fabro"].is_object());

View file

@ -1,14 +1,12 @@
use std::collections::VecDeque;
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use std::time::{Duration, Instant};
use bytes::Bytes;
use chrono::Utc;
use futures::Stream;
use slatedb::{Db, DbRead};
use tokio::sync::{Mutex, broadcast, mpsc};
use tokio::time::sleep;
use tokio_stream::wrappers::UnboundedReceiverStream;
use crate::keys;
@ -17,9 +15,6 @@ use crate::{EventEnvelope, EventPayload, Result, RunProjection, RunSummary, Stor
use fabro_types::{RunBlobId, RunId};
const DEFAULT_EVENT_TAIL_LIMIT: usize = 1024;
const BLOB_VISIBILITY_WAIT_TIMEOUT: Duration = Duration::from_secs(30);
const BLOB_VISIBILITY_POLL_INTERVAL: Duration = Duration::from_millis(50);
#[derive(Clone)]
pub struct RunDatabase {
inner: Arc<RunDatabaseInner>,
@ -302,18 +297,6 @@ impl RunDatabase {
}
let id = RunBlobId::new(data);
self.inner.db.put(keys::blob_key(&id), data).await?;
let deadline = Instant::now() + BLOB_VISIBILITY_WAIT_TIMEOUT;
loop {
if self.read_blob(&id).await?.is_some() {
break;
}
if Instant::now() >= deadline {
return Err(StoreError::Other(format!(
"blob {id} remained unreadable after write"
)));
}
sleep(BLOB_VISIBILITY_POLL_INTERVAL).await;
}
Ok(id)
}

View file

@ -16,7 +16,7 @@ use crate::pipeline::{self, Persisted, TransformOptions, Validated};
use crate::records::RunRecord;
use crate::run_lookup::default_scratch_base;
use crate::transforms::{Transform, expand_vars};
use crate::workflow_bundle::{AcceptedRunDefinition, WorkflowBundle};
use crate::workflow_bundle::{RunDefinition, WorkflowBundle};
use fabro_sandbox::daytona::detect_repo_info;
use fabro_util::json::normalize_json_value;
@ -114,7 +114,7 @@ pub async fn create(store: &Database, request: CreateRunInput) -> Result<Created
let current_dir = resolved.current_dir.clone();
let file_resolver = resolved.file_resolver.clone();
let accepted_definition = match (&workflow_path, &workflow_bundle) {
(Some(workflow_path), Some(workflow_bundle)) => Some(AcceptedRunDefinition::new(
(Some(workflow_path), Some(workflow_bundle)) => Some(RunDefinition::new(
workflow_path.clone(),
workflow_bundle.clone(),
)),
@ -169,7 +169,7 @@ async fn persist_created_run(
workflow_source: &str,
workflow_config: Option<String>,
submitted_manifest_bytes: Option<&[u8]>,
accepted_definition: Option<&AcceptedRunDefinition>,
accepted_definition: Option<&RunDefinition>,
) -> Result<(), FabroError> {
let record = persisted.run_record();
let run_store = match store.create_run(&record.run_id).await {

View file

@ -30,18 +30,12 @@ use crate::run_control::RunControlState;
use crate::run_options::{GitCheckpointOptions, LifecycleOptions, RunOptions};
use crate::run_status::{RunStatus, StatusReason};
use crate::runtime_store::RunStoreHandle;
use crate::workflow_bundle::{
ACCEPTED_RUN_DEFINITION_VERSION, AcceptedRunDefinition, WorkflowBundle,
};
use crate::workflow_bundle::{RunDefinition, WorkflowBundle};
use fabro_config::run::PullRequestSettings;
use fabro_retro::retro::Retro;
use fabro_sandbox::daytona::DaytonaConfig;
use fabro_sandbox::daytona::detect_repo_info;
use tokio::runtime::Handle;
use tokio::time::sleep;
const ACCEPTED_DEFINITION_BLOB_WAIT_TIMEOUT: Duration = Duration::from_secs(5);
const ACCEPTED_DEFINITION_BLOB_POLL_INTERVAL: Duration = Duration::from_millis(50);
struct RunSession {
cancel_token: Option<Arc<AtomicBool>>,
@ -428,32 +422,17 @@ impl RunSession {
async fn load_accepted_run_definition(
run_store: &RunStoreHandle,
blob_id: fabro_types::RunBlobId,
) -> Result<AcceptedRunDefinition, FabroError> {
let deadline = Instant::now() + ACCEPTED_DEFINITION_BLOB_WAIT_TIMEOUT;
let bytes = loop {
let blob = run_store
.read_blob(&blob_id)
.await
.map_err(|err| FabroError::engine(err.to_string()))?;
if let Some(blob) = blob {
break blob;
}
if Instant::now() >= deadline {
return Err(FabroError::engine(format!(
"accepted run definition blob is missing from the run store: {blob_id}"
)));
}
sleep(ACCEPTED_DEFINITION_BLOB_POLL_INTERVAL).await;
};
let definition: AcceptedRunDefinition =
serde_json::from_slice(&bytes).map_err(|err| FabroError::Parse(err.to_string()))?;
if definition.version != ACCEPTED_RUN_DEFINITION_VERSION {
return Err(FabroError::Parse(format!(
"unsupported accepted run definition version: {}",
definition.version
)));
}
Ok(definition)
) -> Result<RunDefinition, FabroError> {
let bytes = run_store
.read_blob(&blob_id)
.await
.map_err(|err| FabroError::engine(err.to_string()))?
.ok_or_else(|| {
FabroError::engine(format!(
"run definition blob is missing from the run store: {blob_id}"
))
})?;
serde_json::from_slice(&bytes).map_err(|err| FabroError::Parse(err.to_string()))
}
fn resolve_sandbox_provider(settings: &Settings) -> Result<SandboxProvider, FabroError> {
@ -1078,9 +1057,10 @@ mod tests {
.unwrap();
let bundle_file = created.run_dir.join("workflow_bundle.json");
if bundle_file.exists() {
std::fs::remove_file(&bundle_file).unwrap();
}
assert!(
!bundle_file.exists(),
"run scratch should not persist workflow_bundle.json"
);
let started = start(
&run_dir,

View file

@ -60,20 +60,16 @@ impl WorkflowBundle {
}
}
pub const ACCEPTED_RUN_DEFINITION_VERSION: u32 = 1;
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct AcceptedRunDefinition {
pub version: u32,
pub struct RunDefinition {
pub workflow_path: PathBuf,
pub workflows: HashMap<PathBuf, BundledWorkflow>,
}
impl AcceptedRunDefinition {
impl RunDefinition {
#[must_use]
pub fn new(workflow_path: PathBuf, bundle: WorkflowBundle) -> Self {
Self {
version: ACCEPTED_RUN_DEFINITION_VERSION,
workflow_path,
workflows: bundle.workflows,
}