From 6ee87dedc705a03752dd9af642ee79a2df2d6eca Mon Sep 17 00:00:00 2001 From: Lewis Date: Wed, 30 Sep 2026 13:48:24 +0300 Subject: [PATCH] knot-xrpc/tests: blobs go thru router & stupid tcp socket Lewis: May this revision serve well! --- knot2/crates/knot-xrpc/tests/blobs.rs | 610 ++++++++++++++++++++++++++ 1 file changed, 610 insertions(+) create mode 100644 knot2/crates/knot-xrpc/tests/blobs.rs diff --git a/knot2/crates/knot-xrpc/tests/blobs.rs b/knot2/crates/knot-xrpc/tests/blobs.rs new file mode 100644 index 000000000..609b94467 --- /dev/null +++ b/knot2/crates/knot-xrpc/tests/blobs.rs @@ -0,0 +1,610 @@ +mod common; + +use axum::body::Body; +use common::post_json; +use futures::StreamExt; +use http::{Request, StatusCode, header}; +use serde_json::json; +use tower::ServiceExt; + +use common::{ + OWNER, World, empty_repo, enable_chain, get, post_as, seeded_feature_branch, service_jwt_for, +}; +use knot_types::RepoDid; + +const PNG: &[u8] = &[ + 0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a, 0x00, 0x00, 0x00, 0x0d, 0x49, 0x48, 0x44, 0x52, + 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x01, 0x08, 0x06, 0x00, 0x00, 0x00, 0x1f, 0x15, 0xc4, + 0x89, +]; + +const SECOND: &[u8] = &[ + 0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a, 0x00, 0x00, 0x00, 0x0d, 0x49, 0x48, 0x44, 0x52, + 0x00, 0x01, +]; +const STRANGER: &str = "did:web:stranger.nel.pet"; + +fn blob_cid_link(bytes: &[u8]) -> knot_types::BlobCid { + knot_types::BlobCid::from_digest(knot_types::Sha256Digest::hash(bytes)) +} +fn blob_ref_json(cid: knot_types::BlobCid, mime: &str, size: u64) -> serde_json::Value { + json!({ + "$type": "blob", + "ref": { "$link": cid.to_string() }, + "mimeType": mime, + "size": size, + }) +} + +async fn upload_as( + world: &World, + repo: &RepoDid, + body: &[u8], + hint: &str, +) -> (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, hint) + .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(); + let value = serde_json::from_slice(&bytes).unwrap_or(serde_json::Value::Null); + (status, value) +} + +async fn upload(world: &World, repo: &RepoDid, body: &[u8]) -> (StatusCode, serde_json::Value) { + upload_as(world, repo, body, "image/png").await +} + +async fn get_json(world: &World, path_and_query: &str) -> (StatusCode, serde_json::Value) { + let (status, _headers, bytes) = get(world, path_and_query).await; + let value = serde_json::from_slice(&bytes).unwrap_or(serde_json::Value::Null); + (status, value) +} + +async fn create_with_blob( + world: &World, + actor: Option<&str>, + repo: &RepoDid, + blob: serde_json::Value, +) -> (StatusCode, serde_json::Value) { + let record = json!({ + "repo": repo.as_str(), + "collection": "sh.tangled.repo.issue", + "record": { + "$type": "sh.tangled.repo.issue", + "repo": repo.as_str(), + "title": "kelp with a picture", + "blobs": [blob], + "createdAt": "2026-09-30T00:00:00Z", + }, + }); + match actor { + Some(actor) => { + post_as( + world, + "/xrpc/com.atproto.repo.createRecord", + actor, + repo.as_str(), + record, + ) + .await + } + None => post_json(world, "/xrpc/com.atproto.repo.createRecord", record).await, + } +} + +fn aged_blob_path(world: &World, repo: &RepoDid, body: &[u8]) -> std::path::PathBuf { + let oid = knot_lfs::LfsOid::from_digest(knot_types::Sha256Digest::hash(body)); + let store = &world.state.lfs.as_ref().unwrap().handle.store; + store + .object_file(knot_lfs::Pool::Attachments, repo, &oid) + .unwrap() + .unwrap() + .1 +} + +#[tokio::test] +async fn attachment_round_trips_from_upload_to_getblob_and_listblobs() { + let world = World::new(); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + enable_chain(&world, &did); + let repo = did.as_str().to_owned(); + + let (status, uploaded) = upload(&world, &did, PNG).await; + assert_eq!(status, StatusCode::OK, "{uploaded}"); + assert_eq!(uploaded["blob"]["$type"], "blob"); + assert_eq!(uploaded["blob"]["mimeType"], "image/png"); + assert_eq!(uploaded["blob"]["size"], PNG.len() as u64); + assert_eq!( + common::blob_cid_of(&uploaded["blob"]["ref"]["$link"]), + blob_cid_link(PNG) + ); + + let blob = uploaded["blob"].clone(); + let (status, written) = create_with_blob(&world, Some(OWNER), &did, blob).await; + assert_eq!(status, StatusCode::OK, "{written}"); + let uri = written["uri"].as_str().unwrap().to_owned(); + let rkey = uri.rsplit('/').next().unwrap().to_owned(); + + let (status, _headers, _bytes) = get( + &world, + &format!("/xrpc/com.atproto.repo.getRecord?repo={repo}&collection=sh.tangled.repo.issue&rkey={rkey}"), + ) + .await; + assert_eq!(status, StatusCode::OK); + + let (status, headers, served) = get( + &world, + &format!( + "/xrpc/com.atproto.sync.getBlob?did={repo}&cid={}", + blob_cid_link(PNG) + ), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!(served.as_ref(), PNG); + assert_eq!(headers.get(header::CONTENT_TYPE).unwrap(), "image/png"); + assert_eq!(headers.get("x-content-type-options").unwrap(), "nosniff"); + assert!( + headers + .get("content-security-policy") + .unwrap() + .to_str() + .unwrap() + .starts_with("default-src 'none'") + ); + + let (status, listed) = get_json( + &world, + &format!("/xrpc/com.atproto.sync.listBlobs?did={repo}"), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + common::listed_cids(&listed), + vec![blob_cid_link(PNG)], + "{listed}" + ); + + futures::stream::iter([ + ( + blob_ref_json(blob_cid_link(b"never uploaded anywhere"), "image/png", 9), + "isn't in this repo's blob store", + ), + ( + blob_ref_json(blob_cid_link(PNG), "image/png", 3), + "is stored as", + ), + ( + blob_ref_json(blob_cid_link(PNG), "image/webp", PNG.len() as u64), + "reads as", + ), + ]) + .for_each(|(blob, refused)| { + let (world, did) = (&world, &did); + async move { + let (status, response) = create_with_blob(world, Some(OWNER), did, blob).await; + assert_eq!(status, StatusCode::BAD_REQUEST, "{response}"); + assert!( + response["message"].as_str().unwrap().contains(refused), + "{response}" + ); + } + }) + .await; +} + +#[tokio::test] +async fn uploads_without_lfs_store_return_503_with_the_setting() { + let world = World::lfs_less(); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + let (status, refused) = upload(&world, &did, PNG).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{refused}"); + assert!( + refused["message"] + .as_str() + .unwrap() + .contains("lfs.store_path"), + "The 503 includes `lfs.store_path`: {refused}" + ); + + enable_chain(&world, &did); + let referenced = blob_ref_json(blob_cid_link(PNG), "image/png", PNG.len() as u64); + let (status, refused) = create_with_blob(&world, Some(OWNER), &did, referenced).await; + assert_eq!(status, StatusCode::SERVICE_UNAVAILABLE, "{refused}"); + assert!( + refused["message"] + .as_str() + .unwrap() + .contains("lfs.store_path"), + "Creating a record with an attachment is refused like the upload: {refused}" + ); +} + +#[tokio::test] +async fn upload_bounds_take_any_bytes_under_one_mime_budget() { + let world = World::new(); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + let mut boundary = vec![0x89, b'P', b'N', b'G', 0x0d, 0x0a, 0x1a, 0x0a]; + boundary.resize(1_000_000, 0); + let (status, accepted) = upload(&world, &did, &boundary).await; + assert_eq!(status, StatusCode::OK, "{accepted}"); + assert_eq!(accepted["blob"]["size"], 1_000_000_u64, "{accepted}"); + assert_eq!(accepted["blob"]["mimeType"], "image/png", "{accepted}"); +} + +#[tokio::test] +async fn unsniffable_bytes_take_the_hint_except_html() { + let world = World::new(); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + let mystery = b"an unsniffable body"; + let (status, trusted) = upload_as(&world, &did, mystery, "application/x-whelk-report").await; + assert_eq!(status, StatusCode::OK, "{trusted}"); + assert_eq!( + trusted["blob"]["mimeType"], "application/x-whelk-report", + "{trusted}" + ); + let (status, refused_html) = upload_as(&world, &did, mystery, "text/html; charset=utf-8").await; + assert_eq!(status, StatusCode::OK, "{refused_html}"); + assert_eq!( + refused_html["blob"]["mimeType"], "application/octet-stream", + "{refused_html}" + ); + let (status, bare) = upload_as(&world, &did, b"mystery bytes", "").await; + assert_eq!(status, StatusCode::OK, "{bare}"); + assert_eq!( + bare["blob"]["mimeType"], "application/octet-stream", + "{bare}" + ); +} + +#[tokio::test] +async fn pdf_svg_and_mp4_uploads_take_the_sniffed_mime() { + async fn upload_returns_mime(world: &World, repo: &RepoDid, body: &[u8], mime: &str) { + let (status, response) = upload(world, repo, body).await; + assert_eq!(status, StatusCode::OK, "{response}"); + assert_eq!(response["blob"]["mimeType"], mime, "{response}"); + } + let world = World::new(); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + upload_returns_mime( + &world, + &did, + b"%PDF-1.7\n%an academic paper\n", + "application/pdf", + ) + .await; + upload_returns_mime( + &world, + &did, + b"", + "image/svg+xml", + ) + .await; + upload_returns_mime( + &world, + &did, + &[0, 0, 0, 32, b'f', b't', b'y', b'p', b'i', b's', b'o', b'm'], + "video/mp4", + ) + .await; +} + +#[tokio::test] +async fn getblob_returns_blob_not_found_for_unknown_cid() { + let world = World::new(); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + enable_chain(&world, &did); + let repo = did.as_str().to_owned(); + let (status, response) = get_json( + &world, + &format!( + "/xrpc/com.atproto.sync.getBlob?did={repo}&cid={}", + blob_cid_link(b"never uploaded anywhere") + ), + ) + .await; + assert_eq!(status, StatusCode::NOT_FOUND, "{response}"); + assert_eq!(response["error"], "BlobNotFound", "{response}"); +} + +#[tokio::test] +async fn listblobs_pages_and_applies_since_at_each_records_rev() { + let world = World::new(); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + enable_chain(&world, &did); + let repo = did.as_str().to_owned(); + + let links: Vec<(jacquard_common::types::tid::Tid, knot_types::BlobCid)> = + futures::stream::iter([PNG, SECOND]) + .then(|body| { + let (world, did) = (&world, &did); + async move { + let (status, uploaded) = upload(world, did, body).await; + assert_eq!(status, StatusCode::OK, "{uploaded}"); + let (status, written) = + create_with_blob(world, Some(OWNER), did, uploaded["blob"].clone()).await; + assert_eq!(status, StatusCode::OK, "{written}"); + (common::commit_rev_of(&written), blob_cid_link(body)) + } + }) + .collect() + .await; + + let (status, both) = get_json( + &world, + &format!("/xrpc/com.atproto.sync.listBlobs?did={repo}"), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + common::listed_cids(&both), + vec![links[0].1, links[1].1], + "The full listing is oldest reference first." + ); + + let (status, page) = get_json( + &world, + &format!("/xrpc/com.atproto.sync.listBlobs?did={repo}&limit=1"), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!(common::listed_cids(&page), vec![links[0].1], "{page}"); + let cursor = page["cursor"].as_str().unwrap().to_owned(); + let (status, rest) = get_json( + &world, + &format!("/xrpc/com.atproto.sync.listBlobs?did={repo}&limit=1&cursor={cursor}"), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!(common::listed_cids(&rest), vec![links[1].1], "{rest}"); + assert!(rest["cursor"].is_null(), "The last page has no cursor."); + + let (status, since) = get_json( + &world, + &format!( + "/xrpc/com.atproto.sync.listBlobs?did={repo}&since={}", + links[0].0 + ), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + common::listed_cids(&since), + vec![links[1].1], + "`since` at the first write's own rev excludes that rev, matching the firehose cursor." + ); + let (status, empty) = get_json( + &world, + &format!( + "/xrpc/com.atproto.sync.listBlobs?did={repo}&since={}", + links[1].0 + ), + ) + .await; + assert_eq!(status, StatusCode::OK); + assert_eq!( + empty["cids"], + json!([]), + "`since` at the newest rev returns an empty page." + ); +} + +async fn tcp_call( + addr: std::net::SocketAddr, + method: &str, + path_and_query: &str, + bearer: Option<&str>, + content_type: Option<&str>, + body: &[u8], +) -> (u16, Vec<(String, String)>, Vec) { + use tokio::io::{AsyncReadExt, AsyncWriteExt}; + + let mut request = + format!("{method} {path_and_query} HTTP/1.1\r\nHost: {addr}\r\nConnection: close\r\n"); + if let Some(token) = bearer { + request.push_str(&format!("Authorization: Bearer {token}\r\n")); + } + if let Some(mime) = content_type { + request.push_str(&format!("Content-Type: {mime}\r\n")); + } + request.push_str(&format!("Content-Length: {}\r\n\r\n", body.len())); + + let mut stream = tokio::net::TcpStream::connect(addr).await.unwrap(); + stream.write_all(request.as_bytes()).await.unwrap(); + stream.write_all(body).await.unwrap(); + let mut raw = Vec::new(); + stream.read_to_end(&mut raw).await.unwrap(); + + let split = raw + .windows(4) + .position(|window| window == b"\r\n\r\n") + .expect("the response starts with a header block"); + let head = String::from_utf8_lossy(&raw[..split]).to_string(); + let mut lines = head.split("\r\n"); + let status: u16 = lines + .next() + .unwrap() + .split_whitespace() + .nth(1) + .unwrap() + .parse() + .unwrap(); + let headers = lines + .map(|line| { + let (name, value) = line.split_once(':').expect("a header line"); + (name.trim().to_ascii_lowercase(), value.trim().to_owned()) + }) + .collect(); + (status, headers, raw[split + 4..].to_vec()) +} + +fn header_of<'a>(headers: &'a [(String, String)], name: &str) -> &'a str { + &headers + .iter() + .find(|(key, _)| key == name) + .unwrap_or_else(|| panic!("the {name} header is present")) + .1 +} + +#[tokio::test] +async fn attachments_round_trip_over_real_tcp_against_live_listener() { + let world = World::new(); + let (did, _bare, _work) = empty_repo(&world, "conch"); + enable_chain(&world, &did); + let repo = did.as_str(); + let addr = common::serve(&world).await; + + let (status, _headers, uploaded) = tcp_call( + addr, + "POST", + "/xrpc/com.atproto.repo.uploadBlob", + Some(&service_jwt_for("com.atproto.repo.uploadBlob", OWNER, repo)), + Some("image/png"), + PNG, + ) + .await; + assert_eq!(status, 200); + let uploaded: serde_json::Value = serde_json::from_slice(&uploaded).unwrap(); + let link = common::blob_cid_of(&uploaded["blob"]["ref"]["$link"]); + assert_eq!(uploaded["blob"]["mimeType"], "image/png"); + assert_eq!(link, blob_cid_link(PNG)); + + let (status, headers, served) = tcp_call( + addr, + "GET", + &format!("/xrpc/com.atproto.sync.getBlob?did={repo}&cid={link}"), + None, + None, + b"", + ) + .await; + assert_eq!(status, 200); + assert_eq!(served, PNG); + assert_eq!(header_of(&headers, "content-type"), "image/png"); + assert_eq!(header_of(&headers, "x-content-type-options"), "nosniff"); + assert_eq!( + header_of(&headers, "content-security-policy"), + "default-src 'none'; sandbox" + ); +} + +#[tokio::test] +async fn refused_writers_leave_stored_blobs_grace_window_untouched() { + let world = World::under(knot_types::ContributionPolicy::Collaborators); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + enable_chain(&world, &did); + let (status, uploaded) = upload(&world, &did, PNG).await; + assert_eq!(status, StatusCode::OK, "{uploaded}"); + + let path = aged_blob_path(&world, &did, PNG); + let aged = std::time::UNIX_EPOCH + std::time::Duration::from_secs(3_600); + std::fs::OpenOptions::new() + .write(true) + .open(&path) + .unwrap() + .set_modified(aged) + .unwrap(); + + let blob = uploaded["blob"].clone(); + let (status, refused) = create_with_blob(&world, None, &did, blob.clone()).await; + assert_eq!(status, StatusCode::UNAUTHORIZED, "{refused}"); + let (status, refused) = create_with_blob(&world, Some(STRANGER), &did, blob).await; + assert_eq!(status, StatusCode::FORBIDDEN, "{refused}"); + assert_eq!( + std::fs::metadata(&path).unwrap().modified().unwrap(), + aged, + "Refused writes never touch a stored blob's mtime." + ); +} + +#[tokio::test] +async fn listblobs_keeps_shared_blob_at_latest_referencing_record() { + let world = World::new(); + let (did, _bare, _work) = empty_repo(&world, "anemone"); + enable_chain(&world, &did); + let repo = did.as_str().to_owned(); + + let (status, uploaded) = upload(&world, &did, PNG).await; + assert_eq!(status, StatusCode::OK, "{uploaded}"); + let blob = uploaded["blob"].clone(); + let (status, first) = create_with_blob(&world, Some(OWNER), &did, blob.clone()).await; + assert_eq!(status, StatusCode::OK, "{first}"); + let (status, second) = create_with_blob(&world, Some(OWNER), &did, blob).await; + assert_eq!(status, StatusCode::OK, "{second}"); + + let at_first = common::commit_rev_of(&first); + let (status, listed) = get_json( + &world, + &format!("/xrpc/com.atproto.sync.listBlobs?did={repo}&since={at_first}"), + ) + .await; + assert_eq!(status, StatusCode::OK, "{listed}"); + assert_eq!( + common::listed_cids(&listed), + vec![blob_cid_link(PNG)], + "Writing the record after the cursor keeps the blob listed." + ); +} + +#[tokio::test] +async fn pull_lists_blobs_as_issues_do() { + let world = World::new(); + let (did, base, head) = seeded_feature_branch(&world, "whelk"); + enable_chain(&world, &did); + let repo = did.as_str().to_owned(); + let (status, uploaded) = upload(&world, &did, PNG).await; + assert_eq!(status, StatusCode::OK, "{uploaded}"); + + let (status, written) = post_as( + &world, + "/xrpc/com.atproto.repo.createRecord", + OWNER, + &repo, + json!({ + "repo": repo, + "collection": "sh.tangled.repo.pull", + "record": { + "$type": "sh.tangled.repo.pull", + "target": { "repo": repo, "branch": "main" }, + "title": "kelp with a picture", + "rounds": [], + "versions": [{ "base": base.to_hex(), "head": head.to_hex() }], + "blobs": [uploaded["blob"].clone()], + "createdAt": "2026-10-01T00:00:00Z", + }, + }), + ) + .await; + assert_eq!(status, StatusCode::OK, "{written}"); + let rkey = written["uri"] + .as_str() + .unwrap() + .rsplit('/') + .next() + .unwrap() + .to_owned(); + + let (status, record) = get_json( + &world, + &format!("/xrpc/com.atproto.repo.getRecord?repo={repo}&collection=sh.tangled.repo.pull&rkey={rkey}"), + ) + .await; + assert_eq!(status, StatusCode::OK, "{record}"); + assert_eq!( + common::blob_cid_of(&record["value"]["blobs"][0]["ref"]["$link"]), + blob_cid_link(PNG), + "The served pull: {record}" + ); +} -- 2.51.2