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::{
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
27/// The default page size.
28pub const DEFAULT_LIMIT: usize = 50;
29
30/// Row-scoped domain functions: authored as free calls that implicitly
31/// reference the current row; rewritten to `$self.fn(…)`.
32const 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";
49/// The reserved key a top-level `values` projection's single item is renamed
50/// to, so it rides through id/path injection, paging and distinct as an
51/// ordinary field and is peeled off at the end.
52const VALUE_KEY: &str = "__oqx_value";
53
54/// `query`'s options.
55#[derive(Clone, Copy, Default)]
56pub struct QueryOptions<'a> {
57    /// The page cap (default 50).
58    pub limit: Option<usize>,
59    /// Resume after a truncated page's cursor.
60    pub cursor: Option<&'a str>,
61    /// The provider behind `semantic(...)`; `None` → `semantic_unavailable`
62    /// when the query names a phrase.
63    pub provider: Option<&'a dyn EmbeddingProvider>,
64    /// Force the pure in-memory engine (skip the tier-3 pushdown planner).
65    /// The default plans; the differential gate runs both and compares.
66    pub in_memory: bool,
67}
68
69/// §1.4 `OqxResult`.
70#[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    /// A top-level `values` projection's bare values, in place of `hits`.
80    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    /// The wire shape: `{ hits, truncated, cursor, consumer, count?, exists?,
98    /// none?, values? }`.
99    #[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
128// ---- the `$self` rewrite -----------------------------------------------------------
129
130fn 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/// Rewrite every row function in a parsed query to a `$self` method call.
286#[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
303// ---- semantic phrases -----------------------------------------------------------------
304
305fn 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/// The distinct phrases `semantic("…")` names (free calls in the raw parse);
398/// empty when the source does not parse.
399#[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
420// ---- name mentions --------------------------------------------------------------------
421
422fn 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
489/// Whether `name` is read anywhere in `q` other than as its source: a bare
490/// identifier or a `^`-escaped one in the `from` steps, `where`, `select`,
491/// `order by`, `follow`, `limit`/`offset`, or any nested block.
492fn 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
504// ---- hits ----------------------------------------------------------------------------
505
506/// A projected row as a hit: `{ id, path, ...rest }` with the injected
507/// columns peeled off (JavaScript's `String()` on the id, `""` for an absent
508/// path). The projection's values are rendered per §1.4 "rows as values": a
509/// store row nested in the result (an empty-projection `collect { }` and
510/// friends) becomes `{ id, path }`.
511fn 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
548/// Top-level `select distinct`: dedup hits by their USER projection (every
549/// field but `id`/`path`), keeping the first. The key is the canonical JSON
550/// of the sorted `[key, value]` pairs (`JSON.stringify` in the reference).
551fn 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
582/// A top-level `limit`/`offset` on the collect path is applied by the runner,
583/// so it must be a plain non-negative integer literal.
584fn 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
608// ---- the run --------------------------------------------------------------------------
609
610/// Run an OQX query against `repo_id` (§1.4).
611pub fn query(
612    store: &Store,
613    repo_id: &str,
614    source: &str,
615    opts: QueryOptions<'_>,
616) -> Result<OqxResult> {
617    // Without a provider the phrases stay unembedded and the context reports
618    // `filter_invalid` ("needs an embedding provider") when one is reached —
619    // after its target check, so `semantic()` on nodes names the targets
620    // (§9; the `query` tool's pre-check is what reports `semantic_unavailable`).
621    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
646/// One query's engine: a fresh store context per run (the planned path gives
647/// the residual a context serving the produced rows as its root). A failure
648/// inside a property read or row function is the engine's own error (the
649/// context's `get` / `call_method` return `Err`); only a failed root scan,
650/// which the `root` seam cannot raise, is read back after the run.
651struct Runner<'a> {
652    store: &'a Store,
653    repo_id: &'a str,
654    semantic: HashMap<String, SemanticVec>,
655    planned: bool,
656}
657
658impl Runner<'_> {
659    /// Tier-3 pushdown reduces the scan in SQL and the in-memory engine
660    /// finishes the residual over the produced rows (a declined plan, or
661    /// `in_memory`, runs the whole query in memory over a full scan), so
662    /// results match a pure scan. This is `oqx::PlannedEngine::run` inlined:
663    /// the store context borrows the connection, so it cannot be the
664    /// `'static` context a `Plan` carries.
665    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            // The residual's source is its one read of the rows root unless
677            // the query text itself names `__oqx_rows__` somewhere else (a
678            // `^`-reach, a `limit`, a nested block…): then every read must
679            // see the rows, so they are cloned out instead of moved.
680            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        // `root` has no error channel: a store failure during a root scan was
694        // served as an empty scan and wins over whatever the run made of it.
695        if let Some(failed) = engine.context().take_root_failure() {
696            return Err(failed.into());
697        }
698        // A root scan projected as a VALUE (`select all: $repo.docs`) was
699        // materialized into its rows by the engine (`DataContext::materialize`).
700        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    // collect / first / single: inject id + path so every hit carries them. A
737    // top-level `select distinct` is applied HERE, not in the engine (the
738    // injected id/path are unique per row and would defeat the engine's
739    // projection dedup). A top-level `values` projection runs as a RECORD
740    // projection whose single item is renamed to VALUE_KEY.
741    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    // On the collect path the query's own limit/offset is taken out of the
780    // engine query and applied after the runner's distinct; first/single keep
781    // theirs (the engine's offset-aware cap is exactly right for them).
782    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    // collect: keyset pagination on (path, id) when the order is the default.
826    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    // Rows become hits ({ id, path, … }) before anything inspects them — the
833    // distinct key and the cursor's (path, id) read the HIT — else only the
834    // page does (the offset/limit slice and the cap see plain rows), so a
835    // scan that projects thousands of rows shapes fifty.
836    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        // A user field named `id` overrides the injected one in place.
940        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}