fix(workflow): make publish failures terminal

This commit is contained in:
Bryan Helmkamp 2026-07-27 11:25:18 -04:00
parent 2bcf94fed8
commit 1c82bd9008
No known key found for this signature in database
27 changed files with 989 additions and 338 deletions

View file

@ -2116,6 +2116,7 @@ These legacy events may appear in older run logs. Current CLI backend runs do no
"properties": {
"pr_url": "https://github.com/org/repo/pull/42",
"pr_number": 42,
"head_sha": "d34db33f",
"draft": true
}
}
@ -2125,6 +2126,7 @@ These legacy events may appear in older run logs. Current CLI backend runs do no
|----------|------|-------------|
| `pr_url` | string | Pull request URL |
| `pr_number` | number | Pull request number |
| `head_sha` | string (optional) | Verified commit SHA at the remote PR head; absent on older events |
| `draft` | boolean | Whether the PR is a draft |
### `pull_request.linked`

View file

@ -8894,6 +8894,7 @@ components:
type: string
enum:
- workflow_error
- publish_failed
- cancelled
- approval_denied
- terminated

View file

@ -55,6 +55,7 @@ When you choose the GitHub App strategy, the CLI opens GitHub with a pre-filled
| Permission | Level | Purpose |
|---|---|---|
| Contents | Write | Clone repos, push run branches and checkpoints |
| Workflows | Write | Push changes under `.github/workflows/` |
| Metadata | Read | Look up repository installation status |
| Pull requests | Write | Create and update PRs from workflows |
| Checks | Write | Report workflow status on commits |
@ -219,7 +220,7 @@ When a workflow runs in a remote sandbox (Daytona or Docker), Fabro clones the c
2. SSH URLs (e.g. `git@github.com:owner/repo.git`) are converted to HTTPS
3. Fabro signs a short-lived JWT using the App ID and private key (RS256, 10-minute validity)
4. Using the JWT, Fabro looks up the GitHub App installation for the repository (`GET /repos/\{owner\}/\{repo\}/installation`)
5. Fabro requests a scoped Installation Access Token with `contents: write` permission on the specific repository
5. Fabro requests a scoped Installation Access Token with `contents: write` and `workflows: write` permissions on the specific repository
6. The sandbox clones via HTTPS using `x-access-token` as the username and the token as the password
For public repositories, the clone works without credentials. The token is still generated because it's needed for pushing checkpoints.
@ -248,7 +249,9 @@ The upper bound on what Fabro will mint is whatever permissions the GitHub App i
### Checkpoint pushing
After each workflow stage, Fabro [checkpoints](/execution/checkpoints) by pushing the run branch and metadata branch to origin. Inside remote sandboxes, the git remote URL is configured with the Installation Access Token for authenticated pushing.
After each workflow stage, Fabro [checkpoints](/execution/checkpoints) by pushing the run branch and metadata branch to origin. Before a successful run becomes terminal, the publish stage pushes the final commit again and treats failure as a run failure. Inside remote sandboxes, the git remote URL is configured with the Installation Access Token for authenticated pushing.
When pull request creation is enabled, Fabro then checks that GitHub reports the run branch at the exact final commit before opening the PR. A failed final push, branch check, or PR creation marks the run as failed with `publish_failed`; the terminal run event is emitted only after this step finishes.
For long-running workflows, Fabro refreshes the token before each push since Installation Access Tokens are short-lived (typically 1 hour).

View file

@ -614,6 +614,7 @@ mod tests {
repo: "widgets".into(),
base_branch: "main".into(),
head_branch: "fabro/run/42".into(),
head_sha: "final-sha".into(),
title: "Ship the server-side PR".into(),
draft: true,
};

View file

@ -1383,6 +1383,7 @@ mod tests {
repo: "fabro".into(),
base_branch: "main".into(),
head_branch: "fabro/run/42".into(),
head_sha: "final-sha".into(),
title: "Ship the change".into(),
draft: true,
});

View file

@ -85,6 +85,7 @@ fn pr_view_reads_pull_request_from_store_without_pull_request_json() {
repo: "fabro".to_string(),
base_branch: "main".to_string(),
head_branch: "fabro/run/demo".to_string(),
head_sha: Some("final-sha".to_string()),
title: "Map the constellations".to_string(),
draft: false,
}),

View file

