From 6b65f2a6af9f963f021a4428fbe805ee39f0d662 Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Sat, 22 Aug 2026 10:29:00 -0400 Subject: [PATCH] Share the create-run pipeline tail between both admission lanes The intent and legacy-manifest create handlers each carried a full copy of the same post-admission sequence: LLM readiness resolution, graph compilation and model pinning, persistence, summary read, managed-run registration, title-generation spawn, and the 201 response. The copies had already drifted on when the run ID is resolved (before compilation in one lane, after in the other). Extract one finalize_created_run tail, with a small CreatedRunErrorStyle carrying each lane's pinned error mapping and log lines so the wire contracts are unchanged. Both lanes now resolve identity before compilation and share the parent-link validation, which lets the PinnedRun copy of PreparedRun's identity accessors be deleted. Co-Authored-By: Claude Fable 5 --- lib/apps/fabro-server/src/run_compiler.rs | 17 - .../fabro-server/src/server/handler/runs.rs | 294 +++++++++--------- 2 files changed, 142 insertions(+), 169 deletions(-) diff --git a/lib/apps/fabro-server/src/run_compiler.rs b/lib/apps/fabro-server/src/run_compiler.rs index 24ce7b36f..c82885bfc 100644 --- a/lib/apps/fabro-server/src/run_compiler.rs +++ b/lib/apps/fabro-server/src/run_compiler.rs @@ -212,23 +212,6 @@ pub(crate) struct PinnedRun { metadata: RunMetadata, } -impl PinnedRun { - pub(crate) fn parent_id(&self) -> Option { - self.metadata.parent_id - } - - pub(crate) fn resolve_run_id(mut self) -> (Self, RunId) { - let run_id = self.metadata.run_id.unwrap_or_default(); - self.metadata.run_id = Some(run_id); - (self, run_id) - } - - pub(crate) fn with_web_url(mut self, web_url: Option) -> Self { - self.metadata.web_url = web_url; - self - } -} - #[derive(Debug, thiserror::Error)] pub(crate) enum RunCompilerError { /// The acquired source bundle is invalid: missing entrypoint or broken diff --git a/lib/apps/fabro-server/src/server/handler/runs.rs b/lib/apps/fabro-server/src/server/handler/runs.rs index 4b2638bf4..11d32feaa 100644 --- a/lib/apps/fabro-server/src/server/handler/runs.rs +++ b/lib/apps/fabro-server/src/server/handler/runs.rs @@ -728,7 +728,130 @@ async fn create_run_from_intent( if let Err(error) = validate_intent_environment(&state, prepared.settings()).await { return run_intent_admission_error(error.into()); } + let (prepared, run_id) = prepared.resolve_run_id(); + if let Err(response) = validate_optional_parent(&state, run_id, prepared.parent_id()).await { + return response; + } + let prepared = prepared.with_web_url(state.run_web_url(&run_id)); + finalize_created_run( + state, + prepared, + intent.title.is_some(), + entrypoint, + CreatedRunErrorStyle::Intent, + ) + .await +} + +/// Shared parent-link validation for both create lanes: a run must not be +/// its own parent, and an explicit parent must pass [`validate_parent_link`]. +async fn validate_optional_parent( + state: &AppState, + run_id: RunId, + parent_id: Option, +) -> Result<(), Response> { + let Some(parent_id) = parent_id else { + return Ok(()); + }; + if parent_id == run_id { + return Err(ApiError::bad_request("A run cannot be its own parent.").into_response()); + } + validate_parent_link(state, run_id, parent_id) + .await + .map_err(IntoResponse::into_response) +} + +/// Which endpoint dialect's pinned error mapping the shared creation tail +/// speaks: the RunIntent admission contract or the legacy manifest wire +/// contract. +enum CreatedRunErrorStyle { + Intent, + LegacyManifest, +} + +impl CreatedRunErrorStyle { + fn compiler_error(&self, error: RunCompilerError) -> Response { + match self { + Self::Intent => run_intent_admission_error(error.into()), + Self::LegacyManifest => run_compiler_error_response(error), + } + } + + fn persist_error(&self, error: &WorkflowError) -> Response { + match self { + Self::Intent => { + tracing::error!( + error = %error, + error_chain = ?error_util::collect_chain(error), + "Failed to persist admitted run intent" + ); + intent_error( + StatusCode::INTERNAL_SERVER_ERROR, + "failed to persist run", + "run_persistence_failed", + ) + } + Self::LegacyManifest => ApiError::new( + StatusCode::INTERNAL_SERVER_ERROR, + format!("Failed to persist run state: {error}"), + ) + .into_response(), + } + } + + fn missing_summary_error(&self) -> Response { + match self { + Self::Intent => intent_error( + StatusCode::INTERNAL_SERVER_ERROR, + "created run summary is unavailable", + "run_persistence_failed", + ), + Self::LegacyManifest => ApiError::not_found("Run not found.").into_response(), + } + } + + fn summary_error(&self, error: &dyn std::fmt::Display) -> Response { + match self { + Self::Intent => { + tracing::error!(error = %error, "Failed to read admitted run summary"); + intent_error( + StatusCode::INTERNAL_SERVER_ERROR, + "failed to read created run", + "run_persistence_failed", + ) + } + Self::LegacyManifest => { + ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, error.to_string()).into_response() + } + } + } + + fn log_created(&self, run_id: RunId) { + // The legacy manifest lane logs its "Run created" line before + // compilation, so only the intent lane logs here. + if matches!(self, Self::Intent) { + info!(run_id = %run_id, "Run created from intent"); + } + } +} + +/// The shared tail of both run-creation lanes: resolve LLM readiness, compile +/// and pin, persist, register the managed run, spawn title generation, and +/// render the 201 response. Identity (run id, parent link, web URL) must +/// already be resolved on `prepared`; only error mapping differs per lane, +/// through [`CreatedRunErrorStyle`]. +async fn finalize_created_run( + state: Arc, + prepared: run_compiler::PreparedRun, + explicit_title_supplied: bool, + title_generation_target: ManifestPath, + style: CreatedRunErrorStyle, +) -> Response { let catalog = state.catalog(); + // Resolve once: we need both the provider IDs (for the run create input + // and ask-fabro-readiness) and the LLM client itself (for the spawned + // title-generation task). `ready_llm_provider_ids` would otherwise call + // `resolve_llm_client` a second time and discard the client. let (llm_result, ready_provider_ids) = state.resolve_llm_client_with_ready_ids().await; let llm_client_for_title = llm_result.ok(); let run_materialization_provider_ids = { @@ -749,18 +872,8 @@ async fn create_run_from_intent( .await { Ok(pinned) => pinned, - Err(error) => return run_intent_admission_error(error.into()), + Err(error) => return style.compiler_error(error), }; - let (pinned, run_id) = pinned.resolve_run_id(); - if let Some(parent_id) = pinned.parent_id() { - if parent_id == run_id { - return ApiError::bad_request("A run cannot be its own parent.").into_response(); - } - if let Err(err) = validate_parent_link(&state, run_id, parent_id).await { - return err.into_response(); - } - } - let pinned = pinned.with_web_url(state.run_web_url(&run_id)); let persistence_input = run_compiler::assemble_run(pinned); let created = match Box::pin(operations::persist_create_run( state.stores.runs.as_ref(), @@ -769,19 +882,9 @@ async fn create_run_from_intent( .await { Ok(created) => created, - Err(err) => { - tracing::error!( - error = %err, - error_chain = ?error_util::collect_chain(&err), - "Failed to persist admitted run intent" - ); - return intent_error( - StatusCode::INTERNAL_SERVER_ERROR, - "failed to persist run", - "run_persistence_failed", - ); - } + Err(error) => return style.persist_error(&error), }; + let created_at = created.run_id.created_at(); let summary = match state .stores .runs @@ -789,24 +892,10 @@ async fn create_run_from_intent( .await { Ok(Some(summary)) => summary, - Ok(None) => { - return intent_error( - StatusCode::INTERNAL_SERVER_ERROR, - "created run summary is unavailable", - "run_persistence_failed", - ); - } - Err(err) => { - tracing::error!(error = %err, "Failed to read admitted run summary"); - return intent_error( - StatusCode::INTERNAL_SERVER_ERROR, - "failed to read created run", - "run_persistence_failed", - ); - } + Ok(None) => return style.missing_summary_error(), + Err(error) => return style.summary_error(&error), }; let deterministic_title = summary.title.clone(); - let created_at = created.run_id.created_at(); { let mut runs = state.runs.lock().expect("runs lock poisoned"); runs.insert( @@ -820,8 +909,8 @@ async fn create_run_from_intent( ), ); } - if intent.title.is_none() && !ready_provider_ids.is_empty() { - if let Some(client) = llm_client_for_title { + if !explicit_title_supplied && !ready_provider_ids.is_empty() { + if let Some(llm_result) = llm_client_for_title { let run_spec = created.persisted.run_spec(); let workflow = run_title_generation::workflow_summary(&run_spec.graph); let run_inputs = run_spec.settings.run.inputs.clone(); @@ -831,16 +920,16 @@ async fn create_run_from_intent( state: Arc::clone(&state), run_id: created.run_id, deterministic_title, - workflow_target: entrypoint.to_string(), + workflow_target: title_generation_target.to_string(), workflow, run_inputs, - client: client.client, + client: llm_result.client, model_id: title_model.id.to_string(), provider_id: title_model.provider.clone(), }); } } - info!(run_id = %created.run_id, "Run created from intent"); + style.log_created(created.run_id); ( StatusCode::CREATED, Json(state.decorate_run_summary(summary).await), @@ -1214,118 +1303,19 @@ pub(crate) async fn create_run_from_manifest( { return ApiError::bad_request(error).into_response(); } - if let Some(parent_id) = prepared.parent_id() { - if parent_id == run_id { - return ApiError::bad_request("A run cannot be its own parent.").into_response(); - } - if let Err(err) = validate_parent_link(&state, run_id, parent_id).await { - return err.into_response(); - } + if let Err(response) = validate_optional_parent(&state, run_id, prepared.parent_id()).await { + return response; } info!(run_id = %run_id, "Run created"); - let catalog = state.catalog(); - // Resolve once: we need both the provider IDs (for the run create input - // and ask-fabro-readiness) and the LLM client itself (for the spawned - // title-generation task). `ready_llm_provider_ids` would otherwise call - // `resolve_llm_client` a second time and discard the client. - let (llm_result, ready_provider_ids) = state.resolve_llm_client_with_ready_ids().await; - let llm_client_for_title = llm_result.ok(); - let run_materialization_provider_ids = { - #[cfg(any(test, feature = "test-support"))] - { - server_test_support::test_run_materialization_provider_ids( - catalog.as_ref(), - &ready_provider_ids, - ) - } - #[cfg(not(any(test, feature = "test-support")))] - { - ready_provider_ids.clone() - } - }; - let pinned = - match run_compiler::compile_and_pin(prepared, run_materialization_provider_ids, catalog) - .await - { - Ok(pinned) => pinned, - Err(err) => return run_compiler_error_response(err), - }; - let persistence_input = run_compiler::assemble_run(pinned); - let created = match Box::pin(operations::persist_create_run( - state.stores.runs.as_ref(), - persistence_input, - )) - .await - { - Ok(created) => created, - Err(err) => { - return ApiError::new( - StatusCode::INTERNAL_SERVER_ERROR, - format!("Failed to persist run state: {err}"), - ) - .into_response(); - } - }; - let created_at = created.run_id.created_at(); - let summary = match state - .stores - .runs - .get_cached_summary(&created.run_id, Utc::now()) - .await - { - Ok(Some(summary)) => summary, - Ok(None) => return ApiError::not_found("Run not found.").into_response(), - Err(err) => { - return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()) - .into_response(); - } - }; - let deterministic_title = summary.title.clone(); - - { - let mut runs = state.runs.lock().expect("runs lock poisoned"); - runs.insert( - created.run_id, - managed_run( - created.persisted.source().to_string(), - RunStatus::Submitted, - created_at, - created.run_dir, - RunExecutionMode::Start, - ), - ); - } - - if !explicit_title_supplied && !ready_provider_ids.is_empty() { - if let Some(llm_result) = llm_client_for_title { - let run_spec = created.persisted.run_spec(); - let workflow = run_title_generation::workflow_summary(&run_spec.graph); - let run_inputs = run_spec.settings.run.inputs.clone(); - let workflow_target = title_generation_target.to_string(); - let title_catalog = state.catalog(); - let title_model = title_catalog.small_default_for_configured_ids(&ready_provider_ids); - let title_model_id = title_model.id.clone(); - let title_provider_id = title_model.provider.clone(); - spawn_generated_title_task(GeneratedTitleTask { - state: Arc::clone(&state), - run_id: created.run_id, - deterministic_title, - workflow_target, - workflow, - run_inputs, - client: llm_result.client, - model_id: title_model_id.to_string(), - provider_id: title_provider_id, - }); - } - } - - ( - StatusCode::CREATED, - Json(state.decorate_run_summary(summary).await), + finalize_created_run( + state, + prepared, + explicit_title_supplied, + title_generation_target, + CreatedRunErrorStyle::LegacyManifest, ) - .into_response() + .await } struct GeneratedTitleTask {