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}
33
34pub(crate) struct ViewState {
35 spec: ViewSpec,
36 mat: Option<MaterializedSet>,
37 needs_rebuild: bool,
38}
39
40impl ShardViews {
41 #[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
53pub 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 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 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 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 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 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 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 #[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 #[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
251fn 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
300pub(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
331pub(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 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}