diff --git a/Cargo.lock b/Cargo.lock index e71612fb2..2b5fa3d2a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3393,6 +3393,7 @@ dependencies = [ "serde_json", "sha2 0.10.9", "shlex", + "strum 0.28.0", "tempfile", "thiserror 2.0.18", "tokio", diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs index d53a591ff..af050d073 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/event.rs @@ -614,7 +614,7 @@ mod tests { repo: "widgets".into(), base_branch: "main".into(), head_branch: "fabro/run/42".into(), - head_sha: "final-sha".into(), + head_sha: Some("final-sha".to_string()), title: "Ship the server-side PR".into(), draft: true, }; diff --git a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs index 4b7a04cbf..38ba83bab 100644 --- a/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/apps/fabro-cli/src/commands/run/run_progress/mod.rs @@ -1383,7 +1383,7 @@ mod tests { repo: "fabro".into(), base_branch: "main".into(), head_branch: "fabro/run/42".into(), - head_sha: "final-sha".into(), + head_sha: Some("final-sha".to_string()), title: "Ship the change".into(), draft: true, }); diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index c3d3113af..0fc2cd014 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -4169,9 +4169,15 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { let mut runs = state.runs.lock().expect("runs lock poisoned"); if let Some(managed_run) = runs.get_mut(&run_id) { match &result { - ExecutionResult::Completed(result) => match result.as_ref() { - Ok(started) => match &started.finalized.outcome { - Ok(_) => { + ExecutionResult::Completed(result) => { + // A run can fail either before it produces a `Started` or in + // its own outcome; both carry the same `WorkflowError`. + let outcome = match result.as_ref() { + Ok(started) => started.finalized.outcome.as_ref().map(|_| ()), + Err(e) => Err(e), + }; + match outcome { + Ok(()) => { info!(run_id = %run_id, "Run completed"); managed_run.status = RunStatus::Succeeded { reason: SuccessReason::Completed, @@ -4183,44 +4189,16 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { reason: FailureReason::Cancelled, }; } - Err(e @ WorkflowError::Publish { .. }) => { + Err(e) => { let detail = e.display_with_causes(); - error!(run_id = %run_id, error = %detail, "Run publish failed"); + error!(run_id = %run_id, error = %detail, "Run failed"); managed_run.status = RunStatus::Failed { - reason: FailureReason::PublishFailed, + reason: e.failure_reason(), }; managed_run.error = Some(detail); } - Err(e) => { - error!(run_id = %run_id, error = %e, "Run failed"); - managed_run.status = RunStatus::Failed { - reason: FailureReason::WorkflowError, - }; - managed_run.error = Some(e.to_string()); - } - }, - Err(WorkflowError::Cancelled) => { - info!(run_id = %run_id, "Run cancelled"); - managed_run.status = RunStatus::Failed { - reason: FailureReason::Cancelled, - }; } - Err(e @ WorkflowError::Publish { .. }) => { - let detail = e.display_with_causes(); - error!(run_id = %run_id, error = %detail, "Run publish failed"); - managed_run.status = RunStatus::Failed { - reason: FailureReason::PublishFailed, - }; - managed_run.error = Some(detail); - } - Err(e) => { - error!(run_id = %run_id, error = %e, "Run failed"); - managed_run.status = RunStatus::Failed { - reason: FailureReason::WorkflowError, - }; - managed_run.error = Some(e.to_string()); - } - }, + } ExecutionResult::CancelledBySignal => { info!(run_id = %run_id, "Run cancelled"); managed_run.status = RunStatus::Failed { diff --git a/lib/apps/fabro-server/src/server/handler/pull_requests.rs b/lib/apps/fabro-server/src/server/handler/pull_requests.rs index a30a55c03..8bf5fa898 100644 --- a/lib/apps/fabro-server/src/server/handler/pull_requests.rs +++ b/lib/apps/fabro-server/src/server/handler/pull_requests.rs @@ -348,15 +348,8 @@ async fn create_run_pull_request( conclusion: Some(inputs.conclusion), run_state: Some(run_state), }; - let created_pull_request = match pull_request::maybe_open_pull_request(request).await { - Ok(Some(created)) => created, - Ok(None) => { - return ApiError::new( - StatusCode::INTERNAL_SERVER_ERROR, - "Pull request creation returned no record unexpectedly.", - ) - .into_response(); - } + let created_pull_request = match pull_request::open_pull_request(request).await { + Ok(created) => created, Err(err) => return ApiError::new(StatusCode::BAD_GATEWAY, err).into_response(), }; @@ -364,7 +357,7 @@ async fn create_run_pull_request( &created_pull_request.link, &created_pull_request.base_branch, &created_pull_request.head_branch, - &created_pull_request.head_sha, + inputs.final_git_sha, &created_pull_request.title, true, ); diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 73c9fcd34..40e7461e4 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -4685,7 +4685,7 @@ channel = "#deploys" repo: "fabro".to_string(), base_branch: "main".to_string(), head_branch: "fabro/run/test".to_string(), - head_sha: "final-sha".to_string(), + head_sha: Some("final-sha".to_string()), title: "Ship & notify".to_string(), draft: false, }, @@ -6426,7 +6426,7 @@ async fn create_run_with_pull_request_record( repo: "widgets".to_string(), base_branch: "main".to_string(), head_branch: "feature".to_string(), - head_sha: "final-sha".to_string(), + head_sha: Some("final-sha".to_string()), title: title.to_string(), draft: false, }, @@ -16488,7 +16488,7 @@ async fn list_runs_includes_live_metadata_from_run_state() { repo: "repo".to_string(), base_branch: "main".to_string(), head_branch: "fabro/run".to_string(), - head_sha: "final-sha".to_string(), + head_sha: Some("final-sha".to_string()), title: "Fix board metadata".to_string(), draft: false, }, diff --git a/lib/components/fabro-github/src/lib.rs b/lib/components/fabro-github/src/lib.rs index d150c4759..df45eca9d 100644 --- a/lib/components/fabro-github/src/lib.rs +++ b/lib/components/fabro-github/src/lib.rs @@ -888,20 +888,6 @@ fn normalize_https_host_path(url: &str) -> String { } } -/// Check whether a branch exists in a GitHub repository. -/// -/// Uses a GitHub App installation token to query the branches API. -/// Returns `true` if the branch exists, `false` if it doesn't (404). -pub async fn branch_exists( - ctx: &GitHubContext<'_>, - owner: &str, - repo: &str, - branch: &str, -) -> anyhow::Result { - let client = ctx.http_client()?; - branch_exists_with_client(&client, ctx, owner, repo, branch).await -} - /// Return the commit SHA at the head of a GitHub branch. /// /// Returns `None` when the branch does not exist. @@ -910,6 +896,17 @@ pub async fn branch_head_sha( owner: &str, repo: &str, branch: &str, +) -> anyhow::Result> { + let client = ctx.http_client()?; + branch_head_sha_with_client(&client, ctx, owner, repo, branch).await +} + +async fn branch_head_sha_with_client( + client: &impl HttpClient, + ctx: &GitHubContext<'_>, + owner: &str, + repo: &str, + branch: &str, ) -> anyhow::Result> { #[derive(Deserialize)] struct BranchResponse { @@ -921,11 +918,10 @@ pub async fn branch_head_sha( sha: String, } - let client = ctx.http_client()?; let token = ctx .creds .resolve_bearer_token( - &client, + client, owner, repo, ctx.base_url, @@ -935,7 +931,8 @@ pub async fn branch_head_sha( let url = format!("{}/repos/{owner}/{repo}/branches/{branch}", ctx.base_url); let auth = format!("Bearer {token}"); - let resp = HttpClient::request(&client, HttpMethod::Get, &url, &github_headers(&auth), None) + let resp = client + .request(HttpMethod::Get, &url, &github_headers(&auth), None) .await .context("Failed to read remote branch head")?; @@ -951,38 +948,6 @@ pub async fn branch_head_sha( } } -async fn branch_exists_with_client( - client: &impl HttpClient, - ctx: &GitHubContext<'_>, - owner: &str, - repo: &str, - branch: &str, -) -> anyhow::Result { - let token = ctx - .creds - .resolve_bearer_token( - client, - owner, - repo, - ctx.base_url, - serde_json::json!({ "contents": "write" }), - ) - .await?; - - let url = format!("{}/repos/{owner}/{repo}/branches/{branch}", ctx.base_url); - let auth = format!("Bearer {token}"); - let resp = client - .request(HttpMethod::Get, &url, &github_headers(&auth), None) - .await - .context("Failed to check branch existence")?; - - match resp.status { - 200 => Ok(true), - 404 => Ok(false), - status => bail!("Unexpected status {status} checking branch '{branch}'"), - } -} - /// Check whether a GitHub App is installed for a specific repository. /// /// Uses the App JWT to query `GET /repos/{owner}/{repo}/installation`. @@ -1092,46 +1057,36 @@ pub async fn resolve_clone_credentials( ctx: &GitHubContext<'_>, owner: &str, repo: &str, -) -> anyhow::Result<(Option, Option)> { - match ctx.creds { - GitHubCredentials::Pat(token) => { - Ok((Some("x-access-token".to_string()), Some(token.clone()))) - } - GitHubCredentials::Installation(token) => Ok(( - Some("x-access-token".to_string()), - Some(token.valid_token()?.to_string()), - )), - GitHubCredentials::App(_) => { - let client = ctx.http_client()?; - resolve_clone_credentials_with_client(&client, ctx, owner, repo).await - } - } -} - -async fn resolve_clone_credentials_with_client( - client: &impl HttpClient, - ctx: &GitHubContext<'_>, - owner: &str, - repo: &str, ) -> anyhow::Result<(Option, Option)> { let token = match ctx.creds { GitHubCredentials::Pat(token) => token.clone(), GitHubCredentials::Installation(token) => token.valid_token()?.to_string(), GitHubCredentials::App(_) => { - ctx.creds - .resolve_bearer_token( - client, - owner, - repo, - ctx.base_url, - serde_json::json!({ "contents": "write", "workflows": "write" }), - ) - .await? + let client = ctx.http_client()?; + mint_git_write_token(&client, ctx, owner, repo).await? } }; Ok((Some("x-access-token".to_string()), Some(token))) } +/// Mint an installation token scoped for git writes, including workflow files. +async fn mint_git_write_token( + client: &impl HttpClient, + ctx: &GitHubContext<'_>, + owner: &str, + repo: &str, +) -> anyhow::Result { + ctx.creds + .resolve_bearer_token( + client, + owner, + repo, + ctx.base_url, + serde_json::json!({ "contents": "write", "workflows": "write" }), + ) + .await +} + /// Embed a token into an HTTPS URL for authenticated git operations. /// /// Converts `https://github.com/owner/repo` to @@ -1983,12 +1938,21 @@ mod tests { } // ----------------------------------------------------------------------- - // branch_exists + // branch_head_sha // ----------------------------------------------------------------------- - #[tokio::test] - async fn branch_exists_returns_true_on_200() { - let mock = MockHttpClient::new() + fn app_creds() -> GitHubCredentials { + GitHubCredentials::App(GitHubAppCredentials { + app_id: "test".to_string(), + private_key_pem: test_rsa_key().to_string(), + slug: None, + }) + } + + /// Mock the installation-token exchange every App-credentialed call makes + /// before it reaches the endpoint under test. + fn mock_with_installation_token() -> MockHttpClient { + MockHttpClient::new() .on( HttpMethod::Get, "/repos/owner/repo/installation", @@ -2001,20 +1965,19 @@ mod tests { 201, r#"{"token": "ghs_test", "expires_at": "2099-01-01T00:00:00Z"}"#, ) - .on( - HttpMethod::Get, - "/repos/owner/repo/branches/my-branch", - 200, - r#"{"name": "my-branch"}"#, - ); + } - let pem = test_rsa_key(); - let creds = GitHubCredentials::App(GitHubAppCredentials { - app_id: "test".to_string(), - private_key_pem: pem.to_string(), - slug: None, - }); - let result = branch_exists_with_client( + #[tokio::test] + async fn branch_head_sha_returns_commit_on_200() { + let mock = mock_with_installation_token().on( + HttpMethod::Get, + "/repos/owner/repo/branches/my-branch", + 200, + r#"{"name": "my-branch", "commit": {"sha": "abc123"}}"#, + ); + + let creds = app_creds(); + let result = branch_head_sha_with_client( &mock, &GitHubContext::new(&creds, ""), "owner", @@ -2022,38 +1985,21 @@ mod tests { "my-branch", ) .await; - assert!(result.unwrap()); + + assert_eq!(result.unwrap(), Some("abc123".to_string())); } #[tokio::test] - async fn branch_exists_returns_false_on_404() { - let mock = MockHttpClient::new() - .on( - HttpMethod::Get, - "/repos/owner/repo/installation", - 200, - r#"{"id": 1}"#, - ) - .on( - HttpMethod::Post, - "/app/installations/1/access_tokens", - 201, - r#"{"token": "ghs_test", "expires_at": "2099-01-01T00:00:00Z"}"#, - ) - .on( - HttpMethod::Get, - "/repos/owner/repo/branches/no-such-branch", - 404, - "", - ); + async fn branch_head_sha_returns_none_on_404() { + let mock = mock_with_installation_token().on( + HttpMethod::Get, + "/repos/owner/repo/branches/no-such-branch", + 404, + "", + ); - let pem = test_rsa_key(); - let creds = GitHubCredentials::App(GitHubAppCredentials { - app_id: "test".to_string(), - private_key_pem: pem.to_string(), - slug: None, - }); - let result = branch_exists_with_client( + let creds = app_creds(); + let result = branch_head_sha_with_client( &mock, &GitHubContext::new(&creds, ""), "owner", @@ -2061,38 +2007,21 @@ mod tests { "no-such-branch", ) .await; - assert!(!result.unwrap()); + + assert_eq!(result.unwrap(), None); } #[tokio::test] - async fn branch_exists_returns_error_on_500() { - let mock = MockHttpClient::new() - .on( - HttpMethod::Get, - "/repos/owner/repo/installation", - 200, - r#"{"id": 1}"#, - ) - .on( - HttpMethod::Post, - "/app/installations/1/access_tokens", - 201, - r#"{"token": "ghs_test", "expires_at": "2099-01-01T00:00:00Z"}"#, - ) - .on( - HttpMethod::Get, - "/repos/owner/repo/branches/broken", - 500, - "", - ); + async fn branch_head_sha_returns_error_on_500() { + let mock = mock_with_installation_token().on( + HttpMethod::Get, + "/repos/owner/repo/branches/broken", + 500, + "", + ); - let pem = test_rsa_key(); - let creds = GitHubCredentials::App(GitHubAppCredentials { - app_id: "test".to_string(), - private_key_pem: pem.to_string(), - slug: None, - }); - let result = branch_exists_with_client( + let creds = app_creds(); + let result = branch_head_sha_with_client( &mock, &GitHubContext::new(&creds, ""), "owner", @@ -2100,22 +2029,23 @@ mod tests { "broken", ) .await; + assert!(result.is_err()); } #[tokio::test] - async fn branch_exists_with_token_uses_direct_bearer_token() { + async fn branch_head_sha_with_token_uses_direct_bearer_token() { let mock = MockHttpClient::new() .on( HttpMethod::Get, "/repos/owner/repo/branches/my-branch", 200, - r#"{"name": "my-branch"}"#, + r#"{"name": "my-branch", "commit": {"sha": "abc123"}}"#, ) .with_req_header("Authorization", "Bearer ghu_test"); let creds = GitHubCredentials::Pat("ghu_test".to_string()); - let result = branch_exists_with_client( + let result = branch_head_sha_with_client( &mock, &GitHubContext::new(&creds, ""), "owner", @@ -2124,7 +2054,7 @@ mod tests { ) .await; - assert!(result.unwrap()); + assert_eq!(result.unwrap(), Some("abc123".to_string())); } // ----------------------------------------------------------------------- @@ -2395,17 +2325,11 @@ mod tests { slug: None, }); let context = GitHubContext::new(&credentials, ""); - let resolved = resolve_clone_credentials_with_client(&mock, &context, "owner", "repo") + let token = mint_git_write_token(&mock, &context, "owner", "repo") .await .unwrap(); - assert_eq!( - resolved, - ( - Some("x-access-token".to_string()), - Some("ghs_xxx".to_string()) - ) - ); + assert_eq!(token, "ghs_xxx"); } #[test] diff --git a/lib/components/fabro-workflow/Cargo.toml b/lib/components/fabro-workflow/Cargo.toml index 932bb0071..535a9f6b5 100644 --- a/lib/components/fabro-workflow/Cargo.toml +++ b/lib/components/fabro-workflow/Cargo.toml @@ -46,6 +46,7 @@ fabro-static.workspace = true fabro-types = { path = "../../foundation/fabro-types" } fabro-http.workspace = true thiserror.workspace = true +strum.workspace = true serde.workspace = true serde_json.workspace = true jsonschema.workspace = true diff --git a/lib/components/fabro-workflow/src/error.rs b/lib/components/fabro-workflow/src/error.rs index c134a5e70..ac7eed9c5 100644 --- a/lib/components/fabro-workflow/src/error.rs +++ b/lib/components/fabro-workflow/src/error.rs @@ -252,6 +252,21 @@ impl FailureSignatureExt for FailureSignature { } } +/// Pipeline stage that produced an [`Error::Stage`]. +/// +/// The three stages share a failure shape — a message, an eagerly classified +/// [`FailureCategory`], an optional command output tail, and an optional +/// source — and differ only in where they run and whether a retry is possible. +#[derive(Debug, Clone, Copy, PartialEq, Eq, strum::Display)] +pub enum ErrorStage { + /// A node handler failed. Retryable: the engine can re-run the node. + Handler, + /// The engine itself failed while driving the graph. Retryable. + Engine, + /// The publish stage failed. Terminal: publish runs once, after execution. + Publish, +} + #[derive(ThisError, Debug, Clone)] pub enum Error { #[error("Parse error: {0}")] @@ -276,26 +291,9 @@ pub enum Error { source: SharedTemplateError, }, - #[error("Engine error: {message}")] - Engine { - message: String, - failure_class: FailureCategory, - exec_output_tail: Option, - #[source] - source: Option, - }, - - #[error("Publish error: {message}")] - Publish { - message: String, - failure_class: FailureCategory, - exec_output_tail: Option, - #[source] - source: Option, - }, - - #[error("Handler error: {message}")] - Handler { + #[error("{stage} error: {message}")] + Stage { + stage: ErrorStage, message: String, failure_class: FailureCategory, exec_output_tail: Option, @@ -334,17 +332,48 @@ pub enum Error { impl Error { /// Smart constructor for Handler errors. Classifies the failure reason /// eagerly. - pub fn handler(message: impl Into) -> Self { + /// Build a stage error, classifying the message eagerly. + fn stage( + stage: ErrorStage, + message: impl Into, + exec_output_tail: Option, + ) -> Self { let message = message.into(); let failure_class = classify_failure_reason(&message); - Self::Handler { + Self::Stage { + stage, message, failure_class, - exec_output_tail: None, + exec_output_tail, source: None, } } + /// Build a stage error from a source, classifying the rendered chain so + /// hints buried in the causes still reach [`Self::failure_category`]. + fn stage_with_source( + stage: ErrorStage, + message: impl Into, + source: impl Into, + exec_output_tail: Option, + ) -> Self { + let message = message.into(); + let source = SharedError::new(source.into()); + let failure_class = + classify_failure_reason(&render_with_causes(&message, &collect_chain(&source))); + Self::Stage { + stage, + message, + failure_class, + exec_output_tail, + source: Some(source), + } + } + + pub fn handler(message: impl Into) -> Self { + Self::stage(ErrorStage::Handler, message, None) + } + pub fn template(message: impl Into, source: TemplateError) -> Self { Self::Template { message: message.into(), @@ -356,14 +385,7 @@ impl Error { message: impl Into, exec_output_tail: Option, ) -> Self { - let message = message.into(); - let failure_class = classify_failure_reason(&message); - Self::Handler { - message, - failure_class, - exec_output_tail, - source: None, - } + Self::stage(ErrorStage::Handler, message, exec_output_tail) } pub fn handler_with_source( @@ -378,51 +400,22 @@ impl Error { source: impl Into, exec_output_tail: Option, ) -> Self { - let message = message.into(); - let source = SharedError::new(source.into()); - let causes = collect_chain(&source); - let rendered = render_with_causes(&message, &causes); - let failure_class = classify_failure_reason(&rendered); - Self::Handler { - message, - failure_class, - exec_output_tail, - source: Some(source), - } + Self::stage_with_source(ErrorStage::Handler, message, source, exec_output_tail) } pub fn handler_with_anyhow(message: impl Into, source: anyhow::Error) -> Self { Self::handler_with_source(message, source) } - /// Smart constructor for Engine errors. Classifies the failure reason - /// eagerly. pub fn engine(message: impl Into) -> Self { - let message = message.into(); - let failure_class = classify_failure_reason(&message); - Self::Engine { - message, - failure_class, - exec_output_tail: None, - source: None, - } + Self::stage(ErrorStage::Engine, message, None) } pub fn engine_with_source( message: impl Into, source: impl Into, ) -> Self { - let message = message.into(); - let source = SharedError::new(source.into()); - let causes = collect_chain(&source); - let rendered = render_with_causes(&message, &causes); - let failure_class = classify_failure_reason(&rendered); - Self::Engine { - message, - failure_class, - exec_output_tail: None, - source: Some(source), - } + Self::stage_with_source(ErrorStage::Engine, message, source, None) } pub fn engine_with_anyhow(message: impl Into, source: anyhow::Error) -> Self { @@ -431,14 +424,7 @@ impl Error { /// Build an error for the required publish stage. pub fn publish(message: impl Into) -> Self { - let message = message.into(); - let failure_class = classify_failure_reason(&message); - Self::Publish { - message, - failure_class, - exec_output_tail: None, - source: None, - } + Self::stage(ErrorStage::Publish, message, None) } pub fn publish_with_source( @@ -453,25 +439,13 @@ impl Error { source: impl Into, exec_output_tail: Option, ) -> Self { - let message = message.into(); - let source = SharedError::new(source.into()); - let causes = collect_chain(&source); - let rendered = render_with_causes(&message, &causes); - let failure_class = classify_failure_reason(&rendered); - Self::Publish { - message, - failure_class, - exec_output_tail, - source: Some(source), - } + Self::stage_with_source(ErrorStage::Publish, message, source, exec_output_tail) } #[must_use] pub fn causes(&self) -> Vec { match self { - Self::Engine { source, .. } - | Self::Publish { source, .. } - | Self::Handler { source, .. } => source + Self::Stage { source, .. } => source .as_ref() .map_or_else(Vec::new, |source| collect_chain(source)), Self::Template { source, .. } => collect_chain(source), @@ -487,16 +461,16 @@ impl Error { /// Whether this error category is retryable (transient) or terminal. /// - /// Retryable: Handler and Engine, I/O, LLM errors when the SDK marks them - /// retryable, and Publish errors classified as transient infrastructure - /// failures. Terminal: Parse, Validation, OutputSchemaValidation, - /// Stylesheet, Checkpoint, and Cancelled. + /// Retryable: Handler and Engine stages (the engine can re-run the node), + /// I/O, and LLM errors the SDK marks retryable. Terminal: the Publish + /// stage (it runs once, after execution), Parse, Validation, + /// OutputSchemaValidation, Stylesheet, Checkpoint, and Cancelled. #[must_use] pub fn is_retryable(&self) -> bool { match self { - Self::Handler { .. } | Self::Engine { .. } | Self::Io(_) => true, - Self::Publish { failure_class, .. } => { - matches!(failure_class, FailureCategory::TransientInfra) + Self::Io(_) => true, + Self::Stage { stage, .. } => { + matches!(stage, ErrorStage::Handler | ErrorStage::Engine) } Self::Llm(sdk_err) => sdk_err.retryable(), Self::Parse(_) @@ -533,9 +507,20 @@ impl Error { | Self::Unsupported(_) | Self::OutputSchemaValidation(_) => FailureCategory::Deterministic, Self::Precondition(_) | Self::RunNotFound(_) => FailureCategory::Structural, - Self::Handler { failure_class, .. } - | Self::Engine { failure_class, .. } - | Self::Publish { failure_class, .. } => *failure_class, + Self::Stage { failure_class, .. } => *failure_class, + } + } + + /// The terminal [`FailureReason`] this error maps to on a run. + #[must_use] + pub fn failure_reason(&self) -> FailureReason { + match self { + Self::Cancelled => FailureReason::Cancelled, + Self::Stage { + stage: ErrorStage::Publish, + .. + } => FailureReason::PublishFailed, + _ => FailureReason::WorkflowError, } } @@ -551,23 +536,13 @@ impl Error { #[must_use] pub fn to_failure_detail(&self) -> FailureDetail { - let message = match self { - Self::Engine { message, .. } - | Self::Publish { message, .. } - | Self::Handler { message, .. } => message.clone(), - _ => self.to_string(), - }; - let explicit_exec_output_tail = match self { - Self::Engine { - exec_output_tail, .. - } - | Self::Publish { - exec_output_tail, .. - } - | Self::Handler { - exec_output_tail, .. - } => exec_output_tail.clone(), - _ => None, + let (message, explicit_exec_output_tail) = match self { + Self::Stage { + message, + exec_output_tail, + .. + } => (message.clone(), exec_output_tail.clone()), + _ => (self.to_string(), None), }; FailureDetail { message, @@ -880,7 +855,10 @@ mod tests { source, }); - assert!(matches!(fabro_error, Error::Engine { .. })); + assert!(matches!(fabro_error, Error::Stage { + stage: ErrorStage::Engine, + .. + })); let message = fabro_error.to_string(); assert!(message.contains("deserialize run spec on branch fabro/meta/run-1")); assert!(message.contains(&source_message)); @@ -2061,10 +2039,30 @@ mod tests { ); } + /// Publish runs once, after execution, so no caller can retry it — even + /// when the message looks transient. The failure category is still + /// classified for reporting. #[test] - fn publish_error_is_only_retryable_for_transient_failures() { - assert!(Error::publish("connection timed out").is_retryable()); + fn publish_errors_are_terminal() { + assert!(!Error::publish("connection timed out").is_retryable()); assert!(!Error::publish("permission denied").is_retryable()); + assert_eq!( + Error::publish("connection timed out").failure_category(), + FailureCategory::TransientInfra + ); + } + + #[test] + fn failure_reason_distinguishes_publish_and_cancelled() { + assert_eq!( + Error::publish("nope").failure_reason(), + FailureReason::PublishFailed + ); + assert_eq!(Error::Cancelled.failure_reason(), FailureReason::Cancelled); + assert_eq!( + Error::engine("boom").failure_reason(), + FailureReason::WorkflowError + ); } #[test] diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index 42a4ffbc4..094f36dd4 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -1321,7 +1321,7 @@ fn event_body_from_event(event: &Event) -> EventBody { repo: repo.clone(), base_branch: base_branch.clone(), head_branch: head_branch.clone(), - head_sha: (!head_sha.is_empty()).then(|| head_sha.clone()), + head_sha: head_sha.clone(), title: title.clone(), draft: *draft, }), diff --git a/lib/components/fabro-workflow/src/event/events.rs b/lib/components/fabro-workflow/src/event/events.rs index f512e80c6..4bc7bca4f 100644 --- a/lib/components/fabro-workflow/src/event/events.rs +++ b/lib/components/fabro-workflow/src/event/events.rs @@ -730,8 +730,9 @@ pub enum Event { repo: String, base_branch: String, head_branch: String, - #[serde(default)] - head_sha: String, + /// Absent on events written before the head SHA was recorded. + #[serde(default, skip_serializing_if = "Option::is_none")] + head_sha: Option, title: String, draft: bool, }, @@ -782,7 +783,7 @@ impl Event { repo: record.repo.clone(), base_branch: base_branch.to_string(), head_branch: head_branch.to_string(), - head_sha: head_sha.to_string(), + head_sha: Some(head_sha.to_string()), title: title.to_string(), draft, } diff --git a/lib/components/fabro-workflow/src/lifecycle/event.rs b/lib/components/fabro-workflow/src/lifecycle/event.rs index c4af83250..9a6ca4b0e 100644 --- a/lib/components/fabro-workflow/src/lifecycle/event.rs +++ b/lib/components/fabro-workflow/src/lifecycle/event.rs @@ -496,7 +496,7 @@ impl RunLifecycle for EventLifecycle { } for push in &result.push_results { self.emitter.emit(&Event::GitPush { - branch: push.refspec.clone(), + branch: push.branch.clone(), success: push.success, exec_output_tail: push.exec_output_tail.clone(), }); diff --git a/lib/components/fabro-workflow/src/lifecycle/git.rs b/lib/components/fabro-workflow/src/lifecycle/git.rs index 4db74c904..656a298a7 100644 --- a/lib/components/fabro-workflow/src/lifecycle/git.rs +++ b/lib/components/fabro-workflow/src/lifecycle/git.rs @@ -69,11 +69,24 @@ pub(crate) struct GitCheckpointResult { #[derive(Debug, Clone)] pub(crate) struct PushResult { - pub refspec: String, + pub branch: String, pub success: bool, pub exec_output_tail: Option, } +/// Push a run branch to its remote counterpart. +/// +/// Owns the refspec convention so the checkpoint push and the terminal publish +/// push cannot drift apart. +pub(crate) async fn push_run_branch( + sandbox: &dyn fabro_sandbox::Sandbox, + branch: &str, +) -> fabro_sandbox::Result<()> { + sandbox + .git_push_ref(&format!("refs/heads/{branch}:refs/heads/{branch}")) + .await +} + /// Sub-lifecycle responsible for git operations (checkpoint commits, pushes, /// diffs). pub(crate) struct GitLifecycle { @@ -307,15 +320,14 @@ impl RunLifecycle for GitLifecycle { .as_ref() .and_then(|g| g.run_branch.as_ref()) { - let refspec = format!("refs/heads/{branch}:refs/heads/{branch}"); let (push_ok, exec_output_tail) = - match self.sandbox.git_push_ref(&refspec).await { + match push_run_branch(self.sandbox.as_ref(), branch).await { Ok(()) => (true, None), Err(err) => { let exec_output_tail = fabro_sandbox::default_redacted_output_tail(&err); tracing::warn!( - refspec = %refspec, + branch = %branch, error = %fabro_sandbox::display_for_log(&err), "git push from run lifecycle failed" ); @@ -329,7 +341,7 @@ impl RunLifecycle for GitLifecycle { } }; git_result.push_results.push(PushResult { - refspec, + branch: branch.clone(), success: push_ok, exec_output_tail, }); diff --git a/lib/components/fabro-workflow/src/operations/start.rs b/lib/components/fabro-workflow/src/operations/start.rs index 58c9e1c20..67143dc4b 100644 --- a/lib/components/fabro-workflow/src/operations/start.rs +++ b/lib/components/fabro-workflow/src/operations/start.rs @@ -928,16 +928,12 @@ impl RunSession { model: self.pr_model, }; - let concluded = match Box::pin(pipeline::conclude(executed, &finalize_opts)).await { - Ok(concluded) => concluded, - Err(err) => { - self.steering_hub.drain_pending_at_run_end(); - store_progress_logger.flush().await; - return Err(err); - } + let concluding = async { + let concluded = Box::pin(pipeline::conclude(executed, &finalize_opts)).await?; + let published = Box::pin(pipeline::publish(concluded, &publish_opts)).await; + Box::pin(pipeline::finalize(published, &finalize_opts)).await }; - let published = Box::pin(pipeline::publish(concluded, &publish_opts)).await; - let finalized = match Box::pin(pipeline::finalize(published, &finalize_opts)).await { + let finalized = match concluding.await { Ok(finalized) => finalized, Err(err) => { self.steering_hub.drain_pending_at_run_end(); diff --git a/lib/components/fabro-workflow/src/pipeline/finalize.rs b/lib/components/fabro-workflow/src/pipeline/finalize.rs index 673989a5c..73f01568d 100644 --- a/lib/components/fabro-workflow/src/pipeline/finalize.rs +++ b/lib/components/fabro-workflow/src/pipeline/finalize.rs @@ -9,7 +9,7 @@ use fabro_types::{BilledTokenCounts, DiffSummary, EventBody, RunFailure, RunProj use fabro_util::error::collect_causes; use fabro_util::time::elapsed_ms; -use super::types::{Concluded, Executed, FinalizeOptions, Finalized, Published}; +use super::types::{Concluded, Executed, FinalizeOptions, Finalized, PublishOutcome, Published}; use crate::error::{Error, run_failure_from_error, run_failure_from_outcome_failure}; use crate::event::{Event, RunNoticeCode, RunNoticeLevel}; use crate::outcome::{Outcome, StageOutcome}; @@ -44,36 +44,16 @@ pub fn classify_engine_result( }; (status, failure, run_status) } - Err(Error::Cancelled) => ( - StageOutcome::Failed { - retry_requested: false, - }, - Some(run_failure_from_error( - &Error::Cancelled, - FailureReason::Cancelled, - )), - RunStatus::Failed { - reason: FailureReason::Cancelled, - }, - ), - Err(err @ Error::Publish { .. }) => ( - StageOutcome::Failed { - retry_requested: false, - }, - Some(run_failure_from_error(err, FailureReason::PublishFailed)), - RunStatus::Failed { - reason: FailureReason::PublishFailed, - }, - ), - Err(err) => ( - StageOutcome::Failed { - retry_requested: false, - }, - Some(run_failure_from_error(err, FailureReason::WorkflowError)), - RunStatus::Failed { - reason: FailureReason::WorkflowError, - }, - ), + Err(err) => { + let reason = err.failure_reason(); + ( + StageOutcome::Failed { + retry_requested: false, + }, + Some(run_failure_from_error(err, reason)), + RunStatus::Failed { reason }, + ) + } } } @@ -489,13 +469,7 @@ pub(crate) fn build_terminal_event( } let failure = match outcome { - Err(Error::Cancelled) => { - run_failure_from_error(&Error::Cancelled, FailureReason::Cancelled) - } - Err(err @ Error::Publish { .. }) => { - run_failure_from_error(err, FailureReason::PublishFailed) - } - Err(err) => run_failure_from_error(err, FailureReason::WorkflowError), + Err(err) => run_failure_from_error(err, err.failure_reason()), Ok(outcome) => { if let Some(failure) = outcome.failure.as_ref() { run_failure_from_outcome_failure(failure, FailureReason::WorkflowError) @@ -616,25 +590,22 @@ pub async fn finalize(published: Published, options: &FinalizeOptions) -> Result let Published { execution_outcome, publish_outcome, + publish_error, mut conclusion, artifact_count, run_options, services, } = published; - let pushed_branch = publish_outcome - .as_ref() - .ok() - .and_then(|outcome| outcome.pushed_branch()) - .map(str::to_string); - let pr_url = publish_outcome - .as_ref() - .ok() - .and_then(|outcome| outcome.pr_url()) - .map(str::to_string); - let outcome = match (execution_outcome, publish_outcome) { - (Err(error), _) | (Ok(_), Err(error)) => Err(error), - (Ok(outcome), Ok(_)) => Ok(outcome), + let PublishOutcome { + pushed_branch, + pr_url, + } = publish_outcome; + // An execution failure outranks a publish failure: publish only runs after + // a successful execution, so the two are never both set. + let outcome = match (execution_outcome, publish_error) { + (Err(error), _) | (Ok(_), Some(error)) => Err(error), + (Ok(outcome), None) => Ok(outcome), }; let (final_status, failure, _run_status) = classify_engine_result(&outcome); @@ -725,6 +696,7 @@ mod tests { use super::*; use crate::context::Context; + use crate::error::ErrorStage; use crate::event::{Emitter, StoreProgressLogger, append_event}; use crate::run_metadata::{RunMetadataRuntime, RunMetadataWriterHandle}; use crate::run_options::{GitCheckpointOptions, RunOptions}; @@ -1417,10 +1389,8 @@ mod tests { }) .await; - assert!(matches!( - &published.publish_outcome, - Ok(crate::pipeline::PublishOutcome::NotRequested) - )); + assert_eq!(published.publish_outcome, PublishOutcome::default()); + assert!(published.publish_error.is_none()); let finalized = finalize(published, &options).await.unwrap(); assert!(finalized.outcome.is_ok()); @@ -1474,12 +1444,21 @@ mod tests { .await; assert!(matches!( - &published.publish_outcome, - Err(Error::Publish { .. }) + &published.publish_error, + Some(Error::Stage { + stage: ErrorStage::Publish, + .. + }) )); let finalized = finalize(published, &options).await.unwrap(); - assert!(matches!(finalized.outcome, Err(Error::Publish { .. }))); + assert!(matches!( + finalized.outcome, + Err(Error::Stage { + stage: ErrorStage::Publish, + .. + }) + )); assert_eq!( finalized .conclusion @@ -1499,6 +1478,71 @@ mod tests { } } + /// An empty diff means there is nothing to open a pull request for. The + /// branch still gets pushed and the run still succeeds. + #[tokio::test] + async fn empty_diff_pushes_branch_without_opening_pull_request() { + let repo_dir = tempfile::tempdir().unwrap(); + init_git_repo(repo_dir.path()); + let emitter = Arc::new(Emitter::new(test_run_id())); + let events = record_events(&emitter); + let services = test_services( + RunStoreHandle::local(seeded_run_store().await), + emitter, + Arc::new(fabro_agent::LocalSandbox::new( + repo_dir.path().to_path_buf(), + )), + Arc::new(RunMetadataRuntime::new()), + None, + ); + let mut run_options = test_run_options(repo_dir.path()); + run_options.base_branch = Some("main".to_string()); + run_options.git = Some(GitCheckpointOptions { + base_sha: None, + run_branch: Some("fabro/run/test".to_string()), + meta_branch: None, + }); + let executed = test_executed( + Graph::new("test"), + Ok(Outcome::success()), + run_options, + 5, + services, + ); + let options = FinalizeOptions { + run_dir: repo_dir.path().to_path_buf(), + run_id: test_run_id(), + workflow_name: "test".to_string(), + preserve_sandbox: false, + stop_on_terminal: true, + last_git_sha: Some("final-sha".to_string()), + }; + let mut concluded = conclude(executed, &options).await.unwrap(); + concluded.conclusion.diff.patch = None; + let published = crate::pipeline::publish(concluded, &crate::pipeline::PublishOptions { + pr_config: Some(fabro_types::settings::run::PullRequestSettings { + enabled: true, + draft: true, + auto_merge: false, + merge_strategy: fabro_types::settings::run::MergeStrategy::Squash, + }), + github_app: None, + origin_url: Some("https://github.com/owner/repo.git".to_string()), + model: "test-model".to_string(), + }) + .await; + + assert!(published.publish_error.is_none()); + let finalized = finalize(published, &options).await.unwrap(); + + assert!(finalized.outcome.is_ok()); + assert_eq!(finalized.pushed_branch.as_deref(), Some("fabro/run/test")); + assert_eq!(finalized.pr_url, None); + let events = events.lock().unwrap(); + let names = events.iter().map(RunEvent::event_name).collect::>(); + assert_eq!(names, vec!["git.push", "run.completed"]); + } + #[tokio::test] async fn pull_request_failure_precedes_terminal_publish_failure() { let repo_dir = tempfile::tempdir().unwrap(); @@ -1553,7 +1597,16 @@ mod tests { .await; let finalized = finalize(published, &options).await.unwrap(); - assert!(matches!(finalized.outcome, Err(Error::Publish { .. }))); + assert!(matches!( + finalized.outcome, + Err(Error::Stage { + stage: ErrorStage::Publish, + .. + }) + )); + // The push landed before the pull request failed, so the branch is + // still reported — that is exactly the run where the user needs it. + assert_eq!(finalized.pushed_branch.as_deref(), Some("fabro/run/test")); let events = events.lock().unwrap(); let names = events.iter().map(RunEvent::event_name).collect::>(); assert_eq!(names, vec!["git.push", "pull_request.failed", "run.failed"]); diff --git a/lib/components/fabro-workflow/src/pipeline/mod.rs b/lib/components/fabro-workflow/src/pipeline/mod.rs index 57fd59e15..98b813b6a 100644 --- a/lib/components/fabro-workflow/src/pipeline/mod.rs +++ b/lib/components/fabro-workflow/src/pipeline/mod.rs @@ -20,7 +20,7 @@ pub(crate) use persist::persist; pub use publish::publish; pub use pull_request::{ AutoMergeOptions, CreatedPullRequest, OpenPullRequestRequest, PrContent, build_pr_content, - maybe_open_pull_request, + open_pull_request, }; pub use transform::transform; pub use types::{ diff --git a/lib/components/fabro-workflow/src/pipeline/publish.rs b/lib/components/fabro-workflow/src/pipeline/publish.rs index d4ba2fd61..10eca9448 100644 --- a/lib/components/fabro-workflow/src/pipeline/publish.rs +++ b/lib/components/fabro-workflow/src/pipeline/publish.rs @@ -1,10 +1,10 @@ use std::sync::Arc; -use super::pull_request::{AutoMergeOptions, OpenPullRequestRequest, maybe_open_pull_request}; +use super::pull_request::{AutoMergeOptions, OpenPullRequestRequest, open_pull_request}; use super::types::{Concluded, PublishOptions, PublishOutcome, Published}; use crate::error::Error; use crate::event::Event; -use crate::outcome::StageOutcome; +use crate::lifecycle::git::push_run_branch; /// PUBLISH phase: push the final run commit and, when configured, open a pull /// request. @@ -12,7 +12,9 @@ use crate::outcome::StageOutcome; /// Publish is always present in the pipeline. It becomes a no-op when the run /// did not succeed, is a dry run, or has no remote branch configured. pub async fn publish(concluded: Concluded, options: &PublishOptions) -> Published { - let publish_outcome = publish_inner(&concluded, options).await; + let mut publish_outcome = PublishOutcome::default(); + let publish_error = concluded.publish(options, &mut publish_outcome).await.err(); + let Concluded { outcome, conclusion, @@ -25,6 +27,7 @@ pub async fn publish(concluded: Concluded, options: &PublishOptions) -> Publishe Published { execution_outcome: outcome, publish_outcome, + publish_error, conclusion, artifact_count, run_options, @@ -32,170 +35,153 @@ pub async fn publish(concluded: Concluded, options: &PublishOptions) -> Publishe } } -async fn publish_inner( - concluded: &Concluded, - options: &PublishOptions, -) -> Result { - let successful_execution = concluded.outcome.as_ref().is_ok_and(|outcome| { - matches!( - outcome.status, - StageOutcome::Succeeded | StageOutcome::PartiallySucceeded - ) - }); - if !successful_execution || concluded.run_options.dry_run_enabled() { - return Ok(PublishOutcome::NotRequested); - } +impl Concluded { + /// Run the publish steps, recording each one into `outcome` as it lands. + /// + /// `outcome` accumulates what actually happened, so a branch that reached + /// the remote is still reported when pull request creation later fails. + async fn publish( + &self, + options: &PublishOptions, + outcome: &mut PublishOutcome, + ) -> Result<(), Error> { + // A run that did not succeed, or that never intended to touch the + // remote, has nothing to publish — even when a pull request was asked + // for. Only a run that got far enough to publish can fail publishing. + if !self + .outcome + .as_ref() + .is_ok_and(|o| o.status.is_successful()) + || self.run_options.dry_run_enabled() + { + return Ok(()); + } - let pull_request_requested = options.pr_config.is_some(); - let Some(origin_url) = options - .origin_url - .as_deref() - .filter(|origin| !origin.trim().is_empty()) - else { - if pull_request_requested { - return Err(pull_request_error( - concluded, - "pull request creation requires a GitHub origin URL", - )); - } - return Ok(PublishOutcome::NotRequested); - }; - let Some(run_branch) = concluded.run_options.run_branch() else { - if pull_request_requested { - return Err(pull_request_error( - concluded, - "pull request creation requires a run branch", - )); - } - return Ok(PublishOutcome::NotRequested); - }; - if !concluded.run_options.settings.run.run_branch.push { - if pull_request_requested { - return Err(pull_request_error( - concluded, - "pull request creation requires run branch pushing", - )); - } - return Ok(PublishOutcome::NotRequested); - } + let pull_request_requested = options.pr_config.is_some(); + let (origin_url, run_branch) = match self.publish_target(options) { + Ok(target) => target, + Err(_) if !pull_request_requested => return Ok(()), + Err(reason) => return Err(self.pull_request_error(reason)), + }; - let final_sha = concluded - .conclusion - .final_git_commit_sha - .as_deref() - .ok_or_else(|| Error::publish("cannot publish a run without a final git commit SHA"))?; - let refspec = format!("refs/heads/{run_branch}:refs/heads/{run_branch}"); - match concluded.services.sandbox.git_push_ref(&refspec).await { - Ok(()) => { - concluded.services.emitter.emit(&Event::GitPush { - branch: run_branch.to_string(), - success: true, - exec_output_tail: None, - }); - } - Err(error) => { - let exec_output_tail = fabro_sandbox::default_redacted_output_tail(&error); - concluded.services.emitter.emit(&Event::GitPush { - branch: run_branch.to_string(), - success: false, - exec_output_tail: exec_output_tail.clone(), - }); - return Err(Error::publish_with_source_and_exec_output_tail( - format!("failed to push final commit {final_sha} to branch '{run_branch}'"), - error, - exec_output_tail, - )); - } - } + let final_sha = self + .conclusion + .final_git_commit_sha + .as_deref() + .ok_or_else(|| Error::publish("cannot publish a run without a final git commit SHA"))?; - let diff = concluded - .conclusion - .diff - .patch - .as_deref() - .unwrap_or_default(); - let Some(pr_config) = options.pr_config.as_ref() else { - return Ok(PublishOutcome::Published { - pushed_branch: run_branch.to_string(), - pr_url: None, - }); - }; - if diff.trim().is_empty() { - return Ok(PublishOutcome::NoChanges { - pushed_branch: run_branch.to_string(), - }); - } + self.push_final_commit(run_branch, final_sha).await?; + outcome.pushed_branch = Some(run_branch.to_string()); - let base_branch = concluded - .run_options - .base_branch - .as_deref() - .ok_or_else(|| { - pull_request_error(concluded, "pull request creation requires a base branch") + let Some(pr_config) = options.pr_config.as_ref() else { + return Ok(()); + }; + let diff = self.conclusion.diff.patch.as_deref().unwrap_or_default(); + if diff.trim().is_empty() { + return Ok(()); + } + + let base_branch = self.run_options.base_branch.as_deref().ok_or_else(|| { + self.pull_request_error("pull request creation requires a base branch") })?; - let credentials = options.github_app.as_ref().ok_or_else(|| { - pull_request_error( - concluded, - "pull request creation requires GitHub credentials", - ) - })?; - let auto_merge = pr_config.auto_merge.then_some(AutoMergeOptions { - merge_strategy: pr_config.merge_strategy, - }); - let github_base_url = fabro_github::github_api_base_url(); + let credentials = options.github_app.as_ref().ok_or_else(|| { + self.pull_request_error("pull request creation requires GitHub credentials") + })?; + let github_base_url = fabro_github::github_api_base_url(); - let created = maybe_open_pull_request(OpenPullRequestRequest { - github: fabro_github::GitHubContext::new(credentials, &github_base_url), - origin_url, - base_branch, - head_branch: run_branch, - expected_head_sha: final_sha, - goal: concluded.graph.goal(), - diff, - model: &options.model, - draft: pr_config.draft, - auto_merge, - run_store: &concluded.services.run_store, - llm_source: concluded.services.llm_source.as_ref(), - catalog: Arc::clone(&concluded.services.catalog), - conclusion: Some(&concluded.conclusion), - run_state: None, - }) - .await - .map_err(|error| { - concluded.services.emitter.emit(&Event::PullRequestFailed { - error: error.clone(), - }); - Error::publish_with_source("failed to create pull request", anyhow::anyhow!(error)) - })? - .ok_or_else(|| { - pull_request_error( - concluded, - "pull request creation found no changes after the stored diff was checked", - ) - })?; + let created = open_pull_request(OpenPullRequestRequest { + github: fabro_github::GitHubContext::new(credentials, &github_base_url), + origin_url, + base_branch, + head_branch: run_branch, + expected_head_sha: final_sha, + goal: self.graph.goal(), + diff, + model: &options.model, + draft: pr_config.draft, + auto_merge: pr_config.auto_merge.then_some(AutoMergeOptions { + merge_strategy: pr_config.merge_strategy, + }), + run_store: &self.services.run_store, + llm_source: self.services.llm_source.as_ref(), + catalog: Arc::clone(&self.services.catalog), + conclusion: Some(&self.conclusion), + run_state: None, + }) + .await + .map_err(|error| { + self.services.emitter.emit(&Event::PullRequestFailed { + error: error.clone(), + }); + Error::publish_with_source("failed to create pull request", anyhow::anyhow!(error)) + })?; - concluded - .services - .emitter - .emit(&Event::pull_request_created( + self.services.emitter.emit(&Event::pull_request_created( &created.link, &created.base_branch, &created.head_branch, - &created.head_sha, + final_sha, &created.title, pr_config.draft, )); + outcome.pr_url = Some(created.link.html_url()); - Ok(PublishOutcome::Published { - pushed_branch: run_branch.to_string(), - pr_url: Some(created.link.html_url()), - }) -} + Ok(()) + } -fn pull_request_error(concluded: &Concluded, message: &str) -> Error { - concluded.services.emitter.emit(&Event::PullRequestFailed { - error: message.to_string(), - }); - Error::publish(message) + /// The origin and run branch to publish to. + /// + /// `Err` carries why there is no target. That is only a failure when a + /// pull request was requested; otherwise publish just has nothing to do. + fn publish_target<'a>( + &'a self, + options: &'a PublishOptions, + ) -> Result<(&'a str, &'a str), &'static str> { + let origin_url = options + .origin_url + .as_deref() + .filter(|origin| !origin.trim().is_empty()) + .ok_or("pull request creation requires a GitHub origin URL")?; + let run_branch = self + .run_options + .run_branch() + .ok_or("pull request creation requires a run branch")?; + if !self.run_options.settings.run.run_branch.push { + return Err("pull request creation requires run branch pushing"); + } + Ok((origin_url, run_branch)) + } + + async fn push_final_commit(&self, run_branch: &str, final_sha: &str) -> Result<(), Error> { + match push_run_branch(self.services.sandbox.as_ref(), run_branch).await { + Ok(()) => { + self.services.emitter.emit(&Event::GitPush { + branch: run_branch.to_string(), + success: true, + exec_output_tail: None, + }); + Ok(()) + } + Err(error) => { + let exec_output_tail = fabro_sandbox::default_redacted_output_tail(&error); + self.services.emitter.emit(&Event::GitPush { + branch: run_branch.to_string(), + success: false, + exec_output_tail: exec_output_tail.clone(), + }); + Err(Error::publish_with_source_and_exec_output_tail( + format!("failed to push final commit {final_sha} to branch '{run_branch}'"), + error, + exec_output_tail, + )) + } + } + } + + fn pull_request_error(&self, message: &str) -> Error { + self.services.emitter.emit(&Event::PullRequestFailed { + error: message.to_string(), + }); + Error::publish(message) + } } diff --git a/lib/components/fabro-workflow/src/pipeline/pull_request.rs b/lib/components/fabro-workflow/src/pipeline/pull_request.rs index 638c87847..d66072923 100644 --- a/lib/components/fabro-workflow/src/pipeline/pull_request.rs +++ b/lib/components/fabro-workflow/src/pipeline/pull_request.rs @@ -441,7 +441,7 @@ pub struct AutoMergeOptions { pub merge_strategy: MergeStrategy, } -/// Inputs for [`maybe_open_pull_request`]. +/// Inputs for [`open_pull_request`]. pub struct OpenPullRequestRequest<'a> { pub github: github_app::GitHubContext<'a>, pub origin_url: &'a str, @@ -468,21 +468,15 @@ pub struct CreatedPullRequest { pub title: String, pub base_branch: String, pub head_branch: String, - pub head_sha: String, } -/// Optionally open a pull request after a successful workflow run. +/// Open a pull request for a completed run. /// -/// Returns `Ok(Some(CreatedPullRequest))` if a PR was created, `Ok(None)` if -/// the diff was empty, or `Err` on failure. -pub async fn maybe_open_pull_request( +/// Callers are responsible for skipping runs with an empty diff; reaching here +/// means a pull request is expected, so every failure is an error. +pub async fn open_pull_request( req: OpenPullRequestRequest<'_>, -) -> Result, String> { - if req.diff.is_empty() { - debug!("Empty diff, skipping pull request creation"); - return Ok(None); - } - +) -> Result { let https_url = ssh_url_to_https(req.origin_url); let (owner, repo) = github_app::parse_github_owner_repo(&https_url).map_err(|err| format!("{err:#}"))?; @@ -565,13 +559,12 @@ pub async fn maybe_open_pull_request( number: created.number, }; - Ok(Some(CreatedPullRequest { + Ok(CreatedPullRequest { link, title, base_branch: req.base_branch.to_string(), head_branch: req.head_branch.to_string(), - head_sha: req.expected_head_sha.to_string(), - })) + }) } #[cfg(test)] @@ -581,7 +574,7 @@ mod tests { use std::time::Duration; use chrono::Utc; - use fabro_auth::{CredentialSource, EnvCredentialSource, VaultCredentialSource}; + use fabro_auth::{CredentialSource, VaultCredentialSource}; use fabro_graphviz::graph::Graph; use fabro_llm::Error as LlmError; use fabro_llm::client::Client; @@ -687,10 +680,6 @@ mod tests { )) } - fn test_catalog() -> Arc { - Arc::new(Catalog::from_builtin().expect("default catalog should build")) - } - fn test_catalog_with_provider_base_url(provider: &str, base_url: &str) -> Arc { let mut settings = LlmCatalogSettings::default(); settings @@ -718,10 +707,6 @@ mod tests { )) } - fn test_llm_source() -> Arc { - Arc::new(EnvCredentialSource::new()) - } - fn test_projection() -> RunProjection { RunProjection::new( "Test run".to_string(), @@ -1400,46 +1385,12 @@ mod tests { ); } - #[tokio::test] - async fn empty_diff_returns_none() { - let store = test_store(); - let run_store = store.create_run(&fixtures::RUN_1).await.unwrap(); - let run_store_handle: RunStoreHandle = run_store.into(); - let llm_source = test_llm_source(); - let creds = fabro_github::GitHubCredentials::App(fabro_github::GitHubAppCredentials { - app_id: "123".to_string(), - private_key_pem: "unused".to_string(), - slug: None, - }); - let base_url = github_app::github_api_base_url(); - let result = maybe_open_pull_request(OpenPullRequestRequest { - github: github_app::GitHubContext::new(&creds, &base_url), - origin_url: "https://github.com/owner/repo.git", - base_branch: "main", - head_branch: "fabro/run/123", - expected_head_sha: "final-sha", - goal: "Fix bug", - diff: "", - model: "claude-sonnet-4-20250514", - draft: false, - auto_merge: None, - run_store: &run_store_handle, - llm_source: llm_source.as_ref(), - catalog: test_catalog(), - conclusion: None, - run_state: None, - }) - .await; - assert!(result.is_ok()); - assert!(result.unwrap().is_none()); - } - #[tokio::test] async fn stale_remote_branch_is_rejected_before_pull_request_creation() { let payload = pr_content_json("Fix bug", "Narrative."); let harness = setup_fallback_test_harness_with_branch_sha(&payload, "stale-sha").await; let github_base_url = harness.github_server.url(""); - let error = maybe_open_pull_request(OpenPullRequestRequest { + let error = open_pull_request(OpenPullRequestRequest { github: fabro_github::GitHubContext::new(&harness.creds, &github_base_url), origin_url: "https://github.com/owner/repo.git", base_branch: "main", @@ -1614,9 +1565,9 @@ mod tests { assert!(body.contains("Generated with [Fabro](https://fabro.sh)")); } - // ── maybe_open_pull_request fallback tests ────────────────────────── + // ── open_pull_request fallback tests ────────────────────────── - /// Set of mock servers and credentials for the `maybe_open_pull_request` + /// Set of mock servers and credentials for the `open_pull_request` /// fallback path. The builder's `Client::from_source` rebuilds the LLM /// client from the credential source, so the in-process MockProvider /// cannot intercept — we mock the OpenAI HTTP endpoint instead. @@ -1812,14 +1763,14 @@ mod tests { /// falls back to `pr_title_from_goal` (first line, decoration stripped) /// and PR creation succeeds with that title. #[tokio::test] - async fn maybe_open_pull_request_falls_back_to_goal_title_when_llm_returns_empty_title() { + async fn open_pull_request_falls_back_to_goal_title_when_llm_returns_empty_title() { let payload = pr_content_json("", "Narrative."); let harness = setup_fallback_test_harness(&payload).await; let github_base_url = harness.github_server.url(""); let github = github_app::GitHubContext::new(&harness.creds, &github_base_url); - let result = maybe_open_pull_request(OpenPullRequestRequest { + let result = open_pull_request(OpenPullRequestRequest { github, origin_url: "https://github.com/owner/repo.git", base_branch: "main", @@ -1839,15 +1790,14 @@ mod tests { .await .expect("PR creation should succeed"); - let record = result.expect("PR record should be Some"); - assert_eq!(record.title, "Fix telemetry leak"); + assert_eq!(result.title, "Fix telemetry leak"); harness.assert_mocks_called_once().await; } /// LLM returns an empty title; the content builder fallback still caps /// the deterministic goal title at 72 chars ending with `…`. #[tokio::test] - async fn maybe_open_pull_request_caps_fallback_title_at_72_chars() { + async fn open_pull_request_caps_fallback_title_at_72_chars() { let payload = pr_content_json("", "Narrative."); let harness = setup_fallback_test_harness(&payload).await; @@ -1857,7 +1807,7 @@ mod tests { // Single ~200-char line, no `Plan:` / heading prefix, no newlines. let goal = "x".repeat(200); - let result = maybe_open_pull_request(OpenPullRequestRequest { + let result = open_pull_request(OpenPullRequestRequest { github, origin_url: "https://github.com/owner/repo.git", base_branch: "main", @@ -1877,8 +1827,7 @@ mod tests { .await .expect("PR creation should succeed"); - let record = result.expect("PR record should be Some"); - let title = record.title; + let title = result.title; assert_eq!(title.chars().count(), 72); assert!(title.ends_with('\u{2026}')); harness.assert_mocks_called_once().await; diff --git a/lib/components/fabro-workflow/src/pipeline/types.rs b/lib/components/fabro-workflow/src/pipeline/types.rs index c7cfcc815..19bbc42cf 100644 --- a/lib/components/fabro-workflow/src/pipeline/types.rs +++ b/lib/components/fabro-workflow/src/pipeline/types.rs @@ -348,41 +348,23 @@ pub struct Concluded { pub services: Arc, } -#[derive(Debug, Clone, PartialEq, Eq)] -pub enum PublishOutcome { - NotRequested, - NoChanges { - pushed_branch: String, - }, - Published { - pushed_branch: String, - pr_url: Option, - }, -} - -impl PublishOutcome { - pub fn pushed_branch(&self) -> Option<&str> { - match self { - Self::NotRequested => None, - Self::NoChanges { pushed_branch } | Self::Published { pushed_branch, .. } => { - Some(pushed_branch) - } - } - } - - pub fn pr_url(&self) -> Option<&str> { - match self { - Self::Published { pr_url, .. } => pr_url.as_deref(), - Self::NotRequested | Self::NoChanges { .. } => None, - } - } +/// What the PUBLISH phase actually accomplished. +/// +/// Recorded separately from the phase's error so a branch that reached the +/// remote is still reported when a later step, such as pull request creation, +/// fails. An all-`None` value means publish had nothing to do. +#[derive(Debug, Clone, Default, PartialEq, Eq)] +pub struct PublishOutcome { + pub pushed_branch: Option, + pub pr_url: Option, } /// Output of the PUBLISH phase. #[non_exhaustive] pub struct Published { pub execution_outcome: Result, - pub publish_outcome: Result, + pub publish_outcome: PublishOutcome, + pub publish_error: Option, pub conclusion: Conclusion, pub artifact_count: usize, pub run_options: RunOptions, diff --git a/lib/components/fabro-workflow/src/pull_request.rs b/lib/components/fabro-workflow/src/pull_request.rs index bb9a74744..78723914d 100644 --- a/lib/components/fabro-workflow/src/pull_request.rs +++ b/lib/components/fabro-workflow/src/pull_request.rs @@ -1,4 +1,4 @@ pub use crate::pipeline::{ AutoMergeOptions, CreatedPullRequest, OpenPullRequestRequest, PrContent, build_pr_content, - maybe_open_pull_request, + open_pull_request, };