mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-09 03:20:56 +00:00
fabro(01KSE2PAVXD56N4TWNK4T5H5VA): simplify_opus (succeeded)
Fabro-Run: 01KSE2PAVXD56N4TWNK4T5H5VA
Fabro-Completed: 8
Fabro-Checkpoint: eea6420ea4
⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
parent
afdd4900fa
commit
35fbcaa791
9 changed files with 256 additions and 253 deletions
1
Cargo.lock
generated
1
Cargo.lock
generated
|
|
@ -1769,7 +1769,6 @@ dependencies = [
|
|||
name = "fabro-automation"
|
||||
version = "0.243.0-nightly.1"
|
||||
dependencies = [
|
||||
"chrono",
|
||||
"croner",
|
||||
"hex",
|
||||
"serde",
|
||||
|
|
|
|||
|
|
@ -27,7 +27,7 @@ fn automation_api_reuses_domain_types() {
|
|||
fn automation_json_matches_openapi_shape() {
|
||||
let automation = Automation {
|
||||
id: "nightly-deps".parse().unwrap(),
|
||||
revision: AutomationRevision::from_str("abc123").unwrap(),
|
||||
revision: AutomationRevision::from_raw("abc123"),
|
||||
name: "Nightly dependency update".to_string(),
|
||||
description: Some("Open a PR for dependency updates.".to_string()),
|
||||
enabled: true,
|
||||
|
|
|
|||
|
|
@ -13,11 +13,11 @@ doctest = false
|
|||
workspace = true
|
||||
|
||||
[dependencies]
|
||||
chrono.workspace = true
|
||||
croner = "3.0.1"
|
||||
hex.workspace = true
|
||||
serde.workspace = true
|
||||
sha2.workspace = true
|
||||
tempfile = "3"
|
||||
thiserror.workspace = true
|
||||
tokio.workspace = true
|
||||
toml.workspace = true
|
||||
|
|
|
|||
|
|
@ -37,8 +37,6 @@ pub enum AutomationStoreError {
|
|||
NotFound(AutomationId),
|
||||
#[error("automation already exists: {0}")]
|
||||
AlreadyExists(AutomationId),
|
||||
#[error("missing automation revision")]
|
||||
MissingRevision,
|
||||
#[error("automation revision mismatch")]
|
||||
RevisionMismatch {
|
||||
expected: AutomationRevision,
|
||||
|
|
@ -51,6 +49,8 @@ pub enum AutomationStoreError {
|
|||
path: PathBuf,
|
||||
source: TomlDeError,
|
||||
},
|
||||
#[error("failed to serialize automation TOML: {0}")]
|
||||
Serialize(String),
|
||||
#[error("I/O error at {}: {source}", path.display())]
|
||||
Io {
|
||||
path: PathBuf,
|
||||
|
|
|
|||
|
|
@ -1,7 +1,6 @@
|
|||
pub mod error;
|
||||
pub mod id;
|
||||
pub mod model;
|
||||
|
||||
mod error;
|
||||
mod id;
|
||||
mod model;
|
||||
mod store;
|
||||
|
||||
pub use error::{AutomationStoreError, AutomationValidationError};
|
||||
|
|
|
|||
|
|
@ -127,6 +127,14 @@ impl AutomationRevision {
|
|||
Self(hex::encode(Sha256::digest(bytes)))
|
||||
}
|
||||
|
||||
/// Wrap a client-supplied revision string (e.g. from an `If-Match`
|
||||
/// header). The value is compared bytewise against a stored revision; no
|
||||
/// validation is performed here.
|
||||
#[must_use]
|
||||
pub fn from_raw(value: impl Into<String>) -> Self {
|
||||
Self(value.into())
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn as_str(&self) -> &str {
|
||||
&self.0
|
||||
|
|
@ -165,28 +173,28 @@ impl Automation {
|
|||
pub fn from_toml_bytes(id: AutomationId, bytes: &[u8]) -> Result<Self, TomlDeError> {
|
||||
let source = std::str::from_utf8(bytes).map_err(TomlDeError::custom)?;
|
||||
let persisted = toml::from_str::<PersistedAutomation>(source)?;
|
||||
// serde has already validated newtypes and trigger shapes. This call
|
||||
// checks cross-field invariants.
|
||||
persisted
|
||||
.into_automation(id, AutomationRevision::from_bytes(bytes))
|
||||
.map_err(TomlDeError::custom)
|
||||
let revision = AutomationRevision::from_bytes(bytes);
|
||||
Self::assemble(id, revision, persisted.into_replace()).map_err(TomlDeError::custom)
|
||||
}
|
||||
|
||||
pub fn from_draft(
|
||||
draft: AutomationDraft,
|
||||
/// Build, validate, and assign a revision to an `Automation` in one
|
||||
/// step. Used by the store immediately after persisting canonical TOML
|
||||
/// bytes so the in-memory revision always matches what is on disk.
|
||||
pub(crate) fn assemble(
|
||||
id: AutomationId,
|
||||
revision: AutomationRevision,
|
||||
replace: AutomationReplace,
|
||||
) -> Result<Self, AutomationValidationError> {
|
||||
let automation = Self {
|
||||
id: draft.id,
|
||||
validate_common(&replace.name, &replace.triggers)?;
|
||||
Ok(Self {
|
||||
id,
|
||||
revision,
|
||||
name: draft.name,
|
||||
description: draft.description,
|
||||
enabled: draft.enabled.unwrap_or(true),
|
||||
target: draft.target,
|
||||
triggers: draft.triggers,
|
||||
};
|
||||
automation.validate()?;
|
||||
Ok(automation)
|
||||
name: replace.name,
|
||||
description: replace.description,
|
||||
enabled: replace.enabled,
|
||||
target: replace.target,
|
||||
triggers: replace.triggers,
|
||||
})
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
|
|
@ -200,15 +208,6 @@ impl Automation {
|
|||
}
|
||||
}
|
||||
|
||||
pub fn to_toml_bytes(&self) -> Result<Vec<u8>, TomlEditSerError> {
|
||||
let persisted = PersistedAutomation::from(self);
|
||||
to_document(&persisted).map(|document| document.to_string().into_bytes())
|
||||
}
|
||||
|
||||
pub fn validate(&self) -> Result<(), AutomationValidationError> {
|
||||
validate_common(&self.name, &self.triggers)
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub fn api_trigger(&self) -> Option<&ApiTrigger> {
|
||||
self.triggers.iter().find_map(|trigger| match trigger {
|
||||
|
|
@ -218,23 +217,29 @@ impl Automation {
|
|||
}
|
||||
}
|
||||
|
||||
impl AutomationReplace {
|
||||
pub(crate) fn into_automation(
|
||||
self,
|
||||
id: AutomationId,
|
||||
revision: AutomationRevision,
|
||||
) -> Result<Automation, AutomationValidationError> {
|
||||
let automation = Automation {
|
||||
id,
|
||||
revision,
|
||||
name: self.name,
|
||||
impl AutomationDraft {
|
||||
/// Drop the `id` (which becomes the storage filename) and surface the
|
||||
/// remaining fields in the canonical replace shape, applying the
|
||||
/// `enabled` default.
|
||||
#[must_use]
|
||||
pub fn into_replace(self) -> AutomationReplace {
|
||||
AutomationReplace {
|
||||
name: self.name,
|
||||
description: self.description,
|
||||
enabled: self.enabled,
|
||||
target: self.target,
|
||||
triggers: self.triggers,
|
||||
};
|
||||
automation.validate()?;
|
||||
Ok(automation)
|
||||
enabled: self.enabled.unwrap_or(true),
|
||||
target: self.target,
|
||||
triggers: self.triggers,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl AutomationReplace {
|
||||
/// Serialize this replace value into canonical TOML bytes. The
|
||||
/// representation matches `PersistedAutomation` so on-disk and in-memory
|
||||
/// shapes stay aligned without an extra clone.
|
||||
pub(crate) fn to_toml_bytes(&self) -> Result<Vec<u8>, TomlEditSerError> {
|
||||
to_document(&PersistedAutomationRef::from(self))
|
||||
.map(|document| document.to_string().into_bytes())
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -253,33 +258,36 @@ impl AutomationPatch {
|
|||
}
|
||||
|
||||
impl PersistedAutomation {
|
||||
pub(crate) fn into_automation(
|
||||
self,
|
||||
id: AutomationId,
|
||||
revision: AutomationRevision,
|
||||
) -> Result<Automation, AutomationValidationError> {
|
||||
let automation = Automation {
|
||||
id,
|
||||
revision,
|
||||
name: self.name,
|
||||
fn into_replace(self) -> AutomationReplace {
|
||||
AutomationReplace {
|
||||
name: self.name,
|
||||
description: self.description,
|
||||
enabled: self.enabled,
|
||||
target: self.target,
|
||||
triggers: self.triggers,
|
||||
};
|
||||
automation.validate()?;
|
||||
Ok(automation)
|
||||
enabled: self.enabled,
|
||||
target: self.target,
|
||||
triggers: self.triggers,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<&Automation> for PersistedAutomation {
|
||||
fn from(value: &Automation) -> Self {
|
||||
#[derive(Debug, Serialize)]
|
||||
struct PersistedAutomationRef<'a> {
|
||||
name: &'a str,
|
||||
#[serde(skip_serializing_if = "Option::is_none")]
|
||||
description: Option<&'a str>,
|
||||
enabled: bool,
|
||||
target: &'a AutomationTarget,
|
||||
#[serde(default, skip_serializing_if = "<[_]>::is_empty")]
|
||||
triggers: &'a [AutomationTrigger],
|
||||
}
|
||||
|
||||
impl<'a> From<&'a AutomationReplace> for PersistedAutomationRef<'a> {
|
||||
fn from(value: &'a AutomationReplace) -> Self {
|
||||
Self {
|
||||
name: value.name.clone(),
|
||||
description: value.description.clone(),
|
||||
name: &value.name,
|
||||
description: value.description.as_deref(),
|
||||
enabled: value.enabled,
|
||||
target: value.target.clone(),
|
||||
triggers: value.triggers.clone(),
|
||||
target: &value.target,
|
||||
triggers: &value.triggers,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -458,14 +466,6 @@ impl fmt::Display for AutomationRevision {
|
|||
}
|
||||
}
|
||||
|
||||
impl FromStr for AutomationRevision {
|
||||
type Err = AutomationValidationError;
|
||||
|
||||
fn from_str(value: &str) -> Result<Self, Self::Err> {
|
||||
Ok(Self(value.to_string()))
|
||||
}
|
||||
}
|
||||
|
||||
impl Serialize for AutomationRevision {
|
||||
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
|
||||
where
|
||||
|
|
@ -680,7 +680,14 @@ expression = "0 3 * * *"
|
|||
"#,
|
||||
))
|
||||
.expect("draft should deserialize");
|
||||
assert!(Automation::from_draft(draft, AutomationRevision::from_bytes(b"")).is_err());
|
||||
assert!(
|
||||
Automation::assemble(
|
||||
draft.id.clone(),
|
||||
AutomationRevision::from_bytes(b""),
|
||||
draft.into_replace(),
|
||||
)
|
||||
.is_err()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -698,7 +705,14 @@ type = "api"
|
|||
"#,
|
||||
))
|
||||
.expect("draft should deserialize");
|
||||
assert!(Automation::from_draft(draft, AutomationRevision::from_bytes(b"")).is_err());
|
||||
assert!(
|
||||
Automation::assemble(
|
||||
draft.id.clone(),
|
||||
AutomationRevision::from_bytes(b""),
|
||||
draft.into_replace(),
|
||||
)
|
||||
.is_err()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -725,7 +739,14 @@ expression = "* * * * * *"
|
|||
"#,
|
||||
))
|
||||
.expect("draft should deserialize");
|
||||
assert!(Automation::from_draft(draft, AutomationRevision::from_bytes(b"")).is_err());
|
||||
assert!(
|
||||
Automation::assemble(
|
||||
draft.id.clone(),
|
||||
AutomationRevision::from_bytes(b""),
|
||||
draft.into_replace(),
|
||||
)
|
||||
.is_err()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
|
|||
|
|
@ -1,11 +1,14 @@
|
|||
use std::collections::BTreeMap;
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
#[expect(
|
||||
clippy::disallowed_types,
|
||||
reason = "atomic_write writes through spawn_blocking + NamedTempFile, which only exposes std::io::Write."
|
||||
)]
|
||||
use std::io::Write as _;
|
||||
use std::path::PathBuf;
|
||||
|
||||
use tokio::fs::{self, OpenOptions};
|
||||
use tokio::io::AsyncWriteExt as _;
|
||||
use tempfile::NamedTempFile;
|
||||
use tokio::sync::RwLock;
|
||||
use tokio::{fs, task};
|
||||
|
||||
use crate::error::{AutomationStoreError, AutomationValidationError};
|
||||
use crate::id::AutomationId;
|
||||
|
|
@ -114,9 +117,7 @@ impl AutomationStore {
|
|||
if items.contains_key(&id) {
|
||||
return Err(AutomationStoreError::AlreadyExists(id));
|
||||
}
|
||||
|
||||
let automation = Automation::from_draft(draft, AutomationRevision::from_bytes(b""))?;
|
||||
let automation = self.persist_with_revision(automation).await?;
|
||||
let automation = self.persist(id.clone(), draft.into_replace()).await?;
|
||||
items.insert(id, automation.clone());
|
||||
Ok(automation)
|
||||
}
|
||||
|
|
@ -132,9 +133,7 @@ impl AutomationStore {
|
|||
.get(id)
|
||||
.ok_or_else(|| AutomationStoreError::NotFound(id.clone()))?;
|
||||
ensure_revision(current, expected)?;
|
||||
|
||||
let automation = draft.into_automation(id.clone(), AutomationRevision::from_bytes(b""))?;
|
||||
let automation = self.persist_with_revision(automation).await?;
|
||||
let automation = self.persist(id.clone(), draft).await?;
|
||||
items.insert(id.clone(), automation.clone());
|
||||
Ok(automation)
|
||||
}
|
||||
|
|
@ -150,10 +149,8 @@ impl AutomationStore {
|
|||
.get(id)
|
||||
.ok_or_else(|| AutomationStoreError::NotFound(id.clone()))?;
|
||||
ensure_revision(current, expected)?;
|
||||
|
||||
let draft = patch.apply_to(current);
|
||||
let automation = draft.into_automation(id.clone(), AutomationRevision::from_bytes(b""))?;
|
||||
let automation = self.persist_with_revision(automation).await?;
|
||||
let replace = patch.apply_to(current);
|
||||
let automation = self.persist(id.clone(), replace).await?;
|
||||
items.insert(id.clone(), automation.clone());
|
||||
Ok(automation)
|
||||
}
|
||||
|
|
@ -179,19 +176,22 @@ impl AutomationStore {
|
|||
Ok(())
|
||||
}
|
||||
|
||||
async fn persist_with_revision(
|
||||
/// Validate the replace value, render canonical TOML, write atomically,
|
||||
/// and return the assembled `Automation` whose revision matches the
|
||||
/// bytes that landed on disk.
|
||||
async fn persist(
|
||||
&self,
|
||||
automation: Automation,
|
||||
id: AutomationId,
|
||||
replace: AutomationReplace,
|
||||
) -> Result<Automation, AutomationStoreError> {
|
||||
let bytes = automation
|
||||
let bytes = replace
|
||||
.to_toml_bytes()
|
||||
.map_err(|err| AutomationValidationError::InvalidWorkflowSelector(err.to_string()))?;
|
||||
atomic_write(&self.dir, &self.path_for(&automation.id), &bytes).await?;
|
||||
.map_err(|err| AutomationStoreError::Serialize(err.to_string()))?;
|
||||
let revision = AutomationRevision::from_bytes(&bytes);
|
||||
Ok(Automation {
|
||||
revision,
|
||||
..automation
|
||||
})
|
||||
let automation = Automation::assemble(id, revision, replace)?;
|
||||
let path = self.path_for(&automation.id);
|
||||
atomic_write(&self.dir, &path, bytes).await?;
|
||||
Ok(automation)
|
||||
}
|
||||
|
||||
fn path_for(&self, id: &AutomationId) -> PathBuf {
|
||||
|
|
@ -214,53 +214,34 @@ fn ensure_revision(
|
|||
}
|
||||
|
||||
async fn atomic_write(
|
||||
dir: &Path,
|
||||
final_path: &Path,
|
||||
bytes: &[u8],
|
||||
dir: &std::path::Path,
|
||||
final_path: &std::path::Path,
|
||||
bytes: Vec<u8>,
|
||||
) -> Result<(), AutomationStoreError> {
|
||||
fs::create_dir_all(dir)
|
||||
.await
|
||||
.map_err(|err| AutomationStoreError::io(dir, err))?;
|
||||
|
||||
let temp_path = temp_path_for(dir, final_path);
|
||||
let mut file = OpenOptions::new()
|
||||
.write(true)
|
||||
.create_new(true)
|
||||
.open(&temp_path)
|
||||
.await
|
||||
.map_err(|err| AutomationStoreError::io(&temp_path, err))?;
|
||||
let write_result = async {
|
||||
file.write_all(bytes).await?;
|
||||
file.flush().await?;
|
||||
file.sync_all().await
|
||||
}
|
||||
.await;
|
||||
if let Err(err) = write_result {
|
||||
let _ = fs::remove_file(&temp_path).await;
|
||||
return Err(AutomationStoreError::io(&temp_path, err));
|
||||
}
|
||||
drop(file);
|
||||
|
||||
if let Err(err) = fs::rename(&temp_path, final_path).await {
|
||||
let _ = fs::remove_file(&temp_path).await;
|
||||
return Err(AutomationStoreError::io(final_path, err));
|
||||
}
|
||||
let dir = dir.to_path_buf();
|
||||
let final_path = final_path.to_path_buf();
|
||||
let join_dir = dir.clone();
|
||||
task::spawn_blocking(move || -> Result<(), AutomationStoreError> {
|
||||
let mut temp =
|
||||
NamedTempFile::new_in(&dir).map_err(|err| AutomationStoreError::io(&dir, err))?;
|
||||
temp.write_all(&bytes)
|
||||
.map_err(|err| AutomationStoreError::io(temp.path(), err))?;
|
||||
temp.as_file()
|
||||
.sync_all()
|
||||
.map_err(|err| AutomationStoreError::io(temp.path(), err))?;
|
||||
temp.persist(&final_path)
|
||||
.map_err(|err| AutomationStoreError::io(final_path, err.error))?;
|
||||
Ok(())
|
||||
})
|
||||
.await
|
||||
.map_err(|err| AutomationStoreError::io(join_dir, std::io::Error::other(err)))??;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn temp_path_for(dir: &Path, final_path: &Path) -> PathBuf {
|
||||
static COUNTER: AtomicU64 = AtomicU64::new(0);
|
||||
let stem = final_path
|
||||
.file_name()
|
||||
.and_then(|name| name.to_str())
|
||||
.unwrap_or("automation.toml");
|
||||
let now = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.map_or(0, |duration| duration.as_nanos());
|
||||
let counter = COUNTER.fetch_add(1, Ordering::Relaxed);
|
||||
dir.join(format!(".{stem}.{now}.{counter}.tmp"))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use tokio::fs;
|
||||
|
|
|
|||
|
|
@ -4,15 +4,16 @@ use std::time::Duration;
|
|||
|
||||
use async_trait::async_trait;
|
||||
use fabro_api::types::RunManifest;
|
||||
use fabro_automation::{AutomationId, AutomationTarget};
|
||||
use fabro_automation::AutomationTarget;
|
||||
use fabro_config::Storage;
|
||||
use fabro_redact::DisplaySafeUrl;
|
||||
use fabro_sandbox::redact::redact_auth_url;
|
||||
use fabro_types::RunId;
|
||||
use tokio::process::Command;
|
||||
use tokio::time::timeout;
|
||||
use tokio::{fs, task};
|
||||
|
||||
pub(crate) struct AutomationRunMaterializeInput {
|
||||
pub automation_id: AutomationId,
|
||||
pub target: AutomationTarget,
|
||||
pub run_id: RunId,
|
||||
pub user_settings_path: PathBuf,
|
||||
|
|
@ -31,8 +32,6 @@ pub(crate) enum AutomationRunMaterializeError {
|
|||
InvalidTarget(String),
|
||||
#[error("failed to clone automation repository: {0}")]
|
||||
CloneFailed(String),
|
||||
#[error("failed to resolve automation workflow: {0}")]
|
||||
WorkflowNotFound(String),
|
||||
#[error("failed to build run manifest: {0}")]
|
||||
Manifest(String),
|
||||
}
|
||||
|
|
@ -80,9 +79,10 @@ impl AutomationRunMaterializer for GitAutomationRunMaterializer {
|
|||
));
|
||||
}
|
||||
let sanitized_clone_url = github_clone_url(owner, repo);
|
||||
let clone_url = self
|
||||
.authenticated_clone_url(owner, repo, &sanitized_clone_url)
|
||||
.await?;
|
||||
let auth_url = self.authenticated_clone_url(&sanitized_clone_url).await?;
|
||||
let clone_url = auth_url
|
||||
.as_ref()
|
||||
.map_or_else(|| sanitized_clone_url.clone(), DisplaySafeUrl::raw_string);
|
||||
|
||||
fs::create_dir_all(&input.temp_root).await.map_err(|err| {
|
||||
AutomationRunMaterializeError::CloneFailed(format!(
|
||||
|
|
@ -91,38 +91,38 @@ impl AutomationRunMaterializer for GitAutomationRunMaterializer {
|
|||
))
|
||||
})?;
|
||||
let checkout_dir = input.temp_root.join(input.run_id.to_string());
|
||||
run_git(
|
||||
git_clone_args(&clone_url, &checkout_dir),
|
||||
self.git_timeout,
|
||||
"git clone",
|
||||
)
|
||||
.await?;
|
||||
run_git(
|
||||
git_remote_set_url_args(&checkout_dir, &sanitized_clone_url),
|
||||
self.git_timeout,
|
||||
"git remote set-url origin",
|
||||
)
|
||||
.await?;
|
||||
run_git(
|
||||
git_checkout_args(&checkout_dir, input.target.ref_.as_str()),
|
||||
self.git_timeout,
|
||||
"git checkout",
|
||||
)
|
||||
.await?;
|
||||
|
||||
build_manifest_from_checkout(input, checkout_dir).await
|
||||
let result = self
|
||||
.run_checkout(
|
||||
&input,
|
||||
&checkout_dir,
|
||||
&clone_url,
|
||||
&sanitized_clone_url,
|
||||
auth_url.as_ref(),
|
||||
)
|
||||
.await;
|
||||
// Always clean up the materialized clone: callers don't need the
|
||||
// working tree after the manifest is built, and a failed clone
|
||||
// (e.g. partial fetch) should not leak gigabytes into scratch.
|
||||
if let Err(err) = fs::remove_dir_all(&checkout_dir).await {
|
||||
if err.kind() != std::io::ErrorKind::NotFound {
|
||||
tracing::warn!(
|
||||
error = %err,
|
||||
path = %checkout_dir.display(),
|
||||
"Failed to clean up automation checkout",
|
||||
);
|
||||
}
|
||||
}
|
||||
result
|
||||
}
|
||||
}
|
||||
|
||||
impl GitAutomationRunMaterializer {
|
||||
async fn authenticated_clone_url(
|
||||
&self,
|
||||
owner: &str,
|
||||
repo: &str,
|
||||
sanitized_clone_url: &str,
|
||||
) -> Result<String, AutomationRunMaterializeError> {
|
||||
) -> Result<Option<DisplaySafeUrl>, AutomationRunMaterializeError> {
|
||||
let Some(credentials) = self.github_credentials.as_ref() else {
|
||||
return Ok(sanitized_clone_url.to_string());
|
||||
return Ok(None);
|
||||
};
|
||||
let ctx = match self.http_client.clone() {
|
||||
Some(client) => fabro_github::GitHubContext::with_http_client(
|
||||
|
|
@ -132,15 +132,42 @@ impl GitAutomationRunMaterializer {
|
|||
),
|
||||
None => fabro_github::GitHubContext::new(credentials, &self.github_api_base_url),
|
||||
};
|
||||
let (_username, token) = fabro_github::resolve_clone_credentials(&ctx, owner, repo)
|
||||
fabro_github::resolve_authenticated_url(&ctx, sanitized_clone_url)
|
||||
.await
|
||||
.map_err(|err| AutomationRunMaterializeError::CloneFailed(err.to_string()))?;
|
||||
match token {
|
||||
Some(token) => fabro_github::embed_token_in_url(sanitized_clone_url, &token)
|
||||
.map(|url| url.raw_string())
|
||||
.map_err(|err| AutomationRunMaterializeError::CloneFailed(err.to_string())),
|
||||
None => Ok(sanitized_clone_url.to_string()),
|
||||
}
|
||||
.map(Some)
|
||||
.map_err(|err| AutomationRunMaterializeError::CloneFailed(err.to_string()))
|
||||
}
|
||||
|
||||
async fn run_checkout(
|
||||
&self,
|
||||
input: &AutomationRunMaterializeInput,
|
||||
checkout_dir: &Path,
|
||||
clone_url: &str,
|
||||
sanitized_clone_url: &str,
|
||||
auth_url: Option<&DisplaySafeUrl>,
|
||||
) -> Result<AutomationRunMaterialized, AutomationRunMaterializeError> {
|
||||
run_git(
|
||||
git_clone_args(clone_url, checkout_dir),
|
||||
self.git_timeout,
|
||||
"git clone",
|
||||
auth_url,
|
||||
)
|
||||
.await?;
|
||||
run_git(
|
||||
git_remote_set_url_args(checkout_dir, sanitized_clone_url),
|
||||
self.git_timeout,
|
||||
"git remote set-url origin",
|
||||
auth_url,
|
||||
)
|
||||
.await?;
|
||||
run_git(
|
||||
git_checkout_args(checkout_dir, input.target.ref_.as_str()),
|
||||
self.git_timeout,
|
||||
"git checkout",
|
||||
auth_url,
|
||||
)
|
||||
.await?;
|
||||
build_manifest_from_checkout(input, checkout_dir).await
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -187,6 +214,7 @@ async fn run_git(
|
|||
args: Vec<OsString>,
|
||||
git_timeout: Duration,
|
||||
label: &'static str,
|
||||
auth_url: Option<&DisplaySafeUrl>,
|
||||
) -> Result<(), AutomationRunMaterializeError> {
|
||||
let mut command = Command::new("git");
|
||||
command.args(&args);
|
||||
|
|
@ -200,45 +228,40 @@ async fn run_git(
|
|||
git_timeout.as_secs()
|
||||
))
|
||||
})?
|
||||
.map_err(|err| AutomationRunMaterializeError::CloneFailed(format!("{label}: {err}")))?;
|
||||
.map_err(|err| {
|
||||
AutomationRunMaterializeError::CloneFailed(redact_auth_url(
|
||||
&format!("{label}: {err}"),
|
||||
auth_url,
|
||||
))
|
||||
})?;
|
||||
if output.status.success() {
|
||||
return Ok(());
|
||||
}
|
||||
let stderr = String::from_utf8_lossy(&output.stderr).trim().to_string();
|
||||
let stdout = String::from_utf8_lossy(&output.stdout).trim().to_string();
|
||||
let detail = if stderr.is_empty() { stdout } else { stderr };
|
||||
Err(AutomationRunMaterializeError::CloneFailed(format!(
|
||||
"{label} exited with status {}: {}",
|
||||
output.status,
|
||||
redact_command_output(&detail)
|
||||
let stderr = String::from_utf8_lossy(&output.stderr);
|
||||
let stdout = String::from_utf8_lossy(&output.stdout);
|
||||
let detail = if stderr.trim().is_empty() {
|
||||
stdout.trim()
|
||||
} else {
|
||||
stderr.trim()
|
||||
};
|
||||
Err(AutomationRunMaterializeError::CloneFailed(redact_auth_url(
|
||||
&format!("{label} exited with status {}: {detail}", output.status),
|
||||
auth_url,
|
||||
)))
|
||||
}
|
||||
|
||||
fn redact_command_output(value: &str) -> String {
|
||||
value
|
||||
.split_whitespace()
|
||||
.map(redact_url_token)
|
||||
.collect::<Vec<_>>()
|
||||
.join(" ")
|
||||
}
|
||||
|
||||
fn redact_url_token(value: &str) -> String {
|
||||
fabro_redact::DisplaySafeUrl::parse(value)
|
||||
.map_or_else(|_| value.to_string(), |url| url.redacted_string())
|
||||
}
|
||||
|
||||
async fn build_manifest_from_checkout(
|
||||
input: AutomationRunMaterializeInput,
|
||||
checkout_dir: PathBuf,
|
||||
input: &AutomationRunMaterializeInput,
|
||||
checkout_dir: &Path,
|
||||
) -> Result<AutomationRunMaterialized, AutomationRunMaterializeError> {
|
||||
let workflow = PathBuf::from(input.target.workflow.as_str());
|
||||
let user_settings_path = input.user_settings_path;
|
||||
let user_settings_path = input.user_settings_path.clone();
|
||||
let run_id = input.run_id;
|
||||
let automation_id = input.automation_id.to_string();
|
||||
let cwd = checkout_dir.to_path_buf();
|
||||
let built = task::spawn_blocking(move || {
|
||||
fabro_manifest::build_run_manifest(fabro_manifest::ManifestBuildInput {
|
||||
workflow,
|
||||
cwd: checkout_dir,
|
||||
cwd,
|
||||
run_id: Some(run_id),
|
||||
user_settings_path: Some(user_settings_path),
|
||||
..fabro_manifest::ManifestBuildInput::default()
|
||||
|
|
@ -246,7 +269,7 @@ async fn build_manifest_from_checkout(
|
|||
})
|
||||
.await
|
||||
.map_err(|err| AutomationRunMaterializeError::Manifest(err.to_string()))?
|
||||
.map_err(|err| classify_manifest_error(&automation_id, &err))?;
|
||||
.map_err(|err| AutomationRunMaterializeError::Manifest(err.to_string()))?;
|
||||
let submitted_manifest_bytes = serde_json::to_vec(&built.manifest)
|
||||
.map_err(|err| AutomationRunMaterializeError::Manifest(err.to_string()))?;
|
||||
Ok(AutomationRunMaterialized {
|
||||
|
|
@ -255,21 +278,6 @@ async fn build_manifest_from_checkout(
|
|||
})
|
||||
}
|
||||
|
||||
fn classify_manifest_error(
|
||||
automation_id: &str,
|
||||
err: &anyhow::Error,
|
||||
) -> AutomationRunMaterializeError {
|
||||
let message = err.to_string();
|
||||
if err
|
||||
.chain()
|
||||
.any(|cause| cause.to_string().contains("workflow") && cause.to_string().contains("not"))
|
||||
{
|
||||
AutomationRunMaterializeError::WorkflowNotFound(format!("{automation_id}: {message}"))
|
||||
} else {
|
||||
AutomationRunMaterializeError::Manifest(message)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(any(test, feature = "test-support"))]
|
||||
pub(crate) struct StaticAutomationRunMaterializer {
|
||||
result: Result<AutomationRunMaterialized, AutomationRunMaterializeError>,
|
||||
|
|
@ -305,7 +313,7 @@ impl AutomationRunMaterializer for StaticAutomationRunMaterializer {
|
|||
mod tests {
|
||||
use std::str::FromStr as _;
|
||||
|
||||
use fabro_automation::{AutomationId, GitRefSelector, RepositorySlug, WorkflowSlug};
|
||||
use fabro_automation::{GitRefSelector, RepositorySlug, WorkflowSlug};
|
||||
|
||||
use super::*;
|
||||
|
||||
|
|
@ -318,12 +326,17 @@ mod tests {
|
|||
}
|
||||
|
||||
#[test]
|
||||
fn redact_command_output_strips_credentials() {
|
||||
let redacted = redact_command_output(
|
||||
"fatal: https://x-access-token:ghs_secret@github.com/acme/widgets.git failed",
|
||||
fn redact_auth_url_strips_credentials_from_stderr() {
|
||||
let auth_url =
|
||||
DisplaySafeUrl::parse("https://x-access-token:ghs_secret@github.com/acme/widgets.git")
|
||||
.expect("auth url should parse");
|
||||
let redacted = redact_auth_url(
|
||||
"fatal: https://x-access-token:ghs_secret@github.com/acme/widgets.git\nremote: denied",
|
||||
Some(&auth_url),
|
||||
);
|
||||
assert!(redacted.contains("https://x-access-token:***@github.com/acme/widgets.git"));
|
||||
assert!(!redacted.contains("ghs_secret"));
|
||||
// Newlines are preserved (unlike a whitespace-collapse redactor).
|
||||
assert!(redacted.contains('\n'));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -359,14 +372,13 @@ mod tests {
|
|||
};
|
||||
let run_id = RunId::new();
|
||||
let input = AutomationRunMaterializeInput {
|
||||
automation_id: AutomationId::from_str("nightly").unwrap(),
|
||||
target,
|
||||
run_id,
|
||||
user_settings_path: dir.path().join("settings.toml"),
|
||||
temp_root: dir.path().join("tmp"),
|
||||
};
|
||||
|
||||
let materialized = build_manifest_from_checkout(input, dir.path().to_path_buf())
|
||||
let materialized = build_manifest_from_checkout(&input, dir.path())
|
||||
.await
|
||||
.expect("manifest should build");
|
||||
|
||||
|
|
|
|||
|
|
@ -1,5 +1,4 @@
|
|||
use std::collections::BTreeMap;
|
||||
use std::str::FromStr as _;
|
||||
use std::sync::Arc;
|
||||
|
||||
use axum::body::Bytes;
|
||||
|
|
@ -149,8 +148,9 @@ fn default_true() -> bool {
|
|||
}
|
||||
|
||||
async fn list_automations(_auth: RequiredUser, State(state): State<Arc<AppState>>) -> Response {
|
||||
let mut automations = state.automation_store().list().await;
|
||||
automations.sort_by(|left, right| left.id.cmp(&right.id));
|
||||
// `AutomationStore::list` already yields entries in `AutomationId` order
|
||||
// (BTreeMap iteration); no additional sort is required.
|
||||
let automations = state.automation_store().list().await;
|
||||
let total = automations.len() as u64;
|
||||
(
|
||||
StatusCode::OK,
|
||||
|
|
@ -367,7 +367,6 @@ async fn create_automation_run(
|
|||
let materialized = match state
|
||||
.automation_materializer()
|
||||
.materialize(AutomationRunMaterializeInput {
|
||||
automation_id: id.clone(),
|
||||
target: automation.target.clone(),
|
||||
run_id,
|
||||
user_settings_path: state.active_config_path().to_path_buf(),
|
||||
|
|
@ -441,8 +440,7 @@ fn parse_if_match(headers: &HeaderMap) -> Result<AutomationRevision, ApiError> {
|
|||
"If-Match revision must not be empty.",
|
||||
));
|
||||
}
|
||||
Ok(AutomationRevision::from_str(revision)
|
||||
.expect("AutomationRevision accepts any non-empty string"))
|
||||
Ok(AutomationRevision::from_raw(revision))
|
||||
}
|
||||
|
||||
fn with_etag(status: StatusCode, automation: Automation) -> Response {
|
||||
|
|
@ -461,29 +459,22 @@ fn store_error(err: AutomationStoreError) -> ApiError {
|
|||
AutomationStoreError::AlreadyExists(_) => {
|
||||
ApiError::new(StatusCode::CONFLICT, "Automation already exists.")
|
||||
}
|
||||
AutomationStoreError::MissingRevision => ApiError::new(
|
||||
StatusCode::PRECONDITION_REQUIRED,
|
||||
"If-Match header is required.",
|
||||
),
|
||||
AutomationStoreError::RevisionMismatch { .. } => {
|
||||
ApiError::new(StatusCode::CONFLICT, "Automation revision mismatch.")
|
||||
}
|
||||
AutomationStoreError::Validation(err) => validation_error(&err),
|
||||
AutomationStoreError::Parse { .. } | AutomationStoreError::Io { .. } => {
|
||||
AutomationStoreError::Parse { .. }
|
||||
| AutomationStoreError::Serialize(_)
|
||||
| AutomationStoreError::Io { .. } => {
|
||||
ApiError::new(StatusCode::INTERNAL_SERVER_ERROR, err.to_string())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn materialize_error(err: &AutomationRunMaterializeError) -> ApiError {
|
||||
match err {
|
||||
AutomationRunMaterializeError::InvalidTarget(_)
|
||||
| AutomationRunMaterializeError::CloneFailed(_)
|
||||
| AutomationRunMaterializeError::WorkflowNotFound(_)
|
||||
| AutomationRunMaterializeError::Manifest(_) => {
|
||||
ApiError::new(StatusCode::UNPROCESSABLE_ENTITY, err.to_string())
|
||||
}
|
||||
}
|
||||
// All current variants surface as 422 — they describe automation
|
||||
// misconfiguration or repository state that the caller can correct.
|
||||
ApiError::new(StatusCode::UNPROCESSABLE_ENTITY, err.to_string())
|
||||
}
|
||||
|
||||
impl TryFrom<RawAutomationTarget> for AutomationTarget {
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue