1#![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#[derive(Debug, Default)]
34pub(crate) struct ViewReg {
35 pub(crate) catalog: RwLock<(u64, ViewCatalog)>,
36}
37
38#[derive(Debug, Default)]
40pub(crate) struct ShardViews {
41 pub(crate) version: u64,
42 pub(crate) views: Vec<ViewState>,
43 #[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 #[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 #[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
90pub 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 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 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 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 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 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 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 #[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 #[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
266fn 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
315pub(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
346pub(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 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}