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