diff --git a/libs/@local/graph/atlas/src/serve/membership/mod.rs b/libs/@local/graph/atlas/src/serve/membership/mod.rs new file mode 100644 index 00000000000..14def78c4f2 --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/membership/mod.rs @@ -0,0 +1,10 @@ +//! Type membership over the rows a request delivers. +//! +//! A request names the entity types it wants marked, and the response reports, for each delivered +//! row, which of those types it belongs to. Answering that needs one membership set per requested +//! type, held for the length of the response assembly. [`OntologySelection`] is the request's list +//! and [`SelectionSlot`] is a position within it. + +mod ontology; + +pub(crate) use self::ontology::{OntologySelection, SelectionSlot}; diff --git a/libs/@local/graph/atlas/src/serve/membership/ontology.rs b/libs/@local/graph/atlas/src/serve/membership/ontology.rs new file mode 100644 index 00000000000..fde2e64e180 --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/membership/ontology.rs @@ -0,0 +1,131 @@ +//! Resolution of requested ontology types into per-row membership tests. +//! +//! A generation records type membership in one of two forms. The closure form is a dense bit set +//! over base positions, materialized for the types worth precomputing. The postings form is the +//! sparse list every other type keeps. [`OntologyMembership`] holds whichever form a type has. The +//! per-row test is therefore one call whatever the generation recorded. +#![expect(clippy::empty_enums, reason = "zerocopy uses them in the derive")] + +use hashql_core::id::IdVec; + +use crate::{ + bitset::DenseBitSlice, + identity::BasePosition, + postgres::id::ArchivedOntologyTypeUuid, + salt::{fit::prepare::IdentityProvider as _, postings::artifact::Membership}, + serve::world::Ontology, +}; + +/// One requested type's membership set, in whichever form the generation recorded. +pub(crate) enum OntologyMembership<'ontology> { + /// The sparse postings list of the positions belonging to the type. + Direct(Membership<'ontology>), + /// The materialized closure, a dense bit set over every base position. + Closure(&'ontology DenseBitSlice), + /// The generation records nothing for the requested type. No row belongs to it. + Unresolved, +} + +impl<'ontology> OntologyMembership<'ontology> { + /// Resolves the type `id` against the generation's recorded membership. + /// + /// A closure is preferred where one exists, because the dense form answers a test without a + /// search. An identifier the generation never recorded resolves to + /// [`Unresolved`](Self::Unresolved) rather than refusing: a request may name a type this + /// generation has no rows for, and the answer is an empty membership. + pub(crate) fn new(ontology: &'ontology Ontology, id: ArchivedOntologyTypeUuid) -> Self { + let Some(id) = ontology.identity().row_of(id) else { + return Self::Unresolved; + }; + + ontology.closure().membership(id).map_or_else( + || { + ontology + .postings() + .membership(id) + .map_or(Self::Unresolved, Self::Direct) + }, + Self::Closure, + ) + } + + /// Returns whether the base row at `position` belongs to this type. + pub(crate) fn contains(&self, position: BasePosition) -> bool { + match self { + Self::Direct(membership) => membership.contains(position), + Self::Closure(membership) => membership.contains(position), + Self::Unresolved => false, + } + } +} + +hashql_core::id::newtype! { + /// A requested type's position, with duplicate requests occupying distinct slots. + pub(crate) struct SelectionSlot(u32) +} + +/// The ontology types a request asks to have marked, in request order. +/// +/// An unsized view over the decoded request's type list, keeping the request's own order and its +/// duplicates. A [`SelectionSlot`] indexes into it, and the response's membership bits are +/// reported against those slots. +#[derive( + Debug, zerocopy::FromBytes, zerocopy::IntoBytes, zerocopy::KnownLayout, zerocopy::Immutable, +)] +#[repr(C)] +pub(crate) struct OntologySelection([ArchivedOntologyTypeUuid]); + +/// A resolved membership set per requested slot, borrowed from one generation's ontology. +pub(crate) struct OntologyMemberships<'context> { + /// The membership set of every requested slot, in request order. + memberships: IdVec>, +} + +impl<'context> OntologyMemberships<'context> { + /// Iterates the resolved memberships beside the slot each answers for. + pub(crate) fn iter_enumerated( + &self, + ) -> impl ExactSizeIterator)> { + self.memberships.iter_enumerated() + } +} + +impl OntologySelection { + /// Returns the number of requested slots, counting a repeated type once per request. + pub(crate) const fn len(&self) -> usize { + self.0.len() + } + + /// Returns whether the request named no types at all. + pub(crate) const fn is_empty(&self) -> bool { + self.0.is_empty() + } + + /// Returns whether `id` occupies at least one requested slot. + pub(crate) fn contains(&self, id: ArchivedOntologyTypeUuid) -> bool { + self.0.contains(&id) + } + + /// Views a decoded request's type list as a selection, without copying it. + pub(crate) fn new(ontology: &[ArchivedOntologyTypeUuid]) -> &Self { + zerocopy::transmute_ref!(ontology) + } + + /// Resolves every requested slot against `ontology`. + /// + /// Resolution happens once per response, and the result borrows the generation's recorded + /// membership rather than copying it. A duplicate request resolves once per slot, which keeps + /// slot indices aligned with the request the client sent. + pub(crate) fn resolve<'context>( + &'context self, + ontology: &'context Ontology, + ) -> OntologyMemberships<'context> { + OntologyMemberships { + memberships: self + .0 + .iter() + .map(|&id| OntologyMembership::new(ontology, id)) + .collect(), + } + } +} diff --git a/libs/@local/graph/atlas/src/serve/mod.rs b/libs/@local/graph/atlas/src/serve/mod.rs index 3c3e4667763..b82d0349fae 100644 --- a/libs/@local/graph/atlas/src/serve/mod.rs +++ b/libs/@local/graph/atlas/src/serve/mod.rs @@ -4,7 +4,7 @@ //! and removal. Without feed options or temporal axes, a generation instead exposes a static //! publication. An active feed can lag or fail. Each scene-backed delivery request //! captures a coherent world-and-delta epoch and obtains a cached or newly resolved visibility -//! scope before constructing a `scene::Scene`. A cached mask and schedule may predate the +//! scope before constructing a [`scene::Scene`]. A cached mask and schedule may predate the //! request's epoch within the same delta lifetime. Scene geometry, identity and topology lookups //! still use only the request's captured publication. //! @@ -21,7 +21,10 @@ pub(crate) mod delta; pub(crate) mod density; pub(crate) mod hydrate; mod intern; +pub(crate) mod membership; +mod neighbourhood; pub(crate) mod runtime; +pub(crate) mod scene; mod schedule; pub(crate) mod secret; #[cfg(test)] diff --git a/libs/@local/graph/atlas/src/serve/neighbourhood/mod.rs b/libs/@local/graph/atlas/src/serve/neighbourhood/mod.rs new file mode 100644 index 00000000000..6629694e256 --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/neighbourhood/mod.rs @@ -0,0 +1,321 @@ +//! Edge delivery around delivered nodes: incident edges and the capped induced edge set. +//! +//! A response delivers nodes and then the edges among them. [`Neighbourhood`] answers both edge +//! questions over a [`NeighbourhoodProvider`]: every visible edge incident to one node, and the +//! best-ranked edges induced by a delivered node set under a size cap. An induced edge ranks at +//! its less prominent endpoint, and [`EdgeSet::complete`] reports whether the cap excluded any. + +use alloc::collections::BinaryHeap; +use core::cmp::Ordering; + +use super::{ + scene::Scene, + visibility::Visible, + world::node_importance::{ImportanceProvider, NodePriority}, +}; +use crate::{ + bitset::CompressedBitSet, + identity::{EdgeRowId, NodeRowId}, + postgres::id::ArchivedEntityId, +}; + +#[cfg(test)] +mod tests; + +/// An admitted edge's visible row, endpoint node rows and stable identity. +#[derive(Debug, Copy, Clone, PartialEq, Eq)] +pub(crate) struct DeliveredEdge { + /// The edge row, carrying the evidence that the mask admitted it. + pub row: Visible, + /// The `[source, target]` node rows. + pub endpoints: [NodeRowId; 2], + /// The edge's entity key. + pub identity: ArchivedEntityId, +} + +impl DeliveredEdge { + /// Returns the endpoint opposite `node`, or [`None`] when `node` is neither endpoint. + pub(crate) fn partner_of(self, node: NodeRowId) -> Option { + let [source, target] = self.endpoints; + + if source == node { + Some(target) + } else if target == node { + Some(source) + } else { + None + } + } +} + +/// A capped selection of delivered edges, and whether the cap excluded any. +#[derive(Debug)] +pub(crate) struct EdgeSet { + /// Whether the cap kept every offered edge. + pub complete: bool, + /// The kept edges, in identity order. + pub edges: Vec, +} + +/// One edge under [`RankCap`] ordering: its priority and identity, then the edge itself. +#[derive(Debug)] +struct Candidate { + /// The offer's priority, then the edge's identity as the tie-break. + key: (NodePriority, ArchivedEntityId), + /// The offered edge. + edge: DeliveredEdge, +} + +impl PartialEq for Candidate { + /// Compares `key` alone, which keeps equality consistent with [`Ord`]. + /// + /// Candidates with equal keys share the edge identity the key embeds. + fn eq(&self, other: &Self) -> bool { + self.key == other.key + } +} + +impl Eq for Candidate {} + +impl PartialOrd for Candidate { + fn partial_cmp(&self, other: &Self) -> Option { + Some(self.cmp(other)) + } +} + +impl Ord for Candidate { + fn cmp(&self, other: &Self) -> Ordering { + self.key.cmp(&other.key) + } +} + +/// A bounded selection of the `capacity` best-ranked edges offered to it. +/// +/// The heap's top is the worst kept edge, the one a better offer displaces. +#[derive(Debug)] +struct RankCap { + /// The most edges the cap keeps. + capacity: usize, + /// The kept candidates, worst on top. + kept: BinaryHeap, + /// Whether an offer has arrived at a full cap. + truncated: bool, +} + +impl RankCap { + /// Creates an empty cap retaining at most `capacity` edges. + const fn new(capacity: usize) -> Self { + Self { + capacity, + kept: BinaryHeap::new(), + truncated: false, + } + } + + /// Offers `edge` at `priority`. + /// + /// The cap keeps it only while it has spare capacity or the edge outranks the current worst + /// kept edge. + fn offer(&mut self, priority: NodePriority, edge: DeliveredEdge) { + let candidate = Candidate { + key: (priority, edge.identity), + edge, + }; + if self.kept.len() < self.capacity { + self.kept.push(candidate); + return; + } + + self.truncated = true; + if let Some(mut worst) = self.kept.peek_mut() + && candidate < *worst + { + *worst = candidate; + } + } + + /// Returns whether the cap can prune `priority` without comparing edge identities. + /// + /// Pruning begins only after an offer reaches an already full cap. A zero-capacity cap then + /// excludes every priority. A nonempty cap excludes priorities worse than its worst kept + /// edge's, leaving equal priorities for the identity comparison in [`offer`](Self::offer). + /// Before truncation even a worse priority needs an offer to establish the incompleteness + /// [`EdgeSet::complete`] reports. + fn excludes(&self, priority: NodePriority) -> bool { + // Only `offer` records a truncation, and `complete` reads it. A full cap that has turned + // nothing away yet therefore still receives its next offer rather than excluding it here. + self.truncated && self.kept.peek().is_none_or(|worst| priority > worst.key.0) + } + + /// Converts the cap into its kept edges sorted by identity, and whether it truncated. + fn into_set(self) -> EdgeSet { + let mut edges: Vec<_> = self + .kept + .into_iter() + .map(|candidate| candidate.edge) + .collect(); + edges.sort_unstable_by_key(|edge| edge.identity); + + EdgeSet { + complete: !self.truncated, + edges, + } + } +} + +/// The capability to answer edge lookups and per-node adjacency against a captured scene. +/// +/// An adjacency can include a row that [`provide_edge`](Self::provide_edge) withholds. The edge +/// queries consider only delivered rows. +pub(crate) trait NeighbourhoodProvider { + /// Returns the delivered edge at `row`, or [`None`] when it is absent, invisible or unplaced. + fn provide_edge(&self, row: EdgeRowId) -> Option; + /// Yields the edge rows whose target is `node`, each exactly once. + /// + /// # Implementation Note + /// + /// Implementations must yield every incoming edge row exactly once. For each yielded row that + /// [`provide_edge`](Self::provide_edge) delivers, the target in [`DeliveredEdge::endpoints`] + /// must be `node`. + fn provide_incoming(&self, node: NodeRowId) -> impl Iterator; + /// Yields the edge rows whose source is `node`, each exactly once, a self-loop included. + /// + /// # Implementation Note + /// + /// Implementations must yield every outgoing edge row exactly once. For each yielded row that + /// [`provide_edge`](Self::provide_edge) delivers, the source in [`DeliveredEdge::endpoints`] + /// must be `node`. + fn provide_outgoing(&self, node: NodeRowId) -> impl Iterator; +} + +impl NeighbourhoodProvider for Scene<'_> { + fn provide_edge(&self, row: EdgeRowId) -> Option { + let endpoints = self.world.topology.endpoints(self.epoch, row)?; + let visible = self.mask.visible_edge(row, endpoints)?; + + let identity = self.world.topology.key_of(self.epoch, row)?; + Some(DeliveredEdge { + row: visible, + endpoints, + identity, + }) + } + + fn provide_incoming(&self, node: NodeRowId) -> impl Iterator { + self.world.topology.incoming(self.epoch, node) + } + + fn provide_outgoing(&self, node: NodeRowId) -> impl Iterator { + self.world.topology.outgoing(self.epoch, node) + } +} + +impl ImportanceProvider for Scene<'_> { + fn provide_priority(&self, node: NodeRowId) -> Option { + self.world.layout.priority(self.epoch, node) + } +} + +/// Edge queries against one adjacency provider, most often a captured [`Scene`]. +#[derive(Debug, Copy, Clone)] +pub(crate) struct Neighbourhood

