1use 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#[derive(Default)]
23pub(crate) struct ViewReg {
24 pub(crate) catalog: RwLock<(u64, ViewCatalog)>,
25}
26
27#[derive(Default)]
29pub(crate) struct ShardViews {
30 pub(crate) version: u64,
31 pub(crate) views: Vec<ViewState>,
32 #[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 #[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 #[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
77pub 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 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 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 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 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 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 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 #[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 #[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
278fn 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
327pub(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
359pub(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 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}