This repository has no description
40 kB
1255 lines
1use std::collections::{BTreeMap, BTreeSet, HashSet};
2use std::hash::{Hash, Hasher};
3use std::marker::PhantomData;
4use std::sync::atomic::{AtomicBool, AtomicI64, AtomicU64, AtomicUsize, Ordering};
5use std::sync::{Arc, Mutex};
6
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};
17use crate::{
18 IndexGeneration, KeyBudget, KeyLease, KeyRecord, KeyReprieve, KeyReprieved, SweepFloor,
19};
20
21#[derive(Debug, Clone, Copy)]
22struct Provenance {
23 added_by: AccountKey,
24 created_at: UnixSeconds,
25}
26
27impl Provenance {
28 fn intern(interner: &Interner, grant: &Grant) -> Self {
29 Self {
30 added_by: interner.intern_account(&grant.added_by),
31 created_at: grant.created_at,
32 }
33 }
34
35 fn grant(self, interner: &Interner, subject: AccountKey) -> Grant {
36 Grant {
37 subject: interner.resolve_account(subject),
38 added_by: interner.resolve_account(self.added_by),
39 created_at: self.created_at,
40 }
41 }
42}
43
44fn decode_change<P: ChangePayload>(change: &Change) -> Result<P, IndexError> {
45 if change.type_name != P::type_name() {
46 return Err(IndexError::UnexpectedType {
47 change: change.id,
48 expected: P::type_name(),
49 found: change.type_name.clone(),
50 });
51 }
52 P::decode(change.payload()).map_err(|error| IndexError::Decode {
53 change: change.id,
54 type_name: P::type_name(),
55 reason: error.to_string(),
56 })
57}
58
59fn decode_delta<P: ChangePayload>(changes: &[Change]) -> Result<Vec<P>, IndexError> {
60 changes.iter().map(decode_change::<P>).collect()
61}
62
63fn group_by<K: Ord, T>(items: Vec<T>, key: impl Fn(&T) -> K) -> BTreeMap<K, Vec<T>> {
64 items.into_iter().fold(BTreeMap::new(), |mut acc, item| {
65 acc.entry(key(&item)).or_default().push(item);
66 acc
67 })
68}
69
70pub(crate) struct GrantSetProjection<Cob>
71where
72 Cob: Evaluate,
73 Cob::Change: ChangePayload + GrantChange,
74{
75 membership: scc::HashMap<AccountKey, Provenance>,
76 coverage: CoverageCell,
77 tip: Mutex<Option<ChangeId>>,
78 _cob: PhantomData<fn() -> Cob>,
79}
80
81impl<Cob> GrantSetProjection<Cob>
82where
83 Cob: Evaluate,
84 Cob::Change: ChangePayload + GrantChange,
85{
86 pub(crate) fn new() -> Self {
87 Self {
88 membership: scc::HashMap::new(),
89 coverage: CoverageCell::new(Coverage::Warming),
90 tip: Mutex::new(None),
91 _cob: PhantomData,
92 }
93 }
94
95 pub(crate) fn coverage(&self) -> Coverage {
96 self.coverage.get()
97 }
98
99 pub(crate) fn reset(&self) {
100 let mut tip = self
101 .tip
102 .lock()
103 .unwrap_or_else(|poisoned| poisoned.into_inner());
104 self.membership.clear_sync();
105 *tip = None;
106 self.coverage.set(Coverage::Ready);
107 }
108
109 pub(crate) fn contains(&self, interner: &Interner, did: &AccountDid) -> Resolved<bool> {
110 match self.coverage.get() {
111 Coverage::Warming => Resolved::Warming,
112 Coverage::Ready => Resolved::Ready(
113 interner
114 .account(did)
115 .is_some_and(|did| self.membership.contains_sync(&did)),
116 ),
117 }
118 }
119
120 pub(crate) fn entries(&self, interner: &Interner) -> Resolved<Vec<Grant>> {
121 match self.coverage.get() {
122 Coverage::Warming => Resolved::Warming,
123 Coverage::Ready => {
124 let mut out = BTreeMap::new();
125 self.membership.iter_sync(|&subject, slot| {
126 let grant = slot.grant(interner, subject);
127 out.insert(grant.subject.clone(), grant);
128 true
129 });
130 Resolved::Ready(out.into_values().collect())
131 }
132 }
133 }
134
135 pub(crate) fn refresh(
136 &self,
137 interner: &Interner,
138 store: &CobStore,
139 object: CobId,
140 ) -> Result<(), IndexError>
141 where
142 Cob: Checkpoint + Evaluate<State = Roster>,
143 {
144 let mut tip = self
145 .tip
146 .lock()
147 .unwrap_or_else(|poisoned| poisoned.into_inner());
148 match *tip {
149 None => {
150 let (roster, seeded) = store
151 .materialize::<Cob>(object)
152 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
153 self.seed(interner, &roster);
154 *tip = Some(seeded);
155 }
156 Some(prev) => {
157 let delta = store
158 .changes_since::<Cob>(object, Some(prev))
159 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
160 let decoded = decode_delta::<Cob::Change>(&delta.changes)
161 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
162 self.apply_delta(interner, decoded);
163 *tip = Some(delta.tip);
164 }
165 }
166 self.coverage.set(Coverage::Ready);
167 Ok(())
168 }
169
170 fn seed(&self, interner: &Interner, roster: &Roster) {
171 self.membership.clear_sync();
172 roster.entries().for_each(|(subject, entry)| {
173 let slot = Provenance {
174 added_by: interner.intern_account(&entry.added_by),
175 created_at: entry.created_at,
176 };
177 let _ = self
178 .membership
179 .insert_sync(interner.intern_account(subject), slot);
180 });
181 }
182
183 fn apply_delta(&self, interner: &Interner, changes: Vec<Cob::Change>) {
184 group_by(changes, |change| change.subject().clone())
185 .into_iter()
186 .for_each(|(did, ops)| {
187 let current = interner
188 .account(&did)
189 .and_then(|key| self.membership.read_sync(&key, |_, slot| *slot));
190 let net = ops
191 .into_iter()
192 .fold(current, |slot, change| match change.as_grant() {
193 Some(grant) => slot.or_else(|| Some(Provenance::intern(interner, grant))),
194 None => None,
195 });
196 match net {
197 Some(slot) => {
198 *self
199 .membership
200 .entry_sync(interner.intern_account(&did))
201 .or_insert(slot)
202 .get_mut() = slot;
203 }
204 None => {
205 if let Some(key) = interner.account(&did) {
206 let _ = self.membership.remove_sync(&key);
207 }
208 }
209 }
210 });
211 }
212}
213
214struct RepoRoster {
215 tip: Option<ChangeId>,
216 entries: BTreeMap<AccountKey, Provenance>,
217}
218
219const REPO_LOCK_STRIPES: usize = 256;
220
221pub(crate) struct CollaboratorsProjection {
222 rosters: scc::HashMap<RepoKey, RepoRoster>,
223 locks: [Mutex<()>; REPO_LOCK_STRIPES],
224}
225
226impl CollaboratorsProjection {
227 pub(crate) fn new() -> Self {
228 Self {
229 rosters: scc::HashMap::new(),
230 locks: std::array::from_fn(|_| Mutex::new(())),
231 }
232 }
233
234 fn repo_lock(&self, repo: RepoKey) -> &Mutex<()> {
235 let mut hasher = std::collections::hash_map::DefaultHasher::new();
236 repo.hash(&mut hasher);
237 &self.locks[(hasher.finish() % REPO_LOCK_STRIPES as u64) as usize]
238 }
239
240 pub(crate) fn coverage(&self) -> Coverage {
241 Coverage::Ready
242 }
243
244 pub(crate) fn is_folded(&self, interner: &Interner, repo: &RepoDid) -> bool {
245 interner
246 .repo(repo)
247 .is_some_and(|repo| self.rosters.contains_sync(&repo))
248 }
249
250 pub(crate) fn contains(
251 &self,
252 interner: &Interner,
253 repo: &RepoDid,
254 did: &AccountDid,
255 ) -> Resolved<bool> {
256 let Some(repo) = interner.repo(repo) else {
257 return Resolved::Warming;
258 };
259 match self.rosters.read_sync(&repo, |_, roster| {
260 interner
261 .account(did)
262 .is_some_and(|account| roster.entries.contains_key(&account))
263 }) {
264 Some(present) => Resolved::Ready(present),
265 None => Resolved::Warming,
266 }
267 }
268
269 pub(crate) fn entries(&self, interner: &Interner, repo: &RepoDid) -> Resolved<Vec<Grant>> {
270 let Some(repo) = interner.repo(repo) else {
271 return Resolved::Warming;
272 };
273 match self
274 .rosters
275 .read_sync(&repo, |_, roster| roster_grants(interner, roster))
276 {
277 Some(grants) => Resolved::Ready(grants),
278 None => Resolved::Warming,
279 }
280 }
281
282 pub(crate) fn mark_repo_empty(&self, repo: RepoKey) {
283 let lock = self.repo_lock(repo);
284 let _guard = lock.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
285 self.install(
286 repo,
287 RepoRoster {
288 tip: None,
289 entries: BTreeMap::new(),
290 },
291 );
292 }
293
294 pub(crate) fn refresh_repo(
295 &self,
296 interner: &Interner,
297 store: &CobStore,
298 repo: RepoKey,
299 object: CobId,
300 ) -> Result<(), IndexError> {
301 let lock = self.repo_lock(repo);
302 let _guard = lock.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
303 self.fold_repo(interner, store, repo, object)
304 .inspect_err(|_| self.purge_repo(repo))
305 }
306
307 fn fold_repo(
308 &self,
309 interner: &Interner,
310 store: &CobStore,
311 repo: RepoKey,
312 object: CobId,
313 ) -> Result<(), IndexError> {
314 let prev = self
315 .rosters
316 .read_sync(&repo, |_, roster| roster.tip)
317 .flatten();
318 match prev {
319 None => {
320 let (roster, tip) = store.materialize::<CollaboratorsCob>(object)?;
321 let entries = roster
322 .entries()
323 .map(|(subject, entry)| {
324 (
325 interner.intern_account(subject),
326 Provenance {
327 added_by: interner.intern_account(&entry.added_by),
328 created_at: entry.created_at,
329 },
330 )
331 })
332 .collect();
333 self.install(
334 repo,
335 RepoRoster {
336 tip: Some(tip),
337 entries,
338 },
339 );
340 }
341 Some(prev) => {
342 let delta = store.changes_since::<CollaboratorsCob>(object, Some(prev))?;
343 let decoded = decode_delta::<CollaboratorsChange>(&delta.changes)?;
344 let mut occupied = self.rosters.entry_sync(repo).or_insert_with(|| RepoRoster {
345 tip: None,
346 entries: BTreeMap::new(),
347 });
348 let roster = occupied.get_mut();
349 roster.entries = decoded
350 .into_iter()
351 .fold(std::mem::take(&mut roster.entries), |entries, change| {
352 apply(interner, entries, change)
353 });
354 roster.tip = Some(delta.tip);
355 }
356 }
357 Ok(())
358 }
359
360 pub(crate) fn drop_repo(&self, repo: RepoKey) {
361 let lock = self.repo_lock(repo);
362 let _guard = lock.lock().unwrap_or_else(|poisoned| poisoned.into_inner());
363 self.purge_repo(repo);
364 }
365
366 fn purge_repo(&self, repo: RepoKey) {
367 let _ = self.rosters.remove_sync(&repo);
368 }
369
370 fn install(&self, repo: RepoKey, roster: RepoRoster) {
371 match self.rosters.entry_sync(repo) {
372 scc::hash_map::Entry::Occupied(mut occupied) => {
373 let _ = occupied.insert(roster);
374 }
375 scc::hash_map::Entry::Vacant(vacant) => {
376 vacant.insert_entry(roster);
377 }
378 }
379 }
380}
381
382fn roster_grants(interner: &Interner, roster: &RepoRoster) -> Vec<Grant> {
383 roster
384 .entries
385 .iter()
386 .map(|(&subject, slot)| {
387 let grant = slot.grant(interner, subject);
388 (grant.subject.clone(), grant)
389 })
390 .collect::<BTreeMap<_, _>>()
391 .into_values()
392 .collect()
393}
394
395fn apply(
396 interner: &Interner,
397 mut entries: BTreeMap<AccountKey, Provenance>,
398 change: CollaboratorsChange,
399) -> BTreeMap<AccountKey, Provenance> {
400 match change {
401 CollaboratorsChange::Add(grant) => {
402 entries
403 .entry(interner.intern_account(&grant.subject))
404 .or_insert_with(|| Provenance::intern(interner, &grant));
405 }
406 CollaboratorsChange::Remove(removal) => {
407 if let Some(key) = interner.account(&removal.subject) {
408 entries.remove(&key);
409 }
410 }
411 }
412 entries
413}
414
415struct RecordSlot {
416 owner: OwnerKey,
417 rkey: RkeyKey,
418 name: NameKey,
419 created_at: UnixSeconds,
420}
421
422pub(crate) struct RegistryProjection {
423 aliases: scc::HashMap<(OwnerKey, RkeyKey), RepoKey>,
424 names: scc::HashMap<(OwnerKey, NameKey), BTreeSet<(UnixSeconds, RepoKey)>>,
425 records: scc::HashMap<RepoKey, RecordSlot>,
426 coverage: CoverageCell,
427 tip: Mutex<Option<ChangeId>>,
428}
429
430impl RegistryProjection {
431 pub(crate) fn new() -> Self {
432 Self {
433 aliases: scc::HashMap::new(),
434 names: scc::HashMap::new(),
435 records: scc::HashMap::new(),
436 coverage: CoverageCell::new(Coverage::Warming),
437 tip: Mutex::new(None),
438 }
439 }
440
441 pub(crate) fn coverage(&self) -> Coverage {
442 self.coverage.get()
443 }
444
445 pub(crate) fn reset(&self, interner: &Interner) -> Vec<RepoDid> {
446 let mut tip = self
447 .tip
448 .lock()
449 .unwrap_or_else(|poisoned| poisoned.into_inner());
450 let evacuated = self.hosted_repos(interner);
451 self.aliases.clear_sync();
452 self.names.clear_sync();
453 self.records.clear_sync();
454 *tip = None;
455 self.coverage.set(Coverage::Ready);
456 evacuated
457 }
458
459 pub(crate) fn resolve(
460 &self,
461 interner: &Interner,
462 owner: &OwnerDid,
463 rkey: &RepoRkey,
464 ) -> Resolved<Option<RepoDid>> {
465 match self.coverage.get() {
466 Coverage::Warming => Resolved::Warming,
467 Coverage::Ready => Resolved::Ready(
468 interner
469 .owner(owner)
470 .zip(interner.rkey(rkey))
471 .and_then(|key| self.aliases.read_sync(&key, |_, repo| *repo))
472 .map(|repo| interner.resolve_repo(repo)),
473 ),
474 }
475 }
476
477 pub(crate) fn resolve_clone_path(
478 &self,
479 interner: &Interner,
480 owner: &OwnerDid,
481 path: &ClonePath,
482 ) -> Resolved<Option<RepoDid>> {
483 if self.coverage.get() == Coverage::Warming {
484 return Resolved::Warming;
485 }
486 let Some(owner) = interner.owner(owner) else {
487 return Resolved::Ready(None);
488 };
489 let by_rkey = path
490 .rkeys()
491 .filter_map(|rkey| interner.rkey(rkey))
492 .find_map(|rkey| self.aliases.read_sync(&(owner, rkey), |_, repo| *repo));
493 if let Some(repo) = by_rkey {
494 return Resolved::Ready(Some(interner.resolve_repo(repo)));
495 }
496 Resolved::Ready(
497 path.names()
498 .filter_map(|name| interner.name(name))
499 .find_map(|name| self.oldest_registration_for_name(interner, owner, name)),
500 )
501 }
502
503 fn oldest_registration_for_name(
504 &self,
505 interner: &Interner,
506 owner: OwnerKey,
507 name: NameKey,
508 ) -> Option<RepoDid> {
509 // Notice how set keys sort by whichever string the interner saw first,
510 // and a cold rebuild will see them in a different order than live replay,
511 // so 2 entries with the same timestamp will settle by comparing
512 // DIDs instead.
513 self.names
514 .read_sync(&(owner, name), |_, registered| {
515 let earliest = registered.first()?.0;
516 registered
517 .iter()
518 .take_while(|(created_at, _)| *created_at == earliest)
519 .map(|(_, repo)| interner.resolve_repo(*repo))
520 .min()
521 })
522 .flatten()
523 }
524
525 pub(crate) fn owner_of(
526 &self,
527 interner: &Interner,
528 repo: &RepoDid,
529 ) -> Resolved<Option<OwnerDid>> {
530 if self.coverage.get() == Coverage::Warming {
531 return Resolved::Warming;
532 }
533 let Some(target) = interner.repo(repo) else {
534 return Resolved::Ready(None);
535 };
536 Resolved::Ready(
537 self.records
538 .read_sync(&target, |_, slot| slot.owner)
539 .map(|owner| interner.resolve_owner(owner)),
540 )
541 }
542
543 pub(crate) fn rkey_of(
544 &self,
545 interner: &Interner,
546 repo: &RepoDid,
547 ) -> Resolved<Option<RepoRkey>> {
548 if self.coverage.get() == Coverage::Warming {
549 return Resolved::Warming;
550 }
551 let Some(target) = interner.repo(repo) else {
552 return Resolved::Ready(None);
553 };
554 Resolved::Ready(
555 self.records
556 .read_sync(&target, |_, slot| slot.rkey)
557 .map(|rkey| interner.resolve_rkey(rkey)),
558 )
559 }
560
561 pub(crate) fn hosted_repos(&self, interner: &Interner) -> Vec<RepoDid> {
562 let mut repos = BTreeSet::new();
563 self.records.iter_sync(|repo, _| {
564 repos.insert(interner.resolve_repo(*repo));
565 true
566 });
567 repos.into_iter().collect()
568 }
569
570 pub(crate) fn refresh(
571 &self,
572 interner: &Interner,
573 store: &CobStore,
574 object: CobId,
575 ) -> Result<Vec<RepoDid>, IndexError> {
576 let mut tip = self
577 .tip
578 .lock()
579 .unwrap_or_else(|poisoned| poisoned.into_inner());
580 match *tip {
581 None => {
582 let (registry, seeded) = store
583 .materialize::<RepoRegistryCob>(object)
584 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
585 self.seed(interner, ®istry);
586 *tip = Some(seeded);
587 self.coverage.set(Coverage::Ready);
588 Ok(Vec::new())
589 }
590 Some(prev) => {
591 let delta = store
592 .changes_since::<RepoRegistryCob>(object, Some(prev))
593 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
594 let decoded = decode_delta::<RegistryChange>(&delta.changes)
595 .inspect_err(|_| self.coverage.set(Coverage::Warming))?;
596 let displaced = self.apply_delta(interner, decoded);
597 *tip = Some(delta.tip);
598 self.coverage.set(Coverage::Ready);
599 Ok(self.evacuated(interner, displaced))
600 }
601 }
602 }
603
604 fn seed(&self, interner: &Interner, registry: &Registry) {
605 self.aliases.clear_sync();
606 self.names.clear_sync();
607 self.records.clear_sync();
608 registry.records().for_each(|(repo, record)| {
609 let repo = interner.intern_repo(repo);
610 let owner = interner.intern_owner(&record.owner);
611 let name = interner.intern_name(&record.name);
612 self.upsert_record(
613 repo,
614 owner,
615 interner.intern_rkey(&record.rkey),
616 name,
617 record.created_at,
618 );
619 self.bind_name(owner, name, record.created_at, repo);
620 });
621 registry.aliases().for_each(|(owner, rkey, repo)| {
622 self.upsert_alias(
623 interner.intern_owner(owner),
624 interner.intern_rkey(rkey),
625 interner.intern_repo(repo),
626 );
627 });
628 }
629
630 fn evacuated(&self, interner: &Interner, displaced: Vec<RepoKey>) -> Vec<RepoDid> {
631 if displaced.is_empty() {
632 return Vec::new();
633 }
634 let live = self.live_repos();
635 displaced
636 .into_iter()
637 .collect::<BTreeSet<_>>()
638 .into_iter()
639 .filter(|repo| !live.contains(repo))
640 .map(|repo| interner.resolve_repo(repo))
641 .collect()
642 }
643
644 fn live_repos(&self) -> BTreeSet<RepoKey> {
645 let mut live = BTreeSet::new();
646 self.records.iter_sync(|repo, _| {
647 live.insert(*repo);
648 true
649 });
650 live
651 }
652
653 fn apply_delta(&self, interner: &Interner, changes: Vec<RegistryChange>) -> Vec<RepoKey> {
654 changes
655 .into_iter()
656 .fold(Vec::new(), |displaced, change| match change {
657 RegistryChange::Register(registration) => {
658 self.apply_register(interner, registration, displaced)
659 }
660 RegistryChange::Rename(rename) => self.apply_rename(interner, rename, displaced),
661 RegistryChange::Deregister(target) => {
662 self.apply_deregister(interner, target, displaced)
663 }
664 })
665 }
666
667 fn apply_register(
668 &self,
669 interner: &Interner,
670 registration: Registration,
671 mut displaced: Vec<RepoKey>,
672 ) -> Vec<RepoKey> {
673 let repo = interner.intern_repo(®istration.repo);
674 let owner = interner.intern_owner(®istration.owner);
675 let rkey = interner.intern_rkey(®istration.rkey);
676 let name = interner.intern_name(®istration.name);
677 if self.records.contains_sync(&repo) {
678 self.drop_record(repo);
679 displaced.push(repo);
680 }
681 displaced.extend(self.steal_alias(owner, rkey, repo));
682 self.upsert_record(repo, owner, rkey, name, registration.created_at);
683 self.upsert_alias(owner, rkey, repo);
684 self.bind_name(owner, name, registration.created_at, repo);
685 displaced
686 }
687
688 fn apply_rename(
689 &self,
690 interner: &Interner,
691 rename: Rename,
692 mut displaced: Vec<RepoKey>,
693 ) -> Vec<RepoKey> {
694 let repo = interner.intern_repo(&rename.repo);
695 let owner = interner.intern_owner(&rename.owner);
696 let rkey = interner.intern_rkey(&rename.rkey);
697 let name = interner.intern_name(&rename.name);
698 let held = self
699 .records
700 .read_sync(&repo, |_, slot| {
701 (slot.owner == owner).then_some((slot.created_at, slot.name))
702 })
703 .flatten();
704 let Some((created_at, previous)) = held else {
705 return displaced;
706 };
707 displaced.extend(self.steal_alias(owner, rkey, repo));
708 self.upsert_record(repo, owner, rkey, name, created_at);
709 self.upsert_alias(owner, rkey, repo);
710 self.bind_name(owner, name, created_at, repo);
711 if previous != name {
712 self.unbind_name(owner, previous, created_at, repo);
713 }
714 displaced
715 }
716
717 fn apply_deregister(
718 &self,
719 interner: &Interner,
720 target: RepoRef,
721 mut displaced: Vec<RepoKey>,
722 ) -> Vec<RepoKey> {
723 let Some(owner) = interner.owner(&target.owner) else {
724 return displaced;
725 };
726 let Some(rkey) = interner.rkey(&target.rkey) else {
727 return displaced;
728 };
729 let Some(repo) = self.aliases.read_sync(&(owner, rkey), |_, repo| *repo) else {
730 return displaced;
731 };
732 self.drop_record(repo);
733 displaced.push(repo);
734 displaced
735 }
736
737 fn steal_alias(&self, owner: OwnerKey, rkey: RkeyKey, target: RepoKey) -> Option<RepoKey> {
738 let holder = self.aliases.read_sync(&(owner, rkey), |_, repo| *repo)?;
739 if holder == target {
740 return None;
741 }
742 let canonical = self
743 .records
744 .read_sync(&holder, |_, slot| slot.rkey == rkey)
745 .unwrap_or(false);
746 if canonical {
747 self.drop_record(holder);
748 Some(holder)
749 } else {
750 let _ = self.aliases.remove_sync(&(owner, rkey));
751 None
752 }
753 }
754
755 fn drop_record(&self, repo: RepoKey) {
756 if let Some((_, slot)) = self.records.remove_sync(&repo) {
757 self.unbind_name(slot.owner, slot.name, slot.created_at, repo);
758 }
759 self.aliases.retain_sync(|_, holder| *holder != repo);
760 }
761
762 fn bind_name(&self, owner: OwnerKey, name: NameKey, created_at: UnixSeconds, repo: RepoKey) {
763 match self.names.entry_sync((owner, name)) {
764 scc::hash_map::Entry::Occupied(mut occupied) => {
765 occupied.get_mut().insert((created_at, repo));
766 }
767 scc::hash_map::Entry::Vacant(vacant) => {
768 vacant.insert_entry(BTreeSet::from([(created_at, repo)]));
769 }
770 }
771 }
772
773 fn unbind_name(&self, owner: OwnerKey, name: NameKey, created_at: UnixSeconds, repo: RepoKey) {
774 let _ = self.names.remove_if_sync(&(owner, name), |registered| {
775 registered.remove(&(created_at, repo));
776 registered.is_empty()
777 });
778 }
779
780 fn upsert_record(
781 &self,
782 repo: RepoKey,
783 owner: OwnerKey,
784 rkey: RkeyKey,
785 name: NameKey,
786 created_at: UnixSeconds,
787 ) {
788 if self
789 .records
790 .update_sync(&repo, |_, slot| {
791 slot.owner = owner;
792 slot.rkey = rkey;
793 slot.name = name;
794 slot.created_at = created_at;
795 })
796 .is_none()
797 {
798 let _ = self.records.insert_sync(
799 repo,
800 RecordSlot {
801 owner,
802 rkey,
803 name,
804 created_at,
805 },
806 );
807 }
808 }
809
810 fn upsert_alias(&self, owner: OwnerKey, rkey: RkeyKey, repo: RepoKey) {
811 let key = (owner, rkey);
812 if self
813 .aliases
814 .update_sync(&key, |_, slot| *slot = repo)
815 .is_none()
816 {
817 let _ = self.aliases.insert_sync(key, repo);
818 }
819 }
820}
821
822enum Reading {
823 Published(Vec<Arc<OfferedKey>>),
824 Unheld,
825 Unread,
826}
827
828impl Reading {
829 fn keys(&self) -> &[Arc<OfferedKey>] {
830 match self {
831 Reading::Published(keys) => keys,
832 Reading::Unheld | Reading::Unread => &[],
833 }
834 }
835
836 fn is_published(&self) -> bool {
837 matches!(self, Reading::Published(_))
838 }
839}
840
841struct Held {
842 reading: Reading,
843 lease: KeyLease,
844}
845
846const PER_KEY_OVERHEAD: usize = 192;
847
848const PER_ACCOUNT_OVERHEAD: usize = 128;
849
850impl Held {
851 fn bytes(&self) -> usize {
852 PER_ACCOUNT_OVERHEAD
853 + self
854 .reading
855 .keys()
856 .iter()
857 .map(|key| key.as_bytes().len() + PER_KEY_OVERHEAD)
858 .sum::<usize>()
859 }
860
861 fn answers(&self, now: UnixSeconds) -> bool {
862 self.reading.is_published() && self.lease.is_live(now)
863 }
864
865 fn settled(&self, now: UnixSeconds) -> bool {
866 !matches!(self.reading, Reading::Unread) && self.lease.is_live(now)
867 }
868
869 fn renewal_due(&self, now: UnixSeconds) -> bool {
870 match self.reading {
871 Reading::Published(_) => self.lease.renewal_due(now),
872 Reading::Unheld | Reading::Unread => !self.lease.is_live(now),
873 }
874 }
875
876 fn publishes(&self, key: &OfferedKey, now: UnixSeconds) -> bool {
877 self.answers(now)
878 && self
879 .reading
880 .keys()
881 .iter()
882 .any(|published| published.as_ref() == key)
883 }
884
885 fn is_unheld(&self) -> bool {
886 matches!(self.reading, Reading::Unheld)
887 }
888}
889
890#[derive(Default)]
891struct Publishers(Vec<AccountKey>);
892
893impl Publishers {
894 fn add(&mut self, account: AccountKey) {
895 if let Err(at) = self.0.binary_search(&account) {
896 self.0.insert(at, account);
897 }
898 }
899
900 fn remove(&mut self, account: AccountKey) -> bool {
901 if let Ok(at) = self.0.binary_search(&account) {
902 self.0.remove(at);
903 }
904 self.0.is_empty()
905 }
906
907 fn to_vec(&self) -> Vec<AccountKey> {
908 self.0.clone()
909 }
910}
911
912const NEVER: u64 = u64::MAX;
913
914pub(crate) struct KeyProjection {
915 owners: scc::HashMap<Arc<OfferedKey>, Publishers>,
916 held: scc::HashMap<AccountKey, Held>,
917 budget: KeyBudget,
918 tracked_bytes: AtomicUsize,
919 unheld: AtomicUsize,
920 suspect: AtomicBool,
921 swept_at: AtomicI64,
922 ready_at: AtomicU64,
923}
924
925impl KeyProjection {
926 pub(crate) fn new(budget: KeyBudget) -> Self {
927 Self {
928 owners: scc::HashMap::new(),
929 held: scc::HashMap::new(),
930 budget,
931 tracked_bytes: AtomicUsize::new(0),
932 unheld: AtomicUsize::new(0),
933 suspect: AtomicBool::new(false),
934 swept_at: AtomicI64::new(i64::MIN),
935 ready_at: AtomicU64::new(NEVER),
936 }
937 }
938
939 pub(crate) fn coverage(&self, generation: IndexGeneration) -> Coverage {
940 match self.ready_at.load(Ordering::Acquire) == generation.get() {
941 true => Coverage::Ready,
942 false => Coverage::Warming,
943 }
944 }
945
946 pub(crate) fn suspect_now(&self) {
947 self.suspect.store(true, Ordering::Release);
948 }
949
950 pub(crate) fn take_suspicion(&self, now: UnixSeconds, floor: SweepFloor) -> bool {
951 let held_until = |swept: i64| swept.saturating_add(crate::whole_secs(floor.get()));
952 match self.suspect.load(Ordering::Acquire) {
953 false => false,
954 true => {
955 self.swept_at
956 .fetch_update(Ordering::AcqRel, Ordering::Acquire, |swept| {
957 (now.get() >= held_until(swept)).then_some(now.get())
958 })
959 .is_ok()
960 && self.suspect.swap(false, Ordering::AcqRel)
961 }
962 }
963 }
964
965 pub(crate) fn mark_ready(&self, generation: IndexGeneration) {
966 self.ready_at.store(generation.get(), Ordering::Release);
967 }
968
969 pub(crate) fn mark_warming(&self) {
970 self.ready_at.store(NEVER, Ordering::Release);
971 }
972
973 pub(crate) fn owner(
974 &self,
975 interner: &Interner,
976 key: &OfferedKey,
977 now: UnixSeconds,
978 ) -> Resolved<Option<AccountDid>> {
979 let publishers = self
980 .owners
981 .read_sync(key, |_, publishers| publishers.to_vec())
982 .unwrap_or_default();
983 Resolved::Ready(
984 publishers
985 .into_iter()
986 .find(|account| {
987 self.held_by(*account, |held| held.publishes(key, now)) == Some(true)
988 })
989 .map(|account| interner.resolve_account(account)),
990 )
991 }
992
993 pub(crate) fn any_unheld(&self) -> bool {
994 self.unheld.load(Ordering::Acquire) > 0
995 }
996
997 fn track_unheld(&self, lost: bool, gained: bool) {
998 match (lost, gained) {
999 (false, true) => {
1000 self.unheld.fetch_add(1, Ordering::AcqRel);
1001 }
1002 (true, false) => {
1003 self.unheld.fetch_sub(1, Ordering::AcqRel);
1004 }
1005 (true, true) | (false, false) => {}
1006 }
1007 }
1008
1009 pub(crate) fn record(
1010 &self,
1011 interner: &Interner,
1012 did: &AccountDid,
1013 keys: Vec<OfferedKey>,
1014 lease: KeyLease,
1015 ) -> KeyRecord {
1016 let account = interner.intern_account(did);
1017 let incoming = Held {
1018 reading: Reading::Published(keys.into_iter().map(Arc::new).collect()),
1019 lease,
1020 };
1021 let entry = self.held.entry_sync(account);
1022 let held = match &entry {
1023 scc::hash_map::Entry::Occupied(slot) => slot.get().bytes(),
1024 scc::hash_map::Entry::Vacant(_) => 0,
1025 };
1026 if self.reserve_bytes(incoming.bytes() as isize - held as isize) {
1027 self.settle(entry, account, incoming);
1028 return KeyRecord::Stored;
1029 }
1030 let unheld = Held {
1031 reading: Reading::Unheld,
1032 lease,
1033 };
1034 if !self.reserve_bytes(unheld.bytes() as isize - held as isize) {
1035 return KeyRecord::Saturated;
1036 }
1037 self.settle(entry, account, unheld);
1038 KeyRecord::Unheld
1039 }
1040
1041 fn settle(
1042 &self,
1043 entry: scc::hash_map::Entry<'_, AccountKey, Held>,
1044 account: AccountKey,
1045 incoming: Held,
1046 ) {
1047 let published: HashSet<Arc<OfferedKey>> = incoming.reading.keys().iter().cloned().collect();
1048 let gained = incoming.is_unheld();
1049 match entry {
1050 scc::hash_map::Entry::Occupied(mut slot) => {
1051 let replaced = slot.insert(incoming);
1052 self.track_unheld(replaced.is_unheld(), gained);
1053 self.rewire(account, &published, replaced.reading.keys());
1054 }
1055 scc::hash_map::Entry::Vacant(slot) => {
1056 let _locked = slot.insert_entry(incoming);
1057 self.track_unheld(false, gained);
1058 self.rewire(account, &published, &[]);
1059 }
1060 }
1061 }
1062
1063 fn rewire(
1064 &self,
1065 account: AccountKey,
1066 published: &HashSet<Arc<OfferedKey>>,
1067 replaced: &[Arc<OfferedKey>],
1068 ) {
1069 published
1070 .iter()
1071 .for_each(|key| self.claim(Arc::clone(key), account));
1072 replaced
1073 .iter()
1074 .filter(|key| !published.contains(*key))
1075 .for_each(|key| self.disown(key, account));
1076 }
1077
1078 pub(crate) fn reprieve(
1079 &self,
1080 interner: &Interner,
1081 did: &AccountDid,
1082 now: UnixSeconds,
1083 grace: KeyReprieve,
1084 exhausted: KeyLease,
1085 ) -> KeyReprieved {
1086 let account = interner.intern_account(did);
1087 let mut slot = match self.held.entry_sync(account) {
1088 scc::hash_map::Entry::Vacant(slot) => {
1089 let pending = Held {
1090 reading: Reading::Unread,
1091 lease: grace.first_failure(now),
1092 };
1093 if self.reserve_bytes(pending.bytes() as isize) {
1094 let _locked = slot.insert_entry(pending);
1095 }
1096 return KeyReprieved::Pending;
1097 }
1098 scc::hash_map::Entry::Occupied(slot) => slot,
1099 };
1100 match grace.extend(slot.get().lease, now) {
1101 Some(lease) => {
1102 let carried = slot.get().reading.is_published();
1103 slot.get_mut().lease = lease;
1104 match carried {
1105 true => KeyReprieved::Extended,
1106 false => KeyReprieved::Pending,
1107 }
1108 }
1109 None => {
1110 let given_up = Held {
1111 reading: Reading::Published(Vec::new()),
1112 lease: exhausted,
1113 };
1114 let _ = self.reserve_bytes(given_up.bytes() as isize - slot.get().bytes() as isize);
1115 self.track_unheld(slot.get().is_unheld(), false);
1116 self.disown_all(&slot, account);
1117 *slot.get_mut() = given_up;
1118 KeyReprieved::Exhausted
1119 }
1120 }
1121 }
1122
1123 pub(crate) fn on_file(&self, interner: &Interner, did: &AccountDid) -> bool {
1124 interner
1125 .account(did)
1126 .is_some_and(|account| self.held.contains_sync(&account))
1127 }
1128
1129 fn held_by(&self, account: AccountKey, ready: impl Fn(&Held) -> bool) -> Option<bool> {
1130 self.held.read_sync(&account, |_, held| ready(held))
1131 }
1132
1133 fn held_for(
1134 &self,
1135 interner: &Interner,
1136 did: &AccountDid,
1137 ready: impl Fn(&Held) -> bool,
1138 ) -> Option<bool> {
1139 interner
1140 .account(did)
1141 .and_then(|account| self.held_by(account, ready))
1142 }
1143
1144 pub(crate) fn is_fresh(&self, interner: &Interner, did: &AccountDid, now: UnixSeconds) -> bool {
1145 self.held_for(interner, did, |held| held.answers(now))
1146 .unwrap_or(false)
1147 }
1148
1149 pub(crate) fn renewal_due(
1150 &self,
1151 interner: &Interner,
1152 did: &AccountDid,
1153 now: UnixSeconds,
1154 ) -> bool {
1155 self.held_for(interner, did, |held| held.renewal_due(now))
1156 .unwrap_or(true)
1157 }
1158
1159 pub(crate) fn all_live(
1160 &self,
1161 interner: &Interner,
1162 subjects: &[AccountDid],
1163 now: UnixSeconds,
1164 ) -> bool {
1165 subjects.iter().all(|did| {
1166 self.held_for(interner, did, |held| held.settled(now))
1167 .unwrap_or(false)
1168 })
1169 }
1170
1171 pub(crate) fn publisher_among(
1172 &self,
1173 interner: &Interner,
1174 candidates: &[AccountDid],
1175 key: &OfferedKey,
1176 now: UnixSeconds,
1177 ) -> Option<AccountDid> {
1178 candidates
1179 .iter()
1180 .find(|did| {
1181 self.held_for(interner, did, |held| held.publishes(key, now))
1182 .unwrap_or(false)
1183 })
1184 .cloned()
1185 }
1186
1187 pub(crate) fn retain(&self, interner: &Interner, kept: &[AccountDid]) {
1188 let keep: BTreeSet<AccountKey> = kept
1189 .iter()
1190 .filter_map(|did| interner.account(did))
1191 .collect();
1192 let mut released = Vec::new();
1193 self.held.iter_sync(|account, _| {
1194 if !keep.contains(account) {
1195 released.push(*account);
1196 }
1197 true
1198 });
1199 released
1200 .into_iter()
1201 .for_each(|account| self.release(account));
1202 }
1203
1204 fn release(&self, account: AccountKey) {
1205 if let scc::hash_map::Entry::Occupied(slot) = self.held.entry_sync(account) {
1206 self.evict(&slot, account);
1207 let _ = slot.remove();
1208 }
1209 }
1210
1211 fn evict(
1212 &self,
1213 slot: &scc::hash_map::OccupiedEntry<'_, AccountKey, Held>,
1214 account: AccountKey,
1215 ) {
1216 self.track_unheld(slot.get().is_unheld(), false);
1217 self.disown_all(slot, account);
1218 let _ = self.reserve_bytes(-(slot.get().bytes() as isize));
1219 }
1220
1221 fn disown_all(
1222 &self,
1223 slot: &scc::hash_map::OccupiedEntry<'_, AccountKey, Held>,
1224 account: AccountKey,
1225 ) {
1226 slot.get()
1227 .reading
1228 .keys()
1229 .iter()
1230 .for_each(|key| self.disown(key, account));
1231 }
1232
1233 fn reserve_bytes(&self, growth: isize) -> bool {
1234 self.tracked_bytes
1235 .fetch_update(Ordering::AcqRel, Ordering::Acquire, |tracked| {
1236 let next = tracked.saturating_add_signed(growth);
1237 (growth <= 0 || next <= self.budget.get()).then_some(next)
1238 })
1239 .is_ok()
1240 }
1241
1242 fn claim(&self, key: Arc<OfferedKey>, account: AccountKey) {
1243 self.owners
1244 .entry_sync(key)
1245 .or_default()
1246 .get_mut()
1247 .add(account);
1248 }
1249
1250 fn disown(&self, key: &OfferedKey, account: AccountKey) {
1251 let _ = self
1252 .owners
1253 .remove_if_sync(key, |publishers| publishers.remove(account));
1254 }
1255}