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::{BinaryOp, Expr, Query, SelectItem};
17use oqx::walk::{Clause, Node, VisitContext, Visitor, transform, visit};
18use oqx::{Consumer, Engine, InMemoryEngine, Value, build, resolve_aliases};
19use serde_json::{Map, Value as Json};
20
21use crate::context::{SemanticVec, StoreContext, Target, render_row_values, target_of};
22use crate::cursor::{decode_path_cursor, encode_cursor};
23use crate::error::{Result, SurfaceError};
24use crate::paths::reference_path;
25use crate::planner::{SqlitePlanner, root_target};
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(…)`. `refs(x)` is not
32/// one: it reads no row, so it stays a free function of the context. The
33/// planner treats every call alike (declined by the translator, flagged by
34/// `expr_may_raise`), `has_edge` and `refs` included.
35const ROW_FNS: [&str; 12] = [
36    "text",
37    "semantic",
38    "under",
39    "under_heading",
40    "within",
41    "under_kind",
42    "yaml_path",
43    "json_pointer",
44    "has_edge",
45    "has_anchor",
46    "child_count",
47    "parent_type",
48];
49
50const ID_KEY: &str = "__oqx_id";
51const PATH_KEY: &str = "__oqx_path";
52/// A hit is a store row (§1.4, 1.5). A bare root scan's rows are store rows by
53/// construction; a `follow` destination (`follow before` over a frontmatter
54/// list of paths), a `from E` re-projection or any other source can reach a
55/// scalar, which the injected `$id`/`$path` reads would render as the junk hit
56/// `{ id: "undefined", path: "" }`. For those queries the row itself is
57/// projected too (`$it`) under this key and [`to_hit`] fails the query when it
58/// is not a row. Only then: the clone of every projected row is a real cost
59/// here, and a root scan cannot need it. A `values` projection returns no
60/// hits, so it is exempt. Same rule as the reference's `SELF_ITEM`.
61const SELF_KEY: &str = "__oqx_self";
62/// The reserved key a top-level `values` projection's single item is renamed
63/// to, so it rides through id/path injection, paging and distinct as an
64/// ordinary field and is peeled off at the end.
65const VALUE_KEY: &str = "__oqx_value";
66
67/// `query`'s options.
68#[derive(Clone, Copy, Default)]
69pub struct QueryOptions<'a> {
70    /// The page cap (default 50).
71    pub limit: Option<usize>,
72    /// Resume after a truncated page's cursor.
73    pub cursor: Option<&'a str>,
74    /// The provider behind `semantic(...)`; `None` → `semantic_unavailable`
75    /// when the query names a phrase.
76    pub provider: Option<&'a dyn EmbeddingProvider>,
77    /// Force the pure in-memory engine (skip the tier-3 pushdown planner).
78    /// The default plans; the differential gate runs both and compares.
79    pub in_memory: bool,
80}
81
82/// §1.4 `OqxResult`.
83#[derive(Clone, Debug, PartialEq)]
84pub struct OqxResult {
85    pub hits: Vec<Json>,
86    pub truncated: bool,
87    pub cursor: Option<String>,
88    pub consumer: Consumer,
89    pub count: Option<f64>,
90    pub exists: Option<bool>,
91    pub none: Option<bool>,
92    /// A top-level `values` projection's bare values, in place of `hits`.
93    pub values: Option<Vec<Json>>,
94}
95
96impl OqxResult {
97    fn scalar(consumer: Consumer) -> Self {
98        Self {
99            hits: Vec::new(),
100            truncated: false,
101            cursor: None,
102            consumer,
103            count: None,
104            exists: None,
105            none: None,
106            values: None,
107        }
108    }
109
110    /// The wire shape: `{ hits, truncated, cursor, consumer, count?, exists?,
111    /// none?, values? }`.
112    #[must_use]
113    pub fn to_json(&self) -> Json {
114        let mut m = Map::new();
115        m.insert("hits".to_owned(), Json::Array(self.hits.clone()));
116        m.insert("truncated".to_owned(), Json::Bool(self.truncated));
117        m.insert(
118            "cursor".to_owned(),
119            self.cursor.clone().map_or(Json::Null, Json::String),
120        );
121        m.insert(
122            "consumer".to_owned(),
123            Json::String(self.consumer.as_str().to_owned()),
124        );
125        if let Some(n) = self.count {
126            m.insert("count".to_owned(), Value::Number(n).to_canonical_json());
127        }
128        if let Some(b) = self.exists {
129            m.insert("exists".to_owned(), Json::Bool(b));
130        }
131        if let Some(b) = self.none {
132            m.insert("none".to_owned(), Json::Bool(b));
133        }
134        if let Some(v) = &self.values {
135            m.insert("values".to_owned(), Json::Array(v.clone()));
136        }
137        Json::Object(m)
138    }
139}
140
141// ---- the `$self` rewrite -----------------------------------------------------------
142
143/// Paths on the surface are the reference form (`spec/surface` §1 "Paths",
144/// 2.0): `$path` and `$dst_path` read `/a/b.md`. A string literal a query
145/// compares with one of them — `$path == "a.md"`, `"a.md" != $dst_path`,
146/// `$path.startsWith("lab/")` — is rooted first, so both spellings match; a
147/// property or binding compared with a path is never touched (it is the
148/// author's value). The planner sees the rewritten tree too, so the pushed
149/// SQL and the in-memory engine agree by construction.
150const PATH_INTRINSICS: [&str; 2] = ["$path", "$dst_path"];
151
152fn is_path_read(e: &Expr) -> bool {
153    match e {
154        Expr::Ident { name, .. } | Expr::Outer { name, .. } | Expr::Member { name, .. } => {
155            PATH_INTRINSICS.contains(&name.as_str())
156        }
157        _ => false,
158    }
159}
160
161fn root_literal(e: Expr) -> Expr {
162    match e {
163        Expr::Lit {
164            value: Value::Str(s),
165            span,
166        } => Expr::Lit {
167            value: Value::Str(reference_path(&s)),
168            span,
169        },
170        other => other,
171    }
172}
173
174/// The literal-rooting rewrite of one expression (see [`rewrite_query`]).
175#[must_use]
176pub fn root_path_literals(e: Expr) -> Expr {
177    match e {
178        Expr::Binary {
179            op: op @ (BinaryOp::Eq | BinaryOp::Ne),
180            left,
181            right,
182            span,
183        } => {
184            if is_path_read(&left) {
185                Expr::Binary {
186                    op,
187                    left,
188                    right: Box::new(root_literal(*right)),
189                    span,
190                }
191            } else if is_path_read(&right) {
192                Expr::Binary {
193                    op,
194                    left: Box::new(root_literal(*left)),
195                    right,
196                    span,
197                }
198            } else {
199                Expr::Binary {
200                    op,
201                    left,
202                    right,
203                    span,
204                }
205            }
206        }
207        Expr::Call {
208            recv: Some(recv),
209            name,
210            mut args,
211            span,
212        } if name == "startsWith" && args.len() == 1 && is_path_read(&recv) => {
213            let arg = root_literal(args.remove(0));
214            Expr::Call {
215                recv: Some(recv),
216                name,
217                args: vec![arg],
218                span,
219            }
220        }
221        other => other,
222    }
223}
224
225/// Rewrite every row function in a parsed query to a `$self` method call, and
226/// root every string literal compared with a path intrinsic — one `transform`
227/// over the AST (`spec/oqx/AST.md` §5), every block included.
228#[must_use]
229pub fn rewrite_query(q: &Query) -> Query {
230    transform(q, &mut |e, _| match e {
231        Expr::Call {
232            recv: None,
233            name,
234            args,
235            span,
236        } if ROW_FNS.contains(&name.as_str()) => Expr::Call {
237            recv: Some(Box::new(build::ident("$self"))),
238            name,
239            args,
240            span,
241        },
242        other => root_path_literals(other),
243    })
244}
245
246// ---- the root row (spec/surface §1.1, 2.0) ----------------------------------------------
247
248/// The `$repo` error (removed in surface 2.0): the same text in both engines.
249pub const REPO_REMOVED_MESSAGE: &str = "`$repo` was removed in surface 2.0 — reach the repository's collections through the root row: `^docs` from a top-level row (one caret per enclosing block, or the absolute `0^docs`), and the repository id as `^$id` / `0^$id`";
250
251/// The bare-target error. `depth` is the scope depth the name was read at when
252/// known statically (the carets in the hint count it); `None` when raised at
253/// read time by the context, where the hint spells out the rule instead.
254#[must_use]
255pub fn bare_target_message(t: Target, depth: Option<usize>) -> String {
256    let name = t.as_str();
257    let noun = match t {
258        Target::Docs => "documents",
259        Target::Blocks => "blocks",
260        Target::Nodes => "nodes",
261        Target::Edges => "edges",
262    };
263    let carets = "^".repeat(depth.unwrap_or(1));
264    let how = if depth.is_none() {
265        format!("; one caret per enclosing block, or the absolute `0^{name}`")
266    } else {
267        String::new()
268    };
269    format!(
270        "`{name}` inside a block reads a property of the current row, which has none — did you mean `{carets}{name}` (the repository's {noun}{how})?"
271    )
272}
273
274/// The past-the-root error: `^docs` at the root scope (more carets than there
275/// are enclosing scopes) would read absent and count 0 — refused, naming the
276/// bare spelling. `levels` is the caret count written.
277#[must_use]
278pub fn past_root_message(t: Target, levels: usize) -> String {
279    let name = t.as_str();
280    let noun = match t {
281        Target::Docs => "documents",
282        Target::Blocks => "blocks",
283        Target::Nodes => "nodes",
284        Target::Edges => "edges",
285    };
286    let carets = "^".repeat(levels);
287    format!(
288        "`{carets}{name}` reaches past the root — there is no enclosing row at this depth; at the top level the repository's {noun} are the bare `{name}` (`{name} count {{ … }}`, `from {name}`, `entries({name})`)"
289    )
290}
291
292/// The repository is OQX's root row (`spec/oqx` 0.18): from a top-level row
293/// `^docs` reaches the documents (one caret per enclosing block, or the
294/// absolute `0^docs`) and `^$id` the repository id. Two spellings are refused
295/// before evaluation, by one walk over the parsed tree, so the error does not
296/// depend on which rows the data holds and the planned and in-memory paths
297/// agree trivially: `$repo` anywhere (bare, as a member head, `^$repo`,
298/// `0^$repo`), and a bare `docs`/`edges` at scope depth ≥ 1 — a property read
299/// of a row that has none. `blocks`/`nodes` are relations of some rows
300/// (`docs.blocks`, `docs.nodes`, `blocks.nodes`, `section.blocks`) and are
301/// judged by the context at read time with the same message. A `^<target>`
302/// with more carets than enclosing scopes (`^docs count { }` at the top level)
303/// would read absent and count 0: refused too, naming the bare spelling. Port
304/// of the reference's `checkRootSpellings`.
305pub fn check_root_spellings(q: &Query) -> Result<()> {
306    struct Check(Option<SurfaceError>);
307    impl Visitor for Check {
308        fn enter(&mut self, node: Node<'_>, ctx: &VisitContext<'_>) -> bool {
309            if self.0.is_some() {
310                return false;
311            }
312            match node {
313                Node::Expr(Expr::Ident { name, .. } | Expr::Outer { name, .. })
314                    if name == "$repo" =>
315                {
316                    self.0 = Some(SurfaceError::filter_invalid(REPO_REMOVED_MESSAGE, "OQX"));
317                }
318                Node::Expr(Expr::Ident { name, .. })
319                    if ctx.depth >= 1 && matches!(name.as_str(), "docs" | "edges") =>
320                {
321                    let t = Target::parse(name).expect("a target name");
322                    self.0 = Some(SurfaceError::filter_invalid(
323                        bare_target_message(t, Some(ctx.depth)),
324                        "OQX",
325                    ));
326                }
327                Node::Expr(Expr::Outer { name, levels, .. }) if *levels > ctx.depth => {
328                    if let Some(t) = Target::parse(name) {
329                        self.0 = Some(SurfaceError::filter_invalid(
330                            past_root_message(t, *levels),
331                            "OQX",
332                        ));
333                    }
334                }
335                _ => {}
336            }
337            self.0.is_none()
338        }
339    }
340    let mut check = Check(None);
341    visit(Node::Query(q), &mut check);
342    check.0.map_or(Ok(()), Err)
343}
344
345// ---- semantic phrases -----------------------------------------------------------------
346
347/// Collects the distinct `semantic("…")` phrases of a tree.
348struct Phrases(Vec<String>);
349
350impl Visitor for Phrases {
351    fn enter(&mut self, node: Node<'_>, _ctx: &VisitContext<'_>) -> bool {
352        if let Node::Expr(Expr::Call {
353            recv: None,
354            name,
355            args,
356            ..
357        }) = node
358            && name == "semantic"
359            && let Some(Expr::Lit {
360                value: Value::Str(s),
361                ..
362            }) = args.first()
363            && !self.0.contains(s)
364        {
365            self.0.push(s.clone());
366        }
367        true
368    }
369}
370
371/// The distinct phrases `semantic("…")` names (free calls in the raw parse);
372/// empty when the source does not parse.
373#[must_use]
374pub fn collect_semantic_phrases(source: &str) -> Vec<String> {
375    let Ok(q) = oqx::parse_string(source) else {
376        return Vec::new();
377    };
378    let mut phrases = Phrases(Vec::new());
379    visit(Node::Query(&q), &mut phrases);
380    phrases.0
381}
382
383// ---- name mentions --------------------------------------------------------------------
384
385/// Whether `name` is read anywhere in `q` other than as its source: a bare
386/// identifier or a `^`-escaped one in the `from` steps, `where`, `select`,
387/// `order by`, `follow`, `limit`/`offset`, or any nested block.
388fn mentions_outside_source(q: &Query, name: &str) -> bool {
389    struct Mentions<'n> {
390        name: &'n str,
391        found: bool,
392    }
393    impl Visitor for Mentions<'_> {
394        fn enter(&mut self, node: Node<'_>, ctx: &VisitContext<'_>) -> bool {
395            if self.found {
396                return false;
397            }
398            // The source itself (the root's `source` slot) is not a mention.
399            if ctx.clause == Some(Clause::Source) && ctx.path.len() == 1 {
400                return false;
401            }
402            if let Node::Expr(Expr::Ident { name, .. } | Expr::Outer { name, .. }) = node
403                && name == self.name
404            {
405                self.found = true;
406            }
407            !self.found
408        }
409    }
410    let mut m = Mentions { name, found: false };
411    visit(Node::Query(q), &mut m);
412    m.found
413}
414
415// ---- hits ----------------------------------------------------------------------------
416
417/// Whether the engine's top-level rows can be anything but store rows: a
418/// `follow` (a destination may be a property's value), a `from E`
419/// re-projection, or a source that is not a bare root scan.
420fn may_reach_non_rows(q: &Query) -> bool {
421    q.follow.is_some() || !q.from.is_empty() || root_target(&q.source).is_none()
422}
423
424/// The reference's `describeValue`: what a non-row hit was, for the error.
425fn describe_value(v: &Value) -> String {
426    match v {
427        Value::Undefined | Value::Null => "an absent value".to_owned(),
428        Value::Str(s) => format!(
429            "a string ({})",
430            serde_json::to_string(s).unwrap_or_default()
431        ),
432        Value::Number(_) => format!("a number ({v})"),
433        Value::Bool(b) => format!("a boolean ({b})"),
434        Value::Array(_) => "an array".to_owned(),
435        Value::Range(_) => "a range".to_owned(),
436        Value::Object(_) => "an object".to_owned(),
437    }
438}
439
440fn not_a_store_row(v: &Value) -> SurfaceError {
441    SurfaceError::filter_invalid(
442        format!(
443            "a hit must be a document, block, node or edge row — the query reached {}; to follow document references held in a property use refs(<field>)",
444            describe_value(v)
445        ),
446        "OQX",
447    )
448}
449
450/// A projected row as a hit: `{ id, path, ...rest }` with the injected
451/// columns peeled off (JavaScript's `String()` on the id, `""` for an absent
452/// path). The projection's values are rendered per §1.4 "rows as values": a
453/// store row nested in the result (an empty-projection `collect { }` and
454/// friends) becomes `{ id, path }`. When the row itself was projected under
455/// [`SELF_KEY`] and is not a store row, the query fails (§1.4: a hit is a
456/// store row).
457fn to_hit(row: Value) -> Result<Value> {
458    let Value::Object(mut o) = row else {
459        return Ok(Value::Object(oqx::Object::new()));
460    };
461    if let Some(me) = o.remove(SELF_KEY) {
462        if target_of(&me).is_none() {
463            return Err(not_a_store_row(&me));
464        }
465    }
466    let Value::Object(o) = render_row_values(Value::Object(o)) else {
467        return Ok(Value::Object(oqx::Object::new()));
468    };
469    let mut id = Value::Undefined;
470    let mut path = Value::Undefined;
471    let mut rest = Vec::new();
472    for (k, v) in o {
473        match k.as_str() {
474            ID_KEY => id = v,
475            PATH_KEY => path = v,
476            _ => rest.push((k, v)),
477        }
478    }
479    let mut hit = oqx::Object::with_capacity(rest.len() + 2);
480    hit.insert("id", Value::Str(id.to_string()));
481    hit.insert(
482        "path",
483        Value::Str(if path.is_absent() {
484            String::new()
485        } else {
486            path.to_string()
487        }),
488    );
489    for (k, v) in rest {
490        hit.insert(k, v);
491    }
492    Ok(Value::Object(hit))
493}
494
495fn hit_str(hit: &Value, key: &str) -> String {
496    hit.as_object()
497        .and_then(|o| o.get(key))
498        .map(|v| v.to_string())
499        .unwrap_or_default()
500}
501
502/// Top-level `select distinct`: dedup hits by their USER projection (every
503/// field but `id`/`path`), keeping the first. The key is the canonical JSON
504/// of the sorted `[key, value]` pairs (`JSON.stringify` in the reference).
505fn dedup_hits_by_projection(hits: Vec<Value>) -> Vec<Value> {
506    let mut seen: Vec<String> = Vec::new();
507    let mut out = Vec::new();
508    for h in hits {
509        let mut pairs: Vec<(String, Value)> = h
510            .as_object()
511            .map(|o| {
512                o.iter()
513                    .filter(|(k, _)| *k != "id" && *k != "path")
514                    .map(|(k, v)| (k.to_owned(), v.clone()))
515                    .collect()
516            })
517            .unwrap_or_default();
518        pairs.sort_by(|a, b| a.0.cmp(&b.0));
519        let key = Value::Array(
520            pairs
521                .into_iter()
522                .map(|(k, v)| Value::Array(vec![Value::Str(k), v]))
523                .collect(),
524        )
525        .to_canonical_json()
526        .to_string();
527        if seen.contains(&key) {
528            continue;
529        }
530        seen.push(key);
531        out.push(h);
532    }
533    out
534}
535
536/// A top-level `limit`/`offset` on the collect path is applied by the runner,
537/// so it must be a plain non-negative integer literal.
538fn const_bound(e: Option<&Expr>, word: &str) -> Result<Option<usize>> {
539    match e {
540        None => Ok(None),
541        Some(Expr::Lit {
542            value: Value::Number(n),
543            ..
544        }) if n.fract() == 0.0 && *n >= 0.0 && n.is_finite() => Ok(Some(*n as usize)),
545        Some(_) => Err(SurfaceError::filter_invalid(
546            format!("top-level {word} must be a non-negative integer literal"),
547            "OQX",
548        )),
549    }
550}
551
552fn value_of(hit: &Value) -> Value {
553    hit.as_object()
554        .and_then(|o| o.get(VALUE_KEY))
555        .cloned()
556        .unwrap_or(Value::Undefined)
557}
558
559fn without_value_key(hit: Value) -> Json {
560    hit.to_canonical_json()
561}
562
563// ---- the run --------------------------------------------------------------------------
564
565/// Run an OQX query against `repo_id` (§1.4).
566pub fn query(
567    store: &Store,
568    repo_id: &str,
569    source: &str,
570    opts: QueryOptions<'_>,
571) -> Result<OqxResult> {
572    // Without a provider the phrases stay unembedded and the context reports
573    // `filter_invalid` ("needs an embedding provider") when one is reached —
574    // after its target check, so `semantic()` on nodes names the targets
575    // (§9; the `query` tool's pre-check is what reports `semantic_unavailable`).
576    let phrases = collect_semantic_phrases(source);
577    let mut semantic: HashMap<String, SemanticVec> = HashMap::new();
578    if let Some(provider) = opts.provider.filter(|_| !phrases.is_empty()) {
579        for phrase in phrases {
580            let vec = provider
581                .embed_query(&phrase)
582                .map_err(|e| SurfaceError::new(e.code(), e.to_string()))?;
583            semantic.insert(
584                phrase,
585                SemanticVec {
586                    model: provider.model().to_owned(),
587                    vec: f32_to_blob(&vec),
588                },
589            );
590        }
591    }
592    let runner = Runner {
593        store,
594        repo_id,
595        semantic,
596        planned: !opts.in_memory,
597    };
598    run_inner(&runner, source, opts)
599}
600
601/// One query's engine: a fresh store context per run (the planned path gives
602/// the residual a context serving the produced rows as its root). A failure
603/// inside a property read or row function is the engine's own error (the
604/// context's `get` / `call_method` return `Err`); only a failed root scan,
605/// which the `root` seam cannot raise, is read back after the run.
606struct Runner<'a> {
607    store: &'a Store,
608    repo_id: &'a str,
609    semantic: HashMap<String, SemanticVec>,
610    planned: bool,
611}
612
613impl Runner<'_> {
614    /// Tier-3 pushdown reduces the scan in SQL and the in-memory engine
615    /// finishes the residual over the produced rows (a declined plan, or
616    /// `in_memory`, runs the whole query in memory over a full scan), so
617    /// results match a pure scan. This is `oqx::PlannedEngine::run` inlined:
618    /// the store context borrows the connection, so it cannot be the
619    /// `'static` context a `Plan` carries.
620    fn run(&self, q: &Query) -> Result<oqx::OqxResult> {
621        let conn = self.store.conn();
622        let ctx = StoreContext::new(conn, self.repo_id, self.semantic.clone());
623        let plan = if self.planned {
624            SqlitePlanner::new(conn, self.repo_id)
625                .try_plan(q, &[])
626                .map_err(|e| SurfaceError::other(format!("sqlite: {e}")))?
627        } else {
628            None
629        };
630        let (ctx, residual) = match plan {
631            // The residual's source is its one read of the rows root unless
632            // the query text itself names `__oqx_rows__` somewhere else (a
633            // `^`-reach, a `limit`, a nested block…): then every read must
634            // see the rows, so they are cloned out instead of moved.
635            Some(plan) => {
636                let once = !mentions_outside_source(&plan.residual, oqx::ROWS_ROOT);
637                let ctx = if once {
638                    ctx.with_rows_root_once(plan.rows)
639                } else {
640                    ctx.with_rows_root(plan.rows)
641                };
642                (ctx, Some(plan.residual))
643            }
644            None => (ctx, None),
645        };
646        let engine = InMemoryEngine::new(ctx);
647        let out = engine.run(residual.as_ref().unwrap_or(q), &[]);
648        // `root` has no error channel: a store failure during a root scan was
649        // served as an empty scan and wins over whatever the run made of it.
650        if let Some(failed) = engine.context().take_root_failure() {
651            return Err(failed.into());
652        }
653        // A root scan projected as a VALUE (`select all: ^docs`) was
654        // materialized into its rows by the engine (`DataContext::materialize`).
655        Ok(out?)
656    }
657}
658
659fn run_inner(engine: &Runner<'_>, source: &str, opts: QueryOptions<'_>) -> Result<OqxResult> {
660    // The query's `select` aliases are resolved HERE, once, before the runner
661    // renames/injects items (an engine evaluates the query it is given).
662    let raw = oqx::parse_string(source)?;
663    check_root_spellings(&raw)?;
664    let parsed = rewrite_query(&resolve_aliases(&raw)?);
665    let consumer = parsed.consumer;
666
667    match consumer {
668        Consumer::Exists => {
669            let res = engine.run(&parsed)?;
670            let mut r = OqxResult::scalar(consumer);
671            r.exists = Some(matches!(res, oqx::OqxResult::Exists(true)));
672            return Ok(r);
673        }
674        Consumer::Count => {
675            let res = engine.run(&parsed)?;
676            let mut r = OqxResult::scalar(consumer);
677            r.count = Some(match res {
678                oqx::OqxResult::Count(n) => n,
679                _ => 0.0,
680            });
681            return Ok(r);
682        }
683        Consumer::None => {
684            let res = engine.run(&parsed)?;
685            let mut r = OqxResult::scalar(consumer);
686            r.none = Some(match res {
687                oqx::OqxResult::None(b) => b,
688                _ => true,
689            });
690            return Ok(r);
691        }
692        Consumer::Collect | Consumer::First | Consumer::Single => {}
693    }
694
695    // collect / first / single: inject id + path so every hit carries them. A
696    // top-level `select distinct` is applied HERE, not in the engine (the
697    // injected id/path are unique per row and would defeat the engine's
698    // projection dedup). A top-level `values` projection runs as a RECORD
699    // projection whose single item is renamed to VALUE_KEY.
700    let top_distinct = parsed.distinct;
701    let top_values = parsed.values;
702    let user_select: Vec<SelectItem> = if top_values {
703        parsed
704            .select
705            .first()
706            .map(|it| match it {
707                SelectItem::Field {
708                    expr, lift, span, ..
709                } => SelectItem::Field {
710                    name: VALUE_KEY.to_owned(),
711                    expr: expr.clone(),
712                    lift: *lift,
713                    span: *span,
714                },
715                SelectItem::Collect { op, span, .. } => SelectItem::Collect {
716                    name: VALUE_KEY.to_owned(),
717                    op: op.clone(),
718                    span: *span,
719                },
720            })
721            .into_iter()
722            .collect()
723    } else {
724        parsed.select.clone()
725    };
726    let id_item = build::field(ID_KEY, build::ident("$id"));
727    let path_item = build::field(PATH_KEY, build::ident("$path"));
728    let mut select = vec![id_item, path_item];
729    if !top_values && may_reach_non_rows(&parsed) {
730        select.push(build::field(SELF_KEY, build::ident("$it")));
731    }
732    select.extend(user_select);
733    // On the collect path the query's own limit/offset is taken out of the
734    // engine query and applied after the runner's distinct; first/single keep
735    // theirs (the engine's offset-aware cap is exactly right for them).
736    let (top_limit, top_offset) = (parsed.limit.clone(), parsed.offset.clone());
737    let q = Query {
738        distinct: false,
739        values: false,
740        select,
741        limit: if consumer == Consumer::Collect {
742            None
743        } else {
744            parsed.limit.clone()
745        },
746        offset: if consumer == Consumer::Collect {
747            None
748        } else {
749            parsed.offset.clone()
750        },
751        ..parsed.clone()
752    };
753    let res = engine.run(&q)?;
754
755    if matches!(consumer, Consumer::First | Consumer::Single) {
756        let row = match res {
757            oqx::OqxResult::First(r) | oqx::OqxResult::Single(r) => r,
758            _ => None,
759        };
760        let mut out = OqxResult::scalar(consumer);
761        match row {
762            None => {
763                if top_values {
764                    out.values = Some(Vec::new());
765                }
766            }
767            Some(r) => {
768                let hit = to_hit(r)?;
769                if top_values {
770                    out.values = Some(vec![value_of(&hit).to_canonical_json()]);
771                } else {
772                    out.hits = vec![without_value_key(hit)];
773                }
774            }
775        }
776        return Ok(out);
777    }
778
779    // collect: keyset pagination on (path, id) when the order is the default.
780    let mut rows: Vec<Value> = match res {
781        oqx::OqxResult::Collect(rows) => rows,
782        _ => Vec::new(),
783    };
784    let custom = parsed.order_by.as_ref().is_some_and(|o| !o.is_empty());
785    let cursor = opts.cursor.filter(|c| !c.is_empty() && !custom);
786    // Rows become hits ({ id, path, … }) before anything inspects them — the
787    // distinct key and the cursor's (path, id) read the HIT — else only the
788    // page does (the offset/limit slice and the cap see plain rows), so a
789    // scan that projects thousands of rows shapes fifty.
790    let eager = top_distinct || cursor.is_some();
791    if eager {
792        rows = rows.into_iter().map(to_hit).collect::<Result<_>>()?;
793    }
794    if top_distinct {
795        rows = dedup_hits_by_projection(rows);
796    }
797    let offset = const_bound(top_offset.as_ref(), "offset")?.unwrap_or(0);
798    let limit = const_bound(top_limit.as_ref(), "limit")?;
799    if offset > 0 || limit.is_some() {
800        let end = limit.map_or(rows.len(), |l| (offset + l).min(rows.len()));
801        if offset >= rows.len() {
802            rows.clear();
803        } else {
804            rows.truncate(end);
805            rows.drain(..offset);
806        }
807    }
808    let cap = opts.limit.unwrap_or(DEFAULT_LIMIT);
809    let mut page = rows;
810    if let Some(cursor) = cursor {
811        // The cursor carries the hit's (rooted) path; a 1.x cursor is refused (§1.4).
812        let parts = decode_path_cursor(cursor, "query", 2)?;
813        let (path, id) = (&parts[0], &parts[1]);
814        page.retain(|h| {
815            let hp = hit_str(h, "path");
816            let hi = hit_str(h, "id");
817            hp > *path || (hp == *path && hi > *id)
818        });
819    }
820    let truncated = page.len() > cap;
821    page.truncate(cap);
822    if !eager {
823        page = page.into_iter().map(to_hit).collect::<Result<_>>()?;
824    }
825    let cursor = if truncated && !custom {
826        page.last()
827            .map(|last| encode_cursor(&[&hit_str(last, "path"), &hit_str(last, "id")]))
828    } else {
829        None
830    };
831    let mut out = OqxResult::scalar(Consumer::Collect);
832    out.truncated = truncated;
833    out.cursor = cursor;
834    if top_values {
835        out.values = Some(
836            page.iter()
837                .map(|h| value_of(h).to_canonical_json())
838                .collect(),
839        );
840    } else {
841        out.hits = page.into_iter().map(without_value_key).collect();
842    }
843    Ok(out)
844}
845
846#[cfg(test)]
847mod tests {
848    use super::*;
849    use oqx::ast::Where;
850
851    #[test]
852    fn row_functions_become_self_methods() {
853        let q = oqx::parse_string(
854            "from blocks where text(\"x\") && under_heading(\"h\") && size(attrs) > 0 && doc.$path.startsWith(\"a\")",
855        )
856        .unwrap();
857        let r = rewrite_query(&q);
858        let Some(Where::And { parts, .. }) = &r.r#where else {
859            panic!("and")
860        };
861        let Where::Scalar { expr, .. } = &parts[0] else {
862            panic!("scalar")
863        };
864        assert!(
865            matches!(expr, Expr::Call { recv: Some(r), name, .. } if name == "text" && **r == build::ident("$self"))
866        );
867        let Where::Scalar { expr, .. } = &parts[2] else {
868            panic!("scalar")
869        };
870        assert!(
871            matches!(expr, Expr::Binary { left, .. } if matches!(&**left, Expr::Call { recv: None, name, .. } if name == "size"))
872        );
873    }
874
875    #[test]
876    fn path_literals_are_rooted_only_against_path_reads() {
877        let q = oqx::parse_string(
878            "select n: ^docs collect { $path where $path == ^after || \"x.md\" != $dst_path } from docs where $path.startsWith(\"lab/\") && doc.$path == \"a.md\" && type == \"a.md\" && ^$path == \"b.md\"",
879        )
880        .unwrap();
881        let printed = oqx::print_query(&rewrite_query(&q)).unwrap();
882        assert!(printed.contains("$path.startsWith(\"/lab/\")"), "{printed}");
883        assert!(printed.contains("doc.$path == \"/a.md\""), "{printed}");
884        assert!(printed.contains("type == \"a.md\""), "{printed}");
885        assert!(printed.contains("^$path == \"/b.md\""), "{printed}");
886        assert!(printed.contains("\"/x.md\" != $dst_path"), "{printed}");
887        // a property compared with a path is the author's value: untouched
888        assert!(printed.contains("$path == ^after"), "{printed}");
889        // `contains`/`endsWith`/`matches` and the empty string are literal
890        let q = oqx::parse_string(
891            "from docs where $path.contains(\"a/\") && $path.endsWith(\".md\") && $path == \"\"",
892        )
893        .unwrap();
894        let printed = oqx::print_query(&rewrite_query(&q)).unwrap();
895        assert!(printed.contains("$path.contains(\"a/\")"), "{printed}");
896        assert!(printed.contains("$path.endsWith(\".md\")"), "{printed}");
897        assert!(printed.contains("$path == \"/\""), "{printed}");
898    }
899
900    #[test]
901    fn repo_and_bare_targets_are_refused_statically() {
902        let parse = |s: &str| oqx::parse_string(s).unwrap();
903        for q in [
904            "from docs where $repo",
905            "select r: size($repo.docs) from docs",
906            "select r: $repo.$id from docs",
907            "select r: nodes first { select x: ^$repo.$id values } from docs",
908            "select r: nodes first { select x: 0^$repo values } from docs",
909            "$repo.docs count { }",
910        ] {
911            let e = check_root_spellings(&parse(q)).unwrap_err();
912            assert_eq!(e.code, "filter_invalid", "{q}");
913            assert_eq!(e.message, REPO_REMOVED_MESSAGE, "{q}");
914        }
915        // a bare `docs`/`edges` at depth ≥ 1; the hint counts the carets
916        let e = check_root_spellings(&parse("select x: docs collect { } from docs")).unwrap_err();
917        assert_eq!(e.message, bare_target_message(Target::Docs, Some(1)));
918        assert!(
919            e.message
920                .contains("did you mean `^docs` (the repository's documents)?")
921        );
922        let e = check_root_spellings(&parse(
923            "select x: nodes collect { select y: size(edges) } from docs",
924        ))
925        .unwrap_err();
926        assert!(
927            e.message
928                .contains("did you mean `^^edges` (the repository's edges)?")
929        );
930        let e = check_root_spellings(&parse("docs collect { from docs }")).unwrap_err();
931        assert!(e.message.contains("`^docs`"));
932        // the root scope is depth 0: a bare target there is the scan
933        for q in [
934            "docs count { }",
935            "from docs where nodes exists { }",
936            "entries(docs) first { }",
937            "select n: size(^docs), id: ^$id, k: size(0^edges) from docs",
938            "from blocks where blocks exists { }", // judged by the context at read time
939        ] {
940            check_root_spellings(&parse(q)).unwrap_or_else(|e| panic!("{q}: {e}"));
941        }
942        // a `^target` reaching past the root
943        let e = check_root_spellings(&parse("^docs count { }")).unwrap_err();
944        assert_eq!(e.message, past_root_message(Target::Docs, 1));
945        assert_eq!(
946            e.message,
947            "`^docs` reaches past the root — there is no enclosing row at this depth; at the top level the repository's documents are the bare `docs` (`docs count { … }`, `from docs`, `entries(docs)`)"
948        );
949        let e = check_root_spellings(&parse("select $path from docs where ^^edges exists { }"))
950            .unwrap_err();
951        assert!(
952            e.message.contains("`^^edges` reaches past the root"),
953            "{}",
954            e.message
955        );
956        // the read-time message spells out the rule
957        let m = bare_target_message(Target::Blocks, None);
958        assert!(m.contains("did you mean `^blocks` (the repository's blocks; one caret per enclosing block, or the absolute `0^blocks`)?"), "{m}");
959    }
960
961    #[test]
962    fn semantic_phrases_are_collected_distinct() {
963        let phrases = collect_semantic_phrases(
964            "select s: semantic(\"alpha\") from docs where semantic(\"alpha\") > 0.5 || nodes exists { where semantic(\"beta\") > 0 } order by semantic(\"gamma\") desc",
965        );
966        assert_eq!(phrases, ["alpha", "beta", "gamma"]);
967        assert!(collect_semantic_phrases("not a query {{").is_empty());
968        assert!(collect_semantic_phrases("from docs").is_empty());
969    }
970
971    #[test]
972    fn hits_peel_the_injected_columns() {
973        let mut o = oqx::Object::new();
974        o.insert(ID_KEY, Value::Str("d_1".into()));
975        o.insert(PATH_KEY, Value::Null);
976        o.insert("layer", Value::Str("canon".into()));
977        let hit = to_hit(Value::Object(o)).unwrap();
978        let ho = hit.as_object().unwrap();
979        assert_eq!(ho.keys().collect::<Vec<_>>(), ["id", "path", "layer"]);
980        assert_eq!(ho.get("path"), Some(&Value::Str(String::new())));
981        // A user field named `id` overrides the injected one in place.
982        let mut o = oqx::Object::new();
983        o.insert(ID_KEY, Value::Str("d_1".into()));
984        o.insert(PATH_KEY, Value::Str("a.md".into()));
985        o.insert("id", Value::Number(7.0));
986        let hit = to_hit(Value::Object(o)).unwrap();
987        let ho = hit.as_object().unwrap();
988        assert_eq!(ho.keys().collect::<Vec<_>>(), ["id", "path"]);
989        assert_eq!(ho.get("id"), Some(&Value::Number(7.0)));
990    }
991
992    #[test]
993    fn a_hit_that_is_not_a_store_row_fails_the_query() {
994        // The row itself rides under SELF_KEY when the query can reach a
995        // non-row; a string there is the `follow before` shape (§1.4, 1.5).
996        let mut o = oqx::Object::new();
997        o.insert(ID_KEY, Value::Undefined);
998        o.insert(PATH_KEY, Value::Undefined);
999        o.insert(SELF_KEY, Value::Str("/timeline/kickoff.md".into()));
1000        let e = to_hit(Value::Object(o)).unwrap_err();
1001        assert_eq!(e.code, "filter_invalid");
1002        assert_eq!(
1003            e.message,
1004            "a hit must be a document, block, node or edge row — the query reached a string (\"/timeline/kickoff.md\"); to follow document references held in a property use refs(<field>)"
1005        );
1006        // A tagged row under SELF_KEY passes and the key is peeled off.
1007        let mut row = oqx::Object::new();
1008        row.insert("doc_id", Value::Str("d_1".into()));
1009        row.insert("path", Value::Str("a.md".into()));
1010        let row = crate::context::tag_row(row, crate::context::Target::Docs);
1011        let mut o = oqx::Object::new();
1012        o.insert(ID_KEY, Value::Str("d_1".into()));
1013        o.insert(PATH_KEY, Value::Str("a.md".into()));
1014        o.insert(SELF_KEY, row);
1015        let hit = to_hit(Value::Object(o)).unwrap();
1016        assert_eq!(
1017            hit.as_object().unwrap().keys().collect::<Vec<_>>(),
1018            ["id", "path"]
1019        );
1020        // Only queries that can reach a non-row project the row itself.
1021        let parse = |s: &str| oqx::parse_string(s).unwrap();
1022        assert!(!may_reach_non_rows(&parse(
1023            "from docs where layer == \"canon\""
1024        )));
1025        assert!(may_reach_non_rows(&parse("from docs follow before")));
1026        assert!(may_reach_non_rows(&parse("refs(\"/index.md\") first { }")));
1027    }
1028
1029    #[test]
1030    fn distinct_dedups_by_user_projection_first_wins() {
1031        let mk = |id: &str, t: &str| {
1032            let mut o = oqx::Object::new();
1033            o.insert("id", Value::Str(id.into()));
1034            o.insert("path", Value::Str("p".into()));
1035            o.insert("type", Value::Str(t.into()));
1036            Value::Object(o)
1037        };
1038        let out = dedup_hits_by_projection(vec![mk("1", "a"), mk("2", "b"), mk("3", "a")]);
1039        assert_eq!(out.len(), 2);
1040        assert_eq!(hit_str(&out[0], "id"), "1");
1041        assert_eq!(hit_str(&out[1], "id"), "2");
1042    }
1043
1044    #[test]
1045    fn top_level_bounds_must_be_literals() {
1046        assert_eq!(const_bound(None, "limit").unwrap(), None);
1047        assert_eq!(
1048            const_bound(
1049                Some(&Expr::Lit {
1050                    value: Value::Number(3.0),
1051                    span: oqx::Span::EMPTY
1052                }),
1053                "limit"
1054            )
1055            .unwrap(),
1056            Some(3)
1057        );
1058        let e = const_bound(
1059            Some(&Expr::Lit {
1060                value: Value::Number(-1.0),
1061                span: oqx::Span::EMPTY,
1062            }),
1063            "offset",
1064        )
1065        .unwrap_err();
1066        assert_eq!(e.code, "filter_invalid");
1067        assert!(e.message.contains("top-level offset"));
1068    }
1069}