Skip to main content

omgbase_surface/
query.rs

1//! The runner (`spec/surface/README.md` §1.4): parse the source, rewrite the
2//! row functions to `$self` methods, run it through the engine over the
3//! store context, and shape the engine's result into the surface's
4//! `OqxResult` (lean `{ id, path, … }` hits, keyset paging, the consumer
5//! scalars). Port of `packages/core/src/oqx-js/run.ts`.
6//!
7//! By default the tier-3 planner ([`crate::planner`]) pre-filters the scan in
8//! SQL and the in-memory engine finishes the residual over the produced rows;
9//! [`QueryOptions::in_memory`] forces the pure in-memory engine (the
10//! differential gate runs every case both ways — a planner must be invisible).
11
12use 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
25/// The default page size.
26pub const DEFAULT_LIMIT: usize = 50;
27
28/// Row-scoped domain functions: authored as free calls that implicitly
29/// reference the current row; rewritten to `$self.fn(…)`.
30const 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";
47/// The reserved key a top-level `values` projection's single item is renamed
48/// to, so it rides through id/path injection, paging and distinct as an
49/// ordinary field and is peeled off at the end.
50const VALUE_KEY: &str = "__oqx_value";
51
52/// `query`'s options.
53#[derive(Clone, Copy, Default)]
54pub struct QueryOptions<'a> {
55    /// The page cap (default 50).
56    pub limit: Option<usize>,
57    /// Resume after a truncated page's cursor.
58    pub cursor: Option<&'a str>,
59    /// The provider behind `semantic(...)`; `None` → `semantic_unavailable`
60    /// when the query names a phrase.
61    pub provider: Option<&'a dyn EmbeddingProvider>,
62    /// Force the pure in-memory engine (skip the tier-3 pushdown planner).
63    /// The default plans; the differential gate runs both and compares.
64    pub in_memory: bool,
65}
66
67/// §1.4 `OqxResult`.
68#[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    /// A top-level `values` projection's bare values, in place of `hits`.
78    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    /// The wire shape: `{ hits, truncated, cursor, consumer, count?, exists?,
96    /// none?, values? }`.
97    #[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
126// ---- the `$self` rewrite -----------------------------------------------------------
127
128fn 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/// Rewrite every row function in a parsed query to a `$self` method call.
277#[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
294// ---- semantic phrases -----------------------------------------------------------------
295
296fn 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/// The distinct phrases `semantic("…")` names (free calls in the raw parse);
384/// empty when the source does not parse.
385#[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
406// ---- name mentions --------------------------------------------------------------------
407
408fn 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
473/// Whether `name` is read anywhere in `q` other than as its source: a bare
474/// identifier or a `^`-escaped one in the `from` steps, `where`, `select`,
475/// `order by`, `follow`, `limit`/`offset`, or any nested block.
476fn 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
488// ---- hits ----------------------------------------------------------------------------
489
490/// A projected row as a hit: `{ id, path, ...rest }` with the injected
491/// columns peeled off (JavaScript's `String()` on the id, `""` for an absent
492/// path). The projection's values are rendered per §1.4 "rows as values": a
493/// store row nested in the result (an empty-projection `collect { }` and
494/// friends) becomes `{ id, path }`.
495fn 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
532/// Top-level `select distinct`: dedup hits by their USER projection (every
533/// field but `id`/`path`), keeping the first. The key is the canonical JSON
534/// of the sorted `[key, value]` pairs (`JSON.stringify` in the reference).
535fn 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
566/// A top-level `limit`/`offset` on the collect path is applied by the runner,
567/// so it must be a plain non-negative integer literal.
568fn 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
592// ---- the run --------------------------------------------------------------------------
593
594/// Run an OQX query against `repo_id` (§1.4).
595pub fn query(
596    store: &Store,
597    repo_id: &str,
598    source: &str,
599    opts: QueryOptions<'_>,
600) -> Result<OqxResult> {
601    // Without a provider the phrases stay unembedded and the context reports
602    // `filter_invalid` ("needs an embedding provider") when one is reached —
603    // after its target check, so `semantic()` on nodes names the targets
604    // (§9; the `query` tool's pre-check is what reports `semantic_unavailable`).
605    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
630/// One query's engine: a fresh store context per run (the planned path gives
631/// the residual a context serving the produced rows as its root). A failure
632/// inside a property read or row function is the engine's own error (the
633/// context's `get` / `call_method` return `Err`); only a failed root scan,
634/// which the `root` seam cannot raise, is read back after the run.
635struct Runner<'a> {
636    store: &'a Store,
637    repo_id: &'a str,
638    semantic: HashMap<String, SemanticVec>,
639    planned: bool,
640}
641
642impl Runner<'_> {
643    /// Tier-3 pushdown reduces the scan in SQL and the in-memory engine
644    /// finishes the residual over the produced rows (a declined plan, or
645    /// `in_memory`, runs the whole query in memory over a full scan), so
646    /// results match a pure scan. This is `oqx::PlannedEngine::run` inlined:
647    /// the store context borrows the connection, so it cannot be the
648    /// `'static` context a `Plan` carries.
649    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            // The residual's source is its one read of the rows root unless
661            // the query text itself names `__oqx_rows__` somewhere else (a
662            // `^`-reach, a `limit`, a nested block…): then every read must
663            // see the rows, so they are cloned out instead of moved.
664            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        // `root` has no error channel: a store failure during a root scan was
678        // served as an empty scan and wins over whatever the run made of it.
679        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    // collect / first / single: inject id + path so every hit carries them. A
719    // top-level `select distinct` is applied HERE, not in the engine (the
720    // injected id/path are unique per row and would defeat the engine's
721    // projection dedup). A top-level `values` projection runs as a RECORD
722    // projection whose single item is renamed to VALUE_KEY.
723    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    // On the collect path the query's own limit/offset is taken out of the
762    // engine query and applied after the runner's distinct; first/single keep
763    // theirs (the engine's offset-aware cap is exactly right for them).
764    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    // collect: keyset pagination on (path, id) when the order is the default.
808    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    // Rows become hits ({ id, path, … }) before anything inspects them — the
815    // distinct key and the cursor's (path, id) read the HIT — else only the
816    // page does (the offset/limit slice and the cap see plain rows), so a
817    // scan that projects thousands of rows shapes fifty.
818    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        // A user field named `id` overrides the injected one in place.
922        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}