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