Publish a successful run from its run_finished hook

Pushing the run branch and opening the pull request now happen in the
run's worker, in Fabro's run_finished hook, after the last stage and
before the run's terminal record, as the legacy publish step did. A
failed push or pull request fails the run with publish_failed instead of
leaving a warning on a run that already succeeded.

fabro-petri gains a RunPublisher the hooks call for a successful run with
its run branch, final commit, snapshot repository and patch; the worker's
GitHub publisher pushes from the snapshot repository with a push token it
mints at that moment, opens the pull request its settings ask for, and
records it. The worker resolves the server's GitHub credentials itself
for both the read-only checkout token and the push token, so the server
no longer hands it a clone credential or publishes after the run.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-09-30 13:46:42 -04:00
parent 5c6195a352
commit 0250b40586
19 changed files with 711 additions and 674 deletions

View file

@ -33,13 +33,14 @@ macOS note: if `cargo nextest run` fails with `Too many open files (os error 24)
- A GitHub target's workspace is checked out by Fabro's hooks, not by the
sandbox driver or Petri's `start` checkout: when a fresh run's scope is
acquired, `fabro-petri`'s `RunWorkspaces::check_out_source` fetches the
target's revision inside the sandbox with a read-only token the server
resolves at each worker launch (`FABRO_RUN_GIT_CREDENTIAL`, scrubbed at
worker startup), and seeds the workspace's snapshot repository with that
commit so checkpoint bundles from a shallow clone import. When the run
ends, the server pushes the final checkpoint to `fabro/run/<id>` from the
snapshot repository and requests the pull request
(`fabro-server/src/server/run_publication.rs`). `CloneRequest` still
target's revision inside the sandbox with a read-only token the worker
mints from the server's GitHub credentials, and seeds the workspace's
snapshot repository with that commit so checkpoint bundles from a shallow
clone import. In the `run_finished` hook, before the terminal record, a
successful run's worker pushes the final checkpoint to `fabro/run/<id>`
from the snapshot repository and opens the pull request its settings ask
for (`fabro-cli/src/commands/run/publish.rs`); a failure fails the run
with `publish_failed`. `CloneRequest` still
travels beside the sandbox spec so the run record names the origin and
branch; the sandbox layer refuses a request that asks it to clone.
Preflight and `fabro exec` initialize sandboxes with `CloneRequest::none()`,
@ -152,13 +153,14 @@ Fabro is an AI-powered workflow orchestration platform. Workflows are defined as
- A GitHub target's workspace is checked out by Fabro's hooks, not by the
sandbox driver or Petri's `start` checkout: when a fresh run's scope is
acquired, `fabro-petri`'s `RunWorkspaces::check_out_source` fetches the
target's revision inside the sandbox with a read-only token the server
resolves at each worker launch (`FABRO_RUN_GIT_CREDENTIAL`, scrubbed at
worker startup), and seeds the workspace's snapshot repository with that
commit so checkpoint bundles from a shallow clone import. When the run
ends, the server pushes the final checkpoint to `fabro/run/<id>` from the
snapshot repository and requests the pull request
(`fabro-server/src/server/run_publication.rs`). `CloneRequest` still
target's revision inside the sandbox with a read-only token the worker
mints from the server's GitHub credentials, and seeds the workspace's
snapshot repository with that commit so checkpoint bundles from a shallow
clone import. In the `run_finished` hook, before the terminal record, a
successful run's worker pushes the final checkpoint to `fabro/run/<id>`
from the snapshot repository and opens the pull request its settings ask
for (`fabro-cli/src/commands/run/publish.rs`); a failure fails the run
with `publish_failed`. `CloneRequest` still
travels beside the sandbox spec so the run record names the origin and
branch; the sandbox layer refuses a request that asks it to clone.
Preflight and `fabro exec` initialize sandboxes with `CloneRequest::none()`,

View file

@ -86,7 +86,7 @@ This means your original working directory stays untouched while the agent makes
If the working directory has uncommitted changes, the worktree starts from committed `HEAD` and those uncommitted changes are not included. Fabro logs a warning so you can commit, stash, or run explicitly in place when that is what you want.
</Note>
For Docker and Daytona sandboxes, a GitHub target is checked out inside the sandbox and checkpoint Git operations run there; each checkpoint commit reaches the server as a Git bundle. When a successful run ends, the server pushes the run branch to origin when pushing is configured.
For Docker and Daytona sandboxes, a GitHub target is checked out inside the sandbox and checkpoint Git operations run there; each checkpoint commit reaches the server as a Git bundle. Before a successful run finishes, Fabro pushes the run branch to origin when pushing is configured.
## Resuming a run
@ -123,7 +123,7 @@ After a node completes, Fabro:
3. Collects the code diff.
4. Emits a checkpoint event with execution state and the code commit SHA. The run store persists this event and updates the projection.
A checkpoint commit failure stops execution. A diff failure emits a warning notice. When the run ends, the server pushes the run branch; a failed push is a warning notice on the run.
A checkpoint commit failure stops execution. A diff failure emits a warning notice. A failed final push or pull request marks a successful run as failed with `publish_failed`.
## Inspecting run history

View file

@ -29,7 +29,7 @@ The rest of this page describes the `app` strategy, which is required for browse
|---|---|
| **OAuth login** | Users sign in to the web UI with their GitHub account |
| **Private repo cloning** | Daytona and Docker sandboxes fetch private repositories using short-lived, read-only Installation Access Tokens |
| **Run branch pushing** | When a run ends, the Fabro server pushes the run branch to origin |
| **Run branch pushing** | Before a successful run finishes, Fabro pushes the run branch to origin |
| **Auto-PR** | When `[run.pull_request] enabled = true` in the [run config](/execution/run-configuration#runpull_request), Fabro opens a PR from the agent's working branch after a successful run |
| **Auto-merge** | When `[run.pull_request] auto_merge = true`, Fabro enables GitHub's auto-merge on created PRs so they merge automatically once required checks pass |
| **Sandbox GITHUB_TOKEN** | When `[run.integrations.github.permissions]` are declared at any layer (workflow, project, or user settings), Fabro mints a scoped Installation Access Token and injects it as `GITHUB_TOKEN` in the sandbox |
@ -218,9 +218,9 @@ An empty `allowed_usernames` list rejects all users.
When a run targets a GitHub repository, its workspace is checked out inside the sandbox (Docker or Daytona) before the first stage runs:
1. When the server launches the run's worker, it signs a short-lived JWT using the App ID and private key (RS256, 10-minute validity)
1. When the run's worker starts, it signs a short-lived JWT using the App ID and private key (RS256, 10-minute validity)
2. Using the JWT, Fabro looks up the GitHub App installation for the repository (`GET /repos/\{owner\}/\{repo\}/installation`)
3. Fabro requests a scoped Installation Access Token with `contents: read` permission on the specific repository and hands it to the worker, which removes it from its environment at startup
3. Fabro requests a scoped Installation Access Token with `contents: read` permission on the specific repository
4. Inside the sandbox, Fabro fetches the selected revision from `https://github.com/<owner>/<repo>` at the run's `[run.clone] depth`, presenting the token as an HTTP header on that one command, and checks out the working branch
5. The workspace's `origin` is the plain HTTPS URL: the token is never written into the repository, its configuration, or its remote
@ -318,9 +318,9 @@ Every commit a run creates is authored and committed by the run's GitHub credent
### Checkpoint pushing
After each workflow stage, Fabro [checkpoints](/execution/checkpoints) the workspace on the run branch, `fabro/run/<run-id>`, and moves the commit to the server. When a successful run ends, the server pushes the run's final commit to that branch on origin with its own credentials (an Installation Access Token with `contents: write`, minted for the push). The sandbox never holds a credential that can push. `[run.run_branch] push = false` keeps the branch on the server.
After each workflow stage, Fabro [checkpoints](/execution/checkpoints) the workspace on the run branch, `fabro/run/<run-id>`, and moves the commit to the server. After a successful run's last stage and before the run finishes, Fabro pushes the final commit to that branch on origin with an Installation Access Token with `contents: write`, minted for the push, so a long run never pushes with an expired token. The sandbox never holds a credential that can push. `[run.run_branch] push = false` keeps the branch on the server.
When pull request creation is enabled and the run changed files, the server then requests the pull request, and Fabro checks that GitHub reports the run branch at the exact final commit before opening it. A failed push or pull request request is recorded on the run as a warning notice; the run's own outcome is unchanged.
When pull request creation is enabled and the run changed files, Fabro then checks that GitHub reports the run branch at the exact final commit and opens the pull request. A failed 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.
## Troubleshooting

View file

