//! Adapter regressions through the native HTTP client, using loopback only.
use super::*;
use crate::native_s3::tests::{Reply, configure_repo, fixture, fixture_handler};
use md5::Md5;
use sha2::{Digest, Sha256};
fn repo() -> (tempfile::TempDir, GitRepo) {
let directory = tempfile::tempdir().unwrap();
let repo = GitRepo {
root: directory.path().to_path_buf(),
};
(directory, repo)
}
fn entry(object: &str) -> Entry {
Entry {
metadata: StorageEntry {
pointer: format!("{object}.dvc"),
object: object.to_owned(),
md5: Some("900150983cd24fb0d6963f7d28e17f72".to_owned()),
size: Some(3),
version_id: Some("v1".to_owned()),
etag: Some("abc".to_owned()),
hash_name: "md5".to_owned(),
},
key: format!("root/{object}"),
version: "v1".to_owned(),
etag: Some("abc".to_owned()),
}
}
fn head(version: &str, tag: &str, size: u64) -> Reply {
Reply {
status: 200,
headers: vec![
("x-amz-version-id", version.to_owned()),
("ETag", format!("\"{tag}\"")),
("Content-Length", size.to_string()),
],
body: Vec::new(),
}
}
fn missing(code: &str, status: u16) -> Reply {
Reply {
status,
headers: vec![("x-amz-error-code", code.to_owned())],
body: format!("<Error><Code>{code}</Code><Message>fixture</Message></Error>").into_bytes(),
}
}
fn copied_receipt() -> Value {
json!({"schema_version":1,"status":"copied","remote":"workspace-mgr","bucket":"fixture-bucket","remote_prefix":"root","source":"task","destination":"archive/task","transaction_id":"fixture-transaction","versions":[{"source_object":"task/a","destination_object":"archive/task/a","source_version_id":"v1","source_last_modified":"2026-10-07T00:00:00+00:00","source_is_latest":true,"source_list_order":0,"delete_marker":false,"size":3,"source_etag":"abc","destination_version_id":"dst","destination_etag":"copied","destination_last_modified":"2026-10-07T00:00:01+00:00"}]})
}
fn registry_object() -> String {
format!(
"root/.workspace-mgr/archive/{}.json",
crate::hex::encode_lower(Sha256::digest(b"task"))
)
}
fn registry_listing() -> Reply {
Reply::xml(&format!(
"<ListVersionsResult><Version><Key>{}</Key><VersionId>registry-version</VersionId><Size>1</Size><ETag>"r"</ETag><IsLatest>true</IsLatest><LastModified>2026-10-07T00:00:00Z</LastModified></Version><IsTruncated>false</IsTruncated></ListVersionsResult>",
registry_object()
))
}
fn registry_body(receipt: &Value) -> Reply {
Reply {
status: 200,
headers: vec![("x-amz-version-id", "registry-version".into())],
body: serde_json::to_vec(receipt).unwrap(),
}
}
#[test]
fn only_missing_exact_head_resolves_verified_historical_mapping() {
let (_directory, repo) = repo();
let receipt = copied_receipt();
let (client, worker) = fixture_handler(5, move |request| {
if request.method == "HEAD" {
if request.target.contains("/archive/task/a?") {
head("dst", "copied", 3)
} else {
missing("NoSuchVersion", 404)
}
} else if request.target.contains("versions=") {
registry_listing()
} else {
registry_body(&receipt)
}
});
verify_head(&client, &repo, &entry("task/a")).unwrap();
let requests = worker.join().unwrap();
assert_eq!(
requests
.iter()
.filter(|request| request.method == "HEAD")
.count(),
2
);
assert!(
requests
.iter()
.any(|request| request.target == "/fixture-bucket/root/archive/task/a?versionId=dst")
);
}
#[test]
fn historical_mapping_mismatched_size_or_etag_never_heads_destination() {
let (_directory, repo) = repo();
for field in ["size", "source_etag"] {
let mut receipt = copied_receipt();
receipt["versions"][0][field] = if field == "size" {
4.into()
} else {
"different".into()
};
let (client, worker) = fixture_handler(4, move |request| {
if request.method == "HEAD" {
missing("NoSuchVersion", 404)
} else if request.target.contains("versions=") {
registry_listing()
} else {
registry_body(&receipt)
}
});
assert!(
verify_head(&client, &repo, &entry("task/a"))
.unwrap_err()
.to_string()
.contains("mismatched")
);
let requests = worker.join().unwrap();
assert_eq!(
requests
.iter()
.filter(|request| request.method == "HEAD")
.count(),
1
);
}
}
#[test]
fn dense_listing_stops_at_two_pages_and_falls_back_to_exact_heads() {
let (_directory, repo) = repo();
let (client, worker) = fixture_handler(10, |request| {
if request.method == "HEAD" {
head("v1", "abc", 3)
} else {
Reply::xml(
"<ListVersionsResult><IsTruncated>true</IsTruncated><NextKeyMarker>root/task/nonadvancing</NextKeyMarker><NextVersionIdMarker>same</NextVersionIdMarker></ListVersionsResult>",
)
}
});
let entries = (0..8)
.map(|index| entry(&format!("task/a{index}")))
.collect::<Vec<_>>();
verify_entries(&client, &repo, &entries).unwrap();
let requests = worker.join().unwrap();
assert_eq!(
requests
.iter()
.filter(|request| request.method == "GET")
.count(),
2
);
assert_eq!(
requests
.iter()
.filter(|request| request.method == "HEAD")
.count(),
8
);
}
#[test]
fn denied_dense_listing_uses_readable_exact_heads() {
let (_directory, repo) = repo();
let (client, worker) = fixture_handler(9, |request| {
if request.method == "HEAD" {
head("v1", "abc", 3)
} else {
missing("AccessDenied", 403)
}
});
let entries = (0..8)
.map(|index| entry(&format!("task/a{index}")))
.collect::<Vec<_>>();
verify_entries(&client, &repo, &entries).unwrap();
let requests = worker.join().unwrap();
assert_eq!(
requests
.iter()
.filter(|request| request.method == "GET")
.count(),
1
);
assert_eq!(
requests
.iter()
.filter(|request| request.method == "HEAD")
.count(),
8
);
}
#[test]
fn network_read_workers_are_bounded_to_sixteen() {
use std::sync::{
Arc,
atomic::{AtomicUsize, Ordering},
};
let (_directory, repo) = repo();
let active = Arc::new(AtomicUsize::new(0));
let maximum = Arc::new(AtomicUsize::new(0));
let worker_active = active.clone();
let worker_maximum = maximum.clone();
let (client, worker) = fixture_handler(40, move |request| {
assert_eq!(request.method, "HEAD");
let count = worker_active.fetch_add(1, Ordering::SeqCst) + 1;
worker_maximum.fetch_max(count, Ordering::SeqCst);
std::thread::sleep(std::time::Duration::from_millis(35));
worker_active.fetch_sub(1, Ordering::SeqCst);
head("v1", "abc", 3)
});
// Sparse prefixes require HEAD; no listing can conceal an unbounded pool.
let entries = (0..40)
.map(|index| entry(&format!("task{index}/a")))
.collect::<Vec<_>>();
verify_entries(&client, &repo, &entries).unwrap();
assert_eq!(worker.join().unwrap().len(), 40);
assert!(maximum.load(Ordering::SeqCst) <= 16);
assert!(maximum.load(Ordering::SeqCst) > 1);
}
#[test]
fn sparse_versions_use_exact_heads_and_metadata_mismatch_fails() {
let (_dir, repo) = repo();
let (client, worker) = fixture(vec![head("v1", "abc", 3)]);
verify_entries(&client, &repo, &[entry("task/a")]).unwrap();
let requests = worker.join().unwrap();
assert_eq!(requests.len(), 1);
assert_eq!(requests[0].method, "HEAD");
assert!(requests[0].target.ends_with("versionId=v1"));
for response in [
head("wrong", "abc", 3),
head("v1", "wrong", 3),
head("v1", "abc", 4),
] {
let (client, worker) = fixture(vec![response]);
assert!(
verify_head(&client, &repo, &entry("task/a"))
.unwrap_err()
.to_string()
.contains("mismatched")
);
worker.join().unwrap();
}
}
#[test]
fn permission_and_named_bucket_errors_never_trigger_registry_lookup() {
let (_dir, repo) = repo();
for (status, code) in [
(403, "AccessDenied"),
(404, "NoSuchBucket"),
(403, "NoSuchBucket"),
] {
let (client, worker) = fixture(vec![missing(code, status)]);
assert!(verify_head(&client, &repo, &entry("task/a")).is_err());
let requests = worker.join().unwrap();
assert_eq!(requests.len(), 1);
assert_eq!(requests[0].method, "HEAD");
}
}
#[test]
fn cache_content_requires_exact_physical_size_and_md5() {
let (dir, _repo) = repo();
let path = dir.path().join("bytes");
let entry = entry("task/a");
fs::write(&path, b"abc").unwrap();
assert!(content_matches(&entry, &path).unwrap());
fs::write(&path, b"abd").unwrap();
assert!(!content_matches(&entry, &path).unwrap());
fs::write(&path, b"abcx").unwrap();
assert!(!content_matches(&entry, &path).unwrap());
}
#[test]
fn pending_aliases_validate_exact_storage_size_and_identity() {
let (client, _worker) = fixture(Vec::new());
let receipt = json!({"schema_version":1,"status":"planned","remote":"workspace-mgr","bucket":"fixture-bucket","remote_prefix":"root","source":"task","destination":"archive/task","versions":[{"source_object":"task/a","destination_object":"archive/task/a","source_version_id":"v1","delete_marker":false,"size":3,"source_etag":"abc"}]});
let mut target = entry("archive/task/a");
pending_aliases(
&client,
std::slice::from_mut(&mut target),
std::slice::from_ref(&receipt),
)
.unwrap();
assert_eq!(target.key, "root/task/a");
assert_eq!(target.version, "v1");
let mut invalid = receipt.clone();
invalid["bucket"] = "another".into();
assert!(pending_aliases(&client, &mut [entry("archive/task/a")], &[invalid]).is_err());
let mut wrong = entry("archive/task/a");
wrong.metadata.size = Some(4);
assert!(pending_aliases(&client, &mut [wrong], &[receipt]).is_err());
}
#[test]
fn historical_directory_manifest_flattens_without_reading_payload() {
let (dir, repo) = repo();
repo.run(["init"]).unwrap();
repo.run(["config", "user.name", "Fixture"]).unwrap();
repo.run(["config", "user.email", "fixture@example.invalid"])
.unwrap();
fs::create_dir_all(dir.path().join("task")).unwrap();
// DVC's directory hash is MD5 of its sorted logical [{md5,relpath}] tree.
let tree = b"[{\"md5\": \"900150983cd24fb0d6963f7d28e17f72\", \"relpath\": \"a\"}]";
let digest = format!("{}.dir", crate::hex::encode_lower(Md5::digest(tree)));
let pointer = format!(
"outs:\n- md5: {digest}\n size: 3\n path: data\n files:\n - md5: 900150983cd24fb0d6963f7d28e17f72\n size: 3\n relpath: a\n cloud:\n workspace-mgr:\n version_id: v1\n etag: abc\n"
);
fs::write(dir.path().join("task/data.dvc"), pointer).unwrap();
repo.run(["add", "task/data.dvc"]).unwrap();
repo.run(["commit", "-m", "fixture"]).unwrap();
let payload = json!([{"revision":"HEAD","pointers":["task/data.dvc"]}]);
let listed = purge(&repo, "list", &payload).unwrap();
assert_eq!(
listed,
json!([{"pointer":"task/data.dvc","object":"task/data/a","version_id":"v1"}])
);
assert!(!dir.path().join("task/data").exists());
}
fn version_row(key: &str, id: &str) -> String {
format!(
"<Version><Key>{key}</Key><VersionId>{id}</VersionId><IsLatest>false</IsLatest><LastModified>2026-10-07T00:00:00Z</LastModified><ETag>"abc"</ETag><Size>3</Size></Version>"
)
}
fn marker_row(key: &str, id: &str) -> String {
format!(
"<DeleteMarker><Key>{key}</Key><VersionId>{id}</VersionId><IsLatest>true</IsLatest><LastModified>2026-10-07T00:00:01Z</LastModified></DeleteMarker>"
)
}
fn history(rows: &str) -> Reply {
Reply::xml(&format!(
"<ListVersionsResult>{rows}<IsTruncated>false</IsTruncated></ListVersionsResult>"
))
}
fn deleted() -> Reply {
Reply {
status: 204,
headers: Vec::new(),
body: Vec::new(),
}
}
#[test]
fn generic_purge_deletes_complete_exact_object_history_including_markers() {
let (_directory, repo) = repo();
let first = Reply::xml(&format!(
"<ListVersionsResult>{}<IsTruncated>true</IsTruncated><NextKeyMarker>root/task/a</NextKeyMarker><NextVersionIdMarker>v1</NextVersionIdMarker></ListVersionsResult>",
version_row("root/task/a", "v1")
));
let (client, worker) = fixture(vec![
history(""),
history(""),
first,
history(&marker_row("root/task/a", "d1")),
deleted(),
deleted(),
history(&version_row("root/task/ab", "sibling")),
]);
let result = delete_candidates(
&client,
&repo,
&json!([{"pointer":"task/a.dvc","object":"task/a","version_id":"v1"}]),
)
.unwrap();
assert_eq!(
result["deleted"][0]["deleted_version_ids"],
json!(["d1", "v1"])
);
let requests = worker.join().unwrap();
let deleted = requests
.iter()
.filter(|request| request.method == "DELETE")
.collect::<Vec<_>>();
assert_eq!(deleted.len(), 2);
assert!(deleted.iter().all(|request| {
request
.target
.starts_with("/fixture-bucket/root/task/a?versionId=")
}));
assert!(
deleted
.iter()
.any(|request| request.target.ends_with("versionId=d1"))
);
assert!(
!deleted
.iter()
.any(|request| request.target.contains("sibling"))
);
}
fn oversized_purge_candidates() -> Value {
// Real candidate fields carry the size: every distinct, exact version ID
// is 480 ASCII bytes, rather than padding an ignored metadata field.
let suffix = "v".repeat(475);
let candidates = (0..6_869)
.map(|index| {
json!({"pointer":"task/a.dvc","object":"task/a","version_id":format!("{index:04}-{suffix}")})
})
.collect::<Vec<_>>();
let payload = Value::Array(candidates);
assert!(serde_json::to_vec(&payload).unwrap().len() > 3 * 1024 * 1024);
payload
}
#[test]
fn large_purge_adapter_deletes_retry_history_without_argv_or_neighbor_deletion() {
use std::sync::{
Arc, Mutex,
atomic::{AtomicUsize, Ordering},
};
let (_directory, repo) = repo();
let payload = oversized_purge_candidates();
let candidates = payload.as_array().unwrap();
assert_eq!(candidates.len(), 6_869);
assert_eq!(
candidates
.iter()
.map(|candidate| candidate["version_id"].as_str().unwrap())
.collect::<BTreeSet<_>>()
.len(),
candidates.len()
);
// All other requested versions have already disappeared during a prior
// attempt. Two payload versions and one delete marker remain to retire.
let present =
[0, 3_000, 6_868].map(|index| candidates[index]["version_id"].as_str().unwrap().to_owned());
let expected_deleted = present.iter().cloned().collect::<BTreeSet<_>>();
let retired = Arc::new(Mutex::new(BTreeSet::new()));
let inventory_reads = Arc::new(AtomicUsize::new(0));
let handler_retired = retired.clone();
let handler_reads = inventory_reads.clone();
let handler_present = present.clone();
// Public adapter registry inspection: versioning + two empty listings (3).
// Native purge: versioning (1), generic registry lookup (2), history (1),
// three exact DELETEs (3), then complete post-delete history verification (1).
let (client, worker) = fixture_handler(11, move |request| {
let url = url::Url::parse(&format!("http://fixture{}", request.target)).unwrap();
let query = url.query_pairs().into_owned().collect::<BTreeMap<_, _>>();
if request.method == "DELETE" {
assert_eq!(url.path(), "/fixture-bucket/root/task/a");
let version = query.get("versionId").unwrap();
assert!(
handler_present.contains(version),
"unrequested deletion: {version}"
);
assert!(handler_retired.lock().unwrap().insert(version.clone()));
return deleted();
}
assert_eq!(request.method, "GET");
if query.contains_key("versioning") {
return versioning();
}
assert!(query.contains_key("versions"));
let prefix = query.get("prefix").unwrap();
if *prefix == registry_object() {
return history("");
}
assert_eq!(prefix, "root/task/a");
let scan = handler_reads.fetch_add(1, Ordering::SeqCst);
let retired = handler_retired.lock().unwrap();
assert_eq!(retired.len(), if scan == 0 { 0 } else { 3 });
let mut rows = String::new();
for version in &handler_present[..2] {
if !retired.contains(version) {
rows.push_str(&version_row("root/task/a", version));
}
}
if !retired.contains(&handler_present[2]) {
rows.push_str(&marker_row("root/task/a", &handler_present[2]));
}
// S3 prefix listings also return this neighboring key. It remains in
// both inventories while the exact target's final history is empty.
rows.push_str(&version_row("root/task/ab", "neighbor-version"));
history(&rows)
});
configure_repo(&client, &repo);
let result = crate::dvc::version_purge_adapter(&repo, "delete", &payload).unwrap();
assert_eq!(result["mode"], "permanent-version-deletion");
assert_eq!(result["remote"], "workspace-mgr");
assert_eq!(result["deleted"].as_array().unwrap().len(), 1);
assert_eq!(result["deleted"][0]["object"], "task/a");
assert_eq!(result["deleted"][0]["pointer"], "task/a.dvc");
assert_eq!(
result["deleted"][0]["deleted_version_ids"],
json!(expected_deleted.iter().cloned().collect::<Vec<_>>())
);
assert_eq!(*retired.lock().unwrap(), expected_deleted);
assert_eq!(inventory_reads.load(Ordering::SeqCst), 2);
// Generic retirement reports per object, including retries with absent IDs.
assert_eq!(result["already_absent"], json!([]));
assert_eq!(result["retained_unmapped"], json!([]));
let requests = worker.join().unwrap();
assert_eq!(requests.len(), 11);
assert_eq!(
requests
.iter()
.filter(|request| request.method == "DELETE")
.count(),
3
);
assert!(
requests
.iter()
.filter(|request| request.method == "DELETE")
.all(|request| {
request
.target
.starts_with("/fixture-bucket/root/task/a?versionId=")
&& !request.target.contains("neighbor-version")
})
);
assert_eq!(requests.last().unwrap().method, "GET");
assert!(
requests
.last()
.unwrap()
.target
.contains("prefix=root%2Ftask%2Fa")
);
}
#[test]
fn large_purge_adapter_validates_last_candidate_before_any_deletion() {
let (_directory, repo) = repo();
let mut payload = oversized_purge_candidates();
payload.as_array_mut().unwrap().last_mut().unwrap()["version_id"] = json!("");
assert_eq!(payload.as_array().unwrap().len(), 6_869);
assert!(serde_json::to_vec(&payload).unwrap().len() > 3 * 1024 * 1024);
// Public adapter inspect (3 requests) and native bucket versioning (1)
// precede validation. No object history or DELETE may be reached.
let (client, worker) = fixture(vec![versioning(), history(""), history(""), versioning()]);
configure_repo(&client, &repo);
let error = crate::dvc::version_purge_adapter(&repo, "delete", &payload)
.unwrap_err()
.to_string();
assert!(error.contains("missing or invalid version_id"), "{error}");
let requests = worker.join().unwrap();
assert_eq!(requests.len(), 4);
assert!(requests.iter().all(|request| request.method == "GET"));
for request in requests {
let url = url::Url::parse(&format!("http://fixture{}", request.target)).unwrap();
let query = url.query_pairs().into_owned().collect::<BTreeMap<_, _>>();
if query.contains_key("versions") {
assert_eq!(query.get("prefix").unwrap(), ®istry_object());
} else {
assert!(query.contains_key("versioning"));
}
}
}
#[test]
fn archive_purge_rejects_missing_coordination_and_unpublished_private_proof() {
let (_directory, repo) = repo();
for proof in [
None,
Some(
json!({"state_path":"/private/journal","mode":"git-cas","transaction_id":"fixture-transaction"}),
),
] {
let receipt = copied_receipt();
let mut payload = json!({"candidates":[{"pointer":"task/.workspace-mgr-archive.json","object":"task/a","version_id":"v1"}]});
if let Some(proof) = proof {
payload["coordination"] = json!([{"receipt":receipt,"coordination":proof}]);
}
let (client, worker) = fixture(vec![
history(&version_row("root/task/a", "v1")),
registry_listing(),
registry_body(&receipt),
registry_listing(),
]);
let error = delete_candidates(&client, &repo, &payload)
.unwrap_err()
.to_string();
assert!(
error.contains("coordination") || error.contains("published"),
"{error}"
);
let requests = worker.join().unwrap();
assert!(!requests.iter().any(|request| request.method == "DELETE"));
}
}
fn file_pointer(repo: &GitRepo) {
fs::create_dir_all(repo.root.join("task")).unwrap();
fs::write(repo.root.join("task/data.dvc"),"outs:\n- md5: 900150983cd24fb0d6963f7d28e17f72\n hash: md5\n size: 3\n path: data\n cloud:\n workspace-mgr:\n version_id: v1\n etag: abc\n").unwrap();
}
fn versioning() -> Reply {
Reply::xml("<VersioningConfiguration><Status>Enabled</Status></VersioningConfiguration>")
}
fn downloaded(bytes: &[u8], version: &str, tag: &str) -> Reply {
Reply {
status: 200,
headers: vec![
("x-amz-version-id", version.to_owned()),
("ETag", format!("\"{tag}\"")),
],
body: bytes.to_vec(),
}
}
#[test]
fn fetch_streams_exact_version_with_ifmatch_and_installs_only_verified_bytes() {
let (_directory, repo) = repo();
let (client, worker) = fixture(vec![versioning(), downloaded(b"abc", "v1", "abc")]);
configure_repo(&client, &repo);
file_pointer(&repo);
let result = read(&repo, &["task/data.dvc".into()], "--fetch", &[]).unwrap();
assert_eq!(result["checked_objects"], json!(["task/data"]));
let path = native_engine::cache_path(&repo, "900150983cd24fb0d6963f7d28e17f72").unwrap();
assert_eq!(fs::read(path).unwrap(), b"abc");
let requests = worker.join().unwrap();
assert_eq!(
requests[1].target,
"/fixture-bucket/root/task/data?versionId=v1"
);
assert_eq!(requests[1].headers["if-match"], "\"abc\"");
}
#[test]
fn exact_fetch_preserves_legacy_binary_and_chunk_boundary_hash_semantics() {
// Fixed digests from DVC 3.67.1's fobj_md5(..., name="md5-dos2unix").
for (body, digest) in [
(b"a\0\r\n".to_vec(), "a200e344e12b35719025ffdb8b428ee8"),
(
[vec![b'a'; (1 << 20) - 1], b"\r\nb".to_vec()].concat(),
"8c5619e93d78797e95a4233d0df49978",
),
] {
let (_directory, repo) = repo();
let (client, worker) = fixture(vec![
versioning(),
downloaded(&body, "v1", "abc"),
versioning(),
head("v1", "abc", body.len() as u64),
]);
configure_repo(&client, &repo);
file_pointer(&repo);
let raw = format!(
"outs:\n- md5: {digest}\n hash: md5-dos2unix\n size: {}\n path: data\n cloud:\n workspace-mgr:\n version_id: v1\n etag: abc\n",
body.len()
);
let pointer = repo.root.join("task/data.dvc");
fs::write(&pointer, &raw).unwrap();
for _ in 0..2 {
read(&repo, &["task/data.dvc".into()], "--fetch", &[]).unwrap();
let cache =
native_engine::cache_path_with_algorithm(&repo, digest, "md5-dos2unix").unwrap();
assert_eq!(fs::read(cache).unwrap(), body);
assert_eq!(fs::read_to_string(&pointer).unwrap(), raw);
}
let requests = worker.join().unwrap();
assert_eq!(requests.len(), 4);
assert_eq!(requests[1].method, "GET");
assert_eq!(requests[3].method, "HEAD");
assert_eq!(requests[1].headers["if-match"], "\"abc\"");
}
}
#[test]
fn fetch_preserves_corrupt_cache_when_new_get_metadata_or_hash_disagrees() {
for reply in [
downloaded(b"abc", "different", "abc"),
downloaded(b"abc", "v1", "wrong"),
downloaded(b"abd", "v1", "abc"),
] {
let (_directory, repo) = repo();
let (client, worker) = fixture(vec![versioning(), reply]);
configure_repo(&client, &repo);
file_pointer(&repo);
let path = native_engine::cache_path(&repo, "900150983cd24fb0d6963f7d28e17f72").unwrap();
fs::create_dir_all(path.parent().unwrap()).unwrap();
fs::write(&path, b"old-corrupt").unwrap();
assert!(read(&repo, &["task/data.dvc".into()], "--fetch", &[]).is_err());
assert_eq!(fs::read(&path).unwrap(), b"old-corrupt");
worker.join().unwrap();
}
}
#[test]
fn valid_cache_bytes_still_require_remote_exact_version() {
let (_directory, repo) = repo();
let (client, worker) = fixture(vec![versioning(), missing("AccessDenied", 403)]);
configure_repo(&client, &repo);
file_pointer(&repo);
native_engine::install_cache(&repo, "900150983cd24fb0d6963f7d28e17f72", b"abc").unwrap();
assert!(
read(&repo, &["task/data.dvc".into()], "--fetch", &[])
.unwrap_err()
.to_string()
.contains("AccessDenied")
);
let requests = worker.join().unwrap();
assert_eq!(requests[1].method, "HEAD");
assert_eq!(requests.len(), 2);
}
#[test]
fn legacy_cache_layout_needs_exact_head_verification_without_a_get() {
let (_directory, repo) = repo();
let (client, worker) = fixture(vec![versioning(), head("v1", "abc", 3)]);
configure_repo(&client, &repo);
file_pointer(&repo);
let digest = "900150983cd24fb0d6963f7d28e17f72";
let legacy = native_engine::cache_root(&repo)
.unwrap()
.join("90/0150983cd24fb0d6963f7d28e17f72");
fs::create_dir_all(legacy.parent().unwrap()).unwrap();
fs::write(&legacy, b"abc").unwrap();
assert!(!native_engine::cache_path(&repo, digest).unwrap().exists());
read(&repo, &["task/data.dvc".into()], "--fetch", &[]).unwrap();
assert_eq!(fs::read(&legacy).unwrap(), b"abc");
let requests = worker.join().unwrap();
assert_eq!(requests.len(), 2);
assert_eq!(requests[1].method, "HEAD");
assert_eq!(
requests[1].target,
"/fixture-bucket/root/task/data?versionId=v1"
);
}
#[test]
fn omitted_legacy_hash_preserves_cache_namespace_and_exact_payload_bytes() {
let (_directory, repo) = repo();
let body = b"a\r\nb\r\n";
let digest = "dd8c6a395b5dd36c56d23275028f526c";
let (client, worker) = fixture(vec![
versioning(),
downloaded(body, "v1", "abc"),
versioning(),
head("v1", "abc", 6),
]);
configure_repo(&client, &repo);
file_pointer(&repo);
// Real DVC 2 pointers omit hash. Their normalized digest also names a
// different raw-MD5 cache object; preserve both byte representations.
let raw = format!(
"outs:\n- md5: {digest}\n size: 6\n path: data\n cloud:\n workspace-mgr:\n version_id: v1\n etag: abc\n"
);
let pointer = repo.root.join("task/data.dvc");
fs::write(&pointer, &raw).unwrap();
native_engine::install_cache(&repo, digest, b"a\nb\n").unwrap();
let canonical = native_engine::cache_path(&repo, digest).unwrap();
let legacy = native_engine::cache_path_with_algorithm(&repo, digest, "md5-dos2unix").unwrap();
for _ in 0..2 {
read(&repo, &["task/data.dvc".into()], "--fetch", &[]).unwrap();
assert_eq!(fs::read(&canonical).unwrap(), b"a\nb\n");
assert_eq!(fs::read(&legacy).unwrap(), body);
assert_eq!(fs::read_to_string(&pointer).unwrap(), raw);
}
let requests = worker.join().unwrap();
assert_eq!(requests.len(), 4);
assert_eq!(requests[1].method, "GET");
assert_eq!(requests[3].method, "HEAD");
}
fn published_proof(repo: &GitRepo, remote: &Path, receipt: &Value) -> Value {
repo.run(["init", "-q", "-b", "main"]).unwrap();
repo.run(["config", "user.name", "Fixture"]).unwrap();
repo.run(["config", "user.email", "fixture@example.invalid"])
.unwrap();
repo.run(["init", "-q", "--bare", remote.to_str().unwrap()])
.unwrap();
repo.run(["remote", "add", "origin", remote.to_str().unwrap()])
.unwrap();
fs::write(
repo.root.join(".workspace-mgr.toml"),
"[git]\nremote='origin'\nbranch='main'\n",
)
.unwrap();
fs::create_dir_all(repo.root.join("archive/task")).unwrap();
fs::write(
repo.root.join("archive/task/.workspace-mgr-archive.json"),
serde_json::to_vec_pretty(receipt).unwrap(),
)
.unwrap();
repo.run(["add", "."]).unwrap();
repo.run(["commit", "-q", "-m", "published fixture"])
.unwrap();
repo.run(["push", "-q", "origin", "main"]).unwrap();
crate::archive_registry::coordinate_published(repo, receipt).unwrap()
}
#[test]
fn published_archive_purge_deletes_only_mapped_versions_and_retains_new_writes() {
use std::sync::{
Arc,
atomic::{AtomicBool, Ordering},
};
let (directory, repo) = repo();
let mut receipt = copied_receipt();
receipt["versions"][0]["source_is_latest"] = false.into();
receipt["versions"].as_array_mut().unwrap().push(json!({"source_object":"task/a","destination_object":"archive/task/a","source_version_id":"d1","source_last_modified":"2026-10-07T00:00:01+00:00","source_is_latest":true,"source_list_order":1,"delete_marker":true,"size":null,"source_etag":null,"destination_version_id":"dst-delete","destination_etag":null,"destination_last_modified":"2026-10-07T00:00:02+00:00"}));
let proof = published_proof(&repo, &directory.path().join("remote.git"), &receipt);
let deleted_payload = Arc::new(AtomicBool::new(false));
let deleted_marker = Arc::new(AtomicBool::new(false));
let payload_flag = deleted_payload.clone();
let marker_flag = deleted_marker.clone();
let stored_receipt = receipt.clone();
// Initial list + canonical/history checks before each of two DELETEs +
// final complete rescan. Every request targets this private loopback.
let (client, worker) = fixture_handler(19, move |request| {
if request.method == "DELETE" {
if request.target.ends_with("versionId=v1") {
payload_flag.store(true, Ordering::SeqCst);
} else if request.target.ends_with("versionId=d1") {
marker_flag.store(true, Ordering::SeqCst);
} else {
panic!("attempted unmapped version deletion: {}", request.target);
}
return deleted();
}
if request.method == "HEAD" {
return head("dst", "copied", 3);
}
if request.target.contains("versions=") {
if request
.target
.contains("prefix=root%2F.workspace-mgr%2Farchive%2F")
{
return registry_listing();
}
if request.target.contains("prefix=root%2Farchive%2Ftask%2F") {
return history(&format!(
"{}{}",
version_row("root/archive/task/a", "dst")
.replace(""abc"", ""copied""),
marker_row("root/archive/task/a", "dst-delete")
));
}
let mut rows = String::new();
if !payload_flag.load(Ordering::SeqCst) {
rows.push_str(&version_row("root/task/a", "v1"));
}
if !marker_flag.load(Ordering::SeqCst) {
rows.push_str(&marker_row("root/task/a", "d1"));
}
rows.push_str(&version_row("root/task/b", "independent-new-version"));
rows.push_str(&marker_row("root/task/b", "independent-new-marker"));
return history(&rows);
}
registry_body(&stored_receipt)
});
let payload = json!({"candidates":[{"pointer":"task/.workspace-mgr-archive.json","object":"task/a","version_id":"v1"},{"pointer":"task/.workspace-mgr-archive.json","object":"task/a","version_id":"d1"}],"prefixes":[receipt],"coordination":[{"receipt":receipt,"coordination":proof}]});
let result = delete_candidates(&client, &repo, &payload).unwrap();
assert!(deleted_payload.load(Ordering::SeqCst));
assert!(deleted_marker.load(Ordering::SeqCst));
assert_eq!(result["retained_unmapped"].as_array().unwrap().len(), 2);
assert_eq!(result["cleaned_prefixes"], json!([]));
let requests = worker.join().unwrap();
assert_eq!(
requests
.iter()
.filter(|request| request.method == "DELETE")
.count(),
2
);
}