{ + /// The scene or other provider the queries read edges from. + pub provider: P, +} + +impl Neighbourhood

{ + /// Iterates every delivered edge incident to `node`, outgoing first. + /// + /// Each edge occurs once, a self-loop included: a self-loop occurs in the outgoing run, and + /// the incoming run drops every edge whose source is `node`. An edge the provider withholds + /// from [`provide_edge`](NeighbourhoodProvider::provide_edge) does not occur. + pub(crate) fn incident(&self, node: NodeRowId) -> impl Iterator + '_ { + let outgoing = self + .provider + .provide_outgoing(node) + .filter_map(|row| self.provider.provide_edge(row)); + + // Self-loops already occur in the outgoing run. + let incoming = self + .provider + .provide_incoming(node) + .filter_map(|row| self.provider.provide_edge(row)) + .filter(move |edge| edge.endpoints[0] != node); + + outgoing.chain(incoming) + } + + /// Offers every candidate edge between two nodes in `delivered` to `cap`. + /// + /// The scan skips a delivered source with no priority before reading its outgoing edges, an + /// edge the provider withholds from [`provide_edge`](NeighbourhoodProvider::provide_edge), and + /// an edge whose target is not delivered. Each remaining edge ranks at the less prominent of + /// its two endpoints, the larger [`NodePriority`]. + /// + /// # Panics + /// + /// Panics if the delivered target of a scanned edge has no priority. + fn offer_induced(&self, delivered: &CompressedBitSet, cap: &mut RankCap) + where + P: ImportanceProvider, + { + for source in delivered.iter() { + let Some(source_priority) = self.provider.provide_priority(source) else { + continue; + }; + // An edge ranks at its less prominent endpoint, and `source` is one of the two. Every + // edge `source` offers therefore ranks at or below `source`, and excluding `source` + // excludes them all. + if cap.excludes(source_priority) { + continue; + } + + // Each edge occurs in exactly one source's outgoing run. + for row in self.provider.provide_outgoing(source) { + let Some(edge) = self.provider.provide_edge(row) else { + continue; + }; + + let [_, target] = edge.endpoints; + if !delivered.contains(target) { + continue; + } + + let target_priority = self + .provider + .provide_priority(target) + .expect("should have a priority for every placed endpoint"); + + let priority = source_priority.max(target_priority); + if !cap.excludes(priority) { + cap.offer(priority, edge); + } + } + } + } + + /// Selects the `capacity` best-ranked edges between nodes in `delivered`. + /// + /// Candidates are outgoing edges of a node in `delivered` that has a priority, provided that + /// [`provide_edge`](NeighbourhoodProvider::provide_edge) delivers the edge and its target also + /// belongs to `delivered`. A source without a priority contributes no candidates. Each + /// candidate ranks at the less prominent of its endpoints, the larger [`NodePriority`], and the + /// lower edge identity wins a tie. The set returns in identity order, with + /// [`EdgeSet::complete`] false when the cap excluded an edge. + /// + /// # Panics + /// + /// Panics when the scan reaches an edge whose delivered target has no priority. + pub(crate) fn induced( + &self, + delivered: &CompressedBitSet, + capacity: usize, + ) -> EdgeSet + where + P: ImportanceProvider, + { + let mut cap = RankCap::new(capacity); + self.offer_induced(delivered, &mut cap); + cap.into_set() + } +} diff --git a/libs/@local/graph/atlas/src/serve/neighbourhood/tests.rs b/libs/@local/graph/atlas/src/serve/neighbourhood/tests.rs new file mode 100644 index 00000000000..68d5e621820 --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/neighbourhood/tests.rs @@ -0,0 +1,506 @@ +//! Cases covering incident and induced edge delivery over a scene and over a synthetic graph. + +use alloc::sync::Arc; + +use arc_swap::Guard; +use hashql_core::id::{Id as _, IdVec}; +use proptest::{arbitrary::any, collection, prop_assert_eq, property_test}; +use rand::{SeedableRng as _, rngs::StdRng}; +use type_system::principal::actor::{ActorId, ActorType}; +use uuid::Uuid; + +use super::{DeliveredEdge, EdgeSet, Neighbourhood, NeighbourhoodProvider, RankCap}; +use crate::{ + bitset::CompressedBitSet, + identity::{EdgeRowId, ImportanceRank, NodeRowId}, + morton::Zoom, + postgres::id::ArchivedEntityId, + serve::{ + delta::{Delta, epoch::Epoch}, + scene::Scene, + schedule::ViewSchedule, + tests::fixture::{TamperFixture, secret}, + visibility::{VisibilityActor, VisibilityMask, Visible}, + world::{ + World, + node_importance::{ImportanceProvider, NodePriority}, + }, + }, +}; + +/// The captured scene the [`Scene`]-backed cases resolve against. +/// +/// An opened synthetic generation with a captured epoch, actor, visibility mask and delivery +/// schedule. +struct SceneFixture { + /// The opened world. + world: Arc, + /// The captured publication. + epoch: Epoch, + /// The principal the mask resolves for. + actor: VisibilityActor, + /// The current visibility mask. + mask: VisibilityMask, + /// The view schedule captured under `mask`. + schedule: ViewSchedule, + /// The published files, kept until the fixture drops. + _files: TamperFixture, +} + +impl SceneFixture { + /// Publishes and opens the synthetic generation under a root named `name`. + /// + /// The fixture starts with a fresh delta identity and a full-visibility mask for one fixed + /// test actor. + /// + /// # Panics + /// + /// Panics if publishing or opening the generation fails, or allocating the delta does. + fn new(name: &str) -> Self { + let files = TamperFixture::publish(name); + let world = Arc::new( + World::open(files.generation().clone(), &secret()) + .expect("should open the synthetic generation"), + ); + let delta = Delta::new(Arc::clone(&world), StdRng::seed_from_u64(17)) + .expect("should allocate a delta identity"); + let epoch = Epoch::from(Guard::from_inner(Arc::new(delta))); + let actor = VisibilityActor { + id: ActorId::new(Uuid::from_u128(1), ActorType::User), + instance_admin: false, + }; + let mask = VisibilityMask::full(actor); + let schedule = ViewSchedule::of(Arc::clone(&world), &epoch, &mask); + Self { + world, + epoch, + actor, + mask, + schedule, + _files: files, + } + } + + /// Narrows the visibility mask and schedule to exactly the supplied node and edge rows. + /// + /// A later `scene()` call sees only those rows. + fn restrict( + &mut self, + nodes: impl IntoIterator, + edges: impl IntoIterator, + ) { + self.mask = VisibilityMask::partial( + self.actor, + CompressedBitSet::from_rows(nodes.into_iter().map(NodeRowId::new)), + CompressedBitSet::from_rows(edges.into_iter().map(EdgeRowId::new)), + ); + self.schedule = ViewSchedule::of(Arc::clone(&self.world), &self.epoch, &self.mask); + } + + /// Assembles the current world, epoch, mask and schedule as one [`Scene`] at offset zero. + fn scene(&self) -> Scene<'_> { + Scene { + world: &self.world, + epoch: &self.epoch, + mask: &self.mask, + schedule: &self.schedule, + delivery: self + .schedule + .cut(Zoom::MIN) + .expect("should bind the zero offset"), + } + } +} + +/// A scene-backed neighbourhood returns only edges both endpoints' visibility admits. +/// +/// This holds for `incident` and `induced` alike. +#[test] +fn scene_admission() { + let mut fixture = SceneFixture::new("neighbourhood-scene-admission"); + let source = NodeRowId::new(2); + assert_eq!( + Neighbourhood { + provider: fixture.scene() + } + .incident(source) + .count(), + 2 + ); + let cases: [(&[u64], &[u64], &[u64]); 4] = [ + (&[1, 2], &[1], &[1]), + (&[2], &[1, 2], &[2]), + (&[1], &[1, 2], &[]), + (&[1, 2], &[], &[]), + ]; + let delivered = CompressedBitSet::from_rows([NodeRowId::new(1), source]); + for (nodes, edges, expected) in cases { + fixture.restrict(nodes.iter().copied(), edges.iter().copied()); + let neighbourhood = Neighbourhood { + provider: fixture.scene(), + }; + let actual: Vec<_> = neighbourhood + .incident(source) + .map(|edge| edge.row.get()) + .collect(); + assert_eq!(actual, expected); + let actual = neighbourhood.induced(&delivered, 2); + assert!(actual.complete); + assert_eq!( + actual + .edges + .iter() + .map(|edge| edge.row.get()) + .collect::>(), + expected + ); + } +} + +/// A synthetic [`NeighbourhoodProvider`] over explicit node priorities and linked edges. +/// +/// Hidden edges stand in for visibility masking without a real [`Scene`]. +struct Graph { + /// Every node's priority, in row order. + priorities: IdVec, + /// Every linked edge, in row order. + edges: IdVec, + /// The edge rows the provider withholds. + hidden: CompressedBitSet, +} + +impl Graph { + /// Builds a graph with one node per rank in `ranks`, in row order, and no edges. + fn new(ranks: impl IntoIterator) -> Self { + Self { + priorities: ranks + .into_iter() + .map(|rank| NodePriority::Rank(ImportanceRank::new(rank))) + .collect(), + edges: IdVec::new(), + hidden: CompressedBitSet::new(), + } + } + + /// Adds and returns an edge between `endpoints` with a synthetic identity from `identity`. + /// + /// The endpoints may coincide in a self-loop. + fn link(&mut self, endpoints: [u64; 2], identity: u128) -> DeliveredEdge { + let edge = DeliveredEdge { + row: Visible::new(EdgeRowId::from_usize(self.edges.len()).get()), + endpoints: endpoints.map(NodeRowId::new), + identity: ArchivedEntityId { + web_id: Uuid::from_u128(1).into(), + entity_uuid: Uuid::from_u128(identity).into(), + }, + }; + self.edges.push(edge); + edge + } + + /// Selects the `capacity`-capped induced edge set among `rows` over this graph. + /// + /// [`Neighbourhood::induced`] computes it. + /// + /// # Panics + /// + /// Panics when the scan reaches an edge whose delivered target has no priority. + fn select(&self, rows: impl IntoIterator, capacity: usize) -> EdgeSet { + let delivered = CompressedBitSet::from_rows(rows.into_iter().map(NodeRowId::new)); + Neighbourhood { provider: self }.induced(&delivered, capacity) + } + + /// Selects the `capacity`-capped induced edge set by a full sort. + /// + /// The sort orders by the larger endpoint priority and edge identity, truncates to + /// `capacity`, then returns the selected edges in identity order. This is the reference + /// result for [`Neighbourhood::induced`]'s incremental selection. + /// + /// # Panics + /// + /// Panics if an edge between delivered nodes refers to a missing priority. + fn reference(&self, delivered: &CompressedBitSet, capacity: usize) -> EdgeSet { + let mut edges: Vec<_> = self + .edges + .iter() + .copied() + .filter(|edge| { + let [source, target] = edge.endpoints; + !self.hidden.contains(EdgeRowId::new(edge.row.get())) + && delivered.contains(source) + && delivered.contains(target) + }) + .collect(); + edges.sort_unstable_by_key(|edge| { + let [source, target] = edge.endpoints; + ( + self.priorities[source].max(self.priorities[target]), + edge.identity, + ) + }); + let complete = edges.len() <= capacity; + edges.truncate(capacity); + edges.sort_unstable_by_key(|edge| edge.identity); + EdgeSet { complete, edges } + } +} + +impl NeighbourhoodProvider for &Graph { + fn provide_edge(&self, row: EdgeRowId) -> Option { + if self.hidden.contains(row) { + return None; + } + self.edges.get(row).copied() + } + + fn provide_incoming(&self, node: NodeRowId) -> impl Iterator { + self.edges + .iter_enumerated() + .filter_map(move |(row, edge)| (edge.endpoints[1] == node).then_some(row)) + } + + fn provide_outgoing(&self, node: NodeRowId) -> impl Iterator { + self.edges + .iter_enumerated() + .filter_map(move |(row, edge)| (edge.endpoints[0] == node).then_some(row)) + } +} + +impl ImportanceProvider for &Graph { + fn provide_priority(&self, node: NodeRowId) -> Option { + self.priorities.get(node).copied() + } +} + +/// `incident` returns each parallel edge and the self-loop exactly once. +/// +/// An unknown node yields no edges. +#[test] +fn incident_parallel_and_self_loop() { + let mut graph = Graph::new([0, 1, 2]); + let expected = [ + graph.link([0, 1], 0), + graph.link([0, 1], 1), + graph.link([1, 0], 2), + graph.link([1, 1], 3), + graph.link([1, 2], 4), + ]; + let neighbourhood = Neighbourhood { provider: &graph }; + let mut actual: Vec<_> = neighbourhood.incident(NodeRowId::new(1)).collect(); + actual.sort_unstable_by_key(|edge| edge.identity); + assert_eq!(actual, expected); + assert_eq!(neighbourhood.incident(NodeRowId::new(99)).count(), 0); +} + +/// `incident` omits a hidden incoming edge and a hidden self-loop. +/// +/// Only the visible outgoing edge remains. +#[test] +fn incident_hidden_edge() { + let mut graph = Graph::new([0, 1]); + let incoming = graph.link([0, 1], 0); + let outgoing = graph.link([1, 0], 1); + let self_loop = graph.link([1, 1], 2); + for edge in [incoming, self_loop] { + graph.hidden.insert(EdgeRowId::new(edge.row.get())); + } + assert_eq!( + Neighbourhood { provider: &graph } + .incident(NodeRowId::new(1)) + .collect::>(), + [outgoing], + ); +} + +/// A row repeated in the delivered set does not duplicate its edges in the induced result. +#[test] +fn induced_duplicate_rows() { + let mut graph = Graph::new([0, 1]); + let expected = [ + graph.link([0, 1], 0), + graph.link([1, 0], 1), + graph.link([1, 1], 2), + ]; + let actual = graph.select([0, 1, 0, 1], expected.len()); + assert_eq!(actual.edges, expected); + assert!(actual.complete); +} + +/// Under a tight capacity, `induced` ranks each edge at its less prominent endpoint. +/// +/// It keeps the edge whose less prominent endpoint outranks the other candidate's, and marks the +/// result incomplete. +#[test] +fn induced_worse_endpoint() { + let mut graph = Graph::new([0, 9, 4, 5]); + graph.link([0, 1], 0); + let expected = graph.link([2, 3], 1); + let actual = graph.select(0..4, 1); + assert_eq!(actual.edges, [expected]); + assert!(!actual.complete); +} + +/// `induced` omits hidden edges and edges with an undelivered endpoint without charging the cap. +/// +/// A zero capacity is complete only when no qualifying edge exists. +#[test] +fn induced_nonqualifying_tail() { + let mut graph = Graph::new([0, 1, 2, 3, 4]); + let expected = graph.link([0, 1], 0); + let hidden = graph.link([2, 3], 1); + graph.hidden.insert(EdgeRowId::new(hidden.row.get())); + graph.link([3, 4], 2); + + let actual = graph.select([0, 1, 2, 3, 99], 1); + assert_eq!(actual.edges, [expected]); + assert!(actual.complete); + let actual = graph.select([0, 1, 2, 3], 0); + assert!(actual.edges.is_empty()); + assert!(!actual.complete); + let actual = graph.select([2, 3], 0); + assert!(actual.edges.is_empty()); + assert!(actual.complete); +} + +/// `partner_of` returns the other endpoint from either side of a plain edge. +/// +/// It returns [`None`] for a node not on the edge, and the same node for a self-loop. +#[test] +fn partner_direction() { + let mut graph = Graph::new([0, 1]); + let edge = graph.link([0, 1], 0); + assert_eq!(edge.partner_of(NodeRowId::new(0)), Some(NodeRowId::new(1))); + assert_eq!(edge.partner_of(NodeRowId::new(1)), Some(NodeRowId::new(0))); + assert_eq!(edge.partner_of(NodeRowId::new(2)), None); + let self_loop = graph.link([0, 0], 1); + assert_eq!( + self_loop.partner_of(NodeRowId::new(0)), + Some(NodeRowId::new(0)) + ); +} + +/// A zero-capacity [`RankCap`] excludes nothing before an offer and everything after one. +/// +/// The resulting set is empty and incomplete, and an untouched zero-capacity cap is complete. +#[test] +fn cap_zero() { + let mut graph = Graph::new([0, 1]); + let edge = graph.link([0, 1], 0); + let priority = graph.priorities[NodeRowId::new(1)]; + let empty = RankCap::new(0); + assert!(!empty.excludes(priority)); + assert!(empty.into_set().complete); + + let mut cap = RankCap::new(0); + cap.offer(priority, edge); + assert!(cap.excludes(priority)); + let actual = cap.into_set(); + assert!(!actual.complete); + assert!(actual.edges.is_empty()); +} + +/// A full cap is not yet truncated and remains complete until it turns an offer away. +/// +/// At that point it reports both the exclusion and the incompleteness. +#[test] +fn cap_full_before_truncation() { + let mut graph = Graph::new([0, 1]); + let best = graph.priorities[NodeRowId::new(0)]; + let worst = graph.priorities[NodeRowId::new(1)]; + let first = graph.link([0, 0], 0); + let second = graph.link([0, 1], 1); + let mut cap = RankCap::new(1); + cap.offer(best, first); + assert!( + !cap.excludes(worst), + "a full selection can still be complete" + ); + assert!(cap.into_set().complete); + + let mut cap = RankCap::new(1); + cap.offer(best, first); + cap.offer(worst, second); + assert!(cap.excludes(worst)); + assert!(!cap.into_set().complete); +} + +/// Equal priorities still compare by identity. +/// +/// The cap keeps the lower-identity edge and excludes on that comparison even when priorities +/// alone would not decide it. +#[test] +fn cap_identity_tie() { + let mut graph = Graph::new([7, 8]); + let priority = graph.priorities[NodeRowId::new(0)]; + let worse = graph.priorities[NodeRowId::new(1)]; + let first = graph.link([0, 0], 9); + let loser = graph.link([0, 1], 10); + let winner = graph.link([0, 0], 1); + let mut cap = RankCap::new(1); + cap.offer(priority, first); + cap.offer(worse, loser); + assert!( + !cap.excludes(priority), + "equal priorities still compare identities" + ); + assert!(cap.excludes(worse)); + cap.offer(priority, winner); + let actual = cap.into_set(); + assert_eq!(actual.edges, [winner]); + assert!(!actual.complete); +} + +/// A ranked priority always outranks an identity-only priority, whatever the identity values. +#[test] +fn cap_priority_domains() { + let mut graph = Graph::new([0, 1]); + let unranked = graph.link([0, 1], 0); + let ranked = graph.link([0, 1], 100); + let mut cap = RankCap::new(1); + cap.offer(NodePriority::Identity(unranked.identity), unranked); + cap.offer(NodePriority::Rank(ImportanceRank::MAX), ranked); + let actual = cap.into_set(); + assert_eq!(actual.edges, [ranked]); + assert!(!actual.complete); +} + +/// [`Neighbourhood::induced`]'s incremental selection agrees with a full sort. +/// +/// [`Graph::reference`] is the sort, over randomized ranks, membership, edges and capacities. +#[property_test] +fn induced_reference( + ranks: [u8; 8], + members: [bool; 8], + #[strategy = collection::vec((0_u8..8, 0_u8..8, any::(), any::()), 0..40)] + inputs: Vec<(u8, u8, u16, bool)>, + #[strategy = 0_usize..45] capacity: usize, +) { + let mut graph = Graph::new(ranks.map(u32::from)); + for (node, priority) in graph.priorities.iter_enumerated_mut() { + let rank = ranks[node.as_usize()]; + if rank & 0x80 != 0 { + *priority = NodePriority::Identity(ArchivedEntityId { + web_id: Uuid::from_u128(1).into(), + entity_uuid: Uuid::from_u128(u128::from(rank & 0x7F)).into(), + }); + } + } + for (index, (source, target, key, admitted)) in inputs.into_iter().enumerate() { + let index = u64::try_from(index).expect("should fit the generated edge index"); + let identity = (u128::from(key) << 64) | u128::from(index); + let edge = graph.link([u64::from(source), u64::from(target)], identity); + if !admitted { + graph.hidden.insert(EdgeRowId::new(edge.row.get())); + } + } + let delivered = CompressedBitSet::from_rows( + members + .into_iter() + .enumerate() + .filter_map(|(index, member)| member.then_some(NodeRowId::from_usize(index))), + ); + let expected = graph.reference(&delivered, capacity); + let actual = Neighbourhood { provider: &graph }.induced(&delivered, capacity); + prop_assert_eq!(actual.complete, expected.complete); + prop_assert_eq!(actual.edges, expected.edges); +} diff --git a/libs/@local/graph/atlas/src/serve/scene.rs b/libs/@local/graph/atlas/src/serve/scene.rs new file mode 100644 index 00000000000..a72552904fe --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/scene.rs @@ -0,0 +1,86 @@ +//! Request-bound document assembly over a publication and cached visibility scope. + +use error_stack::{Report, ResultExt as _}; + +use super::{ + delta::epoch::Epoch, + schedule::{DeliverySchedule, ViewSchedule}, + visibility::{VisibilityMask, cache::CacheEntry}, + world::World, +}; +use crate::morton::Zoom; + +/// The reason [`Scene::of`] could not assemble a scene. +#[derive(Debug)] +pub(crate) enum SceneError { + /// The requested density offset puts a scoped delivery cut beyond the Morton key width. + Delivery, +} + +impl core::fmt::Display for SceneError { + fn fmt(&self, fmt: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + match self { + Self::Delivery => fmt.write_str("unable to build delivery schedule"), + } + } +} + +impl core::error::Error for SceneError {} + +/// A request publication bound to one cached visibility scope and delivery schedule. +/// +/// The world and epoch belong to the admitted request. The mask and delivery schedule belong to a +/// cache entry for the same generation and delta lifetime, but the cache may have resolved them at +/// an older revision. Identity, geometry and topology reads use `epoch`. Row admission uses the +/// cached mask, while delivery order and Morton keys come from the cached schedule. Scoped delivery +/// applies the requested density offset across the generation's served tile-zoom range. Corpus +/// delivery keeps its recorded cuts for every offset. Every field shares the scene's borrow +/// lifetime. +#[derive(Debug, Copy, Clone)] +pub(crate) struct Scene<'scene> { + /// The requested generation's opened serving artifacts. + pub world: &'scene World, + /// The immutable delta publication captured for this request. + pub epoch: &'scene Epoch, + /// The cached authorization decision applied to delivered rows. + pub mask: &'scene VisibilityMask, + /// The cached row assignment and keys from visibility resolution. + pub schedule: &'scene ViewSchedule, + /// The offset-bound scoped cuts or the recorded corpus cuts. + pub delivery: DeliverySchedule<'scene>, +} + +impl<'scope> Scene<'scope> { + /// Binds a cached visibility scope to the request's publication and delivery schedule. + /// + /// `zoom` deepens the scoped delivery cuts while preserving the generation's served tile-zoom + /// range. `world` and `epoch` must come from the same requested + /// [`Universe`](crate::serve::runtime::registry::Universe). `entry` must be the scope returned + /// for that request's generation and delta lifetime. This method does not verify either + /// association. The publication used to build `entry` need not have the same revision because + /// visibility caching intentionally reuses a scope across publications within one lifetime. + /// + /// # Errors + /// + /// Returns [`SceneError::Delivery`] when `zoom` puts a scoped delivery cut beyond the Morton + /// key width. + pub(crate) fn of( + world: &'scope World, + epoch: &'scope Epoch, + entry: &'scope CacheEntry, + zoom: Zoom, + ) -> Result> { + let delivery = entry + .schedule + .cut(zoom) + .change_context(SceneError::Delivery)?; + + Ok(Scene { + world, + epoch, + mask: &entry.mask, + schedule: &entry.schedule, + delivery, + }) + } +} diff --git a/libs/@local/graph/atlas/src/serve/visibility/cache/error.rs b/libs/@local/graph/atlas/src/serve/visibility/cache/error.rs new file mode 100644 index 00000000000..bbfaece337a --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/visibility/cache/error.rs @@ -0,0 +1,22 @@ +//! Failures while preparing a visibility cache entry. + +use core::fmt; + +/// A failure preparing a visibility cache entry. +#[derive(Debug)] +pub(crate) enum VisibilityCacheError { + /// The offloaded schedule construction panicked or its worker returned no value. + /// + /// Both [`OffloadError`](crate::offload::OffloadError) cases collapse into this context. + Panic, +} + +impl fmt::Display for VisibilityCacheError { + fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Panic => write!(fmt, "visibility cache panicked"), + } + } +} + +impl core::error::Error for VisibilityCacheError {} diff --git a/libs/@local/graph/atlas/src/serve/visibility/cache/filter.rs b/libs/@local/graph/atlas/src/serve/visibility/cache/filter.rs new file mode 100644 index 00000000000..68d5499c90e --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/visibility/cache/filter.rs @@ -0,0 +1,55 @@ +//! Content identity for a request's visibility filter. +#![expect(clippy::empty_enums, reason = "zerocopy uses them in the derive")] + +use core::hash::{Hash, Hasher}; + +use crate::integrity::{Sha256, Sha256Digest, Update as _}; + +/// A SHA-256 token derived from a filter document's exact bytes. +/// +/// Hashing the filter document gives the cache key and fixed-width authority token a shared +/// content identity. No JSON canonicalization occurs. Formatting differences are distinct hash +/// inputs. Equality compares only SHA-256 outputs: collisions are not detected or resolved by +/// comparing retained documents. +#[derive( + Debug, + Copy, + Clone, + zerocopy::IntoBytes, + zerocopy::FromBytes, + zerocopy::Immutable, + zerocopy::Unaligned, + zerocopy::KnownLayout, +)] +#[repr(transparent)] +pub(crate) struct FilterDigest(Sha256Digest); + +impl FilterDigest { + /// Digests a filter document's exact bytes. + /// + /// Before hashing, this method prefixes the bytes, distinguishing filter documents from other + /// content hashed by this crate without parsing or canonicalizing the document. + pub(crate) fn of(filter: &[u8]) -> Self { + let mut hasher = Sha256::new(); + hasher.update(b"filter"); + hasher.update(filter); + + Self(hasher.finalize()) + } +} + +const impl PartialEq for FilterDigest { + #[inline] + fn eq(&self, other: &Self) -> bool { + self.0 == other.0 + } +} + +const impl Eq for FilterDigest {} + +impl Hash for FilterDigest { + #[inline] + fn hash(&self, state: &mut H) { + self.0.hash(state); + } +} diff --git a/libs/@local/graph/atlas/src/serve/visibility/cache/mod.rs b/libs/@local/graph/atlas/src/serve/visibility/cache/mod.rs new file mode 100644 index 00000000000..89a97810b78 --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/visibility/cache/mod.rs @@ -0,0 +1,478 @@ +//! Retention and refresh of resolved visibility scopes. +//! +//! Caching avoids repeating a permission-store round trip and schedule construction over every +//! visible row for each request. A [`CacheKey`] contains the generation, sampled delta-lifetime +//! tag, actor and filter digest, but not the delta revision. A hit may carry an older authorization +//! mask and delivery schedule while request data uses a newer publication from the same lifetime. +//! +//! Freshness starts at the resolving request's admission time and does not restart when resolution +//! completes. An entry becomes stale at the soft age, but a lookup still receives it. A lookup +//! whose current epoch matches the key starts a detached refresh only when it acquires the refresh +//! claim. At the hard age, that matching lookup waits for foreground resolution. A lookup under +//! another generation or lifetime removes the expired entry and returns [`None`]. +//! +//! Foreground resolution belongs to the caller's future. A stale refresh runs in a detached task +//! and may finish after that caller or the observed generation ceases to be present. Entry weights +//! measure capacity. The cache enforces that capacity as a best-effort target rather than an +//! instantaneous memory ceiling. + +use alloc::sync::Arc; +use core::{ + sync::atomic::{self, Atomic}, + time::Duration, +}; +use std::time::Instant; + +use error_stack::{Report, ResultExt as _}; +use moka::ops::compute::{CompResult, Op}; +use serde_json::value::RawValue; +use tracing::Instrument as _; +use type_system::principal::actor::ActorId; + +use self::error::VisibilityCacheError; +use crate::{ + allocator::HeapMemoryUsage as _, + file::generation::GenerationId, + offload, + serve::{ + delta::{DeltaId, epoch::Epoch}, + density::ViewOccupancy, + schedule::ViewSchedule, + visibility::VisibilityMask, + world::World, + }, +}; + +mod error; +mod filter; + +#[cfg(test)] +mod tests; + +pub(crate) use self::filter::FilterDigest; + +/// Computes the cache charge for an entry retaining `retained` heap bytes under `filter`. +/// +/// The charge covers the entry's own inline bytes, its key, the heap the mask and schedule hold, +/// and the retained filter document. Saturating arithmetic caps an oversized charge at +/// [`u32::MAX`], the largest weight accepted by the cache. +fn weight_of(retained: u64, filter: Option<&RawValue>) -> u32 { + let inline = size_of::() as u64 + size_of::() as u64; + let total = retained + .saturating_add(inline) + .saturating_add(filter.map_or(0, |document| document.get().len() as u64)); + + total.saturating_cast() +} + +/// A resolved scope that has not been admitted to the cache yet. +/// +/// The mask and schedule capture one epoch. The cache may later reuse them at newer revisions of +/// the same delta lifetime. This type separates finished, weighed resolution from cache admission. +/// Foreground misses and expiries construct it while holding the per-key compute lock. Detached +/// refreshes construct it before taking that lock for conditional replacement. Insertion turns it +/// into a [`CacheEntry`] by adding refresh state, the resolving lookup's admission time and a cache +/// publication number allocated only then. +#[derive(Debug)] +pub(crate) struct PendingCacheEntry { + mask: VisibilityMask, + schedule: ViewSchedule, + filter: Option>, + occupancy: Option, + weight: u32, +} + +impl PendingCacheEntry { + /// Builds the delivery schedule for `mask` over `world` at `epoch`. + /// + /// The schedule construction runs on the offload pool, because it walks every visible row and + /// would otherwise block the request's async worker. Forking `epoch` gives that work its own + /// owned handle on the same captured publication, independent of the caller's guard. Dropping + /// this future abandons the result but does not cancel schedule work already submitted to + /// Rayon. + /// + /// # Errors + /// + /// Returns [`VisibilityCacheError::Panic`] when the offloaded construction panics or its + /// worker disappears without returning a value. + pub(crate) async fn new( + world: Arc, + epoch: &Epoch, + mask: VisibilityMask, + filter: Option>, + ) -> Result> { + let epoch = epoch.fork(); + let (schedule, occupancy, mask) = offload::run(move || { + let schedule = ViewSchedule::of(world, &epoch, &mask); + let occupancy = schedule.occupancy(); + + (schedule, occupancy, mask) + }) + .await + .change_context(VisibilityCacheError::Panic)?; + + let weight = weight_of( + mask.heap_memory_usage() + schedule.heap_memory_usage(), + filter.as_deref(), + ); + + Ok(Self { + mask, + schedule, + filter, + occupancy, + weight, + }) + } +} + +hashql_core::id::newtype! { + /// A cache insertion token checked before a refresh replaces an entry. + /// + /// A refresh compares the publication it claimed against the one currently held to detect an intervening replacement. + /// + /// # Warning + /// + /// [`PublicationProducer`] uses a u32 counter, which wraps after 2³² allocations across the cache. Equality distinguishes insertions only while their tokens have not repeated. An outstanding refresh can mistake a later insertion with the same token for its claimed entry. + struct Publication(u64) +} + +hashql_core::id::newtype_producer!(struct PublicationProducer(Publication)); + +/// One refresh's exclusive claim on a cache entry. +/// +/// Dropping the guard clears that entry's claim. The refresh future owns the guard and releases +/// it on completion, unwinding or future drop. Later refreshes remain subject to the retired-epoch +/// checks in [`VisibilityCache::resolve`] and to the entry's expiry. +struct RefreshClaim { + entry: Arc, +} + +impl RefreshClaim { + /// Returns the claimed entry's publication. + fn publication(&self) -> Publication { + self.entry.publication + } +} + +impl Drop for RefreshClaim { + fn drop(&mut self) { + self.entry + .refreshing + .store(false, atomic::Ordering::Release); + } +} + +/// One actor's resolved scope, retained for reuse across requests. +/// +/// The entry is shared behind an [`Arc`], letting a request that took it keep reading a +/// coherent scope while a refresh replaces the cache's copy. +#[derive(Debug)] +pub(crate) struct CacheEntry { + /// The rows the actor may receive. + pub mask: VisibilityMask, + /// The delivery schedule built over those rows. + pub schedule: ViewSchedule, + /// The filter document retained from this scope's resolution. + /// + /// A later request naming only its digest can reuse this document. + pub filter: Option>, + /// Distinct occupied-cell counts by Morton depth, used to resolve a density offset. + /// + /// Corpus policy has no occupancy profile. Every scoped schedule has one, including the zero + /// profile of an empty scope. + pub occupancy: Option, + resolved_at: Instant, + publication: Publication, + refreshing: Atomic, + weight: u32, +} + +impl CacheEntry { + /// Stamps a resolved scope with its admission time and publication. + fn new( + PendingCacheEntry { + mask, + schedule, + filter, + occupancy, + weight, + }: PendingCacheEntry, + resolved_at: Instant, + publication: Publication, + ) -> Self { + Self { + mask, + schedule, + filter, + occupancy, + resolved_at, + publication, + refreshing: Atomic::::new(false), + weight, + } + } + + /// Returns whether the entry has reached the `soft` age. + fn is_stale(&self, now: Instant, soft: Duration) -> bool { + now.saturating_duration_since(self.resolved_at) >= soft + } + + /// Returns whether the entry has reached the `hard` age, past which it is no longer served. + fn is_expired(&self, now: Instant, hard: Duration) -> bool { + now.saturating_duration_since(self.resolved_at) >= hard + } + + /// Claims the exclusive right to refresh this entry until the returned guard drops. + /// + /// Returns [`None`] while another refresh holds the claim. A request that still holds an entry + /// a later refresh has already replaced can claim it again. The publication comparison in + /// [`VisibilityCache::resolve`] rejects that redundant resolution's result while + /// [`Publication`] tokens have not repeated. + fn claim_refresh(entry: &Arc) -> Option { + entry + .refreshing + .compare_exchange( + false, + true, + atomic::Ordering::AcqRel, + atomic::Ordering::Acquire, + ) + .is_ok() + .then(|| RefreshClaim { + entry: Arc::clone(entry), + }) + } +} + +/// An actor and filter scope within one generation and [delta lifetime](DeltaId). +/// +/// A generation can reopen with another [`DeltaId`]. Including both values normally separates its +/// cache scopes. Each lifetime draws its own sampled 64-bit tag, but the cache does not detect when +/// separate samples produce equal tags. The key excludes the revision to permit reuse and refresh +/// across publications within a lifetime. Freshness limits rather than revision changes bound scope +/// staleness. +#[derive(Debug, PartialEq, Eq, Hash)] +pub(crate) struct CacheKey { + generation: GenerationId, + delta: DeltaId, + actor: ActorId, + filter: Option, +} + +impl CacheKey { + /// Names the scope `actor` reads `epoch` under, optionally through `filter`. + pub(crate) fn new(epoch: &Epoch, actor: ActorId, filter: Option) -> Self { + Self { + generation: epoch.generation(), + delta: epoch.reference().id, + actor, + filter, + } + } + + /// Returns whether `epoch` carries this key's generation and lifetime tag. + /// + /// A match makes `epoch` eligible for resolution. It cannot distinguish a sampled [`DeltaId`] + /// collision. + fn matches(&self, epoch: &Epoch) -> bool { + self.generation == epoch.generation() && self.delta == epoch.reference().id + } +} + +/// The operator's capacity and freshness limits for resolved visibility scopes. +/// +/// Ages start at the lookup admission time recorded on insertion. `soft` permits one matching +/// lookup to refresh in the background while receiving the held entry. `hard` prevents serving an +/// expired entry: a matching lifetime resolves synchronously, while a non-present lifetime receives +/// no entry. Construction does not require `soft < hard`. When `soft ≥ hard`, there is no interval +/// in which a still-servable entry can start a stale refresh. +#[derive(Debug, Copy, Clone, PartialEq, Eq)] +pub struct VisibilityLimits { + /// The best-effort weighted capacity in bytes, not a strict instantaneous heap ceiling. + pub bytes: u64, + /// The age at which an entry becomes stale and a matching lookup may claim its refresh. + pub soft: Duration, + /// The age at or beyond which serving requires foreground re-resolution. + pub hard: Duration, +} + +/// One process's retained visibility scopes. +#[derive(Debug)] +pub(crate) struct VisibilityCache { + entries: moka::future::Cache>, + publications: Arc, + limits: VisibilityLimits, +} + +impl VisibilityCache { + /// Builds a cache configured with `limits`. + /// + /// Eviction charges each entry by footprint and uses the `tiny_lfu` policy to account for + /// access frequency. Capacity enforcement is best effort. The backing cache also applies a + /// time-to-live of [`VisibilityLimits::hard`] from physical insertion. Explicit age checks use + /// the entry's earlier logical admission time. + /// + /// # Panics + /// + /// Panics when [`VisibilityLimits::hard`] exceeds 31,536,000,000 seconds (1,000 times 365 + /// days), including by a fractional second. + pub(crate) fn new(limits: VisibilityLimits) -> Self { + Self { + entries: moka::future::Cache::builder() + .max_capacity(limits.bytes) + .weigher(|_key, entry: &Arc| entry.weight) + .eviction_policy(moka::policy::EvictionPolicy::tiny_lfu()) + .time_to_live(limits.hard) + .build(), + publications: Arc::new(PublicationProducer::new()), + limits, + } + } + + /// Returns the filter document a retained entry was resolved under. + /// + /// A request may present a filter digest without the document it names, and this is how the + /// resolver recovers the document rather than refusing the request. + pub(super) async fn filter_document(&self, key: &CacheKey) -> Option> { + self.entries.get(key).await?.filter.as_ref().map(Arc::clone) + } + + /// Returns the entry under `key`, resolving it where no live entry exists. + /// + /// The key's compute lock serializes the decision. The cache reuses any held entry younger than + /// [`VisibilityLimits::hard`], including one whose generation or lifetime does not match + /// `epoch`. When the held entry is missing or expired, a matching key resolves and + /// publishes a replacement. A nonmatching key removes only an expired held entry and otherwise + /// remains absent. With a zero maximum age, every held entry is logically expired. + /// + /// # Errors + /// + /// Returns the resolver's error without publishing a replacement. A later request may retry. + async fn get_or_insert_with( + &self, + epoch: &Epoch, + key: CacheKey, + now: Instant, + resolver: R, + ) -> Result>, E> + where + R: AsyncFnOnce(&Epoch) -> Result, + E: Send + Sync + 'static, + { + let eligible = key.matches(epoch); + self.entries + .entry(key) + .and_try_compute_with(async |held| { + if held.is_some_and(|held| !held.value().is_expired(now, self.limits.hard)) { + return Ok::<_, E>(Op::Nop); + } + + if !eligible { + return Ok(Op::Remove); + } + + Ok(Op::Put(Arc::new(CacheEntry::new( + resolver(epoch).await?, + now, + self.publications.next(), + )))) + }) + .await + .map(|result| match result { + CompResult::StillNone(_) | CompResult::Removed(_) => None, + CompResult::Unchanged(entry) + | CompResult::Inserted(entry) + | CompResult::ReplacedWith(entry) => Some(entry.into_value()), + }) + } + + /// Reuses a cached scope or resolves an eligible epoch. + /// + /// The caller's future resolves a missing or expired matching entry. Cancelling that future + /// drops the foreground resolver, although schedule work already submitted to Rayon + /// continues without a receiver. A stale matching entry returns immediately and starts a + /// detached Tokio refresh when its claim is free. The detached task outlives caller + /// cancellation. A nonmatching lookup serves the entry without refresh until its maximum age, + /// then removes it without resolution. + /// + /// The detached task uses the current tracing span. On a returned resolution error, it releases + /// the entry's refresh claim before calling `on_refresh_error`. Unwinding or task cancellation + /// also releases the claim, but does not call the error callback. The held entry keeps its + /// original expiry after any failed refresh. + /// + /// # Errors + /// + /// Returns the resolver's error when a matching entry is missing or expired. + /// + /// # Panics + /// + /// Panics if a stale matching lookup claims a refresh unless the caller has entered a + /// [`Runtime`](tokio::runtime::Runtime). + pub(crate) async fn resolve( + &self, + epoch: &Epoch, + key: CacheKey, + now: Instant, + resolver: R, + on_refresh_error: impl FnOnce(E) + Send + 'static, + ) -> Result>, E> + where + R: for<'epoch> AsyncFnOnce(&'epoch Epoch) -> Result + Send + 'static, + for<'epoch> >::CallOnceFuture: Send, + E: Send + Sync + 'static, + { + let Some(entry) = self.entries.get(&key).await else { + return self.get_or_insert_with(epoch, key, now, resolver).await; + }; + + if entry.is_expired(now, self.limits.hard) { + return self.get_or_insert_with(epoch, key, now, resolver).await; + } + + if !key.matches(epoch) { + return Ok(Some(entry)); + } + + if entry.is_stale(now, self.limits.soft) + && let Some(claim) = CacheEntry::claim_refresh(&entry) + { + let entries = self.entries.clone(); + let publications = Arc::clone(&self.publications); + let epoch = epoch.fork(); + + let _handle = tokio::spawn( + async move { + let resolution = match resolver(&epoch).await { + Ok(resolution) => resolution, + Err(error) => { + drop(claim); + on_refresh_error(error); + return; + } + }; + + let _result = entries + .entry(key) + .and_compute_with(async |held| { + if held + .is_none_or(|held| held.value().publication != claim.publication()) + { + return Op::Nop; + } + + Op::Put(Arc::new(CacheEntry::new( + resolution, + now, + publications.next(), + ))) + }) + .await; + } + .in_current_span(), + ); + } + + Ok(Some(entry)) + } +} diff --git a/libs/@local/graph/atlas/src/serve/visibility/cache/tests.rs b/libs/@local/graph/atlas/src/serve/visibility/cache/tests.rs new file mode 100644 index 00000000000..14b0edf79d9 --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/visibility/cache/tests.rs @@ -0,0 +1,752 @@ +//! Logical expiry and refresh eligibility with backing-cache retention held fixed. +//! +//! The backing cache has no TTL. Backdated entries use a fixed supplied time. + +use alloc::sync::Arc; +use core::{sync::atomic::Ordering, time::Duration}; +use std::{fs, time::Instant}; + +use arc_swap::Guard; +use rand::{SeedableRng as _, rngs::StdRng}; +use tokio::{sync::oneshot, time::timeout}; +use tracing::{Dispatch, Instrument as _}; +use tracing_subscriber::Registry; +use type_system::principal::actor::{ActorId, ActorType}; +use uuid::Uuid; + +use super::{ + CacheEntry, CacheKey, PendingCacheEntry, PublicationProducer, VisibilityCache, + VisibilityLimits, weight_of, +}; +use crate::{ + allocator::HeapMemoryUsage as _, + file::{generation::GenerationId, repository::Artifact as _, salt::artifact}, + serve::{ + delta::{Delta, epoch::Epoch}, + schedule::ViewSchedule, + tests::fixture::{TamperFixture, secret}, + visibility::{VisibilityActor, VisibilityMask}, + world::World, + }, +}; + +/// The hard-expiry threshold and expired-entry age. +const HARD: Duration = Duration::from_secs(60); + +/// Synthetic serving artifacts paired with a cache whose backing entries do not expire. +struct Fixture { + files: TamperFixture, + world: Arc, + epoch: Epoch, + retired: GenerationId, + cache: VisibilityCache, + actor: ActorId, + now: Instant, +} + +impl Fixture { + /// Builds a visibility-cache fixture at a fixed logical time. + /// + /// Its backing entries never expire. The fixture includes an active lifetime, an unrelated + /// retired generation, one actor and serving artifacts published under `name`. + /// + /// # Panics + /// + /// Panics when the synthetic generation cannot be published or opened. + fn new(name: &str) -> Self { + let files = TamperFixture::publish(name); + let world = Arc::new( + World::open(files.generation().clone(), &secret()) + .expect("the synthetic generation should open"), + ); + let delta = Delta::new(Arc::clone(&world), StdRng::seed_from_u64(0x5EED)) + .expect("the seeded RNG should allocate a delta identity"); + let epoch = Epoch::from(Guard::from_inner(Arc::new(delta))); + let retired: GenerationId = "ab" + .repeat(32) + .parse() + .expect("64 hexadecimal digits should name a generation"); + assert_ne!(epoch.generation(), retired); + + Self { + files, + world, + epoch, + retired, + cache: VisibilityCache { + entries: moka::future::Cache::builder().build(), + publications: Arc::new(PublicationProducer::new()), + limits: VisibilityLimits { + bytes: u64::MAX, + soft: Duration::from_secs(30), + hard: HARD, + }, + }, + actor: ActorId::new(Uuid::from_u128(1), ActorType::User), + now: Instant::now(), + } + } + + /// Returns this fixture's cache key with `generation` substituted. + fn key(&self, generation: GenerationId) -> CacheKey { + CacheKey { + generation, + ..CacheKey::new(&self.epoch, self.actor, None) + } + } + + /// Asserts the backing cache holds no entry under `key`. + /// + /// # Panics + /// + /// Panics when an entry is present. + async fn absent(&self, key: &CacheKey) { + assert!( + self.cache.entries.get(key).await.is_none(), + "the key should be absent" + ); + } + + /// Inserts an entry under `key`, backdated by `age`, and returns it. + /// + /// # Panics + /// + /// Panics when `age` predates the instant range the fixture's clock reading allows. + async fn seed(&self, key: CacheKey, age: Duration) -> Arc { + let resolved_at = self + .now + .checked_sub(age) + .expect("the age should fit the instant's range"); + let entry = Arc::new(CacheEntry::new( + pending(&self.world, &self.epoch, self.actor), + resolved_at, + self.cache.publications.next(), + )); + self.cache.entries.insert(key, Arc::clone(&entry)).await; + entry + } +} + +/// Constructs a Corpus resolution without the async scheduling offload. +fn pending(world: &Arc, epoch: &Epoch, actor: ActorId) -> PendingCacheEntry { + let mask = VisibilityMask::full(VisibilityActor { + id: actor, + instance_admin: false, + }); + let schedule = ViewSchedule::of(Arc::clone(world), epoch, &mask); + let weight = weight_of( + mask.heap_memory_usage() + schedule.heap_memory_usage(), + None, + ); + + PendingCacheEntry { + mask, + schedule, + filter: None, + occupancy: None, + weight, + } +} + +/// A resolver over the fixture's world and actor for any supplied lifetime. +/// +/// A macro rather than a function because the closure captures the fixture's borrows by +/// value and each caller needs its own. +macro_rules! resolving { + ($fixture:expr) => {{ + let world = Arc::clone(&$fixture.world); + let actor = $fixture.actor; + async move |epoch: &Epoch| Ok::<_, ()>(pending(&world, epoch, actor)) + }}; +} + +/// Guards cache-only paths with a resolver that must never run. +/// +/// # Panics +/// +/// Always, because the cache alone must answer these lookups. +/// +/// # Errors +/// +/// Never: the resolution this stands in for is one the lookup must not reach. +async fn forbidden(_epoch: &Epoch) -> Result { + panic!("this lookup should not invoke the resolver") +} + +/// Fails every resolution attempt with a fixed error. +/// +/// # Errors +/// +/// Always, with the message the case asserts on. +async fn refusing(_epoch: &Epoch) -> Result { + Err("permission resolution failed") +} + +/// Panics instead of returning a resolution. +/// +/// # Panics +/// +/// Always, when invoked. +/// +/// # Errors +/// +/// Never: it unwinds instead of returning. +#[expect( + clippy::panic_in_result_fn, + reason = "this resolver exercises cleanup after a task panic" +)] +fn panicking(_epoch: &Epoch) -> Result { + panic!("a refresh resolver panic should unwind out of its task") +} + +/// Builds a delta lifetime over a tampered generation distinct from the fixture's. +/// +/// # Panics +/// +/// Panics when the placeholder cannot be rewritten or the tampered generation cannot +/// be opened. +fn other_epoch(files: &TamperFixture) -> Epoch { + let generation = files.tamper(&artifact::Representations::NAME, |path| { + fs::remove_file(path).expect("the staged placeholder should be removable"); + fs::write(path, b"alternate representation placeholder") + .expect("the alternate placeholder should write"); + }); + let world = + Arc::new(World::open(generation, &secret()).expect("the other generation should open")); + let delta = Delta::new(world, StdRng::seed_from_u64(0xBEEF)) + .expect("the seeded RNG should allocate the other delta identity"); + Epoch::from(Guard::from_inner(Arc::new(delta))) +} + +/// Reopens the fixture's generation as another world and delta lifetime. +/// +/// # Panics +/// +/// Panics when the generation cannot be reopened. +fn reopened(files: &TamperFixture) -> (Arc, Epoch) { + let world = Arc::new( + World::open(files.generation().clone(), &secret()) + .expect("the same generation should reopen"), + ); + let delta = Delta::new(Arc::clone(&world), StdRng::seed_from_u64(0xA11CE)) + .expect("the seeded RNG should allocate a new delta identity"); + (world, Epoch::from(Guard::from_inner(Arc::new(delta)))) +} + +/// Separates in-flight resolutions by their full lifetime keys. +/// +/// A resolution admitted before a same-generation reopen remains under its original key. +#[tokio::test(start_paused = true)] +async fn resolve_reopened_inflight() { + let fixture = Fixture::new("cache-resolve-reopened-inflight"); + let (world, epoch) = reopened(&fixture.files); + assert_eq!(epoch.generation(), fixture.epoch.generation()); + assert_ne!(epoch.reference().id, fixture.epoch.reference().id); + assert!(!Arc::ptr_eq(&world, &fixture.world)); + + let (announce, started) = oneshot::channel(); + let (release, released) = oneshot::channel(); + let previous_world = Arc::clone(&fixture.world); + let actor = fixture.actor; + let old = fixture.cache.resolve( + &fixture.epoch, + CacheKey::new(&fixture.epoch, actor, None), + fixture.now, + async move |epoch: &Epoch| { + announce + .send(()) + .expect("the replacement should await resolution startup"); + released + .await + .expect("the replacement should release the prior resolution"); + Ok::<_, ()>(pending(&previous_world, epoch, actor)) + }, + |()| panic!("should return foreground errors to the caller"), + ); + let new = async { + started + .await + .expect("the prior resolution should announce startup"); + let entry = fixture + .cache + .resolve( + &epoch, + CacheKey::new(&epoch, actor, None), + fixture.now, + async move |epoch: &Epoch| Ok::<_, ()>(pending(&world, epoch, actor)), + |()| panic!("should return foreground errors to the caller"), + ) + .await + .expect("the new lifetime should resolve independently") + .expect("the new lifetime should have a publication"); + release + .send(()) + .expect("the prior resolution should still await release"); + entry + }; + let (old, new) = timeout(Duration::from_secs(1), async { tokio::join!(old, new) }) + .await + .expect("independent lifetime keys should not block each other's resolution"); + let old = old + .expect("the prior admitted resolution should finish") + .expect("the prior resolution should return its publication"); + assert!(!Arc::ptr_eq(&old, &new)); + for (epoch, expected) in [(&fixture.epoch, old), (&epoch, new)] { + let held = fixture + .cache + .entries + .get(&CacheKey::new(epoch, actor, None)) + .await + .expect("each lifetime should retain its own publication"); + assert!(Arc::ptr_eq(&held, &expected)); + assert_eq!(held.resolved_at, fixture.now); + } +} + +/// Keeps a previous delta lifetime retired after its generation reopens. +/// +/// Soft-stale, hard-expired and missing entries under the previous key do not resolve again. +#[tokio::test] +async fn resolve_reopened_retired() { + let fixture = Fixture::new("cache-resolve-reopened-retired"); + let (_world, epoch) = reopened(&fixture.files); + assert_eq!(epoch.generation(), fixture.epoch.generation()); + assert_ne!(epoch.reference().id, fixture.epoch.reference().id); + let key = || CacheKey::new(&fixture.epoch, fixture.actor, None); + let held = fixture.seed(key(), fixture.cache.limits.soft).await; + let answer = fixture + .cache + .resolve(&epoch, key(), fixture.now, forbidden, |()| { + panic!("should not refresh a retired entry") + }) + .await + .expect("a previous lifetime lookup should not fail") + .expect("the unexpired publication should remain usable"); + assert!(Arc::ptr_eq(&answer, &held)); + assert!(!held.refreshing.load(Ordering::Acquire)); + + fixture.seed(key(), HARD).await; + let expired = fixture + .cache + .resolve(&epoch, key(), fixture.now, forbidden, |()| { + panic!("should not refresh a retired entry") + }) + .await + .expect("the expired lifetime should not resolve again"); + assert!(expired.is_none()); + fixture.absent(&key()).await; + let missing = fixture + .cache + .resolve(&epoch, key(), fixture.now, forbidden, |()| { + panic!("should not refresh a retired entry") + }) + .await + .expect("a previous lifetime miss should not resolve again"); + assert!(missing.is_none()); +} + +/// Reuses an unexpired retired entry without refreshing it. +/// +/// The current lifetime can neither invoke the resolver nor take a refresh claim for the retired +/// generation. +#[tokio::test] +async fn resolve_stale_retired() { + let fixture = Fixture::new("cache-resolve-stale-retired"); + let key = || fixture.key(fixture.retired); + let held = fixture.seed(key(), fixture.cache.limits.soft).await; + let answer = fixture + .cache + .resolve(&fixture.epoch, key(), fixture.now, forbidden, |()| { + panic!("should not refresh a retired entry") + }) + .await + .expect("the retired lookup should not fail") + .expect("the unexpired entry should remain usable"); + + assert!( + Arc::ptr_eq(&answer, &held), + "the lookup should preserve its publication" + ); + assert!( + !held.refreshing.load(Ordering::Acquire), + "an ineligible lookup should leave refresh unclaimed" + ); +} + +/// Refreshes a stale entry only after its lifetime becomes active again. +/// +/// A stale entry remains reusable without refresh under a foreign lifetime. Returning to the +/// entry's own unchanged delta lifetime returns the entry and starts a refresh. A failed refresh +/// preserves the immediate answer, reports the error to the callback and releases the claim. +#[tokio::test(start_paused = true)] +async fn resolve_reactivated() { + let fixture = Fixture::new("cache-resolve-reactivated"); + let other = other_epoch(&fixture.files); + assert_ne!( + other.generation(), + fixture.epoch.generation(), + "the epochs should name different generations" + ); + let key = || fixture.key(fixture.epoch.generation()); + let held = fixture.seed(key(), fixture.cache.limits.soft).await; + + let retired = fixture + .cache + .resolve(&other, key(), fixture.now, forbidden, |()| { + panic!("should not refresh a retired entry") + }) + .await + .expect("the retired lookup should not fail") + .expect("the unexpired entry should remain usable under the other epoch"); + assert!( + Arc::ptr_eq(&retired, &held), + "the other generation should reuse the held publication" + ); + + let (announce, finished) = oneshot::channel(); + let active = fixture + .cache + .resolve(&fixture.epoch, key(), fixture.now, refusing, move |error| { + announce + .send(error) + .expect("should await the refresh failure"); + }) + .await + .expect("the active lookup should not fail") + .expect("the active lookup should return the held entry"); + assert!( + Arc::ptr_eq(&active, &held), + "a refresh should preserve the immediate answer" + ); + assert_eq!( + timeout(Duration::from_secs(1), finished) + .await + .expect("should finish the refresh") + .expect("should report the refresh failure"), + "permission resolution failed" + ); + assert!( + !held.refreshing.load(Ordering::Acquire), + "a failed refresh should release its claim" + ); +} + +/// Retains stale data after a detached refresh fails. +/// +/// The error callback runs after releasing the refresh claim and under the scheduling span. +#[tokio::test(start_paused = true)] +async fn resolve_refresh_failure() { + let fixture = Fixture::new("cache-resolve-refresh-failure"); + let key = || fixture.key(fixture.epoch.generation()); + let held = fixture.seed(key(), fixture.cache.limits.soft).await; + let dispatch = Dispatch::new(Registry::default()); + let _default = tracing::dispatcher::set_default(&dispatch); + + let requests = [ + tracing::info_span!("request"), + tracing::info_span!("request"), + tracing::Span::none(), + ]; + assert_ne!(requests[0].id(), requests[1].id()); + for request in requests { + let expected_id = request.id(); + let (announce, finished) = oneshot::channel(); + let claimed = Arc::clone(&held); + let error = vec![17_u8, 29]; + let answer = fixture + .cache + .resolve( + &fixture.epoch, + key(), + fixture.now, + async move |_: &Epoch| Err(error), + move |error| { + announce + .send(( + error, + claimed.refreshing.load(Ordering::Acquire), + tracing::Span::current().id(), + )) + .expect("should await the error callback"); + }, + ) + .instrument(request) + .await + .expect("should answer from the stale entry") + .expect("should retain the stale entry"); + assert!(Arc::ptr_eq(&answer, &held)); + + let (error, refreshing, actual_id) = timeout(Duration::from_secs(1), finished) + .await + .expect("should complete the refresh") + .expect("should call the error observer"); + assert_eq!(error, [17, 29]); + assert!( + !refreshing, + "should release the claim before reporting the failure" + ); + assert_eq!(actual_id, expected_id); + assert_eq!(tracing::Span::current().id(), None); + let retained = fixture + .cache + .entries + .get(&key()) + .await + .expect("should retain the failed refresh's publication"); + assert!(Arc::ptr_eq(&retained, &held)); + } +} + +/// Releases a refresh claim after resolver panic and admits the next attempt. +#[tokio::test] +async fn resolve_refresh_panic() { + let fixture = Fixture::new("cache-resolve-refresh-panic"); + let key = || fixture.key(fixture.epoch.generation()); + let held = fixture.seed(key(), fixture.cache.limits.soft).await; + + let (announce, started) = oneshot::channel(); + let resolver = async move |epoch: &Epoch| { + announce + .send(()) + .expect("the case should await refresh startup"); + panicking(epoch) + }; + let answer = fixture + .cache + .resolve(&fixture.epoch, key(), fixture.now, resolver, |()| { + panic!("should not report an unwinding panic as a resolver error") + }) + .await + .expect("the stale lookup should not fail") + .expect("the stale lookup should return the held entry"); + + assert!( + Arc::ptr_eq(&answer, &held), + "a refresh should preserve the immediate answer" + ); + // The refresh runs on this case's current-thread scheduler, and its announcement precedes + // its panic within one poll. Awaiting the announcement returns after the task has unwound + // and released the claim it holds. + timeout(Duration::from_secs(1), started) + .await + .expect("refresh startup should not stall") + .expect("the stale lookup should launch the eligible refresh"); + assert!( + !held.refreshing.load(Ordering::Acquire), + "a panicking refresh should release its claim" + ); + + let (announce, finished) = oneshot::channel(); + let later = fixture + .cache + .resolve(&fixture.epoch, key(), fixture.now, refusing, move |error| { + announce + .send(error) + .expect("should await the later refresh failure"); + }) + .await + .expect("the later lookup should not fail") + .expect("the later lookup should return the held entry"); + + assert!( + Arc::ptr_eq(&later, &held), + "the later lookup should answer from the same publication" + ); + assert_eq!( + timeout(Duration::from_secs(1), finished) + .await + .expect("should finish the later refresh") + .expect("should report the later refresh failure"), + "permission resolution failed" + ); +} + +/// Removes a retired entry at `HARD` age and returns no value. +#[tokio::test] +async fn resolve_expired_retired() { + let fixture = Fixture::new("cache-resolve-expired-retired"); + let key = || fixture.key(fixture.retired); + fixture.absent(&key()).await; + fixture.seed(key(), HARD).await; + + let answer = fixture + .cache + .resolve(&fixture.epoch, key(), fixture.now, forbidden, |()| { + panic!("should not refresh a retired entry") + }) + .await + .expect("a retired lookup should not fail"); + + assert!( + answer.is_none(), + "a retired expired entry should not escape the cache" + ); + fixture.absent(&key()).await; +} + +/// Reuses a soft-stale entry one nanosecond before expiration. +#[tokio::test] +async fn compute_unexpired_active() { + let fixture = Fixture::new("cache-compute-unexpired-active"); + let key = || fixture.key(fixture.epoch.generation()); + fixture.absent(&key()).await; + let age = HARD + .checked_sub(Duration::from_nanos(1)) + .expect("the expiry interval should exceed one nanosecond"); + let held = fixture.seed(key(), age).await; + + let answer = fixture + .cache + .get_or_insert_with(&fixture.epoch, key(), fixture.now, forbidden) + .await + .expect("an unexpired lookup should not fail") + .expect("the held entry should remain usable"); + + assert!( + Arc::ptr_eq(&answer, &held), + "reuse should preserve the held publication" + ); +} + +/// Reuses an unexpired entry under a different generation's epoch. +#[tokio::test] +async fn compute_unexpired_retired() { + let fixture = Fixture::new("cache-compute-unexpired-retired"); + let key = || fixture.key(fixture.retired); + fixture.absent(&key()).await; + let age = HARD + .checked_sub(Duration::from_nanos(1)) + .expect("the expiry interval should exceed one nanosecond"); + let held = fixture.seed(key(), age).await; + + let answer = fixture + .cache + .get_or_insert_with(&fixture.epoch, key(), fixture.now, forbidden) + .await + .expect("an unexpired lookup should not fail") + .expect("the held entry should remain usable before expiration"); + + assert!( + Arc::ptr_eq(&answer, &held), + "reuse should preserve the held publication" + ); +} + +/// Resolves an active-generation miss and retains the supplied timestamp. +#[tokio::test] +async fn resolve_missing_active() { + let fixture = Fixture::new("cache-resolve-missing-active"); + let key = || fixture.key(fixture.epoch.generation()); + fixture.absent(&key()).await; + + let answer = fixture + .cache + .resolve( + &fixture.epoch, + key(), + fixture.now, + resolving!(fixture), + |()| panic!("should return foreground errors to the caller"), + ) + .await + .expect("an active miss should not fail") + .expect("an active miss should resolve a new entry"); + + assert_eq!(answer.resolved_at, fixture.now); + let retained = fixture + .cache + .entries + .get(&key()) + .await + .expect("the new entry should remain cached"); + assert!(Arc::ptr_eq(&answer, &retained)); +} + +/// Skips resolution and returns no entry for a retired-generation miss. +#[tokio::test] +async fn resolve_missing_retired() { + let fixture = Fixture::new("cache-resolve-missing-retired"); + let key = || fixture.key(fixture.retired); + fixture.absent(&key()).await; + + let answer = fixture + .cache + .resolve(&fixture.epoch, key(), fixture.now, forbidden, |()| { + panic!("should not refresh a retired entry") + }) + .await + .expect("a retired miss should not fail"); + + assert!(answer.is_none(), "a retired miss should remain absent"); + fixture.absent(&key()).await; +} + +/// Propagates a foreground resolution error instead of returning an expired fallback. +#[tokio::test] +async fn resolve_resolver_failure() { + let fixture = Fixture::new("cache-resolve-resolver-failure"); + let key = || fixture.key(fixture.epoch.generation()); + fixture.absent(&key()).await; + fixture.seed(key(), HARD).await; + + let error = fixture + .cache + .resolve(&fixture.epoch, key(), fixture.now, refusing, |_| { + panic!("should return foreground errors to the caller") + }) + .await + .expect_err("resolution failure should propagate without an expired fallback"); + + assert_eq!(error, "permission resolution failed"); +} + +/// Replaces an expired active-generation entry with a new resolution. +#[tokio::test] +async fn compute_expired_active() { + let fixture = Fixture::new("cache-compute-expired-active"); + let key = || fixture.key(fixture.epoch.generation()); + fixture.absent(&key()).await; + let expired = fixture.seed(key(), HARD).await; + + let answer = fixture + .cache + .get_or_insert_with(&fixture.epoch, key(), fixture.now, resolving!(fixture)) + .await + .expect("an active compute should not fail") + .expect("the expired entry should yield a new resolution"); + + assert!( + !Arc::ptr_eq(&answer, &expired), + "resolution should replace the expired publication" + ); + assert_ne!(answer.publication, expired.publication); + assert_eq!(answer.resolved_at, fixture.now); + let retained = fixture + .cache + .entries + .get(&key()) + .await + .expect("the replacement should remain cached"); + assert!(Arc::ptr_eq(&answer, &retained)); +} + +/// Removes an expired retired-generation entry without resolving a replacement. +#[tokio::test] +async fn compute_expired_retired() { + let fixture = Fixture::new("cache-compute-expired-retired"); + let key = || fixture.key(fixture.retired); + fixture.absent(&key()).await; + fixture.seed(key(), HARD).await; + + let answer = fixture + .cache + .get_or_insert_with(&fixture.epoch, key(), fixture.now, forbidden) + .await + .expect("a retired compute should not fail"); + + assert!( + answer.is_none(), + "an expired entry should not escape the retired-generation compute" + ); + fixture.absent(&key()).await; +} diff --git a/libs/@local/graph/atlas/src/serve/visibility/mod.rs b/libs/@local/graph/atlas/src/serve/visibility/mod.rs index 4db5518f72f..40253ecd1a2 100644 --- a/libs/@local/graph/atlas/src/serve/visibility/mod.rs +++ b/libs/@local/graph/atlas/src/serve/visibility/mod.rs @@ -5,7 +5,7 @@ //! the node and edge rows the actor may receive, together with the principal and instance- //! administrator status used for resolution. Filter construction receives property-protection //! configuration as a separate input. The mask freezes the resulting permission-store decision. -//! `cache` defines when requests reuse, refresh or refuse it. +//! [`cache`] defines when requests reuse, refresh or refuse it. //! //! A mask is either the whole corpus or an explicit row set, and //! [`VisibilityKind`] names which. That distinction is the actor's admitted policy rather than an @@ -24,6 +24,9 @@ use crate::{ identity::{EdgeRowId, NodeRowId}, }; +pub(crate) mod cache; +pub(crate) mod resolver; + /// A visible row set over one row domain. #[derive(Debug, Clone, PartialEq, Eq)] enum Rows { diff --git a/libs/@local/graph/atlas/src/serve/visibility/resolver/mod.rs b/libs/@local/graph/atlas/src/serve/visibility/resolver/mod.rs new file mode 100644 index 00000000000..3dac5c9e2bd --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/visibility/resolver/mod.rs @@ -0,0 +1,129 @@ +//! Store-backed visibility resolution for a request's captured generations. +//! +//! A shared [`VisibilityCache`] reuses retained scopes while new resolutions use the present +//! [`World`](crate::serve::world::World) and its paired +//! [`Epoch`]. + +use alloc::sync::Arc; + +use error_stack::{Report, ResultExt as _}; +use hash_graph_postgres_store::store::PostgresStorePool; +use hash_graph_store::{filter::Filter, pool::StorePool as _}; +use serde_json::value::RawValue; +use type_system::{knowledge::Entity, principal::actor::ActorId}; + +use super::cache::{ + CacheEntry, CacheKey, FilterDigest, PendingCacheEntry, VisibilityCache, VisibilityLimits, +}; +use crate::serve::{ + delta::epoch::Epoch, + hydrate::visibility::{VisibilityProofError, visibility_proof}, + runtime::registry::Observation, +}; + +#[cfg(test)] +mod tests; + +/// One process's visibility cache and permission store. +pub(crate) struct ScopeResolver { + pool: Arc, + cache: VisibilityCache, +} + +impl ScopeResolver { + /// Builds a resolver over `pool`, retaining resolved scopes within `limits`. + /// + /// # Panics + /// + /// Panics when [`VisibilityLimits::hard`] exceeds 31,536,000,000 seconds (1,000 times 365 + /// days), including by a fractional second. + pub(crate) fn new(pool: Arc, limits: VisibilityLimits) -> Self { + Self { + pool, + cache: VisibilityCache::new(limits), + } + } + + /// Resolves the requested scope using the observation's admission time. + /// + /// `filter` is the caller-validated cache identity. A supplied `document` must come from the + /// same submitted representation as that digest. This method neither validates the pairing nor + /// recomputes the digest. [`RawValue`] may omit surrounding whitespace that the digest + /// includes. `document` must be absent when `filter` is absent. Without a document, a filtered + /// request can proceed only while its keyed cache entry still retains one. + /// + /// Returns [`None`] when a filtered scope has no supplied or cached document, or when the + /// requested [delta lifetime](crate::serve::delta::DeltaId) has no reusable entry and is not + /// the observation's present lifetime. Registry retention alone does not guarantee that a + /// retired generation remains resolvable after its visibility entry reaches the maximum cache + /// age. A resolution admitted while a lifetime is present may finish after promotion and + /// remains keyed to the originally observed lifetime. + /// + /// # Errors + /// + /// Returns [`VisibilityProofError`] when the retained filter document does not parse, + /// store-backed visibility resolution fails, or schedule construction fails. + /// + /// # Panics + /// + /// Panics if a stale matching entry claims a refresh unless the caller has entered a + /// [`Runtime`](tokio::runtime::Runtime). + pub(crate) async fn resolve( + &self, + observation: &Observation, + actor: ActorId, + filter: Option, + document: Option>, + ) -> Result>, Report> { + let key = CacheKey::new(observation.requested().epoch(), actor, filter); + let document = match (filter, document) { + (Some(_), None) => { + let Some(document) = self.cache.filter_document(&key).await else { + return Ok(None); + }; + Some(document) + } + (_, document) => document, + }; + + let pool = Arc::clone(&self.pool); + let world = Arc::clone(observation.present().world()); + self.cache + .resolve( + observation.present().epoch(), + key, + observation.admitted_at(), + async move |epoch: &Epoch| { + let store = pool + .acquire(None) + .await + .change_context(VisibilityProofError::Connect)?; + + let filter = document + .as_deref() + .map(|document| serde_json::from_str::>(document.get())) + .transpose() + .change_context(VisibilityProofError::Document)?; + + let mask = visibility_proof( + &world, + epoch, + actor, + filter.as_ref(), + &store.settings.filter_protection, + &store, + ) + .await?; + drop(store); + + PendingCacheEntry::new(world, epoch, mask, document) + .await + .change_context(VisibilityProofError::ComputeView) + }, + |error: Report| { + tracing::warn!(?error, "failed to refresh the cached visibility scope"); + }, + ) + .await + } +} diff --git a/libs/@local/graph/atlas/src/serve/visibility/resolver/tests.rs b/libs/@local/graph/atlas/src/serve/visibility/resolver/tests.rs new file mode 100644 index 00000000000..bd658acacdf --- /dev/null +++ b/libs/@local/graph/atlas/src/serve/visibility/resolver/tests.rs @@ -0,0 +1,822 @@ +//! Cache reuse and miss eligibility for captured requests. +//! +//! A missing PostgreSQL socket makes attempted store resolution observable as a connection error. + +use alloc::{collections::BTreeMap, sync::Arc}; +use core::{ + assert_matches, fmt, + future::{Future, poll_fn}, + pin::{Pin, pin}, + task::Poll, + time::Duration, +}; +use std::{fs, time::Instant}; + +use hash_graph_postgres_store::store::{ + DatabaseConnectionInfo, DatabasePoolConfig, DatabaseType, PostgresStorePool, + PostgresStoreSettings, +}; +use hash_graph_store::pool::StorePool as _; +use serde_json::value::RawValue; +use tokio::{sync::mpsc, time::timeout}; +use tokio_postgres::NoTls; +use tokio_util::sync::CancellationToken; +use tracing::{ + Dispatch, Event, Instrument as _, Level, Subscriber, + field::{Field, Visit}, + span::Id, +}; +use tracing_subscriber::{ + Layer, Registry, + layer::{Context, SubscriberExt as _}, + registry::LookupSpan, +}; +use type_system::principal::actor::{ActorId, ActorType}; +use uuid::Uuid; + +use super::ScopeResolver; +use crate::{ + file::{ + generation::{Generation, GenerationRoot}, + repository::Artifact as _, + salt::artifact, + }, + math::nz, + serve::{ + delta::epoch::Epoch, + hydrate::visibility::VisibilityProofError, + runtime::{ + manager::{GenerationManager, ManagerOptions, source::RuntimeSource}, + registry::ObserveError, + }, + tests::fixture::{TamperFixture, secret}, + visibility::{ + VisibilityActor, VisibilityMask, + cache::{ + CacheEntry, CacheKey, FilterDigest, PendingCacheEntry, VisibilityCache, + VisibilityLimits, + }, + }, + world::World, + }, +}; + +/// One event's fields, each rendered to the string its debug shape prints. +struct Fields(BTreeMap); + +impl Visit for Fields { + fn record_debug(&mut self, field: &Field, value: &dyn fmt::Debug) { + self.0.insert(field.name().to_owned(), format!("{value:?}")); + } +} + +/// Removes terminal styling from a recorded field before text assertions. +/// +/// An escape alone introduces a two-character sequence. An escape followed by `[` introduces +/// a control sequence whose parameters run up to its final byte. +fn plain(rendered: &str) -> String { + let mut text = String::with_capacity(rendered.len()); + let mut characters = rendered.chars(); + while let Some(character) = characters.next() { + if character != '\u{1b}' { + text.push(character); + continue; + } + if characters.next() == Some('[') { + for parameter in characters.by_ref() { + if ('@'..='~').contains(¶meter) { + break; + } + } + } + } + + text +} + +/// A subscriber layer forwarding resolver diagnostics to a test channel. +/// +/// Each forwarded event preserves its level, span scope and fields. +struct Diagnostics(mpsc::UnboundedSender<(Level, Vec, BTreeMap)>); + +impl LookupSpan<'lookup>> Layer for Diagnostics { + /// Forwards the event, keeping only events whose target is the resolver module. + /// + /// The span scope is collected from the root outward, which is what lets a case assert + /// the diagnostic was recorded under the request span rather than beside it. + /// + /// # Panics + /// + /// Panics when the receiving end has been dropped, which would otherwise lose the + /// event a case is waiting for. + fn on_event(&self, event: &Event<'_>, context: Context<'_, S>) { + if event.metadata().target() != "hash_graph_atlas::serve::visibility::resolver" { + return; + } + let mut fields = Fields(BTreeMap::new()); + event.record(&mut fields); + let scope = context.event_scope(event).map_or_else(Vec::new, |scope| { + scope.from_root().map(|span| span.id()).collect() + }); + self.0 + .send((*event.metadata().level(), scope, fields.0)) + .expect("should retain the diagnostic receiver"); + } +} + +/// How often the manager under test looks for a newly activated generation. +const POLL_INTERVAL: Duration = Duration::from_millis(5); +/// The wall-clock budget before an awaited step is considered stalled. +const BUDGET: Duration = Duration::from_secs(5); +/// Keeps a replaced generation observable throughout each replacement test. +const RETENTION: Duration = Duration::from_secs(30); +/// A retention interval that expires the replaced generation during a test. +const SHORT_RETENTION: Duration = Duration::from_millis(80); +/// Cache limits under which only the soft and hard ages decide a lookup, with no byte ceiling. +const LIMITS: VisibilityLimits = VisibilityLimits { + bytes: u64::MAX, + soft: Duration::from_secs(30), + hard: Duration::from_secs(60), +}; + +/// Returns the root containing the fixture's generation. +/// +/// # Panics +/// +/// Panics when the fixture's generation has no parent directory or the root refuses to +/// open. +fn root_of(files: &TamperFixture) -> GenerationRoot { + GenerationRoot::new( + files + .generation() + .path() + .parent() + .expect("the fixture generation has a root"), + ) + .expect("the fixture root should open") +} + +/// Builds a second generation differing only by `marker` in one artifact. +/// +/// Activating it changes the generation. +/// +/// # Panics +/// +/// Panics when the placeholder cannot be rewritten. +fn variant_of(files: &TamperFixture, marker: &[u8]) -> Generation { + files.tamper(&artifact::Representations::NAME, |path| { + fs::remove_file(path).expect("the staged placeholder should be removable"); + fs::write(path, marker).expect("the variant placeholder should write"); + }) +} + +/// Builds a store pool whose connection attempts fail. +/// +/// A store pool over a socket that does not exist, which makes any attempt to resolve +/// from the store observable as a connection failure. +/// +/// # Panics +/// +/// Panics when the pool cannot be constructed. +async fn pool() -> Arc { + Arc::new( + PostgresStorePool::new( + &DatabaseConnectionInfo::new( + DatabaseType::Postgres, + "resolver-test".to_owned(), + String::new(), + "/no-resolver-test-postgres".to_owned(), + 5432, + "resolver-test".to_owned(), + ), + &DatabasePoolConfig { + max_connections: nz!(1), + }, + NoTls, + PostgresStoreSettings::default(), + ) + .await + .expect("an unconnected pool should construct"), + ) +} + +/// Builds and activates a manager fixture with `retention`. +/// +/// Returns the serving-artifact fixture, its pool and the initialized manager. +/// +/// # Panics +/// +/// Panics when the generation cannot be activated or the manager refuses its options. +async fn boot( + name: &str, + retention: Duration, +) -> (TamperFixture, Arc, GenerationManager) { + let files = TamperFixture::publish(name); + root_of(&files) + .activate(files.generation().id()) + .expect("the fixture generation should activate"); + let pool = pool().await; + let source = RuntimeSource { + root: root_of(&files), + secret: secret(), + pool: Arc::clone(&pool), + feed: None, + }; + let manager = GenerationManager::new( + source, + ManagerOptions { + poll_interval: POLL_INTERVAL, + .. + }, + retention, + ) + .expect("a non-zero interval should construct a manager"); + + (files, pool, manager) +} + +// Paused Tokio time can exhaust the timeout while Rayon is still opening a generation. +/// Drives the manager's run loop until `probe` answers, and returns that answer. +/// +/// # Panics +/// +/// Panics when the run loop finishes before its shutdown signal, and when the probe +/// has not answered inside the budget. +async fn advance( + mut running: Pin<&mut impl Future>, + mut probe: impl FnMut() -> Option, +) -> T { + timeout( + BUDGET, + poll_fn(|context| { + assert!( + running.as_mut().poll(context).is_pending(), + "the run loop should not finish before its shutdown signal" + ); + probe().map_or(Poll::Pending, Poll::Ready) + }), + ) + .await + .expect("the maintenance pass should not stall") +} + +/// Returns a user actor identified by `id`. +fn actor_of(id: u128) -> ActorId { + ActorId::new(Uuid::from_u128(id), ActorType::User) +} + +/// Seeds and returns the unfiltered publication cached for `actor` at `now`. +/// +/// A fresh same-key hit before any refresh or replacement is eligible must return this pointer. +/// +/// # Panics +/// +/// Panics when the resolution fails or caches nothing. +async fn seed_full( + cache: &VisibilityCache, + world: Arc, + epoch: &Epoch, + now: Instant, + actor: ActorId, +) -> Arc { + let key = CacheKey::new(epoch, actor, None); + cache + .resolve( + epoch, + key, + now, + async move |epoch: &Epoch| { + PendingCacheEntry::new( + world, + epoch, + VisibilityMask::full(VisibilityActor { + id: actor, + instance_admin: false, + }), + None, + ) + .await + }, + |_| panic!("should return foreground errors to the caller"), + ) + .await + .expect("seeding an eligible key should not fail") + .expect("an eligible key should seed an entry") +} + +/// Seeds and returns the publication cached for `actor` under `digest`, retaining `document`. +/// +/// # Panics +/// +/// Panics when the resolution fails or caches nothing. +async fn seed_filtered( + cache: &VisibilityCache, + world: Arc, + epoch: &Epoch, + now: Instant, + actor: ActorId, + digest: FilterDigest, + document: Arc, +) -> Arc { + let key = CacheKey::new(epoch, actor, Some(digest)); + cache + .resolve( + epoch, + key, + now, + async move |epoch: &Epoch| { + PendingCacheEntry::new( + world, + epoch, + VisibilityMask::full(VisibilityActor { + id: actor, + instance_admin: false, + }), + Some(document), + ) + .await + }, + |_| panic!("should return foreground errors to the caller"), + ) + .await + .expect("seeding an eligible key should not fail") + .expect("an eligible key should seed an entry") +} + +/// Reuses a fresh unfiltered publication without store access. +/// +/// The unconnectable pool proves that the hit performs no store query. +#[tokio::test] +async fn resolve_fresh_unfiltered_hit() { + let (_files, pool, mut manager) = boot("resolver-fresh-unfiltered-hit", RETENTION).await; + let registry = Arc::clone(manager.registry()); + let shutdown = CancellationToken::new(); + + let (seeded, answer) = { + let mut running = pin!(manager.run(shutdown.clone().cancelled_owned())); + let observation = advance(running.as_mut(), || registry.observe(None).ok()).await; + + let resolver = ScopeResolver::new(Arc::clone(&pool), LIMITS); + let actor = actor_of(1); + let seeded = seed_full( + &resolver.cache, + Arc::clone(observation.present().world()), + observation.present().epoch(), + observation.admitted_at(), + actor, + ) + .await; + + let answer = resolver + .resolve(&observation, actor, None, None) + .await + .expect("a cache hit should not touch the unconnected pool") + .expect("the fresh entry should remain reachable"); + + shutdown.cancel(); + timeout(BUDGET, running.as_mut()) + .await + .expect("the cancelled run should finish within its budget"); + + (seeded, answer) + }; + manager.shutdown().await; + + assert!( + Arc::ptr_eq(&seeded, &answer), + "the hit should return the seeded publication" + ); +} + +/// Keeps stale data available while reporting a failed detached refresh. +/// +/// The warning remains under the request span and contains the whole report in its `error` field: +/// the resolver's context opens it, with the pool's connection failure retained inside. +#[tokio::test] +async fn resolve_refresh_diagnostic() { + let (_files, pool, mut manager) = boot("resolver-refresh-diagnostic", RETENTION).await; + let registry = Arc::clone(manager.registry()); + let shutdown = CancellationToken::new(); + + let (seeded, answer, event, expected_id) = { + let mut running = pin!(manager.run(shutdown.clone().cancelled_owned())); + let observation = advance(running.as_mut(), || registry.observe(None).ok()).await; + let resolver = ScopeResolver::new(pool, LIMITS); + let actor = actor_of(1); + let seeded = seed_full( + &resolver.cache, + Arc::clone(observation.present().world()), + observation.present().epoch(), + observation + .admitted_at() + .checked_sub(LIMITS.soft) + .expect("should backdate the cached scope"), + actor, + ) + .await; + let (events, mut received) = mpsc::unbounded_channel(); + let dispatch = Dispatch::new(Registry::default().with(Diagnostics(events))); + // this current-thread runtime polls both resolution and its refresh with this subscriber. + let _default = tracing::dispatcher::set_default(&dispatch); + let request = tracing::info_span!("request"); + let expected_id = request.id().expect("should enable the request span"); + let answer = resolver + .resolve(&observation, actor, None, None) + .instrument(request) + .await + .expect("should return the stale scope before the store failure") + .expect("should retain the stale scope"); + let event = timeout(BUDGET, received.recv()) + .await + .expect("should report the store refresh failure") + .expect("should retain the diagnostic sender"); + core::assert_matches!(received.try_recv(), Err(mpsc::error::TryRecvError::Empty)); + assert_eq!(tracing::Span::current().id(), None); + + shutdown.cancel(); + timeout(BUDGET, running.as_mut()) + .await + .expect("should finish the cancelled run"); + (seeded, answer, event, expected_id) + }; + manager.shutdown().await; + + // production logs the full report through `?error`, including every context. The store driver's + // wording is outside the resolver contract. This case obtains a control connection failure from + // an identically configured pool and verifies the text that the diagnostic must retain. + let control = self::pool().await; + let refused = timeout(BUDGET, control.acquire(None)) + .await + .expect("the unconnectable pool should refuse within the budget"); + let Err(refused) = refused else { + panic!("the unconnectable pool should refuse a connection") + }; + let connection_failure = plain(&refused.current_context().to_string()); + assert!( + !connection_failure.is_empty(), + "the control refusal should carry text for the report to have kept" + ); + + let (level, scope, fields) = event; + assert!(Arc::ptr_eq(&seeded, &answer)); + assert_eq!(level, Level::WARN); + assert_eq!(scope, vec![expected_id]); + assert_eq!( + fields.keys().map(String::as_str).collect::>(), + ["error", "message"], + "the diagnostic should carry the report and its message, and nothing else" + ); + assert_eq!( + fields["message"], + "failed to refresh the cached visibility scope" + ); + let error = plain(&fields["error"]); + assert!( + error.contains("the resolution reached no store connection"), + "the resolver's own context should open the report:\n{error}" + ); + assert!( + error.contains(&connection_failure), + "the connection failure should survive into the report, and this one did not:\n{error}" + ); + assert!( + error.lines().count() > 1, + "the field should hold the whole report rather than one context's text:\n{error}" + ); +} + +/// Reuses a cached filter document and publication when the caller supplies only its digest. +#[tokio::test] +async fn resolve_filtered_cached_document() { + let (_files, pool, mut manager) = boot("resolver-filtered-cached-document", RETENTION).await; + let registry = Arc::clone(manager.registry()); + let shutdown = CancellationToken::new(); + + let (seeded, answer) = { + let mut running = pin!(manager.run(shutdown.clone().cancelled_owned())); + let observation = advance(running.as_mut(), || registry.observe(None).ok()).await; + + let resolver = ScopeResolver::new(Arc::clone(&pool), LIMITS); + let actor = actor_of(1); + let document: Arc = Arc::from( + RawValue::from_string(r#"{"path":["fixture"]}"#.to_owned()) + .expect("the fixture filter document should parse"), + ); + let digest = FilterDigest::of(document.get().as_bytes()); + let seeded = seed_filtered( + &resolver.cache, + Arc::clone(observation.present().world()), + observation.present().epoch(), + observation.admitted_at(), + actor, + digest, + document, + ) + .await; + + let answer = resolver + .resolve(&observation, actor, Some(digest), None) + .await + .expect("a cached document should not touch the unconnected pool") + .expect("the cached document should reach its publication"); + + shutdown.cancel(); + timeout(BUDGET, running.as_mut()) + .await + .expect("the cancelled run should finish within its budget"); + + (seeded, answer) + }; + manager.shutdown().await; + + assert!( + Arc::ptr_eq(&seeded, &answer), + "an unresent filter should still reach the retained document's publication" + ); +} + +/// Refuses filtered resolution when no matching document is available. +/// +/// A filtered lookup naming a digest nothing retained resolves to nothing rather than +/// to an unfiltered scope, and does not reach the store for it. +#[tokio::test] +async fn resolve_filtered_missing_document() { + let (_files, pool, mut manager) = boot("resolver-filtered-missing-document", RETENTION).await; + let registry = Arc::clone(manager.registry()); + let shutdown = CancellationToken::new(); + + let answer = { + let mut running = pin!(manager.run(shutdown.clone().cancelled_owned())); + let observation = advance(running.as_mut(), || registry.observe(None).ok()).await; + + let resolver = ScopeResolver::new(Arc::clone(&pool), LIMITS); + let actor = actor_of(1); + let digest = FilterDigest::of(b"never cached"); + let answer = resolver + .resolve(&observation, actor, Some(digest), None) + .await + .expect("an uncached digest should not touch the unconnected pool"); + + shutdown.cancel(); + timeout(BUDGET, running.as_mut()) + .await + .expect("the cancelled run should finish within its budget"); + + answer + }; + manager.shutdown().await; + + assert!( + answer.is_none(), + "an uncached filtered document should resolve to nothing" + ); +} + +/// Keeps cache publications separate for two actors in the same observation. +#[tokio::test] +async fn resolve_actor_separation() { + let (_files, pool, mut manager) = boot("resolver-actor-separation", RETENTION).await; + let registry = Arc::clone(manager.registry()); + let shutdown = CancellationToken::new(); + + let (entry_a, entry_b, answer_a, answer_b) = { + let mut running = pin!(manager.run(shutdown.clone().cancelled_owned())); + let observation = advance(running.as_mut(), || registry.observe(None).ok()).await; + + let resolver = ScopeResolver::new(Arc::clone(&pool), LIMITS); + let actor_a = actor_of(1); + let actor_b = actor_of(2); + let world = observation.present().world(); + let epoch = observation.present().epoch(); + let now = observation.admitted_at(); + let entry_a = seed_full(&resolver.cache, Arc::clone(world), epoch, now, actor_a).await; + let entry_b = seed_full(&resolver.cache, Arc::clone(world), epoch, now, actor_b).await; + + let answer_a = resolver + .resolve(&observation, actor_a, None, None) + .await + .expect("a cache hit should not touch the unconnected pool") + .expect("actor a's entry should remain reachable"); + let answer_b = resolver + .resolve(&observation, actor_b, None, None) + .await + .expect("a cache hit should not touch the unconnected pool") + .expect("actor b's entry should remain reachable"); + + shutdown.cancel(); + timeout(BUDGET, running.as_mut()) + .await + .expect("the cancelled run should finish within its budget"); + + (entry_a, entry_b, answer_a, answer_b) + }; + manager.shutdown().await; + + assert!( + !Arc::ptr_eq(&entry_a, &entry_b), + "distinct actors should seed distinct publications" + ); + assert!( + Arc::ptr_eq(&answer_a, &entry_a), + "actor a should reach its own publication" + ); + assert!( + Arc::ptr_eq(&answer_b, &entry_b), + "actor b should reach its own publication" + ); +} + +/// Resolves a retained generation from its captured observation. +/// +/// After a replacement generation activates, a lookup against an observation of the +/// still-retained original answers the generation the caller asked for and returns the +/// publication seeded under it. +#[tokio::test] +async fn resolve_retained_hit() { + let (files, pool, mut manager) = boot("resolver-retained-hit", RETENTION).await; + let registry = Arc::clone(manager.registry()); + let shutdown = CancellationToken::new(); + let root = root_of(&files); + + let (seeded, answer, retained_generation, first_id) = { + let mut running = pin!(manager.run(shutdown.clone().cancelled_owned())); + let first = advance(running.as_mut(), || registry.observe(None).ok()).await; + let first_id = first.present().epoch().generation(); + + let resolver = ScopeResolver::new(Arc::clone(&pool), LIMITS); + let actor = actor_of(1); + let seeded = seed_full( + &resolver.cache, + Arc::clone(first.present().world()), + first.present().epoch(), + first.admitted_at(), + actor, + ) + .await; + + let replacement = variant_of(&files, b"resolver-retained-hit variant"); + root.activate(replacement.id()) + .expect("the replacement generation should activate"); + let retained = advance(running.as_mut(), || { + let observation = registry.observe(Some(first_id)).ok()?; + (observation.present().epoch().generation() != first_id).then_some(observation) + }) + .await; + let retained_generation = retained.requested().epoch().generation(); + + let answer = resolver + .resolve(&retained, actor, None, None) + .await + .expect("a retained hit should not touch the unconnected pool") + .expect("the retained entry should remain reachable"); + + shutdown.cancel(); + timeout(BUDGET, running.as_mut()) + .await + .expect("the cancelled run should finish within its budget"); + + (seeded, answer, retained_generation, first_id) + }; + manager.shutdown().await; + + assert_eq!( + retained_generation, first_id, + "the retained observation should answer the originally requested generation" + ); + assert!( + Arc::ptr_eq(&seeded, &answer), + "the retained lookup should return the seeded publication" + ); +} + +/// Returns `None` for a retained miss without querying the store. +#[tokio::test] +async fn resolve_retained_missing() { + let (files, pool, mut manager) = boot("resolver-retained-missing", RETENTION).await; + let registry = Arc::clone(manager.registry()); + let shutdown = CancellationToken::new(); + let root = root_of(&files); + + let answer = { + let mut running = pin!(manager.run(shutdown.clone().cancelled_owned())); + let first = advance(running.as_mut(), || registry.observe(None).ok()).await; + let first_id = first.present().epoch().generation(); + + let resolver = ScopeResolver::new(Arc::clone(&pool), LIMITS); + let actor = actor_of(1); + + let replacement = variant_of(&files, b"resolver-retained-missing variant"); + root.activate(replacement.id()) + .expect("the replacement generation should activate"); + let retained = advance(running.as_mut(), || { + let observation = registry.observe(Some(first_id)).ok()?; + (observation.present().epoch().generation() != first_id).then_some(observation) + }) + .await; + + let answer = resolver + .resolve(&retained, actor, None, None) + .await + .expect("a retained miss should not touch the unconnected pool"); + + shutdown.cancel(); + timeout(BUDGET, running.as_mut()) + .await + .expect("the cancelled run should finish within its budget"); + + answer + }; + manager.shutdown().await; + + assert!( + answer.is_none(), + "an unseeded retained epoch should resolve to nothing" + ); +} + +/// Refuses an entry from a reactivated generation's previous delta lifetime. +/// +/// Reactivation preserves the generation identity but creates a new delta lifetime. The resulting +/// cache miss reaches the store and fails on the absent test socket. +#[tokio::test] +async fn resolve_reopened_missing() { + let (files, pool, mut manager) = boot("resolver-reopened-missing", SHORT_RETENTION).await; + let registry = Arc::clone(manager.registry()); + let shutdown = CancellationToken::new(); + let root = root_of(&files); + + let (error, reopened_generation, first_id, delta_changed) = { + let mut running = pin!(manager.run(shutdown.clone().cancelled_owned())); + let first = advance(running.as_mut(), || registry.observe(None).ok()).await; + let first_id = first.present().epoch().generation(); + let first_delta = first.present().epoch().reference().id; + + let resolver = ScopeResolver::new(Arc::clone(&pool), LIMITS); + let actor = actor_of(1); + let _seeded = seed_full( + &resolver.cache, + Arc::clone(first.present().world()), + first.present().epoch(), + first.admitted_at(), + actor, + ) + .await; + + let replacement = variant_of(&files, b"resolver-reopened-missing variant"); + root.activate(replacement.id()) + .expect("the replacement generation should activate"); + advance(running.as_mut(), || { + let observation = registry.observe(None).ok()?; + (observation.present().epoch().generation() != first_id).then_some(()) + }) + .await; + + advance(running.as_mut(), || { + matches!( + registry.observe(Some(first_id)), + Err(ObserveError::Unavailable(refused)) if refused == first_id + ) + .then_some(()) + }) + .await; + + root.activate(files.generation().id()) + .expect("reactivating the original generation should succeed"); + let reopened = advance(running.as_mut(), || { + let observation = registry.observe(None).ok()?; + let epoch = observation.present().epoch(); + (epoch.generation() == first_id && epoch.reference().id != first_delta) + .then_some(observation) + }) + .await; + let reopened_generation = reopened.present().epoch().generation(); + let delta_changed = reopened.present().epoch().reference().id != first_delta; + + let error = timeout(BUDGET, resolver.resolve(&reopened, actor, None, None)) + .await + .expect("the missing socket should fail within the timeout") + .expect_err("a fresh delta lifetime should carry no seeded entry"); + + shutdown.cancel(); + timeout(BUDGET, running.as_mut()) + .await + .expect("the cancelled run should finish within its budget"); + + (error, reopened_generation, first_id, delta_changed) + }; + manager.shutdown().await; + + assert_eq!( + reopened_generation, first_id, + "reopening should keep the original generation identity" + ); + assert!( + delta_changed, + "reopening should draw a fresh delta lifetime" + ); + assert_matches!( + error.current_context(), + VisibilityProofError::Connect, + "a cache miss should reach the unconnected pool: {error:?}" + ); +}