@ -23,6 +23,7 @@ pub(crate) mod overrides;
pub(crate) mod petri_stream;
mod petri_worker;
pub(crate) mod preview;
mod publish;
mod remote_workflow;
mod resolution;
pub(crate) mod resume;
@ -39,20 +40,10 @@ pub(crate) mod test_support;
pub(crate) mod timeline;
pub(crate) mod wait;
/// The credentials the server hands a run worker through its environment,
/// captured and scrubbed from the process before anything is spawned.
#[derive(Default)]
pub(crate) struct WorkerSecrets {
/// The worker's bearer for the server's API.
pub(crate) token: Option<String>,
/// The read-only credential the run's GitHub target is fetched with.
pub(crate) git_credential: Option<String>,
}
pub(crate) async fn dispatch(
cmd: RunCommands,
base_ctx: &CommandContext,
worker_secrets: WorkerSecrets,
worker_token: Option<String>,
) -> Result<()> {
let printer = base_ctx.printer();
@ -116,8 +107,7 @@ pub(crate) async fn dispatch(
mode,
fabro_home,
}) => {
let worker_token = worker_secrets
.token
let worker_token = worker_token
.filter(|token| !token.trim().is_empty())
.ok_or_else(|| {
anyhow!("FABRO_WORKER_TOKEN is required for worker subprocess auth")
@ -132,7 +122,6 @@ pub(crate) async fn dispatch(
mode,
fabro_home,
&worker_token,
worker_secrets.git_credential,
)
.instrument(run_span),
)

View file

@ -75,14 +75,14 @@ use fabro_petri::artifacts::ClientArtifactWriter;
use fabro_petri::blobs::ClientBlobs;
use fabro_petri::controls::{RunControls, SteerError};
use fabro_petri::engine::{self, Conclusion, Execution, RunRequest};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::hooks::{HooksSpec, RunPublisher};
use fabro_petri::interview::{Approval, FabroInterviewer};
use fabro_petri::petri::OwnerId;
use fabro_petri::platform_records::{HttpPlatformRecords, PlatformRecords};
use fabro_petri::providers::{DaytonaCredentials, SandboxProviderConfig};
use fabro_petri::runtime::{self, RuntimeSpec};
use fabro_petri::secrets::VaultSecrets;
use fabro_petri::source::{RunSource, SourceCredential};
use fabro_petri::source::RunSource;
use fabro_petri::{HttpRunStore, admission};
use fabro_static::EnvVars;
use fabro_store::RunProjection;
@ -99,26 +99,24 @@ use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use tracing::{info, warn};
use super::publish;
use super::runner::{self, WorkerTitlePhase};
use crate::args::RunWorkerMode;
use crate::command_context;
/// What the worker holds when it hands a run to Petri.
pub(super) struct PetriWorker<'a> {
pub(super) run_id: RunId,
pub(super) target: ServerTarget,
pub(super) client: Client,
pub(super) run_state: RunProjection,
pub(super) storage_dir: &'a Path,
pub(super) run_dir: PathBuf,
pub(super) mode: RunWorkerMode,
pub(super) run_id: RunId,
pub(super) target: ServerTarget,
pub(super) client: Client,
pub(super) run_state: RunProjection,
pub(super) storage_dir: &'a Path,
pub(super) run_dir: PathBuf,
pub(super) mode: RunWorkerMode,
/// The Fabro home the server named; `None` falls back to Petri's own
/// lookup of the worker's environment.
pub(super) fabro_home: Option<PathBuf>,
pub(super) worker_token: &'a str,
/// The read-only credential the run's GitHub target is fetched with,
/// when the server resolved one.
pub(super) git_credential: Option<SourceCredential>,
pub(super) fabro_home: Option<PathBuf>,
pub(super) worker_token: &'a str,
}
/// Execute the run to its end. `Ok` when the record says it succeeded;
@ -203,17 +201,41 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
}
runner::set_worker_title(&run_id, WorkerTitlePhase::Running);
// A GitHub target is fetched into its sandbox with a read-only token and
// published with a push token, both minted from the server's
// credentials here.
let github = match publish::github_credentials(&*vault.read().await) {
Ok(credentials) => credentials,
Err(err) => {
warn!(run_id = %run_id, error = %err, "GitHub credentials are unavailable to the worker");
None
}
};
let source = RunSource::for_run(
worker.run_state.spec.target.as_ref(),
&worker.run_state.spec.settings.run,
worker.git_credential.clone(),
publish::source_credential(&worker.run_state.spec, github.as_ref()).await,
);
let publisher = publish::GitHubPublisher::for_run(
run_id,
&worker.run_state.spec,
github,
Arc::new(VaultCredentialSource::new(Arc::clone(&vault))),
Arc::new(
command_context::load_cli_catalog()
.context("failed to build the worker's pull request catalog")?,
),
Arc::clone(&records),
worker.client.clone_for_reuse(),
)
.map(|publisher| Arc::new(publisher) as Arc<dyn RunPublisher>);
let hooks = HooksSpec::for_run(
Arc::clone(&records),
&worker.run_state.spec.settings.run,
Arc::new(ClientArtifactWriter::new(worker.client.clone_for_reuse())),
)
.with_source(source)
.with_publisher(publisher)
.with_test_gates(test_checkpoint_gates());
let request = RunRequest {
run_id: run_id.to_string(),

View file

@ -0,0 +1,382 @@
//! A GitHub-target run's repository work in its worker: the read credential
//! its workspaces are fetched with, and its publication when it succeeds.
//!
//! The worker resolves the server's GitHub credentials itself, as the
//! legacy worker did: the strategy and App id from the server settings the
//! server named (`FABRO_CONFIG`), the App key the server hands the worker,
//! or `GITHUB_TOKEN` from the worker's vault snapshot. From them it mints a
//! read-only token for the checkout inside the sandbox when the run starts,
//! and a push token when the run ends, so a long run never pushes with an
//! expired token.
//!
//! Publication runs in Fabro's `run_finished` hook, after the last stage and
//! before the run's terminal record, as the legacy publish step did: the
//! final checkpoint is pushed from the run's snapshot repository to the run
//! branch on GitHub, and, when the run changed files and its settings ask
//! for one, a pull request is opened and recorded. A failure fails the run
//! with `publish_failed`.
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{Context, Result};
use base64::Engine as _;
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
use fabro_config::ServerSettingsBuilder;
use fabro_github::{GitCloneCredentials, GitHubContext, GitHubCredentials};
use fabro_llm::credentials::{CredentialProvider, readiness};
use fabro_llm::lithos_catalog::Catalog;
use fabro_petri::hooks::{Publication, RunPublisher};
use fabro_petri::platform_records::PlatformRecords;
use fabro_petri::source::SourceCredential;
use fabro_static::EnvVars;
use fabro_store::platform_records::{PlatformRecord, PullRequestCreatedRecord};
use fabro_types::settings::run::{PullRequestSettings, RunMode};
use fabro_types::settings::server::GithubIntegrationStrategy;
use fabro_types::{GitHubRepositorySlug, RunId, RunSpec, RunTarget};
use fabro_vault::Vault;
use fabro_workflow::pull_request::{self, AutoMergeOptions, OpenPullRequestRequest};
use tokio::process::Command;
use tokio::time;
use tracing::warn;
/// How long one push to the repository may take.
const PUSH_TIMEOUT: Duration = Duration::from_mins(5);
/// Attempts at the push: a freshly minted token can take a moment to reach
/// every GitHub replica.
const PUSH_ATTEMPTS: u32 = 3;
const PUSH_RETRY_DELAY: Duration = Duration::from_secs(2);
/// The server's GitHub credentials as the worker reaches them, or `None`
/// when none are configured.
pub(super) fn github_credentials(vault: &Vault) -> Result<Option<GitHubCredentials>> {
let settings = ServerSettingsBuilder::load_default().context("loading the server settings")?;
let github = &settings.server.integrations.github;
match github.strategy {
GithubIntegrationStrategy::App => {
GitHubCredentials::from_env_with_slug(github.app_id.as_deref(), github.slug.as_deref())
.map_err(anyhow::Error::msg)
}
GithubIntegrationStrategy::Token => {
let Some(token) = vault
.get(EnvVars::GITHUB_TOKEN)
.map(str::trim)
.filter(|token| !token.is_empty())
else {
return Ok(None);
};
fabro_github::validate_static_github_token(token)?;
Ok(Some(GitHubCredentials::Pat(token.to_string())))
}
}
}
/// The run's GitHub repository, when its target names one.
fn repository(spec: &RunSpec) -> Option<GitHubRepositorySlug> {
let Some(RunTarget::Git(target)) = spec.target.as_ref() else {
return None;
};
Some(target.clone().validate().ok()?.repository().clone())
}
/// The read-only credential the run's workspaces are fetched with. `None`
/// when the run has no GitHub target or no credentials resolve; a public
/// repository is then fetched anonymously.
pub(super) async fn source_credential(
spec: &RunSpec,
credentials: Option<&GitHubCredentials>,
) -> Option<SourceCredential> {
let repository = repository(spec)?;
let credentials = credentials?;
let base_url = fabro_github::github_api_base_url();
let context = GitHubContext::new(credentials, &base_url);
match fabro_github::resolve_read_only_clone_credentials(
&context,
repository.owner(),
repository.repo(),
)
.await
{
Ok(credentials) => SourceCredential::from_encoded(encode(&credentials)),
Err(err) => {
warn!(repository = %repository, error = %err, "no read credential for the run's repository; it is fetched anonymously");
None
}
}
}
fn encode(credentials: &GitCloneCredentials) -> String {
BASE64_STANDARD.encode(format!(
"{}:{}",
credentials.username(),
credentials.password()
))
}
/// A successful GitHub-target run's publication: the run branch pushed, and
/// the pull request its settings ask for opened.
pub(super) struct GitHubPublisher {
run_id: RunId,
spec: RunSpec,
repository: GitHubRepositorySlug,
credentials: Option<GitHubCredentials>,
pull_request: Option<PullRequestSettings>,
llm_source: Arc<dyn CredentialProvider>,
catalog: Arc<Catalog>,
records: Arc<dyn PlatformRecords>,
client: fabro_client::Client,
}
impl GitHubPublisher {
/// The publisher of a run whose target is a GitHub repository and whose
/// run branch is pushed; `None` for a dry run or any other run.
pub(super) fn for_run(
run_id: RunId,
spec: &RunSpec,
credentials: Option<GitHubCredentials>,
llm_source: Arc<dyn CredentialProvider>,
catalog: Arc<Catalog>,
records: Arc<dyn PlatformRecords>,
client: fabro_client::Client,
) -> Option<Self> {
let settings = &spec.settings.run;
if settings.execution.mode == RunMode::DryRun
|| !settings.run_branch.enabled
|| !settings.run_branch.push
{
return None;
}
let repository = repository(spec)?;
Some(Self {
run_id,
spec: spec.clone(),
repository,
credentials,
pull_request: settings
.pull_request
.clone()
.filter(|pull_request| pull_request.enabled),
llm_source,
catalog,
records,
client,
})
}
/// The model that writes the pull request: the run's, or the catalog's
/// default among the providers whose credentials resolve.
async fn model(&self) -> Option<String> {
if let Some(model) = self.spec.settings.run.model.name.clone() {
return Some(model);
}
let ready = readiness(self.catalog.enabled_providers(), self.llm_source.as_ref()).await;
self.catalog
.default_offering_for(&ready.ready)
.map(|entry| entry.model.id().to_string())
}
async fn open_pull_request(
&self,
context: GitHubContext<'_>,
settings: &PullRequestSettings,
publication: &Publication,
) -> Result<(), String> {
let Some(RunTarget::Git(target)) = self.spec.target.as_ref() else {
return Err("pull request creation requires a GitHub target".to_string());
};
let model = self
.model()
.await
.ok_or_else(|| "no LLM model is available to write the pull request".to_string())?;
// The run so far, for the pull request's details: best effort.
let run_state = self.client.get_run_state(&self.run_id).await.ok();
let origin_url = self.repository.https_url();
let created = pull_request::open_pull_request(OpenPullRequestRequest {
github: context,
origin_url: &origin_url,
base_branch: &target.branch,
head_branch: &publication.run_branch,
expected_head_sha: &publication.head_sha,
goal: &self.spec.graph.goal,
diff: &publication.patch,
model: &model,
draft: settings.draft,
auto_merge: settings.auto_merge.then_some(AutoMergeOptions {
merge_strategy: settings.merge_strategy,
}),
llm_source: Arc::clone(&self.llm_source),
catalog: Arc::clone(&self.catalog),
conclusion: None,
run_state: run_state.as_ref(),
})
.await
.map_err(|err| format!("failed to create pull request: {err}"))?;
let link = &created.link;
let record = PlatformRecord::PullRequestCreated(PullRequestCreatedRecord {
number: link.number,
owner: link.owner.clone(),
repo: link.repo.clone(),
html_url: link.html_url(),
head_sha: Some(publication.head_sha.clone()),
draft: settings.draft,
operation: None,
});
self.records
.append(&self.run_id, &record, None)
.await
.map_err(|err| format!("the pull request was opened but not recorded: {err}"))?;
Ok(())
}
}
#[async_trait::async_trait]
impl RunPublisher for GitHubPublisher {
async fn publish(&self, publication: &Publication) -> Result<(), String> {
let credentials = self.credentials.as_ref().ok_or_else(|| {
"pushing the run branch requires the server's GitHub credentials".to_string()
})?;
let base_url = fabro_github::github_api_base_url();
let context = GitHubContext::new(credentials, &base_url);
let push_credentials = fabro_github::resolve_clone_credentials(
&context,
self.repository.owner(),
self.repository.repo(),
)
.await
.map_err(|err| format!("no push credential for {}: {err:#}", self.repository))?;
push(&self.repository, &push_credentials, publication).await?;
let Some(settings) = &self.pull_request else {
return Ok(());
};
if publication.patch.trim().is_empty() {
return Ok(());
}
Box::pin(self.open_pull_request(context, settings, publication)).await
}
}
/// Push the run's final commit from its snapshot repository to the run
/// branch on GitHub, retrying a failure that may be a token still
/// replicating. The credential reaches `git` as an HTTP header for the
/// repository alone and never appears in the error.
async fn push(
repository: &GitHubRepositorySlug,
credentials: &GitCloneCredentials,
publication: &Publication,
) -> Result<(), String> {
let mut url = repository.https_url();
url.push_str(".git");
let header = format!("AUTHORIZATION: basic {}", encode(credentials));
let env = [
("GIT_CONFIG_COUNT", "1".to_string()),
("GIT_CONFIG_KEY_0", format!("http.{url}.extraheader")),
("GIT_CONFIG_VALUE_0", header),
];
let refspec = format!(
"{}:refs/heads/{}",
publication.head_sha, publication.run_branch
);
let mut last = String::new();
for attempt in 1..=PUSH_ATTEMPTS {
let output = git(
&publication.snapshot_repository,
&["push", "--quiet", &url, &refspec],
&env,
)
.await?;
if output.status.success() {
return Ok(());
}
last = String::from_utf8_lossy(&output.stderr)
.trim()
.replace(credentials.password(), "***");
warn!(
attempt,
branch = publication.run_branch,
error = last,
"pushing the run branch failed"
);
if attempt < PUSH_ATTEMPTS {
time::sleep(PUSH_RETRY_DELAY).await;
}
}
Err(format!(
"the run branch {} could not be pushed to {repository}: {last}",
publication.run_branch
))
}
/// `git` in `directory` with `env` added, non-interactive, bounded.
async fn git(
directory: &Path,
args: &[&str],
env: &[(&str, String)],
) -> Result<std::process::Output, String> {
let mut command = Command::new("git");
command
.args(args)
.current_dir(directory)
.envs(env.iter().map(|(key, value)| (*key, value)))
.env("GIT_TERMINAL_PROMPT", "0")
.stdin(std::process::Stdio::null())
.kill_on_drop(true);
match time::timeout(PUSH_TIMEOUT, command.output()).await {
Ok(Ok(output)) => Ok(output),
Ok(Err(err)) => Err(format!("git could not run: {err}")),
Err(_) => Err(format!("git push timed out after {PUSH_TIMEOUT:?}")),
}
}
#[cfg(test)]
mod tests {
use fabro_llm::credentials::NoCredentials;
use fabro_llm::test_support;
use fabro_petri::test_support::MemoryPlatformRecords;
use fabro_types::GitRunTarget;
use fabro_types::test_support::test_run_spec;
use super::*;
fn spec() -> RunSpec {
let mut spec = test_run_spec();
spec.target = Some(RunTarget::Git(GitRunTarget {
repo: "acme/widgets".to_string(),
branch: "main".to_string(),
tag: None,
sha: None,
}));
spec.settings.run.run_branch.enabled = true;
spec.settings.run.run_branch.push = true;
spec
}
#[test]
fn only_a_pushed_github_target_run_is_published() {
let publishes = |spec: &RunSpec| {
GitHubPublisher::for_run(
RunId::new(),
spec,
None,
Arc::new(NoCredentials),
Arc::new(test_support::test_catalog()),
Arc::new(MemoryPlatformRecords::new()),
fabro_client::Client::new_no_proxy("http://127.0.0.1:9").unwrap(),
)
.is_some()
};
assert!(publishes(&spec()));
let mut dry = spec();
dry.settings.run.execution.mode = RunMode::DryRun;
assert!(!publishes(&dry));
let mut unpushed = spec();
unpushed.settings.run.run_branch.push = false;
assert!(!publishes(&unpushed));
let mut empty = spec();
empty.target = Some(RunTarget::None {});
assert!(!publishes(&empty));
}
}

View file

@ -14,7 +14,6 @@ use fabro_interview::{
};
use fabro_manifest::SuppliedWorkflowVersionPackager;
use fabro_petri::controls::RunControls;
use fabro_petri::source::SourceCredential;
use fabro_tool::fabro_client::ClientBackend;
use fabro_types::RunId;
use fabro_vault::{SecretStore, Vault};
@ -63,7 +62,6 @@ pub(crate) async fn execute(
mode: RunWorkerMode,
fabro_home: Option<PathBuf>,
worker_token: &str,
git_credential: Option<String>,
) -> Result<()> {
let _ = fabro_proc::title_init();
set_worker_title(&run_id, initial_worker_title_phase(mode));
@ -91,7 +89,6 @@ pub(crate) async fn execute(
mode,
fabro_home,
worker_token,
git_credential: git_credential.and_then(SourceCredential::from_encoded),
}))
.await
}

