Skip to main content

kevy_embedded/
ops_view.rs

1//! Embedded views (server parity minus VIA/FIELDS hydration:
2//! in-process callers dereference and read fields directly).
3//!
4//! Mirrors the embedded index architecture: per-shard view states in
5//! `Inner` (maintained inside `commit_write` right after index
6//! maintenance, under the same shard lock), a store-level registry,
7//! synchronous builds, typed API (`Tree` passed directly — no text
8//! grammar in-process).
9
10// The sidecar IS the catalog's persistence — `boot` reads it and a
11// directory without one "boots empty". So a rename that fails loses
12// the index definitions at the next start, after the command that
13// created them has already replied OK. That is a gap, not a
14// non-event, and it is written up as an open question rather than
15// silently accepted here: .claude/OPEN-QUESTIONS-6.4.md §3.
16#![expect(
17    clippy::let_underscore_must_use,
18    reason = "the catalog has no other home; see .claude/OPEN-QUESTIONS-6.4.md"
19)]
20
21use crate::{KevyError, KevyResult};
22use std::io;
23use std::sync::RwLock;
24
25use kevy_index::{
26    IndexValue, MaterializedSet, Tree, ViewCatalog, ViewMode, ViewSpec, eval_tree, key_in_tree,
27};
28
29use crate::ops_index::ShardSegs;
30use crate::store::{Store, lock_write};
31
32/// Store-level registry.
33#[derive(Debug, Default)]
34pub(crate) struct ViewReg {
35    pub(crate) catalog: RwLock<(u64, ViewCatalog)>,
36}
37
38/// One shard's view states (inside `Inner`, guarded by the shard lock).
39#[derive(Debug, Default)]
40pub(crate) struct ShardViews {
41    pub(crate) version: u64,
42    pub(crate) views: Vec<ViewState>,
43    /// `reserved_bytes` generation cache — see
44    /// `ShardSegs::stats_dirty`; same contract, view half. Tier-only,
45    /// like its twin.
46    #[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
47    pub(crate) stats_dirty: bool,
48    #[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
49    pub(crate) reserved_cache: u64,
50}
51
52#[derive(Debug)]
53
54pub(crate) struct ViewState {
55    spec: ViewSpec,
56    mat: Option<MaterializedSet>,
57    needs_rebuild: bool,
58}
59
60impl ShardViews {
61    /// Invalidate the cache — no-op without the tier backend; see
62    /// `ShardSegs::mark_stats_dirty`.
63    #[inline]
64    pub(crate) fn mark_stats_dirty(&mut self) {
65        #[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
66        {
67            self.stats_dirty = true;
68        }
69    }
70
71    /// Σ approximate heap bytes of the materialized view sets — the
72    /// view half of the tier's `reserved_bytes` feed.
73    /// Virtual views hold no set and contribute nothing.
74    #[cfg(all(feature = "tier", not(target_arch = "wasm32")))]
75    pub(crate) fn reserved_bytes(&mut self) -> u64 {
76        if !self.stats_dirty {
77            return self.reserved_cache;
78        }
79        let sum = self
80            .views
81            .iter()
82            .map(|v| v.mat.as_ref().map_or(0, MaterializedSet::approx_bytes))
83            .sum();
84        self.reserved_cache = sum;
85        self.stats_dirty = false;
86        sum
87    }
88}
89
90/// One page of view members plus the resume cursor.
91pub type ViewPage = (Vec<(Vec<u8>, IndexValue)>, Option<(IndexValue, Vec<u8>)>);
92
93#[cfg(feature = "persist")]
94const SIDECAR: &str = "view-catalog.meta";
95
96impl Store {
97    /// Declare a view (typed tree; `via` is not supported embedded —
98    /// read fields in-process). Builds synchronously.
99    pub fn view_create(
100        &self,
101        name: &[u8],
102        tree: Tree,
103        order_by: &[u8],
104        desc: bool,
105        mode: ViewMode,
106    ) -> KevyResult<()> {
107        self.check_view_refs(&tree, order_by)?;
108        let spec = ViewSpec {
109            name: name.to_vec(),
110            tree,
111            order_by: order_by.to_vec(),
112            desc,
113            mode,
114            via: None,
115        };
116        {
117            let mut g =
118                self.views.catalog.write().unwrap_or_else(std::sync::PoisonError::into_inner);
119            let (ver, cat) = &mut *g;
120            cat.create(spec).map_err(|e| io::Error::new(io::ErrorKind::InvalidInput, e))?;
121            *ver += 1;
122        }
123        self.persist_view_sidecar();
124        for shard in self.shards.iter() {
125            let mut g = lock_write(shard);
126            let inner = &mut *g;
127            crate::ops_index::sync_segs(&self.indexes, &mut inner.idx_segs, &mut inner.store);
128            sync_views(&self.views, &mut inner.view_segs, &inner.idx_segs);
129        }
130        Ok(())
131    }
132
133    /// Every index a view references (its leaves + ORDER BY) must
134    /// already be declared.
135    fn check_view_refs(&self, tree: &Tree, order_by: &[u8]) -> KevyResult<()> {
136        let mut names: Vec<Vec<u8>> = vec![order_by.to_vec()];
137        tree.each_leaf(&mut |l| names.push(l.index.clone()));
138        let g = self.indexes.catalog.read().unwrap_or_else(std::sync::PoisonError::into_inner);
139        for n in &names {
140            if g.1.get(n).is_none() {
141                return Err(KevyError::InvalidInput("view references unknown index".into()));
142            }
143        }
144        Ok(())
145    }
146
147    /// Drop a view; `false` if absent.
148    pub fn view_drop(&self, name: &[u8]) -> bool {
149        let hit = {
150            let mut g =
151                self.views.catalog.write().unwrap_or_else(std::sync::PoisonError::into_inner);
152            let (ver, cat) = &mut *g;
153            let hit = cat.drop_view(name);
154            if hit {
155                *ver += 1;
156            }
157            hit
158        };
159        if hit {
160            self.persist_view_sidecar();
161        }
162        hit
163    }
164
165    /// Ordered page across shards (`after` resumes exclusively; DESC
166    /// views page from the large end).
167    pub fn view_query(
168        &self,
169        name: &[u8],
170        after: Option<&(IndexValue, Vec<u8>)>,
171        limit: usize,
172    ) -> KevyResult<ViewPage> {
173        let limit = limit.clamp(1, 100_000);
174        let mut desc = false;
175        let mut all: Vec<(IndexValue, Vec<u8>)> = Vec::new();
176        let mut found = false;
177        for shard in self.shards.iter() {
178            let mut g = lock_write(shard);
179            let inner = &mut *g;
180            crate::ops_index::sync_segs(&self.indexes, &mut inner.idx_segs, &mut inner.store);
181            sync_views(&self.views, &mut inner.view_segs, &inner.idx_segs);
182            let Some(i) = inner.view_segs.views.iter().position(|v| v.spec.name == name) else {
183                continue;
184            };
185            found = true;
186            if inner.view_segs.views[i].needs_rebuild {
187                inner.view_segs.mark_stats_dirty();
188                rebuild(&mut inner.view_segs.views[i], &inner.idx_segs);
189            }
190            let vs = &inner.view_segs.views[i];
191            desc = vs.spec.desc;
192            match &vs.mat {
193                Some(m) => all.extend(m.page(after, limit, vs.spec.desc)),
194                None => stream_virtual(&vs.spec, &inner.idx_segs, after, limit, &mut all),
195            }
196        }
197        if !found {
198            return Err(KevyError::NotFound("no such view".into()));
199        }
200        all.sort();
201        if desc {
202            all.reverse();
203        }
204        all.truncate(limit);
205        let next = if all.len() == limit { all.last().cloned() } else { None };
206        Ok((all.into_iter().map(|(v, k)| (k, v)).collect(), next))
207    }
208
209    /// Declared views (name, mode, leaves).
210    pub fn view_list(&self) -> Vec<(Vec<u8>, ViewMode, usize)> {
211        let g = self.views.catalog.read().unwrap_or_else(std::sync::PoisonError::into_inner);
212        g.1.iter().map(|s| (s.name.clone(), s.mode, s.tree.leaves())).collect()
213    }
214
215    /// Summed member count across shards.
216    pub fn view_count(&self, name: &[u8]) -> KevyResult<u64> {
217        Ok(self.view_query(name, None, 100_000)?.0.len() as u64)
218    }
219
220    #[cfg(feature = "persist")]
221    fn persist_view_sidecar(&self) {
222        let Some(dir) = &self.config.data_dir else { return };
223        let g = self.views.catalog.read().unwrap_or_else(std::sync::PoisonError::into_inner);
224        let tmp = dir.join("view-catalog.meta.tmp");
225        if std::fs::write(&tmp, g.1.to_sidecar()).is_ok() {
226            let _ = std::fs::rename(&tmp, dir.join(SIDECAR));
227        }
228    }
229
230    /// Without `persist` there is no data dir — no sidecar to write
231    /// or load; both halves are no-ops.
232    #[cfg(not(feature = "persist"))]
233    fn persist_view_sidecar(&self) {}
234
235    #[cfg(not(feature = "persist"))]
236    pub(crate) fn view_boot(&self) {}
237
238    /// Boot half — load the persisted view catalog.
239    #[cfg(feature = "persist")]
240    pub(crate) fn view_boot(&self) {
241        let Some(dir) = &self.config.data_dir else { return };
242        if let Ok(text) = std::fs::read_to_string(dir.join(SIDECAR))
243            && let Some(cat) = ViewCatalog::from_sidecar(&text)
244            && !cat.is_empty()
245        {
246            let mut g =
247                self.views.catalog.write().unwrap_or_else(std::sync::PoisonError::into_inner);
248            *g = (g.0 + 1, cat);
249        }
250    }
251}
252
253fn resolver<'a>(segs: &'a ShardSegs) -> impl Fn(&[u8]) -> Option<&'a kevy_index::Segment> {
254    move |name: &[u8]| segs.segs.iter().find(|(s, _)| s.name == name).map(|(_, seg)| seg)
255}
256
257fn eval_shard(spec: &ViewSpec, segs: &ShardSegs) -> Vec<(IndexValue, Vec<u8>)> {
258    let r = resolver(segs);
259    let members = eval_tree(&spec.tree, &&r);
260    members
261        .into_iter()
262        .filter_map(|k| r(&spec.order_by).and_then(|s| s.verify_entry(&k)).map(|v| (v.clone(), k)))
263        .collect()
264}
265
266/// Virtual-mode page: order-driven streaming over the ORDER BY index,
267/// probing the tree per key (same clamp rationale as the server
268/// pager).
269fn stream_virtual(
270    spec: &ViewSpec,
271    segs: &ShardSegs,
272    after: Option<&(IndexValue, Vec<u8>)>,
273    limit: usize,
274    all: &mut Vec<(IndexValue, Vec<u8>)>,
275) {
276    let r = resolver(segs);
277    if let Some(order_seg) = r(&spec.order_by) {
278        let cursor = after.map(|(v, k)| kevy_index::Cursor { value: v.clone(), key: k.clone() });
279        let mut got = 0usize;
280        for (v, k) in order_seg.scan(cursor.as_ref(), spec.desc) {
281            if key_in_tree(&spec.tree, k, &&r) {
282                all.push((v.clone(), k.to_vec()));
283                got += 1;
284                if got == limit {
285                    break;
286                }
287            }
288        }
289    }
290}
291
292fn rebuild(vs: &mut ViewState, segs: &ShardSegs) {
293    let spec = vs.spec.clone();
294    let Some(mat) = &mut vs.mat else {
295        vs.needs_rebuild = false;
296        return;
297    };
298    mat.clear();
299    let mut rows = eval_shard(&spec, segs);
300    rows.sort();
301    if let ViewMode::Materialized { top_k } = spec.mode
302        && top_k > 0
303    {
304        if spec.desc {
305            rows.reverse();
306        }
307        rows.truncate((top_k + top_k / 4) as usize);
308    }
309    for (v, k) in rows {
310        mat.apply(&k, true, Some(v));
311    }
312    vs.needs_rebuild = false;
313}
314
315/// Reconcile with the catalog (under the shard lock).
316pub(crate) fn sync_views(reg: &ViewReg, sv: &mut ShardViews, segs: &ShardSegs) {
317    let g = reg.catalog.read().unwrap_or_else(std::sync::PoisonError::into_inner);
318    let (ver, cat) = &*g;
319    if sv.version == *ver {
320        return;
321    }
322    sv.mark_stats_dirty();
323    let mut next = Vec::new();
324    for spec in cat.iter() {
325        match sv.views.iter().position(|v| v.spec == *spec) {
326            Some(i) => next.push(sv.views.swap_remove(i)),
327            None => {
328                let mat = match spec.mode {
329                    ViewMode::Virtual => None,
330                    ViewMode::Materialized { top_k } => {
331                        Some(MaterializedSet::new(top_k, spec.desc))
332                    }
333                };
334                let mut vs = ViewState { spec: spec.clone(), needs_rebuild: mat.is_some(), mat };
335                if vs.needs_rebuild {
336                    rebuild(&mut vs, segs);
337                }
338                next.push(vs);
339            }
340        }
341    }
342    sv.views = next;
343    sv.version = *ver;
344}
345
346/// Write hook — call AFTER `ops_index::on_commit` (same shard lock).
347pub(crate) fn on_commit(reg: &ViewReg, sv: &mut ShardViews, segs: &ShardSegs, parts: &[&[u8]]) {
348    {
349        let g = reg.catalog.read().unwrap_or_else(std::sync::PoisonError::into_inner);
350        if g.1.is_empty() {
351            return;
352        }
353    }
354    sync_views(reg, sv, segs);
355    let verb = parts.first().copied().unwrap_or(b"");
356    if verb.eq_ignore_ascii_case(b"FLUSHALL") || verb.eq_ignore_ascii_case(b"FLUSHDB") {
357        for vs in &mut sv.views {
358            if let Some(m) = &mut vs.mat {
359                m.clear();
360            }
361        }
362        sv.mark_stats_dirty();
363        return;
364    }
365    // Same exact written-key walk as the index hook.
366    let mut touched = false;
367    let views = &mut sv.views;
368    crate::ops_index::each_written_key_pub(verb, parts, |key| {
369        for vs in &mut *views {
370            let Some(mat) = &mut vs.mat else { continue };
371            touched = true;
372            let r = resolver(segs);
373            let member = key_in_tree(&vs.spec.tree, key, &&r);
374            let order = r(&vs.spec.order_by).and_then(|s| s.verify_entry(key)).cloned();
375            if mat.apply(key, member, order) {
376                vs.needs_rebuild = true;
377            }
378        }
379    });
380    if touched {
381        sv.mark_stats_dirty();
382    }
383}