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