View file

@ -58,23 +58,19 @@ async fn main() {
// inherits a process env that no longer contains this credential, so an
// unscrubbed spawn site cannot leak it. The token flows to `runner::execute`
// through explicit function arguments instead of the environment.
let worker_secrets = if subcommand == Some("__run-worker") {
let secrets = commands::run::WorkerSecrets {
token: process_env_var(EnvVars::FABRO_WORKER_TOKEN),
git_credential: process_env_var(EnvVars::FABRO_RUN_GIT_CREDENTIAL),
};
let worker_token = if subcommand == Some("__run-worker") {
let worker_token = process_env_var(EnvVars::FABRO_WORKER_TOKEN);
#[expect(
clippy::disallowed_methods,
reason = "Scrub the worker's credentials from this process's env before any \
child process is spawned, so no descendant can inherit them."
reason = "Scrub the worker bearer from this process's env before any \
child process is spawned, so no descendant can inherit it."
)]
{
std::env::remove_var(EnvVars::FABRO_WORKER_TOKEN);
std::env::remove_var(EnvVars::FABRO_RUN_GIT_CREDENTIAL);
}
secrets
worker_token
} else {
commands::run::WorkerSecrets::default()
None
};
install_miette_hook();
@ -84,7 +80,7 @@ async fn main() {
let start = std::time::Instant::now();
let (command_name, result) = Box::pin(main_inner(worker_secrets)).await;
let (command_name, result) = Box::pin(main_inner(worker_token)).await;
let duration_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX);
let exit_code = result.as_ref().err().map_or(0, exit::exit_code_for);
@ -192,7 +188,7 @@ pub(crate) fn process_env_var(name: &str) -> Option<String> {
std::env::var(name).ok()
}
async fn main_inner(worker_secrets: commands::run::WorkerSecrets) -> (String, Result<()>) {
async fn main_inner(worker_token: Option<String>) -> (String, Result<()>) {
let _ = default_provider().install_default();
let cli = Cli::parse();
@ -255,7 +251,7 @@ async fn main_inner(worker_secrets: commands::run::WorkerSecrets) -> (String, Re
commands::exec::execute(args, &base_ctx).await?;
}
Commands::RunCmd(cmd) => {
Box::pin(commands::run::dispatch(cmd, &base_ctx, worker_secrets)).await?;
Box::pin(commands::run::dispatch(cmd, &base_ctx, worker_token)).await?;
}
Commands::Preflight(args) => {
commands::preflight::execute(args, &base_ctx).await?;

View file

@ -294,7 +294,7 @@ impl GitAuthConfig {
}
}
pub(crate) fn git_env(&self, clone_url: &str) -> Vec<(String, String)> {
fn git_env(&self, clone_url: &str) -> Vec<(String, String)> {
vec![
("GIT_CONFIG_COUNT".to_string(), "1".to_string()),
(
@ -305,7 +305,7 @@ impl GitAuthConfig {
]
}
pub(crate) fn sensitive_values(&self) -> &[String] {
fn sensitive_values(&self) -> &[String] {
&self.sensitive_values
}
}

View file

@ -169,7 +169,6 @@ mod handler;
pub(crate) mod petri_runs;
mod pull_request_supervisor;
pub(crate) mod resource_sampler;
pub(crate) mod run_publication;
pub(crate) mod run_records;
mod session_runtime;
pub(crate) mod stream_follower;
@ -3828,7 +3827,6 @@ fn worker_launch_spec(
run_dir: &std::path::Path,
agent_fabro_tools_enabled: bool,
github_app_private_key: Option<String>,
run_git_credential: Option<String>,
) -> anyhow::Result<WorkerLaunchSpec> {
let current_exe = std::env::current_exe().context("reading current executable path")?;
let executable =
@ -3867,7 +3865,6 @@ fn worker_launch_spec(
fabro_log,
active_config_path: state.active_config_path().to_path_buf(),
github_app_private_key,
run_git_credential,
fabro_home: fabro_config::Home::from_env().root().to_path_buf(),
})
}
@ -4185,9 +4182,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
return;
}
// The read-only credential the worker fetches a GitHub target with,
// resolved for this launch.
let run_git_credential = run_publication::clone_credential(&state, &run_state.spec).await;
// The worker reads the Daytona key from the vault itself; only the
// GitHub App key crosses on its command.
let github_app_private_key = match state.vault_secret(EnvVars::GITHUB_APP_PRIVATE_KEY).await {
@ -4214,7 +4208,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
&run_dir_for_build,
agent_fabro_tools_enabled,
github_app_private_key,
run_git_credential,
)
})
.await
@ -4324,7 +4317,6 @@ async fn execute_run_subprocess(state: Arc<AppState>, run_id: RunId) {
};
accumulate_concluded_run_usage(&state, &final_state);
run_publication::spawn(Arc::clone(&state), run_id);
let mut runs = state.runs.lock().expect("runs lock poisoned");
if let Some(managed_run) = runs.get_mut(&run_id) {

View file

@ -314,36 +314,6 @@ fn available_pull_request_response(
}
}
/// Record that a pull request should be created for the run, unless one
/// exists or a creation is already pending. The caller holds the run's
/// pull request create lock. `Ok(true)` when this call recorded the request.
pub(in crate::server) async fn request_pull_request_creation(
state: &Arc<AppState>,
id: RunId,
model: String,
force: bool,
) -> Result<bool, ApiError> {
// Under the create lock, the projection is the latest word on whether a
// pull request exists or a creation is already pending.
let run_state = state.load_run_projection(&id).await?;
let appended = run_state.pull_request.is_none()
&& !run_state
.pull_request_creation
.as_ref()
.is_some_and(fabro_types::PullRequestCreation::is_pending);
if appended {
let record = PlatformRecord::PullRequestRequested(PullRequestRequestedRecord {
creation_id: fabro_types::PullRequestCreationId::new(),
model,
force,
});
run_records::append(state, id, record)
.await
.map_err(|err| ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string()))?;
}
Ok(appended)
}
async fn create_run_pull_request(
RequireRunScoped(id): RequireRunScoped,
State(state): State<Arc<AppState>>,
@ -383,10 +353,29 @@ async fn create_run_pull_request(
}
};
let _create_guard = state.pull_request_create_locks.lock(id).await;
let appended = match request_pull_request_creation(&state, id, model, body.force).await {
Ok(appended) => appended,
let creation_id = fabro_types::PullRequestCreationId::new();
// Under the create lock, the projection is the latest word on whether a
// pull request exists or a creation is already pending.
let run_state = match state.load_run_projection(&id).await {
Ok(run_state) => run_state,
Err(err) => return err.into_response(),
};
let appended = run_state.pull_request.is_none()
&& !run_state
.pull_request_creation
.as_ref()
.is_some_and(fabro_types::PullRequestCreation::is_pending);
if appended {
let record = PlatformRecord::PullRequestRequested(PullRequestRequestedRecord {
creation_id,
model,
force: body.force,
});
if let Err(err) = run_records::append(&state, id, record).await {
return ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
.into_response();
}
}
let run_state = match state.load_run_projection(&id).await {
Ok(run_state) => run_state,

View file

@ -51,7 +51,7 @@ use fabro_petri::providers::SandboxProviderConfig;
use fabro_petri::recovery::{self, Recovery, RecoveryRequest};
use fabro_petri::runtime::{self, RuntimeSpec};
use fabro_petri::secrets::VaultSecrets;
use fabro_petri::source::{RunSource, SourceCredential};
use fabro_petri::source::RunSource;
use fabro_petri::{SqliteRunStore, admission, projection, run_graph};
use fabro_static::EnvVars;
use fabro_store::platform_records::{RunLifecycleKind, RunLifecycleRecord};
@ -461,12 +461,12 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
let observers = vec![petri_interviewer.observer()];
let (_, eligible) = state.resolve_llm_client_with_ready_ids().await;
let dry_run = run_state.spec.settings.run.execution.mode == RunMode::DryRun;
// The in-process path serves the server's tests: a Git target is
// fetched without a credential, and nothing is published.
let source = RunSource::for_run(
run_state.spec.target.as_ref(),
&run_state.spec.settings.run,
super::run_publication::clone_credential(&state, &run_state.spec)
.await
.and_then(SourceCredential::from_encoded),
None,
);
let hooks = HooksSpec::for_run(
Arc::new(SqlitePlatformRecords::new(Arc::clone(
@ -546,7 +546,6 @@ pub(crate) async fn execute(state: Arc<AppState>, run_id: RunId) {
warn!(run_id = %run_id, error = ?err, "the run's final state could not be read for the usage aggregate");
}
}
super::run_publication::spawn(Arc::clone(&state), run_id);
}
/// Bring a Petri run the server left in flight back to its worker after a

View file

@ -1,499 +0,0 @@
//! A GitHub-target run's work leaves for its repository from the server.
//!
//! The run's workspaces are checked out from the repository inside their
//! sandboxes, and every checkpoint reaches the run's snapshot repository on
//! this host (`fabro_petri::checkpoint`). When the run ends, the server
//! pushes the final checkpoint to the run branch, `fabro/run/<id>`, on the
//! repository with its own GitHub credentials, then, for a successful run
//! with changes whose settings ask for one, records the pull request request
//! the creation supervisor opens. The sandbox never holds a credential that
//! can push.
//!
//! A run that did not succeed, a dry run, a run whose run branch is not
//! pushed, and a run with no GitHub target publish nothing. A push that
//! fails is a warning notice on the run, as is a pull request that cannot
//! be requested; neither changes the run's outcome.
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
use base64::Engine as _;
use base64::engine::general_purpose::STANDARD as BASE64_STANDARD;
use fabro_config::Storage;
use fabro_store::platform_records::{PlatformRecord, RunNoticeRecord};
use fabro_types::settings::run::RunMode;
use fabro_types::{RunId, RunNoticeLevel, RunSpec, RunTarget};
use tokio::process::Command;
use tokio::{fs, time};
use tracing::{info, warn};
use super::handler::pull_requests::request_pull_request_creation;
use super::{AppState, run_records};
use crate::git_checkout::{self, GitAuthConfig};
/// How long one push to the repository may take.
const PUSH_TIMEOUT: Duration = Duration::from_mins(5);
/// Attempts at the push: a fresh token can take a moment to reach every
/// GitHub replica.
const PUSH_ATTEMPTS: u32 = 3;
const PUSH_RETRY_DELAY: Duration = Duration::from_secs(2);
/// The read-only credential a run's worker fetches its GitHub target with,
/// as the base64 of `username:password`. `None` when the run checks nothing
/// out, or the server has no GitHub credentials for the repository (a
/// public repository is then fetched anonymously).
pub(crate) async fn clone_credential(state: &AppState, spec: &RunSpec) -> Option<String> {
let Some(RunTarget::Git(target)) = spec.target.as_ref() else {
return None;
};
let settings = &spec.settings.run;
if !settings.clone.enabled || settings.execution.mode == RunMode::DryRun {
return None;
}
let validated = target.clone().validate().ok()?;
let github = &state.server_settings().server.integrations.github;
let credentials = match state.github_credentials(github).await {
Ok(Some(credentials)) => credentials,
Ok(None) => return None,
Err(err) => {
warn!(error = %err, "GitHub credentials are unavailable; the run's repository is fetched anonymously");
return None;
}
};
let context = match state.http_client.clone() {
Some(client) => fabro_github::GitHubContext::with_http_client(
&credentials,
&state.github_api_base_url,
client,
),
None => fabro_github::GitHubContext::new(&credentials, &state.github_api_base_url),
};
let repository = validated.repository();
match fabro_github::resolve_read_only_clone_credentials(
&context,
repository.owner(),
repository.repo(),
)
.await
{
Ok(credentials) => Some(BASE64_STANDARD.encode(format!(
"{}:{}",
credentials.username(),
credentials.password()
))),
Err(err) => {
warn!(repository = %repository, error = %err, "no read credential for the run's repository; it is fetched anonymously");
None
}
}
}
/// Publish the run's work once it has ended, in the background.
pub(crate) fn spawn(state: Arc<AppState>, run_id: RunId) {
tokio::spawn(async move {
if let Err(err) = publish(&state, run_id).await {
warn!(run_id = %run_id, error = %err, "the run's work was not published");
notice(&state, run_id, "run_publish_failed", &format!("{err:#}")).await;
}
});
}
/// What an ended run pushes: the repository, the branch, the commit, and
/// whether a pull request follows.
struct Publication {
repository: fabro_types::GitHubRepositorySlug,
run_branch: String,
sha: String,
pull_request: bool,
model: Option<String>,
}
async fn publish(state: &Arc<AppState>, run_id: RunId) -> anyhow::Result<()> {
state.petri_projector.settle(run_id).await;
let Some(projection) = run_records::projection(state, run_id).await? else {
return Ok(());
};
let Some(publication) = publication(&projection) else {
return Ok(());
};
let snapshots = Storage::new(state.server_storage_dir())
.run_scratch(&run_id)
.root()
.join("petri")
.join("snapshots");
let Some(repository) = snapshot_holding(&snapshots, &publication.sha).await else {
anyhow::bail!(
"the run's final commit {} is in none of its snapshot repositories",
publication.sha
);
};
push(state, &repository, &publication).await?;
info!(
run_id = %run_id,
repository = %publication.repository,
branch = publication.run_branch,
sha = publication.sha,
"run branch pushed"
);
if publication.pull_request {
request_pull_request(state, run_id, publication.model).await;
}
Ok(())
}
/// What the ended run publishes, or `None` when it publishes nothing.
fn publication(projection: &fabro_store::RunProjection) -> Option<Publication> {
let spec = &projection.spec;
let Some(RunTarget::Git(target)) = spec.target.as_ref() else {
return None;
};
let settings = &spec.settings.run;
if settings.execution.mode == RunMode::DryRun
|| !settings.run_branch.enabled
|| !settings.run_branch.push
{
return None;
}
let conclusion = projection.conclusion.as_ref()?;
if !conclusion.status.is_successful() {
return None;
}
let sha = conclusion
.final_git_commit_sha
.clone()
.filter(|sha| !sha.trim().is_empty())?;
let run_branch = projection
.start
.as_ref()
.and_then(|start| start.run_branch.clone())?;
let repository = target.clone().validate().ok()?.repository().clone();
let has_changes = conclusion
.diff
.patch
.as_deref()
.is_some_and(|patch| !patch.trim().is_empty());
let pull_request = has_changes
&& settings
.pull_request
.as_ref()
.is_some_and(|pull_request| pull_request.enabled);
Some(Publication {
repository,
run_branch,
sha,
pull_request,
model: settings.model.name.clone(),
})
}
/// The snapshot repository of the run that holds `sha`.
async fn snapshot_holding(snapshots: &Path, sha: &str) -> Option<PathBuf> {
let mut entries = fs::read_dir(snapshots).await.ok()?;
let mut repositories = Vec::new();
while let Ok(Some(entry)) = entries.next_entry().await {
let path = entry.path();
if path.extension().is_some_and(|extension| extension == "git") {
repositories.push(path);
}
}
repositories.sort();
for repository in repositories {
let held = git(
&repository,
&["cat-file", "-e", &format!("{sha}^{{commit}}")],
&[],
)
.await
.is_ok_and(|output| output.status.success());
if held {
return Some(repository);
}
}
None
}
/// Push `sha` from the snapshot repository to the run branch on GitHub,
/// with the server's credentials, retrying a failure that may be a token
/// still replicating.
async fn push(
state: &AppState,
repository: &Path,
publication: &Publication,
) -> anyhow::Result<()> {
let github = &state.server_settings().server.integrations.github;
let credentials = state
.github_credentials(github)
.await?
.ok_or_else(|| anyhow::anyhow!("the server has no GitHub credentials to push with"))?;
let context = match state.http_client.clone() {
Some(client) => fabro_github::GitHubContext::with_http_client(
&credentials,
&state.github_api_base_url,
client,
),
None => fabro_github::GitHubContext::new(&credentials, &state.github_api_base_url),
};
let push_credentials = fabro_github::resolve_clone_credentials(
&context,
publication.repository.owner(),
publication.repository.repo(),
)
.await?;
let auth = GitAuthConfig::new(&push_credentials);
let url = git_checkout::github_clone_url(&publication.repository);
let env = auth.git_env(&url);
let refspec = format!("{}:refs/heads/{}", publication.sha, publication.run_branch);
let mut last = String::new();
for attempt in 1..=PUSH_ATTEMPTS {
let output = git(repository, &["push", "--quiet", &url, &refspec], &env).await?;
if output.status.success() {
return Ok(());
}
last = redact(
String::from_utf8_lossy(&output.stderr).trim(),
auth.sensitive_values(),
);
warn!(
attempt,
branch = publication.run_branch,
error = last,
"pushing the run branch failed"
);
if attempt < PUSH_ATTEMPTS {
time::sleep(PUSH_RETRY_DELAY).await;
}
}
anyhow::bail!(
"the run branch {} could not be pushed to {}: {last}",
publication.run_branch,
publication.repository
)
}
/// Record the pull request request for the run, with the run's model or
/// the catalog's default, and hand it to the creation supervisor.
async fn request_pull_request(state: &Arc<AppState>, run_id: RunId, model: Option<String>) {
let model = if let Some(model) = model {
model
} else {
let configured = state.ready_llm_provider_ids().await;
let catalog = state.catalog();
let Some(entry) = catalog.default_offering_for(&configured) else {
notice(
state,
run_id,
"pull_request_not_requested",
"no LLM model is available to write the pull request",
)
.await;
return;
};
entry.model.id().to_string()
};
let guard = state.pull_request_create_locks.lock(run_id).await;
let requested = request_pull_request_creation(state, run_id, model, false).await;
drop(guard);
match requested {
Ok(true) => {
if let Ok(projection) = state.load_run_projection(&run_id).await {
if let Some(creation) = projection.pull_request_creation.as_ref() {
state.enqueue_pull_request_creation(run_id, creation.requested_at);
state.notify_pull_request_scheduler();
}
}
}
Ok(false) => {}
Err(err) => {
warn!(run_id = %run_id, error = ?err, "the pull request was not requested");
notice(
state,
run_id,
"pull_request_not_requested",
"the pull request could not be requested",
)
.await;
}
}
}
async fn notice(state: &AppState, run_id: RunId, code: &str, message: &str) {
let record = PlatformRecord::RunNotice(RunNoticeRecord {
level: RunNoticeLevel::Warn,
code: code.to_string(),
message: message.to_string(),
});
if let Err(err) = run_records::append(state, run_id, record).await {
warn!(run_id = %run_id, error = %err, "the run's publication notice was not recorded");
}
}
/// `git` in `directory` with `env` added, non-interactive, bounded.
async fn git(
directory: &Path,
args: &[&str],
env: &[(String, String)],
) -> anyhow::Result<std::process::Output> {
let mut command = Command::new("git");
command
.args(args)
.current_dir(directory)
.envs(env.iter().map(|(key, value)| (key, value)))
.env("GIT_TERMINAL_PROMPT", "0")
.stdin(std::process::Stdio::null())
.kill_on_drop(true);
Ok(time::timeout(PUSH_TIMEOUT, command.output())
.await
.map_err(|_| anyhow::anyhow!("git {} timed out", args.first().unwrap_or(&"")))??)
}
fn redact(text: &str, secrets: &[String]) -> String {
secrets
.iter()
.filter(|secret| !secret.is_empty())
.fold(text.to_string(), |text, secret| text.replace(secret, "***"))
}
#[cfg(test)]
mod tests {
use chrono::Utc;
use fabro_types::settings::run::PullRequestSettings;
use fabro_types::{
Conclusion, GitRunTarget, RunProjection, RunTarget, StageOutcome, StartRecord, test_support,
};
use super::*;
const SHA: &str = "0123456789abcdef0123456789abcdef01234567";
/// A GitHub-target run that succeeded with a change and a pull request
/// asked for, pushed to its run branch.
fn ended() -> RunProjection {
let mut spec = test_support::test_run_spec();
spec.target = Some(RunTarget::Git(GitRunTarget {
repo: "acme/widgets".to_string(),
branch: "main".to_string(),
tag: None,
sha: None,
}));
spec.settings.run.run_branch.enabled = true;
spec.settings.run.run_branch.push = true;
spec.settings.run.pull_request = Some(PullRequestSettings {
enabled: true,
..PullRequestSettings::default()
});
let mut projection = RunProjection::new(String::new(), spec, Utc::now());
projection.start = Some(StartRecord {
start_time: Utc::now(),
run_branch: Some("fabro/run/1".to_string()),
base_sha: None,
});
let mut conclusion = Conclusion::outcome_only(Utc::now(), StageOutcome::Succeeded, None);
conclusion.final_git_commit_sha = Some(SHA.to_string());
conclusion.diff.patch = Some("diff --git a/README.md b/README.md\n".to_string());
projection.conclusion = Some(conclusion);
projection
}
#[test]
fn a_successful_github_run_pushes_its_run_branch_and_asks_for_a_pull_request() {
let publication = publication(&ended()).expect("the run publishes");
assert_eq!(publication.repository.to_string(), "acme/widgets");
assert_eq!(publication.run_branch, "fabro/run/1");
assert_eq!(publication.sha, SHA);
assert!(publication.pull_request);
}
#[test]
fn a_run_without_changes_or_a_pull_request_setting_pushes_without_one() {
let mut unchanged = ended();
unchanged.conclusion.as_mut().unwrap().diff.patch = Some(" \n".to_string());
assert!(!publication(&unchanged).unwrap().pull_request);
let mut not_asked = ended();
not_asked.spec.settings.run.pull_request = None;
assert!(!publication(&not_asked).unwrap().pull_request);
}
#[test]
fn nothing_is_published_for_a_failed_dry_unpushed_or_non_github_run() {
let mut failed = ended();
failed.conclusion.as_mut().unwrap().status = StageOutcome::Failed {
retry_requested: false,
};
assert!(publication(&failed).is_none());
let mut dry = ended();
dry.spec.settings.run.execution.mode = RunMode::DryRun;
assert!(publication(&dry).is_none());
let mut unpushed = ended();
unpushed.spec.settings.run.run_branch.push = false;
assert!(publication(&unpushed).is_none());
let mut folder = ended();
folder.spec.target = Some(RunTarget::None {});
assert!(publication(&folder).is_none());
let mut no_commit = ended();
no_commit.conclusion.as_mut().unwrap().final_git_commit_sha = None;
assert!(publication(&no_commit).is_none());
}
#[test]
fn a_push_failure_never_prints_the_credential() {
assert_eq!(
redact("fatal: token s3cret rejected", &["s3cret".to_string()]),
"fatal: token *** rejected"
);
}
#[tokio::test]
#[expect(
clippy::disallowed_methods,
reason = "the test builds its fixture repositories with synchronous git"
)]
async fn the_snapshot_repository_holding_the_final_commit_is_found() {
let root = tempfile::tempdir().unwrap();
let snapshots = root.path().join("snapshots");
let work = root.path().join("work");
std::fs::create_dir_all(&snapshots).unwrap();
std::fs::create_dir_all(&work).unwrap();
let run = |directory: &Path, args: &[&str]| {
let output = std::process::Command::new("git")
.current_dir(directory)
.args(args)
.output()
.unwrap();
assert!(output.status.success(), "{args:?}");
String::from_utf8(output.stdout).unwrap().trim().to_string()
};
run(&work, &["init", "-q"]);
run(&work, &[
"-c",
"user.name=t",
"-c",
"user.email=t@example.com",
"commit",
"-q",
"--allow-empty",
"-m",
"one",
]);
let sha = run(&work, &["rev-parse", "HEAD"]);
for name in ["a.git", "b.git"] {
run(&snapshots, &["init", "-q", "--bare", name]);
}
run(&work, &[
"push",
"-q",
&snapshots.join("b.git").to_string_lossy(),
"HEAD:refs/checkpoints/0/1/1",
]);
assert_eq!(
snapshot_holding(&snapshots, &sha).await,
Some(snapshots.join("b.git"))
);
assert_eq!(snapshot_holding(&snapshots, SHA).await, None);
}
}

View file

@ -2274,7 +2274,6 @@ fn worker_command_forwards_github_app_private_key_from_vault() {
storage_dir.path(),
false,
Some("test-private-key".to_string()),
None,
)
.unwrap();
let cmd = LocalWorkerRuntime::command_for_spec(&spec);
@ -2289,41 +2288,6 @@ fn worker_command_forwards_github_app_private_key_from_vault() {
);
}
#[cfg(unix)]
#[test]
fn worker_command_carries_the_run_git_credential_only_when_resolved() {
let storage_dir = tempfile::tempdir().unwrap();
let state = worker_command_test_state(storage_dir.path(), &["dev-token"], Some(TEST_DEV_TOKEN));
let spec = worker_launch_spec(
state.as_ref(),
RunId::new(),
RunExecutionMode::Start,
storage_dir.path(),
false,
None,
Some("eC1hY2Nlc3MtdG9rZW46c2VjcmV0".to_string()),
)
.unwrap();
let cmd = LocalWorkerRuntime::command_for_spec(&spec);
assert_eq!(
command_env_value(&cmd, EnvVars::FABRO_RUN_GIT_CREDENTIAL),
EnvOverride::Set("eC1hY2Nlc3MtdG9rZW46c2VjcmV0".to_string())
);
let cmd = worker_command(
state.as_ref(),
RunId::new(),
RunExecutionMode::Start,
storage_dir.path(),
false,
)
.unwrap();
assert_eq!(
command_env_value(&cmd, EnvVars::FABRO_RUN_GIT_CREDENTIAL),
EnvOverride::Unchanged
);
}
#[cfg(unix)]
#[test]
fn worker_command_omits_github_app_private_key_when_unset() {
@ -2526,7 +2490,6 @@ fn worker_command(
run_dir,
agent_fabro_tools_enabled,
None,
None,
)?;
Ok(LocalWorkerRuntime::command_for_spec(&spec))
}

View file

@ -48,8 +48,6 @@ pub(crate) struct WorkerLaunchSpec {
pub(crate) fabro_log: Option<String>,
pub(crate) active_config_path: PathBuf,
pub(crate) github_app_private_key: Option<String>,
/// The read-only credential the run's GitHub target is fetched with.
pub(crate) run_git_credential: Option<String>,
/// The Fabro home the server resolved, so a Petri run's skills step
/// reads the same home whatever the worker's environment says.
pub(crate) fabro_home: PathBuf,
@ -111,9 +109,6 @@ impl LocalWorkerRuntime {
if let Some(pem) = spec.github_app_private_key.as_deref() {
cmd.env(EnvVars::GITHUB_APP_PRIVATE_KEY, pem);
}
if let Some(credential) = spec.run_git_credential.as_deref() {
cmd.env(EnvVars::FABRO_RUN_GIT_CREDENTIAL, credential);
}
#[cfg(unix)]
fabro_proc::pre_exec_setpgid(cmd.as_std_mut());

View file

@ -131,14 +131,16 @@ pub enum RunStatus {
/// What the durable record says about the run once it ended.
#[derive(Clone, Debug)]
pub struct RunOutcome {
pub status: RunStatus,
pub status: RunStatus,
/// The root invocation's failure message, when it failed.
pub failure: Option<String>,
pub failure: Option<String>,
/// Whether the record is whole: the run recorded its finish and every
/// log replays byte for byte.
pub complete: bool,
pub complete: bool,
/// Every reason `complete` is false.
pub incomplete: Vec<String>,
pub incomplete: Vec<String>,
/// Whether the run failed publishing its work after its last stage.
pub publish_failed: bool,
}
/// Why the run could not be executed or its outcome read.
@ -301,6 +303,16 @@ pub async fn run(request: RunRequest) -> Result<RunOutcome, RunError> {
outcome.status = RunStatus::Failed;
outcome.failure = Some(failure);
}
// A successful run whose publication failed is a failed run: its work
// did not reach where the settings sent it.
if let Some(failure) = fabro_hooks
.as_ref()
.and_then(|hooks| hooks.publish_failure())
{
outcome.status = RunStatus::Failed;
outcome.failure = Some(failure);
outcome.publish_failed = true;
}
Ok(outcome)
}
@ -360,6 +372,7 @@ pub fn conclusion(result: &Result<RunOutcome, RunError>) -> Conclusion {
Ok(outcome) => {
let reason = match outcome.status {
RunStatus::Cancelled => FailureReason::Cancelled,
RunStatus::Failed if outcome.publish_failed => FailureReason::PublishFailed,
RunStatus::Success | RunStatus::Failed => FailureReason::WorkflowError,
};
Conclusion::Failed {
@ -493,6 +506,7 @@ fn outcome(
failure,
complete: inspection.complete,
incomplete: inspection.incomplete,
publish_failed: false,
})
}
@ -531,9 +545,21 @@ mod tests {
} else {
vec!["execution 0 did not finish".to_string()]
},
publish_failed: false,
}
}
#[test]
fn a_failed_publication_concludes_publish_failed() {
let mut outcome = outcome_with(RunStatus::Failed, Some("the push was rejected"), true);
outcome.publish_failed = true;
let Conclusion::Failed { reason, message } = conclusion(&Ok(outcome)) else {
panic!("the run failed");
};
assert_eq!(reason, FailureReason::PublishFailed);
assert!(message.contains("the push was rejected"), "{message}");
}
#[test]
fn a_whole_successful_record_concludes_succeeded() {
assert_eq!(

View file

@ -30,9 +30,11 @@
//! was already collected earlier in the run. A failed write is a recorded
//! problem on the transition, never a blocked route.
//! - `run_finished`: the run's diff, its run branch against its base commit, as
//! the `run.diff` platform record with the patch as a blob; then the
//! forwarded point, so the local service runs `run_complete` and `run_failed`
//! with the sandbox in place.
//! the `run.diff` platform record with the patch as a blob; for a successful
//! run, its publication ([`RunPublisher`]: the platform pushes the run branch
//! and opens a pull request), whose failure fails the run before its terminal
//! record; then the forwarded point, so the local service runs `run_complete`
//! and `run_failed` with the sandbox in place.
//! - `scope_acquired`: a fresh run's Git target checked out into the workspace
//! from inside the scope, with its snapshot repository seeded with the
//! starting commit ([`crate::source`]); a resumed run's workspace brought to
@ -99,7 +101,7 @@ use petri_runtime::driver::lifecycle::{
ScopeReleased, Transition, TransitionError, TransitionReport,
};
use petri_runtime::executor::{EnvError, ExecEnv};
use petri_runtime::ir::{ExecutionId, FailureInfo, ScopeId, Status};
use petri_runtime::ir::{ExecutionId, FailureInfo, RunStatus, ScopeId, Status};
use serde_json::json;
use tokio::sync::{Mutex as AsyncMutex, OnceCell};
use tokio::{fs, time};
@ -234,6 +236,29 @@ pub struct HooksSpec {
/// Where a Git target's workspaces are checked out from, when a fresh
/// run first acquires them. `None` for a run with no remote repository.
pub source: Option<RunSource>,
/// What a successful run's work does when it ends; `None` publishes
/// nothing.
pub publisher: Option<Arc<dyn RunPublisher>>,
}
/// What a successful run hands its publisher when it ends: the run branch,
/// the commit it ends on, the snapshot repository that holds it, and the
/// run's patch against the commit the branch started from.
#[derive(Clone, Debug)]
pub struct Publication {
pub run_branch: String,
pub head_sha: String,
pub snapshot_repository: PathBuf,
pub patch: String,
}
/// The platform's end-of-run publication: what a successful run's work does
/// after its last stage and before its terminal record, such as pushing the
/// run branch and opening a pull request. An `Err` fails the run with the
/// message, as a failed publish did on the legacy executor.
#[async_trait::async_trait]
pub trait RunPublisher: Send + Sync {
async fn publish(&self, publication: &Publication) -> Result<(), String>;
}
impl HooksSpec {
@ -252,9 +277,17 @@ impl HooksSpec {
test_gates: None,
artifact_writer,
source: None,
publisher: None,
}
}
/// Publish a successful run's work through `publisher` when it ends.
#[must_use]
pub fn with_publisher(mut self, publisher: Option<Arc<dyn RunPublisher>>) -> Self {
self.publisher = publisher;
self
}
/// Check a Git target's workspaces out from `source`.
#[must_use]
pub fn with_source(mut self, source: Option<RunSource>) -> Self {
@ -411,6 +444,9 @@ pub struct FabroHooks {
scopes: ScopeEnvs,
/// The checkpoint failure that ended the run, when one did.
failure: Mutex<Option<String>>,
publisher: Option<Arc<dyn RunPublisher>>,
/// Why the run's publication failed, when it did.
publish_failure: Mutex<Option<String>>,
/// Whether the run continues from its records: a sandbox workspace is
/// then brought to its snapshot when its scope is first acquired.
resumed: bool,
@ -470,6 +506,8 @@ impl FabroHooks {
},
scopes: ScopeEnvs::default(),
failure: Mutex::default(),
publisher: spec.publisher,
publish_failure: Mutex::default(),
resumed,
restore: OnceCell::new(),
store,
@ -491,6 +529,13 @@ impl FabroHooks {
sync::lock(&self.failure).clone()
}
/// Why the run's publication failed, when it did: the run then fails
/// with this message.
#[must_use]
pub fn publish_failure(&self) -> Option<String> {
sync::lock(&self.publish_failure).clone()
}
/// The run's workspaces on this host, as the hooks reach them.
#[must_use]
pub fn workspaces(&self) -> &RunWorkspaces {
@ -1159,7 +1204,7 @@ impl FabroHooks {
/// the branch started from, in the snapshot repository on this host.
/// Nothing is recorded for a run that never created its branch or
/// never checkpointed.
async fn record_run_diff(&self) -> Result<(), HookError> {
async fn record_run_diff(&self) -> Result<Option<Publication>, HookError> {
self.recorded_checkpoints().await?;
let branch = match self.checkpoints.branch.get() {
Some(branch) => Some(branch.clone()),
@ -1167,14 +1212,14 @@ impl FabroHooks {
};
let Some(branch) = branch else {
debug!(run_id = %self.run_id, "no run branch is recorded; no run diff");
return Ok(());
return Ok(None);
};
let Some(base_sha) = branch.base_sha.clone() else {
return Ok(());
return Ok(None);
};
let Some((workspace, head_sha)) = self.checkpoints.last() else {
debug!(run_id = %self.run_id, "no checkpoint is recorded; no run diff");
return Ok(());
return Ok(None);
};
// The run's diff is measured in the workspace the branch started
// in; a last checkpoint elsewhere (a nested invocation's workspace)
@ -1189,6 +1234,12 @@ impl FabroHooks {
source,
})?;
let patch_blob = self.patch_blob(&diff).await?;
let publication = branch.run_branch.clone().map(|run_branch| Publication {
run_branch,
head_sha: head_sha.clone(),
snapshot_repository: self.workspaces.snapshot_repository(&workspace),
patch: diff.patch.clone(),
});
let record = PlatformRecord::RunDiff(RunDiffRecord {
base_sha: Some(base_sha),
head_sha: Some(head_sha),
@ -1209,7 +1260,31 @@ impl FabroHooks {
deletions = diff.summary.deletions,
"run diff recorded"
);
Ok(())
Ok(publication)
}
/// Hand a successful run's work to the publisher, before the run's
/// terminal record. A failure fails the run with its reason.
async fn publish(&self, publication: Option<Publication>) {
let Some(publisher) = &self.publisher else {
return;
};
let Some(publication) = publication else {
debug!(run_id = %self.run_id, "the run has no run branch head; nothing to publish");
return;
};
match publisher.publish(&publication).await {
Ok(()) => info!(
run_id = %self.run_id,
branch = publication.run_branch,
sha = publication.head_sha,
"run published"
),
Err(message) => {
warn!(run_id = %self.run_id, error = %message, "the run's publication failed");
*sync::lock(&self.publish_failure) = Some(message);
}
}
}
/// Hold at a test gate when one is set for this point and node.
@ -1425,8 +1500,15 @@ impl ExecutionHooks for FabroHooks {
failure = finished.failure.as_deref().unwrap_or(""),
"Petri run finished; recording the run's diff and running the run-end hooks"
);
if let Err(error) = self.record_run_diff().await {
warn!(run_id = %self.run_id, error = %error.render(), "the run's diff was not recorded");
let publication = match self.record_run_diff().await {
Ok(publication) => publication,
Err(error) => {
warn!(run_id = %self.run_id, error = %error.render(), "the run's diff was not recorded");
None
}
};
if finished.status == RunStatus::Success && self.checkpoint_failure().is_none() {
self.publish(publication).await;
}
self.inner.run_finished(context, finished).await
}
@ -1580,6 +1662,7 @@ mod tests {
test_gates: None,
artifact_writer,
source: None,
publisher: None,
},
Arc::new(NoHooks),
run_id,

View file

@ -27,7 +27,7 @@ use fabro_petri::checkpoint::{
};
use fabro_petri::controls::RunControls;
use fabro_petri::engine::{self, Execution, RunRequest, RunStatus};
use fabro_petri::hooks::HooksSpec;
use fabro_petri::hooks::{HooksSpec, Publication, RunPublisher};
use fabro_petri::platform_records::PlatformRecords;
use fabro_petri::providers::{DaytonaCredentials, SandboxProviderConfig};
use fabro_petri::recovery::{self, Recovery, RecoveryRequest};
@ -96,6 +96,8 @@ struct Harness {
/// Treat the local provider's workspaces as a sandbox's: `git` runs
/// through the scope's environment and checkpoints leave as bundles.
sandboxed: bool,
/// What a successful run's work does when it ends.
publisher: Option<Arc<dyn RunPublisher>>,
_root: tempfile::TempDir,
}
@ -121,6 +123,7 @@ impl Harness {
artifacts: Vec::new(),
source: None,
sandboxed: false,
publisher: None,
_root: root,
}
}
@ -136,6 +139,7 @@ impl Harness {
test_gates: None,
artifact_writer: Arc::new(StoreArtifactWriter::new(self.artifact_store.clone())),
source: self.source.clone(),
publisher: self.publisher.clone(),
}
}
@ -1239,3 +1243,104 @@ async fn an_unavailable_revision_fails_the_checkout() {
.expect_err("the branch does not exist");
assert!(error.to_string().contains("git fetch failed"), "{error}");
}
/// A publisher that records what it was handed and answers as told.
struct RecordingPublisher {
published: std::sync::Mutex<Vec<Publication>>,
fail: Option<String>,
}
#[async_trait::async_trait]
impl RunPublisher for RecordingPublisher {
async fn publish(&self, publication: &Publication) -> Result<(), String> {
self.published
.lock()
.expect("the ledger locks")
.push(publication.clone());
self.fail.clone().map_or(Ok(()), Err)
}
}
/// A run checked out from `origin` on the local provider, whose one stage
/// has the node attributes `attributes`, published through `publisher`.
async fn published_run(
attributes: &str,
publisher: &Arc<RecordingPublisher>,
) -> (Harness, engine::RunOutcome) {
let mut harness = Harness::new();
let (origin, _) = upstream(&harness.run_dir.with_file_name("upstream"), 2).await;
harness.source = Some(RunSource {
origin: format!("file://{}", origin.display()),
revision: SourceRevision::Branch("main".to_string()),
branch: "main".to_string(),
depth: Some(1),
credential: None,
});
harness.publisher = Some(Arc::clone(publisher) as Arc<dyn RunPublisher>);
let workflow = workflow(
&format!(" edit [shape=parallelogram, {attributes}]"),
" start -> edit -> exit",
);
let outcome = harness
.run_on(SandboxProviderKind::LOCAL, &workflow, SETTINGS)
.await;
(harness, outcome)
}
/// A successful run hands its publisher the run branch, the commit it ends
/// on (held by the snapshot repository) and its patch, before the run ends.
#[tokio::test]
async fn a_successful_run_is_published_with_its_branch_head_and_patch() {
let publisher = Arc::new(RecordingPublisher {
published: std::sync::Mutex::default(),
fail: None,
});
let (harness, outcome) = published_run("script=\"echo edited >> README.md\"", &publisher).await;
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
assert!(!outcome.publish_failed);
let published = publisher.published.lock().unwrap().clone();
assert_eq!(published.len(), 1, "published once");
let publication = &published[0];
assert_eq!(
publication.run_branch,
format!("fabro/run/{}", harness.run_id)
);
let (_, last) = harness.checkpoints().last().cloned().expect("a checkpoint");
assert_eq!(publication.head_sha, last);
assert_eq!(
git(&publication.snapshot_repository, &["cat-file", "-t", &last]).await,
"commit"
);
assert!(
publication.patch.contains("+edited"),
"{}",
publication.patch
);
}
/// A publication that fails fails the run, with the reason.
#[tokio::test]
async fn a_failed_publication_fails_the_run() {
let publisher = Arc::new(RecordingPublisher {
published: std::sync::Mutex::default(),
fail: Some("the push was rejected".to_string()),
});
let (_, outcome) = published_run("script=\"echo edited >> README.md\"", &publisher).await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
assert!(outcome.publish_failed);
assert_eq!(outcome.failure.as_deref(), Some("the push was rejected"));
}
/// A run that fails (here, at a goal gate) is not published.
#[tokio::test]
async fn a_failed_run_is_not_published() {
let publisher = Arc::new(RecordingPublisher {
published: std::sync::Mutex::default(),
fail: None,
});
let (_, outcome) = published_run("script=\"exit 3\", goal_gate=true", &publisher).await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
assert!(!outcome.publish_failed);
assert!(publisher.published.lock().unwrap().is_empty());
}

View file

@ -26,9 +26,6 @@ impl EnvVars {
pub const FABRO_PUSH_CRED_REFRESH_INTERVAL_SECONDS: &'static str =
"FABRO_PUSH_CRED_REFRESH_INTERVAL_SECONDS";
pub const FABRO_QUIET: &'static str = "FABRO_QUIET";
/// The read-only credential a run's worker fetches its GitHub target
/// with: the base64 of `username:password`, scrubbed at worker startup.
pub const FABRO_RUN_GIT_CREDENTIAL: &'static str = "FABRO_RUN_GIT_CREDENTIAL";
pub const FABRO_SERVER: &'static str = "FABRO_SERVER";
pub const FABRO_SERVER_MAX_CONCURRENT_RUNS: &'static str = "FABRO_SERVER_MAX_CONCURRENT_RUNS";
pub const FABRO_SLACK_APP_TOKEN: &'static str = "FABRO_SLACK_APP_TOKEN";
@ -231,7 +228,6 @@ mod tests {
EnvVars::FABRO_PUSH_CRED_REFRESH_AHEAD,
EnvVars::FABRO_PUSH_CRED_REFRESH_INTERVAL_SECONDS,
EnvVars::FABRO_QUIET,
EnvVars::FABRO_RUN_GIT_CREDENTIAL,
EnvVars::FABRO_SERVER,
EnvVars::FABRO_SERVER_MAX_CONCURRENT_RUNS,
EnvVars::FABRO_SLACK_APP_TOKEN,