diff --git a/knot2/crates/knot-xrpc/src/social.rs b/knot2/crates/knot-xrpc/src/social.rs index b7735476d..fc5431b90 100644 --- a/knot2/crates/knot-xrpc/src/social.rs +++ b/knot2/crates/knot-xrpc/src/social.rs @@ -42,7 +42,7 @@ use knot_record::reaction::{ReactionRecord, reaction_collection}; use knot_resource::StructureBytes; use knot_runtime::{Clock, HttpTransport, Signer}; use knot_types::{ - AccountDid, AtUri, BranchName, ChangeId, CobId, ContributionPermission, MstKey, Oid, + AccountDid, AtUri, BranchName, ChangeId, CobId, ContributionPermission, Decision, MstKey, Oid, RecordAddress, RecordBody, RecordCid, RecordCollection, RecordRkey, RepoDid, SubjectNumber, Tid, TypeName, UnixSeconds, }; @@ -60,7 +60,7 @@ pub(crate) const LIST_COB_VERSIONS_ROUTE: &str = "/xrpc/sh.tangled.repo.listCobV const NOT_STATE_AUTHORITY: &str = "State changes belong to the issue's author or a maintainer"; -const COLLABORATE_REMEDY: &str = "Ask the repository owner to add you as a collaborator, or to open the contribution policy to members or anyone"; +pub(crate) const COLLABORATE_REMEDY: &str = "Please ask the repo owner to add you as a collaborator, or to open the contribution policy to members or anyone, thanks."; pub struct RecordKeys(Mutex); @@ -179,6 +179,8 @@ struct IssueDraft { mentions: Option>, #[serde(default)] references: Option>>, + #[serde(default)] + blobs: Option>, #[serde(rename = "createdAt")] _created_at: serde::de::IgnoredAny, } @@ -188,10 +190,14 @@ impl IssueDraft { self, written_to: &RepoDid, collection: &RecordCollection, - ) -> Result { + ) -> Result<(Prose, Option>), XrpcError> { kind_matches(self.kind.as_ref(), collection)?; on_this_knot(&self.repo, written_to, "record")?; - prose(self.title, self.body, self.mentions, self.references) + let blobs = self.blobs; + Ok(( + prose(self.title, self.body, self.mentions, self.references)?, + blobs, + )) } } @@ -252,7 +258,7 @@ struct MarkupInput { #[serde(default)] original: Option, #[serde(default)] - blobs: Option, + blobs: Option>, } impl MarkupInput { @@ -261,11 +267,6 @@ impl MarkupInput { self.kind.as_ref(), &knot_record::comment::markup_collection(), )?; - if self.blobs.is_some() { - return Err(XrpcError::invalid_request( - "blobs are attachments, and this knot won't take them; please leave blobs out", - )); - } let unusable = |what: &'static str| { move |error| { XrpcError::invalid_request(format!("Comment body's {what} is unusable: {error}")) @@ -277,7 +278,7 @@ impl MarkupInput { .map(MarkdownBody::new) .transpose() .map_err(unusable("original"))?; - Ok(CommentMarkup::new(text, original)) + Ok(CommentMarkup::new(text, original, self.blobs)) } } @@ -416,6 +417,8 @@ struct PullDraft { references: Option>>, #[serde(rename = "dependentOn", default)] dependent_on: Option>, + #[serde(default)] + blobs: Option>, #[serde(rename = "createdAt")] _created_at: serde::de::IgnoredAny, } @@ -426,6 +429,13 @@ struct ParsedPull { source: Option, versions: Vec, dependent_on: Option, + attachments: Option>, +} + +impl ParsedPull { + fn attachments(&self) -> &[knot_record::attachment::Attachment] { + self.attachments.as_deref().unwrap_or(&[]) + } } impl PullDraft { @@ -498,6 +508,7 @@ impl PullDraft { source, versions, dependent_on, + attachments: self.blobs, }) } } @@ -1028,14 +1039,34 @@ enum Doing<'a> { }, } -struct Session { - actor: AccountDid, - repo: RepoDid, +pub(crate) struct Session { + pub(crate) actor: AccountDid, + pub(crate) repo: RepoDid, permission: ContributionPermission, - now: UnixSeconds, + pub(crate) now: UnixSeconds, +} + +impl Session { + pub(crate) fn permission(&self) -> ContributionPermission { + self.permission + } + + pub(crate) fn require_contribution(self, act: &'static str) -> Result { + match self.permission.contributes() { + Decision::Allow => Ok(Contributor { session: self }), + Decision::Deny => Err(XrpcError::forbidden(format!( + "`{act}` belong to contributors this repo's policy allows, and it leaves \ + you out; {COLLABORATE_REMEDY}" + ))), + } + } } -async fn session( +pub(crate) struct Contributor { + pub(crate) session: Session, +} + +pub(crate) async fn session( state: &XrpcState, headers: &HeaderMap, method: &crate::Method, @@ -1217,11 +1248,19 @@ async fn issue_open( record: serde_json::Value, ) -> Result { let collection = issue_collection(); - let prose = typed::(record, &collection)?.parse(&repo, &collection)?; - let session = session(state, headers, method, repo).await?; + let (prose, blobs) = typed::(record, &collection)?.parse(&repo, &collection)?; + let session = crate::blobs::session_admitting( + state, + headers, + method, + repo, + blobs.as_deref().unwrap_or(&[]), + ) + .await?; let record = IssueRecord::new( session.repo.clone(), prose, + blobs, session.now, session.actor.clone(), )? @@ -1241,7 +1280,9 @@ async fn comment_open( ) -> Result { let collection = comment_collection(); let draft = typed::(record, &collection)?.parse(&repo)?; - let session = session(state, headers, method, repo).await?; + let session = + crate::blobs::session_admitting(state, headers, method, repo, draft.markup.attachments()) + .await?; let subject = draft.subject.clone(); let reply = draft.reply_to.clone(); let round = draft.round; @@ -1536,14 +1577,22 @@ async fn issue_put( ) -> Result { let (repo, rkey, from) = swap.edit_of(&issue_collection())?; let collection = issue_collection(); - let prose = typed::(record, &collection)?.parse(&repo, &collection)?; + let (prose, blobs) = typed::(record, &collection)?.parse(&repo, &collection)?; + let session = crate::blobs::session_admitting( + state, + headers, + method, + repo, + blobs.as_deref().unwrap_or(&[]), + ) + .await?; let address = RecordAddress::new(collection, rkey); - let session = session(state, headers, method, repo).await?; on_repo(state, session, move |ledger, cx| { let (object, live) = standing(ledger, &address)?; let record = IssueRecord::new( ledger.repo.clone(), prose, + blobs, live.created_at(), cx.actor().clone(), )? @@ -1586,8 +1635,10 @@ async fn comment_put( let (repo, rkey, from) = swap.edit_of(&comment_collection())?; let collection = comment_collection(); let draft = typed::(record, &collection)?.parse(&repo)?; + let session = + crate::blobs::session_admitting(state, headers, method, repo, draft.markup.attachments()) + .await?; let address = RecordAddress::new(collection, rkey); - let session = session(state, headers, method, repo).await?; on_repo(state, session, move |ledger, cx| { let (object, live) = standing(ledger, &address)?; if live.parent() != Some(draft.subject.address()) { @@ -1866,6 +1917,7 @@ fn pull_body( draft.source.clone(), dependent, versions, + draft.attachments.clone(), created_at, editor.clone(), )? @@ -2000,7 +2052,8 @@ async fn pull_open( ) -> Result { let collection = pull_collection(); let draft = typed::(record, &collection)?.parse(&repo)?; - let session = session(state, headers, method, repo).await?; + let session = + crate::blobs::session_admitting(state, headers, method, repo, draft.attachments()).await?; let Session { ref repo, now, @@ -2071,7 +2124,8 @@ async fn pull_put( let collection = pull_collection(); let draft = typed::(record, &collection)?.parse(&repo)?; let address = RecordAddress::new(collection, rkey); - let session = session(state, headers, method, repo).await?; + let session = + crate::blobs::session_admitting(state, headers, method, repo, draft.attachments()).await?; on_repo(state, session, move |ledger, cx| { let (object, live) = standing(ledger, &address)?; let opened = standing_pull(ledger, &address)?; diff --git a/knot2/crates/knot-xrpc/src/tests.rs b/knot2/crates/knot-xrpc/src/tests.rs index 0c4fcb72e..df4e14ecb 100644 --- a/knot2/crates/knot-xrpc/src/tests.rs +++ b/knot2/crates/knot-xrpc/src/tests.rs @@ -334,6 +334,7 @@ fn state_from( ); secrets.ensure(&knot).unwrap(); let state = Arc::new(XrpcState { + blobs: Default::default(), imports: crate::imports::ImportQueue::channel().0, layout, index, @@ -1343,7 +1344,7 @@ async fn delete_issue( fn rkey_of(uri: &AtUri) -> String { uri.rkey() - .expect("An at-uri of a record has its record key") + .expect("An AT-URI ends in a record key") .as_ref() .to_owned() } @@ -1357,7 +1358,7 @@ fn record_of( repo: crate::atproto::RepoSubject::Did(repo.clone()), collection: knot_types::RecordCollection::new( uri.collection() - .expect("An at-uri of a record has its collection") + .expect("An AT-URI has a collection") .as_str(), ) .unwrap(), @@ -1831,6 +1832,211 @@ async fn rewrites_chain_editors_and_erasure_tombstones_record_and_bytes() { ); } +async fn car_blocks_of(response: Response) -> std::collections::BTreeMap { + let body = axum::body::to_bytes(response.into_body(), usize::MAX) + .await + .unwrap(); + knot_record::car_blocks(&body) + .unwrap() + .blocks + .into_iter() + .map(|(cid, bytes)| (cid.to_string(), bytes)) + .collect() +} + +async fn sync_record_blocks( + world: &World, + repo: &RepoDid, + uri: &AtUri, +) -> std::collections::BTreeMap { + let response = crate::atproto::sync_get_record( + world.state(), + crate::query::ValidatedQuery(crate::atproto::SyncRecordParams { + did: crate::atproto::RepoSubject::Did(repo.clone()), + collection: knot_types::RecordCollection::new( + uri.collection() + .expect("An AT-URI has a collection") + .as_str(), + ) + .unwrap(), + rkey: knot_types::UriRkey::new(rkey_of(uri)).unwrap(), + }), + ) + .await + .unwrap(); + car_blocks_of(response).await +} + +async fn repo_export_blocks( + world: &World, + repo: &RepoDid, + since: Option, +) -> std::collections::BTreeMap { + let since = since.map(|sent| { + serde_json::from_value::(serde_json::Value::String(sent)) + .expect("a rev string parses as a TID") + }); + let response = crate::atproto::sync_get_repo( + world.state(), + crate::query::ValidatedQuery(crate::atproto::RepoExportParams { + did: crate::atproto::RepoSubject::Did(repo.clone()), + since, + }), + ) + .await + .unwrap(); + car_blocks_of(response).await +} + +fn tip_rev(world: &World, repo: &RepoDid) -> String { + let git = world.layout.open(repo).unwrap(); + knot_record::chain::tip(&git) + .unwrap() + .expect("issue_world enables the chain") + .rev + .to_string() +} + +fn cob_version_bytes( + world: &World, + repo: &RepoDid, + uri: &AtUri, +) -> Vec<(Option, Option>)> { + let git = world.layout.open(repo).unwrap(); + let address = knot_types::RecordAddress::new( + knot_types::RecordCollection::new( + uri.collection() + .expect("An AT-URI has a collection") + .as_str(), + ) + .unwrap(), + knot_types::RecordRkey::new(rkey_of(uri)).unwrap(), + ); + let object = match world.state.index.social_object_at(repo, &address) { + knot_index::Resolved::Ready(Some(object)) => object, + other => panic!("the written issue ought to project to an object, got {other:?}"), + }; + let type_name = knot_types::TypeName::new(ISSUE).unwrap(); + let delta = knot_cob::CobStore::new(&git) + .changes_of(&type_name, object, knot_cob::HistoryModel::Convergent, None) + .unwrap(); + delta + .changes + .iter() + .map(|change| { + let payload = knot_cobs::decode_social(change.payload()).unwrap(); + let emission = knot_cobs::Social::emission(&payload); + ( + emission.produces().map(|cid| cid.to_string()), + change.record.bytes().map(<[u8]>::to_vec), + ) + }) + .collect() +} + +#[tokio::test] +async fn deleted_prose_stays_off_every_presenting_route_and_reachable_from_history() { + let (world, repo) = issue_world().await; + let before = tip_rev(&world, &repo); + let first = ok(open_issue(&world, &world.stranger, &repo, "kelp", "a forest").await); + let uri = uri_of(&first); + let v1 = cid_of(&first); + let birth = first["commit"]["rev"].as_str().unwrap().to_owned(); + let second = ok(edit_issue(&world, &world.stranger, &repo, &uri, &v1, "kelp, again").await); + let v2 = cid_of(&second); + ok(delete_issue(&world, &world.stranger, &repo, &uri, &v2).await); + + let gone = failed( + record_at(&world, &repo, &uri, None).await, + StatusCode::BAD_REQUEST, + ); + assert_eq!(gone["error"], "RecordNotFound"); + let by_old_cid = failed( + record_at(&world, &repo, &uri, Some(&v1)).await, + StatusCode::BAD_REQUEST, + ); + assert_eq!( + by_old_cid["error"], "RecordNotFound", + "The deleted version's own CID mustn't resurrect it." + ); + assert!( + records_in(&world, &repo, ISSUE).await.is_empty(), + "listRecords must drop the row." + ); + let proof = sync_record_blocks(&world, &repo, &uri).await; + assert!( + !proof.is_empty(), + "sync.getRecord returns a proof of absence." + ); + assert!( + !proof.contains_key(&v1.to_string()) && !proof.contains_key(&v2.to_string()), + "sync.getRecord mustn't return either version's block." + ); + let chain = versions_page(&world, &uri, None).await; + let rows = chain["versions"].as_array().unwrap(); + assert_eq!(rows.len(), 3, "one create, one edit, one tombstone"); + rows.iter().for_each(|row| { + let keys = row.as_object().unwrap(); + assert!( + keys.keys().all(|key| key == "cid" + || key == "changeId" + || key == "editor" + || key == "createdAt"), + "Version row is metadata alone, not prose: {row}" + ); + }); + assert!( + rows[2].get("cid").is_none(), + "The tombstone row has no `cid`." + ); + + let mut served = repo_export_blocks(&world, &repo, None).await; + assert!( + !served.contains_key(&v1.to_string()) && !served.contains_key(&v2.to_string()), + "The current export omits the deleted blocks." + ); + served = repo_export_blocks(&world, &repo, Some(before)).await; + assert!( + served.contains_key(&v1.to_string()) && served.contains_key(&v2.to_string()), + "Diffing the export from before the create still returns both versions, the way \ + reflog keeps a reverted commit." + ); + let replay_since_create = repo_export_blocks(&world, &repo, Some(birth)).await; + assert!( + !replay_since_create.contains_key(&v1.to_string()), + "`since` at the create's own rev excludes that rev, matching the firehose cursor." + ); + + assert!( + block_served(&world, &repo, &v1).await && block_served(&world, &repo, &v2).await, + "getBlocks keeps serving both versions by CID." + ); + let blocks = crate::atproto::sync_get_blocks( + world.state(), + crate::atproto::BlocksQuery { + did: crate::atproto::RepoSubject::Did(repo.clone()), + cids: vec![v1.cid()], + }, + ) + .await + .unwrap(); + let served_v1 = car_blocks_of(blocks).await; + let change_tree_v1 = cob_version_bytes(&world, &repo, &uri) + .into_iter() + .find_map(|(cid, bytes)| (cid == Some(v1.to_string())).then_some(bytes)) + .flatten() + .expect("the change graph keeps the version's record blob"); + assert_eq!( + served_v1.get(&v1.to_string()).map(Bytes::as_ref), + Some(change_tree_v1.as_slice()), + "`refs/cobs/*` serves the same bytes `getBlocks` serves, so a mirror fetch reads them." + ); + assert!( + replayed_blocks(&world).contains(&v1.cid()), + "Firehose replay from before the create still includes the leaf." + ); +} + #[tokio::test] async fn state_transition_belongs_to_author_or_maintainer_and_repeat_answers_409() { let (world, repo) = issue_world().await; @@ -1879,7 +2085,7 @@ async fn state_transition_belongs_to_author_or_maintainer_and_repeat_answers_409 shut_out["message"] .as_str() .unwrap() - .ends_with("or to open the contribution policy to members or anyone"), + .ends_with("or to open the contribution policy to members or anyone, thanks."), "{shut_out}" ); assert_eq!( @@ -1952,6 +2158,7 @@ fn issue_body_bytes(world: &World, actor: &Actor, repo: &RepoDid, title: &str, t let record = knot_record::issue::IssueRecord::new( repo.clone(), prose, + None, world.state.now(), actor.did.clone(), ) @@ -4186,7 +4393,7 @@ mod fork_endpoints { .find_ref(&main_ref()) .unwrap(); let new = advance(&source, &main_ref(), "tide.txt", "rock pool\n", 1_001); - assert_ne!(old, Some(new.clone())); + assert_ne!(old, Some(new)); rewrite_sidecar(&world, &did, |job| { job["state"] = json!("failed"); @@ -5746,7 +5953,7 @@ async fn comment_opens_edits_and_deletes_along_its_version_chain() { assert_eq!(record["value"]["$type"], COMMENT); assert_eq!(record["value"]["subject"]["uri"], issue_uri.as_str()); assert_eq!( - record["value"]["subject"]["cid"], + record["value"]["subject"]["cid"]["$link"], issue_cid.to_string(), "subject strongRef must name the exact issue versoon the comment was made on" ); @@ -6180,12 +6387,14 @@ async fn comment_rejects_extras_it_cant_store() { &["embed is an attachment", "leave embed out"], ); - let mut blobbed = comment_record(&issue_uri, &issue_cid, "with a blob"); - blobbed["body"]["blobs"] = json!([]); - message_contains( - &rejected_with(blobbed).await, - &["blobs are attachments", "leave blobs out"], - ); + let mut malformed = comment_record(&issue_uri, &issue_cid, "with a malformed blob"); + malformed["body"]["blobs"] = json!([{ + "$type": "blob", + "ref": { "$link": "not-a-cid" }, + "mimeType": "image/png", + "size": 5 + }]); + message_contains(&rejected_with(malformed).await, &["unusable"]); } #[tokio::test]