mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
refactor: dedupe artifact entry adapters, retry URL helper, query-param check
Extract run_artifact_entry_from / artifact_entry_from in fabro-server so the two list-artifact handlers share a single conversion site. Push the ?retry=... query append into stage_artifacts_url in fabro-client so upload callers don't repeat it. Replace required_filename and required_retry with one generic required_query_param<T> helper. (From impls were the cleaner shape but the orphan rule blocks them: NodeArtifact lives in fabro-store, RunArtifactEntry in fabro-api, neither is in fabro-server. Free fns achieve the same dedup without adding a fabro-store -> fabro-api coupling.) Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
f628b91c04
commit
8f4c12580c
2 changed files with 41 additions and 48 deletions
|
|
@ -1228,7 +1228,12 @@ impl Client {
|
|||
clippy::disallowed_types,
|
||||
reason = "Client builds raw server API request URLs for wire transit; logging redaction is handled at log boundaries."
|
||||
)]
|
||||
fn stage_artifacts_url(&self, run_id: &RunId, stage_id: &StageId) -> Result<fabro_http::Url> {
|
||||
fn stage_artifacts_url(
|
||||
&self,
|
||||
run_id: &RunId,
|
||||
stage_id: &StageId,
|
||||
retry: u32,
|
||||
) -> Result<fabro_http::Url> {
|
||||
let base_url = self.base_url();
|
||||
let mut url = fabro_http::Url::parse(&base_url)
|
||||
.with_context(|| format!("invalid server base URL {base_url}"))?;
|
||||
|
|
@ -1243,6 +1248,8 @@ impl Client {
|
|||
&stage_id.to_string(),
|
||||
"artifacts",
|
||||
]);
|
||||
url.query_pairs_mut()
|
||||
.append_pair("retry", &retry.to_string());
|
||||
Ok(url)
|
||||
}
|
||||
|
||||
|
|
@ -1255,10 +1262,8 @@ impl Client {
|
|||
path: &Path,
|
||||
bearer_token: &str,
|
||||
) -> Result<()> {
|
||||
let mut url = self.stage_artifacts_url(run_id, stage_id)?;
|
||||
let mut url = self.stage_artifacts_url(run_id, stage_id, retry)?;
|
||||
url.query_pairs_mut().append_pair("filename", filename);
|
||||
url.query_pairs_mut()
|
||||
.append_pair("retry", &retry.to_string());
|
||||
|
||||
let file = File::open(path)
|
||||
.await
|
||||
|
|
@ -1296,9 +1301,7 @@ impl Client {
|
|||
artifacts: &[ArtifactUpload],
|
||||
bearer_token: &str,
|
||||
) -> Result<()> {
|
||||
let mut url = self.stage_artifacts_url(run_id, stage_id)?;
|
||||
url.query_pairs_mut()
|
||||
.append_pair("retry", &retry.to_string());
|
||||
let url = self.stage_artifacts_url(run_id, stage_id, retry)?;
|
||||
let mut manifest_entries = Vec::with_capacity(artifacts.len());
|
||||
let mut file_parts = Vec::with_capacity(artifacts.len());
|
||||
|
||||
|
|
|
|||
|
|
@ -63,8 +63,8 @@ use fabro_slack::threads::ThreadRegistry;
|
|||
use fabro_slack::{blocks as slack_blocks, connection as slack_connection};
|
||||
use fabro_static::EnvVars;
|
||||
use fabro_store::{
|
||||
ArtifactKey, ArtifactStore, Database, EventEnvelope, EventPayload, PendingInterviewRecord,
|
||||
StageId,
|
||||
ArtifactKey, ArtifactStore, Database, EventEnvelope, EventPayload, NodeArtifact,
|
||||
PendingInterviewRecord, StageArtifactEntry, StageId,
|
||||
};
|
||||
#[cfg(test)]
|
||||
use fabro_types::BlockedReason;
|
||||
|
|
@ -3353,24 +3353,12 @@ pub(crate) fn parse_blob_id_path(blob_id: &str) -> Result<RunBlobId, Response> {
|
|||
|
||||
#[allow(
|
||||
clippy::result_large_err,
|
||||
reason = "Missing filename validation returns HTTP 400 responses directly."
|
||||
reason = "Missing query parameter validation returns HTTP 400 responses directly."
|
||||
)]
|
||||
fn required_filename(params: &ArtifactFilenameParams) -> Result<String, Response> {
|
||||
match params.filename.as_ref() {
|
||||
Some(filename) if !filename.is_empty() => Ok(filename.clone()),
|
||||
_ => Err(ApiError::bad_request("Missing filename query parameter.").into_response()),
|
||||
}
|
||||
}
|
||||
|
||||
#[allow(
|
||||
clippy::result_large_err,
|
||||
reason = "Missing retry validation returns HTTP 400 responses directly."
|
||||
)]
|
||||
fn required_retry(params: &ArtifactFilenameParams) -> Result<u32, Response> {
|
||||
match params.retry {
|
||||
Some(retry) => Ok(retry),
|
||||
None => Err(ApiError::bad_request("Missing retry query parameter.").into_response()),
|
||||
}
|
||||
fn required_query_param<T: Clone>(value: Option<&T>, name: &str) -> Result<T, Response> {
|
||||
value.cloned().ok_or_else(|| {
|
||||
ApiError::bad_request(format!("Missing {name} query parameter.")).into_response()
|
||||
})
|
||||
}
|
||||
|
||||
#[allow(
|
||||
|
|
@ -6191,16 +6179,7 @@ async fn list_run_artifacts(
|
|||
|
||||
match state.artifact_store.list_for_run(&id).await {
|
||||
Ok(entries) => Json(RunArtifactListResponse {
|
||||
data: entries
|
||||
.into_iter()
|
||||
.map(|entry| RunArtifactEntry {
|
||||
stage_id: entry.node.to_string(),
|
||||
node_slug: entry.node.node_id().to_string(),
|
||||
retry: entry.retry.cast_signed(),
|
||||
relative_path: entry.filename,
|
||||
size: entry.size.cast_signed(),
|
||||
})
|
||||
.collect(),
|
||||
data: entries.into_iter().map(run_artifact_entry_from).collect(),
|
||||
})
|
||||
.into_response(),
|
||||
Err(err) => {
|
||||
|
|
@ -6209,6 +6188,24 @@ async fn list_run_artifacts(
|
|||
}
|
||||
}
|
||||
|
||||
fn run_artifact_entry_from(entry: NodeArtifact) -> RunArtifactEntry {
|
||||
RunArtifactEntry {
|
||||
stage_id: entry.node.to_string(),
|
||||
node_slug: entry.node.node_id().to_string(),
|
||||
retry: entry.retry.cast_signed(),
|
||||
relative_path: entry.filename,
|
||||
size: entry.size.cast_signed(),
|
||||
}
|
||||
}
|
||||
|
||||
fn artifact_entry_from(entry: StageArtifactEntry) -> ArtifactEntry {
|
||||
ArtifactEntry {
|
||||
filename: entry.filename,
|
||||
retry: entry.retry.cast_signed(),
|
||||
size: entry.size.cast_signed(),
|
||||
}
|
||||
}
|
||||
|
||||
async fn list_stage_artifacts(
|
||||
_auth: AuthenticatedService,
|
||||
State(state): State<Arc<AppState>>,
|
||||
|
|
@ -6228,14 +6225,7 @@ async fn list_stage_artifacts(
|
|||
|
||||
match state.artifact_store.list_for_node(&id, &stage_id).await {
|
||||
Ok(entries) => Json(ArtifactListResponse {
|
||||
data: entries
|
||||
.into_iter()
|
||||
.map(|entry| ArtifactEntry {
|
||||
filename: entry.filename,
|
||||
retry: entry.retry.cast_signed(),
|
||||
size: entry.size.cast_signed(),
|
||||
})
|
||||
.collect(),
|
||||
data: entries.into_iter().map(artifact_entry_from).collect(),
|
||||
})
|
||||
.into_response(),
|
||||
Err(err) => {
|
||||
|
|
@ -6614,7 +6604,7 @@ async fn put_stage_artifact(
|
|||
if let Err(response) = load_run_spec(state.as_ref(), &id).await.map(|_| ()) {
|
||||
return response;
|
||||
}
|
||||
let retry = match required_retry(¶ms) {
|
||||
let retry = match required_query_param(params.retry.as_ref(), "retry") {
|
||||
Ok(retry) => retry,
|
||||
Err(response) => return response,
|
||||
};
|
||||
|
|
@ -6625,7 +6615,7 @@ async fn put_stage_artifact(
|
|||
};
|
||||
match artifact_upload_content_type(&parts.headers) {
|
||||
Ok(ArtifactUploadContentType::OctetStream) => {
|
||||
let filename = match required_filename(¶ms) {
|
||||
let filename = match required_query_param(params.filename.as_ref(), "filename") {
|
||||
Ok(filename) => filename,
|
||||
Err(response) => return response,
|
||||
};
|
||||
|
|
@ -6667,11 +6657,11 @@ async fn get_stage_artifact(
|
|||
Ok(stage_id) => stage_id,
|
||||
Err(response) => return response,
|
||||
};
|
||||
let filename = match required_filename(¶ms) {
|
||||
let filename = match required_query_param(params.filename.as_ref(), "filename") {
|
||||
Ok(filename) => filename,
|
||||
Err(response) => return response,
|
||||
};
|
||||
let retry = match required_retry(¶ms) {
|
||||
let retry = match required_query_param(params.retry.as_ref(), "retry") {
|
||||
Ok(retry) => retry,
|
||||
Err(response) => return response,
|
||||
};
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue