This repository has no description
0

Configure Feed

Select the types of activity you want to include in your feed.

core / knot2 / crates / knot-index / src / projections.rs
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, &registry); 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(&registration.repo); 673 let owner = interner.intern_owner(&registration.owner); 674 let rkey = interner.intern_rkey(&registration.rkey); 675 let name = interner.intern_name(&registration.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}