@ -1280,6 +1280,7 @@ mod runs {
fn parse_failure_reason(reason: &str) -> Option<FailureReason> {
match reason {
"workflow_error" => Some(FailureReason::WorkflowError),
"publish_failed" => Some(FailureReason::PublishFailed),
"cancelled" => Some(FailureReason::Cancelled),
"approval_denied" => Some(FailureReason::ApprovalDenied),
"terminated" => Some(FailureReason::Terminated),

View file

@ -2053,6 +2053,7 @@ fn build_github_app_manifest(
"public": false,
"default_permissions": {
"contents": "write",
"workflows": "write",
"metadata": "read",
"pull_requests": "write",
"checks": "write",
@ -2354,11 +2355,27 @@ mod tests {
InstallObjectStoreCredentialMode, InstallObjectStoreInput, InstallObjectStoreProvider,
InstallObjectStoreState, InstallSandboxProviderState, InstallSandboxState,
InstallTokenQuery, LlmProvidersInput, PendingInstall, ServerConfigInput, ServerSecrets,
classify_object_store_validation_error, detect_canonical_url, install_object_store_lookup,
lock_unpoisoned, post_install_finish, provider_base_url_override,
resolve_install_object_store_state, token_is_valid, write_artifact_store_metadata,
build_github_app_manifest, classify_object_store_validation_error, detect_canonical_url,
install_object_store_lookup, lock_unpoisoned, post_install_finish,
provider_base_url_override, resolve_install_object_store_state, token_is_valid,
write_artifact_store_metadata,
};
#[test]
fn github_app_manifest_allows_workflow_file_writes() {
let manifest = build_github_app_manifest(
"Fabro Test",
"https://fabro.example/setup",
"https://fabro.example/auth/callback/github",
"https://fabro.example/setup",
);
assert_eq!(
manifest["default_permissions"]["workflows"],
serde_json::Value::String("write".to_string())
);
}
#[test]
fn token_validation_accepts_any_matching_source() {
let state = InstallAppState::for_test("expected");

View file

@ -4183,6 +4183,14 @@ async fn execute_run_in_process(state: Arc<AppState>, run_id: RunId) {
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 {
@ -4197,6 +4205,14 @@ async fn execute_run_in_process(state: Arc<AppState>, run_id: RunId) {
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 {

View file

@ -165,6 +165,7 @@ struct RunPrInputs<'a> {
goal: &'a str,
base_branch: &'a str,
run_branch: &'a str,
final_git_sha: &'a str,
diff: &'a str,
conclusion: &'a fabro_types::Conclusion,
normalized_origin: String,
@ -224,6 +225,17 @@ impl<'a> RunPrInputs<'a> {
"run_not_finished",
)
})?;
let final_git_sha = conclusion
.final_git_commit_sha
.as_deref()
.filter(|sha| !sha.trim().is_empty())
.ok_or_else(|| {
ApiError::with_code(
StatusCode::BAD_REQUEST,
"Run has no final git commit SHA — the remote branch cannot be verified.",
"missing_final_git_commit",
)
})?;
if !force && !conclusion.status.is_successful() {
return Err(ApiError::with_code(
StatusCode::BAD_REQUEST,
@ -240,6 +252,7 @@ impl<'a> RunPrInputs<'a> {
goal: run_spec.graph.goal(),
base_branch,
run_branch,
final_git_sha,
diff,
conclusion,
normalized_origin,
@ -323,6 +336,7 @@ async fn create_run_pull_request(
origin_url: &inputs.normalized_origin,
base_branch: inputs.base_branch,
head_branch: inputs.run_branch,
expected_head_sha: inputs.final_git_sha,
goal: inputs.goal,
diff: inputs.diff,
model: &model,
@ -350,6 +364,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,
&created_pull_request.title,
true,
);

View file

@ -4685,6 +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(),
title: "Ship <prod> & notify".to_string(),
draft: false,
},
@ -6425,6 +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(),
title: title.to_string(),
draft: false,
},
@ -6519,7 +6521,7 @@ async fn create_completed_run_ready_for_pull_request(
status: "succeeded".to_string(),
reason: SuccessReason::Completed,
total_usd_micros: None,
final_git_commit_sha: None,
final_git_commit_sha: Some("final-sha".to_string()),
final_patch: Some(final_patch.to_string()),
diff_summary: None,
billing: None,
@ -9058,6 +9060,14 @@ async fn get_run_pull_request_returns_stored_github_association_when_github_pr_i
#[tokio::test]
async fn create_run_pull_request_creates_and_persists_record() {
let github = MockServer::start();
let branch_mock = github.mock(|when, then| {
when.method("GET")
.path("/repos/acme/widgets/branches/fabro/run/42")
.header("authorization", "Bearer ghu_test");
then.status(200)
.header("content-type", "application/json")
.body(json!({ "commit": { "sha": "final-sha" } }).to_string());
});
let create_mock = github.mock(|when, then| {
when.method("POST")
.path("/repos/acme/widgets/pulls")
@ -9155,6 +9165,7 @@ async fn create_run_pull_request_creates_and_persists_record() {
assert_eq!(state_body["pull_request"]["repo"], "widgets");
response_mock.assert_async().await;
branch_mock.assert();
create_mock.assert();
}
@ -16477,6 +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(),
title: "Fix board metadata".to_string(),
draft: false,
},

View file

@ -587,7 +587,10 @@ async fn mint_installation_token_with_jwt(
})
}
/// Request a scoped Installation Access Token with `contents: write`.
/// Request a scoped Installation Access Token for git writes.
///
/// The `workflows` permission is required when a pushed commit creates or
/// updates files under `.github/workflows/`.
pub async fn create_installation_access_token(
client: &impl HttpClient,
jwt: &str,
@ -601,7 +604,7 @@ pub async fn create_installation_access_token(
owner,
repo,
base_url,
serde_json::json!({ "contents": "write" }),
serde_json::json!({ "contents": "write", "workflows": "write" }),
)
.await
}
@ -899,6 +902,55 @@ pub async fn branch_exists(
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.
pub async fn branch_head_sha(
ctx: &GitHubContext<'_>,
owner: &str,
repo: &str,
branch: &str,
) -> anyhow::Result<Option<String>> {
#[derive(Deserialize)]
struct BranchResponse {
commit: BranchCommit,
}
#[derive(Deserialize)]
struct BranchCommit {
sha: String,
}
let client = ctx.http_client()?;
let token = ctx
.creds
.resolve_bearer_token(
&client,
owner,
repo,
ctx.base_url,
serde_json::json!({ "contents": "read" }),
)
.await?;
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)
.await
.context("Failed to read remote branch head")?;
match resp.status {
200 => {
let branch: BranchResponse = resp
.json()
.context("Failed to parse remote branch response")?;
Ok(Some(branch.commit.sha))
}
404 => Ok(None),
status => bail!("Unexpected status {status} reading branch '{branch}'"),
}
}
async fn branch_exists_with_client(
client: &impl HttpClient,
ctx: &GitHubContext<'_>,
@ -1032,26 +1084,47 @@ pub async fn update_app_webhook_config(
/// Resolve git clone credentials for a GitHub repository.
///
/// Returns `(username, password)` for authenticated cloning.
/// Returns `(username, password)` for authenticated cloning and pushing.
/// Always generates a token regardless of repo visibility, since the token
/// is needed for pushing from the sandbox.
/// is needed for pushing from the sandbox. The token includes `workflows:
/// write` so a run can publish workflow-file changes.
pub async fn resolve_clone_credentials(
ctx: &GitHubContext<'_>,
owner: &str,
repo: &str,
) -> anyhow::Result<(Option<String>, Option<String>)> {
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<String>, Option<String>)> {
let token = match ctx.creds {
GitHubCredentials::Pat(token) => token.clone(),
GitHubCredentials::Installation(token) => token.valid_token()?.to_string(),
GitHubCredentials::App(_) => {
let client = ctx.http_client()?;
ctx.creds
.resolve_bearer_token(
&client,
client,
owner,
repo,
ctx.base_url,
serde_json::json!({ "contents": "write" }),
serde_json::json!({ "contents": "write", "workflows": "write" }),
)
.await?
}
@ -1775,7 +1848,9 @@ mod tests {
r#"{"token": "ghs_xxx", "expires_at": "2099-01-01T00:00:00Z"}"#,
)
.with_req_header("Authorization", "Bearer test-jwt")
.with_req_body(r#"{"permissions":{"contents":"write"},"repositories":["repo"]}"#);
.with_req_body(
r#"{"permissions":{"contents":"write","workflows":"write"},"repositories":["repo"]}"#,
);
let token = create_installation_access_token(&mock, "test-jwt", "owner", "repo", "")
.await
@ -2296,6 +2371,43 @@ mod tests {
);
}
#[tokio::test]
async fn resolve_clone_credentials_requests_workflow_write_permission() {
let mock = MockHttpClient::new()
.on(
HttpMethod::Get,
"/repos/owner/repo/installation",
200,
r#"{"id": 123}"#,
)
.on(
HttpMethod::Post,
"/app/installations/123/access_tokens",
201,
r#"{"token": "ghs_xxx", "expires_at": "2099-01-01T00:00:00Z"}"#,
)
.with_req_body(
r#"{"permissions":{"contents":"write","workflows":"write"},"repositories":["repo"]}"#,
);
let credentials = GitHubCredentials::App(GitHubAppCredentials {
app_id: "test".to_string(),
private_key_pem: test_rsa_key().to_string(),
slug: None,
});
let context = GitHubContext::new(&credentials, "");
let resolved = resolve_clone_credentials_with_client(&mock, &context, "owner", "repo")
.await
.unwrap();
assert_eq!(
resolved,
(
Some("x-access-token".to_string()),
Some("ghs_xxx".to_string())
)
);
}
#[test]
fn installation_token_valid_token_rejects_expired_tokens() {
let expired = InstallationToken {

View file

@ -3791,6 +3791,7 @@ mod tests {
repo: "fabro".to_string(),
base_branch: "main".to_string(),
head_branch: "fabro/run/demo".to_string(),
head_sha: Some("final-sha".to_string()),
title: "Add run PR chip".to_string(),
draft: false,
}),
@ -3840,6 +3841,7 @@ mod tests {
repo: github_pull_request.repo.clone(),
base_branch: "main".to_string(),
head_branch: "fabro/run/demo".to_string(),
head_sha: Some("final-sha".to_string()),
title: "Add run PR chip".to_string(),
draft: false,
}),

View file

@ -67,7 +67,8 @@ assert_eq!(graph.goal(), "Run tests");
use fabro_workflow::operations::start;
use fabro_workflow::pipeline;
// Use `operations::start(...)` for the full initialize -> execute -> finalize flow.
// Use `operations::start(...)` for the full
// initialize -> execute -> conclude -> publish -> finalize flow.
// Use `pipeline::initialize(...)` + `pipeline::execute(...)` when you need partial lifecycle control.
```

View file

@ -285,6 +285,15 @@ pub enum Error {
source: Option<SharedError>,
},
#[error("Publish error: {message}")]
Publish {
message: String,
failure_class: FailureCategory,
exec_output_tail: Option<ExecOutputTail>,
#[source]
source: Option<SharedError>,
},
#[error("Handler error: {message}")]
Handler {
message: String,
@ -420,10 +429,49 @@ impl Error {
Self::engine_with_source(message, source)
}
/// Build an error for the required publish stage.
pub fn publish(message: impl Into<String>) -> Self {
let message = message.into();
let failure_class = classify_failure_reason(&message);
Self::Publish {
message,
failure_class,
exec_output_tail: None,
source: None,
}
}
pub fn publish_with_source(
message: impl Into<String>,
source: impl Into<anyhow::Error>,
) -> Self {
Self::publish_with_source_and_exec_output_tail(message, source, None)
}
pub fn publish_with_source_and_exec_output_tail(
message: impl Into<String>,
source: impl Into<anyhow::Error>,
exec_output_tail: Option<ExecOutputTail>,
) -> 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),
}
}
#[must_use]
pub fn causes(&self) -> Vec<String> {
match self {
Self::Engine { source, .. } | Self::Handler { source, .. } => source
Self::Engine { source, .. }
| Self::Publish { source, .. }
| Self::Handler { source, .. } => source
.as_ref()
.map_or_else(Vec::new, |source| collect_chain(source)),
Self::Template { source, .. } => collect_chain(source),
@ -439,15 +487,17 @@ impl Error {
/// Whether this error category is retryable (transient) or terminal.
///
/// Retryable: Handler (transient handler failures), Engine (could be
/// transient), Io (network/disk issues are often transient),
/// Llm (delegates to SdkError). Terminal: Parse, Validation,
/// OutputSchemaValidation, Stylesheet (configuration errors), Checkpoint
/// (storage integrity), Cancelled (explicit cancellation).
/// 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.
#[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::Llm(sdk_err) => sdk_err.retryable(),
Self::Parse(_)
| Self::Validation(_)
@ -483,9 +533,9 @@ impl Error {
| Self::Unsupported(_)
| Self::OutputSchemaValidation(_) => FailureCategory::Deterministic,
Self::Precondition(_) | Self::RunNotFound(_) => FailureCategory::Structural,
Self::Handler { failure_class, .. } | Self::Engine { failure_class, .. } => {
*failure_class
}
Self::Handler { failure_class, .. }
| Self::Engine { failure_class, .. }
| Self::Publish { failure_class, .. } => *failure_class,
}
}
@ -502,13 +552,18 @@ impl Error {
#[must_use]
pub fn to_failure_detail(&self) -> FailureDetail {
let message = match self {
Self::Engine { message, .. } | Self::Handler { message, .. } => message.clone(),
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(),
@ -1974,6 +2029,7 @@ mod tests {
}],
},
Error::engine("engine err"),
Error::publish("publish err"),
Error::handler("handler err"),
Error::Llm(SdkError::Network {
message: "refused".into(),
@ -2005,6 +2061,12 @@ mod tests {
);
}
#[test]
fn publish_error_is_only_retryable_for_transient_failures() {
assert!(Error::publish("connection timed out").is_retryable());
assert!(!Error::publish("permission denied").is_retryable());
}
#[test]
fn failure_class_stability() {
let messages = [

View file

@ -1311,6 +1311,7 @@ fn event_body_from_event(event: &Event) -> EventBody {
repo,
base_branch,
head_branch,
head_sha,
title,
draft,
} => EventBody::PullRequestCreated(fabro_types::PullRequestCreatedProps {
@ -1320,6 +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()),
title: title.clone(),
draft: *draft,
}),

View file

@ -730,6 +730,8 @@ pub enum Event {
repo: String,
base_branch: String,
head_branch: String,
#[serde(default)]
head_sha: String,
title: String,
draft: bool,
},
@ -769,6 +771,7 @@ impl Event {
record: &PullRequestLink,
base_branch: &str,
head_branch: &str,
head_sha: &str,
title: &str,
draft: bool,
) -> Self {
@ -779,6 +782,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(),
title: title.to_string(),
draft,
}

View file

@ -37,8 +37,8 @@ use crate::event::{
use crate::handler::HandlerRegistry;
use crate::outcome::{Outcome, StageOutcome};
use crate::pipeline::{
self, FinalizeOptions, Finalized, InitOptions, LlmSpec, Persisted, PullRequestOptions,
ResumeState, SandboxEnvSpec, build_conclusion_from_store, classify_engine_result,
self, FinalizeOptions, Finalized, InitOptions, LlmSpec, Persisted, PublishOptions, ResumeState,
SandboxEnvSpec, build_conclusion_from_store, classify_engine_result,
};
#[cfg(test)]
use crate::records::Checkpoint;
@ -794,7 +794,7 @@ fn runtime_setup_commands(
}
impl RunSession {
/// Shared engine: initialize, execute, finalize, pull_request.
/// Shared engine: initialize, execute, conclude, publish, finalize.
async fn run(
self,
persisted: Persisted,
@ -921,14 +921,14 @@ impl RunSession {
.expect("last_git_sha mutex should not be poisoned: no code panics while holding this lock")
.clone(),
};
let pr_opts = PullRequestOptions {
let publish_opts = PublishOptions {
pr_config: self.pr_config,
github_app: self.pr_github_app,
origin_url: self.pr_origin_url,
model: self.pr_model,
};
let concluded = match Box::pin(pipeline::finalize(executed, &finalize_opts)).await {
let concluded = match Box::pin(pipeline::conclude(executed, &finalize_opts)).await {
Ok(concluded) => concluded,
Err(err) => {
self.steering_hub.drain_pending_at_run_end();
@ -936,7 +936,15 @@ impl RunSession {
return Err(err);
}
};
let finalized = Box::pin(pipeline::pull_request(concluded, &pr_opts)).await;
let published = Box::pin(pipeline::publish(concluded, &publish_opts)).await;
let finalized = match Box::pin(pipeline::finalize(published, &finalize_opts)).await {
Ok(finalized) => finalized,
Err(err) => {
self.steering_hub.drain_pending_at_run_end();
store_progress_logger.flush().await;
return Err(err);
}
};
// Emit `agent.steer.dropped { reason: run_ended }` for any
// unconsumed pending steers on the success path, then flush. The
// scopeguard above re-runs as a no-op (drain is idempotent on an

View file

@ -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};
use super::types::{Concluded, Executed, FinalizeOptions, Finalized, 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};
@ -56,6 +56,15 @@ pub fn classify_engine_result(
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,
@ -483,6 +492,9 @@ pub(crate) fn build_terminal_event(
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),
Ok(outcome) => {
if let Some(failure) = outcome.failure.as_ref() {
@ -521,16 +533,13 @@ async fn stop_sandbox_on_terminal(
Ok(())
}
/// FINALIZE phase: build conclusion, write the meta branch, emit the terminal
/// `WorkflowRunCompleted`/`WorkflowRunFailed` event.
///
/// The terminal event is emitted here (not from `on_run_end`) so observers
/// can't act on "done" before the meta branch writes are flushed.
/// CONCLUDE phase: collect the execution result, final commit, and diff.
///
/// # Errors
///
/// Returns `Error` if persisting terminal state fails.
pub async fn finalize(executed: Executed, options: &FinalizeOptions) -> Result<Concluded, Error> {
/// Returns `Error` if the run state needed to build the conclusion cannot be
/// collected.
pub async fn conclude(executed: Executed, options: &FinalizeOptions) -> Result<Concluded, Error> {
let Executed {
graph,
outcome,
@ -561,20 +570,78 @@ pub async fn finalize(executed: Executed, options: &FinalizeOptions) -> Result<C
let checkpoint = projection
.as_ref()
.and_then(|state| state.current_checkpoint());
let conclusion = build_conclusion_from_parts(
let final_git_commit_sha = options.last_git_sha.clone().or_else(|| {
run_options
.git
.as_ref()
.and_then(|git| git.base_sha.clone())
});
let mut conclusion = build_conclusion_from_parts(
checkpoint,
&projection_billing,
&projection_order,
final_status,
failure_reason,
wall_time_ms,
options.last_git_sha.clone(),
final_git_commit_sha,
);
let ((final_patch, diff_summary), ()) = tokio::join!(
compute_final_patch(&run_options, &services, final_status),
write_finalize_commit(&run_options, &services, &conclusion),
);
let (final_patch, diff_summary) =
compute_final_patch(&run_options, &services, final_status).await;
conclusion.diff = fabro_types::RunDiff {
patch: final_patch,
summary: diff_summary,
};
Ok(Concluded {
outcome,
conclusion,
artifact_count,
graph,
run_options,
services,
})
}
/// FINALIZE phase: persist the final conclusion, emit the terminal event, and
/// clean up the sandbox.
///
/// This runs after PUBLISH so a required push or pull-request failure becomes
/// the terminal run result.
///
/// # Errors
///
/// Returns `Error` if persisting terminal state fails.
pub async fn finalize(published: Published, options: &FinalizeOptions) -> Result<Finalized, Error> {
let Published {
execution_outcome,
publish_outcome,
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 (final_status, failure, _run_status) = classify_engine_result(&outcome);
conclusion.status = final_status;
conclusion.failure = failure;
write_finalize_commit(&run_options, &services, &conclusion).await;
if services.metadata_runtime.metadata_degraded() {
services.emitter.notice(
@ -588,9 +655,9 @@ pub async fn finalize(executed: Executed, options: &FinalizeOptions) -> Result<C
&outcome,
conclusion.timing,
artifact_count,
options.last_git_sha.clone(),
final_patch,
diff_summary,
conclusion.final_git_commit_sha.clone(),
conclusion.diff.patch.clone(),
conclusion.diff.summary,
conclusion.billing.clone(),
);
services.emitter.emit(&terminal_event);
@ -626,12 +693,12 @@ pub async fn finalize(executed: Executed, options: &FinalizeOptions) -> Result<C
);
}
Ok(Concluded {
Ok(Finalized {
run_id: run_options.run_id,
outcome,
conclusion,
graph,
run_options,
services,
pushed_branch,
pr_url,
})
}
@ -716,6 +783,21 @@ mod tests {
}
}
async fn finalize_executed(
executed: Executed,
options: &FinalizeOptions,
) -> Result<Finalized, Error> {
let concluded = conclude(executed, options).await?;
let published = crate::pipeline::publish(concluded, &crate::pipeline::PublishOptions {
pr_config: None,
github_app: None,
origin_url: None,
model: "test-model".to_string(),
})
.await;
finalize(published, options).await
}
fn test_store() -> Arc<Database> {
Arc::new(Database::new(
Arc::new(InMemory::new()),
@ -869,6 +951,26 @@ mod tests {
use crate::test_support::test_usage;
#[test]
fn publish_error_builds_publish_failed_terminal_event() {
let event = build_terminal_event(
&Err(Error::publish("GitHub rejected pull request creation")),
fabro_types::RunTiming::wall_only(10),
0,
Some("final-sha".to_string()),
Some("diff".to_string()),
None,
None,
);
match event {
Event::WorkflowRunFailed { failure, .. } => {
assert_eq!(failure.reason, FailureReason::PublishFailed);
}
other => panic!("expected run failure, got {other:?}"),
}
}
#[test]
fn conclusion_stage_order_follows_projection_first_event_order() {
let mut projection = test_projection();
@ -1068,7 +1170,7 @@ mod tests {
services,
);
let concluded = finalize(executed, &FinalizeOptions {
let concluded = finalize_executed(executed, &FinalizeOptions {
run_dir: run_dir.clone(),
run_id: test_run_id(),
workflow_name: "test".to_string(),
@ -1249,7 +1351,7 @@ mod tests {
services,
);
finalize(executed, &FinalizeOptions {
finalize_executed(executed, &FinalizeOptions {
run_dir: repo_dir.path().to_path_buf(),
run_id: test_run_id(),
workflow_name: "test".to_string(),
@ -1273,6 +1375,196 @@ mod tests {
]);
}
#[tokio::test]
async fn configured_run_branch_without_remote_is_not_reported_as_pushed() {
let repo_dir = tempfile::tempdir().unwrap();
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(MockSandbox::linux()),
Arc::new(RunMetadataRuntime::new()),
None,
);
let mut run_options = test_run_options(repo_dir.path());
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 concluded = conclude(executed, &options).await.unwrap();
let published = crate::pipeline::publish(concluded, &crate::pipeline::PublishOptions {
pr_config: None,
github_app: None,
origin_url: None,
model: "test-model".to_string(),
})
.await;
assert!(matches!(
&published.publish_outcome,
Ok(crate::pipeline::PublishOutcome::NotRequested)
));
let finalized = finalize(published, &options).await.unwrap();
assert!(finalized.outcome.is_ok());
assert_eq!(finalized.pushed_branch, None);
let events = events.lock().unwrap();
let names = events.iter().map(RunEvent::event_name).collect::<Vec<_>>();
assert_eq!(names, vec!["run.completed"]);
}
#[tokio::test]
async fn final_push_failure_becomes_terminal_publish_failure() {
let repo_dir = tempfile::tempdir().unwrap();
let sandbox = Arc::new(MockSandbox::linux());
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,
sandbox,
Arc::new(RunMetadataRuntime::new()),
None,
);
let mut run_options = test_run_options(repo_dir.path());
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 concluded = conclude(executed, &options).await.unwrap();
let published = crate::pipeline::publish(concluded, &crate::pipeline::PublishOptions {
pr_config: None,
github_app: None,
origin_url: Some("https://github.com/owner/repo.git".to_string()),
model: "test-model".to_string(),
})
.await;
assert!(matches!(
&published.publish_outcome,
Err(Error::Publish { .. })
));
let finalized = finalize(published, &options).await.unwrap();
assert!(matches!(finalized.outcome, Err(Error::Publish { .. })));
assert_eq!(
finalized
.conclusion
.failure
.as_ref()
.map(|failure| failure.reason),
Some(FailureReason::PublishFailed)
);
let events = events.lock().unwrap();
let names = events.iter().map(RunEvent::event_name).collect::<Vec<_>>();
assert_eq!(names, vec!["git.push", "run.failed"]);
match &events.last().unwrap().body {
EventBody::RunFailed(props) => {
assert_eq!(props.failure.reason, FailureReason::PublishFailed);
}
other => panic!("expected run.failed, got {other:?}"),
}
}
#[tokio::test]
async fn pull_request_failure_precedes_terminal_publish_failure() {
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 =
Some("diff --git a/a b/a\n+published change\n".to_string());
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;
let finalized = finalize(published, &options).await.unwrap();
assert!(matches!(finalized.outcome, Err(Error::Publish { .. })));
let events = events.lock().unwrap();
let names = events.iter().map(RunEvent::event_name).collect::<Vec<_>>();
assert_eq!(names, vec!["git.push", "pull_request.failed", "run.failed"]);
match &events.last().unwrap().body {
EventBody::RunFailed(props) => {
assert_eq!(props.failure.reason, FailureReason::PublishFailed);
}
other => panic!("expected run.failed, got {other:?}"),
}
}
#[tokio::test]
async fn finalize_stops_sandbox_on_terminal_without_deleting() {
let repo_dir = tempfile::tempdir().unwrap();
@ -1292,7 +1584,7 @@ mod tests {
services,
);
finalize(executed, &FinalizeOptions {
finalize_executed(executed, &FinalizeOptions {
run_dir: repo_dir.path().to_path_buf(),
run_id: test_run_id(),
workflow_name: "test".to_string(),
@ -1326,7 +1618,7 @@ mod tests {
services,
);
finalize(executed, &FinalizeOptions {
finalize_executed(executed, &FinalizeOptions {
run_dir: repo_dir.path().to_path_buf(),
run_id: test_run_id(),
workflow_name: "test".to_string(),
@ -1379,7 +1671,7 @@ mod tests {
services,
);
finalize(executed, &FinalizeOptions {
finalize_executed(executed, &FinalizeOptions {
run_dir: repo.to_path_buf(),
run_id: test_run_id(),
workflow_name: "test".to_string(),

View file

@ -3,6 +3,7 @@ mod finalize;
mod initialize;
mod parse;
mod persist;
mod publish;
mod pull_request;
mod transform;
pub(crate) mod types;
@ -12,18 +13,19 @@ pub use execute::execute;
pub(crate) use finalize::build_conclusion_from_store;
#[cfg(any(test, feature = "test-support"))]
pub(crate) use finalize::{billing_from_projection, build_terminal_event};
pub use finalize::{classify_engine_result, finalize, write_finalize_commit};
pub use finalize::{classify_engine_result, conclude, finalize, write_finalize_commit};
pub use initialize::initialize;
pub use parse::parse;
pub(crate) use persist::persist;
pub use publish::publish;
pub use pull_request::{
AutoMergeOptions, CreatedPullRequest, OpenPullRequestRequest, PrContent, build_pr_content,
maybe_open_pull_request, pull_request,
maybe_open_pull_request,
};
pub use transform::transform;
pub use types::{
Concluded, Executed, FinalizeOptions, Finalized, InitOptions, Initialized, LlmSpec, Parsed,
Persisted, PullRequestOptions, ResumeState, SandboxEnvSpec, TEMPLATE_UNDEFINED_VARIABLE_RULE,
TransformOptions, Transformed, Validated,
Persisted, PublishOptions, PublishOutcome, Published, ResumeState, SandboxEnvSpec,
TEMPLATE_UNDEFINED_VARIABLE_RULE, TransformOptions, Transformed, Validated,
};
pub use validate::validate;

View file

@ -0,0 +1,201 @@
use std::sync::Arc;
use super::pull_request::{AutoMergeOptions, OpenPullRequestRequest, maybe_open_pull_request};
use super::types::{Concluded, PublishOptions, PublishOutcome, Published};
use crate::error::Error;
use crate::event::Event;
use crate::outcome::StageOutcome;
/// PUBLISH phase: push the final run commit and, when configured, open a pull
/// request.
///
/// 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 Concluded {
outcome,
conclusion,
artifact_count,
graph: _,
run_options,
services,
} = concluded;
Published {
execution_outcome: outcome,
publish_outcome,
conclusion,
artifact_count,
run_options,
services,
}
}
async fn publish_inner(
concluded: &Concluded,
options: &PublishOptions,
) -> Result<PublishOutcome, Error> {
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);
}
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 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 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(),
});
}
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 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 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",
)
})?;
concluded
.services
.emitter
.emit(&Event::pull_request_created(
&created.link,
&created.base_branch,
&created.head_branch,
&created.head_sha,
&created.title,
pr_config.draft,
));
Ok(PublishOutcome::Published {
pushed_branch: run_branch.to_string(),
pr_url: Some(created.link.html_url()),
})
}
fn pull_request_error(concluded: &Concluded, message: &str) -> Error {
concluded.services.emitter.emit(&Event::PullRequestFailed {
error: message.to_string(),
});
Error::publish(message)
}

View file

@ -13,9 +13,7 @@ use fabro_types::settings::run::MergeStrategy;
use fabro_util::text::strip_goal_decoration;
use tracing::{debug, info, warn};
use super::types::{Concluded, Finalized, PullRequestOptions};
use crate::event::{Event, RunNoticeCode, RunNoticeLevel};
use crate::outcome::{StageOutcome, format_cost as outcome_format_cost};
use crate::outcome::format_cost as outcome_format_cost;
use crate::records::{Conclusion, RunSpec};
use crate::runtime_store::RunStoreHandle;
@ -327,22 +325,6 @@ fn assemble_pr_body(
parts.join("\n")
}
async fn load_pull_request_diff(run_store: &RunStoreHandle) -> String {
run_store
.state()
.await
.inspect_err(|err| {
tracing::warn!(error = %err, "Failed to load final patch from store for PR");
})
.ok()
.and_then(|state| {
state
.conclusion
.and_then(|conclusion| conclusion.diff.patch)
})
.unwrap_or_default()
}
/// Build complete PR content by combining LLM-generated narrative with
/// deterministic fallbacks and programmatic sections.
pub async fn build_pr_content(
@ -461,20 +443,23 @@ pub struct AutoMergeOptions {
/// Inputs for [`maybe_open_pull_request`].
pub struct OpenPullRequestRequest<'a> {
pub github: github_app::GitHubContext<'a>,
pub origin_url: &'a str,
pub base_branch: &'a str,
pub head_branch: &'a str,
pub goal: &'a str,
pub diff: &'a str,
pub model: &'a str,
pub draft: bool,
pub auto_merge: Option<AutoMergeOptions>,
pub run_store: &'a RunStoreHandle,
pub llm_source: &'a dyn CredentialSource,
pub catalog: Arc<Catalog>,
pub conclusion: Option<&'a Conclusion>,
pub run_state: Option<&'a RunProjection>,
pub github: github_app::GitHubContext<'a>,
pub origin_url: &'a str,
pub base_branch: &'a str,
pub head_branch: &'a str,
/// Commit that must be visible at the remote branch before the PR is
/// opened.
pub expected_head_sha: &'a str,
pub goal: &'a str,
pub diff: &'a str,
pub model: &'a str,
pub draft: bool,
pub auto_merge: Option<AutoMergeOptions>,
pub run_store: &'a RunStoreHandle,
pub llm_source: &'a dyn CredentialSource,
pub catalog: Arc<Catalog>,
pub conclusion: Option<&'a Conclusion>,
pub run_state: Option<&'a RunProjection>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
@ -483,6 +468,7 @@ 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.
@ -516,6 +502,25 @@ pub async fn maybe_open_pull_request(
let body = truncate_pr_body(&content.body);
let title = content.title;
let remote_head = github_app::branch_head_sha(&req.github, &owner, &repo, req.head_branch)
.await
.map_err(|err| format!("failed to verify remote branch head: {err:#}"))?;
match remote_head {
Some(remote_head) if remote_head == req.expected_head_sha => {}
Some(remote_head) => {
return Err(format!(
"remote branch '{}' points to commit {remote_head}, expected final commit {}",
req.head_branch, req.expected_head_sha
));
}
None => {
return Err(format!(
"remote branch '{}' does not exist; expected final commit {}",
req.head_branch, req.expected_head_sha
));
}
}
let created = github_app::create_pull_request(
&req.github,
&owner,
@ -565,111 +570,10 @@ pub async fn maybe_open_pull_request(
title,
base_branch: req.base_branch.to_string(),
head_branch: req.head_branch.to_string(),
head_sha: req.expected_head_sha.to_string(),
}))
}
/// PULL_REQUEST phase: optionally create a pull request after finalize.
///
/// This stage is infallible: failures are emitted and logged, but the pipeline
/// completes.
pub async fn pull_request(concluded: Concluded, options: &PullRequestOptions) -> Finalized {
let Concluded {
outcome,
conclusion,
graph,
run_options,
services,
} = concluded;
let mut pr_url = None;
if let Some(pr_cfg) = &options.pr_config {
if run_options.dry_run_enabled() {
tracing::debug!("Skipping PR creation: run is in dry-run mode");
} else if let Err(ref e) = outcome {
tracing::debug!(error = %e, "Skipping PR creation: engine returned an error");
} else if let Ok(ref result) = outcome {
if matches!(
result.status,
StageOutcome::Succeeded | StageOutcome::PartiallySucceeded
) {
let diff = load_pull_request_diff(&services.run_store).await;
if let (Some(base_branch), Some(run_branch), Some(creds), Some(origin)) = (
&run_options.base_branch,
run_options.run_branch(),
&options.github_app,
&options.origin_url,
) {
let auto_merge = if pr_cfg.auto_merge {
Some(AutoMergeOptions {
merge_strategy: pr_cfg.merge_strategy,
})
} else {
None
};
match maybe_open_pull_request(OpenPullRequestRequest {
github: github_app::GitHubContext::new(
creds,
&github_app::github_api_base_url(),
),
origin_url: origin,
base_branch,
head_branch: run_branch,
goal: graph.goal(),
diff: &diff,
model: &options.model,
draft: pr_cfg.draft,
auto_merge,
run_store: &services.run_store,
llm_source: services.llm_source.as_ref(),
catalog: Arc::clone(&services.catalog),
conclusion: Some(&conclusion),
run_state: None,
})
.await
{
Ok(Some(created)) => {
services.emitter.emit(&Event::pull_request_created(
&created.link,
&created.base_branch,
&created.head_branch,
&created.title,
pr_cfg.draft,
));
pr_url = Some(created.link.html_url());
}
Ok(None) => {}
Err(e) => {
services
.emitter
.emit(&Event::PullRequestFailed { error: e.clone() });
services.emitter.notice(
RunNoticeLevel::Warn,
RunNoticeCode::PullRequestFailed,
format!("PR creation failed: {e}"),
);
}
}
}
}
}
}
Finalized {
run_id: run_options.run_id,
outcome,
conclusion,
pushed_branch: run_options
.settings
.run
.run_branch
.push
.then(|| run_options.run_branch().map(str::to_string))
.flatten(),
pr_url,
}
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
@ -691,18 +595,14 @@ mod tests {
};
use fabro_vault::{SecretType, Vault};
use futures::stream;
use httpmock::Method::POST;
use httpmock::Method::{GET, POST};
use httpmock::MockServer;
use object_store::memory::InMemory;
use tokio::sync::RwLock as AsyncRwLock;
use tokio_util::sync::CancellationToken;
use super::*;
use crate::event::{Event, append_event};
use crate::outcome::Outcome;
use crate::records::StageSummary;
use crate::run_options::{GitCheckpointOptions, RunOptions};
use crate::services::EngineServices;
struct MockProvider {
name: String,
@ -917,48 +817,6 @@ mod tests {
}
}
#[tokio::test]
async fn pull_request_omits_pushed_branch_when_run_branch_push_disabled() {
let temp = tempfile::tempdir().unwrap();
let mut settings = WorkflowSettings::default();
settings.run.run_branch.push = false;
let run_options = RunOptions {
settings,
run_dir: temp.path().to_path_buf(),
cancel_token: CancellationToken::new(),
run_id: fixtures::RUN_1,
labels: HashMap::new(),
workflow_slug: None,
github_app: None,
pre_run_git: None,
fork_source_ref: None,
base_branch: None,
display_base_sha: None,
git: Some(GitCheckpointOptions {
base_sha: None,
run_branch: Some("fabro/run/test".to_string()),
meta_branch: None,
}),
};
let concluded = Concluded {
outcome: Ok(Outcome::success()),
conclusion: make_test_conclusion(),
graph: Graph::new("test"),
run_options,
services: EngineServices::test_default().run,
};
let finalized = pull_request(concluded, &PullRequestOptions {
pr_config: None,
github_app: None,
origin_url: None,
model: "test-model".to_string(),
})
.await;
assert_eq!(finalized.pushed_branch, None);
}
// ── format_arc_details_section tests ────────────────────────────────
#[test]
@ -1555,20 +1413,21 @@ mod tests {
});
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",
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,
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());
@ -1576,79 +1435,41 @@ mod tests {
}
#[tokio::test]
async fn load_pull_request_diff_uses_store_without_disk_patch() {
let tmp = tempfile::tempdir().unwrap();
let store = test_store();
let run_store = store.create_run(&fixtures::RUN_1).await.unwrap();
let run_spec = RunSpec {
run_id: fixtures::RUN_1,
settings: fabro_types::WorkflowSettings::default(),
graph: Graph::new("test"),
graph_source: None,
workflow_slug: None,
automation: None,
source_directory: Some(tmp.path().display().to_string()),
git: None,
labels: std::collections::HashMap::new(),
provenance: test_support::test_run_provenance(),
manifest_blob: None,
definition_blob: None,
fork_source_ref: None,
};
append_event(&run_store, &fixtures::RUN_1, &Event::RunCreated {
run_id: fixtures::RUN_1,
title: None,
settings: serde_json::to_value(&run_spec.settings).unwrap(),
graph: serde_json::to_value(&run_spec.graph).unwrap(),
workflow_source: None,
workflow_config: None,
labels: run_spec.labels.clone().into_iter().collect(),
run_dir: tmp.path().display().to_string(),
source_directory: run_spec.source_directory.clone(),
workflow_slug: None,
automation: None,
db_prefix: None,
provenance: run_spec.provenance.clone(),
manifest_blob: None,
git: None,
fork_source_ref: None,
retried_from: None,
parent_id: None,
web_url: None,
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 {
github: fabro_github::GitHubContext::new(&harness.creds, &github_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: "diff --git a/src/lib.rs b/src/lib.rs\n+fn x() {}\n",
model: "claude-sonnet-4-20250514",
draft: false,
auto_merge: None,
run_store: &harness.run_store,
llm_source: harness.llm_source.as_ref(),
catalog: harness.catalog.clone(),
conclusion: None,
run_state: None,
})
.await
.unwrap();
append_event(&run_store, &fixtures::RUN_1, &Event::RunRunnable {
source: fabro_types::RunRunnableSource::StartRequested,
actor: None,
})
.await
.unwrap();
append_event(&run_store, &fixtures::RUN_1, &Event::RunStarting)
.await
.unwrap();
append_event(&run_store, &fixtures::RUN_1, &Event::RunRunning)
.await
.unwrap();
append_event(&run_store, &fixtures::RUN_1, &Event::WorkflowRunCompleted {
timing: fabro_types::RunTiming::wall_only(1),
artifact_count: 0,
status: "succeeded".to_string(),
reason: SuccessReason::Completed,
total_usd_micros: None,
final_git_commit_sha: None,
final_patch: Some(
"diff --git a/src/lib.rs b/src/lib.rs\n+fn from_store() {}\n".to_string(),
),
diff_summary: None,
billing: None,
})
.await
.unwrap();
.expect_err("stale remote branch must prevent PR creation");
let diff = load_pull_request_diff(&run_store.clone().into()).await;
assert!(diff.contains("from_store"));
assert!(error.contains("stale-sha"));
assert!(error.contains("final-sha"));
httpmock::Mock::new(harness.openai_mock_id, &harness.openai_server)
.assert_async()
.await;
httpmock::Mock::new(harness.branch_mock_id, &harness.github_server)
.assert_async()
.await;
httpmock::Mock::new(harness.github_mock_id, &harness.github_server)
.assert_calls_async(0)
.await;
}
// ── Structured-output PR content tests ──────────────────────────────
@ -1807,6 +1628,7 @@ mod tests {
openai_server: MockServer,
github_server: MockServer,
openai_mock_id: usize,
branch_mock_id: usize,
github_mock_id: usize,
llm_source: Arc<dyn CredentialSource>,
catalog: Arc<Catalog>,
@ -1819,6 +1641,9 @@ mod tests {
httpmock::Mock::new(self.openai_mock_id, &self.openai_server)
.assert_async()
.await;
httpmock::Mock::new(self.branch_mock_id, &self.github_server)
.assert_async()
.await;
httpmock::Mock::new(self.github_mock_id, &self.github_server)
.assert_async()
.await;
@ -1830,6 +1655,13 @@ mod tests {
/// credential source, and a run store seeded with a non-empty
/// `final_patch`.
async fn setup_fallback_test_harness(openai_payload_text: &str) -> FallbackHarness {
setup_fallback_test_harness_with_branch_sha(openai_payload_text, "final-sha").await
}
async fn setup_fallback_test_harness_with_branch_sha(
openai_payload_text: &str,
branch_sha: &str,
) -> FallbackHarness {
let openai_server = MockServer::start_async().await;
let openai_mock = openai_server
.mock_async(|when, then| {
@ -1843,6 +1675,19 @@ mod tests {
.await;
let github_server = MockServer::start_async().await;
let branch_sha = branch_sha.to_string();
let branch_mock = github_server
.mock_async(move |when, then| {
when.method(GET)
.path("/repos/owner/repo/branches/fabro/run/123")
.header("authorization", "Bearer test-token");
then.status(200)
.header("content-type", "application/json")
.json_body(serde_json::json!({
"commit": { "sha": branch_sha }
}));
})
.await;
let github_mock = github_server
.mock_async(|when, then| {
when.method(POST)
@ -1878,8 +1723,7 @@ mod tests {
let store = test_store();
let run_store = store.create_run(&fixtures::RUN_1).await.unwrap();
// Seed a non-empty `final_patch` so `load_pull_request_diff` returns
// diff content and the early-return for empty diffs does not fire.
// Seed a completed run so the PR body can include run details.
let run_spec = RunSpec {
run_id: fixtures::RUN_1,
settings: fabro_types::WorkflowSettings::default(),
@ -1947,6 +1791,7 @@ mod tests {
.unwrap();
let openai_mock_id = openai_mock.id;
let branch_mock_id = branch_mock.id;
let github_mock_id = github_mock.id;
FallbackHarness {
@ -1954,6 +1799,7 @@ mod tests {
openai_server,
github_server,
openai_mock_id,
branch_mock_id,
github_mock_id,
llm_source,
catalog,
@ -1978,6 +1824,7 @@ mod tests {
origin_url: "https://github.com/owner/repo.git",
base_branch: "main",
head_branch: "fabro/run/123",
expected_head_sha: "final-sha",
goal: "Fix telemetry leak\n\ndetails...",
diff: "diff --git a/src/lib.rs b/src/lib.rs\n+fn x() {}\n",
model: "gpt-5.4",
@ -2015,6 +1862,7 @@ mod tests {
origin_url: "https://github.com/owner/repo.git",
base_branch: "main",
head_branch: "fabro/run/123",
expected_head_sha: "final-sha",
goal: &goal,
diff: "diff --git a/src/lib.rs b/src/lib.rs\n+fn x() {}\n",
model: "gpt-5.4",

View file

@ -337,17 +337,59 @@ pub struct Executed {
pub model: String,
}
/// Output of the FINALIZE phase.
/// Output of the CONCLUDE phase.
#[non_exhaustive]
pub struct Concluded {
pub outcome: Result<Outcome, Error>,
pub conclusion: Conclusion,
pub graph: Graph,
pub run_options: RunOptions,
pub services: Arc<RunServices>,
pub outcome: Result<Outcome, Error>,
pub conclusion: Conclusion,
pub artifact_count: usize,
pub graph: Graph,
pub run_options: RunOptions,
pub services: Arc<RunServices>,
}
/// Output of the PULL_REQUEST phase.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum PublishOutcome {
NotRequested,
NoChanges {
pushed_branch: String,
},
Published {
pushed_branch: String,
pr_url: Option<String>,
},
}
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,
}
}
}
/// Output of the PUBLISH phase.
#[non_exhaustive]
pub struct Published {
pub execution_outcome: Result<Outcome, Error>,
pub publish_outcome: Result<PublishOutcome, Error>,
pub conclusion: Conclusion,
pub artifact_count: usize,
pub run_options: RunOptions,
pub services: Arc<RunServices>,
}
/// Output of the FINALIZE phase.
#[non_exhaustive]
pub struct Finalized {
pub run_id: RunId,
@ -383,8 +425,8 @@ pub struct FinalizeOptions {
pub last_git_sha: Option<String>,
}
/// Options for the PULL_REQUEST phase.
pub struct PullRequestOptions {
/// Options for the PUBLISH phase.
pub struct PublishOptions {
pub pr_config: Option<PullRequestSettings>,
pub github_app: Option<fabro_github::GitHubCredentials>,
pub origin_url: Option<String>,

View file

@ -113,6 +113,7 @@ fn success_reason_json_tokens_match_openapi() {
#[test]
fn failure_reason_json_tokens_match_openapi() {
assert_string_json(FailureReason::WorkflowError, "workflow_error");
assert_string_json(FailureReason::PublishFailed, "publish_failed");
assert_string_json(FailureReason::Cancelled, "cancelled");
assert_string_json(FailureReason::ApprovalDenied, "approval_denied");
assert_string_json(FailureReason::Terminated, "terminated");

View file

@ -247,6 +247,8 @@ pub struct PullRequestCreatedProps {
pub repo: String,
pub base_branch: String,
pub head_branch: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub head_sha: Option<String>,
pub title: String,
pub draft: bool,
}

View file

@ -290,6 +290,7 @@ pub enum SuccessReason {
#[strum(serialize_all = "snake_case")]
pub enum FailureReason {
WorkflowError,
PublishFailed,
Cancelled,
ApprovalDenied,
Terminated,

View file

@ -20,6 +20,7 @@
export const FailureReason = {
WORKFLOW_ERROR: 'workflow_error',
PUBLISH_FAILED: 'publish_failed',
CANCELLED: 'cancelled',
APPROVAL_DENIED: 'approval_denied',
TERMINATED: 'terminated',