use std::collections::BTreeSet;
use std::path::Path;
use anyhow::{Context, Result};
use serde::{Deserialize, Serialize};
use tracing::{debug, info, warn};
use crate::reconciler::pond::{DEFAULT_MINIO_API_PORT, DEFAULT_MINIO_PASSWORD, DEFAULT_MINIO_USER};
use crate::reconciler::static_asset::WORKLOAD_KIND;
use crate::{MirrorProviderSlot, Provider, ServiceConfig};
use local_driver::s3_sign::{sign_s3_get_with_query, sign_s3_no_body, uri_encode_key};
use workload_spec::StaticAssetWorkload;
const CATALOG_MANIFEST_KEY: &str = "_yah-asset-catalog.json";
const R2_REGION: &str = "auto";
const MINIO_REGION: &str = "us-east-1";
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PruneCandidate {
pub filename: String,
pub size: u64,
pub last_modified: String,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct PruneReport {
pub service: String,
pub env: String,
pub bucket: String,
pub live_set: Vec<String>,
pub candidates: Vec<PruneCandidate>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct PruneOutcome {
pub service: String,
pub env: String,
pub bucket: String,
pub requested: Vec<String>,
pub deleted: Vec<String>,
pub errors: Vec<(String, String)>,
}
pub async fn compute_prune_candidates(
workspace_root: &Path,
service: &ServiceConfig,
mirror: &crate::MirrorConfig,
env: &str,
) -> Result<PruneReport> {
let live_set = compute_live_set(workspace_root, service)?;
let backend = resolve_backend(workspace_root, mirror, &service.name, env)?;
let listing = list_bucket_objects(&backend).await?;
let candidates: Vec<PruneCandidate> = listing
.into_iter()
.filter(|obj| obj.filename != CATALOG_MANIFEST_KEY)
.filter(|obj| !live_set.contains(&obj.filename))
.collect();
Ok(PruneReport {
service: service.name.clone(),
env: env.to_string(),
bucket: backend.bucket,
live_set: live_set.into_iter().collect(),
candidates,
})
}
pub async fn execute_prune(
workspace_root: &Path,
service: &ServiceConfig,
mirror: &crate::MirrorConfig,
env: &str,
filenames: &[String],
) -> Result<PruneOutcome> {
let backend = resolve_backend(workspace_root, mirror, &service.name, env)?;
let client = reqwest::Client::new();
let mut outcome = PruneOutcome {
service: service.name.clone(),
env: env.to_string(),
bucket: backend.bucket.clone(),
requested: filenames.to_vec(),
..PruneOutcome::default()
};
for filename in filenames {
if filename == CATALOG_MANIFEST_KEY {
outcome.errors.push((
filename.clone(),
format!(
"refusing to delete catalog manifest sidecar `{CATALOG_MANIFEST_KEY}` — \
it's how the reconciler tracks prune candidates across runs"
),
));
continue;
}
let url = format!(
"{}/{}/{}",
backend.endpoint.trim_end_matches('/'),
backend.bucket,
uri_encode_key(filename)
);
match delete_object(&client, &url, &backend).await {
Ok(()) => {
info!(filename, bucket = %backend.bucket, "deleted");
outcome.deleted.push(filename.clone());
}
Err(e) => {
warn!(filename, error = %format!("{e:#}"), "delete failed");
outcome.errors.push((filename.clone(), format!("{e:#}")));
}
}
}
Ok(outcome)
}
pub fn compute_live_set(
workspace_root: &Path,
service: &ServiceConfig,
) -> Result<BTreeSet<String>> {
let mut live = BTreeSet::new();
for component in &service.components {
if component.kind != WORKLOAD_KIND {
continue;
}
let workload_path = workspace_root.join(&component.path).join("workload.toml");
let workload = load_static_asset_workload(&workload_path).with_context(|| {
format!(
"service `{}` component `{}`: loading {}",
service.name,
component.id,
workload_path.display()
)
})?;
for entry in &workload.assets {
live.insert(entry.filename.clone());
}
}
Ok(live)
}
fn load_static_asset_workload(path: &Path) -> Result<StaticAssetWorkload> {
let src =
std::fs::read_to_string(path).with_context(|| format!("reading {}", path.display()))?;
let envelope: workload_spec::Workload =
toml::from_str(&src).with_context(|| format!("parsing {}", path.display()))?;
match envelope {
workload_spec::Workload::StaticAsset(w) => Ok(w),
other => anyhow::bail!(
"{}: expected kind=\"static-asset\", got {:?}",
path.display(),
other.kind_str()
),
}
}
struct Backend {
endpoint: String,
bucket: String,
region: &'static str,
access_key: String,
secret_key: String,
}
fn resolve_backend(
workspace_root: &Path,
mirror: &crate::MirrorConfig,
service_name: &str,
env: &str,
) -> Result<Backend> {
let slot = mirror.providers.get("object_store").with_context(|| {
format!(
"mirror has no `providers.object_store` slot — required for prune \
(service={service_name}, env={env})"
)
})?;
match slot {
MirrorProviderSlot::Reference {
provider_id,
fields,
} => {
let cf = super::cf_creds::CfProvider::resolve(workspace_root, provider_id)?;
anyhow::ensure!(
matches!(cf.cfg.kind, Provider::Cloudflare),
"providers.object_store.use = {provider_id:?} (kind={:?}) not supported for \
prune — only cloudflare-kind reference providers are",
cf.cfg.kind,
);
let bucket = fields
.get("bucket")
.and_then(|v| v.as_str())
.context("providers.object_store missing `bucket` for cloudflare")?
.to_string();
let account_id = cf.account_id.clone();
let (access_key, secret_key) = cf.r2_keys()?;
Ok(Backend {
endpoint: format!("https://{account_id}.r2.cloudflarestorage.com"),
bucket,
region: R2_REGION,
access_key,
secret_key,
})
}
MirrorProviderSlot::Inline {
kind: Provider::MinioContainer,
fields,
} => {
let api_port = fields
.get("api_port")
.and_then(|v| v.as_integer())
.and_then(|n| u16::try_from(n).ok())
.unwrap_or(DEFAULT_MINIO_API_PORT);
let bucket = fields
.get("bucket")
.and_then(|v| v.as_str())
.context("providers.object_store missing `bucket` for minio-container")?
.to_string();
Ok(Backend {
endpoint: format!("http://127.0.0.1:{api_port}"),
bucket,
region: MINIO_REGION,
access_key: DEFAULT_MINIO_USER.to_string(),
secret_key: DEFAULT_MINIO_PASSWORD.to_string(),
})
}
MirrorProviderSlot::Inline { kind, .. } => {
anyhow::bail!(
"providers.object_store.kind = {kind:?} not supported for prune \
(expected `minio-container` or a `cloudflare` reference)"
)
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct BucketObject {
filename: String,
size: u64,
last_modified: String,
}
impl From<BucketObject> for PruneCandidate {
fn from(o: BucketObject) -> Self {
Self {
filename: o.filename,
size: o.size,
last_modified: o.last_modified,
}
}
}
async fn list_bucket_objects(backend: &Backend) -> Result<Vec<PruneCandidate>> {
let client = reqwest::Client::new();
let mut out = Vec::new();
let mut continuation: Option<String> = None;
loop {
let canonical_query = match &continuation {
Some(token) => format!("continuation-token={}&list-type=2", urlencode(token)),
None => "list-type=2".to_string(),
};
let url = format!(
"{}/{}?{}",
backend.endpoint.trim_end_matches('/'),
backend.bucket,
canonical_query
);
let headers = sign_s3_get_with_query(
&url,
&canonical_query,
backend.region,
&backend.access_key,
&backend.secret_key,
)
.with_context(|| format!("signing LIST {url}"))?;
let resp = client
.get(&url)
.headers(headers)
.send()
.await
.with_context(|| format!("LIST {url}"))?;
if !resp.status().is_success() {
let status = resp.status();
let body = resp.text().await.unwrap_or_default();
anyhow::bail!("LIST {url} → {status}: {}", body.trim());
}
let body = resp.text().await.context("reading LIST response body")?;
let (objects, next_token) = parse_list_response(&body)?;
debug!(
count = objects.len(),
has_more = next_token.is_some(),
"list-objects page"
);
out.extend(objects.into_iter().map(PruneCandidate::from));
match next_token {
Some(t) => continuation = Some(t),
None => break,
}
}
Ok(out)
}
fn parse_list_response(body: &str) -> Result<(Vec<BucketObject>, Option<String>)> {
let mut objects = Vec::new();
for chunk in split_tags(body, "Contents") {
let filename = inner_text(chunk, "Key").context("missing <Key>")?;
let size = inner_text(chunk, "Size")
.context("missing <Size>")?
.parse::<u64>()
.context("parsing <Size> as u64")?;
let last_modified = inner_text(chunk, "LastModified")
.context("missing <LastModified>")?
.to_string();
objects.push(BucketObject {
filename: filename.to_string(),
size,
last_modified,
});
}
let is_truncated = inner_text(body, "IsTruncated")
.map(|s| s.eq_ignore_ascii_case("true"))
.unwrap_or(false);
let next_token = if is_truncated {
inner_text(body, "NextContinuationToken").map(|s| s.to_string())
} else {
None
};
Ok((objects, next_token))
}
fn split_tags<'a>(body: &'a str, tag: &str) -> Vec<&'a str> {
let open = format!("<{tag}>");
let close = format!("</{tag}>");
let mut out = Vec::new();
let mut rest = body;
while let Some(start) = rest.find(&open) {
let after = &rest[start + open.len()..];
match after.find(&close) {
Some(end) => {
out.push(&after[..end]);
rest = &after[end + close.len()..];
}
None => break,
}
}
out
}
fn inner_text<'a>(body: &'a str, tag: &str) -> Option<&'a str> {
split_tags(body, tag).into_iter().next().map(|s| s.trim())
}
async fn delete_object(client: &reqwest::Client, url: &str, backend: &Backend) -> Result<()> {
let headers = sign_s3_no_body(
"DELETE",
url,
"",
backend.region,
&backend.access_key,
&backend.secret_key,
)
.with_context(|| format!("signing DELETE {url}"))?;
let resp = client
.delete(url)
.headers(headers)
.send()
.await
.with_context(|| format!("DELETE {url}"))?;
if !resp.status().is_success() {
let status = resp.status();
let body = resp.text().await.unwrap_or_default();
anyhow::bail!("DELETE {url} → {status}: {}", body.trim());
}
Ok(())
}
fn urlencode(s: &str) -> String {
let mut out = String::with_capacity(s.len());
for b in s.bytes() {
let unreserved =
b.is_ascii_alphanumeric() || b == b'-' || b == b'.' || b == b'_' || b == b'~';
if unreserved {
out.push(b as char);
} else {
out.push_str(&format!("%{b:02X}"));
}
}
out
}
pub fn load_service_and_mirror(
workspace_root: &Path,
service_name: &str,
env: &str,
) -> Result<(ServiceConfig, crate::MirrorConfig)> {
let service_toml = crate::paths::service_toml(workspace_root, service_name);
let service = ServiceConfig::load(&service_toml).with_context(|| {
format!(
"loading service `{service_name}` from {}",
service_toml.display()
)
})?;
let mirror_toml = crate::paths::service_mirror_toml(workspace_root, service_name, env);
let mirror = crate::MirrorConfig::load(&mirror_toml).with_context(|| {
format!(
"loading mirror `{env}` for service `{service_name}` from {}",
mirror_toml.display()
)
})?;
Ok((service, mirror))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::ServiceComponent;
use tempfile::tempdir;
const HASH_64: &str = "abcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890";
fn write_static_asset_workload(dir: &Path, body: &str) {
std::fs::create_dir_all(dir).unwrap();
std::fs::write(dir.join("workload.toml"), body).unwrap();
}
fn svc_with_components(name: &str, components: Vec<ServiceComponent>) -> ServiceConfig {
ServiceConfig {
schema_version: 1,
name: name.into(),
address: crate::config::ServiceAddress::front_door("releases.example"),
description: None,
components,
}
}
fn asset_component(id: &str, path: &str) -> ServiceComponent {
ServiceComponent {
mount: None,
id: id.into(),
kind: "static-asset".into(),
path: path.into(),
role: "assets".into(),
publishes: None,
wave: 0,
git: None,
deploy: Default::default(),
}
}
#[test]
fn live_set_unions_filenames_across_components() {
let tmp = tempdir().unwrap();
let root = tmp.path();
write_static_asset_workload(
&root.join("comp-a"),
&format!(
"kind = \"static-asset\"\nschema_version = \"V1\"\n\n\
[[asset]]\nfilename = \"a/one.bin\"\nsource = \"src/one.bin\"\nblake3 = \"{HASH_64}\"\n\n\
[[asset]]\nfilename = \"a/two.bin\"\nsource = \"src/two.bin\"\nblake3 = \"{HASH_64}\"\n"
),
);
write_static_asset_workload(
&root.join("comp-b"),
&format!(
"kind = \"static-asset\"\nschema_version = \"V1\"\n\n\
[[asset]]\nfilename = \"b/three.bin\"\nsource = \"src/three.bin\"\nblake3 = \"{HASH_64}\"\n"
),
);
let svc = svc_with_components(
"demo",
vec![
asset_component("a", "comp-a"),
asset_component("b", "comp-b"),
],
);
let live = compute_live_set(root, &svc).unwrap();
let want: BTreeSet<String> = ["a/one.bin", "a/two.bin", "b/three.bin"]
.iter()
.map(|s| s.to_string())
.collect();
assert_eq!(live, want);
}
#[test]
fn live_set_ignores_non_static_asset_components() {
let tmp = tempdir().unwrap();
let root = tmp.path();
write_static_asset_workload(
&root.join("assets"),
&format!(
"kind = \"static-asset\"\nschema_version = \"V1\"\n\n\
[[asset]]\nfilename = \"x.bin\"\nsource = \"x.bin\"\nblake3 = \"{HASH_64}\"\n"
),
);
let mut svc = svc_with_components("mixed", vec![asset_component("assets", "assets")]);
svc.components.push(ServiceComponent {
mount: None,
id: "site".into(),
kind: "mesofact-static".into(),
path: "site-dir-doesnt-exist".into(),
role: "static".into(),
publishes: None,
wave: 1,
git: None,
deploy: Default::default(),
});
let live = compute_live_set(root, &svc).unwrap();
assert_eq!(live.len(), 1);
assert!(live.contains("x.bin"));
}
#[test]
fn live_set_empty_when_no_static_asset_components() {
let svc = svc_with_components("none", vec![]);
let tmp = tempdir().unwrap();
let live = compute_live_set(tmp.path(), &svc).unwrap();
assert!(live.is_empty());
}
#[test]
fn live_set_errors_on_wrong_kind_workload() {
let tmp = tempdir().unwrap();
let root = tmp.path();
std::fs::create_dir_all(root.join("oops")).unwrap();
std::fs::write(
root.join("oops/workload.toml"),
"kind = \"almanac\"\nschema_version = \"V1\"\ncommand = \"true\"\ncadence = \"once\"\n",
)
.unwrap();
let svc = svc_with_components("oops", vec![asset_component("oops", "oops")]);
let err = compute_live_set(root, &svc).unwrap_err();
let msg = format!("{err:#}");
assert!(msg.contains("static-asset"), "got: {msg}");
}
#[test]
fn parse_list_response_single_page() {
let xml = r#"<?xml version="1.0"?>
<ListBucketResult>
<Name>yah-dev</Name>
<IsTruncated>false</IsTruncated>
<Contents>
<Key>whisper/v1.bin</Key>
<LastModified>2026-01-01T00:00:00.000Z</LastModified>
<ETag>"abc"</ETag>
<Size>100</Size>
</Contents>
<Contents>
<Key>whisper/v2.bin</Key>
<LastModified>2026-02-01T00:00:00.000Z</LastModified>
<ETag>"def"</ETag>
<Size>200</Size>
</Contents>
</ListBucketResult>"#;
let (objects, next) = parse_list_response(xml).unwrap();
assert_eq!(objects.len(), 2);
assert_eq!(objects[0].filename, "whisper/v1.bin");
assert_eq!(objects[0].size, 100);
assert_eq!(objects[0].last_modified, "2026-01-01T00:00:00.000Z");
assert_eq!(objects[1].filename, "whisper/v2.bin");
assert_eq!(objects[1].size, 200);
assert!(next.is_none());
}
#[test]
fn parse_list_response_truncated_yields_token() {
let xml = r#"<ListBucketResult>
<IsTruncated>true</IsTruncated>
<NextContinuationToken>opaque-token-123</NextContinuationToken>
<Contents>
<Key>k</Key>
<LastModified>2026-01-01T00:00:00.000Z</LastModified>
<Size>1</Size>
</Contents>
</ListBucketResult>"#;
let (objects, next) = parse_list_response(xml).unwrap();
assert_eq!(objects.len(), 1);
assert_eq!(next.as_deref(), Some("opaque-token-123"));
}
#[test]
fn parse_list_response_empty_bucket() {
let xml = r#"<ListBucketResult>
<IsTruncated>false</IsTruncated>
</ListBucketResult>"#;
let (objects, next) = parse_list_response(xml).unwrap();
assert!(objects.is_empty());
assert!(next.is_none());
}
#[test]
fn parse_list_response_truncated_without_token_returns_none() {
let xml = r#"<ListBucketResult>
<IsTruncated>true</IsTruncated>
</ListBucketResult>"#;
let (_, next) = parse_list_response(xml).unwrap();
assert!(next.is_none());
}
#[test]
fn urlencode_preserves_unreserved() {
assert_eq!(urlencode("abcXYZ_.-~0"), "abcXYZ_.-~0");
}
#[test]
fn urlencode_escapes_reserved() {
assert_eq!(urlencode("a/b c?d"), "a%2Fb%20c%3Fd");
}
#[test]
fn candidates_exclude_catalog_manifest_sidecar() {
let live: BTreeSet<String> = BTreeSet::new();
let bucket = vec![
PruneCandidate {
filename: CATALOG_MANIFEST_KEY.to_string(),
size: 64,
last_modified: "2026-01-01T00:00:00Z".into(),
},
PruneCandidate {
filename: "old.bin".to_string(),
size: 1024,
last_modified: "2026-01-01T00:00:00Z".into(),
},
];
let candidates: Vec<_> = bucket
.into_iter()
.filter(|o| o.filename != CATALOG_MANIFEST_KEY)
.filter(|o| !live.contains(&o.filename))
.collect();
assert_eq!(candidates.len(), 1);
assert_eq!(candidates[0].filename, "old.bin");
}
#[test]
fn candidates_exclude_live_filenames() {
let mut live = BTreeSet::new();
live.insert("keep.bin".to_string());
let bucket = vec![
PruneCandidate {
filename: "keep.bin".to_string(),
size: 1,
last_modified: "x".into(),
},
PruneCandidate {
filename: "drop.bin".to_string(),
size: 1,
last_modified: "x".into(),
},
];
let candidates: Vec<_> = bucket
.into_iter()
.filter(|o| !live.contains(&o.filename))
.collect();
assert_eq!(candidates.len(), 1);
assert_eq!(candidates[0].filename, "drop.bin");
}
}