1use std::collections::HashMap;
13
14use omgbase_search::{EmbeddingProvider, f32_to_blob};
15use omgbase_store::Store;
16use oqx::ast::{Expr, Follow, OpNode, OrderSpec, Query, SelectItem, Subquery, Where};
17use oqx::{Consumer, Engine, InMemoryEngine, Value};
18use serde_json::{Map, Value as Json};
19
20use crate::context::{SemanticVec, StoreContext, render_row_values};
21use crate::cursor::{decode_cursor, encode_cursor};
22use crate::error::{Result, SurfaceError};
23use crate::planner::SqlitePlanner;
24
25pub const DEFAULT_LIMIT: usize = 50;
27
28const ROW_FNS: [&str; 12] = [
31 "text",
32 "semantic",
33 "under",
34 "under_heading",
35 "within",
36 "under_kind",
37 "yaml_path",
38 "json_pointer",
39 "has_edge",
40 "has_anchor",
41 "child_count",
42 "parent_type",
43];
44
45const ID_KEY: &str = "__oqx_id";
46const PATH_KEY: &str = "__oqx_path";
47const VALUE_KEY: &str = "__oqx_value";
51
52#[derive(Clone, Copy, Default)]
54pub struct QueryOptions<'a> {
55 pub limit: Option<usize>,
57 pub cursor: Option<&'a str>,
59 pub provider: Option<&'a dyn EmbeddingProvider>,
62 pub in_memory: bool,
65}
66
67#[derive(Clone, Debug, PartialEq)]
69pub struct OqxResult {
70 pub hits: Vec<Json>,
71 pub truncated: bool,
72 pub cursor: Option<String>,
73 pub consumer: Consumer,
74 pub count: Option<f64>,
75 pub exists: Option<bool>,
76 pub none: Option<bool>,
77 pub values: Option<Vec<Json>>,
79}
80
81impl OqxResult {
82 fn scalar(consumer: Consumer) -> Self {
83 Self {
84 hits: Vec::new(),
85 truncated: false,
86 cursor: None,
87 consumer,
88 count: None,
89 exists: None,
90 none: None,
91 values: None,
92 }
93 }
94
95 #[must_use]
98 pub fn to_json(&self) -> Json {
99 let mut m = Map::new();
100 m.insert("hits".to_owned(), Json::Array(self.hits.clone()));
101 m.insert("truncated".to_owned(), Json::Bool(self.truncated));
102 m.insert(
103 "cursor".to_owned(),
104 self.cursor.clone().map_or(Json::Null, Json::String),
105 );
106 m.insert(
107 "consumer".to_owned(),
108 Json::String(self.consumer.as_str().to_owned()),
109 );
110 if let Some(n) = self.count {
111 m.insert("count".to_owned(), Value::Number(n).to_canonical_json());
112 }
113 if let Some(b) = self.exists {
114 m.insert("exists".to_owned(), Json::Bool(b));
115 }
116 if let Some(b) = self.none {
117 m.insert("none".to_owned(), Json::Bool(b));
118 }
119 if let Some(v) = &self.values {
120 m.insert("values".to_owned(), Json::Array(v.clone()));
121 }
122 Json::Object(m)
123 }
124}
125
126fn self_ref() -> Expr {
129 Expr::Ident {
130 name: "$self".to_owned(),
131 }
132}
133
134fn rewrite_expr(e: &Expr) -> Expr {
135 match e {
136 Expr::Member { recv, name } => Expr::Member {
137 recv: Box::new(rewrite_expr(recv)),
138 name: name.clone(),
139 },
140 Expr::Index { recv, index } => Expr::Index {
141 recv: Box::new(rewrite_expr(recv)),
142 index: Box::new(rewrite_expr(index)),
143 },
144 Expr::Unary { op, expr } => Expr::Unary {
145 op: *op,
146 expr: Box::new(rewrite_expr(expr)),
147 },
148 Expr::Binary { op, left, right } => Expr::Binary {
149 op: *op,
150 left: Box::new(rewrite_expr(left)),
151 right: Box::new(rewrite_expr(right)),
152 },
153 Expr::Logical { op, left, right } => Expr::Logical {
154 op: *op,
155 left: Box::new(rewrite_expr(left)),
156 right: Box::new(rewrite_expr(right)),
157 },
158 Expr::In { left, right } => Expr::In {
159 left: Box::new(rewrite_expr(left)),
160 right: Box::new(rewrite_expr(right)),
161 },
162 Expr::Range {
163 lo,
164 hi,
165 exclusive_end,
166 } => Expr::Range {
167 lo: lo.as_ref().map(|x| Box::new(rewrite_expr(x))),
168 hi: hi.as_ref().map(|x| Box::new(rewrite_expr(x))),
169 exclusive_end: *exclusive_end,
170 },
171 Expr::Call { recv, name, args } => {
172 let args = args.iter().map(rewrite_expr).collect();
173 match recv {
174 None if ROW_FNS.contains(&name.as_str()) => Expr::Call {
175 recv: Some(Box::new(self_ref())),
176 name: name.clone(),
177 args,
178 },
179 None => Expr::Call {
180 recv: None,
181 name: name.clone(),
182 args,
183 },
184 Some(r) => Expr::Call {
185 recv: Some(Box::new(rewrite_expr(r))),
186 name: name.clone(),
187 args,
188 },
189 }
190 }
191 Expr::Lit(_) | Expr::Ident { .. } | Expr::Outer { .. } | Expr::Binding { .. } => e.clone(),
192 }
193}
194
195fn rewrite_where(w: &Where) -> Where {
196 match w {
197 Where::And { parts } => Where::And {
198 parts: parts.iter().map(rewrite_where).collect(),
199 },
200 Where::Or { parts } => Where::Or {
201 parts: parts.iter().map(rewrite_where).collect(),
202 },
203 Where::Not { expr } => Where::Not {
204 expr: Box::new(rewrite_where(expr)),
205 },
206 Where::Scalar { expr } => Where::Scalar {
207 expr: rewrite_expr(expr),
208 },
209 Where::Op(op) => Where::Op(Box::new(rewrite_op(op))),
210 }
211}
212
213fn rewrite_op(op: &OpNode) -> OpNode {
214 OpNode {
215 receiver: rewrite_expr(&op.receiver),
216 op: op.op,
217 sub: rewrite_sub(&op.sub),
218 count_cmp: op.count_cmp.clone(),
219 distinct: op.distinct,
220 }
221}
222
223fn rewrite_follow(f: &Follow) -> Follow {
224 Follow {
225 receiver: rewrite_expr(&f.receiver),
226 distinct: f.distinct,
227 r#where: f.r#where.as_ref().map(rewrite_expr),
228 frontier: f.frontier.as_ref().map(rewrite_expr),
229 depth: f.depth,
230 by: f.by.as_ref().map(rewrite_expr),
231 }
232}
233
234fn rewrite_select(items: &[SelectItem]) -> Vec<SelectItem> {
235 items
236 .iter()
237 .map(|it| match it {
238 SelectItem::Field { name, expr, lift } => SelectItem::Field {
239 name: name.clone(),
240 expr: rewrite_expr(expr),
241 lift: *lift,
242 },
243 SelectItem::Collect { name, op } => SelectItem::Collect {
244 name: name.clone(),
245 op: Box::new(rewrite_op(op)),
246 },
247 })
248 .collect()
249}
250
251fn rewrite_order(o: Option<&Vec<OrderSpec>>) -> Option<Vec<OrderSpec>> {
252 o.map(|specs| {
253 specs
254 .iter()
255 .map(|s| OrderSpec {
256 expr: rewrite_expr(&s.expr),
257 desc: s.desc,
258 })
259 .collect()
260 })
261}
262
263fn rewrite_sub(s: &Subquery) -> Subquery {
264 Subquery {
265 from: s.from.iter().map(rewrite_expr).collect(),
266 r#where: s.r#where.as_ref().map(rewrite_where),
267 select: rewrite_select(&s.select),
268 order_by: rewrite_order(s.order_by.as_ref()),
269 follow: s.follow.as_ref().map(rewrite_follow),
270 values: s.values,
271 limit: s.limit.as_ref().map(rewrite_expr),
272 offset: s.offset.as_ref().map(rewrite_expr),
273 }
274}
275
276#[must_use]
278pub fn rewrite_query(q: &Query) -> Query {
279 Query {
280 source: rewrite_expr(&q.source),
281 from: q.from.iter().map(rewrite_expr).collect(),
282 r#where: q.r#where.as_ref().map(rewrite_where),
283 select: rewrite_select(&q.select),
284 order_by: rewrite_order(q.order_by.as_ref()),
285 consumer: q.consumer,
286 follow: q.follow.as_ref().map(rewrite_follow),
287 distinct: q.distinct,
288 values: q.values,
289 limit: q.limit.as_ref().map(rewrite_expr),
290 offset: q.offset.as_ref().map(rewrite_expr),
291 }
292}
293
294fn visit_expr(e: &Expr, out: &mut Vec<String>) {
297 match e {
298 Expr::Call { recv, name, args } => {
299 if recv.is_none() && name == "semantic" {
300 if let Some(Expr::Lit(Value::Str(s))) = args.first() {
301 if !out.contains(s) {
302 out.push(s.clone());
303 }
304 }
305 }
306 if let Some(r) = recv {
307 visit_expr(r, out);
308 }
309 for a in args {
310 visit_expr(a, out);
311 }
312 }
313 Expr::Member { recv, .. } => visit_expr(recv, out),
314 Expr::Index { recv, index } => {
315 visit_expr(recv, out);
316 visit_expr(index, out);
317 }
318 Expr::Unary { expr, .. } => visit_expr(expr, out),
319 Expr::Binary { left, right, .. }
320 | Expr::Logical { left, right, .. }
321 | Expr::In { left, right } => {
322 visit_expr(left, out);
323 visit_expr(right, out);
324 }
325 Expr::Range { lo, hi, .. } => {
326 if let Some(l) = lo {
327 visit_expr(l, out);
328 }
329 if let Some(h) = hi {
330 visit_expr(h, out);
331 }
332 }
333 Expr::Lit(_) | Expr::Ident { .. } | Expr::Outer { .. } | Expr::Binding { .. } => {}
334 }
335}
336
337fn visit_where(w: &Where, out: &mut Vec<String>) {
338 match w {
339 Where::And { parts } | Where::Or { parts } => {
340 parts.iter().for_each(|p| visit_where(p, out))
341 }
342 Where::Not { expr } => visit_where(expr, out),
343 Where::Scalar { expr } => visit_expr(expr, out),
344 Where::Op(op) => visit_op(op, out),
345 }
346}
347
348fn visit_op(op: &OpNode, out: &mut Vec<String>) {
349 visit_expr(&op.receiver, out);
350 visit_sub(&op.sub, out);
351}
352
353fn visit_select(items: &[SelectItem], out: &mut Vec<String>) {
354 for it in items {
355 match it {
356 SelectItem::Field { expr, .. } => visit_expr(expr, out),
357 SelectItem::Collect { op, .. } => visit_op(op, out),
358 }
359 }
360}
361
362fn visit_follow(f: &Follow, out: &mut Vec<String>) {
363 visit_expr(&f.receiver, out);
364 for x in [&f.r#where, &f.frontier, &f.by].into_iter().flatten() {
365 visit_expr(x, out);
366 }
367}
368
369fn visit_sub(s: &Subquery, out: &mut Vec<String>) {
370 s.from.iter().for_each(|e| visit_expr(e, out));
371 if let Some(w) = &s.r#where {
372 visit_where(w, out);
373 }
374 visit_select(&s.select, out);
375 if let Some(o) = &s.order_by {
376 o.iter().for_each(|spec| visit_expr(&spec.expr, out));
377 }
378 if let Some(f) = &s.follow {
379 visit_follow(f, out);
380 }
381}
382
383#[must_use]
386pub fn collect_semantic_phrases(source: &str) -> Vec<String> {
387 let Ok(q) = oqx::parse_string(source) else {
388 return Vec::new();
389 };
390 let mut out = Vec::new();
391 visit_expr(&q.source, &mut out);
392 q.from.iter().for_each(|e| visit_expr(e, &mut out));
393 if let Some(w) = &q.r#where {
394 visit_where(w, &mut out);
395 }
396 visit_select(&q.select, &mut out);
397 if let Some(o) = &q.order_by {
398 o.iter().for_each(|spec| visit_expr(&spec.expr, &mut out));
399 }
400 if let Some(f) = &q.follow {
401 visit_follow(f, &mut out);
402 }
403 out
404}
405
406fn expr_mentions(e: &Expr, name: &str) -> bool {
409 match e {
410 Expr::Ident { name: n } | Expr::Outer { name: n, .. } => n == name,
411 Expr::Member { recv, .. } => expr_mentions(recv, name),
412 Expr::Index { recv, index } => expr_mentions(recv, name) || expr_mentions(index, name),
413 Expr::Unary { expr, .. } => expr_mentions(expr, name),
414 Expr::Binary { left, right, .. }
415 | Expr::Logical { left, right, .. }
416 | Expr::In { left, right } => expr_mentions(left, name) || expr_mentions(right, name),
417 Expr::Range { lo, hi, .. } => [lo, hi]
418 .into_iter()
419 .flatten()
420 .any(|x| expr_mentions(x, name)),
421 Expr::Call { recv, args, .. } => {
422 recv.as_deref().is_some_and(|r| expr_mentions(r, name))
423 || args.iter().any(|a| expr_mentions(a, name))
424 }
425 Expr::Lit(_) | Expr::Binding { .. } => false,
426 }
427}
428
429fn where_mentions(w: &Where, name: &str) -> bool {
430 match w {
431 Where::And { parts } | Where::Or { parts } => parts.iter().any(|p| where_mentions(p, name)),
432 Where::Not { expr } => where_mentions(expr, name),
433 Where::Scalar { expr } => expr_mentions(expr, name),
434 Where::Op(op) => op_mentions(op, name),
435 }
436}
437
438fn op_mentions(op: &OpNode, name: &str) -> bool {
439 expr_mentions(&op.receiver, name) || sub_mentions(&op.sub, name)
440}
441
442fn select_mentions(items: &[SelectItem], name: &str) -> bool {
443 items.iter().any(|it| match it {
444 SelectItem::Field { expr, .. } => expr_mentions(expr, name),
445 SelectItem::Collect { op, .. } => op_mentions(op, name),
446 })
447}
448
449fn follow_mentions(f: &Follow, name: &str) -> bool {
450 expr_mentions(&f.receiver, name)
451 || [&f.r#where, &f.frontier, &f.by]
452 .into_iter()
453 .flatten()
454 .any(|x| expr_mentions(x, name))
455}
456
457fn order_mentions(o: Option<&Vec<OrderSpec>>, name: &str) -> bool {
458 o.is_some_and(|specs| specs.iter().any(|s| expr_mentions(&s.expr, name)))
459}
460
461fn sub_mentions(s: &Subquery, name: &str) -> bool {
462 s.from.iter().any(|e| expr_mentions(e, name))
463 || s.r#where.as_ref().is_some_and(|w| where_mentions(w, name))
464 || select_mentions(&s.select, name)
465 || order_mentions(s.order_by.as_ref(), name)
466 || s.follow.as_ref().is_some_and(|f| follow_mentions(f, name))
467 || [&s.limit, &s.offset]
468 .into_iter()
469 .flatten()
470 .any(|x| expr_mentions(x, name))
471}
472
473fn mentions_outside_source(q: &Query, name: &str) -> bool {
477 q.from.iter().any(|e| expr_mentions(e, name))
478 || q.r#where.as_ref().is_some_and(|w| where_mentions(w, name))
479 || select_mentions(&q.select, name)
480 || order_mentions(q.order_by.as_ref(), name)
481 || q.follow.as_ref().is_some_and(|f| follow_mentions(f, name))
482 || [&q.limit, &q.offset]
483 .into_iter()
484 .flatten()
485 .any(|x| expr_mentions(x, name))
486}
487
488fn to_hit(row: Value) -> Value {
496 let Value::Object(o) = render_row_values(row) else {
497 return Value::Object(oqx::Object::new());
498 };
499 let mut id = Value::Undefined;
500 let mut path = Value::Undefined;
501 let mut rest = Vec::new();
502 for (k, v) in o {
503 match k.as_str() {
504 ID_KEY => id = v,
505 PATH_KEY => path = v,
506 _ => rest.push((k, v)),
507 }
508 }
509 let mut hit = oqx::Object::with_capacity(rest.len() + 2);
510 hit.insert("id", Value::Str(id.to_string()));
511 hit.insert(
512 "path",
513 Value::Str(if path.is_absent() {
514 String::new()
515 } else {
516 path.to_string()
517 }),
518 );
519 for (k, v) in rest {
520 hit.insert(k, v);
521 }
522 Value::Object(hit)
523}
524
525fn hit_str(hit: &Value, key: &str) -> String {
526 hit.as_object()
527 .and_then(|o| o.get(key))
528 .map(|v| v.to_string())
529 .unwrap_or_default()
530}
531
532fn dedup_hits_by_projection(hits: Vec<Value>) -> Vec<Value> {
536 let mut seen: Vec<String> = Vec::new();
537 let mut out = Vec::new();
538 for h in hits {
539 let mut pairs: Vec<(String, Value)> = h
540 .as_object()
541 .map(|o| {
542 o.iter()
543 .filter(|(k, _)| *k != "id" && *k != "path")
544 .map(|(k, v)| (k.to_owned(), v.clone()))
545 .collect()
546 })
547 .unwrap_or_default();
548 pairs.sort_by(|a, b| a.0.cmp(&b.0));
549 let key = Value::Array(
550 pairs
551 .into_iter()
552 .map(|(k, v)| Value::Array(vec![Value::Str(k), v]))
553 .collect(),
554 )
555 .to_canonical_json()
556 .to_string();
557 if seen.contains(&key) {
558 continue;
559 }
560 seen.push(key);
561 out.push(h);
562 }
563 out
564}
565
566fn const_bound(e: Option<&Expr>, word: &str) -> Result<Option<usize>> {
569 match e {
570 None => Ok(None),
571 Some(Expr::Lit(Value::Number(n))) if n.fract() == 0.0 && *n >= 0.0 && n.is_finite() => {
572 Ok(Some(*n as usize))
573 }
574 Some(_) => Err(SurfaceError::filter_invalid(
575 format!("top-level {word} must be a non-negative integer literal"),
576 "OQX",
577 )),
578 }
579}
580
581fn value_of(hit: &Value) -> Value {
582 hit.as_object()
583 .and_then(|o| o.get(VALUE_KEY))
584 .cloned()
585 .unwrap_or(Value::Undefined)
586}
587
588fn without_value_key(hit: Value) -> Json {
589 hit.to_canonical_json()
590}
591
592pub fn query(
596 store: &Store,
597 repo_id: &str,
598 source: &str,
599 opts: QueryOptions<'_>,
600) -> Result<OqxResult> {
601 let phrases = collect_semantic_phrases(source);
606 let mut semantic: HashMap<String, SemanticVec> = HashMap::new();
607 if let Some(provider) = opts.provider.filter(|_| !phrases.is_empty()) {
608 for phrase in phrases {
609 let vec = provider
610 .embed_query(&phrase)
611 .map_err(|e| SurfaceError::new(e.code(), e.to_string()))?;
612 semantic.insert(
613 phrase,
614 SemanticVec {
615 model: provider.model().to_owned(),
616 vec: f32_to_blob(&vec),
617 },
618 );
619 }
620 }
621 let runner = Runner {
622 store,
623 repo_id,
624 semantic,
625 planned: !opts.in_memory,
626 };
627 run_inner(&runner, source, opts)
628}
629
630struct Runner<'a> {
636 store: &'a Store,
637 repo_id: &'a str,
638 semantic: HashMap<String, SemanticVec>,
639 planned: bool,
640}
641
642impl Runner<'_> {
643 fn run(&self, q: &Query) -> Result<oqx::OqxResult> {
650 let conn = self.store.conn();
651 let ctx = StoreContext::new(conn, self.repo_id, self.semantic.clone());
652 let plan = if self.planned {
653 SqlitePlanner::new(conn, self.repo_id)
654 .try_plan(q, &[])
655 .map_err(|e| SurfaceError::other(format!("sqlite: {e}")))?
656 } else {
657 None
658 };
659 let (ctx, residual) = match plan {
660 Some(plan) => {
665 let once = !mentions_outside_source(&plan.residual, oqx::ROWS_ROOT);
666 let ctx = if once {
667 ctx.with_rows_root_once(plan.rows)
668 } else {
669 ctx.with_rows_root(plan.rows)
670 };
671 (ctx, Some(plan.residual))
672 }
673 None => (ctx, None),
674 };
675 let engine = InMemoryEngine::new(ctx);
676 let out = engine.run(residual.as_ref().unwrap_or(q), &[]);
677 if let Some(failed) = engine.context().take_root_failure() {
680 return Err(failed.into());
681 }
682 Ok(out?)
683 }
684}
685
686fn run_inner(engine: &Runner<'_>, source: &str, opts: QueryOptions<'_>) -> Result<OqxResult> {
687 let parsed = rewrite_query(&oqx::parse_string(source)?);
688 let consumer = parsed.consumer;
689
690 match consumer {
691 Consumer::Exists => {
692 let res = engine.run(&parsed)?;
693 let mut r = OqxResult::scalar(consumer);
694 r.exists = Some(matches!(res, oqx::OqxResult::Exists(true)));
695 return Ok(r);
696 }
697 Consumer::Count => {
698 let res = engine.run(&parsed)?;
699 let mut r = OqxResult::scalar(consumer);
700 r.count = Some(match res {
701 oqx::OqxResult::Count(n) => n,
702 _ => 0.0,
703 });
704 return Ok(r);
705 }
706 Consumer::None => {
707 let res = engine.run(&parsed)?;
708 let mut r = OqxResult::scalar(consumer);
709 r.none = Some(match res {
710 oqx::OqxResult::None(b) => b,
711 _ => true,
712 });
713 return Ok(r);
714 }
715 Consumer::Collect | Consumer::First | Consumer::Single => {}
716 }
717
718 let top_distinct = parsed.distinct;
724 let top_values = parsed.values;
725 let user_select: Vec<SelectItem> = if top_values {
726 parsed
727 .select
728 .first()
729 .map(|it| match it {
730 SelectItem::Field { expr, lift, .. } => SelectItem::Field {
731 name: VALUE_KEY.to_owned(),
732 expr: expr.clone(),
733 lift: *lift,
734 },
735 SelectItem::Collect { op, .. } => SelectItem::Collect {
736 name: VALUE_KEY.to_owned(),
737 op: op.clone(),
738 },
739 })
740 .into_iter()
741 .collect()
742 } else {
743 parsed.select.clone()
744 };
745 let id_item = SelectItem::Field {
746 name: ID_KEY.to_owned(),
747 expr: Expr::Ident {
748 name: "$id".to_owned(),
749 },
750 lift: 0,
751 };
752 let path_item = SelectItem::Field {
753 name: PATH_KEY.to_owned(),
754 expr: Expr::Ident {
755 name: "$path".to_owned(),
756 },
757 lift: 0,
758 };
759 let mut select = vec![id_item, path_item];
760 select.extend(user_select);
761 let (top_limit, top_offset) = (parsed.limit.clone(), parsed.offset.clone());
765 let q = Query {
766 distinct: false,
767 values: false,
768 select,
769 limit: if consumer == Consumer::Collect {
770 None
771 } else {
772 parsed.limit.clone()
773 },
774 offset: if consumer == Consumer::Collect {
775 None
776 } else {
777 parsed.offset.clone()
778 },
779 ..parsed.clone()
780 };
781 let res = engine.run(&q)?;
782
783 if matches!(consumer, Consumer::First | Consumer::Single) {
784 let row = match res {
785 oqx::OqxResult::First(r) | oqx::OqxResult::Single(r) => r,
786 _ => None,
787 };
788 let mut out = OqxResult::scalar(consumer);
789 match row {
790 None => {
791 if top_values {
792 out.values = Some(Vec::new());
793 }
794 }
795 Some(r) => {
796 let hit = to_hit(r);
797 if top_values {
798 out.values = Some(vec![value_of(&hit).to_canonical_json()]);
799 } else {
800 out.hits = vec![without_value_key(hit)];
801 }
802 }
803 }
804 return Ok(out);
805 }
806
807 let mut rows: Vec<Value> = match res {
809 oqx::OqxResult::Collect(rows) => rows,
810 _ => Vec::new(),
811 };
812 let custom = parsed.order_by.as_ref().is_some_and(|o| !o.is_empty());
813 let cursor = opts.cursor.filter(|c| !c.is_empty() && !custom);
814 let eager = top_distinct || cursor.is_some();
819 if eager {
820 rows = rows.into_iter().map(to_hit).collect();
821 }
822 if top_distinct {
823 rows = dedup_hits_by_projection(rows);
824 }
825 let offset = const_bound(top_offset.as_ref(), "offset")?.unwrap_or(0);
826 let limit = const_bound(top_limit.as_ref(), "limit")?;
827 if offset > 0 || limit.is_some() {
828 let end = limit.map_or(rows.len(), |l| (offset + l).min(rows.len()));
829 if offset >= rows.len() {
830 rows.clear();
831 } else {
832 rows.truncate(end);
833 rows.drain(..offset);
834 }
835 }
836 let cap = opts.limit.unwrap_or(DEFAULT_LIMIT);
837 let mut page = rows;
838 if let Some(cursor) = cursor {
839 let parts = decode_cursor(cursor, "query", 2)?;
840 let (path, id) = (&parts[0], &parts[1]);
841 page.retain(|h| {
842 let hp = hit_str(h, "path");
843 let hi = hit_str(h, "id");
844 hp > *path || (hp == *path && hi > *id)
845 });
846 }
847 let truncated = page.len() > cap;
848 page.truncate(cap);
849 if !eager {
850 page = page.into_iter().map(to_hit).collect();
851 }
852 let cursor = if truncated && !custom {
853 page.last()
854 .map(|last| encode_cursor(&[&hit_str(last, "path"), &hit_str(last, "id")]))
855 } else {
856 None
857 };
858 let mut out = OqxResult::scalar(Consumer::Collect);
859 out.truncated = truncated;
860 out.cursor = cursor;
861 if top_values {
862 out.values = Some(
863 page.iter()
864 .map(|h| value_of(h).to_canonical_json())
865 .collect(),
866 );
867 } else {
868 out.hits = page.into_iter().map(without_value_key).collect();
869 }
870 Ok(out)
871}
872
873#[cfg(test)]
874mod tests {
875 use super::*;
876
877 #[test]
878 fn row_functions_become_self_methods() {
879 let q = oqx::parse_string(
880 "from blocks where text(\"x\") && under_heading(\"h\") && size(attrs) > 0 && doc.$path.startsWith(\"a\")",
881 )
882 .unwrap();
883 let r = rewrite_query(&q);
884 let Some(Where::And { parts }) = &r.r#where else {
885 panic!("and")
886 };
887 let Where::Scalar { expr } = &parts[0] else {
888 panic!("scalar")
889 };
890 assert!(
891 matches!(expr, Expr::Call { recv: Some(r), name, .. } if name == "text" && **r == self_ref())
892 );
893 let Where::Scalar { expr } = &parts[2] else {
894 panic!("scalar")
895 };
896 assert!(
897 matches!(expr, Expr::Binary { left, .. } if matches!(&**left, Expr::Call { recv: None, name, .. } if name == "size"))
898 );
899 }
900
901 #[test]
902 fn semantic_phrases_are_collected_distinct() {
903 let phrases = collect_semantic_phrases(
904 "select s: semantic(\"alpha\") from docs where semantic(\"alpha\") > 0.5 || nodes exists { where semantic(\"beta\") > 0 } order by semantic(\"gamma\") desc",
905 );
906 assert_eq!(phrases, ["alpha", "beta", "gamma"]);
907 assert!(collect_semantic_phrases("not a query {{").is_empty());
908 assert!(collect_semantic_phrases("from docs").is_empty());
909 }
910
911 #[test]
912 fn hits_peel_the_injected_columns() {
913 let mut o = oqx::Object::new();
914 o.insert(ID_KEY, Value::Str("d_1".into()));
915 o.insert(PATH_KEY, Value::Null);
916 o.insert("layer", Value::Str("canon".into()));
917 let hit = to_hit(Value::Object(o));
918 let ho = hit.as_object().unwrap();
919 assert_eq!(ho.keys().collect::<Vec<_>>(), ["id", "path", "layer"]);
920 assert_eq!(ho.get("path"), Some(&Value::Str(String::new())));
921 let mut o = oqx::Object::new();
923 o.insert(ID_KEY, Value::Str("d_1".into()));
924 o.insert(PATH_KEY, Value::Str("a.md".into()));
925 o.insert("id", Value::Number(7.0));
926 let hit = to_hit(Value::Object(o));
927 let ho = hit.as_object().unwrap();
928 assert_eq!(ho.keys().collect::<Vec<_>>(), ["id", "path"]);
929 assert_eq!(ho.get("id"), Some(&Value::Number(7.0)));
930 }
931
932 #[test]
933 fn distinct_dedups_by_user_projection_first_wins() {
934 let mk = |id: &str, t: &str| {
935 let mut o = oqx::Object::new();
936 o.insert("id", Value::Str(id.into()));
937 o.insert("path", Value::Str("p".into()));
938 o.insert("type", Value::Str(t.into()));
939 Value::Object(o)
940 };
941 let out = dedup_hits_by_projection(vec![mk("1", "a"), mk("2", "b"), mk("3", "a")]);
942 assert_eq!(out.len(), 2);
943 assert_eq!(hit_str(&out[0], "id"), "1");
944 assert_eq!(hit_str(&out[1], "id"), "2");
945 }
946
947 #[test]
948 fn top_level_bounds_must_be_literals() {
949 assert_eq!(const_bound(None, "limit").unwrap(), None);
950 assert_eq!(
951 const_bound(Some(&Expr::Lit(Value::Number(3.0))), "limit").unwrap(),
952 Some(3)
953 );
954 let e = const_bound(Some(&Expr::Lit(Value::Number(-1.0))), "offset").unwrap_err();
955 assert_eq!(e.code, "filter_invalid");
956 assert!(e.message.contains("top-level offset"));
957 }
958}