This repository has no description
28 kB
851 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, ClonePath, OfferedKey, OwnerDid, RepoDid, RepoRkey, UnixSeconds};
13
14use crate::coverage::{Coverage, CoverageCell, Resolved};
15use crate::error::IndexError;
16use crate::intern::{AccountKey, Interner, NameKey, 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 name: NameKey,
418 created_at: UnixSeconds,
419}
420
421pub(crate) struct RegistryProjection {
422 aliases: scc::HashMap<(OwnerKey, RkeyKey), RepoKey>,
423 names: scc::HashMap<(OwnerKey, NameKey), BTreeSet<(UnixSeconds, RepoKey)>>,
424 records: scc::HashMap<RepoKey, RecordSlot>,
425 coverage: CoverageCell,
426 tip: Mutex<Option<ChangeId>>,
427}
428
429impl RegistryProjection {
430 pub(crate) fn new() -> Self {
431 Self {
432 aliases: scc::HashMap::new(),
433 names: scc::HashMap::new(),
434 records: scc::HashMap::new(),
435 coverage: CoverageCell::new(Coverage::Warming),
436 tip: Mutex::new(None),
437 }
438 }
439
440 pub(crate) fn coverage(&self) -> Coverage {
441 self.coverage.get()
442 }
443
444 pub(crate) fn reset(&self, interner: &Interner) -> Vec<RepoDid> {
445 let mut tip = self
446 .tip
447 .lock()
448 .unwrap_or_else(|poisoned| poisoned.into_inner());
449 let evacuated = self.hosted_repos(interner);
450 self.aliases.clear_sync();
451 self.names.clear_sync();
452 self.records.clear_sync();
453 *tip = None;
454 self.coverage.set(Coverage::Ready);
455 evacuated
456 }
457
458 pub(crate) fn resolve(
459 &self,
460 interner: &Interner,
461 owner: &OwnerDid,
462 rkey: &RepoRkey,
463 ) -> Resolved<Option<RepoDid>> {
464 match self.coverage.get() {
465 Coverage::Warming => Resolved::Warming,
466 Coverage::Ready => Resolved::Ready(
467 interner
468 .owner(owner)
469 .zip(interner.rkey(rkey))
470 .and_then(|key| self.aliases.read_sync(&key, |_, repo| *repo))
471 .map(|repo| interner.resolve_repo(repo)),
472 ),
473 }
474 }
475
476 pub(crate) fn resolve_clone_path(
477 &self,
478 interner: &Interner,
479 owner: &OwnerDid,
480 path: &ClonePath,
481 ) -> Resolved<Option<RepoDid>> {
482 if self.coverage.get() == Coverage::Warming {
483 return Resolved::Warming;
484 }
485 let Some(owner) = interner.owner(owner) else {
486 return Resolved::Ready(None);
487 };
488 let by_rkey = path
489 .rkeys()
490 .filter_map(|rkey| interner.rkey(rkey))
491 .find_map(|rkey| self.aliases.read_sync(&(owner, rkey), |_, repo| *repo));
492 if let Some(repo) = by_rkey {
493 return Resolved::Ready(Some(interner.resolve_repo(repo)));
494 }
495 Resolved::Ready(
496 path.names()
497 .filter_map(|name| interner.name(name))
498 .find_map(|name| self.oldest_registration_for_name(interner, owner, name)),
499 )
500 }
501
502 fn oldest_registration_for_name(
503 &self,
504 interner: &Interner,
505 owner: OwnerKey,
506 name: NameKey,
507 ) -> Option<RepoDid> {
508 // Notice how set keys sort by whichever string the interner saw first,
509 // and a cold rebuild will see them in a different order than live replay,
510 // so 2 entries with the same timestamp will settle by comparing
511 // DIDs instead.
512 self.names
513 .read_sync(&(owner, name), |_, registered| {
514 let earliest = registered.first()?.0;
515 registered
516 .iter()
517 .take_while(|(created_at, _)| *created_at == earliest)
518 .map(|(_, repo)| interner.resolve_repo(*repo))
519 .min()
520 })
521 .flatten()
522 }
523
524 pub(crate) fn owner_of(
525 &self,
526 interner: &Interner,
527 repo: &RepoDid,
528 ) -> Resolved<Option<OwnerDid>> {
529 if self.coverage.get() == Coverage::Warming {
530 return Resolved::Warming;
531 }
532 let Some(target) = interner.repo(repo) else {
533 return Resolved::Ready(None);
534 };
535 Resolved::Ready(
536 self.records
537 .read_sync(&target, |_, slot| slot.owner)
538 .map(|owner| interner.resolve_owner(owner)),
539 )
540 }
541
542 pub(crate) fn rkey_of(
543 &self,
544 interner: &Interner,
545 repo: &RepoDid,
546 ) -> Resolved<Option<RepoRkey>> {
547 if self.coverage.get() == Coverage::Warming {
548 return Resolved::Warming;
549 }
550 let Some(target) = interner.repo(repo) else {
551 return Resolved::Ready(None);
552 };
553 Resolved::Ready(
554 self.records
555 .read_sync(&target, |_, slot| slot.rkey)
556 .map(|rkey| interner.resolve_rkey(rkey)),
557 )
558 }
559
560 pub(crate) fn hosted_repos(&self, interner: &Interner) -> Vec<RepoDid> {
561 let mut repos = BTreeSet::new();
562 self.records.iter_sync(|repo, _| {
563 repos.insert(interner.resolve_repo(*repo));
564 true
565 });
566 repos.into_iter().collect()
567 }
568
569 pub(crate) fn refresh(
570 &self,
571 interner: &Interner,
572 store: &CobStore,
573 object: CobId,
574 ) -> Result<Vec<RepoDid>, IndexError> {
575 let mut tip = self
576 .tip
577 .lock()
578 .unwrap_or_else(|poisoned| poisoned.into_inner());
579 match *tip {
580 None => {
581 let (registry, seeded) = store
582 .materialize::<RepoRegistryCob>(object)
583 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
584 self.seed(interner, ®istry);
585 *tip = Some(seeded);
586 self.coverage.set(Coverage::Ready);
587 Ok(Vec::new())
588 }
589 Some(prev) => {
590 let delta = store
591 .changes_since::<RepoRegistryCob>(object, Some(prev))
592 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
593 let decoded = decode_delta::<RegistryChange>(&delta.changes)
594 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
595 let displaced = self.apply_delta(interner, decoded);
596 *tip = Some(delta.tip);
597 self.coverage.set(Coverage::Ready);
598 Ok(self.evacuated(interner, displaced))
599 }
600 }
601 }
602
603 fn seed(&self, interner: &Interner, registry: &Registry) {
604 self.aliases.clear_sync();
605 self.names.clear_sync();
606 self.records.clear_sync();
607 registry.records().for_each(|(repo, record)| {
608 let repo = interner.intern_repo(repo);
609 let owner = interner.intern_owner(&record.owner);
610 let name = interner.intern_name(&record.name);
611 self.upsert_record(
612 repo,
613 owner,
614 interner.intern_rkey(&record.rkey),
615 name,
616 record.created_at,
617 );
618 self.bind_name(owner, name, record.created_at, repo);
619 });
620 registry.aliases().for_each(|(owner, rkey, repo)| {
621 self.upsert_alias(
622 interner.intern_owner(owner),
623 interner.intern_rkey(rkey),
624 interner.intern_repo(repo),
625 );
626 });
627 }
628
629 fn evacuated(&self, interner: &Interner, displaced: Vec<RepoKey>) -> Vec<RepoDid> {
630 if displaced.is_empty() {
631 return Vec::new();
632 }
633 let live = self.live_repos();
634 displaced
635 .into_iter()
636 .collect::<BTreeSet<_>>()
637 .into_iter()
638 .filter(|repo| !live.contains(repo))
639 .map(|repo| interner.resolve_repo(repo))
640 .collect()
641 }
642
643 fn live_repos(&self) -> BTreeSet<RepoKey> {
644 let mut live = BTreeSet::new();
645 self.records.iter_sync(|repo, _| {
646 live.insert(*repo);
647 true
648 });
649 live
650 }
651
652 fn apply_delta(&self, interner: &Interner, changes: Vec<RegistryChange>) -> Vec<RepoKey> {
653 changes
654 .into_iter()
655 .fold(Vec::new(), |displaced, change| match change {
656 RegistryChange::Register(registration) => {
657 self.apply_register(interner, registration, displaced)
658 }
659 RegistryChange::Rename(rename) => self.apply_rename(interner, rename, displaced),
660 RegistryChange::Deregister(target) => {
661 self.apply_deregister(interner, target, displaced)
662 }
663 })
664 }
665
666 fn apply_register(
667 &self,
668 interner: &Interner,
669 registration: Registration,
670 mut displaced: Vec<RepoKey>,
671 ) -> Vec<RepoKey> {
672 let repo = interner.intern_repo(®istration.repo);
673 let owner = interner.intern_owner(®istration.owner);
674 let rkey = interner.intern_rkey(®istration.rkey);
675 let name = interner.intern_name(®istration.name);
676 if self.records.contains_sync(&repo) {
677 self.drop_record(repo);
678 displaced.push(repo);
679 }
680 displaced.extend(self.steal_alias(owner, rkey, repo));
681 self.upsert_record(repo, owner, rkey, name, registration.created_at);
682 self.upsert_alias(owner, rkey, repo);
683 self.bind_name(owner, name, registration.created_at, repo);
684 displaced
685 }
686
687 fn apply_rename(
688 &self,
689 interner: &Interner,
690 rename: Rename,
691 mut displaced: Vec<RepoKey>,
692 ) -> Vec<RepoKey> {
693 let repo = interner.intern_repo(&rename.repo);
694 let owner = interner.intern_owner(&rename.owner);
695 let rkey = interner.intern_rkey(&rename.rkey);
696 let name = interner.intern_name(&rename.name);
697 let held = self
698 .records
699 .read_sync(&repo, |_, slot| {
700 (slot.owner == owner).then_some((slot.created_at, slot.name))
701 })
702 .flatten();
703 let Some((created_at, previous)) = held else {
704 return displaced;
705 };
706 displaced.extend(self.steal_alias(owner, rkey, repo));
707 self.upsert_record(repo, owner, rkey, name, created_at);
708 self.upsert_alias(owner, rkey, repo);
709 self.bind_name(owner, name, created_at, repo);
710 if previous != name {
711 self.unbind_name(owner, previous, created_at, repo);
712 }
713 displaced
714 }
715
716 fn apply_deregister(
717 &self,
718 interner: &Interner,
719 target: RepoRef,
720 mut displaced: Vec<RepoKey>,
721 ) -> Vec<RepoKey> {
722 let Some(owner) = interner.owner(&target.owner) else {
723 return displaced;
724 };
725 let Some(rkey) = interner.rkey(&target.rkey) else {
726 return displaced;
727 };
728 let Some(repo) = self.aliases.read_sync(&(owner, rkey), |_, repo| *repo) else {
729 return displaced;
730 };
731 self.drop_record(repo);
732 displaced.push(repo);
733 displaced
734 }
735
736 fn steal_alias(&self, owner: OwnerKey, rkey: RkeyKey, target: RepoKey) -> Option<RepoKey> {
737 let holder = self.aliases.read_sync(&(owner, rkey), |_, repo| *repo)?;
738 if holder == target {
739 return None;
740 }
741 let canonical = self
742 .records
743 .read_sync(&holder, |_, slot| slot.rkey == rkey)
744 .unwrap_or(false);
745 if canonical {
746 self.drop_record(holder);
747 Some(holder)
748 } else {
749 let _ = self.aliases.remove_sync(&(owner, rkey));
750 None
751 }
752 }
753
754 fn drop_record(&self, repo: RepoKey) {
755 if let Some((_, slot)) = self.records.remove_sync(&repo) {
756 self.unbind_name(slot.owner, slot.name, slot.created_at, repo);
757 }
758 self.aliases.retain_sync(|_, holder| *holder != repo);
759 }
760
761 fn bind_name(&self, owner: OwnerKey, name: NameKey, created_at: UnixSeconds, repo: RepoKey) {
762 match self.names.entry_sync((owner, name)) {
763 scc::hash_map::Entry::Occupied(mut occupied) => {
764 occupied.get_mut().insert((created_at, repo));
765 }
766 scc::hash_map::Entry::Vacant(vacant) => {
767 vacant.insert_entry(BTreeSet::from([(created_at, repo)]));
768 }
769 }
770 }
771
772 fn unbind_name(&self, owner: OwnerKey, name: NameKey, created_at: UnixSeconds, repo: RepoKey) {
773 let _ = self.names.remove_if_sync(&(owner, name), |registered| {
774 registered.remove(&(created_at, repo));
775 registered.is_empty()
776 });
777 }
778
779 fn upsert_record(
780 &self,
781 repo: RepoKey,
782 owner: OwnerKey,
783 rkey: RkeyKey,
784 name: NameKey,
785 created_at: UnixSeconds,
786 ) {
787 if self
788 .records
789 .update_sync(&repo, |_, slot| {
790 slot.owner = owner;
791 slot.rkey = rkey;
792 slot.name = name;
793 slot.created_at = created_at;
794 })
795 .is_none()
796 {
797 let _ = self.records.insert_sync(
798 repo,
799 RecordSlot {
800 owner,
801 rkey,
802 name,
803 created_at,
804 },
805 );
806 }
807 }
808
809 fn upsert_alias(&self, owner: OwnerKey, rkey: RkeyKey, repo: RepoKey) {
810 let key = (owner, rkey);
811 if self
812 .aliases
813 .update_sync(&key, |_, slot| *slot = repo)
814 .is_none()
815 {
816 let _ = self.aliases.insert_sync(key, repo);
817 }
818 }
819}
820
821pub(crate) struct KeyProjection {
822 cache: Lru<OfferedKey, AccountKey>,
823}
824
825impl KeyProjection {
826 pub(crate) fn new() -> Self {
827 Self {
828 cache: Lru::by_count(EntryCount::new(KEY_CACHE_CAPACITY as u64)),
829 }
830 }
831
832 pub(crate) fn coverage(&self) -> Coverage {
833 Coverage::Ready
834 }
835
836 pub(crate) fn cache(&self, interner: &Interner, key: OfferedKey, did: &AccountDid) {
837 self.cache.insert(key, interner.intern_account(did));
838 }
839
840 pub(crate) fn owner(
841 &self,
842 interner: &Interner,
843 key: &OfferedKey,
844 ) -> Resolved<Option<AccountDid>> {
845 Resolved::Ready(
846 self.cache
847 .get(key)
848 .map(|account| interner.resolve_account(account)),
849 )
850 }
851}