Simplify GitHub checkout and run publication

- Share one credential header helper between the in-sandbox fetch and
  the run branch push
- Keep only the target branch, goal, and model on the publisher instead
  of a full run spec copy
- Load the worker's LLM catalog once, and mint the read token only when
  the run checks something out
- Pass the source explicitly to checkpoint fetch helpers, dropping
  unreachable branches, and reuse has_object in has_commit
- Move the run patch into the publication instead of cloning it, and
  build it only when a publisher exists
- Add test fixture helpers for file sources and recording publishers

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-09-30 14:50:00 -04:00
parent 0250b40586
commit 886065957c
6 changed files with 121 additions and 134 deletions

View file

@ -71,6 +71,7 @@ use fabro_auth::VaultCredentialSource;
use fabro_client::{Client, ServerTarget};
use fabro_interview::{ControlInterviewer, WorkerControlMessage, WorkerControlOutcome};
use fabro_llm::credentials::{CredentialProvider, readiness};
use fabro_llm::lithos_catalog::Catalog;
use fabro_petri::artifacts::ClientArtifactWriter;
use fabro_petri::blobs::ClientBlobs;
use fabro_petri::controls::{RunControls, SteerError};
@ -168,7 +169,10 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
let vault = runner::load_worker_vault(worker.storage_dir).await?;
let secrets = VaultSecrets::from_vault(&*vault.read().await);
let run_tools = run_tool_services(&worker);
let catalog =
command_context::load_cli_catalog().context("failed to build worker LLM catalog")?;
let runtime = runtime_spec(
catalog.clone(),
&vault,
&worker.run_state,
worker.fabro_home.clone(),
@ -211,20 +215,21 @@ pub(super) async fn execute(worker: PetriWorker<'_>) -> Result<()> {
None
}
};
let source = RunSource::for_run(
let mut source = RunSource::for_run(
worker.run_state.spec.target.as_ref(),
&worker.run_state.spec.settings.run,
publish::source_credential(&worker.run_state.spec, github.as_ref()).await,
None,
);
if let Some(source) = &mut source {
source.credential =
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::new(catalog),
Arc::clone(&records),
worker.client.clone_for_reuse(),
)
@ -639,13 +644,12 @@ fn run_tool_services(worker: &PetriWorker<'_>) -> Option<FabroRunToolServices> {
/// the providers whose credentials resolve, the run's mode, the Fabro
/// home the server named, and the run tools when the run has them.
async fn runtime_spec(
catalog: Catalog,
vault: &Arc<AsyncRwLock<Vault>>,
run_state: &RunProjection,
fabro_home: Option<PathBuf>,
run_tools: Option<FabroRunToolServices>,
) -> Result<RuntimeSpec> {
let catalog =
command_context::load_cli_catalog().context("failed to build worker LLM catalog")?;
let credentials: Arc<dyn CredentialProvider> =
Arc::new(VaultCredentialSource::new(Arc::clone(vault)));
let ready = readiness(catalog.enabled_providers(), credentials.as_ref()).await;

View file

@ -98,7 +98,7 @@ pub(super) async fn source_credential(
)
.await
{
Ok(credentials) => SourceCredential::from_encoded(encode(&credentials)),
Ok(credentials) => encode(&credentials),
Err(err) => {
warn!(repository = %repository, error = %err, "no read credential for the run's repository; it is fetched anonymously");
None
@ -106,20 +106,24 @@ pub(super) async fn source_credential(
}
}
fn encode(credentials: &GitCloneCredentials) -> String {
BASE64_STANDARD.encode(format!(
fn encode(credentials: &GitCloneCredentials) -> Option<SourceCredential> {
SourceCredential::from_encoded(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,
/// The branch the target names: the pull request's base.
base_branch: String,
goal: String,
/// The run's model, when its settings name one.
model: Option<String>,
credentials: Option<GitHubCredentials>,
pull_request: Option<PullRequestSettings>,
llm_source: Arc<dyn CredentialProvider>,
@ -147,11 +151,16 @@ impl GitHubPublisher {
{
return None;
}
let Some(RunTarget::Git(target)) = spec.target.as_ref() else {
return None;
};
let repository = repository(spec)?;
Some(Self {
run_id,
spec: spec.clone(),
repository,
base_branch: target.branch.clone(),
goal: spec.graph.goal.clone(),
model: settings.model.name.clone(),
credentials,
pull_request: settings
.pull_request
@ -167,8 +176,8 @@ impl GitHubPublisher {
/// 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);
if let Some(model) = &self.model {
return Some(model.clone());
}
let ready = readiness(self.catalog.enabled_providers(), self.llm_source.as_ref()).await;
self.catalog
@ -182,9 +191,6 @@ impl GitHubPublisher {
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
@ -195,10 +201,10 @@ impl GitHubPublisher {
let created = pull_request::open_pull_request(OpenPullRequestRequest {
github: context,
origin_url: &origin_url,
base_branch: &target.branch,
base_branch: &self.base_branch,
head_branch: &publication.run_branch,
expected_head_sha: &publication.head_sha,
goal: &self.spec.graph.goal,
goal: &self.goal,
diff: &publication.patch,
model: &model,
draft: settings.draft,
@ -267,24 +273,16 @@ async fn push(
) -> 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 env = encode(credentials)
.map(|credential| credential.header_env(&url))
.unwrap_or_default();
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?;
let output = git_push(&publication.snapshot_repository, &url, &refspec, &env).await?;
if output.status.success() {
return Ok(());
}
@ -307,17 +305,19 @@ async fn push(
))
}
/// `git` in `directory` with `env` added, non-interactive, bounded.
async fn git(
/// `git push` of `refspec` to `url` from `directory` with `env` added,
/// non-interactive, bounded.
async fn git_push(
directory: &Path,
args: &[&str],
env: &[(&str, String)],
url: &str,
refspec: &str,
env: &[(String, String)],
) -> Result<std::process::Output, String> {
let mut command = Command::new("git");
command
.args(args)
.args(["push", "--quiet", url, refspec])
.current_dir(directory)
.envs(env.iter().map(|(key, value)| (*key, value)))
.envs(env.iter().map(|(key, value)| (key, value)))
.env("GIT_TERMINAL_PROMPT", "0")
.stdin(std::process::Stdio::null())
.kill_on_drop(true);

View file

@ -464,7 +464,7 @@ impl RunWorkspaces {
&source.origin,
])
.await?;
self.fetch_source(site, "origin", &source.revision.refspec())
self.fetch_source(source, site, "origin", &source.revision.refspec())
.await?;
self.git(site, "checkout", &[
"checkout",
@ -484,20 +484,23 @@ impl RunWorkspaces {
.await?;
}
let sha = self.git(site, "rev-parse", &["rev-parse", "HEAD"]).await?;
self.seed_snapshot(workspace, &sha).await?;
self.seed_snapshot(source, workspace, &sha).await?;
Ok(Some(sha))
}
/// Put the commit a workspace was checked out at into its snapshot
/// repository, under [`SOURCE_REF`], fetched from the origin at the run's
/// depth.
async fn seed_snapshot(&self, workspace: &str, sha: &str) -> Result<(), CheckpointError> {
let Some(source) = &self.source else {
return Ok(());
};
async fn seed_snapshot(
&self,
source: &RunSource,
workspace: &str,
sha: &str,
) -> Result<(), CheckpointError> {
let repository = Site::Host(self.ensure_snapshot_repository(workspace).await?);
if !self.has_object(&repository, sha).await? {
self.fetch_source(&repository, &source.origin, sha).await?;
self.fetch_source(source, &repository, &source.origin, sha)
.await?;
}
self.git(&repository, "update-ref", &["update-ref", SOURCE_REF, sha])
.await?;
@ -521,20 +524,14 @@ impl RunWorkspaces {
}
/// Fetch `refspec` from `remote` (a remote name, or the origin's URL) at
/// `site`, with the source's credential and depth.
/// `site`, with `source`'s credential and depth.
async fn fetch_source(
&self,
source: &RunSource,
site: &Site,
remote: &str,
refspec: &str,
) -> Result<(), CheckpointError> {
let Some(source) = &self.source else {
return Err(CheckpointError::Command {
action: "fetch".to_string(),
status: "no source".to_string(),
detail: "the run has no repository to fetch from".to_string(),
});
};
let depth = source.depth_arg();
let mut args = vec!["fetch", "-q", "--no-tags"];
if let Some(depth) = depth.as_deref() {
@ -796,14 +793,7 @@ impl RunWorkspaces {
{
return Ok(false);
}
Ok(self
.git_status(site, "cat-file", &[
"cat-file",
"-e",
&format!("{sha}^{{commit}}"),
])
.await?
.is_some())
self.has_object(site, sha).await
}
/// Bring the workspace back to `sha`: tracked files reset, untracked
@ -979,7 +969,7 @@ impl RunWorkspaces {
])
.await?;
}
self.fetch_source(site, "origin", base).await
self.fetch_source(source, site, "origin", base).await
}
async fn verify_restored(&self, site: &Site, sha: &str) -> Result<(), CheckpointError> {

View file

@ -1234,12 +1234,16 @@ 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 publication = branch
.run_branch
.clone()
.filter(|_| self.publisher.is_some())
.map(|run_branch| Publication {
run_branch,
head_sha: head_sha.clone(),
snapshot_repository: self.workspaces.snapshot_repository(&workspace),
patch: diff.patch,
});
let record = PlatformRecord::RunDiff(RunDiffRecord {
base_sha: Some(base_sha),
head_sha: Some(head_sha),
@ -1265,15 +1269,8 @@ impl FabroHooks {
/// 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 {
async fn publish(&self, publisher: &dyn RunPublisher, publication: &Publication) {
match publisher.publish(publication).await {
Ok(()) => info!(
run_id = %self.run_id,
branch = publication.run_branch,
@ -1507,8 +1504,10 @@ impl ExecutionHooks for FabroHooks {
None
}
};
if finished.status == RunStatus::Success && self.checkpoint_failure().is_none() {
self.publish(publication).await;
if let (Some(publisher), Some(publication)) = (&self.publisher, &publication) {
if finished.status == RunStatus::Success && self.checkpoint_failure().is_none() {
self.publish(publisher.as_ref(), publication).await;
}
}
self.inner.run_finished(context, finished).await
}

View file

@ -64,6 +64,23 @@ impl SourceCredential {
pub fn encoded(&self) -> &str {
&self.0
}
/// The environment that has `git` present this credential as an
/// `Authorization` header to `url` alone.
#[must_use]
pub fn header_env(&self, url: &str) -> Vec<(String, String)> {
vec![
("GIT_CONFIG_COUNT".to_string(), "1".to_string()),
(
"GIT_CONFIG_KEY_0".to_string(),
format!("http.{url}.extraheader"),
),
(
"GIT_CONFIG_VALUE_0".to_string(),
format!("AUTHORIZATION: basic {}", self.0),
),
]
}
}
impl fmt::Debug for SourceCredential {
@ -142,20 +159,10 @@ impl RunSource {
/// the credential as an `Authorization` header for the origin alone.
#[must_use]
pub fn fetch_env(&self) -> Vec<(String, String)> {
let Some(credential) = &self.credential else {
return Vec::new();
};
vec![
("GIT_CONFIG_COUNT".to_string(), "1".to_string()),
(
"GIT_CONFIG_KEY_0".to_string(),
format!("http.{}.extraheader", self.origin),
),
(
"GIT_CONFIG_VALUE_0".to_string(),
format!("AUTHORIZATION: basic {}", credential.encoded()),
),
]
self.credential
.as_ref()
.map(|credential| credential.header_env(&self.origin))
.unwrap_or_default()
}
/// The `--depth` argument of a fetch, when the history is limited.

View file

@ -1121,13 +1121,7 @@ async fn a_git_source_is_checked_out_shallow_and_checkpoints_build_on_its_commit
let mut harness = Harness::new();
let (origin, head) = upstream(&harness.run_dir.with_file_name("upstream"), 3).await;
harness.sandboxed = sandboxed;
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.source = Some(file_source(&origin, "main", Some(1)));
let workflow = workflow(
" edit [shape=parallelogram, script=\"test \\\"$(cat README.md)\\\" = 'revision 3' && test \\\"$(git rev-parse --is-shallow-repository)\\\" = true && git rev-parse origin/main && echo edited >> README.md && git -c user.name=Agent -c \
user.email=agent@example.com commit -q -am 'agent edit' && echo uncommitted > \
@ -1194,13 +1188,7 @@ async fn a_prepared_workspace_is_left_as_it_is() {
GitAuthor::default(),
&RunCheckpointSettings::default(),
)
.with_source(Some(RunSource {
origin: format!("file://{}", origin.display()),
revision: SourceRevision::Branch("main".to_string()),
branch: "main".to_string(),
depth: None,
credential: None,
}));
.with_source(Some(file_source(&origin, "main", None)));
let path = root.path().join("prepared");
fs::create_dir_all(&path)
.await
@ -1228,13 +1216,7 @@ async fn an_unavailable_revision_fails_the_checkout() {
GitAuthor::default(),
&RunCheckpointSettings::default(),
)
.with_source(Some(RunSource {
origin: format!("file://{}", origin.display()),
revision: SourceRevision::Branch("missing".to_string()),
branch: "missing".to_string(),
depth: Some(1),
credential: None,
}));
.with_source(Some(file_source(&origin, "missing", Some(1))));
let path = root.path().join("fresh");
let site = Site::Host(path);
let error = workspaces
@ -1250,6 +1232,26 @@ struct RecordingPublisher {
fail: Option<String>,
}
impl RecordingPublisher {
fn new(fail: Option<&str>) -> Arc<Self> {
Arc::new(Self {
published: std::sync::Mutex::default(),
fail: fail.map(str::to_string),
})
}
}
/// The source of a run checked out from the local repository `origin`.
fn file_source(origin: &Path, branch: &str, depth: Option<u32>) -> RunSource {
RunSource {
origin: format!("file://{}", origin.display()),
revision: SourceRevision::Branch(branch.to_string()),
branch: branch.to_string(),
depth,
credential: None,
}
}
#[async_trait::async_trait]
impl RunPublisher for RecordingPublisher {
async fn publish(&self, publication: &Publication) -> Result<(), String> {
@ -1269,13 +1271,7 @@ async fn published_run(
) -> (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.source = Some(file_source(&origin, "main", Some(1)));
harness.publisher = Some(Arc::clone(publisher) as Arc<dyn RunPublisher>);
let workflow = workflow(
&format!(" edit [shape=parallelogram, {attributes}]"),
@ -1291,10 +1287,7 @@ async fn published_run(
/// 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 publisher = RecordingPublisher::new(None);
let (harness, outcome) = published_run("script=\"echo edited >> README.md\"", &publisher).await;
assert_eq!(outcome.status, RunStatus::Success, "{outcome:?}");
assert!(!outcome.publish_failed);
@ -1322,10 +1315,7 @@ async fn a_successful_run_is_published_with_its_branch_head_and_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 publisher = RecordingPublisher::new(Some("the push was rejected"));
let (_, outcome) = published_run("script=\"echo edited >> README.md\"", &publisher).await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
assert!(outcome.publish_failed);
@ -1335,10 +1325,7 @@ async fn a_failed_publication_fails_the_run() {
/// 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 publisher = RecordingPublisher::new(None);
let (_, outcome) = published_run("script=\"exit 3\", goal_gate=true", &publisher).await;
assert_eq!(outcome.status, RunStatus::Failed, "{outcome:?}");
assert!(!outcome.publish_failed);