diff --git a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs index 2a40842bb..446fbf2e5 100644 --- a/lib/apps/fabro-cli/src/commands/run/petri_worker.rs +++ b/lib/apps/fabro-cli/src/commands/run/petri_worker.rs @@ -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 { /// 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>, run_state: &RunProjection, fabro_home: Option, run_tools: Option, ) -> Result { - let catalog = - command_context::load_cli_catalog().context("failed to build worker LLM catalog")?; let credentials: Arc = Arc::new(VaultCredentialSource::new(Arc::clone(vault))); let ready = readiness(catalog.enabled_providers(), credentials.as_ref()).await; diff --git a/lib/apps/fabro-cli/src/commands/run/publish.rs b/lib/apps/fabro-cli/src/commands/run/publish.rs index 0533095d0..0d716803a 100644 --- a/lib/apps/fabro-cli/src/commands/run/publish.rs +++ b/lib/apps/fabro-cli/src/commands/run/publish.rs @@ -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::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, credentials: Option, pull_request: Option, llm_source: Arc, @@ -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 { - 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 { 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); diff --git a/lib/components/fabro-petri/src/checkpoint.rs b/lib/components/fabro-petri/src/checkpoint.rs index bfebe1b82..0c02c0a28 100644 --- a/lib/components/fabro-petri/src/checkpoint.rs +++ b/lib/components/fabro-petri/src/checkpoint.rs @@ -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> { diff --git a/lib/components/fabro-petri/src/hooks.rs b/lib/components/fabro-petri/src/hooks.rs index 7ef25af69..999176e5b 100644 --- a/lib/components/fabro-petri/src/hooks.rs +++ b/lib/components/fabro-petri/src/hooks.rs @@ -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) { - 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 } diff --git a/lib/components/fabro-petri/src/source.rs b/lib/components/fabro-petri/src/source.rs index 1c79ccb8c..3e6c99d07 100644 --- a/lib/components/fabro-petri/src/source.rs +++ b/lib/components/fabro-petri/src/source.rs @@ -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. diff --git a/lib/components/fabro-petri/tests/hooks.rs b/lib/components/fabro-petri/tests/hooks.rs index 634342815..5fe0ab2e7 100644 --- a/lib/components/fabro-petri/tests/hooks.rs +++ b/lib/components/fabro-petri/tests/hooks.rs @@ -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, } +impl RecordingPublisher { + fn new(fail: Option<&str>) -> Arc { + 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) -> 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); 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);