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