This repository has no description
24 kB
745 lines
1use std::collections::{BTreeMap, BTreeSet};
2use std::hash::{Hash, Hasher};
3use std::marker::PhantomData;
4use std::sync::Mutex;
5
6use knot_cache::{Cache, EntryCount, Lru};
7use knot_cob::{Change, ChangeId, ChangePayload, Checkpoint, CobId, CobStore, Evaluate};
8use knot_cobs::{
9 CollaboratorsChange, CollaboratorsCob, Grant, GrantChange, Registration, Registry,
10 RegistryChange, Rename, RepoRef, RepoRegistryCob, Roster,
11};
12use knot_types::{AccountDid, OfferedKey, OwnerDid, RepoDid, RepoRkey, UnixSeconds};
13
14use crate::coverage::{Coverage, CoverageCell, Resolved};
15use crate::error::IndexError;
16use crate::intern::{AccountKey, Interner, OwnerKey, RepoKey, RkeyKey};
17
18const KEY_CACHE_CAPACITY: usize = 16_384;
19
20#[derive(Debug, Clone, Copy)]
21struct Provenance {
22 added_by: AccountKey,
23 created_at: UnixSeconds,
24}
25
26impl Provenance {
27 fn intern(interner: &Interner, grant: &Grant) -> Self {
28 Self {
29 added_by: interner.intern_account(&grant.added_by),
30 created_at: grant.created_at,
31 }
32 }
33
34 fn grant(self, interner: &Interner, subject: AccountKey) -> Grant {
35 Grant {
36 subject: interner.resolve_account(subject),
37 added_by: interner.resolve_account(self.added_by),
38 created_at: self.created_at,
39 }
40 }
41}
42
43fn decode_change<P: ChangePayload>(change: &Change) -> Result<P, IndexError> {
44 if change.type_name != P::type_name() {
45 return Err(IndexError::UnexpectedType {
46 change: change.id,
47 expected: P::type_name(),
48 found: change.type_name.clone(),
49 });
50 }
51 P::decode(change.payload()).map_err(|error| IndexError::Decode {
52 change: change.id,
53 type_name: P::type_name(),
54 reason: error.to_string(),
55 })
56}
57
58fn decode_delta<P: ChangePayload>(changes: &[Change]) -> Result<Vec<P>, IndexError> {
59 changes.iter().map(decode_change::<P>).collect()
60}
61
62fn group_by<K: Ord, T>(items: Vec<T>, key: impl Fn(&T) -> K) -> BTreeMap<K, Vec<T>> {
63 items.into_iter().fold(BTreeMap::new(), |mut acc, item| {
64 acc.entry(key(&item)).or_default().push(item);
65 acc
66 })
67}
68
69pub(crate) struct GrantSetProjection<Cob>
70where
71 Cob: Evaluate,
72 Cob::Change: ChangePayload + GrantChange,
73{
74 membership: scc::HashMap<AccountKey, Provenance>,
75 coverage: CoverageCell,
76 tip: Mutex<Option<ChangeId>>,
77 _cob: PhantomData<fn() -> Cob>,
78}
79
80impl<Cob> GrantSetProjection<Cob>
81where
82 Cob: Evaluate,
83 Cob::Change: ChangePayload + GrantChange,
84{
85 pub(crate) fn new() -> Self {
86 Self {
87 membership: scc::HashMap::new(),
88 coverage: CoverageCell::new(Coverage::Warming),
89 tip: Mutex::new(None),
90 _cob: PhantomData,
91 }
92 }
93
94 pub(crate) fn coverage(&self) -> Coverage {
95 self.coverage.get()
96 }
97
98 pub(crate) fn reset(&self) {
99 let mut tip = self
100 .tip
101 .lock()
102 .unwrap_or_else(|poisoned| poisoned.into_inner());
103 self.membership.clear_sync();
104 *tip = None;
105 self.coverage.set(Coverage::Ready);
106 }
107
108 pub(crate) fn contains(&self, interner: &Interner, did: &AccountDid) -> Resolved<bool> {
109 match self.coverage.get() {
110 Coverage::Warming => Resolved::Warming,
111 Coverage::Ready => Resolved::Ready(
112 interner
113 .account(did)
114 .is_some_and(|did| self.membership.contains_sync(&did)),
115 ),
116 }
117 }
118
119 pub(crate) fn entries(&self, interner: &Interner) -> Resolved<Vec<Grant>> {
120 match self.coverage.get() {
121 Coverage::Warming => Resolved::Warming,
122 Coverage::Ready => {
123 let mut out = BTreeMap::new();
124 self.membership.iter_sync(|&subject, slot| {
125 let grant = slot.grant(interner, subject);
126 out.insert(grant.subject.clone(), grant);
127 true
128 });
129 Resolved::Ready(out.into_values().collect())
130 }
131 }
132 }
133
134 pub(crate) fn refresh(
135 &self,
136 interner: &Interner,
137 store: &CobStore,
138 object: CobId,
139 ) -> Result<(), IndexError>
140 where
141 Cob: Checkpoint + Evaluate<State = Roster>,
142 {
143 let mut tip = self
144 .tip
145 .lock()
146 .unwrap_or_else(|poisoned| poisoned.into_inner());
147 match *tip {
148 None => {
149 let (roster, seeded) = store
150 .materialize::<Cob>(object)
151 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
152 self.seed(interner, &roster);
153 *tip = Some(seeded);
154 }
155 Some(prev) => {
156 let delta = store
157 .changes_since::<Cob>(object, Some(prev))
158 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
159 let decoded = decode_delta::<Cob::Change>(&delta.changes)
160 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
161 self.apply_delta(interner, decoded);
162 *tip = Some(delta.tip);
163 }
164 }
165 self.coverage.set(Coverage::Ready);
166 Ok(())
167 }
168
169 fn seed(&self, interner: &Interner, roster: &Roster) {
170 self.membership.clear_sync();
171 roster.entries().for_each(|(subject, entry)| {
172 let slot = Provenance {
173 added_by: interner.intern_account(&entry.added_by),
174 created_at: entry.created_at,
175 };
176 let _ = self
177 .membership
178 .insert_sync(interner.intern_account(subject), slot);
179 });
180 }
181
182 fn apply_delta(&self, interner: &Interner, changes: Vec<Cob::Change>) {
183 group_by(changes, |change| change.subject().clone())
184 .into_iter()
185 .for_each(|(did, ops)| {
186 let current = interner
187 .account(&did)
188 .and_then(|key| self.membership.read_sync(&key, |_, slot| *slot));
189 let net = ops
190 .into_iter()
191 .fold(current, |slot, change| match change.as_grant() {
192 Some(grant) => slot.or_else(|| Some(Provenance::intern(interner, grant))),
193 None => None,
194 });
195 match net {
196 Some(slot) => {
197 *self
198 .membership
199 .entry_sync(interner.intern_account(&did))
200 .or_insert(slot)
201 .get_mut() = slot;
202 }
203 None => {
204 if let Some(key) = interner.account(&did) {
205 let _ = self.membership.remove_sync(&key);
206 }
207 }
208 }
209 });
210 }
211}
212
213struct RepoRoster {
214 tip: Option<ChangeId>,
215 entries: BTreeMap<AccountKey, Provenance>,
216}
217
218const REPO_LOCK_STRIPES: usize = 256;
219
220pub(crate) struct CollaboratorsProjection {
221 rosters: scc::HashMap<RepoKey, RepoRoster>,
222 locks: [Mutex<()>; REPO_LOCK_STRIPES],
223}
224
225impl CollaboratorsProjection {
226 pub(crate) fn new() -> Self {
227 Self {
228 rosters: scc::HashMap::new(),
229 locks: std::array::from_fn(|_| Mutex::new(())),
230 }
231 }
232
233 fn repo_lock(&self, repo: RepoKey) -> &Mutex<()> {
234 let mut hasher = std::collections::hash_map::DefaultHasher::new();
235 repo.hash(&mut hasher);
236 &self.locks[(hasher.finish() % REPO_LOCK_STRIPES as u64) as usize]
237 }
238
239 pub(crate) fn coverage(&self) -> Coverage {
240 Coverage::Ready
241 }
242
243 pub(crate) fn is_folded(&self, interner: &Interner, repo: &RepoDid) -> bool {
244 interner
245 .repo(repo)
246 .is_some_and(|repo| self.rosters.contains_sync(&repo))
247 }
248
249 pub(crate) fn contains(
250 &self,
251 interner: &Interner,
252 repo: &RepoDid,
253 did: &AccountDid,
254 ) -> Resolved<bool> {
255 let Some(repo) = interner.repo(repo) else {
256 return Resolved::Warming;
257 };
258 match self.rosters.read_sync(&repo, |_, roster| {
259 interner
260 .account(did)
261 .is_some_and(|account| roster.entries.contains_key(&account))
262 }) {
263 Some(present) => Resolved::Ready(present),
264 None => Resolved::Warming,
265 }
266 }
267
268 pub(crate) fn entries(&self, interner: &Interner, repo: &RepoDid) -> Resolved<Vec<Grant>> {
269 let Some(repo) = interner.repo(repo) else {
270 return Resolved::Warming;
271 };
272 match self
273 .rosters
274 .read_sync(&repo, |_, roster| roster_grants(interner, roster))
275 {
276 Some(grants) => Resolved::Ready(grants),
277 None => Resolved::Warming,
278 }
279 }
280
281 pub(crate) fn mark_repo_empty(&self, repo: RepoKey) {
282 let lock = self.repo_lock(repo);
283 let _guard = lock.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
284 self.install(
285 repo,
286 RepoRoster {
287 tip: None,
288 entries: BTreeMap::new(),
289 },
290 );
291 }
292
293 pub(crate) fn refresh_repo(
294 &self,
295 interner: &Interner,
296 store: &CobStore,
297 repo: RepoKey,
298 object: CobId,
299 ) -> Result<(), IndexError> {
300 let lock = self.repo_lock(repo);
301 let _guard = lock.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
302 self.fold_repo(interner, store, repo, object)
303 .inspect_err(|_| self.purge_repo(repo))
304 }
305
306 fn fold_repo(
307 &self,
308 interner: &Interner,
309 store: &CobStore,
310 repo: RepoKey,
311 object: CobId,
312 ) -> Result<(), IndexError> {
313 let prev = self
314 .rosters
315 .read_sync(&repo, |_, roster| roster.tip)
316 .flatten();
317 match prev {
318 None => {
319 let (roster, tip) = store.materialize::<CollaboratorsCob>(object)?;
320 let entries = roster
321 .entries()
322 .map(|(subject, entry)| {
323 (
324 interner.intern_account(subject),
325 Provenance {
326 added_by: interner.intern_account(&entry.added_by),
327 created_at: entry.created_at,
328 },
329 )
330 })
331 .collect();
332 self.install(
333 repo,
334 RepoRoster {
335 tip: Some(tip),
336 entries,
337 },
338 );
339 }
340 Some(prev) => {
341 let delta = store.changes_since::<CollaboratorsCob>(object, Some(prev))?;
342 let decoded = decode_delta::<CollaboratorsChange>(&delta.changes)?;
343 let mut occupied = self.rosters.entry_sync(repo).or_insert_with(|| RepoRoster {
344 tip: None,
345 entries: BTreeMap::new(),
346 });
347 let roster = occupied.get_mut();
348 roster.entries = decoded
349 .into_iter()
350 .fold(std::mem::take(&mut roster.entries), |entries, change| {
351 apply(interner, entries, change)
352 });
353 roster.tip = Some(delta.tip);
354 }
355 }
356 Ok(())
357 }
358
359 pub(crate) fn drop_repo(&self, repo: RepoKey) {
360 let lock = self.repo_lock(repo);
361 let _guard = lock.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
362 self.purge_repo(repo);
363 }
364
365 fn purge_repo(&self, repo: RepoKey) {
366 let _ = self.rosters.remove_sync(&repo);
367 }
368
369 fn install(&self, repo: RepoKey, roster: RepoRoster) {
370 match self.rosters.entry_sync(repo) {
371 scc::hash_map::Entry::Occupied(mut occupied) => {
372 let _ = occupied.insert(roster);
373 }
374 scc::hash_map::Entry::Vacant(vacant) => {
375 vacant.insert_entry(roster);
376 }
377 }
378 }
379}
380
381fn roster_grants(interner: &Interner, roster: &RepoRoster) -> Vec<Grant> {
382 roster
383 .entries
384 .iter()
385 .map(|(&subject, slot)| {
386 let grant = slot.grant(interner, subject);
387 (grant.subject.clone(), grant)
388 })
389 .collect::<BTreeMap<_, _>>()
390 .into_values()
391 .collect()
392}
393
394fn apply(
395 interner: &Interner,
396 mut entries: BTreeMap<AccountKey, Provenance>,
397 change: CollaboratorsChange,
398) -> BTreeMap<AccountKey, Provenance> {
399 match change {
400 CollaboratorsChange::Add(grant) => {
401 entries
402 .entry(interner.intern_account(&grant.subject))
403 .or_insert_with(|| Provenance::intern(interner, &grant));
404 }
405 CollaboratorsChange::Remove(removal) => {
406 if let Some(key) = interner.account(&removal.subject) {
407 entries.remove(&key);
408 }
409 }
410 }
411 entries
412}
413
414struct RecordSlot {
415 owner: OwnerKey,
416 rkey: RkeyKey,
417}
418
419pub(crate) struct RegistryProjection {
420 aliases: scc::HashMap<(OwnerKey, RkeyKey), RepoKey>,
421 records: scc::HashMap<RepoKey, RecordSlot>,
422 coverage: CoverageCell,
423 tip: Mutex<Option<ChangeId>>,
424}
425
426impl RegistryProjection {
427 pub(crate) fn new() -> Self {
428 Self {
429 aliases: scc::HashMap::new(),
430 records: scc::HashMap::new(),
431 coverage: CoverageCell::new(Coverage::Warming),
432 tip: Mutex::new(None),
433 }
434 }
435
436 pub(crate) fn coverage(&self) -> Coverage {
437 self.coverage.get()
438 }
439
440 pub(crate) fn reset(&self, interner: &Interner) -> Vec<RepoDid> {
441 let mut tip = self
442 .tip
443 .lock()
444 .unwrap_or_else(|poisoned| poisoned.into_inner());
445 let evacuated = self.hosted_repos(interner);
446 self.aliases.clear_sync();
447 self.records.clear_sync();
448 *tip = None;
449 self.coverage.set(Coverage::Ready);
450 evacuated
451 }
452
453 pub(crate) fn resolve(
454 &self,
455 interner: &Interner,
456 owner: &OwnerDid,
457 rkey: &RepoRkey,
458 ) -> Resolved<Option<RepoDid>> {
459 match self.coverage.get() {
460 Coverage::Warming => Resolved::Warming,
461 Coverage::Ready => Resolved::Ready(
462 interner
463 .owner(owner)
464 .zip(interner.rkey(rkey))
465 .and_then(|key| self.aliases.read_sync(&key, |_, repo| *repo))
466 .map(|repo| interner.resolve_repo(repo)),
467 ),
468 }
469 }
470
471 pub(crate) fn owner_of(
472 &self,
473 interner: &Interner,
474 repo: &RepoDid,
475 ) -> Resolved<Option<OwnerDid>> {
476 if self.coverage.get() == Coverage::Warming {
477 return Resolved::Warming;
478 }
479 let Some(target) = interner.repo(repo) else {
480 return Resolved::Ready(None);
481 };
482 Resolved::Ready(
483 self.records
484 .read_sync(&target, |_, slot| slot.owner)
485 .map(|owner| interner.resolve_owner(owner)),
486 )
487 }
488
489 pub(crate) fn rkey_of(
490 &self,
491 interner: &Interner,
492 repo: &RepoDid,
493 ) -> Resolved<Option<RepoRkey>> {
494 if self.coverage.get() == Coverage::Warming {
495 return Resolved::Warming;
496 }
497 let Some(target) = interner.repo(repo) else {
498 return Resolved::Ready(None);
499 };
500 Resolved::Ready(
501 self.records
502 .read_sync(&target, |_, slot| slot.rkey)
503 .map(|rkey| interner.resolve_rkey(rkey)),
504 )
505 }
506
507 pub(crate) fn hosted_repos(&self, interner: &Interner) -> Vec<RepoDid> {
508 let mut repos = BTreeSet::new();
509 self.records.iter_sync(|repo, _| {
510 repos.insert(interner.resolve_repo(*repo));
511 true
512 });
513 repos.into_iter().collect()
514 }
515
516 pub(crate) fn refresh(
517 &self,
518 interner: &Interner,
519 store: &CobStore,
520 object: CobId,
521 ) -> Result<Vec<RepoDid>, IndexError> {
522 let mut tip = self
523 .tip
524 .lock()
525 .unwrap_or_else(|poisoned| poisoned.into_inner());
526 match *tip {
527 None => {
528 let (registry, seeded) = store
529 .materialize::<RepoRegistryCob>(object)
530 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
531 self.seed(interner, ®istry);
532 *tip = Some(seeded);
533 self.coverage.set(Coverage::Ready);
534 Ok(Vec::new())
535 }
536 Some(prev) => {
537 let delta = store
538 .changes_since::<RepoRegistryCob>(object, Some(prev))
539 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
540 let decoded = decode_delta::<RegistryChange>(&delta.changes)
541 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
542 let displaced = self.apply_delta(interner, decoded);
543 *tip = Some(delta.tip);
544 self.coverage.set(Coverage::Ready);
545 Ok(self.evacuated(interner, displaced))
546 }
547 }
548 }
549
550 fn seed(&self, interner: &Interner, registry: &Registry) {
551 self.aliases.clear_sync();
552 self.records.clear_sync();
553 registry.records().for_each(|(repo, record)| {
554 self.upsert_record(
555 interner.intern_repo(repo),
556 interner.intern_owner(&record.owner),
557 interner.intern_rkey(&record.rkey),
558 );
559 });
560 registry.aliases().for_each(|(owner, rkey, repo)| {
561 self.upsert_alias(
562 interner.intern_owner(owner),
563 interner.intern_rkey(rkey),
564 interner.intern_repo(repo),
565 );
566 });
567 }
568
569 fn evacuated(&self, interner: &Interner, displaced: Vec<RepoKey>) -> Vec<RepoDid> {
570 if displaced.is_empty() {
571 return Vec::new();
572 }
573 let live = self.live_repos();
574 displaced
575 .into_iter()
576 .collect::<BTreeSet<_>>()
577 .into_iter()
578 .filter(|repo| !live.contains(repo))
579 .map(|repo| interner.resolve_repo(repo))
580 .collect()
581 }
582
583 fn live_repos(&self) -> BTreeSet<RepoKey> {
584 let mut live = BTreeSet::new();
585 self.records.iter_sync(|repo, _| {
586 live.insert(*repo);
587 true
588 });
589 live
590 }
591
592 fn apply_delta(&self, interner: &Interner, changes: Vec<RegistryChange>) -> Vec<RepoKey> {
593 changes
594 .into_iter()
595 .fold(Vec::new(), |displaced, change| match change {
596 RegistryChange::Register(registration) => {
597 self.apply_register(interner, registration, displaced)
598 }
599 RegistryChange::Rename(rename) => self.apply_rename(interner, rename, displaced),
600 RegistryChange::Deregister(target) => {
601 self.apply_deregister(interner, target, displaced)
602 }
603 })
604 }
605
606 fn apply_register(
607 &self,
608 interner: &Interner,
609 registration: Registration,
610 mut displaced: Vec<RepoKey>,
611 ) -> Vec<RepoKey> {
612 let repo = interner.intern_repo(®istration.repo);
613 let owner = interner.intern_owner(®istration.owner);
614 let rkey = interner.intern_rkey(®istration.rkey);
615 if self.records.contains_sync(&repo) {
616 self.drop_record(repo);
617 displaced.push(repo);
618 }
619 displaced.extend(self.steal_alias(owner, rkey, repo));
620 self.upsert_record(repo, owner, rkey);
621 self.upsert_alias(owner, rkey, repo);
622 displaced
623 }
624
625 fn apply_rename(
626 &self,
627 interner: &Interner,
628 rename: Rename,
629 mut displaced: Vec<RepoKey>,
630 ) -> Vec<RepoKey> {
631 let repo = interner.intern_repo(&rename.repo);
632 let owner = interner.intern_owner(&rename.owner);
633 let rkey = interner.intern_rkey(&rename.rkey);
634 let held = self
635 .records
636 .read_sync(&repo, |_, slot| slot.owner == owner)
637 .unwrap_or(false);
638 if !held {
639 return displaced;
640 }
641 displaced.extend(self.steal_alias(owner, rkey, repo));
642 self.upsert_record(repo, owner, rkey);
643 self.upsert_alias(owner, rkey, repo);
644 displaced
645 }
646
647 fn apply_deregister(
648 &self,
649 interner: &Interner,
650 target: RepoRef,
651 mut displaced: Vec<RepoKey>,
652 ) -> Vec<RepoKey> {
653 let Some(owner) = interner.owner(&target.owner) else {
654 return displaced;
655 };
656 let Some(rkey) = interner.rkey(&target.rkey) else {
657 return displaced;
658 };
659 let Some(repo) = self.aliases.read_sync(&(owner, rkey), |_, repo| *repo) else {
660 return displaced;
661 };
662 self.drop_record(repo);
663 displaced.push(repo);
664 displaced
665 }
666
667 fn steal_alias(&self, owner: OwnerKey, rkey: RkeyKey, target: RepoKey) -> Option<RepoKey> {
668 let holder = self.aliases.read_sync(&(owner, rkey), |_, repo| *repo)?;
669 if holder == target {
670 return None;
671 }
672 let canonical = self
673 .records
674 .read_sync(&holder, |_, slot| slot.rkey == rkey)
675 .unwrap_or(false);
676 if canonical {
677 self.drop_record(holder);
678 Some(holder)
679 } else {
680 let _ = self.aliases.remove_sync(&(owner, rkey));
681 None
682 }
683 }
684
685 fn drop_record(&self, repo: RepoKey) {
686 let _ = self.records.remove_sync(&repo);
687 self.aliases.retain_sync(|_, holder| *holder != repo);
688 }
689
690 fn upsert_record(&self, repo: RepoKey, owner: OwnerKey, rkey: RkeyKey) {
691 if self
692 .records
693 .update_sync(&repo, |_, slot| {
694 slot.owner = owner;
695 slot.rkey = rkey;
696 })
697 .is_none()
698 {
699 let _ = self.records.insert_sync(repo, RecordSlot { owner, rkey });
700 }
701 }
702
703 fn upsert_alias(&self, owner: OwnerKey, rkey: RkeyKey, repo: RepoKey) {
704 let key = (owner, rkey);
705 if self
706 .aliases
707 .update_sync(&key, |_, slot| *slot = repo)
708 .is_none()
709 {
710 let _ = self.aliases.insert_sync(key, repo);
711 }
712 }
713}
714
715pub(crate) struct KeyProjection {
716 cache: Lru<OfferedKey, AccountKey>,
717}
718
719impl KeyProjection {
720 pub(crate) fn new() -> Self {
721 Self {
722 cache: Lru::by_count(EntryCount::new(KEY_CACHE_CAPACITY as u64)),
723 }
724 }
725
726 pub(crate) fn coverage(&self) -> Coverage {
727 Coverage::Ready
728 }
729
730 pub(crate) fn cache(&self, interner: &Interner, key: OfferedKey, did: &AccountDid) {
731 self.cache.insert(key, interner.intern_account(did));
732 }
733
734 pub(crate) fn owner(
735 &self,
736 interner: &Interner,
737 key: &OfferedKey,
738 ) -> Resolved<Option<AccountDid>> {
739 Resolved::Ready(
740 self.cache
741 .get(key)
742 .map(|account| interner.resolve_account(account)),
743 )
744 }
745}