diff --git a/knot2/crates/knot-index/src/intern.rs b/knot2/crates/knot-index/src/intern.rs index a6fad4a8e..c0e35b9d0 100644 --- a/knot2/crates/knot-index/src/intern.rs +++ b/knot2/crates/knot-index/src/intern.rs @@ -38,7 +38,7 @@ macro_rules! interned { impl Interner { pub(crate) fn $resolve(&self, key: $key) -> $value { <$value>::new(self.0.resolve(&key.0)) - .expect(concat!("interned ", $label, " is valid ", $label)) + .expect(concat!("Interned ", $label, " is valid ", $label)) } } )? @@ -57,18 +57,38 @@ interned! { const TID_BYTES: usize = 13; #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)] -pub(crate) struct RecordKey([u8; TID_BYTES]); +pub(crate) enum RecordKey { + Tid([u8; TID_BYTES]), + Slug(Spur), +} -impl RecordKey { - pub(crate) fn of(value: &RecordRkey) -> Self { - Self( - <[u8; TID_BYTES]>::try_from(value.as_str().as_bytes()) - .expect("a record key is a TID of thirteen bytes"), - ) +impl Interner { + pub(crate) fn key_of(&self, value: &RecordRkey) -> RecordKey { + match value { + RecordRkey::Minted(tid) => RecordKey::Tid( + <[u8; TID_BYTES]>::try_from(tid.as_str().as_bytes()) + .expect("A minted record key is a TID of thirteen bytes"), + ), + RecordRkey::Composed(key) => RecordKey::Slug(self.0.get_or_intern(key.as_str())), + } } - pub(crate) fn rkey(self) -> RecordRkey { - RecordRkey::new(std::str::from_utf8(&self.0).expect("record key bytes are TID ASCII")) - .expect("record key bytes spell a TID") + pub(crate) fn known_key(&self, value: &RecordRkey) -> Option { + match value { + RecordRkey::Minted(tid) => Some(RecordKey::Tid( + <[u8; TID_BYTES]>::try_from(tid.as_str().as_bytes()) + .expect("A minted record key is a TID of thirteen bytes"), + )), + RecordRkey::Composed(key) => self.0.get(key.as_str()).map(RecordKey::Slug), + } + } + + pub(crate) fn key_str(&self, key: RecordKey) -> String { + match key { + RecordKey::Tid(bytes) => std::str::from_utf8(&bytes) + .expect("Record key bytes are TID ASCII") + .to_owned(), + RecordKey::Slug(spur) => self.0.resolve(&spur).to_owned(), + } } } diff --git a/knot2/crates/knot-index/src/lib.rs b/knot2/crates/knot-index/src/lib.rs index f045a076d..8b2a3342f 100644 --- a/knot2/crates/knot-index/src/lib.rs +++ b/knot2/crates/knot-index/src/lib.rs @@ -404,7 +404,7 @@ impl Index { tracing::warn!( repo = repo.as_str(), %error, - "collaborators unread, the knot can't open the repo" + "Collaborators unread, the knot can't open the repo" ); true } @@ -622,6 +622,15 @@ impl Index { .reaction_of(&self.interner, repo, subject, actor, kind) } + pub fn social_labels_of( + &self, + repo: &RepoDid, + subject: &RecordAddress, + ) -> Resolved>> + { + self.social.labels_of(&self.interner, repo, subject) + } + pub fn social_entries(&self, repo: &RepoDid) -> Resolved> { self.social.entries(&self.interner, repo) } diff --git a/knot2/crates/knot-index/src/social.rs b/knot2/crates/knot-index/src/social.rs index 1e0c77163..87bccbc63 100644 --- a/knot2/crates/knot-index/src/social.rs +++ b/knot2/crates/knot-index/src/social.rs @@ -7,8 +7,8 @@ use std::sync::{LazyLock, Mutex}; use knot_cob::{Change, ChangePayload, CobId, CobStore, Delta, HistoryModel}; use knot_resource::StructureBytes; use knot_types::{ - AccountDid, ChangeId, RecordAddress, RecordCid, RecordCollection, RepoDid, SubjectNumber, - TypeName, UnixSeconds, + AccountDid, ChangeId, RecordAddress, RecordCid, RecordCollection, RecordRkey, RepoDid, + SubjectNumber, TypeName, UnixSeconds, }; use crate::coverage::Resolved; @@ -35,25 +35,25 @@ impl Slot { fn of(interner: &Interner, address: &RecordAddress) -> Self { Self { collection: interner.intern_collection(address.collection()), - rkey: RecordKey::of(address.rkey()), + rkey: interner.key_of(address.rkey()), } } fn known(interner: &Interner, address: &RecordAddress) -> Option { Some(Self { collection: interner.collection(address.collection())?, - rkey: RecordKey::of(address.rkey()), + rkey: interner.known_key(address.rkey())?, }) } fn address(self, interner: &Interner) -> RecordAddress { RecordAddress::new( interner.resolve_collection(self.collection), - self.rkey.rkey(), + RecordRkey::new(interner.key_str(self.rkey)) + .expect("An interned record key spells a valid rkey"), ) } } - #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum SocialEntryState { Serving(RecordCid), @@ -77,12 +77,15 @@ pub struct SocialEntry { pub state: SocialEntryState, } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +type Labelings = Box<[knot_cobs::OpOperand]>; + +#[derive(Debug, Clone, PartialEq, Eq)] struct Objected { slot: Slot, number: SubjectNumber, parent: Option, reacted: Option, + labeled: Option, created_at: UnixSeconds, version: Option, } @@ -175,8 +178,9 @@ impl RepoSocial { Self::reaction_key(&objected).iter().for_each(|key| { self.reacted.entry(*key).or_default().insert(object); }); + let slot = objected.slot; self.objects.insert(object, objected); - objected.slot + slot }); left.into_iter() .chain(arrived.filter(|slot| Some(*slot) != left)) @@ -219,12 +223,56 @@ impl RepoSocial { if *claimant != winner { tracing::warn!( object = %claimant, - "second object under one record address never answers lookups" + "Second object under one record address never answers lookups" ); } }); } + fn labels_of( + &self, + interner: &Interner, + subject: &Slot, + ) -> BTreeMap> { + let mut labels: BTreeMap> = + BTreeMap::new(); + interner + .collection(&knot_cobs::LabelOpCob::collection()) + .and_then(|ops| self.children.get(subject)?.get(&ops)) + .into_iter() + .flatten() + .map(|ordered| ordered.object()) + .filter(|object| self.is_serving(object)) + .filter_map(|object| self.objects.get(&object)?.labeled.as_ref()) + .flatten() + .for_each(|operand| { + let set = labels.entry(operand.key.clone()).or_default(); + match (operand.verb, operand.multiple) { + (knot_cobs::LabelVerb::Add, true) if !set.contains(&operand.value) => { + set.push(operand.value.clone()); + } + (knot_cobs::LabelVerb::Add, true) => {} + (knot_cobs::LabelVerb::Add, false) => { + set.clear(); + set.push(operand.value.clone()); + } + (knot_cobs::LabelVerb::Delete, true) => { + set.retain(|standing| *standing != operand.value); + } + (knot_cobs::LabelVerb::Delete, false) + if set + .first() + .is_some_and(|standing| *standing == operand.value) => + { + set.clear(); + } + (knot_cobs::LabelVerb::Delete, false) => {} + } + }); + labels.retain(|_, values| !values.is_empty()); + labels + } + fn unlist(&mut self, objected: &Objected, ordered: Ordered) { objected.parent.iter().for_each(|parent| { let emptied = self.children.get_mut(parent).is_some_and(|kinds| { @@ -270,7 +318,7 @@ impl RepoSocial { self.answers(&object).then_some(objected) } - fn answers_live(&self, object: &CobId) -> bool { + fn is_serving(&self, object: &CobId) -> bool { self.objects.get(object).is_some_and(|objected| { !self.contests(object, &objected.slot) && objected.version.is_some() }) @@ -279,7 +327,7 @@ impl RepoSocial { fn live_at(&self, slot: &Slot) -> bool { self.addressed .get(slot) - .is_some_and(|object| self.answers_live(object)) + .is_some_and(|object| self.is_serving(object)) } fn live_reaction( @@ -292,7 +340,7 @@ impl RepoSocial { .get(&(*subject, actor, kind))? .iter() .copied() - .find(|object| self.answers_live(object)) + .find(|object| self.is_serving(object)) } fn listed(&self, ordered: &BTreeSet) -> Vec { @@ -616,6 +664,19 @@ impl SocialProjection { }) } + pub(crate) fn labels_of( + &self, + interner: &Interner, + repo: &RepoDid, + subject: &RecordAddress, + ) -> Resolved>> { + let slot = Slot::known(interner, subject); + self.read(interner, repo, |social| { + slot.filter(|slot| social.live_at(slot)) + .map_or_else(BTreeMap::new, |slot| social.labels_of(interner, &slot)) + }) + } + pub(crate) fn address_of( &self, interner: &Interner, @@ -768,7 +829,7 @@ impl SocialProjection { } } -static PROJECTED: LazyLock<[RecordCollection; 4]> = LazyLock::new(knot_cobs::social_collections); +static PROJECTED: LazyLock<[RecordCollection; 6]> = LazyLock::new(knot_cobs::social_collections); fn is_social(type_name: &TypeName) -> bool { PROJECTED @@ -776,6 +837,53 @@ fn is_social(type_name: &TypeName) -> bool { .any(|collection| collection.as_str() == type_name.as_str()) } +fn with_opened_root( + store: &CobStore, + type_name: &TypeName, + object: CobId, + root: Option<&Change>, + state: &knot_cobs::SocialState, + place: impl FnOnce(&Change) -> Result, +) -> Result, IndexError> { + if state.object().is_none() { + return Ok(None); + } + let fetched = root + .is_none() + .then(|| store.changes_of(type_name, object, HistoryModel::Convergent, None)) + .transpose()?; + root.or_else(|| fetched.as_ref().and_then(|fetched| fetched.changes.first())) + .map(place) + .transpose() +} + +fn label_placement( + store: &CobStore, + type_name: &TypeName, + object: CobId, + root: Option<&Change>, + state: &knot_cobs::SocialState, +) -> Result, IndexError> { + if type_name.as_str() != knot_cobs::LabelOpCob::collection().as_str() { + return Ok(None); + } + with_opened_root(store, type_name, object, root, state, |root| { + let decoded = knot_cobs::LabelOpChange::decode(root.payload()).map_err(|error| { + IndexError::Decode { + change: root.id, + type_name: type_name.clone(), + reason: error.to_string(), + } + })?; + Ok(match decoded { + knot_cobs::SocialChange::Open { body, .. } => body.operands.into(), + knot_cobs::SocialChange::Edit { .. } | knot_cobs::SocialChange::Erase { .. } => { + Box::default() + } + }) + }) +} + fn project_object( interner: &Interner, store: &CobStore, @@ -817,6 +925,7 @@ fn replay_from( interner, &state, reaction_placement(interner, store, type_name, object, None, &state).ok()?, + label_placement(store, type_name, object, None, &state).ok()?, ), )) } @@ -868,6 +977,7 @@ fn rebuild_whole( delta.changes.first(), &state, )?, + label_placement(store, type_name, object, delta.changes.first(), &state)?, ), )) } @@ -883,45 +993,39 @@ fn reaction_placement( if type_name.as_str() != knot_cobs::ReactionCob::collection().as_str() { return Ok(None); } - let fetched = root - .is_none() - .then(|| store.changes_of(type_name, object, HistoryModel::Convergent, None)) - .transpose()?; - let root = root.or_else(|| fetched.as_ref().and_then(|fetched| fetched.changes.first())); - let live = state.object(); - root.zip(live) - .map(|(root, live)| { - let decoded = knot_cobs::ReactionChange::decode(root.payload()).map_err(|error| { - IndexError::Decode { - change: root.id, - type_name: type_name.clone(), - reason: error.to_string(), - } - })?; - Ok(match decoded { - knot_cobs::SocialChange::Open { body, .. } => Some(ReactionPlacement { + with_opened_root(store, type_name, object, root, state, |root| { + let decoded = knot_cobs::ReactionChange::decode(root.payload()).map_err(|error| { + IndexError::Decode { + change: root.id, + type_name: type_name.clone(), + reason: error.to_string(), + } + })?; + Ok(match decoded { + knot_cobs::SocialChange::Open { body, .. } => { + state.object().map(|live| ReactionPlacement { actor: interner.intern_account(live.author()), kind: body.kind, - }), - knot_cobs::SocialChange::Edit { .. } | knot_cobs::SocialChange::Erase { .. } => { - None - } - }) + }) + } + knot_cobs::SocialChange::Edit { .. } | knot_cobs::SocialChange::Erase { .. } => None, }) - .transpose() - .map(Option::flatten) + }) + .map(Option::flatten) } fn objected( interner: &Interner, state: &knot_cobs::SocialState, reacted: Option, + labeled: Option, ) -> Option { state.object().map(|object| Objected { slot: Slot::of(interner, object.address()), number: object.number(), parent: object.parent().map(|parent| Slot::of(interner, parent)), reacted, + labeled, created_at: object.created_at(), version: object.version(), }) diff --git a/knot2/crates/knot-index/tests/common/mod.rs b/knot2/crates/knot-index/tests/common/mod.rs index b564c5e0e..8a7b9521a 100644 --- a/knot2/crates/knot-index/tests/common/mod.rs +++ b/knot2/crates/knot-index/tests/common/mod.rs @@ -385,6 +385,72 @@ impl World { ) } + #[allow(clippy::too_many_arguments)] + pub fn opened_as( + &self, + repo: &RepoDid, + key: &str, + author: &AccountDid, + number: u32, + parent: Option<&RecordAddress>, + body: B, + prose: &str, + seconds: i64, + role: RepoRole, + ) -> CobId { + let (change, record) = + contribution_of(key, author, number, parent, body, prose, at(seconds)); + let git = self + .layout + .open(repo) + .unwrap_or_else(|_| self.layout.create(repo).unwrap()); + let cx = WriteContext::new( + author, + repo, + ContributionPermission::new(Blocked::No, role, ContributionPolicy::Anyone), + ); + CobStore::new(&git) + .create_approved::>( + &CobHome::from(repo), + &Approved::signed_by( + &cx, + Materialized { + change: change.clone(), + record: Record::version(record), + }, + ) + .unwrap(), + &self.signer, + at(seconds), + ) + .unwrap() + .object + } + + #[allow(clippy::too_many_arguments)] + pub fn label_op( + &self, + repo: &RepoDid, + key: &str, + author: &AccountDid, + number: u32, + parent: &RecordAddress, + body: knot_cobs::LabelOp, + seconds: i64, + ) -> CobId { + self.opened_as( + repo, + key, + author, + number, + Some(parent), + body, + key, + seconds, + RepoRole::Collaborator, + ) + } + #[allow(clippy::too_many_arguments)] pub fn opened( &self, @@ -401,6 +467,7 @@ impl World { contribution_of(key, author, number, parent, body, prose, at(seconds)); self.open_cob::>(repo, &change, Record::version(record), seconds) } + pub fn erase( &self, repo: &RepoDid, @@ -531,6 +598,10 @@ addressed! { comment_type_name of knot_cobs::CommentChange reaction_collection of knot_record::reaction::reaction_collection(), reaction_address, reaction_type_name of knot_cobs::ReactionChange + label_def_collection of knot_record::label::label_def_collection(), + label_def_address + label_op_collection of knot_record::label::label_op_collection(), label_op_address, + label_op_type_name of knot_cobs::LabelOpChange } pub fn note_record(text: &str) -> (RecordBody, RecordCid) { diff --git a/knot2/crates/knot-index/tests/projections.rs b/knot2/crates/knot-index/tests/projections.rs index 0415998b9..ddf92bf18 100644 --- a/knot2/crates/knot-index/tests/projections.rs +++ b/knot2/crates/knot-index/tests/projections.rs @@ -16,9 +16,10 @@ use serde::{Deserialize, Serialize}; mod common; use common::{ NoteBody, NoteCob, World, acc, at, comment_at, comment_collection, contribution_of, erased, - grant, issue_at, issue_collection, issue_editing, issue_type_name, meta_home, - named_registration, note_address, note_record, own, reaction_collection, registration, - repo_did, rkey, state_at, state_collection, state_type_name, tid, + grant, issue_at, issue_collection, issue_editing, issue_type_name, label_def_collection, + label_op_collection, label_op_type_name, meta_home, named_registration, note_address, + note_record, own, reaction_collection, registration, repo_did, rkey, state_at, + state_collection, state_type_name, tid, }; use knot_cobs::ReactionKind; @@ -93,7 +94,7 @@ fn a_concurrent_reader_never_sees_a_net_absent_subject() { assert!( !leaked.load(Ordering::Acquire), - "net-absent subject is never written, so no reader can observe it mid-delta" + "Net-absent subject is never written, so no reader can observe it mid-delta" ); assert_eq!(index.effective_member(&acc("teq")), Resolved::Ready(false)); assert_eq!(index.effective_member(&acc("f0")), Resolved::Ready(true)); @@ -209,7 +210,7 @@ fn a_diverged_collaborators_tip_purges_the_roll_instead_of_serving_it_stale() { assert!( index.refresh_collaborators(&repo).is_err(), - "tip that no longer descends from folded tip is structural error" + "Tip that doesn't descend from folded tip anymore is structural error" ); assert_eq!( index.effective_collaborator(&repo, &acc("lyna")), @@ -248,7 +249,7 @@ fn a_diverged_members_tip_fails_closed_to_warming() { assert!( index.refresh_members().is_err(), - "tip that no longer descends from folded tip is structural error" + "Tip that doesn't descend from folded tip anymore is structural error" ); assert_eq!(index.coverage().members, Coverage::Warming); assert_eq!( @@ -378,7 +379,7 @@ fn deregister_purges_collaborators_fail_closed() { assert_eq!( index.effective_collaborator(&repo, &acc("lyna")), Resolved::Warming, - "deregistered repo's collaborators are purged and fail closed, not served stale" + "Deregistered repo's collaborators are purged and fail closed, not served stale" ); } @@ -433,7 +434,7 @@ fn a_renamed_repo_keeps_both_rkeys_and_its_collaborators() { assert_eq!( index.resolve_repo(&own("nel"), &rkey("anemone")), Resolved::Ready(Some(repo.clone())), - "prior rkey keeps resolving as alias after rename is delta-applied" + "Prior rkey keeps resolving as alias after rename is delta-applied" ); assert_eq!( index.resolve_repo(&own("nel"), &rkey("barnacle")), @@ -442,12 +443,12 @@ fn a_renamed_repo_keeps_both_rkeys_and_its_collaborators() { assert_eq!( index.rkey_of(&repo), Resolved::Ready(Some(rkey("barnacle"))), - "new rkey is canonical" + "New rkey is canonical" ); assert_eq!( index.effective_collaborator(&repo, &acc("lyna")), Resolved::Ready(true), - "rename never evacuates repo, so its collaborators survive" + "Rename never evacuates repo, so its collaborators survive" ); store @@ -466,12 +467,12 @@ fn a_renamed_repo_keeps_both_rkeys_and_its_collaborators() { assert_eq!( index.resolve_repo(&own("nel"), &rkey("barnacle")), Resolved::Ready(None), - "deregistering through retained alias removes repo and every alias" + "Deregistering through retained alias removes repo and every alias" ); assert_eq!( index.effective_collaborator(&repo, &acc("lyna")), Resolved::Warming, - "deregistered repo's collaborators are evacuated" + "Deregistered repo's collaborators are evacuated" ); } @@ -563,12 +564,12 @@ fn recording_and_retaining_track_who_still_publishes_each_key() { assert_eq!( index.owner_of_key(&key(0), at(0)), Resolved::Ready(None), - "a key the account stopped publishing stops resolving to it" + "A key the account stopped publishing stops resolving to it" ); assert_eq!( index.owner_of_key(&key(1), at(0)), Resolved::Ready(Some(acc("nel"))), - "a key present in both sets survives the replacement" + "A key present in both sets survives the replacement" ); index @@ -585,14 +586,14 @@ fn recording_and_retaining_track_who_still_publishes_each_key() { assert_eq!( index.keys().publisher_among(&[acc("teq")], &key(1), at(0)), None, - "the key mustn't answer for an account that doesn't publish it" + "The key mustn't answer for an account that doesn't publish it" ); index.keys().retain(&KeptAccounts::new(vec![acc("nel")])); assert_eq!( index.owner_of_key(&key(3), at(0)), Resolved::Ready(None), - "an account outside the grant set has its keys released" + "An account outside the grant set has its keys released" ); assert_eq!( index.owner_of_key(&key(1), at(0)), @@ -605,7 +606,7 @@ fn recording_and_retaining_track_who_still_publishes_each_key() { assert_eq!( index.owner_of_key(&key(1), at(0)), Resolved::Ready(None), - "with its last publisher no longer offering it, the key stops answering for anybody" + "With its last publisher not offering it anymore, the key stops answering for anybody" ); } @@ -621,7 +622,7 @@ fn full_budget_reports_unheld_readings_and_then_stops_growing_entirely() { ); assert!( !index.keys().any_unheld(), - "a set that kept every reading it took can vouch for its misses" + "A set that kept every reading it took can vouch for its misses" ); assert_eq!( @@ -643,29 +644,29 @@ fn full_budget_reports_unheld_readings_and_then_stops_growing_entirely() { assert_eq!( index.owner_of_key(&fat(2), at(0)), Resolved::Ready(None), - "a key the budget refused mustn't answer for anybody" + "A key the budget refused mustn't answer for anybody" ); assert_eq!( index.keys().record(&acc("teq"), vec![fat(3)], hour(0)), KeyRecord::Saturated, - "a knot whose grant set publishes more key bytes than it budgeted for must stop growing" + "A knot whose grant set publishes more key bytes than it budgeted for must stop growing" ); assert_eq!( index.owner_of_key(&fat(1), at(0)), Resolved::Ready(Some(acc("nel"))), - "refusing the accounts that didn't fit mustn't release the accounts already on file" + "Refusing the accounts that didn't fit mustn't release the accounts already on file" ); index.keys().retain(&KeptAccounts::new(vec![acc("cuttle")])); assert!( index.keys().any_unheld(), - "releasing a different account mustn't clear the report while the unheld reading stays" + "Releasing a different account mustn't clear the report while the unheld reading stays" ); assert_eq!( index.keys().record(&acc("cuttle"), vec![fat(2)], hour(0)), KeyRecord::Stored, - "the budget frees up with the accounts that leave the grant set" + "The budget frees up with the accounts that leave the grant set" ); assert!( !index.keys().any_unheld(), @@ -704,13 +705,13 @@ fn a_pass_renews_at_the_lease_halfway_and_rereads_once_per_miss_outside_the_floo index.keys().note_miss(); let Resolved::Ready(work) = index.keys().work(at(100), sweep()) else { - panic!("the registry is folded, so the knot knows who may push"); + panic!("Registry is folded, so the knot knows who may push"); }; assert_eq!( work.suspected.len(), 1, - "a key the accounts on file don't publish is the only evidence the knot gets that an \ - account published a new key, so the next pass must reread rather than wait out the ttl" + "A key the accounts on file don't publish is the only evidence the knot gets that an \ + account published a new key, so the next pass must reread rather than wait out the TTL" ); assert!( work.due.is_empty(), @@ -739,7 +740,7 @@ fn a_pass_renews_at_the_lease_halfway_and_rereads_once_per_miss_outside_the_floo "an account inside the first half of its lease is left alone" ); let Resolved::Ready(work) = index.keys().work(at(1_800), sweep()) else { - panic!("the registry is folded, so the knot knows who may push"); + panic!("Registry is folded, so the knot knows who may push"); }; assert_eq!( work.due.len(), @@ -768,7 +769,7 @@ fn an_unfolded_repo_keeps_the_pass_partial_and_an_unopenable_repo_never_grants() index.refresh_collaborators(&folded).unwrap(); let Resolved::Ready(work) = index.keys().work(at(0), sweep()) else { - panic!("one repo the knot hasn't folded mustn't stop it filling the keys it can enumerate"); + panic!("One repo the knot hasn't folded mustn't stop it filling the keys it can enumerate"); }; assert_eq!( work.hosted, @@ -792,7 +793,7 @@ fn an_unfolded_repo_keeps_the_pass_partial_and_an_unopenable_repo_never_grants() "the repo registered without ever being created is what the knot can't open" ); let Resolved::Ready(work) = index.keys().work(at(0), sweep()) else { - panic!("the registry is folded, so the knot knows who may push"); + panic!("Registry is folded, so the knot knows who may push"); }; assert_eq!( work.hosted, @@ -823,9 +824,9 @@ fn an_acl_write_puts_the_key_set_back_to_warming_even_mid_pass() { .keys() .record(&acc("nel"), vec![OfferedKey::from_bytes(vec![7])], hour(0)); let Resolved::Ready(work) = index.keys().work(at(0), sweep()) else { - panic!("the registry is folded, so the knot knows who may push"); + panic!("Registry is folded, so the knot knows who may push"); }; - assert!(work.complete, "the one pusher's keys are on file"); + assert!(work.complete, "The one pusher's keys are on file"); index.keys().mark_ready(work.generation); assert_eq!(index.keys().coverage(), Coverage::Ready); @@ -839,7 +840,7 @@ fn an_acl_write_puts_the_key_set_back_to_warming_even_mid_pass() { ); let Resolved::Ready(rereading) = index.keys().work(at(0), sweep()) else { - panic!("the registry is folded, so the knot knows who may push"); + panic!("Registry is folded, so the knot knows who may push"); }; world.add_member(roll, "teq", 5); index.refresh_members().unwrap(); @@ -880,7 +881,7 @@ fn an_account_the_budget_cant_fit_is_checked_against_its_pds_without_delaying_co its life accepting every key offered to it" ); let Resolved::Ready(later) = index.keys().work(at(3_601), sweep()) else { - panic!("the registry is folded, so the knot knows who may push"); + panic!("Registry is folded, so the knot knows who may push"); }; assert!( later.due.as_slice().contains(&acc("cuttle")), @@ -925,9 +926,9 @@ fn a_warming_member_roll_still_yields_the_accounts_that_may_push() { let (_world, index) = folded(); let Resolved::Ready(work) = index.keys().work(at(0), sweep()) else { - panic!("the registry is folded, so the knot knows who may push"); + panic!("Registry is folded, so the knot knows who may push"); }; - assert_eq!(work.pushers.len(), 1, "the one hosted repo has one owner"); + assert_eq!(work.pushers.len(), 1, "The one hosted repo has one owner"); assert_eq!( work.due.len(), 1, @@ -951,7 +952,7 @@ fn the_sweep_retains_keys_a_push_recorded_for_an_invited_collaborator() { let kept = || match index.keys().work(at(0), sweep()).map(|work| work.members) { Resolved::Ready(Resolved::Ready(members)) => members.kept, - _ => panic!("the registry or the member roll is still warming"), + _ => panic!("Registry or member roll is still warming"), }; assert!( !kept().as_slice().contains(&acc("lyna")), @@ -1675,20 +1676,27 @@ fn assert_projections_agree( objects: &[CobId], addresses: &[RecordAddress], ) { - [issue_collection(), state_collection(), comment_collection(), reaction_collection()] - .into_iter() - .for_each(|collection| { - assert_eq!( - incremental.social_of_kind(repo, &collection), - cold.social_of_kind(repo, &collection), - "Incremental replay and cold rebuild mustn't disagree about which {collection} objects exist" - ); - assert_eq!( - incremental.next_subject_number(repo, &collection), - cold.next_subject_number(repo, &collection), - "nor must they disagree about the next number {collection} hands out" - ); - }); + [ + issue_collection(), + state_collection(), + comment_collection(), + reaction_collection(), + label_def_collection(), + label_op_collection(), + ] + .into_iter() + .for_each(|collection| { + assert_eq!( + incremental.social_of_kind(repo, &collection), + cold.social_of_kind(repo, &collection), + "Incremental replay and cold rebuild mustn't disagree about which {collection} objects exist" + ); + assert_eq!( + incremental.next_subject_number(repo, &collection), + cold.next_subject_number(repo, &collection), + "nor must they disagree about the next number {collection} hands out" + ); + }); assert_eq!( incremental.social_structure(repo), cold.social_structure(repo), @@ -1716,6 +1724,11 @@ fn assert_projections_agree( cold.social_children_of(repo, at), "Incremental fold and cold fold disagree about children of {at}" ); + assert_eq!( + incremental.social_labels_of(repo, at), + cold.social_labels_of(repo, at), + "Incremental fold and cold fold disagree about the standing labels on {at}" + ); [ issue_collection(), state_collection(), @@ -1851,13 +1864,13 @@ fn erased_subject_stops_answering_children_and_comment_keeps_serving() { "the newest comment under the issue shoukd probably come back too" ); let before = - state_of(&index, &repo, comment_object).expect("the index should contain the comment"); + state_of(&index, &repo, comment_object).expect("Index should contain the comment"); world.erase::( &repo, issue_object, &acc("nel"), - version_of(&index, &repo, issue_object).expect("the issue should probably have a version"), + version_of(&index, &repo, issue_object).expect("Issue should probably have a version"), 30, ); index.refresh_social(&repo).unwrap(); @@ -1943,7 +1956,7 @@ fn reaction_index_answers_live_reactions_only() { &repo, nels_like, &acc("nel"), - version_of(&index, &repo, nels_like).expect("the reaction should have a version"), + version_of(&index, &repo, nels_like).expect("Reaction should have a version"), 50, ); index.refresh_social(&repo).unwrap(); @@ -2108,6 +2121,338 @@ fn git_rebuild_and_lone_rebuilds_agree_with_incremental_replay() { ); } +fn label_operand(key: &str, value: &str) -> knot_record::label::LabelOperand { + knot_record::label::LabelOperand::new( + knot_record::label::DefRef::parse( + &knot_types::AtUri::new_owned(format!( + "at://did:plc:limpet/sh.tangled.label.definition/{key}" + )) + .unwrap(), + ) + .unwrap(), + knot_record::label::OperandValue::new(value).unwrap(), + ) +} + +fn label_body( + subject: &knot_types::RecordAddress, + adds: Vec, + deletes: Vec, + defs: &[(&str, bool)], +) -> knot_cobs::LabelOp { + let declared: std::collections::BTreeMap< + knot_record::label::DefRef, + knot_record::label::Def, + > = defs + .iter() + .map(|(slug, multiple)| { + ( + label_operand(slug, "null").key().clone(), + knot_record::label::Def::new( + knot_record::label::LabelName::new(*slug).unwrap(), + knot_record::label::LabelValue::Text(knot_record::label::TextConstraint::free( + knot_record::label::LabelFormat::Any, + )), + vec![knot_record::issue::issue_collection()], + None, + *multiple, + ) + .unwrap(), + ) + }) + .collect(); + let record = knot_record::label::LabelOpRecord::new( + knot_types::RepoDid::new("did:plc:squid").unwrap(), + subject.clone(), + adds, + deletes, + knot_types::UnixSeconds::new(20), + acc("nel"), + ) + .unwrap(); + knot_cobs::LabelOp::applications(&record, &declared).unwrap() +} + +fn spelled( + labels: std::collections::BTreeMap< + knot_record::label::DefRef, + Vec, + >, +) -> std::collections::BTreeMap> { + labels + .into_iter() + .map(|(def, values)| { + ( + def.to_string(), + values + .into_iter() + .map(|value| value.as_str().to_owned()) + .collect(), + ) + }) + .collect() +} + +#[test] +fn standing_labels_match_a_consumer_replaying_the_ops() { + let world = World::new(); + let repo = repo_did("squid"); + world.issue(&repo, "3lubrptx57d22", &acc("nel"), 1, "kelp", 10); + + let subject = issue_at("3lubrptx57d22"); + let single = "size"; + let multiple = "reviewer"; + let declared = [(single, false), (multiple, true)]; + let ops = [ + ( + vec![ + label_operand(multiple, "conch"), + label_operand(multiple, "limpet"), + ], + Vec::new(), + ), + ( + vec![ + label_operand(single, "small"), + label_operand(single, "large"), + ], + Vec::new(), + ), + ( + vec![label_operand(multiple, "conch")], + vec![label_operand(multiple, "conch")], + ), + (Vec::new(), vec![label_operand(single, "small")]), + ]; + [ + ("3lubrptx57d33", &ops[0]), + ("3lubrptx57d44", &ops[1]), + ("3lubrptx57d55", &ops[2]), + ("3lubrptx57d66", &ops[3]), + ] + .into_iter() + .enumerate() + .for_each(|(at, (key, (adds, deletes)))| { + world.label_op( + &repo, + key, + &acc("nel"), + (at + 2) as u32, + &subject, + label_body(&subject, adds.clone(), deletes.clone(), &declared), + 20 + at as i64, + ); + }); + + let index = world.index(); + index.refresh_social(&repo).unwrap(); + + fn consumer( + ops: &[( + Vec, + Vec, + )], + multiple_of: impl Fn(&str) -> bool, + ) -> std::collections::BTreeMap> { + let mut standing: std::collections::BTreeMap> = + std::collections::BTreeMap::new(); + ops.iter().for_each(|(adds, deletes)| { + adds.iter() + .map(|operand| (operand, knot_cobs::LabelVerb::Add)) + .chain( + deletes + .iter() + .map(|operand| (operand, knot_cobs::LabelVerb::Delete)), + ) + .for_each(|(operand, verb)| { + let key = operand.key().to_string(); + let value = operand.value().as_str().to_owned(); + let set = standing.entry(key.clone()).or_default(); + let multiple = multiple_of(&key); + if verb == knot_cobs::LabelVerb::Delete { + if multiple { + set.retain(|standing| *standing != value); + } else if set.first().is_some_and(|standing| *standing == value) { + set.clear(); + } + } else if multiple { + if !set.contains(&value) { + set.push(value); + } + } else { + set.clear(); + set.push(value); + } + }); + }); + standing.retain(|_, values| !values.is_empty()); + standing + } + + assert_eq!( + index.social_labels_of(&repo, &subject).map(spelled), + Resolved::Ready(consumer(&ops, |key| key.ends_with("/reviewer"))), + "the projection's standing set is what a listOps consumer computes from the same records" + ); + let expected = Resolved::Ready(std::collections::BTreeMap::from([ + ( + "at://did:plc:limpet/sh.tangled.label.definition/reviewer".to_owned(), + vec!["limpet".to_owned()], + ), + ( + "at://did:plc:limpet/sh.tangled.label.definition/size".to_owned(), + vec!["large".to_owned()], + ), + ])); + assert_eq!( + index.social_labels_of(&repo, &subject).map(spelled), + expected, + "reviewer keeps limpet after conch's add-then-delete in one op, size stands at large \ + after the last add replaced small, and small's delete is a no-op" + ); +} + +#[test] +fn an_erased_op_removes_the_label_from_the_standing_set() { + let world = World::new(); + let repo = repo_did("squid"); + world.issue(&repo, "3lubrptx57d22", &acc("nel"), 1, "kelp", 10); + + let subject = issue_at("3lubrptx57d22"); + let op = world.label_op( + &repo, + "3lubrptx57d33", + &acc("nel"), + 2, + &subject, + label_body( + &subject, + vec![label_operand("wontfix", "null")], + Vec::new(), + &[("wontfix", false)], + ), + 20, + ); + + let index = world.index(); + index.refresh_social(&repo).unwrap(); + let (_, standing) = note_record("3lubrptx57d33"); + assert_eq!( + index.social_labels_of(&repo, &subject).map(spelled), + Resolved::Ready(std::collections::BTreeMap::from([( + "at://did:plc:limpet/sh.tangled.label.definition/wontfix".to_owned(), + vec!["null".to_owned()] + )])), + "before the erase, the label the op added is standing" + ); + + world.erase::(&repo, op, &acc("nel"), standing, 30); + index + .refresh_social_object(&repo, &label_op_type_name(), op) + .unwrap(); + assert_eq!( + index.social_labels_of(&repo, &subject), + Resolved::Ready(std::collections::BTreeMap::new()), + "an erased op takes the label out of the standing set" + ); +} + +#[test] +fn incremental_label_refresh_matches_a_cold_fold() { + let world = World::new(); + let repo = repo_did("squid"); + world.issue(&repo, "3lubrptx57d22", &acc("nel"), 1, "kelp", 10); + let subject = issue_at("3lubrptx57d22"); + + let warmed = world.index(); + warmed.refresh_social(&repo).unwrap(); + let op = world.label_op( + &repo, + "3lubrptx57d33", + &acc("nel"), + 2, + &subject, + label_body( + &subject, + vec![ + label_operand("wontfix", "null"), + label_operand("reviewer", "conch"), + ], + Vec::new(), + &[("wontfix", false), ("reviewer", true)], + ), + 20, + ); + warmed + .refresh_social_object(&repo, &label_op_type_name(), op) + .unwrap(); + let cold = world.index(); + cold.refresh_social(&repo).unwrap(); + assert_eq!( + cold.social_labels_of(&repo, &subject), + warmed.social_labels_of(&repo, &subject), + "a label op folded incrementally onto a warm projection serves the same standing set \ + as a cold fold from the git history" + ); + + let (_, standing) = note_record("3lubrptx57d33"); + world.erase::(&repo, op, &acc("nel"), standing, 30); + warmed + .refresh_social_object(&repo, &label_op_type_name(), op) + .unwrap(); + cold.refresh_social(&repo).unwrap(); + assert_eq!( + cold.social_labels_of(&repo, &subject), + warmed.social_labels_of(&repo, &subject), + "nor do they disagree once the op is erased incrementally" + ); +} + +#[test] +fn deleted_subject_stops_serving_its_labels() { + let world = World::new(); + let repo = repo_did("squid"); + let issue = world.issue(&repo, "3lubrptx57d22", &acc("nel"), 1, "kelp", 10); + + let subject = issue_at("3lubrptx57d22"); + world.label_op( + &repo, + "3lubrptx57d33", + &acc("nel"), + 2, + &subject, + label_body( + &subject, + vec![label_operand("wontfix", "null")], + Vec::new(), + &[("wontfix", false)], + ), + 20, + ); + + let index = world.index(); + index.refresh_social(&repo).unwrap(); + assert_eq!( + index.social_labels_of(&repo, &subject).map(spelled), + Resolved::Ready(std::collections::BTreeMap::from([( + "at://did:plc:limpet/sh.tangled.label.definition/wontfix".to_owned(), + vec!["null".to_owned()] + )])), + "while the issue stands, so does the label" + ); + + let (_, standing) = note_record("kelp"); + world.erase::(&repo, issue, &acc("nel"), standing, 30); + index + .refresh_social_object(&repo, &issue_type_name(), issue) + .unwrap(); + assert_eq!( + index.social_labels_of(&repo, &subject), + Resolved::Ready(std::collections::BTreeMap::new()), + "a deleted subject stops serving the label set along with the children" + ); +} + #[derive(Serialize, Deserialize)] #[serde(tag = "op", content = "data", rename_all = "snake_case")] enum BadIssue { @@ -2200,6 +2545,7 @@ fn contested_claims_stop_counting_toward_serving_bytes() { let world = World::new(); let repo = repo_did("squid"); world.issue(&repo, "3lubrptx57d22", &acc("nel"), 1, "kelp", 10); + world.issue(&repo, "3lubrptx57d22", &acc("olaren"), 2, "urchin", 20); let index = world.index(); index.refresh_social(&repo).unwrap();