diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index f7f78a66f..2b9ef819c 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -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()); diff --git a/lib/crates/fabro-store/src/slate/run_store.rs b/lib/crates/fabro-store/src/slate/run_store.rs index 173109e5c..f36efb5bf 100644 --- a/lib/crates/fabro-store/src/slate/run_store.rs +++ b/lib/crates/fabro-store/src/slate/run_store.rs @@ -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, @@ -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) } diff --git a/lib/crates/fabro-workflow/src/operations/create.rs b/lib/crates/fabro-workflow/src/operations/create.rs index bd42b9243..987440f9a 100644 --- a/lib/crates/fabro-workflow/src/operations/create.rs +++ b/lib/crates/fabro-workflow/src/operations/create.rs @@ -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 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, 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 { diff --git a/lib/crates/fabro-workflow/src/operations/start.rs b/lib/crates/fabro-workflow/src/operations/start.rs index 780363f3c..4ae8122ae 100644 --- a/lib/crates/fabro-workflow/src/operations/start.rs +++ b/lib/crates/fabro-workflow/src/operations/start.rs @@ -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>, @@ -428,32 +422,17 @@ impl RunSession { async fn load_accepted_run_definition( run_store: &RunStoreHandle, blob_id: fabro_types::RunBlobId, -) -> Result { - 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 { + 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 { @@ -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, diff --git a/lib/crates/fabro-workflow/src/workflow_bundle.rs b/lib/crates/fabro-workflow/src/workflow_bundle.rs index 0ab2d4c54..fd7afe078 100644 --- a/lib/crates/fabro-workflow/src/workflow_bundle.rs +++ b/lib/crates/fabro-workflow/src/workflow_bundle.rs @@ -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, } -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, }