diff --git a/Cargo.lock b/Cargo.lock index 9bb6057b5..ef3e723f3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5746,6 +5746,7 @@ dependencies = [ "knot-types", "lexicons", "parking_lot", + "proptest", "serde", "serde_bytes", "serde_ipld_dagcbor", diff --git a/knot2/crates/knot-sim/src/harness.rs b/knot2/crates/knot-sim/src/harness.rs index d4228cadd..0200fded6 100644 --- a/knot2/crates/knot-sim/src/harness.rs +++ b/knot2/crates/knot-sim/src/harness.rs @@ -780,6 +780,7 @@ fn assemble_router(parts: StateParts) -> Router { appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), slots: knot_resource::Slots::testing(8), lfs: None, + blobs: Default::default(), catalog: Arc::new(knot_messages::Catalog::defaults()), rkeys: Default::default(), }); diff --git a/knot2/crates/knot-sim/tests/lfs_chaos.rs b/knot2/crates/knot-sim/tests/lfs_chaos.rs index ba2a61298..6d7e3e196 100644 --- a/knot2/crates/knot-sim/tests/lfs_chaos.rs +++ b/knot2/crates/knot-sim/tests/lfs_chaos.rs @@ -6,11 +6,10 @@ use std::time::{Duration, Instant, SystemTime}; use knot_cob::{CobHome, CobStore}; use knot_cobs::{Registration, RegistryChange, RepoRef, RepoRegistryCob, deregister_repo}; use knot_git::{Layout, Repo}; -use knot_lfs::{ClaimedSize, DiskStore, LfsOid, LfsStore, LfsStorePath}; +use knot_lfs::{ClaimedSize, DiskStore, LfsOid, LfsStore, LfsStorePath, Pool}; use knot_runtime::OsEntropy; use knot_secrets::{MasterKey, SealedStore}; use knot_types::{KnotId, OwnerDid, RepoDid, RepoName, RepoRkey, UnixSeconds}; -use sha2::{Digest, Sha256}; const REPO_DID: &str = "did:plc:squid"; const REPO_NAME: &str = "anemone"; @@ -33,7 +32,7 @@ fn master() -> MasterKey { } fn media_oid() -> LfsOid { - LfsOid::from_digest(Sha256::digest(MEDIA).into()) + LfsOid::from_digest(knot_types::Sha256Digest::hash(MEDIA)) } struct Paths { @@ -83,13 +82,18 @@ fn build_fixture(root: &Path) -> Paths { let store = DiskStore::open(LfsStorePath::new(&paths.store)).unwrap(); store .put( + Pool::Lfs, &did(), &media_oid(), ClaimedSize::new(MEDIA.len() as u64), &mut &MEDIA[..], ) .unwrap(); - let object_path = store.object_file(&did(), &media_oid()).unwrap().unwrap().1; + let object_path = store + .object_file(Pool::Lfs, &did(), &media_oid()) + .unwrap() + .unwrap() + .1; std::fs::OpenOptions::new() .write(true) .open(object_path) @@ -162,7 +166,7 @@ fn kill9_between_delete_steps_never_strands_the_store_prefix() { let full = started.elapsed(); let warm_store = DiskStore::open(LfsStorePath::new(&warm.store)).unwrap(); assert_eq!( - warm_store.probe(&did(), &media_oid()).unwrap(), + warm_store.probe(Pool::Lfs, &did(), &media_oid()).unwrap(), None, "an uninterrupted delete removes the store prefix itself" ); @@ -202,7 +206,7 @@ fn kill9_between_delete_steps_never_strands_the_store_prefix() { .unwrap_or_else(|error| { panic!("trial {trial}: orphan sweep must run after a killed delete: {error}") }); - let present = store.probe(&did(), &media_oid()).unwrap().is_some(); + let present = store.probe(Pool::Lfs, &did(), &media_oid()).unwrap().is_some(); match registered { true => assert!( present, diff --git a/knot2/crates/knot-sim/tests/lfs_roundtrip.rs b/knot2/crates/knot-sim/tests/lfs_roundtrip.rs index e04388b08..6e2aa8b5b 100644 --- a/knot2/crates/knot-sim/tests/lfs_roundtrip.rs +++ b/knot2/crates/knot-sim/tests/lfs_roundtrip.rs @@ -15,7 +15,7 @@ use knot_cob::{CobHome, CobStore}; use knot_cobs::{Registration, RegistryChange}; use knot_edge::RequiresFullHandshake; use knot_git::{Layout, Repo}; -use knot_lfs::{FreeSpaceFloor, LfsHandle, LfsOid, LfsSize, LfsStore, LfsStorePath}; +use knot_lfs::{FreeSpaceFloor, LfsHandle, LfsOid, LfsSize, LfsStore, LfsStorePath, Pool}; use knot_runtime::{ FakeHttp, HttpResponse, K256Signer, ManualClock, OsEntropy, Signer, UnixMicros, }; @@ -24,7 +24,6 @@ use knot_types::{ AccountDid, AdmissionPolicy, AuthorName, Email, KnotHostname, KnotId, OwnerDid, RepoDid, RepoName, RepoRkey, UnixSeconds, }; -use sha2::{Digest, Sha256}; use tempfile::TempDir; use tokio::net::TcpListener; use tower::ServiceExt; @@ -392,6 +391,7 @@ async fn spawn(published_line: String, with_h3: bool) -> World { appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), slots: knot_resource::Slots::testing(8), lfs: Some(knot_xrpc::LfsWeb::new(lfs.clone(), 8)), + blobs: Default::default(), catalog: Arc::new(knot_messages::Catalog::defaults()), rkeys: Default::default(), }); @@ -470,7 +470,7 @@ async fn spawn(published_line: String, with_h3: bool) -> World { ssh_port, http_base, router, - _imports_stop: _imports_stop, + _imports_stop, layout, h3, _certdir: certdir, @@ -561,7 +561,7 @@ async fn the_lfs_round_trip_gate_holds_over_both_transports_and_the_fork() { let world = spawn_world(public_line).await; let media = media_bytes(); - let media_oid = LfsOid::from_digest(Sha256::digest(&media).into()); + let media_oid = LfsOid::from_digest(knot_types::Sha256Digest::hash(&media)); let ssh = format!("ssh -i {key_path} {}", knot_fixtures::UNSHARED_SSH); let path_env = std::env::var("PATH").unwrap_or_default(); let home = scratch.path().to_str().unwrap().to_string(); @@ -589,7 +589,11 @@ async fn the_lfs_round_trip_gate_holds_over_both_transports_and_the_fork() { let source_repo = RepoDid::new(REPO_DID).unwrap(); assert_eq!( - world.lfs.store.probe(&source_repo, &media_oid).unwrap(), + world + .lfs + .store + .probe(Pool::Lfs, &source_repo, &media_oid) + .unwrap(), Some(LfsSize::new(media.len() as u64)), "pushed media must be durable in the store" ); @@ -619,7 +623,12 @@ async fn the_lfs_round_trip_gate_holds_over_both_transports_and_the_fork() { let fork_did = RepoDid::new(created["repoDid"].as_str().unwrap()).unwrap(); let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(30); let forked = loop { - if let Some(size) = world.lfs.store.probe(&fork_did, &media_oid).unwrap() { + if let Some(size) = world + .lfs + .store + .probe(Pool::Lfs, &fork_did, &media_oid) + .unwrap() + { break size; } assert!( @@ -827,8 +836,8 @@ async fn the_lfs_stack_is_conformant_with_the_reference_server_and_client() { let media = media_bytes(); let second = second_media_bytes(); - let media_oid = LfsOid::from_digest(Sha256::digest(&media).into()); - let second_oid = LfsOid::from_digest(Sha256::digest(&second).into()); + let media_oid = LfsOid::from_digest(knot_types::Sha256Digest::hash(&media)); + let second_oid = LfsOid::from_digest(knot_types::Sha256Digest::hash(&second)); let expected: std::collections::BTreeSet<(String, u64)> = [ (media_oid.as_str().to_string(), media.len() as u64), (second_oid.as_str().to_string(), second.len() as u64), @@ -909,7 +918,7 @@ async fn the_lfs_stack_is_conformant_with_the_reference_server_and_client() { let knot_objects: std::collections::BTreeSet<(String, u64)> = world .lfs .store - .enumerate(&source_repo) + .enumerate(Pool::Lfs, &source_repo) .unwrap() .into_iter() .map(|object| (object.oid.as_str().to_string(), object.size.get())) @@ -960,7 +969,7 @@ async fn the_lfs_stack_is_conformant_with_the_reference_server_and_client() { let knot_removed = world .lfs .store - .object_file(&source_repo, &second_oid) + .object_file(Pool::Lfs, &source_repo, &second_oid) .unwrap() .unwrap() .1; @@ -1087,7 +1096,7 @@ async fn lfs_http_push_stores_an_object_with_a_push_token() { let payload: Vec = (0..4096u32) .map(|n| (n.wrapping_mul(17) % 251) as u8) .collect(); - let oid = LfsOid::from_digest(Sha256::digest(&payload).into()); + let oid = LfsOid::from_digest(knot_types::Sha256Digest::hash(&payload)); let size = payload.len() as u64; let bearer = format!( @@ -1118,7 +1127,12 @@ async fn lfs_http_push_stores_an_object_with_a_push_token() { "upload objects mustn't claim authenticated=true, else git-lfs sends the object put with no auth and loops on 401: {body}" ); assert!( - world.lfs.store.probe(&repo, &oid).unwrap().is_none(), + world + .lfs + .store + .probe(Pool::Lfs, &repo, &oid) + .unwrap() + .is_none(), "object must be absent before the put" ); @@ -1129,7 +1143,7 @@ async fn lfs_http_push_stores_an_object_with_a_push_token() { "the same push token must authorize both the batch and the object put" ); assert_eq!( - world.lfs.store.probe(&repo, &oid).unwrap(), + world.lfs.store.probe(Pool::Lfs, &repo, &oid).unwrap(), Some(LfsSize::new(size)), "the put object must be durable in the store" ); @@ -1160,7 +1174,7 @@ async fn lfs_http_push_stores_an_object_with_a_push_token() { let payload2: Vec = (0..2048u32) .map(|n| (n.wrapping_mul(29) % 251) as u8) .collect(); - let oid2 = LfsOid::from_digest(Sha256::digest(&payload2).into()); + let oid2 = LfsOid::from_digest(knot_types::Sha256Digest::hash(&payload2)); let basic = basic_auth(&service_jwt("sh.tangled.repo.push", "lfs-http-push-2")); let (status, body) = lfs_batch(&world, Some(&basic), "upload", &oid2, payload2.len() as u64).await; @@ -1176,7 +1190,7 @@ async fn lfs_http_push_stores_an_object_with_a_push_token() { "a push token presented as the http basic password must authenticate the put" ); assert_eq!( - world.lfs.store.probe(&repo, &oid2).unwrap(), + world.lfs.store.probe(Pool::Lfs, &repo, &oid2).unwrap(), Some(LfsSize::new(payload2.len() as u64)), "the basic-authenticated object must be durable too" ); @@ -1191,7 +1205,7 @@ async fn lfs_http_push_rejects_missing_and_mismatched_credentials() { let payload: Vec = (0..1024u32) .map(|n| (n.wrapping_mul(13) % 251) as u8) .collect(); - let oid = LfsOid::from_digest(Sha256::digest(&payload).into()); + let oid = LfsOid::from_digest(knot_types::Sha256Digest::hash(&payload)); let size = payload.len() as u64; let (status, _) = lfs_batch(&world, None, "upload", &oid, size).await; @@ -1233,7 +1247,7 @@ async fn lfs_http_push_rejects_missing_and_mismatched_credentials() { world .lfs .store - .probe(&RepoDid::new(REPO_DID).unwrap(), &oid) + .probe(Pool::Lfs, &RepoDid::new(REPO_DID).unwrap(), &oid) .unwrap() .is_none(), "no rejected request may leave an object behind" diff --git a/knot2/crates/knot-ssh/tests/ssh_push.rs b/knot2/crates/knot-ssh/tests/ssh_push.rs index 7a9785341..1b4e753c6 100644 --- a/knot2/crates/knot-ssh/tests/ssh_push.rs +++ b/knot2/crates/knot-ssh/tests/ssh_push.rs @@ -424,6 +424,7 @@ async fn spawn_server_core( appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), slots: knot_resource::Slots::testing(8), lfs: None, + blobs: Default::default(), catalog: Arc::new(knot_messages::Catalog::defaults()), firehose: Arc::new(knot_events::EventLog::new( ManualClock::new(UnixMicros::new(1_000_000_000)), @@ -2023,8 +2024,7 @@ fn trickled_lfs_upload( #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn shutdown_drains_an_in_flight_lfs_transfer_before_exit() { - use knot_lfs::LfsStore; - use sha2::Digest; + use knot_lfs::{LfsStore, Pool}; let scratch = tempfile::tempdir().unwrap(); let (key_path, public_line) = keygen(scratch.path(), "drain"); let lfs_dir = scratch.path().join("lfs"); @@ -2046,7 +2046,7 @@ async fn shutdown_drains_an_in_flight_lfs_transfer_before_exit() { .await; let body: Vec = (0..1_048_576u32).map(|n| (n % 251) as u8).collect(); - let oid = knot_lfs::LfsOid::from_digest(sha2::Sha256::digest(&body).into()); + let oid = knot_lfs::LfsOid::from_digest(knot_types::Sha256Digest::hash(&body)); let (midway_tx, midway_rx) = std::sync::mpsc::channel(); let client = { @@ -2092,7 +2092,7 @@ async fn shutdown_drains_an_in_flight_lfs_transfer_before_exit() { assert_eq!( handle .store - .probe(&repo_did, &oid) + .probe(Pool::Lfs, &repo_did, &oid) .unwrap() .map(|size| size.get()), Some(1_048_576), diff --git a/knot2/crates/knot-xrpc/Cargo.toml b/knot2/crates/knot-xrpc/Cargo.toml index 395797c70..128032fd2 100644 --- a/knot2/crates/knot-xrpc/Cargo.toml +++ b/knot2/crates/knot-xrpc/Cargo.toml @@ -65,3 +65,4 @@ sha2 = { workspace = true } tokio-tungstenite = "0.29" tokio = { workspace = true, features = ["process"] } parking_lot = { workspace = true } +proptest = { workspace = true } diff --git a/knot2/crates/knot-xrpc/tests/blob_properties.proptest-regressions b/knot2/crates/knot-xrpc/tests/blob_properties.proptest-regressions new file mode 100644 index 000000000..bd2aeefe4 --- /dev/null +++ b/knot2/crates/knot-xrpc/tests/blob_properties.proptest-regressions @@ -0,0 +1,7 @@ +# Here lie some seeds for failure cases that I have already generated with the proptester. +# The proptester will start by reading these and running them before generating any new cases! +# +# Please don't mess with these existing cases unless they become unnecessary. +cc 54ac38076efa69c847c604527a729aee1b2f4c713d51d99bcfe9647a749c28bf # which condenses to ops = [Upload(13)] +cc c671907a763ce7005394279834692fc14ee75200aa4f1fdc35792e9a2bf3fbcc # condenses to ops = [Upload(16), Upload(6), Upload(0), Upload(1), Upload(2), Upload(3), Upload(4)] +cc 32078ad37be8c22b598242f5235d8c065b64f88097a37bfe9db908fc017ac3b9 # shrinks to ops = [Upload(2), Create(2), Delete(0)] diff --git a/knot2/crates/knot-xrpc/tests/blob_properties.rs b/knot2/crates/knot-xrpc/tests/blob_properties.rs new file mode 100644 index 000000000..89b1690f5 --- /dev/null +++ b/knot2/crates/knot-xrpc/tests/blob_properties.rs @@ -0,0 +1,365 @@ +mod common; + +use std::collections::BTreeMap; + +use axum::body::Body; +use common::{OWNER, World, empty_repo, enable_chain, get, post_as, service_jwt_for}; +use futures::{StreamExt, TryStreamExt}; +use http::{Request, StatusCode, header}; +use jacquard_common::types::tid::Tid; +use knot_types::{BlobCid, RecordRkey, RepoDid, Sha256Digest}; +use proptest::prelude::*; +use serde_json::json; +use tower::ServiceExt; + +const MAGICS: [(&[u8], &str); 8] = [ + ( + &[0x89, b'P', b'N', b'G', 0x0d, 0x0a, 0x1a, 0x0a], + "image/png", + ), + (&[0xFF, 0xD8, 0xFF, 0xE0], "image/jpeg"), + (b"GIF87a", "image/gif"), + ( + &[b'R', b'I', b'F', b'F', 1, 2, 3, 4, b'W', b'E', b'B', b'P'], + "image/webp", + ), + (&[b'B', b'M', 1, 2, 3, 4], "image/bmp"), + (b" Vec { + let (magic, _) = MAGICS[(token % MAGICS.len() as u16) as usize]; + let mut body = magic.to_vec(); + body.extend_from_slice(&token.to_be_bytes()); + body.extend(std::iter::repeat_n(0x5a, (token % 7) as usize)); + body +} + +fn mime_of(token: u16) -> &'static str { + MAGICS[(token % MAGICS.len() as u16) as usize].1 +} + +fn blob_cid(bytes: &[u8]) -> BlobCid { + BlobCid::from_digest(Sha256Digest::hash(bytes)) +} + +fn rkey_of(response: &serde_json::Value) -> RecordRkey { + RecordRkey::new( + response["uri"] + .as_str() + .expect("the write response has a URI") + .rsplit('/') + .next() + .expect("a URI ends in the rkey"), + ) + .expect("a URI's tail is an rkey") +} + +fn commit_cid_of(response: &serde_json::Value) -> cid::Cid { + cid::Cid::try_from( + response["cid"] + .as_str() + .expect("the write response has a commit"), + ) + .expect("the commit CID parses") +} + +#[derive(Clone, Debug)] +enum Op { + Upload(u16), + Create(u16), + Edit(u16, u16), + Delete(u16), +} + +fn op_strategy() -> impl Strategy { + prop_oneof![ + (0u16..24).prop_map(Op::Upload), + (0u16..24, 0u16..24).prop_map(|(a, b)| Op::Create(if b % 4 == 0 { a } else { a % 5 })), + (0u16..24, 0u16..24, 0u16..24).prop_map(|(record, a, b)| { + Op::Edit(record % 6, if b % 4 == 0 { a } else { a % 5 }) + }), + (0u16..24).prop_map(|record| Op::Delete(record % 6)), + ] +} + +#[derive(Clone)] +struct LiveRecord { + swap: cid::Cid, + rev: Tid, + blob: u16, +} + +struct Model { + uploaded: BTreeMap>, + live: Vec<(RecordRkey, LiveRecord)>, +} + +impl Model { + fn blob_json(&self, token: u16) -> serde_json::Value { + json!({ + "$type": "blob", + "ref": { "$link": blob_cid(&self.uploaded[&token]).to_string() }, + "mimeType": mime_of(token), + "size": self.uploaded[&token].len(), + }) + } + + fn at(&self, index: usize) -> Option<(&RecordRkey, &LiveRecord)> { + self.live.get(index).map(|(rkey, record)| (rkey, record)) + } + + fn expected_listing(&self) -> Vec { + let mut latest: BTreeMap = BTreeMap::new(); + self.live + .iter() + .for_each(|(_, record)| match latest.get(&record.blob) { + Some(seen) if seen.as_str() >= record.rev.as_str() => (), + _ => { + latest.insert(record.blob, &record.rev); + } + }); + let mut ordered: Vec<(&Tid, u16)> = latest + .into_iter() + .map(|(token, rev)| (rev, token)) + .collect(); + ordered.sort_by(|(a, _), (b, _)| a.as_str().cmp(b.as_str())); + ordered + .into_iter() + .map(|(_, token)| blob_cid(&self.uploaded[&token])) + .collect() + } +} + +async fn upload(world: &World, repo: &RepoDid, body: &[u8]) -> (StatusCode, serde_json::Value) { + let request = Request::post("/xrpc/com.atproto.repo.uploadBlob") + .header( + header::AUTHORIZATION, + format!( + "Bearer {}", + service_jwt_for("com.atproto.repo.uploadBlob", OWNER, repo.as_str()) + ), + ) + .header(header::CONTENT_TYPE, "application/octet-stream") + .body(Body::from(body.to_vec())) + .unwrap(); + let response = world.router.clone().oneshot(request).await.unwrap(); + let status = response.status(); + let bytes = axum::body::to_bytes(response.into_body(), usize::MAX) + .await + .unwrap(); + ( + status, + serde_json::from_slice(&bytes).unwrap_or(serde_json::Value::Null), + ) +} + +async fn drive(ops: Vec) -> Result<(), TestCaseError> { + let world = World::unshed(); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + enable_chain(&world, &did); + let repo = did.as_str(); + let model = Model { + uploaded: BTreeMap::new(), + live: Vec::new(), + }; + + let model = futures::stream::iter(ops) + .map(Ok) + .try_fold(model, |mut model, op| { + let (world, did, repo) = (&world, &did, repo); + async move { + match op { + Op::Upload(token) => { + let body = body_of(token); + let (status, response) = upload(world, did, &body).await; + prop_assert_eq!(status, StatusCode::OK, "{}", response); + model.uploaded.entry(token).or_insert(body); + } + Op::Create(token) => { + if !model.uploaded.contains_key(&token) { + return Ok(model); + } + let (status, written) = post_as( + world, + "/xrpc/com.atproto.repo.createRecord", + OWNER, + repo, + json!({ + "repo": repo, + "collection": "sh.tangled.repo.issue", + "record": { + "$type": "sh.tangled.repo.issue", + "repo": repo, + "title": "kelp under test", + "blobs": [model.blob_json(token)], + "createdAt": "2026-10-01T00:00:00Z", + }, + }), + ) + .await; + prop_assert_eq!(status, StatusCode::OK, "{}", written); + model.live.push(( + rkey_of(&written), + LiveRecord { + swap: commit_cid_of(&written), + rev: common::commit_rev_of(&written), + blob: token, + }, + )); + } + Op::Edit(index, token) => { + let (rkey, swap) = match model.at(index as usize) { + Some((rkey, record)) => (rkey.clone(), record.swap.clone()), + None => return Ok(model), + }; + if !model.uploaded.contains_key(&token) { + return Ok(model); + } + let (status, written) = post_as( + world, + "/xrpc/com.atproto.repo.putRecord", + OWNER, + repo, + json!({ + "repo": repo, + "collection": "sh.tangled.repo.issue", + "rkey": rkey.to_string(), + "swapRecord": swap.to_string(), + "record": { + "$type": "sh.tangled.repo.issue", + "repo": repo, + "title": "kelp under test, again", + "blobs": [model.blob_json(token)], + "createdAt": "2026-10-01T00:00:00Z", + }, + }), + ) + .await; + prop_assert_eq!(status, StatusCode::OK, "{}", written); + let replaced = rkey_of(&written); + if let Some((_, record)) = + model.live.iter_mut().find(|(rkey, _)| *rkey == replaced) + { + record.swap = commit_cid_of(&written); + record.rev = common::commit_rev_of(&written); + record.blob = token; + } + } + Op::Delete(index) => { + let (rkey, swap) = match model.at(index as usize) { + Some((rkey, record)) => (rkey.clone(), record.swap.clone()), + None => return Ok(model), + }; + let (status, response) = post_as( + world, + "/xrpc/com.atproto.repo.deleteRecord", + OWNER, + repo, + json!({ + "repo": repo, + "collection": "sh.tangled.repo.issue", + "rkey": rkey.to_string(), + "swapRecord": swap.to_string(), + }), + ) + .await; + prop_assert_eq!(status, StatusCode::OK, "{}", response); + model.live.retain(|(key, _)| key != &rkey); + } + } + + let (status, _headers, listed) = get( + world, + &format!("/xrpc/com.atproto.sync.listBlobs?did={repo}"), + ) + .await; + let listed: serde_json::Value = serde_json::from_slice(&listed).unwrap_or_default(); + prop_assert_eq!(status, StatusCode::OK); + let served: Vec = listed["cids"] + .as_array() + .map(|cids| cids.iter().map(common::blob_cid_of).collect()) + .unwrap_or_default(); + prop_assert_eq!( + served, + model.expected_listing(), + "The listing holds one blob per live record, oldest reference first." + ); + Ok(model) + } + }) + .await?; + + futures::stream::iter(&model.uploaded) + .map(Ok) + .try_for_each(|(token, bytes)| { + let (world, repo) = (&world, repo); + async move { + let (status, headers, served) = get( + world, + &format!( + "/xrpc/com.atproto.sync.getBlob?did={repo}&cid={}", + blob_cid(bytes) + ), + ) + .await; + prop_assert_eq!(status, StatusCode::OK, "Blob {} serves.", token); + prop_assert_eq!(served.as_ref(), bytes.as_slice()); + prop_assert_eq!(headers.get(header::CONTENT_TYPE).unwrap(), mime_of(*token)); + Ok(()) + } + }) + .await?; + + futures::stream::iter(&model.live) + .map(Ok) + .try_for_each(|(rkey, record)| { + let (world, repo) = (&world, repo); + async move { + let (status, _headers, response) = get( + world, + &format!( + "/xrpc/com.atproto.repo.getRecord?repo={repo}&collection=sh.tangled.repo.issue&rkey={rkey}" + ), + ) + .await; + let response: serde_json::Value = + serde_json::from_slice(&response).unwrap_or_default(); + prop_assert_eq!(status, StatusCode::OK, "{}", response); + prop_assert_eq!( + response["value"]["blobs"] + .as_array() + .map(|blobs| blobs.len()), + Some(1), + "{} has blob {} until it is deleted.", + rkey, + record.blob + ); + Ok(()) + } + }) + .await?; + + let git = world.layout.open(&did).unwrap(); + let warm = world.state.blobs.get(&git, &did).unwrap(); + let cold = knot_xrpc::BlobViews::default().get(&git, &did).unwrap(); + prop_assert_eq!(&warm.listed, &cold.listed); + prop_assert_eq!(&warm.referenced, &cold.referenced); + prop_assert_eq!(warm.coverage, cold.coverage); + Ok(()) +} + +proptest! { + #![proptest_config(ProptestConfig::with_cases(24))] + #[test] + fn blob_listings_and_serves_follow_random_write_history( + ops in proptest::collection::vec(op_strategy(), 0..24) + ) { + tokio::runtime::Runtime::new().unwrap().block_on(drive(ops))?; + } +} diff --git a/knot2/crates/knot-xrpc/tests/common/mod.rs b/knot2/crates/knot-xrpc/tests/common/mod.rs index 5f825366f..e070d9d42 100644 --- a/knot2/crates/knot-xrpc/tests/common/mod.rs +++ b/knot2/crates/knot-xrpc/tests/common/mod.rs @@ -151,6 +151,27 @@ pub fn fixture_collection() -> knot_types::RecordCollection { knot_types::RecordCollection::new(FIXTURE_COLLECTION).unwrap() } +pub fn blob_cid_of(value: &serde_json::Value) -> knot_types::BlobCid { + let sent = value.as_str().expect("a blob CID is a string"); + knot_types::BlobCid::new(cid::Cid::try_from(sent).expect("the CID parses")) + .expect("the CID is raw sha2-256") +} + +pub fn listed_cids(response: &serde_json::Value) -> Vec { + response["cids"] + .as_array() + .map(|cids| cids.iter().map(blob_cid_of).collect()) + .unwrap_or_default() +} + +pub fn commit_rev_of(response: &serde_json::Value) -> jacquard_common::types::tid::Tid { + jacquard_common::types::tid::Tid::new( + response["commit"]["rev"] + .as_str() + .expect("the write response has a rev"), + ) + .expect("the rev parses as a TID") +} #[derive(serde::Serialize)] struct FixtureRecord { #[serde(rename = "$type")] @@ -239,6 +260,20 @@ impl World { ObjectFormat::SHA1, LimitConfig::default(), BTreeSet::from([AccountDid::new(OWNER).unwrap()]), + true, + knot_types::ContributionPolicy::Anyone, + ) + } + + pub fn under(policy: knot_types::ContributionPolicy) -> Self { + Self::assemble( + true, + ByteLimits::default(), + ObjectFormat::SHA1, + LimitConfig::default(), + BTreeSet::new(), + true, + policy, ) } @@ -246,13 +281,33 @@ impl World { Self::build_with_limits(rebuilt, byte_limits, object_format, LimitConfig::default()) } + pub fn lfs_less() -> Self { + Self::assemble( + true, + ByteLimits::default(), + ObjectFormat::SHA1, + LimitConfig::default(), + BTreeSet::new(), + false, + knot_types::ContributionPolicy::Anyone, + ) + } + fn build_with_limits( rebuilt: bool, byte_limits: ByteLimits, object_format: ObjectFormat, limits: LimitConfig, ) -> Self { - Self::assemble(rebuilt, byte_limits, object_format, limits, BTreeSet::new()) + Self::assemble( + rebuilt, + byte_limits, + object_format, + limits, + BTreeSet::new(), + true, + knot_types::ContributionPolicy::Anyone, + ) } fn assemble( @@ -261,6 +316,8 @@ impl World { object_format: ObjectFormat, limits: LimitConfig, admins: BTreeSet, + lfs_enabled: bool, + policy: knot_types::ContributionPolicy, ) -> Self { let dir = tempfile::tempdir().unwrap(); let scan_path = dir.path().join("repos"); @@ -314,13 +371,18 @@ impl World { secrets.ensure(&knot).unwrap(); let lfs_store = dir.path().join("lfs"); - std::fs::create_dir_all(&lfs_store).unwrap(); - let lfs_handle = knot_lfs::LfsHandle::open( - knot_lfs::LfsStorePath::new(&lfs_store), - knot_lfs::LfsSize::new(64 * 1024 * 1024), - knot_lfs::FreeSpaceFloor::new(0), - ) - .unwrap(); + let lfs = lfs_enabled.then(|| { + std::fs::create_dir_all(&lfs_store).unwrap(); + knot_xrpc::LfsWeb::new( + knot_lfs::LfsHandle::open( + knot_lfs::LfsStorePath::new(&lfs_store), + knot_lfs::LfsSize::new(64 * 1024 * 1024), + knot_lfs::FreeSpaceFloor::new(0), + ) + .unwrap(), + 8, + ) + }); let (imports, import_jobs) = knot_xrpc::imports::ImportQueue::channel(); let shutdown = tokio::sync::watch::channel(false); @@ -335,7 +397,7 @@ impl World { ci_logs: None, admins, admission: knot_types::AdmissionPolicy::Closed, - contribution_policy: knot_types::ContributionPolicy::Anyone, + contribution_policy: policy, knot_did: knot, knot_hostname: KnotHostname::new(KNOT_HOST).unwrap(), meta_path, @@ -376,7 +438,8 @@ impl World { maintenance: knot_maintenance::MaintenanceHandle::disabled(), appview: knot_types::AppviewEndpoint::new("https://tangled.test").unwrap(), slots: knot_resource::Slots::testing(8), - lfs: Some(knot_xrpc::LfsWeb::new(lfs_handle, 8)), + lfs, + blobs: Default::default(), catalog: Arc::new(knot_messages::Catalog::defaults()), firehose: Arc::new(knot_events::EventLog::new( ManualClock::new(UnixMicros::new(1_000_000_000)),