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
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, &registry); 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(&registration.repo); 674 let owner = interner.intern_owner(&registration.owner); 675 let rkey = interner.intern_rkey(&registration.rkey); 676 let name = interner.intern_name(&registration.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}