Skip to main content

nedb_engine/
pgwire.rs

1// SPDX-FileCopyrightText: 2026 INTERCHAINED LLC
2// SPDX-License-Identifier: BUSL-1.1
3// NEDB · © 2026 INTERCHAINED LLC × Eth-Interchained × Vex (Claude Opus 5)
4
5//! A PostgreSQL wire-protocol endpoint for NEDB — reads **and** writes.
6//!
7//! # What this is
8//!
9//! A front door that speaks the PostgreSQL v3 wire protocol well enough that
10//! tools built for Postgres — `psql`, DBeaver, Metabase, Grafana, psycopg, any
11//! libpq client — can use a NEDB store with ordinary SQL. A documented subset
12//! of SQL is translated to NQL and to engine writes; everything else is
13//! refused with an error naming exactly what was not understood.
14//!
15//! It is **not** a claim of Postgres parity. It is a claim that the SQL people
16//! actually type works, and that the boundary is stated rather than discovered.
17//!
18//! # Why writes belong here
19//!
20//! The first cut of this module was read-only, on the reasoning that a NEDB
21//! write carries `caused_by`, valid-time bounds and idempotency, and none of
22//! that has a natural SQL spelling. That reasoning was wrong, and looking at
23//! the mapping is what made it obvious:
24//!
25//! | SQL | NEDB | and therefore |
26//! |---|---|---|
27//! | `INSERT` | a put | — |
28//! | `UPDATE … WHERE` | a NEW VERSION of each match | the prior value stays readable |
29//! | `DELETE … WHERE` | a tombstone | the deleted row stays in history |
30//!
31//! NEDB is append-only, so an `UPDATE` is *already* a versioned write and a
32//! `DELETE` is *already* a tombstone. Nothing is bent to fit. The consequence
33//! is the point of the whole endpoint:
34//!
35//! ```sql
36//! UPDATE orders SET total = 999 WHERE _id = 'o1';
37//! SELECT total FROM orders WHERE _id = 'o1';                  -- 999
38//! SELECT total FROM orders AS OF SYSTEM TIME 0 WHERE _id = 'o1';  -- 120
39//! ```
40//!
41//! Run the SQL you would run against Postgres, and the tamper-evident history
42//! is free. No triggers, no audit table, no application code.
43//!
44//! Provenance is reachable too: `_caused_by`, `_valid_from` and `_valid_to` are
45//! reserved INSERT columns, lifted out of the payload into the write itself.
46//!
47//! Writes are ON by default — that is the parity position. Set
48//! `NEDBD_PG_READ_ONLY=1` for the deployment where this door must never mutate
49//! anything.
50//!
51//! # Supported SQL
52//!
53//! ```sql
54//! SELECT * | col [, col]* | COUNT(*) | <agg>(col)
55//!   FROM <collection>
56//!   [ AS OF SYSTEM TIME <seq> ]     -- bridges to NQL's AS OF
57//!   [ WHERE <predicate> ]           -- the full NQL predicate surface
58//!   [ GROUP BY <col> ] [ HAVING <predicate> ]
59//!   [ ORDER BY <col> [ASC|DESC] (, ...) ] [ LIMIT <n> ] [ OFFSET <n> ]
60//!
61//! INSERT INTO <collection> (c1, c2) VALUES (v1, v2), (…) [RETURNING …]
62//! UPDATE <collection> SET c = v [, …] [WHERE <predicate>] [RETURNING …]
63//! DELETE FROM <collection> [WHERE <predicate>] [RETURNING …]
64//! ```
65//!
66//! Single-quoted SQL literals are rewritten to NQL's double-quoted form and
67//! `<>` to `!=`. Column projection is applied here, after NQL returns whole
68//! documents, because NQL is FROM-first and has no projection clause.
69//!
70//! That translation serves user collections. Statements that read the
71//! catalogue (`pg_catalog.*`, `information_schema.*`) go instead to the real
72//! SQL evaluator in `sqlselect` — joins, subqueries, `EXISTS`, `ARRAY(...)`,
73//! `ANY`/`ALL`, `UNION`, derived tables, `LATERAL`, aggregates, `CASE`, scalar
74//! functions — because that is what psql's `\d` family is written in. Every
75//! psql 17 backslash command that can succeed against an empty-of-features
76//! Postgres exits 0 here, verified by driving the real binary
77//! (`tests/test_psql_introspection.py`).
78//!
79//! Not supported on the user-collection path, each refused by name: JOIN,
80//! subqueries, CTEs, window functions, DDL, `TRUNCATE`, `GRANT`/`REVOKE`.
81//! `INSERT` requires an explicit column list, because NEDB is schemaless and
82//! there is no declared column order to infer.
83//!
84//! # Protocol coverage
85//!
86//! **Both** protocols are implemented:
87//!
88//! * the **simple query protocol** (`Q`) — what `psql` and libpq's `PQexec`
89//!   use, and therefore psycopg2, which interpolates parameters client-side;
90//! * the **extended query protocol** (`Parse`/`Bind`/`Describe`/`Execute`/
91//!   `Close`/`Sync`/`Flush`) — what psycopg3, asyncpg and the JDBC driver use
92//!   for every parameterised statement. Without it those three could not run a
93//!   single query, so "psql works" was a long way from "your framework works".
94//!
95//! Parameters arrive in text *and* binary format, prepared statements and
96//! portals are per-connection, and a row-capped `Execute` suspends its portal
97//! (`PortalSuspended`) so a JDBC `setFetchSize` pages instead of stalling.
98//!
99//! ## Parameter typing in a store with no schema
100//!
101//! The extended protocol needs types for `$1..$n`, which a relational server
102//! reads out of its catalogue. NEDB has none — so the types are sampled from
103//! the documents already stored, and the stored data *is* the schema. Where a
104//! placeholder sits in a clause rather than beside a column
105//! (`AS OF SYSTEM TIME $1`, `LIMIT $1`) the grammar supplies the type instead,
106//! and an aggregate column is typed from what the aggregate means: a `COUNT` is
107//! an integer, an `AVG` fractional.
108//!
109//! This is not polish. A client that declares its own parameter types
110//! (psycopg3, JDBC) is believed and only its unspecified slots are inferred —
111//! but asyncpg declares none, asks, and then **refuses the call client-side**
112//! if the answer is wrong. Advertising "text" for everything does not degrade
113//! gracefully there; it fails with `expected str, got int` before a query is
114//! ever sent.
115//!
116//! SSL is declined (`N`), so connections are cleartext — hence the loopback
117//! default.
118//!
119//! Authentication mirrors the HTTP surface: with `NEDBD_TOKEN` set the password
120//! must equal it; otherwise any connection is accepted.
121//!
122//! Still outside the boundary, and refused by name: SQL-level cursors
123//! (`DECLARE`/`FETCH`), window functions, set operations on the
124//! user-collection path, and binary *result* format for a column whose stored
125//! values disagree about their type across documents.
126//!
127//! # What an ORM needs, and what it cost to learn
128//!
129//! Speaking psql is not speaking to a framework, and the difference was three
130//! defects deep. SQLAlchemy could not CONNECT (its dialect opens with
131//! `select pg_catalog.version()`, which a table of exact spellings missed); its
132//! reflection needed `GROUP BY` and `array_agg(x ORDER BY y)`; and a QUALIFIED
133//! column in a `WHERE` clause returned ZERO ROWS — silently — because NQL
134//! looks a field up flat and no document has a field named `orders.status`.
135//! Every ORM qualifies its predicates, so every filtered query lied.
136//!
137//! None of that was visible to psql, which is why
138//! `tests/pgwire_suite.py` drives asyncpg, SQLAlchemy and node-postgres
139//! against a live daemon on every push.
140
141use std::collections::HashMap;
142use std::sync::Arc;
143
144use serde_json::Value;
145use tokio::io::{AsyncReadExt, AsyncWriteExt};
146use tokio::net::{TcpListener, TcpStream};
147
148use crate::db::Db;
149
150// ── Postgres type OIDs we hand out ──────────────────────────────────────────
151const OID_BOOL: i32 = 16;
152const OID_INT8: i32 = 20;
153const OID_FLOAT8: i32 = 701;
154const OID_TEXT: i32 = 25;
155
156const PROTO_V3: i32 = 196_608; // 3.0 << 16
157const SSL_REQUEST: i32 = 80_877_103;
158const GSS_REQUEST: i32 = 80_877_104;
159const CANCEL_REQUEST: i32 = 80_877_102;
160
161/// How a caller resolves a database name to an open `Db`.
162///
163/// A trait object rather than a concrete handle so this module does not depend
164/// on `server::Manager` — which keeps the protocol code unit-testable against a
165/// plain `Db` with no HTTP stack in the way.
166pub trait DbResolver: Send + Sync + 'static {
167    /// Look up an open database by the name the client connected with.
168    ///
169    /// MAY BLOCK. The implementation is allowed to take a lock, so this is
170    /// always called from `spawn_blocking` — never on an async worker. Taking
171    /// a tokio `RwLock::blocking_read()` on a runtime thread panics outright
172    /// ("Cannot block the current thread from within a runtime"), which is
173    /// exactly how the first cut of this failed.
174    fn resolve(&self, name: &str) -> Option<Arc<Db>>;
175    /// The bearer token, when one is configured. `None` = open access.
176    fn token(&self) -> Option<String> {
177        None
178    }
179}
180
181// ── wire encoding helpers ───────────────────────────────────────────────────
182
183struct Out(Vec<u8>);
184
185impl Out {
186    fn msg(tag: u8) -> Self {
187        // Tag, then a 4-byte length placeholder patched in `finish`.
188        Out(vec![tag, 0, 0, 0, 0])
189    }
190    fn i16(&mut self, v: i16) { self.0.extend_from_slice(&v.to_be_bytes()); }
191    fn i32(&mut self, v: i32) { self.0.extend_from_slice(&v.to_be_bytes()); }
192    fn cstr(&mut self, s: &str) {
193        // A NUL inside an identifier would truncate the field and desynchronise
194        // the stream, so strip rather than trust.
195        self.0.extend_from_slice(s.replace('\0', "").as_bytes());
196        self.0.push(0);
197    }
198    fn bytes(&mut self, b: &[u8]) { self.0.extend_from_slice(b); }
199    /// Patch the length prefix (which covers the length field itself, not the tag).
200    fn finish(mut self) -> Vec<u8> {
201        let len = (self.0.len() - 1) as i32;
202        self.0[1..5].copy_from_slice(&len.to_be_bytes());
203        self.0
204    }
205}
206
207fn err_msg(code: &str, message: &str) -> Vec<u8> {
208    let mut m = Out::msg(b'E');
209    m.bytes(b"S"); m.cstr("ERROR");
210    m.bytes(b"C"); m.cstr(code);
211    m.bytes(b"M"); m.cstr(message);
212    m.0.push(0);
213    m.finish()
214}
215
216fn ready() -> Vec<u8> {
217    let mut m = Out::msg(b'Z');
218    m.bytes(b"I"); // idle, not in a transaction
219    m.finish()
220}
221
222fn command_complete(tag: &str) -> Vec<u8> {
223    let mut m = Out::msg(b'C');
224    m.cstr(tag);
225    m.finish()
226}
227
228// ── SQL → NQL translation ───────────────────────────────────────────────────
229
230/// One output column: the key to read from the row, and the name to show.
231///
232/// The two differ for aggregates. NQL answers `SUM(total)` with a row holding
233/// `sum_total` (plus `count` and a legacy `value` alias), while SQL callers
234/// expect a single column called `sum`. Carrying both halves keeps NEDB's
235/// internal key names off the wire — the first cut leaked `['count','value']`
236/// out of a `SELECT COUNT(*)`, which is two columns where SQL promises one.
237#[derive(Debug, PartialEq, Clone)]
238pub struct Col {
239    pub src: String,
240    pub out: String,
241}
242
243impl Col {
244    fn same(name: &str) -> Self {
245        Col { src: name.to_string(), out: name.to_string() }
246    }
247    fn renamed(src: &str, out: &str) -> Self {
248        Col { src: src.to_string(), out: out.to_string() }
249    }
250}
251
252/// What a translated statement asks for.
253///
254/// The write variants exist because SQL's write semantics and NEDB's storage
255/// model line up almost exactly, which was not obvious until it was written
256/// down:
257///
258/// | SQL | NEDB |
259/// |---|---|
260/// | `INSERT` | a put |
261/// | `UPDATE … WHERE` | a NEW VERSION of each matching document |
262/// | `DELETE … WHERE` | a tombstone |
263///
264/// NEDB is append-only, so an `UPDATE` is *already* a versioned write and a
265/// `DELETE` is *already* a tombstone. Nothing is being bent to fit. The
266/// consequence is the thing worth selling: run the SQL you would run against
267/// Postgres, and the tamper-evident history falls out for free — the prior
268/// value is still readable with `AS OF SYSTEM TIME`.
269#[derive(Debug, PartialEq)]
270pub enum Stmt {
271    /// Run this NQL, then project these columns (empty = all).
272    Query { nql: String, project: Vec<Col> },
273    /// `INSERT INTO coll (cols) VALUES (…), (…) [RETURNING …]`
274    Insert { coll: String, rows: Vec<InsertRow>, returning: Vec<Col> },
275    /// `UPDATE coll SET … [WHERE …] [RETURNING …]` — a new version per match.
276    Update { coll: String, set: Vec<(String, Value)>, nql: String, returning: Vec<Col> },
277    /// `DELETE FROM coll [WHERE …] [RETURNING …]` — a tombstone per match.
278    Delete { coll: String, nql: String, returning: Vec<Col> },
279    /// Answer from a fixed table — the handshake queries clients send on connect.
280    Canned { cols: Vec<String>, row: Vec<String> },
281    /// Nothing to do (empty statement, or a SET the client does not need honoured).
282    Ok(&'static str),
283}
284
285/// One row of an `INSERT`: an explicit id when the statement supplied one, the
286/// document body, and optional provenance lifted out of reserved columns.
287#[derive(Debug, PartialEq, Clone)]
288pub struct InsertRow {
289    /// From an `_id` or `id` column. `None` means the server assigns one.
290    pub id: Option<String>,
291    pub doc: serde_json::Map<String, Value>,
292    /// From a `_caused_by` column — the causal parents, so provenance is
293    /// reachable from SQL rather than only from the HTTP API.
294    pub caused_by: Vec<String>,
295    pub valid_from: Option<String>,
296    pub valid_to: Option<String>,
297}
298
299/// Strip SQL comments and collapse whitespace, so the matchers below can be
300/// simple without being fragile about formatting.
301fn normalise(sql: &str) -> String {
302    let mut out = String::with_capacity(sql.len());
303    let mut chars = sql.chars().peekable();
304    let mut in_s = false;
305    while let Some(c) = chars.next() {
306        if in_s {
307            out.push(c);
308            if c == '\'' { in_s = false; }
309            continue;
310        }
311        match c {
312            '\'' => { in_s = true; out.push(c); }
313            '-' if chars.peek() == Some(&'-') => {
314                // line comment
315                for n in chars.by_ref() { if n == '\n' { break; } }
316                out.push(' ');
317            }
318            '/' if chars.peek() == Some(&'*') => {
319                chars.next();
320                let mut prev = ' ';
321                while let Some(n) = chars.next() {
322                    if prev == '*' && n == '/' { break; }
323                    prev = n;
324                }
325                out.push(' ');
326            }
327            _ => out.push(c),
328        }
329    }
330    out.split_whitespace().collect::<Vec<_>>().join(" ")
331}
332
333/// Rewrite SQL literal/operator spellings into NQL's.
334///
335/// Only `'…'` → `"…"` and `<>` → `!=`. Done with an explicit scan rather than a
336/// regex so a quote inside a string cannot be mistaken for a delimiter: SQL
337/// escapes an embedded quote by doubling it (`'it''s'`), and that has to become
338/// a single character inside the NQL string rather than terminating it.
339/// Drop the table qualifier from every column reference in a clause tail.
340///
341/// # The silent wrong answer this removes
342///
343/// NQL has no notion of a qualifier: `field_value` looks a field up FLAT, in
344/// one map. So `WHERE orders.status = 'paid'` asked for a field literally
345/// named `orders.status`, no document had one, and the query returned ZERO
346/// ROWS — with no error and no warning, an empty result that reads exactly
347/// like "you have no paid orders".
348///
349/// Every ORM qualifies its predicates. SQLAlchemy emits
350/// `SELECT orders._id FROM orders WHERE orders.status = 'paid'` for the most
351/// ordinary filter there is, so EVERY filtered query answered empty, `.get(pk)`
352/// answered `None`, and `filter_by` answered nothing. The select list had
353/// always stripped qualifiers; the tail was "handed to the NQL parser
354/// unchanged", which is right for the clause GRAMMAR and wrong for a name NQL
355/// cannot interpret.
356///
357/// # Why a mismatched qualifier is an ERROR, not a strip
358///
359/// A qualifier naming something other than this statement's own collection
360/// means the query referenced a relation that is not in its FROM clause.
361/// Stripping it would answer with rows from the one relation that IS there —
362/// a different wrong answer wearing the same empty-looking clothes. Aliases
363/// are refused on this path already, so the collection's own name is the only
364/// qualifier that can be correct.
365///
366/// Runs BEFORE `sql_literals_to_nql`, so only SQL's single-quoted strings have
367/// to be skipped — the rewrite to NQL's double-quoted form has not happened
368/// yet, and a qualifier can never appear inside a literal.
369fn strip_column_qualifiers(
370    tail: &str,
371    coll: &str,
372    alias: Option<&str>,
373) -> Result<String, String> {
374    let bare = coll.rsplit('.').next().unwrap_or(coll);
375    let b: Vec<char> = tail.chars().collect();
376    let mut out = String::with_capacity(tail.len());
377    let mut i = 0usize;
378    let ident_start = |c: char| c.is_alphabetic() || c == '_';
379    let ident_char = |c: char| c.is_alphanumeric() || c == '_';
380
381    while i < b.len() {
382        // A single-quoted literal is copied through untouched.
383        if b[i] == '\'' {
384            out.push(b[i]);
385            i += 1;
386            while i < b.len() {
387                out.push(b[i]);
388                if b[i] == '\'' {
389                    // A doubled '' is one literal quote, not a close.
390                    if b.get(i + 1) == Some(&'\'') {
391                        out.push('\'');
392                        i += 2;
393                        continue;
394                    }
395                    i += 1;
396                    break;
397                }
398                i += 1;
399            }
400            continue;
401        }
402        // A double-quoted run is copied through too. NQL reads double quotes
403        // as a STRING delimiter rather than an identifier one, so a SQL
404        // delimited identifier is a genuine divergence — but it already fails
405        // LOUDLY in the NQL parser ("expected field name, got Str"), and a
406        // loud failure is not this function's problem to solve quietly.
407        if b[i] == '"' {
408            out.push(b[i]);
409            i += 1;
410            while i < b.len() {
411                out.push(b[i]);
412                if b[i] == '"' { i += 1; break; }
413                i += 1;
414            }
415            continue;
416        }
417        if !ident_start(b[i]) {
418            // A number like `1.5` starts with a digit, so it never enters the
419            // identifier branch and its dot is never touched.
420            out.push(b[i]);
421            i += 1;
422            continue;
423        }
424
425        let start = i;
426        while i < b.len() && ident_char(b[i]) {
427            i += 1;
428        }
429        let word: String = b[start..i].iter().collect();
430
431        // `qual.field` — a dot followed immediately by another identifier.
432        if b.get(i) == Some(&'.') && b.get(i + 1).is_some_and(|c| ident_start(*c)) {
433            let fstart = i + 1;
434            let mut j = fstart;
435            while j < b.len() && ident_char(b[j]) {
436                j += 1;
437            }
438            let field: String = b[fstart..j].iter().collect();
439            // A qualified FUNCTION call (`pg_catalog.something(`) is left
440            // exactly as written: this path does not implement functions at
441            // all, and NQL's own refusal names the function, which is more use
442            // to the reader than a claim about relations.
443            let is_call = b[j..].iter().find(|c| !c.is_whitespace()) == Some(&'(');
444            if is_call {
445                out.push_str(&word);
446                out.push('.');
447                out.push_str(&field);
448                i = j;
449                continue;
450            }
451            let matches_alias = alias.is_some_and(|a| word.eq_ignore_ascii_case(a));
452            if matches_alias || word.eq_ignore_ascii_case(bare) || word.eq_ignore_ascii_case(coll) {
453                out.push_str(&field);
454                i = j;
455                continue;
456            }
457            return Err(format!(
458                "no table or alias named {:?} in this query — this statement reads \
459                 {:?}{}, and a qualifier naming anything else would have to be \
460                 answered from a relation that is not in its FROM clause",
461                word,
462                bare,
463                alias.map(|a| format!(" (aliased {:?})", a)).unwrap_or_default()
464            ));
465        }
466        out.push_str(&word);
467    }
468    Ok(out)
469}
470
471/// Rewrite `SELECT count(*) FROM (<inner>) [AS] alias` into a flat count over
472/// the inner query's own collection and predicate — or `None` when the shapes
473/// do not permit it.
474///
475/// `None` is a REFUSAL, never a fallback: every caller reports the boundary
476/// rather than trying something else, because the alternative to an exact
477/// count is a wrong one.
478fn flatten_count_of_subquery(projection: &str, rest: &str) -> Option<String> {
479    // The outer select list must be nothing but `count(*)`, optionally
480    // aliased. Any other column would have to come from the derived table's
481    // output, which a flat count does not produce.
482    let (outer_expr, outer_alias) = split_output_alias(projection.trim());
483    let ou = outer_expr.to_uppercase().replace(' ', "");
484    if ou != "COUNT(*)" {
485        return None;
486    }
487
488    // Take the balanced parenthesised span, honouring literals so a `)` inside
489    // a string cannot close it early.
490    let b: Vec<char> = rest.chars().collect();
491    let mut depth = 0i32;
492    let mut in_s = false;
493    let mut end = None;
494    for (i, &c) in b.iter().enumerate() {
495        match c {
496            '\'' => in_s = !in_s,
497            '(' if !in_s => depth += 1,
498            ')' if !in_s => {
499                depth -= 1;
500                if depth == 0 {
501                    end = Some(i);
502                    break;
503                }
504            }
505            _ => {}
506        }
507    }
508    let end = end?;
509    let inner = b[1..end].iter().collect::<String>().trim().to_string();
510
511    // Nothing may follow the derived table but its alias — a join or a second
512    // FROM item changes what is being counted.
513    let trailing = b[end + 1..].iter().collect::<String>();
514    let (_alias, after) = split_table_alias(trailing.trim());
515    if !after.trim().is_empty() {
516        return None;
517    }
518
519    let iu = inner.to_uppercase();
520    if !iu.starts_with("SELECT") {
521        return None;
522    }
523    // Each of these would make the inner row count differ from the flat one.
524    for kw in ["LIMIT", "OFFSET", "GROUP BY", "HAVING", "UNION", "INTERSECT", "EXCEPT", "JOIN"] {
525        if find_kw(&iu, kw).is_some() {
526            return None;
527        }
528    }
529    if find_kw(&iu, "DISTINCT").is_some() {
530        return None;
531    }
532    // An inner aggregate already reduced the rows to one.
533    let inner_from = find_kw(&iu, "FROM")?;
534    let inner_list = inner[..inner_from].to_uppercase();
535    for agg in ["COUNT(", "SUM(", "AVG(", "MIN(", "MAX(", "ARRAY_AGG(", "STRING_AGG("] {
536        if inner_list.contains(agg) {
537            return None;
538        }
539    }
540    // A nested derived table is not walked — one level is the claim.
541    let inner_rest = inner[inner_from + 4..].trim();
542    if inner_rest.starts_with('(') {
543        return None;
544    }
545
546    // `ORDER BY` cannot change a count, so it is dropped rather than refused.
547    let mut tail = inner_rest.to_string();
548    let tu = tail.to_uppercase();
549    if let Some(ob) = find_kw(&tu, "ORDER BY") {
550        tail = tail[..ob].trim_end().to_string();
551    }
552    Some(format!(
553        "SELECT count(*){} FROM {}",
554        outer_alias.map(|a| format!(" AS {}", a)).unwrap_or_default(),
555        tail
556    ))
557}
558
559/// Words that begin a clause and can therefore never be a bare table alias.
560///
561/// `AS` is absent on purpose: it introduces an alias, and `AS OF` is
562/// disambiguated by looking at the word after it.
563const CLAUSE_WORDS: &[&str] = &[
564    "WHERE", "GROUP", "ORDER", "LIMIT", "OFFSET", "HAVING", "FOR", "VALID",
565    "TRACE", "TRAVERSE", "SEARCH", "RETURNING", "UNION", "INTERSECT", "EXCEPT",
566    "JOIN", "LEFT", "RIGHT", "INNER", "FULL", "CROSS", "ON", "USING", "SET",
567];
568
569/// Take a table alias off the front of a clause tail: `FROM orders o WHERE …`.
570///
571/// Returns the alias and the rest of the tail. The alias is REMOVED because
572/// NQL has no table-alias syntax and would report an "unexpected token" on it
573/// — which is how `FROM orders o` used to fail. Removing it here and teaching
574/// `strip_column_qualifiers` to accept it is what makes `SELECT o.status FROM
575/// orders o` work at all.
576///
577/// `AS OF SYSTEM TIME` also starts with `AS`, so the word AFTER `AS` decides:
578/// `AS OF` is a time-travel clause, anything else is an alias.
579fn split_table_alias(tail: &str) -> (Option<String>, &str) {
580    let t = tail.trim_start();
581    let first_end = t.find(char::is_whitespace).unwrap_or(t.len());
582    let first = &t[..first_end];
583    let fu = first.to_uppercase();
584
585    if fu == "AS" {
586        let rest = t[first_end..].trim_start();
587        let end = rest.find(char::is_whitespace).unwrap_or(rest.len());
588        let word = &rest[..end];
589        if word.eq_ignore_ascii_case("OF") {
590            return (None, t); // `AS OF …`, not an alias
591        }
592        if word.is_empty() {
593            return (None, t);
594        }
595        return (Some(word.trim_matches('"').to_string()), rest[end..].trim_start());
596    }
597    if first.is_empty() || CLAUSE_WORDS.contains(&fu.as_str()) {
598        return (None, t);
599    }
600    // A bare identifier here can only be an alias — the collection name was
601    // already consumed by the caller.
602    if first.chars().next().is_some_and(|c| c.is_alphabetic() || c == '_' || c == '"') {
603        return (Some(first.trim_matches('"').to_string()), t[first_end..].trim_start());
604    }
605    (None, t)
606}
607
608/// Split on a delimiter that is at PAREN DEPTH ZERO and outside a literal.
609///
610/// `projection.split(',')` cuts `SUM(a, b)` in half; a select list is not a
611/// flat comma list once it can contain calls.
612fn split_top_level(s: &str, delim: char) -> Vec<String> {
613    let mut out = vec![];
614    let mut cur = String::new();
615    let mut depth = 0i32;
616    let mut in_s = false;
617    let mut in_d = false;
618    for c in s.chars() {
619        match c {
620            '\'' if !in_d => { in_s = !in_s; cur.push(c); }
621            '"' if !in_s => { in_d = !in_d; cur.push(c); }
622            '(' if !in_s && !in_d => { depth += 1; cur.push(c); }
623            ')' if !in_s && !in_d => { depth -= 1; cur.push(c); }
624            c if c == delim && depth == 0 && !in_s && !in_d => {
625                out.push(std::mem::take(&mut cur));
626            }
627            _ => cur.push(c),
628        }
629    }
630    out.push(cur);
631    out
632}
633
634/// Split `expr AS name` / `expr name` into the expression and its output name.
635///
636/// The alias is the name the CLIENT will look the column up by — SQLAlchemy
637/// reads `count(*) AS count_1` back as `count_1`, so dropping the alias and
638/// returning a column called `count` hands it a result it cannot find.
639fn split_output_alias(p: &str) -> (&str, Option<&str>) {
640    let pu = p.to_uppercase();
641    if let Some(at) = find_kw(&pu, "AS") {
642        let alias = p[at + 2..].trim().trim_matches('"');
643        if !alias.is_empty() {
644            return (p[..at].trim(), Some(alias));
645        }
646    }
647    // A bare alias: `count(*) count_1`. Only after a closing paren or a plain
648    // identifier, and never when the tail is itself part of the expression —
649    // so the split point is the LAST whitespace outside any parenthesis.
650    let b: Vec<char> = p.chars().collect();
651    let mut depth = 0i32;
652    let mut in_s = false;
653    let mut cut = None;
654    for (i, &c) in b.iter().enumerate() {
655        match c {
656            '\'' => in_s = !in_s,
657            '(' if !in_s => depth += 1,
658            ')' if !in_s => depth -= 1,
659            c if c.is_whitespace() && depth == 0 && !in_s => cut = Some(i),
660            _ => {}
661        }
662    }
663    match cut {
664        Some(i) => {
665            let alias = p[i..].trim().trim_matches('"');
666            if alias.is_empty() { (p, None) } else { (p[..i].trim(), Some(alias)) }
667        }
668        None => (p, None),
669    }
670}
671
672fn sql_literals_to_nql(s: &str) -> String {
673    let mut out = String::with_capacity(s.len());
674    let mut it = s.chars().peekable();
675    while let Some(c) = it.next() {
676        match c {
677            '\'' => {
678                out.push('"');
679                while let Some(ch) = it.next() {
680                    if ch == '\'' {
681                        if it.peek() == Some(&'\'') {
682                            it.next();
683                            out.push('\''); // doubled '' is one literal quote
684                        } else {
685                            break;
686                        }
687                    } else if ch == '"' {
688                        // A double quote inside a SQL literal must be escaped
689                        // for NQL, whose lexer collapses \" to a literal quote.
690                        out.push('\\');
691                        out.push('"');
692                    } else {
693                        out.push(ch);
694                    }
695                }
696                out.push('"');
697            }
698            '<' if it.peek() == Some(&'>') => { it.next(); out.push_str("!="); }
699            _ => out.push(c),
700        }
701    }
702    out
703}
704
705fn strip_prefix_ci(s: &str, prefix: &str) -> Option<String> {
706    if s.len() >= prefix.len() && s[..prefix.len()].eq_ignore_ascii_case(prefix) {
707        Some(s[prefix.len()..].trim_start().to_string())
708    } else {
709        None
710    }
711}
712
713/// Find a top-level keyword (not inside quotes or parentheses), returning its
714/// byte offset. Case-insensitive, and only matches on word boundaries.
715fn find_kw(s: &str, kw: &str) -> Option<usize> {
716    let bytes = s.as_bytes();
717    let k = kw.as_bytes();
718    let mut depth = 0i32;
719    let mut in_s = false;
720    let mut in_d = false;
721    let mut i = 0usize;
722    while i < bytes.len() {
723        let c = bytes[i];
724        if in_s { if c == b'\'' { in_s = false; } i += 1; continue; }
725        if in_d { if c == b'"' { in_d = false; } i += 1; continue; }
726        match c {
727            b'\'' => { in_s = true; i += 1; continue; }
728            b'"' => { in_d = true; i += 1; continue; }
729            b'(' => { depth += 1; i += 1; continue; }
730            b')' => { depth -= 1; i += 1; continue; }
731            _ => {}
732        }
733        if depth == 0 && i + k.len() <= bytes.len()
734            && bytes[i..i + k.len()].eq_ignore_ascii_case(k)
735        {
736            let before_ok = i == 0 || !(bytes[i - 1] as char).is_alphanumeric() && bytes[i - 1] != b'_';
737            let after = i + k.len();
738            let after_ok = after >= bytes.len()
739                || !(bytes[after] as char).is_alphanumeric() && bytes[after] != b'_';
740            if before_ok && after_ok {
741                return Some(i);
742            }
743        }
744        i += 1;
745    }
746    None
747}
748
749/// Split a comma-separated list at the TOP level, ignoring commas inside
750/// quotes or parentheses — so `VALUES (1, 'a,b'), (2, 'c')` splits into two
751/// groups and not four.
752fn split_top(s: &str, sep: char) -> Vec<String> {
753    let mut out = vec![];
754    let mut cur = String::new();
755    let mut depth = 0i32;
756    let mut in_s = false;
757    let mut it = s.chars().peekable();
758    while let Some(c) = it.next() {
759        if in_s {
760            cur.push(c);
761            if c == '\'' {
762                // A doubled '' is an escaped quote, not the end of the literal.
763                if it.peek() == Some(&'\'') { cur.push(it.next().unwrap()); } else { in_s = false; }
764            }
765            continue;
766        }
767        match c {
768            '\'' => { in_s = true; cur.push(c); }
769            '(' => { depth += 1; cur.push(c); }
770            ')' => { depth -= 1; cur.push(c); }
771            x if x == sep && depth == 0 => { out.push(cur.trim().to_string()); cur.clear(); }
772            _ => cur.push(c),
773        }
774    }
775    if !cur.trim().is_empty() { out.push(cur.trim().to_string()); }
776    out
777}
778
779/// Parse one SQL scalar literal into JSON.
780///
781/// Deliberately narrow: a string, a number, a boolean, or NULL. Anything else
782/// — a function call, an expression, a cast — is refused by name rather than
783/// coerced into a string that would silently store the wrong value.
784fn sql_value(raw: &str) -> Result<Value, String> {
785    let t = raw.trim();
786    if t.is_empty() {
787        return Err("empty value".into());
788    }
789    let up = t.to_uppercase();
790    if up == "NULL" { return Ok(Value::Null); }
791    if up == "TRUE" { return Ok(Value::Bool(true)); }
792    if up == "FALSE" { return Ok(Value::Bool(false)); }
793    if t.starts_with('\'') && t.ends_with('\'') && t.len() >= 2 {
794        // Unwrap, collapsing the SQL '' escape to one quote.
795        let inner = &t[1..t.len() - 1];
796        return Ok(Value::String(inner.replace("''", "'")));
797    }
798    if let Ok(i) = t.parse::<i64>() { return Ok(Value::from(i)); }
799    if let Ok(f) = t.parse::<f64>() { return Ok(Value::from(f)); }
800    Err(format!(
801        "cannot use {:?} as a value — this endpoint accepts string literals, \
802         numbers, TRUE/FALSE and NULL. Expressions, casts and function calls \
803         are not evaluated, because storing an unevaluated expression as text \
804         would be worse than refusing it", t))
805}
806
807/// Pull a trailing `RETURNING …` off a statement, returning (head, columns).
808fn split_returning(tail: &str) -> (String, Vec<Col>) {
809    let tu = tail.to_uppercase();
810    match find_kw(&tu, "RETURNING") {
811        None => (tail.to_string(), vec![]),
812        Some(at) => {
813            let head = tail[..at].trim().to_string();
814            let list = tail[at + "RETURNING".len()..].trim();
815            if list == "*" {
816                return (head, vec![]);   // empty projection = every column
817            }
818            let cols = split_top(list, ',')
819                .into_iter()
820                .map(|p| {
821                    let raw = p.split_whitespace().next().unwrap_or(&p).to_string();
822                    let name = raw.rsplit('.').next().unwrap_or(&raw).trim_matches('"').to_string();
823                    Col::same(&name)
824                })
825                .collect();
826            (head, cols)
827        }
828    }
829}
830
831/// Columns whose names are reserved: they carry provenance rather than data.
832fn take_reserved(doc: &mut serde_json::Map<String, Value>) -> (Option<String>, Vec<String>, Option<String>, Option<String>) {
833    let id = doc.remove("_id").or_else(|| doc.remove("id"))
834        .and_then(|v| match v {
835            Value::String(s) => Some(s),
836            Value::Null => None,
837            other => Some(other.to_string()),   // a numeric key is a fine id
838        });
839    let caused_by = match doc.remove("_caused_by") {
840        Some(Value::String(s)) => vec![s],
841        Some(Value::Array(a)) => a.into_iter()
842            .filter_map(|v| v.as_str().map(str::to_string)).collect(),
843        _ => vec![],
844    };
845    let vf = doc.remove("_valid_from").and_then(|v| v.as_str().map(str::to_string));
846    let vt = doc.remove("_valid_to").and_then(|v| v.as_str().map(str::to_string));
847    (id, caused_by, vf, vt)
848}
849
850/// `INSERT INTO coll (c1, c2) VALUES (v1, v2), (…) [RETURNING …]`
851fn translate_insert(sql: &str) -> Result<Stmt, String> {
852    let rest = strip_prefix_ci(sql, "INSERT")
853        .and_then(|r| strip_prefix_ci(&r, "INTO"))
854        .ok_or("expected INSERT INTO")?;
855    // Locate VALUES first. Everything before it is `coll (col, …)`; searching
856    // for `(` without that bound finds the VALUES parenthesis instead and
857    // swallows the keyword into the collection name.
858    let ru = rest.to_uppercase();
859    let values_at = find_kw(&ru, "VALUES").ok_or(
860        "expected VALUES — `INSERT … SELECT` is not supported on this endpoint")?;
861    let head = rest[..values_at].trim().to_string();
862    let open = head.find('(').ok_or(
863        "INSERT needs an explicit column list — `INSERT INTO t (a, b) VALUES (…)`. \
864         NEDB is schemaless, so there is no declared column order to infer from")?;
865    let coll = head[..open].trim().trim_matches('"');
866    let coll = coll.rsplit('.').next().unwrap_or(coll).to_string();
867    if coll.is_empty() {
868        return Err("expected a collection name after INSERT INTO".into());
869    }
870    let close = head.rfind(')').ok_or("unterminated column list")?;
871    if close < open {
872        return Err("malformed column list".into());
873    }
874    let tail_from_values = rest[values_at..].to_string();
875    let cols: Vec<String> = split_top(&head[open + 1..close], ',')
876        .into_iter()
877        .map(|c| c.trim().trim_matches('"').to_string())
878        .collect();
879    if cols.is_empty() {
880        return Err("the column list is empty".into());
881    }
882
883    let after = strip_prefix_ci(&tail_from_values, "VALUES")
884        .ok_or("expected VALUES after the column list")?;
885    let (values_part, returning) = split_returning(&after);
886
887    let mut rows = vec![];
888    for group in split_top(&values_part, ',') {
889        let g = group.trim();
890        if !(g.starts_with('(') && g.ends_with(')')) {
891            return Err(format!("expected a parenthesised row of values, got {:?}", g));
892        }
893        let vals = split_top(&g[1..g.len() - 1], ',');
894        if vals.len() != cols.len() {
895            return Err(format!(
896                "{} values for {} columns — every row must match the column list",
897                vals.len(), cols.len()));
898        }
899        let mut doc = serde_json::Map::new();
900        for (c, v) in cols.iter().zip(vals.iter()) {
901            doc.insert(c.clone(), sql_value(v)?);
902        }
903        let (id, caused_by, valid_from, valid_to) = take_reserved(&mut doc);
904        rows.push(InsertRow { id, doc, caused_by, valid_from, valid_to });
905    }
906    if rows.is_empty() {
907        return Err("INSERT with no rows".into());
908    }
909    Ok(Stmt::Insert { coll, rows, returning })
910}
911
912/// `UPDATE coll SET a = 1, b = 'x' [WHERE …] [RETURNING …]`
913fn translate_update(sql: &str) -> Result<Stmt, String> {
914    let rest = strip_prefix_ci(sql, "UPDATE").ok_or("expected UPDATE")?;
915    let ru = rest.to_uppercase();
916    let set_at = find_kw(&ru, "SET").ok_or("expected SET in UPDATE")?;
917    // `UPDATE orders o SET …` — Postgres allows an alias here, and taking the
918    // whole span as the collection name made it part of the name ("orders o").
919    let target = rest[..set_at].trim();
920    let mut parts = target.split_whitespace();
921    let coll = parts.next().unwrap_or("").trim_matches('"');
922    let coll = coll.rsplit('.').next().unwrap_or(coll).to_string();
923    let upd_alias: Option<String> = match parts.next() {
924        Some(w) if w.eq_ignore_ascii_case("AS") => {
925            parts.next().map(|a| a.trim_matches('"').to_string())
926        }
927        Some(w) => Some(w.trim_matches('"').to_string()),
928        None => None,
929    };
930    if coll.is_empty() {
931        return Err("expected a collection name after UPDATE".into());
932    }
933    let after_set = rest[set_at + 3..].trim().to_string();
934    let (after_set, returning) = split_returning(&after_set);
935
936    // WHERE ends the assignment list; everything after it is a NQL predicate.
937    let au = after_set.to_uppercase();
938    let (assigns_raw, where_raw) = match find_kw(&au, "WHERE") {
939        Some(at) => (after_set[..at].to_string(), after_set[at..].to_string()),
940        None => (after_set.clone(), String::new()),
941    };
942
943    let mut set = vec![];
944    for a in split_top(&assigns_raw, ',') {
945        let eq = a.find('=').ok_or(format!("expected `col = value` in SET, got {:?}", a))?;
946        let col = a[..eq].trim().trim_matches('"').to_string();
947        if col.is_empty() {
948            return Err("empty column name in SET".into());
949        }
950        set.push((col, sql_value(&a[eq + 1..])?));
951    }
952    if set.is_empty() {
953        return Err("UPDATE with no assignments".into());
954    }
955    // The matching rows are found with an ordinary NQL read, so the whole
956    // predicate surface (IN, BETWEEN, LIKE, OR, …) works in an UPDATE too.
957    let where_raw = strip_column_qualifiers(where_raw.trim(), &coll, upd_alias.as_deref())?;
958    let nql = format!("FROM {} {}", coll, sql_literals_to_nql(&where_raw))
959        .trim().to_string();
960    Ok(Stmt::Update { coll, set, nql, returning })
961}
962
963/// `DELETE FROM coll [WHERE …] [RETURNING …]`
964fn translate_delete(sql: &str) -> Result<Stmt, String> {
965    let rest = strip_prefix_ci(sql, "DELETE")
966        .and_then(|r| strip_prefix_ci(&r, "FROM"))
967        .ok_or("expected DELETE FROM")?;
968    let (rest, returning) = split_returning(&rest);
969    let end = rest.find(' ').unwrap_or(rest.len());
970    let coll = rest[..end].trim().trim_matches('"');
971    let coll = coll.rsplit('.').next().unwrap_or(coll).to_string();
972    if coll.is_empty() {
973        return Err("expected a collection name after DELETE FROM".into());
974    }
975    let (del_alias, where_raw) = split_table_alias(rest[end..].trim());
976    let where_raw = strip_column_qualifiers(where_raw, &coll, del_alias.as_deref())?;
977    let nql = format!("FROM {} {}", coll, sql_literals_to_nql(&where_raw))
978        .trim().to_string();
979    Ok(Stmt::Delete { coll, nql, returning })
980}
981
982/// Translate one SQL statement into something executable, or explain why not.
983pub fn translate(sql_raw: &str) -> Result<Stmt, String> {
984    let sql = normalise(sql_raw);
985    let sql = sql.trim().trim_end_matches(';').trim();
986    if sql.is_empty() {
987        return Ok(Stmt::Ok(""));
988    }
989    let upper = sql.to_uppercase();
990
991    // ── the handshake. Clients issue these before anything useful; answering
992    // them with plausible values is the difference between "connects" and
993    // "hangs on startup". They are canned on purpose — NEDB has no pg_catalog
994    // and pretending otherwise would be worse than a clear boundary.
995    if upper.starts_with("SET ") || upper.starts_with("BEGIN") || upper.starts_with("COMMIT")
996        || upper.starts_with("ROLLBACK") || upper.starts_with("DISCARD")
997        || upper.starts_with("LISTEN ") || upper.starts_with("UNLISTEN ")
998    {
999        // Accepted and ignored: there is one implicit read-only transaction.
1000        return Ok(Stmt::Ok(if upper.starts_with("SET") { "SET" } else { "OK" }));
1001    }
1002    if upper.starts_with("SHOW ") {
1003        let name = sql[5..].trim().to_lowercase();
1004        let val = match name.as_str() {
1005            "transaction_isolation" | "default_transaction_isolation" => "read committed",
1006            "server_version" => SERVER_VERSION,
1007            "server_encoding" | "client_encoding" => "UTF8",
1008            "standard_conforming_strings" => "on",
1009            "is_superuser" => "off",
1010            _ => "",
1011        };
1012        return Ok(Stmt::Canned { cols: vec![name], row: vec![val.to_string()] });
1013    }
1014    if upper == "SELECT VERSION()" {
1015        return Ok(Stmt::Canned {
1016            cols: vec!["version".into()],
1017            row: vec![full_version_string()],
1018        });
1019    }
1020    if upper == "SELECT 1" || upper == "SELECT 1;" {
1021        return Ok(Stmt::Canned { cols: vec!["?column?".into()], row: vec!["1".into()] });
1022    }
1023    if upper.starts_with("SELECT CURRENT_SCHEMA") {
1024        return Ok(Stmt::Canned { cols: vec!["current_schema".into()], row: vec!["public".into()] });
1025    }
1026    if upper.starts_with("SELECT CURRENT_DATABASE") {
1027        return Ok(Stmt::Canned { cols: vec!["current_database".into()], row: vec!["nedb".into()] });
1028    }
1029    if upper.starts_with("SELECT CURRENT_USER") || upper.starts_with("SELECT USER") {
1030        return Ok(Stmt::Canned { cols: vec!["current_user".into()], row: vec!["nedb".into()] });
1031    }
1032
1033    // ── writes ───────────────────────────────────────────────────────────────
1034    // SQL's write semantics and NEDB's append-only model line up, so these are
1035    // first-class rather than refused. See the `Stmt` doc comment.
1036    if upper.starts_with("INSERT") { return translate_insert(sql); }
1037    if upper.starts_with("UPDATE") { return translate_update(sql); }
1038    if upper.starts_with("DELETE") { return translate_delete(sql); }
1039
1040    // ── the refusals that remain, each naming the boundary ──────────────────
1041    for (kw, why) in [
1042        ("CREATE", "DDL is not supported — collections are created implicitly by the first write to them, because NEDB is schemaless"),
1043        ("ALTER", "DDL is not supported — there is no schema to alter"),
1044        ("DROP", "DDL is not supported; drop a database with DELETE /v1/databases/<db>"),
1045        ("TRUNCATE", "not supported, and not an oversight: NEDB is append-only so that history cannot be discarded. That is the product"),
1046        ("COPY", "not supported; use GET /v1/databases/<db>/since for bulk export"),
1047        ("GRANT", "there is no SQL-level privilege system; auth is the bearer token"),
1048        ("REVOKE", "there is no SQL-level privilege system; auth is the bearer token"),
1049    ] {
1050        if upper.starts_with(kw) {
1051            return Err(format!("{} is not supported — {}", kw, why));
1052        }
1053    }
1054    if !upper.starts_with("SELECT") {
1055        return Err(format!(
1056            "only SELECT, INSERT, UPDATE and DELETE are supported on the Postgres \
1057             endpoint (got {:?})",
1058            sql.split_whitespace().next().unwrap_or("")
1059        ));
1060    }
1061    for (kw, why) in [
1062        (" JOIN ", "JOIN is not supported — NQL is single-collection; join in your client or model the relation with LINK/TRAVERSE"),
1063        (" UNION ", "UNION is not supported"),
1064        (" INTERSECT ", "INTERSECT is not supported"),
1065        (" EXCEPT ", "EXCEPT is not supported"),
1066        (" OVER (", "window functions are not supported"),
1067        ("DISTINCT ", "DISTINCT is not supported — GROUP BY <col> gives the distinct values with counts"),
1068    ] {
1069        if upper.contains(kw) {
1070            return Err(why.to_string());
1071        }
1072    }
1073    if find_kw(&upper, "FROM").is_none() {
1074        return Err("SELECT without FROM is not supported on this endpoint".into());
1075    }
1076
1077    // ── SELECT <projection> FROM <rest> ──────────────────────────────────────
1078    let after_select = strip_prefix_ci(sql, "SELECT").ok_or("expected SELECT")?;
1079    let from_at = find_kw(&after_select.to_uppercase(), "FROM")
1080        .ok_or("expected FROM after the select list")?;
1081    let projection = after_select[..from_at].trim().to_string();
1082    let rest = after_select[from_at + 4..].trim().to_string();
1083    if rest.is_empty() {
1084        return Err("expected a collection name after FROM".into());
1085    }
1086    // ── the one derived table with a provable flat equivalent ───────────────
1087    //
1088    // `SELECT count(*) FROM (SELECT … FROM coll WHERE …) AS anon` is what
1089    // EVERY ORM emits for `.count()` — SQLAlchemy's `Query.count()` wraps the
1090    // whole query in a subquery unconditionally. Refusing it means "SQLAlchemy
1091    // works, except counting", which is not a boundary anyone would accept.
1092    //
1093    // Counting a derived table whose rows are exactly the inner query's rows
1094    // is counting the inner query, so the rewrite is an IDENTITY rather than
1095    // an approximation. Each guard below names a construct that would break
1096    // that identity, and anything carrying one is still refused:
1097    //
1098    //   * `LIMIT` / `OFFSET`   — caps the row count before it is counted
1099    //   * `DISTINCT`           — collapses duplicates, so the counts differ
1100    //   * `GROUP BY`           — the inner rows ARE the groups
1101    //   * an inner aggregate   — already one row, counting it answers 1
1102    //   * anything but `count(*)` outside — the outer list would need the
1103    //     inner columns, which a flat count cannot supply
1104    if rest.starts_with('(') {
1105        if let Some(flat) = flatten_count_of_subquery(&projection, &rest) {
1106            // Recurses ONCE at most: the rewrite is only produced when the
1107            // inner FROM names a real collection, so the flat statement can
1108            // never re-enter this branch.
1109            return translate(&flat);
1110        }
1111        return Err("subqueries in FROM are not supported — except \
1112                    `SELECT count(*) FROM (…)`, which is rewritten to a flat \
1113                    count when the inner query has no LIMIT, OFFSET, DISTINCT, \
1114                    GROUP BY or aggregate of its own (any of those would make the \
1115                    two counts different numbers)".into());
1116    }
1117    let coll_end = rest.find(' ').unwrap_or(rest.len());
1118    let coll = &rest[..coll_end];
1119    if coll.contains(',') {
1120        return Err("selecting from more than one collection is not supported (no JOIN)".into());
1121    }
1122    // Postgres clients often qualify as schema.table; NEDB has one namespace,
1123    // so the schema is dropped — EXCEPT for `information_schema`, whose table
1124    // names (`tables`, `columns`) are words a user could plausibly name a
1125    // collection. Keeping the qualifier there is what stops
1126    // `SELECT * FROM information_schema.tables` and a real collection called
1127    // `tables` from resolving to the same thing.
1128    let bare = coll.rsplit('.').next().unwrap_or(coll).trim_matches('"');
1129    let qualified = coll
1130        .split('.')
1131        .map(|p| p.trim_matches('"'))
1132        .collect::<Vec<_>>()
1133        .join(".");
1134    let coll = if qualified.starts_with("information_schema.") {
1135        qualified.as_str()
1136    } else {
1137        bare
1138    };
1139    let tail = rest[coll_end..].trim();
1140
1141    // ── the select list ──────────────────────────────────────────────────────
1142    //
1143    // Parsed ITEM BY ITEM, which is what lets a list MIX plain columns with an
1144    // aggregate — and that mixture is exactly what a `GROUP BY` query is.
1145    // SQLAlchemy writes `SELECT orders.status, count(*) AS count_1 FROM orders
1146    // GROUP BY orders.status` for the most ordinary grouped query there is,
1147    // and the previous check refused any list containing a parenthesis at all,
1148    // so the whole shape was unreachable even though NQL expresses it
1149    // natively.
1150    //
1151    // NQL's grouped row carries the group key, `count`, and at most one NAMED
1152    // aggregate — so `count(*)` is always available and one of SUM/AVG/MIN/MAX
1153    // may join it. A second named aggregate is refused by name rather than
1154    // silently dropped.
1155    let mut agg_clause = String::new();
1156    let mut agg_srcs: Vec<String> = vec![];
1157    let mut project: Vec<Col> = vec![];
1158
1159    if projection == "*" {
1160        // everything
1161    } else {
1162        for part in split_top_level(&projection, ',') {
1163            let p = part.trim();
1164            if p.is_empty() {
1165                return Err("empty column in the select list".into());
1166            }
1167            let (expr, alias) = split_output_alias(p);
1168            let eu = expr.to_uppercase();
1169
1170            // COUNT(*) and COUNT(col) both become NQL's bare COUNT: NQL counts
1171            // the group, and a per-column non-null count is not expressible.
1172            if eu.starts_with("COUNT(") {
1173                if agg_clause.is_empty() {
1174                    agg_clause = " COUNT".to_string();
1175                }
1176                agg_srcs.push("count".to_string());
1177                project.push(Col::renamed("count", alias.unwrap_or("count")));
1178                continue;
1179            }
1180            if let Some(agg) = ["SUM", "AVG", "MIN", "MAX"]
1181                .iter()
1182                .find(|a| eu.starts_with(&format!("{}(", a)))
1183            {
1184                let inner = expr[agg.len() + 1..].trim_end_matches(')').trim();
1185                if inner.is_empty() || inner == "*" {
1186                    return Err(format!("{}() needs a column", agg));
1187                }
1188                let inner = inner.rsplit('.').next().unwrap_or(inner).trim_matches('"');
1189                let named = format!("{} {}", agg, inner);
1190                if !agg_clause.is_empty() && agg_clause.trim() != "COUNT" && agg_clause.trim() != named {
1191                    return Err(format!(
1192                        "only one of SUM/AVG/MIN/MAX is supported per statement \
1193                         (already have {:?}, then {:?}) — NQL's grouped row carries \
1194                         the group key, `count`, and ONE named aggregate",
1195                        agg_clause.trim(), named));
1196                }
1197                agg_clause = format!(" {}", named);
1198                // NQL emits `<agg>_<field>`; SQL names the column after the
1199                // function unless the query aliased it.
1200                let src = format!("{}_{}", agg.to_lowercase(), inner);
1201                project.push(Col::renamed(&src, alias.unwrap_or(&agg.to_lowercase())));
1202                agg_srcs.push(src);
1203                continue;
1204            }
1205            if expr.contains('(') {
1206                return Err(format!(
1207                    "expressions in the select list are not supported ({:?}) — \
1208                     supported: *, a column list, COUNT(*), or SUM/AVG/MIN/MAX(col)", p));
1209            }
1210            let name = expr.rsplit('.').next().unwrap_or(expr).trim_matches('"');
1211            project.push(Col::renamed(name, alias.unwrap_or(name)));
1212        }
1213    }
1214
1215    // ── clause tail: AS OF SYSTEM TIME → AS OF, then pass the rest through ──
1216    //
1217    // The clause keywords NQL shares with SQL (WHERE, GROUP BY, HAVING,
1218    // ORDER BY, LIMIT, OFFSET) are deliberately handed to the NQL parser
1219    // unchanged rather than re-parsed here. NQL is the authority on what is
1220    // valid; re-implementing its grammar would give two parsers to disagree.
1221    // `FROM orders o WHERE …` — the alias is taken off the tail (NQL has no
1222    // alias syntax) and then ACCEPTED as a qualifier on the columns.
1223    let (alias, tail) = split_table_alias(tail);
1224    let mut tail = strip_column_qualifiers(tail, coll, alias.as_deref())?;
1225    let tu = tail.to_uppercase();
1226    if let Some(at) = find_kw(&tu, "AS OF SYSTEM TIME") {
1227        let before = tail[..at].to_string();
1228        let after = tail[at + "AS OF SYSTEM TIME".len()..].trim_start().to_string();
1229        // Take the sequence token; the rest of the tail follows it.
1230        let end = after.find(' ').unwrap_or(after.len());
1231        let seq = after[..end].trim().trim_matches('\'').trim_matches('"').to_string();
1232        if seq.parse::<u64>().is_err() {
1233            return Err(format!(
1234                "AS OF SYSTEM TIME takes a NEDB sequence number here, not a timestamp (got {:?}). \
1235                 NEDB's history is sequence-addressed and never garbage-collected, so a seq is \
1236                 exact where a wall-clock time would be approximate", seq));
1237        }
1238        tail = format!("{} AS OF {} {}", before.trim(), seq, after[end..].trim())
1239            .trim()
1240            .to_string();
1241    }
1242
1243    // ── ORDER BY <ordinal> → ORDER BY <that select-list column> ─────────────
1244    //
1245    // SQL lets a sort key be a POSITION in the select list, and clients write
1246    // it constantly — `ORDER BY 1, 2` is how psql's own catalogue queries sort,
1247    // and node-postgres sent `GROUP BY status ORDER BY 1` in the very first
1248    // run of the driver harness. NQL has no ordinals: it read the `1` as a
1249    // literal and refused with "expected field name, got Num(1.0)".
1250    //
1251    // The projection is already parsed here, so the position resolves to a
1252    // real field name. An ordinal past the end of the select list, or one used
1253    // with `SELECT *` where there is no list to index, is refused with the
1254    // reason — guessing a column would sort by something the query never named.
1255    let tu_ord = tail.to_uppercase();
1256    if let Some(ob_at) = find_kw(&tu_ord, "ORDER BY") {
1257        let start = ob_at + "ORDER BY".len();
1258        // The clause runs to the next one, or to the end of the tail.
1259        let end = ["LIMIT", "OFFSET", "GROUP BY", "TRACE", "TRAVERSE", "SEARCH"]
1260            .iter()
1261            .filter_map(|k| find_kw(&tu_ord[start..], k).map(|at| start + at))
1262            .min()
1263            .unwrap_or(tail.len());
1264        let mut keys = vec![];
1265        for item in split_top_level(&tail[start..end], ',') {
1266            let item = item.trim();
1267            if item.is_empty() {
1268                continue;
1269            }
1270            let mut parts = item.split_whitespace();
1271            let first = parts.next().unwrap_or("");
1272            let rest: Vec<&str> = parts.collect();
1273            match first.parse::<usize>() {
1274                Ok(n) if n >= 1 => {
1275                    let col = project.get(n - 1).ok_or_else(|| {
1276                        if project.is_empty() {
1277                            format!(
1278                                "ORDER BY {} is a select-list POSITION, and `SELECT *` \
1279                                 has no list to index — name the column instead", n)
1280                        } else {
1281                            format!(
1282                                "ORDER BY {} is out of range: the select list has {} \
1283                                 column(s)", n, project.len())
1284                        }
1285                    })?;
1286                    keys.push(
1287                        std::iter::once(col.src.as_str())
1288                            .chain(rest.iter().copied())
1289                            .collect::<Vec<_>>()
1290                            .join(" "),
1291                    );
1292                }
1293                // Not an ordinal — a named column, or `1 + 1`, which NQL will
1294                // judge for itself.
1295                _ => keys.push(item.to_string()),
1296            }
1297        }
1298        tail = format!("{} ORDER BY {} {}", &tail[..ob_at], keys.join(", "), &tail[end..])
1299            .split_whitespace()
1300            .collect::<Vec<_>>()
1301            .join(" ");
1302    }
1303
1304    // ── GROUP BY: refuse a bare column that SQL would refuse ─────────────────
1305    //
1306    // A grouped NQL row holds only the group key, `count` and the aggregate —
1307    // so projecting `total` from `GROUP BY region` found nothing and rendered
1308    // NULL. Silently answering NULL for a column the query cannot produce is
1309    // the exact failure shape this engine keeps getting bitten by, so it is an
1310    // error, using Postgres's own wording so the message is already familiar.
1311    let tu_all = tail.to_uppercase();
1312    if let Some(gb_at) = find_kw(&tu_all, "GROUP BY") {
1313        let head = tail[..gb_at].trim_end().to_string();
1314        let after = tail[gb_at + "GROUP BY".len()..].trim_start();
1315        let key_end = after.find(|c: char| c == ' ' || c == ',').unwrap_or(after.len());
1316        let group_key = after[..key_end].trim().trim_matches('"').to_string();
1317        let after_key = after[key_end..].trim_start();
1318
1319        // NQL groups by ONE field. Taking the first key and leaving the rest
1320        // in the tail would group by something narrower than the query asked
1321        // for — more rows than Postgres returns, each aggregating too much.
1322        if after_key.starts_with(',') {
1323            return Err(format!(
1324                "GROUP BY takes one key here (got {:?} and more) — NQL groups by a \
1325                 single field, and grouping by only the first would aggregate over \
1326                 rows the query meant to keep apart",
1327                group_key));
1328        }
1329
1330        for c in &project {
1331            let ok = c.src == group_key
1332                || c.src == "count"
1333                || agg_srcs.contains(&c.src);
1334            if !ok {
1335                return Err(format!(
1336                    "column {:?} must appear in the GROUP BY clause or be used in an \
1337                     aggregate function — a grouped row carries the group key, `count`, \
1338                     and the aggregate, nothing else",
1339                    c.src));
1340            }
1341        }
1342
1343        // NQL's aggregate belongs IMMEDIATELY AFTER the group key
1344        // (`GROUP BY status COUNT`), not after the collection name. Emitting
1345        // `FROM orders COUNT GROUP BY status` is refused by the NQL parser
1346        // with "only one aggregate per query" — which is how the most
1347        // ordinary grouped query an ORM writes still failed even once its
1348        // select list parsed.
1349        //
1350        // `count` rides along free with a named aggregate — an NQL grouped row
1351        // carries the key, `count` AND the aggregate — so only the named one
1352        // is emitted when both were asked for.
1353        tail = format!("{} GROUP BY {}{} {}", head, group_key, agg_clause, after_key)
1354            .split_whitespace()
1355            .collect::<Vec<_>>()
1356            .join(" ");
1357        agg_clause.clear();
1358    }
1359
1360    let tail = sql_literals_to_nql(&tail);
1361    let nql = format!("FROM {}{}{}", coll,
1362                      if agg_clause.is_empty() { String::new() } else { agg_clause },
1363                      if tail.is_empty() { String::new() } else { format!(" {}", tail) });
1364
1365    Ok(Stmt::Query { nql: nql.trim().to_string(), project })
1366}
1367
1368const SERVER_VERSION: &str = "15.0";
1369
1370/// The `version()` string, for the SQL engine's `version()` function.
1371pub fn version_string() -> String {
1372    full_version_string()
1373}
1374
1375fn full_version_string() -> String {
1376    format!(
1377        "PostgreSQL {} (NEDB {}) — tamper-evident, append-only, permanent \
1378         history. SELECT + INSERT/UPDATE/DELETE; an UPDATE is a new version, \
1379         so prior values stay readable with AS OF SYSTEM TIME.",
1380        SERVER_VERSION,
1381        env!("CARGO_PKG_VERSION")
1382    )
1383}
1384
1385// ── result shaping ──────────────────────────────────────────────────────────
1386
1387/// Pick the column order for a result set.
1388///
1389/// With an explicit projection, that order. Otherwise the union of keys across
1390/// the returned rows — `_`-prefixed provenance columns last, so `psql` shows
1391/// the user's own fields first and `_hash` does not push `status` off screen.
1392fn columns_for(rows: &[Value], project: &[Col]) -> Vec<Col> {
1393    if !project.is_empty() {
1394        return project.to_vec();
1395    }
1396    let mut plain: Vec<String> = vec![];
1397    let mut meta: Vec<String> = vec![];
1398    for r in rows {
1399        if let Value::Object(m) = r {
1400            for k in m.keys() {
1401                let target = if k.starts_with('_') { &mut meta } else { &mut plain };
1402                if !target.contains(k) {
1403                    target.push(k.clone());
1404                }
1405            }
1406        }
1407    }
1408    plain.sort();
1409    meta.sort();
1410    plain.extend(meta);
1411    plain.into_iter().map(|k| Col::same(&k)).collect()
1412}
1413
1414/// The Postgres type of one JSON value.
1415fn oid_of_value(v: &Value) -> Option<i32> {
1416    match v {
1417        Value::Null => None,
1418        Value::Bool(_) => Some(OID_BOOL),
1419        Value::Number(n) => Some(if n.is_i64() || n.is_u64() { OID_INT8 } else { OID_FLOAT8 }),
1420        Value::String(_) => Some(OID_TEXT),
1421        // Arrays and objects render as their JSON text.
1422        _ => Some(OID_TEXT),
1423    }
1424}
1425
1426/// Reconcile two observed types for the same column.
1427///
1428/// A relational column has one type by construction. A NEDB collection does
1429/// not: document 1 may hold `qty: 3` and document 2 `qty: "three"`. Widening
1430/// to `text` on a conflict is the only answer that can carry both, and mixed
1431/// integers and floats widen to float8 for the same reason.
1432fn unify_oid(a: i32, b: i32) -> i32 {
1433    if a == b {
1434        return a;
1435    }
1436    match (a, b) {
1437        (OID_INT8, OID_FLOAT8) | (OID_FLOAT8, OID_INT8) => OID_FLOAT8,
1438        _ => OID_TEXT,
1439    }
1440}
1441
1442/// The type of `col` across EVERY row in the result, not just the first.
1443///
1444/// Taking the first non-null value's type was a latent wrong answer: a column
1445/// holding `3` in row one and `"n/a"` in row two was advertised as `int8`, and
1446/// a client that believes the description then fails parsing `"n/a"` as an
1447/// integer — or, on the binary path, cannot be sent the value at all.
1448/// Public alias so `pgcatalog` types a column EXACTLY as the wire does.
1449///
1450/// The catalogue reporting `bigint` for a column the protocol then sends as
1451/// text would be a self-contradiction a client is entitled to trust, so both
1452/// go through this one function rather than two that agree today.
1453pub fn oid_for_column(rows: &[Value], col: &str) -> i32 {
1454    oid_for(rows, col)
1455}
1456
1457fn oid_for(rows: &[Value], col: &str) -> i32 {
1458    let mut acc: Option<i32> = None;
1459    for r in rows {
1460        if let Some(o) = r.get(col).and_then(oid_of_value) {
1461            acc = Some(match acc {
1462                None => o,
1463                Some(prev) => unify_oid(prev, o),
1464            });
1465            if acc == Some(OID_TEXT) {
1466                break; // text absorbs everything; no need to look further
1467            }
1468        }
1469    }
1470    acc.unwrap_or(OID_TEXT)
1471}
1472
1473/// Render one cell in the text format Postgres clients expect for format 0.
1474fn cell(v: Option<&Value>) -> Option<String> {
1475    match v {
1476        None | Some(Value::Null) => None, // NULL on the wire
1477        Some(Value::String(s)) => Some(s.clone()),
1478        Some(Value::Bool(b)) => Some(if *b { "t".into() } else { "f".into() }),
1479        Some(other) => Some(other.to_string()),
1480    }
1481}
1482
1483/// Render one cell in binary format for the type the column was advertised as.
1484///
1485/// Needed because asyncpg asks for binary results — it is not an optimisation
1486/// there, it is the only format it requests, so without this it cannot read a
1487/// single row. Text-format clients never reach this path.
1488///
1489/// A value that does not fit the advertised type is an error rather than a
1490/// coercion. The advertised type comes from sampling stored documents, so a
1491/// mismatch means the field is genuinely heterogeneous beyond the sample, and
1492/// quietly sending a zero (or the text bytes under a binary header) would
1493/// corrupt the value in a way the client cannot detect.
1494fn cell_binary(v: Option<&Value>, oid: i32) -> Result<Option<Vec<u8>>, String> {
1495    let v = match v {
1496        None | Some(Value::Null) => return Ok(None),
1497        Some(v) => v,
1498    };
1499    let as_f64 = |n: &serde_json::Number| n.as_f64()
1500        .ok_or_else(|| "a number too large to send as float8".to_string());
1501    Ok(Some(match (oid, v) {
1502        (OID_BOOL, Value::Bool(b)) => vec![u8::from(*b)],
1503        (OID_INT2, Value::Number(n)) => {
1504            let i = n.as_i64().ok_or("not an integer")?;
1505            i16::try_from(i).map_err(|_| format!("{} does not fit in int2", i))?
1506                .to_be_bytes().to_vec()
1507        }
1508        (OID_INT4, Value::Number(n)) => {
1509            let i = n.as_i64().ok_or("not an integer")?;
1510            i32::try_from(i).map_err(|_| format!("{} does not fit in int4", i))?
1511                .to_be_bytes().to_vec()
1512        }
1513        (OID_INT8, Value::Number(n)) => {
1514            n.as_i64().ok_or("not an integer")?.to_be_bytes().to_vec()
1515        }
1516        (OID_FLOAT4, Value::Number(n)) => (as_f64(n)? as f32).to_be_bytes().to_vec(),
1517        (OID_FLOAT8, Value::Number(n)) => as_f64(n)?.to_be_bytes().to_vec(),
1518        // For the text family, binary and text are the same bytes.
1519        (OID_TEXT | OID_VARCHAR | OID_NAME | OID_UNKNOWN | OID_JSON, _) => {
1520            cell(Some(v)).unwrap_or_default().into_bytes()
1521        }
1522        // jsonb is a one-byte version header then the JSON text.
1523        (OID_JSONB, _) => {
1524            let mut b = vec![1u8];
1525            b.extend_from_slice(cell(Some(v)).unwrap_or_default().as_bytes());
1526            b
1527        }
1528        (oid, val) => {
1529            let kind = match val {
1530                Value::Bool(_) => "a boolean",
1531                Value::Number(_) => "a number",
1532                Value::String(_) => "a string",
1533                Value::Array(_) => "an array",
1534                _ => "an object",
1535            };
1536            return Err(format!(
1537                "cannot send {} in binary format as type OID {} — the field holds \
1538                 more than one type across documents, so it cannot be described \
1539                 by a single Postgres type. Select it with a text cast, or use a \
1540                 text-format client",
1541                kind, oid
1542            ));
1543        }
1544    }))
1545}
1546
1547/// A `RowDescription`, with a per-column wire format code.
1548fn row_description_fmt(cols: &[Col], oids: &[i32], fmts: &[i16]) -> Vec<u8> {
1549    let mut m = Out::msg(b'T');
1550    m.i16(cols.len() as i16);
1551    for (i, c) in cols.iter().enumerate() {
1552        m.cstr(&c.out);
1553        m.i32(0); // table OID — unknown
1554        m.i16((i + 1) as i16); // column attribute number
1555        m.i32(oids.get(i).copied().unwrap_or(OID_TEXT));
1556        m.i16(-1); // variable length
1557        m.i32(-1); // no type modifier
1558        m.i16(fmts.get(i).copied().unwrap_or(0));
1559    }
1560    m.finish()
1561}
1562
1563fn row_description(cols: &[Col], oids: &[i32]) -> Vec<u8> {
1564    row_description_fmt(cols, oids, &[])
1565}
1566
1567fn data_row_bytes(vals: &[Option<Vec<u8>>]) -> Vec<u8> {
1568    let mut m = Out::msg(b'D');
1569    m.i16(vals.len() as i16);
1570    for v in vals {
1571        match v {
1572            None => m.i32(-1),
1573            Some(b) => {
1574                m.i32(b.len() as i32);
1575                m.bytes(b);
1576            }
1577        }
1578    }
1579    m.finish()
1580}
1581
1582fn data_row(vals: &[Option<String>]) -> Vec<u8> {
1583    let owned: Vec<Option<Vec<u8>>> =
1584        vals.iter().map(|v| v.as_ref().map(|s| s.as_bytes().to_vec())).collect();
1585    data_row_bytes(&owned)
1586}
1587
1588/// Encode just the rows: `T` followed by one `D` per row, and NO
1589/// `CommandComplete`.
1590///
1591/// Split out because a write with `RETURNING` must emit `T`/`D`* and then its
1592/// OWN tag (`INSERT 0 3`, `UPDATE 1`). The first cut called `encode_result`
1593/// there, which appends `CommandComplete("SELECT n")` — so one statement sent
1594/// TWO CommandComplete messages. That is a protocol violation, and the visible
1595/// symptom was `RETURNING` silently yielding no rows at all: the client took
1596/// the first tag as the end of the statement and discarded the description.
1597pub fn encode_rows(rows: &[Value], project: &[Col]) -> Vec<u8> {
1598    let cols = columns_for(rows, project);
1599    let oids: Vec<i32> = cols.iter().map(|c| oid_for(rows, &c.src)).collect();
1600    let mut out = row_description(&cols, &oids);
1601    for r in rows {
1602        let vals: Vec<Option<String>> = cols.iter().map(|c| cell(r.get(&c.src))).collect();
1603        out.extend_from_slice(&data_row(&vals));
1604    }
1605    out
1606}
1607
1608/// A complete SELECT response: rows plus `CommandComplete("SELECT n")`.
1609pub fn encode_result(rows: &[Value], project: &[Col]) -> Vec<u8> {
1610    let mut out = encode_rows(rows, project);
1611    out.extend_from_slice(&command_complete(&format!("SELECT {}", rows.len())));
1612    out
1613}
1614
1615// ── the extended query protocol: Parse / Bind / Describe / Execute ──────────
1616//
1617// Why this exists at all: psycopg3, asyncpg and the JDBC driver do not speak
1618// the simple query protocol for parameterised statements. Without these six
1619// messages they cannot run a single query — psycopg3 hangs waiting for a
1620// `ParseComplete`, and asyncpg refuses before it ever sends a `Bind`. "psql
1621// works" is not the same as "the drivers your evaluators use work".
1622//
1623// Two facts about real drivers shaped everything below, and both were read off
1624// a wire transcript rather than assumed:
1625//
1626//   1. psycopg3 sends parameters in a MIXED format — a `str` as OID 0 in text
1627//      format, but an `int` as int2/int4/int8 in BINARY, a float as float8
1628//      binary, a bool as a single binary byte. A text-only decoder gets `\x00*`
1629//      where it expected `42`.
1630//
1631//   2. asyncpg declares NO parameter types in `Parse` and then asks
1632//      `Describe(statement)`, encoding its arguments from whatever OIDs come
1633//      back. Answering "text" for all of them does not degrade gracefully — it
1634//      makes asyncpg REFUSE the call client-side ("expected str, got int").
1635//
1636// (2) is the reason `infer_param_oids` exists. NEDB is schemaless, so there is
1637// no catalogue to read a column's type out of — the only honest source of truth
1638// is the data already stored, so the type is sampled from it.
1639
1640/// Parameter/result type OIDs handled on the binary path.
1641const OID_INT2: i32 = 21;
1642const OID_INT4: i32 = 23;
1643const OID_OID: i32 = 26;
1644const OID_FLOAT4: i32 = 700;
1645const OID_VARCHAR: i32 = 1043;
1646const OID_NAME: i32 = 19;
1647const OID_UNKNOWN: i32 = 705;
1648const OID_JSON: i32 = 114;
1649const OID_JSONB: i32 = 3802;
1650
1651/// How many `$n` placeholders a statement carries, and the highest index used.
1652///
1653/// Scans outside string literals so a `'$1'` inside a value is not mistaken for
1654/// a placeholder. Dollar-quoted bodies (`$tag$…$tag$`) are not recognised —
1655/// they need a procedural language NEDB does not have.
1656fn param_count(sql: &str) -> usize {
1657    let b = sql.as_bytes();
1658    let mut i = 0usize;
1659    let mut in_s = false;
1660    let mut max = 0usize;
1661    while i < b.len() {
1662        let c = b[i];
1663        if in_s {
1664            if c == b'\'' {
1665                in_s = false;
1666            }
1667            i += 1;
1668            continue;
1669        }
1670        if c == b'\'' {
1671            in_s = true;
1672            i += 1;
1673            continue;
1674        }
1675        if c == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
1676            let mut j = i + 1;
1677            let mut n = 0usize;
1678            while j < b.len() && b[j].is_ascii_digit() {
1679                n = n * 10 + (b[j] - b'0') as usize;
1680                j += 1;
1681            }
1682            max = max.max(n);
1683            i = j;
1684            continue;
1685        }
1686        i += 1;
1687    }
1688    max
1689}
1690
1691/// The JSON-shaped type of `field` as it is actually stored, sampled from the
1692/// collection, mapped onto the nearest Postgres OID.
1693///
1694/// This is the schemaless answer to "what type is this column?". A relational
1695/// server reads its catalogue; NEDB has none, so it reads the data. Sampling a
1696/// bounded number of rows keeps a `Describe` cheap, and the first row that
1697/// actually carries the field decides — a field missing from row one but
1698/// present in row nine still types correctly.
1699fn infer_field_oid(db: Option<&Arc<Db>>, coll: &str, field: &str) -> i32 {
1700    // `_`-prefixed names are engine metadata, not stored document fields, so
1701    // they type from the engine's own contract — no sampling, and no database
1702    // handle needed.
1703    match field {
1704        "_seq" => return OID_INT8,
1705        "_id" | "_hash" | "_prev" | "_collection" | "_valid_from" | "_valid_to" => return OID_TEXT,
1706        _ => {}
1707    }
1708    // A catalogue relation types its own columns. Sampling a USER collection
1709    // named `pg_type` finds nothing and falls back to text — and asyncpg,
1710    // which declares parameter types client-side and refuses the call when
1711    // the server's answer is wrong, then rejected `WHERE oid = $1` with
1712    // "expected str, got int" before a single byte was sent.
1713    if !field.is_empty() && crate::pgcatalog::is_catalog(coll) {
1714        if let Some(rows) = crate::pgcatalog::rows(coll, db) {
1715            return oid_for(&rows, field);
1716        }
1717    }
1718    let db = match db {
1719        Some(db) => db,
1720        None => return OID_TEXT,
1721    };
1722    if coll.is_empty() || field.is_empty() {
1723        return OID_TEXT;
1724    }
1725    let rows = match crate::nql::query(db, &format!("FROM {} LIMIT {}", coll, TYPE_SAMPLE)) {
1726        Ok((rows, _)) => rows,
1727        Err(_) => return OID_TEXT,
1728    };
1729    // Unified over the sample, not taken from the first hit: a field that is a
1730    // number in one document and a string in another has to be advertised as
1731    // text or a client cannot decode every row of it.
1732    oid_for(&rows, field)
1733}
1734
1735/// The type of an aggregate output column, which no document holds.
1736///
1737/// Sampling stored documents cannot type these: `COUNT(*)` produces a column
1738/// called `count` that exists in no document, so the sampler finds nothing and
1739/// falls back to text. A text-format client papers over that, but a binary
1740/// client is then handed the digits of a number under a text header and
1741/// `COUNT(*)` comes back as the string `"2"` instead of the integer `2`.
1742///
1743/// So aggregates are typed from what the aggregate MEANS: a count is always an
1744/// integer, an average is always fractional, and min/max/sum inherit the type
1745/// of the field they were computed over.
1746fn aggregate_oid(src: &str, db: Option<&Arc<Db>>, coll: &str) -> Option<i32> {
1747    if src == "count" {
1748        return Some(OID_INT8);
1749    }
1750    for (prefix, fixed) in [
1751        ("count_", Some(OID_INT8)),
1752        ("avg_", Some(OID_FLOAT8)),
1753        ("sum_", None),
1754        ("min_", None),
1755        ("max_", None),
1756    ] {
1757        if let Some(field) = src.strip_prefix(prefix) {
1758            return Some(match fixed {
1759                Some(oid) => oid,
1760                // SUM/MIN/MAX of an integer field is an integer; of a
1761                // fractional field, fractional.
1762                None => match infer_field_oid(db, coll, field) {
1763                    OID_INT8 => OID_INT8,
1764                    OID_FLOAT8 => OID_FLOAT8,
1765                    // Summing or ordering a non-numeric field is not
1766                    // meaningful; let the row-derived type answer.
1767                    other => other,
1768                },
1769            });
1770        }
1771    }
1772    None
1773}
1774
1775/// How many documents to sample when typing a column.
1776///
1777/// Bounded so a `Describe` stays cheap. It is a sample, so a field that only
1778/// turns heterogeneous outside it can still surprise us — which is exactly why
1779/// `cell_binary` refuses a mismatch loudly instead of coercing.
1780const TYPE_SAMPLE: usize = 200;
1781
1782/// The collection a statement reads from or writes to, for type sampling.
1783fn stmt_collection(sql: &str) -> String {
1784    let s = normalise(sql);
1785    let up = s.to_uppercase();
1786    let after = if let Some(at) = find_kw(&up, "FROM") {
1787        &s[at + 4..]
1788    } else if let Some(rest) = strip_prefix_ci(&s, "UPDATE") {
1789        return rest
1790            .split_whitespace()
1791            .next()
1792            .unwrap_or("")
1793            .rsplit('.')
1794            .next()
1795            .unwrap_or("")
1796            .trim_matches('"')
1797            .to_string();
1798    } else if let Some(rest) = strip_prefix_ci(&s, "INSERT INTO") {
1799        return rest
1800            .split(|c: char| c.is_whitespace() || c == '(')
1801            .find(|t| !t.is_empty())
1802            .unwrap_or("")
1803            .rsplit('.')
1804            .next()
1805            .unwrap_or("")
1806            .trim_matches('"')
1807            .to_string();
1808    } else {
1809        return String::new();
1810    };
1811    after
1812        .trim()
1813        .split(|c: char| c.is_whitespace())
1814        .find(|t| !t.is_empty())
1815        .unwrap_or("")
1816        .rsplit('.')
1817        .next()
1818        .unwrap_or("")
1819        .trim_matches('"')
1820        .to_string()
1821}
1822
1823/// Which document field each `$n` is being compared against.
1824///
1825/// Three shapes cover essentially all driver-generated SQL:
1826///   `WHERE qty > $1`        → the identifier immediately left of the operator
1827///   `SET status = $1`       → same shape, inside the SET list
1828///   `INSERT INTO t (a,b) VALUES ($1,$2)` → positional against the column list
1829///
1830/// Anything it cannot read returns `None`, which types as `text`. Guessing
1831/// wrong here would make a driver encode a value the engine then fails to
1832/// match, so an unknown is left unknown on purpose.
1833fn param_fields(sql: &str, n_params: usize) -> Vec<Option<String>> {
1834    let s = normalise(sql);
1835    let mut out = vec![None; n_params];
1836
1837    // The INSERT column list maps positionally, which is more reliable than
1838    // scanning leftwards through a VALUES tuple.
1839    let up = s.to_uppercase();
1840    if up.starts_with("INSERT") {
1841        if let (Some(open), Some(vals_at)) = (s.find('('), find_kw(&up, "VALUES")) {
1842            if open < vals_at {
1843                if let Some(close) = s[open..vals_at].rfind(')') {
1844                    let cols: Vec<String> = split_top(&s[open + 1..open + close], ',')
1845                        .into_iter()
1846                        .map(|c| c.trim().trim_matches('"').to_string())
1847                        .collect();
1848                    // `$1` is the first placeholder in the first tuple, and so on.
1849                    let tail = &s[vals_at..];
1850                    let mut seen = 0usize;
1851                    let b = tail.as_bytes();
1852                    let mut i = 0usize;
1853                    let mut in_s = false;
1854                    while i < b.len() {
1855                        if in_s {
1856                            if b[i] == b'\'' { in_s = false; }
1857                            i += 1;
1858                            continue;
1859                        }
1860                        if b[i] == b'\'' { in_s = true; i += 1; continue; }
1861                        if b[i] == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
1862                            let mut j = i + 1;
1863                            let mut num = 0usize;
1864                            while j < b.len() && b[j].is_ascii_digit() {
1865                                num = num * 10 + (b[j] - b'0') as usize;
1866                                j += 1;
1867                            }
1868                            if num >= 1 && num <= n_params {
1869                                if let Some(c) = cols.get(seen % cols.len().max(1)) {
1870                                    out[num - 1] = Some(c.clone());
1871                                }
1872                            }
1873                            seen += 1;
1874                            i = j;
1875                            continue;
1876                        }
1877                        i += 1;
1878                    }
1879                    return out;
1880                }
1881            }
1882        }
1883    }
1884
1885    // Otherwise: for each `$n`, walk left past the operator to the identifier.
1886    let b = s.as_bytes();
1887    let mut i = 0usize;
1888    let mut in_s = false;
1889    while i < b.len() {
1890        if in_s {
1891            if b[i] == b'\'' { in_s = false; }
1892            i += 1;
1893            continue;
1894        }
1895        if b[i] == b'\'' { in_s = true; i += 1; continue; }
1896        if b[i] == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
1897            let mut j = i + 1;
1898            let mut num = 0usize;
1899            while j < b.len() && b[j].is_ascii_digit() {
1900                num = num * 10 + (b[j] - b'0') as usize;
1901                j += 1;
1902            }
1903            if num >= 1 && num <= n_params {
1904                let left = &s[..i];
1905                // Skip the operator characters and whitespace sitting between
1906                // the identifier and the placeholder.
1907                let trimmed = left.trim_end_matches(|c: char| {
1908                    c.is_whitespace() || "=<>!+-*/%(,".contains(c)
1909                });
1910                // A word operator (`LIKE`, `IN`, `BETWEEN`, `AND`) also sits
1911                // between them; step over it to reach the real identifier.
1912                let mut tok = trimmed
1913                    .rsplit(|c: char| c.is_whitespace() || c == '(' || c == ',')
1914                    .find(|t| !t.is_empty())
1915                    .unwrap_or("")
1916                    .trim_matches('"');
1917                let mut before = trimmed;
1918                for _ in 0..4 {
1919                    let upper_tok = tok.to_uppercase();
1920                    // `BETWEEN $1 AND $2` puts BOTH a word operator and an
1921                    // earlier placeholder between `$2` and the column it
1922                    // constrains, so a placeholder has to be stepped over too —
1923                    // otherwise the upper bound of every range query types as
1924                    // text while the lower bound types correctly.
1925                    if upper_tok.starts_with('$')
1926                        || matches!(upper_tok.as_str(),
1927                        "LIKE" | "ILIKE" | "IN" | "BETWEEN" | "AND" | "OR" | "NOT" | "IS") {
1928                        before = before[..before.len() - tok.len()].trim_end_matches(|c: char| {
1929                            c.is_whitespace() || "=<>!(,".contains(c)
1930                        });
1931                        tok = before
1932                            .rsplit(|c: char| c.is_whitespace() || c == '(' || c == ',')
1933                            .find(|t| !t.is_empty())
1934                            .unwrap_or("")
1935                            .trim_matches('"');
1936                    } else {
1937                        break;
1938                    }
1939                }
1940                if !tok.is_empty()
1941                    && tok.chars().all(|c| c.is_alphanumeric() || c == '_' || c == '.')
1942                    && !tok.chars().next().map(|c| c.is_ascii_digit()).unwrap_or(true)
1943                {
1944                    out[num - 1] = Some(tok.rsplit('.').next().unwrap_or(tok).to_string());
1945                }
1946            }
1947            i = j;
1948            continue;
1949        }
1950        i += 1;
1951    }
1952    out
1953}
1954
1955/// The type of a placeholder sitting in a CLAUSE position rather than beside a
1956/// column.
1957///
1958/// `AS OF SYSTEM TIME $1` has no column to sample — the token to its left is
1959/// the word `TIME`. Its type comes from the grammar instead, which is both
1960/// cheaper and more certain than any inference: a system-time bound is a
1961/// sequence number, a valid-time bound is a date string, and a page bound is an
1962/// integer. Without this, a parameterised time-travel query typed as text and
1963/// asyncpg refused to send the integer at all.
1964fn clause_param_oids(sql: &str, n_params: usize) -> Vec<Option<i32>> {
1965    let s = normalise(sql);
1966    let mut out = vec![None; n_params];
1967    let b = s.as_bytes();
1968    let mut i = 0usize;
1969    let mut in_s = false;
1970    while i < b.len() {
1971        if in_s {
1972            if b[i] == b'\'' { in_s = false; }
1973            i += 1;
1974            continue;
1975        }
1976        if b[i] == b'\'' { in_s = true; i += 1; continue; }
1977        if b[i] == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
1978            let mut j = i + 1;
1979            let mut num = 0usize;
1980            while j < b.len() && b[j].is_ascii_digit() {
1981                num = num * 10 + (b[j] - b'0') as usize;
1982                j += 1;
1983            }
1984            if num >= 1 && num <= n_params {
1985                let left = s[..i].trim_end().to_uppercase();
1986                // VALID AS OF is checked FIRST: it ends with "AS OF" too, and
1987                // its argument is a DATE STRING, not a sequence number.
1988                out[num - 1] = if left.ends_with("VALID AS OF") {
1989                    Some(OID_TEXT)
1990                } else if left.ends_with("AS OF SYSTEM TIME")
1991                    || left.ends_with("FOR SYSTEM_TIME AS OF")
1992                    || left.ends_with("AS OF")
1993                    || left.ends_with("LIMIT")
1994                    || left.ends_with("OFFSET")
1995                {
1996                    Some(OID_INT8)
1997                } else {
1998                    None
1999                };
2000            }
2001            i = j;
2002            continue;
2003        }
2004        i += 1;
2005    }
2006    out
2007}
2008
2009/// The OIDs to advertise for `$1..$n`, sampled from stored data.
2010///
2011/// `declared` is what the client itself put in `Parse`. A client that states a
2012/// type is believed — it is about to encode its arguments that way, and second
2013///-guessing it would break the decode. Only the unspecified slots are inferred.
2014fn infer_param_oids(sql: &str, declared: &[i32], db: Option<&Arc<Db>>) -> Vec<i32> {
2015    let n = param_count(sql).max(declared.len());
2016    if n == 0 {
2017        return vec![];
2018    }
2019    let coll = stmt_collection(sql);
2020    let fields = param_fields(sql, n);
2021    let clauses = clause_param_oids(sql, n);
2022    (0..n)
2023        .map(|i| match declared.get(i) {
2024            Some(&oid) if oid != 0 => oid,
2025            // A clause position knows its own type from the grammar, so it
2026            // outranks sampling a column that is not even there.
2027            _ => match clauses[i] {
2028                Some(oid) => oid,
2029                None => match &fields[i] {
2030                    Some(f) => infer_field_oid(db, &coll, f),
2031                    None => OID_TEXT,
2032                },
2033            },
2034        })
2035        .collect()
2036}
2037
2038/// Decode one bound parameter into the SQL literal text to splice into the
2039/// statement.
2040///
2041/// `None` means SQL NULL. Format 1 is binary — see the module note on psycopg3
2042/// sending small integers as int2.
2043fn decode_param(raw: Option<&[u8]>, oid: i32, format: i16) -> Result<Option<String>, String> {
2044    let bytes = match raw {
2045        None => return Ok(None),
2046        Some(b) => b,
2047    };
2048    let quote = |s: &str| format!("'{}'", s.replace('\'', "''"));
2049
2050    if format == 0 {
2051        let s = String::from_utf8_lossy(bytes).to_string();
2052        return Ok(Some(match oid {
2053            OID_BOOL => {
2054                let t = matches!(s.as_str(), "t" | "true" | "TRUE" | "1" | "yes" | "on");
2055                if t { "TRUE".into() } else { "FALSE".into() }
2056            }
2057            OID_INT2 | OID_INT4 | OID_INT8 | OID_OID | OID_FLOAT4 | OID_FLOAT8 => {
2058                // Validate rather than trust: an unparseable "number" spliced
2059                // in bare would become a bare identifier in the NQL text and
2060                // produce a baffling error far from its cause.
2061                if s.parse::<f64>().is_ok() { s } else { quote(&s) }
2062            }
2063            // OID 0 with text format is psycopg3's `str`. Confirmed on the
2064            // wire: it declares a real numeric OID whenever the value is a
2065            // number, so an unspecified text parameter is genuinely a string
2066            // and quoting it is right rather than a guess.
2067            _ => quote(&s),
2068        }));
2069    }
2070    if format != 1 {
2071        return Err(format!("unsupported parameter format code {}", format));
2072    }
2073
2074    // ── binary ──────────────────────────────────────────────────────────────
2075    let need = |n: usize| -> Result<(), String> {
2076        if bytes.len() == n {
2077            Ok(())
2078        } else {
2079            Err(format!(
2080                "binary parameter of type OID {} should be {} bytes, got {}",
2081                oid, n, bytes.len()
2082            ))
2083        }
2084    };
2085    Ok(Some(match oid {
2086        OID_BOOL => {
2087            need(1)?;
2088            if bytes[0] != 0 { "TRUE".into() } else { "FALSE".into() }
2089        }
2090        OID_INT2 => {
2091            need(2)?;
2092            i16::from_be_bytes([bytes[0], bytes[1]]).to_string()
2093        }
2094        OID_INT4 => {
2095            need(4)?;
2096            i32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]).to_string()
2097        }
2098        OID_OID => {
2099            need(4)?;
2100            u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]).to_string()
2101        }
2102        OID_INT8 => {
2103            need(8)?;
2104            i64::from_be_bytes(bytes[..8].try_into().unwrap()).to_string()
2105        }
2106        OID_FLOAT4 => {
2107            need(4)?;
2108            let f = f32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
2109            fmt_float(f as f64)
2110        }
2111        OID_FLOAT8 => {
2112            need(8)?;
2113            fmt_float(f64::from_be_bytes(bytes[..8].try_into().unwrap()))
2114        }
2115        OID_TEXT | OID_VARCHAR | OID_NAME | OID_UNKNOWN | OID_JSON | 0 => {
2116            quote(&String::from_utf8_lossy(bytes))
2117        }
2118        OID_JSONB => {
2119            // jsonb binary is a 1-byte version header followed by the JSON text.
2120            let body = if bytes.first() == Some(&1) { &bytes[1..] } else { bytes };
2121            quote(&String::from_utf8_lossy(body))
2122        }
2123        other => {
2124            return Err(format!(
2125                "parameter type OID {} is not supported in binary format — \
2126                 the supported set is bool, int2/int4/int8, float4/float8, \
2127                 text/varchar/json/jsonb. Send it as text, or cast it in the \
2128                 statement",
2129                other
2130            ))
2131        }
2132    }))
2133}
2134
2135/// Render a float without Rust's `inf`/`NaN` spellings leaking into SQL text.
2136fn fmt_float(f: f64) -> String {
2137    if f.is_nan() {
2138        "'NaN'".into()
2139    } else if f.is_infinite() {
2140        if f > 0.0 { "'Infinity'".into() } else { "'-Infinity'".into() }
2141    } else if f.fract() == 0.0 && f.abs() < 1e15 {
2142        format!("{:.0}", f)
2143    } else {
2144        f.to_string()
2145    }
2146}
2147
2148/// Splice decoded parameters into the statement text.
2149///
2150/// Textual substitution, deliberately: the whole SQL surface is already a text
2151/// translation into NQL, so one representation is simpler and cannot disagree
2152/// with itself. Every value arrives already rendered as a SQL literal by
2153/// `decode_param`, with embedded quotes doubled, so a parameter cannot break
2154/// out of its literal and alter the statement's shape.
2155fn substitute_params(sql: &str, params: &[Option<String>]) -> Result<String, String> {
2156    let b = sql.as_bytes();
2157    let mut out = String::with_capacity(sql.len() + 16);
2158    let mut i = 0usize;
2159    let mut in_s = false;
2160    while i < b.len() {
2161        let c = b[i];
2162        if in_s {
2163            out.push(c as char);
2164            if c == b'\'' { in_s = false; }
2165            i += 1;
2166            continue;
2167        }
2168        if c == b'\'' {
2169            in_s = true;
2170            out.push('\'');
2171            i += 1;
2172            continue;
2173        }
2174        if c == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
2175            let mut j = i + 1;
2176            let mut n = 0usize;
2177            while j < b.len() && b[j].is_ascii_digit() {
2178                n = n * 10 + (b[j] - b'0') as usize;
2179                j += 1;
2180            }
2181            match params.get(n.wrapping_sub(1)) {
2182                Some(Some(lit)) => out.push_str(lit),
2183                Some(None) => out.push_str("NULL"),
2184                None => {
2185                    return Err(format!(
2186                        "bind message supplies {} parameter(s) but the statement uses ${}",
2187                        params.len(), n
2188                    ))
2189                }
2190            }
2191            i = j;
2192            continue;
2193        }
2194        out.push(c as char);
2195        i += 1;
2196    }
2197    Ok(out)
2198}
2199
2200/// A parsed statement, held for the life of the connection (or until `Close`).
2201struct Prepared {
2202    sql: String,
2203    /// OIDs advertised for `$1..$n` — what `ParameterDescription` reports and
2204    /// what `Bind` values are decoded as.
2205    param_oids: Vec<i32>,
2206    /// The advertised output shape, computed on demand and then reused.
2207    ///
2208    /// Lazy because working it out samples stored documents, and a text-format
2209    /// client that never sends `Describe(statement)` should not pay for a scan
2210    /// on every `Parse` — psycopg3 parses once per query.
2211    ///
2212    /// `Some(None)` means "computed, and this statement returns no rows".
2213    out_shape: Option<Option<(Vec<Col>, Vec<i32>)>>,
2214}
2215
2216/// The output columns and types a statement advertises, computed once.
2217fn prepared_shape<'a>(
2218    p: &'a mut Prepared,
2219    db: Option<&Arc<Db>>,
2220) -> &'a Option<(Vec<Col>, Vec<i32>)> {
2221    if p.out_shape.is_none() {
2222        p.out_shape = Some(describe_shape(&p.sql, db, p.param_oids.len()));
2223    }
2224    p.out_shape.as_ref().expect("just filled")
2225}
2226
2227/// A bound statement: fully substituted SQL plus, once run, its result.
2228struct Portal {
2229    sql: String,
2230    /// Filled by the first `Describe` or `Execute` and reused afterwards.
2231    ///
2232    /// Executing once and streaming from the buffer is what makes a suspended
2233    /// portal safe: a second `Execute` on a partially-drained `INSERT` must
2234    /// continue the row stream, not perform the insert again.
2235    result: Option<PortalResult>,
2236    /// The output shape, frozen at the first `Describe`/`Execute`.
2237    ///
2238    /// A schemaless store derives `SELECT *`'s columns from the rows it found,
2239    /// which would let a `Describe` and a later `Execute` disagree about the
2240    /// column count — and a driver that was told three fields and handed two
2241    /// mis-decodes the row rather than failing loudly. Freezing the shape and
2242    /// projecting every row onto it makes the result set rectangular, as SQL
2243    /// promises. The simple protocol keeps the dynamic behaviour, where there
2244    /// is no `Describe` to contradict.
2245    frozen: Option<Vec<Col>>,
2246    /// Result-column format codes requested by `Bind`. Empty = all text.
2247    formats: Vec<i16>,
2248    /// The shape this portal's statement advertised, carried over from the
2249    /// prepared statement when any column is to be sent in BINARY.
2250    ///
2251    /// It has to be the ADVERTISED shape rather than one derived from the rows
2252    /// in hand: asyncpg built its decoders from `Describe`, so re-deriving a
2253    /// different type here would hand it bytes it cannot read.
2254    declared: Option<(Vec<Col>, Vec<i32>)>,
2255}
2256
2257impl Portal {
2258    /// The format code for column `i`, following the protocol's shorthands:
2259    /// no codes means all-text, one code applies to every column.
2260    fn format_of(&self, i: usize) -> i16 {
2261        match self.formats.len() {
2262            0 => 0,
2263            1 => self.formats[0],
2264            _ => self.formats.get(i).copied().unwrap_or(0),
2265        }
2266    }
2267    /// The columns and types to advertise and encode with.
2268    fn shape(&self, r: &PortalResult) -> (Vec<Col>, Vec<i32>) {
2269        match &self.declared {
2270            Some((cols, oids)) if self.formats.iter().any(|f| *f == 1) => {
2271                (cols.clone(), oids.clone())
2272            }
2273            _ => {
2274                let cols = columns_for(&r.rows, &r.project);
2275                let oids = cols.iter().map(|c| oid_for(&r.rows, &c.src)).collect();
2276                (cols, oids)
2277            }
2278        }
2279    }
2280}
2281
2282struct PortalResult {
2283    rows: Vec<Value>,
2284    project: Vec<Col>,
2285    has_rows: bool,
2286    tag: String,
2287    tag_counts_rows: bool,
2288    /// How many rows have gone out across all `Execute`s on this portal.
2289    sent: usize,
2290}
2291
2292fn parse_complete() -> Vec<u8> { Out::msg(b'1').finish() }
2293fn bind_complete() -> Vec<u8> { Out::msg(b'2').finish() }
2294fn close_complete() -> Vec<u8> { Out::msg(b'3').finish() }
2295fn no_data() -> Vec<u8> { Out::msg(b'n').finish() }
2296fn portal_suspended() -> Vec<u8> { Out::msg(b's').finish() }
2297
2298fn parameter_description(oids: &[i32]) -> Vec<u8> {
2299    let mut m = Out::msg(b't');
2300    m.i16(oids.len() as i16);
2301    for o in oids {
2302        m.i32(*o);
2303    }
2304    m.finish()
2305}
2306
2307/// Split a NUL-terminated string off the front of a message body.
2308fn take_cstr(body: &[u8], at: &mut usize) -> String {
2309    let start = *at;
2310    while *at < body.len() && body[*at] != 0 {
2311        *at += 1;
2312    }
2313    let s = String::from_utf8_lossy(&body[start..*at]).to_string();
2314    if *at < body.len() {
2315        *at += 1; // step over the NUL
2316    }
2317    s
2318}
2319
2320fn take_i16(body: &[u8], at: &mut usize) -> Result<i16, String> {
2321    if *at + 2 > body.len() {
2322        return Err("truncated message".into());
2323    }
2324    let v = i16::from_be_bytes([body[*at], body[*at + 1]]);
2325    *at += 2;
2326    Ok(v)
2327}
2328
2329fn take_i32(body: &[u8], at: &mut usize) -> Result<i32, String> {
2330    if *at + 4 > body.len() {
2331        return Err("truncated message".into());
2332    }
2333    let v = i32::from_be_bytes([body[*at], body[*at + 1], body[*at + 2], body[*at + 3]]);
2334    *at += 4;
2335    Ok(v)
2336}
2337
2338/// The field names a collection actually holds, sampled from stored documents.
2339///
2340/// The answer to `SELECT *` on a store with no schema. Sorted, because
2341/// `serde_json`'s map is ordered and both this and the row encoder must agree
2342/// on column order or the values land under the wrong headings.
2343fn sample_columns(db: Option<&Arc<Db>>, coll: &str) -> Vec<Col> {
2344    let db = match db {
2345        Some(db) => db,
2346        None => return vec![],
2347    };
2348    let rows = match crate::nql::query(db, &format!("FROM {} LIMIT 25", coll)) {
2349        Ok((rows, _)) => rows,
2350        Err(_) => return vec![],
2351    };
2352    let mut names: Vec<String> = vec![];
2353    for r in &rows {
2354        if let Value::Object(m) = r {
2355            for k in m.keys() {
2356                if !names.iter().any(|n| n == k) {
2357                    names.push(k.clone());
2358                }
2359            }
2360        }
2361    }
2362    names.sort();
2363    names.iter().map(|n| Col::same(n)).collect()
2364}
2365
2366/// The result shape of a statement, worked out WITHOUT running it.
2367///
2368/// Needed for `Describe(statement)`, which arrives before any `Bind` — asyncpg
2369/// builds its row decoders from the answer. Only the select list is read off
2370/// the result; nothing touches storage except the type sampling.
2371///
2372/// Returns `None` when the statement returns no rows at all (`NoData`).
2373fn describe_shape(
2374    sql: &str,
2375    db: Option<&Arc<Db>>,
2376    n_params: usize,
2377) -> Option<(Vec<Col>, Vec<i32>)> {
2378    let probe = probe_sql(sql, n_params);
2379
2380    // The SQL evaluator describes its own output. It has to: `translate`
2381    // cannot parse a catalogue join at all, so without this a `Describe`
2382    // answered `NoData` — and a client told a SELECT has no output never
2383    // reads its rows.
2384    //
2385    // The probe is EXECUTED here, which is affordable precisely because this
2386    // path only serves catalogue relations and relation-free select lists.
2387    // Column types come from the values it actually produced, unified across
2388    // the rows by the same `oid_for` every other path uses — so a column
2389    // advertised `int8` is one the wire really encodes as int8.
2390    if sql_engine_owns(&probe) {
2391        if let Ok(Some((done, _))) = try_catalog_select(&probe, db) {
2392            if done.project.is_empty() {
2393                return None;
2394            }
2395            let oids = done
2396                .project
2397                .iter()
2398                .map(|c| oid_for(&done.rows, &c.src))
2399                .collect();
2400            return Some((done.project, oids));
2401        }
2402    }
2403
2404    let stmt = translate(&probe).ok()?;
2405    let coll = stmt_collection(sql);
2406
2407    let cols = match stmt {
2408        Stmt::Ok(_) => return None,
2409        Stmt::Canned { cols, .. } => cols.iter().map(|c| Col::same(c)).collect(),
2410        Stmt::Query { project, .. } => {
2411            if project.is_empty() { sample_columns(db, &coll) } else { project }
2412        }
2413        Stmt::Insert { returning, .. } | Stmt::Update { returning, .. } | Stmt::Delete { returning, .. } => {
2414            if !wants_returning(sql) {
2415                return None;
2416            }
2417            if returning.is_empty() { sample_columns(db, &coll) } else { returning }
2418        }
2419    };
2420    if cols.is_empty() {
2421        // Nothing could be determined. `NoData` is a lie for a SELECT, but a
2422        // RowDescription with zero columns is a worse one — it tells the client
2423        // the query definitively has no output.
2424        return None;
2425    }
2426    let oids = cols
2427        .iter()
2428        .map(|c| {
2429            aggregate_oid(&c.src, db, &coll)
2430                .unwrap_or_else(|| infer_field_oid(db, &coll, &c.src))
2431        })
2432        .collect();
2433    Some((cols, oids))
2434}
2435
2436/// A parse-only stand-in for a parameterised statement.
2437///
2438/// Substituting `NULL` was the obvious choice and the wrong one: a clause that
2439/// validates its argument rejects it, so `AS OF SYSTEM TIME $1` failed at
2440/// `Parse` — before the client ever bound a real sequence number. `0` parses
2441/// everywhere a literal can appear, and since only the SELECT list is read back
2442/// out, the stub's value never reaches an answer.
2443fn probe_sql(sql: &str, n_params: usize) -> String {
2444    let stub: Vec<Option<String>> = vec![Some("0".to_string()); n_params];
2445    substitute_params(sql, &stub).unwrap_or_else(|_| sql.to_string())
2446}
2447
2448/// Run a portal's statement if it has not run yet, then report its shape.
2449fn ensure_executed(
2450    portal: &mut Portal,
2451    db_name: &str,
2452    db: Option<&Arc<Db>>,
2453    read_only: bool,
2454) -> Result<(), Vec<u8>> {
2455    if portal.result.is_some() {
2456        return Ok(());
2457    }
2458    let ex = execute_stmt(&portal.sql, db_name, db, read_only)?;
2459    // Freeze the output shape on first sight so `Describe` and every later
2460    // `Execute` describe the same rectangle.
2461    let project = if let Some(f) = &portal.frozen {
2462        f.clone()
2463    } else {
2464        let p = if ex.project.is_empty() {
2465            columns_for(&ex.rows, &[])
2466        } else {
2467            ex.project.clone()
2468        };
2469        portal.frozen = Some(p.clone());
2470        p
2471    };
2472    portal.result = Some(PortalResult {
2473        rows: ex.rows,
2474        project,
2475        has_rows: ex.has_rows,
2476        tag: ex.tag,
2477        tag_counts_rows: ex.tag_counts_rows,
2478        sent: 0,
2479    });
2480    Ok(())
2481}
2482
2483// ── connection handling ─────────────────────────────────────────────────────
2484
2485async fn read_exact(sock: &mut TcpStream, n: usize) -> std::io::Result<Vec<u8>> {
2486    let mut buf = vec![0u8; n];
2487    sock.read_exact(&mut buf).await?;
2488    Ok(buf)
2489}
2490
2491async fn read_i32(sock: &mut TcpStream) -> std::io::Result<i32> {
2492    let b = read_exact(sock, 4).await?;
2493    Ok(i32::from_be_bytes([b[0], b[1], b[2], b[3]]))
2494}
2495
2496fn parse_startup_params(body: &[u8]) -> HashMap<String, String> {
2497    let mut out = HashMap::new();
2498    let mut parts = body.split(|b| *b == 0).map(|s| String::from_utf8_lossy(s).to_string());
2499    while let (Some(k), Some(v)) = (parts.next(), parts.next()) {
2500        if k.is_empty() {
2501            break;
2502        }
2503        out.insert(k, v);
2504    }
2505    out
2506}
2507
2508/// Serve one client connection to completion.
2509async fn handle(mut sock: TcpStream, resolver: Arc<dyn DbResolver>, read_only: bool) -> std::io::Result<()> {
2510    // ── startup, including the SSL negotiation clients try first ────────────
2511    let params = loop {
2512        let len = read_i32(&mut sock).await?;
2513        if len < 8 || len > 1 << 20 {
2514            return Ok(()); // nonsense framing — drop the connection
2515        }
2516        let code = read_i32(&mut sock).await?;
2517        let body = read_exact(&mut sock, (len - 8) as usize).await?;
2518        match code {
2519            SSL_REQUEST | GSS_REQUEST => {
2520                // Decline and let the client retry in the clear.
2521                sock.write_all(b"N").await?;
2522                continue;
2523            }
2524            CANCEL_REQUEST => return Ok(()), // nothing cancellable: reads are synchronous
2525            PROTO_V3 => break parse_startup_params(&body),
2526            other => {
2527                let major = other >> 16;
2528                sock.write_all(&err_msg(
2529                    "0A000",
2530                    &format!("unsupported frontend protocol {}.{} — this endpoint speaks 3.0",
2531                             major, other & 0xffff),
2532                )).await?;
2533                return Ok(());
2534            }
2535        }
2536    };
2537
2538    let db_name = params.get("database").cloned().unwrap_or_default();
2539
2540    // Resolve the database ONCE, here, on a blocking thread.
2541    //
2542    // A Postgres connection is bound to one database for its whole life, so
2543    // per-connection resolution is both correct and simpler than resolving per
2544    // statement — and it keeps the lock acquisition off the async worker.
2545    let resolved: Option<Arc<Db>> = {
2546        let r = Arc::clone(&resolver);
2547        let name = db_name.clone();
2548        tokio::task::spawn_blocking(move || r.resolve(&name))
2549            .await
2550            .unwrap_or(None)
2551    };
2552
2553    // ── auth: mirror the HTTP surface ───────────────────────────────────────
2554    if let Some(expected) = resolver.token() {
2555        // AuthenticationCleartextPassword (3)
2556        let mut m = Out::msg(b'R');
2557        m.i32(3);
2558        sock.write_all(&m.finish()).await?;
2559
2560        let tag = read_exact(&mut sock, 1).await?;
2561        if tag[0] != b'p' {
2562            sock.write_all(&err_msg("28000", "expected a password message")).await?;
2563            return Ok(());
2564        }
2565        let len = read_i32(&mut sock).await?;
2566        if len < 4 || len > 1 << 16 {
2567            return Ok(());
2568        }
2569        let body = read_exact(&mut sock, (len - 4) as usize).await?;
2570        let supplied = String::from_utf8_lossy(&body).trim_end_matches('\0').to_string();
2571        // Constant-time-ish: compare lengths and bytes without early return.
2572        let ok = supplied.len() == expected.len()
2573            && supplied.bytes().zip(expected.bytes()).fold(0u8, |a, (x, y)| a | (x ^ y)) == 0;
2574        if !ok {
2575            sock.write_all(&err_msg("28P01", "password authentication failed")).await?;
2576            return Ok(());
2577        }
2578    }
2579
2580    let mut m = Out::msg(b'R');
2581    m.i32(0); // AuthenticationOk
2582    sock.write_all(&m.finish()).await?;
2583
2584    for (k, v) in [
2585        ("server_version", SERVER_VERSION),
2586        ("server_encoding", "UTF8"),
2587        ("client_encoding", "UTF8"),
2588        ("DateStyle", "ISO, MDY"),
2589        ("integer_datetimes", "on"),
2590        ("standard_conforming_strings", "on"),
2591        ("application_name", "nedbd"),
2592    ] {
2593        let mut p = Out::msg(b'S');
2594        p.cstr(k);
2595        p.cstr(v);
2596        sock.write_all(&p.finish()).await?;
2597    }
2598    let mut k = Out::msg(b'K');
2599    k.i32(std::process::id() as i32);
2600    k.i32(0);
2601    sock.write_all(&k.finish()).await?;
2602    sock.write_all(&ready()).await?;
2603
2604    // ── message loop ────────────────────────────────────────────────────────
2605    //
2606    // Prepared statements and portals live for the connection. `""` is the
2607    // unnamed statement/portal, which every driver reuses constantly — it is an
2608    // ordinary entry in the map rather than a special case.
2609    let mut prepared: HashMap<String, Prepared> = HashMap::new();
2610    let mut portals: HashMap<String, Portal> = HashMap::new();
2611    // After an error inside an extended-protocol sequence, everything up to the
2612    // next `Sync` is discarded. Skipping this is how a server ends up answering
2613    // a Bind the client has already abandoned, and the stream desynchronises.
2614    let mut failed = false;
2615
2616    loop {
2617        let mut tag = [0u8; 1];
2618        if sock.read_exact(&mut tag).await.is_err() {
2619            return Ok(()); // client hung up
2620        }
2621        let len = read_i32(&mut sock).await?;
2622        if len < 4 || len > 64 << 20 {
2623            return Ok(());
2624        }
2625        let body = read_exact(&mut sock, (len - 4) as usize).await?;
2626
2627        // `Sync` always clears the error state; `Terminate` always applies.
2628        if failed && tag[0] != b'S' && tag[0] != b'X' {
2629            continue;
2630        }
2631
2632        match tag[0] {
2633            b'X' => return Ok(()), // Terminate
2634
2635            b'Q' => {
2636                let sql = String::from_utf8_lossy(&body).trim_end_matches('\0').to_string();
2637                let out = run_simple_query(&sql, &db_name, resolved.as_ref(), read_only);
2638                sock.write_all(&out).await?;
2639                sock.write_all(&ready()).await?;
2640                // A simple query closes the unnamed portal, per the protocol.
2641                portals.remove("");
2642            }
2643
2644            // ── Parse: name, SQL, declared parameter type OIDs ─────────────
2645            b'P' => {
2646                let mut at = 0usize;
2647                let name = take_cstr(&body, &mut at);
2648                let sql = take_cstr(&body, &mut at);
2649                let n = take_i16(&body, &mut at).unwrap_or(0).max(0) as usize;
2650                let mut declared = Vec::with_capacity(n);
2651                let mut bad = false;
2652                for _ in 0..n {
2653                    match take_i32(&body, &mut at) {
2654                        Ok(o) => declared.push(o),
2655                        Err(_) => { bad = true; break; }
2656                    }
2657                }
2658                if bad {
2659                    sock.write_all(&err_msg("08P01", "malformed Parse message")).await?;
2660                    failed = true;
2661                    continue;
2662                }
2663                // Reject unsupported SQL here rather than at Execute, so the
2664                // client learns at the point it asked — which is also where
2665                // Postgres reports it.
2666                //
2667                // The SQL evaluator gets asked first, or a catalogue query
2668                // would be refused at `Parse` by the NQL path that was never
2669                // going to run it — and the extended protocol is where every
2670                // ORM and async driver lives, so refusing here refuses them
2671                // all.
2672                let probe = probe_sql(&sql, param_count(&sql));
2673                if !sql_engine_owns(&probe) {
2674                    if let Err(why) = translate(&probe) {
2675                        sock.write_all(&err_msg("0A000", &why)).await?;
2676                        failed = true;
2677                        continue;
2678                    }
2679                }
2680                let param_oids = infer_param_oids(&sql, &declared, resolved.as_ref());
2681                prepared.insert(name, Prepared { sql, param_oids, out_shape: None });
2682                sock.write_all(&parse_complete()).await?;
2683            }
2684
2685            // ── Bind: portal, statement, formats, values, result formats ───
2686            b'B' => {
2687                let mut at = 0usize;
2688                let portal_name = take_cstr(&body, &mut at);
2689                let stmt_name = take_cstr(&body, &mut at);
2690                if !prepared.contains_key(&stmt_name) {
2691                    sock.write_all(&err_msg("26000", &format!(
2692                        "prepared statement {:?} does not exist", stmt_name))).await?;
2693                    failed = true;
2694                    continue;
2695                }
2696                let p = &prepared[&stmt_name];
2697                let mut want_formats: Vec<i16> = vec![];
2698                let res: Result<String, String> = (|| {
2699                    let nfmt = take_i16(&body, &mut at)? .max(0) as usize;
2700                    let mut fmts = Vec::with_capacity(nfmt);
2701                    for _ in 0..nfmt {
2702                        fmts.push(take_i16(&body, &mut at)?);
2703                    }
2704                    let nparam = take_i16(&body, &mut at)?.max(0) as usize;
2705                    let mut vals: Vec<Option<String>> = Vec::with_capacity(nparam);
2706                    for i in 0..nparam {
2707                        let l = take_i32(&body, &mut at)?;
2708                        let raw: Option<Vec<u8>> = if l < 0 {
2709                            None
2710                        } else {
2711                            let l = l as usize;
2712                            if at + l > body.len() {
2713                                return Err("truncated Bind parameter".into());
2714                            }
2715                            let v = body[at..at + l].to_vec();
2716                            at += l;
2717                            Some(v)
2718                        };
2719                        // Zero format codes means "all text"; one means "this
2720                        // format for every parameter"; otherwise one per value.
2721                        let f = match fmts.len() {
2722                            0 => 0,
2723                            1 => fmts[0],
2724                            _ => *fmts.get(i).unwrap_or(&0),
2725                        };
2726                        let oid = *p.param_oids.get(i).unwrap_or(&OID_TEXT);
2727                        vals.push(decode_param(raw.as_deref(), oid, f)?);
2728                    }
2729                    // Result format codes. asyncpg asks for binary on every
2730                    // column, so honouring these is not an optimisation — it
2731                    // is the difference between asyncpg reading rows and
2732                    // refusing the result outright.
2733                    let nres = take_i16(&body, &mut at)?.max(0) as usize;
2734                    for _ in 0..nres {
2735                        let f = take_i16(&body, &mut at)?;
2736                        if f != 0 && f != 1 {
2737                            return Err(format!("unknown result format code {}", f));
2738                        }
2739                        want_formats.push(f);
2740                    }
2741                    substitute_params(&p.sql, &vals)
2742                })();
2743                match res {
2744                    Ok(sql) => {
2745                        // Binary encoding must use the types the client was
2746                        // TOLD about, so pull the advertised shape across.
2747                        let declared = if want_formats.iter().any(|f| *f == 1) {
2748                            let p = prepared.get_mut(&stmt_name).expect("checked above");
2749                            prepared_shape(p, resolved.as_ref()).clone()
2750                        } else {
2751                            None
2752                        };
2753                        portals.insert(portal_name, Portal {
2754                            sql, result: None, frozen: None,
2755                            formats: want_formats, declared,
2756                        });
2757                        sock.write_all(&bind_complete()).await?;
2758                    }
2759                    Err(why) => {
2760                        sock.write_all(&err_msg("08P01", &why)).await?;
2761                        failed = true;
2762                    }
2763                }
2764            }
2765
2766            // ── Describe: 'S' statement, or 'P' portal ─────────────────────
2767            b'D' => {
2768                let kind = body.first().copied().unwrap_or(b'S');
2769                let mut at = 1usize;
2770                let name = take_cstr(&body, &mut at);
2771                if kind == b'S' {
2772                    if !prepared.contains_key(&name) {
2773                        sock.write_all(&err_msg("26000", &format!(
2774                            "prepared statement {:?} does not exist", name))).await?;
2775                        failed = true;
2776                        continue;
2777                    }
2778                    let p = prepared.get_mut(&name).expect("checked above");
2779                    let oids = p.param_oids.clone();
2780                    // asyncpg encodes its arguments from this, so the count has
2781                    // to be right or it refuses the call before sending a Bind.
2782                    sock.write_all(&parameter_description(&oids)).await?;
2783                    // Describe(statement) happens before Bind, so the requested
2784                    // result format is not known yet; Postgres reports text
2785                    // here too and the client's own Bind decides the encoding.
2786                    let out = match prepared_shape(p, resolved.as_ref()) {
2787                        Some((cols, col_oids)) => row_description(cols, col_oids),
2788                        None => no_data(),
2789                    };
2790                    sock.write_all(&out).await?;
2791                } else {
2792                    let portal = match portals.get_mut(&name) {
2793                        Some(p) => p,
2794                        None => {
2795                            sock.write_all(&err_msg("34000", &format!(
2796                                "portal {:?} does not exist", name))).await?;
2797                            failed = true;
2798                            continue;
2799                        }
2800                    };
2801                    // A bound portal can be run: doing it here means the
2802                    // RowDescription reports the columns and types actually
2803                    // present, which is strictly better than a guess. psycopg3
2804                    // takes this path on every query.
2805                    match ensure_executed(portal, &db_name, resolved.as_ref(), read_only) {
2806                        Err(encoded) => {
2807                            sock.write_all(&encoded).await?;
2808                            failed = true;
2809                        }
2810                        Ok(()) => {
2811                            let r = portal.result.as_ref().expect("just executed");
2812                            if !r.has_rows {
2813                                sock.write_all(&no_data()).await?;
2814                            } else {
2815                                let (cols, oids) = portal.shape(r);
2816                                let fmts: Vec<i16> =
2817                                    (0..cols.len()).map(|i| portal.format_of(i)).collect();
2818                                sock.write_all(&row_description_fmt(&cols, &oids, &fmts)).await?;
2819                            }
2820                        }
2821                    }
2822                }
2823            }
2824
2825            // ── Execute: portal, maximum rows (0 = all) ────────────────────
2826            b'E' => {
2827                let mut at = 0usize;
2828                let name = take_cstr(&body, &mut at);
2829                let max_rows = take_i32(&body, &mut at).unwrap_or(0);
2830                let portal = match portals.get_mut(&name) {
2831                    Some(p) => p,
2832                    None => {
2833                        sock.write_all(&err_msg("34000", &format!(
2834                            "portal {:?} does not exist", name))).await?;
2835                        failed = true;
2836                        continue;
2837                    }
2838                };
2839                if let Err(encoded) = ensure_executed(portal, &db_name, resolved.as_ref(), read_only) {
2840                    sock.write_all(&encoded).await?;
2841                    failed = true;
2842                    continue;
2843                }
2844                let r = portal.result.as_ref().expect("just executed");
2845                if !r.has_rows {
2846                    let tag = r.tag.clone();
2847                    sock.write_all(&command_complete(&tag)).await?;
2848                    continue;
2849                }
2850                let (cols, oids) = portal.shape(r);
2851                let limit = if max_rows > 0 {
2852                    (r.sent + max_rows as usize).min(r.rows.len())
2853                } else {
2854                    r.rows.len()
2855                };
2856                // Encode the whole batch BEFORE writing any of it. A value that
2857                // cannot be sent in the advertised binary type has to become an
2858                // error instead of a truncated row stream — half a result set
2859                // followed by an error is far harder to diagnose than an error.
2860                let mut encoded: Vec<Vec<u8>> = Vec::with_capacity(limit - r.sent);
2861                let mut fail: Option<String> = None;
2862                for row in &r.rows[r.sent..limit] {
2863                    let mut vals: Vec<Option<Vec<u8>>> = Vec::with_capacity(cols.len());
2864                    for (i, c) in cols.iter().enumerate() {
2865                        let v = row.get(&c.src);
2866                        let got = if portal.format_of(i) == 1 {
2867                            cell_binary(v, oids.get(i).copied().unwrap_or(OID_TEXT))
2868                                .map_err(|e| format!("column {:?}: {}", c.out, e))
2869                        } else {
2870                            Ok(cell(v).map(|s| s.into_bytes()))
2871                        };
2872                        match got {
2873                            Ok(b) => vals.push(b),
2874                            Err(e) => { fail = Some(e); break; }
2875                        }
2876                    }
2877                    if fail.is_some() {
2878                        break;
2879                    }
2880                    encoded.push(data_row_bytes(&vals));
2881                }
2882                if let Some(why) = fail {
2883                    sock.write_all(&err_msg("22P03", &why)).await?;
2884                    failed = true;
2885                    continue;
2886                }
2887                let mut out = vec![];
2888                for e in &encoded {
2889                    out.extend_from_slice(e);
2890                }
2891                let r = portal.result.as_mut().expect("just executed");
2892                r.sent = limit;
2893                // More rows left and the client capped the batch: suspend the
2894                // portal instead of completing it. This is what a JDBC
2895                // `setFetchSize` and a psycopg3 server-side cursor rely on.
2896                if max_rows > 0 && r.sent < r.rows.len() {
2897                    out.extend_from_slice(&portal_suspended());
2898                } else {
2899                    let tag = if r.tag_counts_rows {
2900                        format!("{} {}", r.tag, r.sent)
2901                    } else {
2902                        r.tag.clone()
2903                    };
2904                    out.extend_from_slice(&command_complete(&tag));
2905                }
2906                sock.write_all(&out).await?;
2907            }
2908
2909            // ── Close: 'S' statement, or 'P' portal ───────────────────────
2910            b'C' => {
2911                let kind = body.first().copied().unwrap_or(b'S');
2912                let mut at = 1usize;
2913                let name = take_cstr(&body, &mut at);
2914                if kind == b'S' {
2915                    prepared.remove(&name);
2916                } else {
2917                    portals.remove(&name);
2918                }
2919                // Closing something that was never open is explicitly not an
2920                // error in the protocol.
2921                sock.write_all(&close_complete()).await?;
2922            }
2923
2924            // Flush: everything is written unbuffered already, so this is a
2925            // no-op — but it must NOT produce a ReadyForQuery, or a client that
2926            // flushes mid-sequence (asyncpg does, after Describe) loses sync.
2927            b'H' => {}
2928
2929            b'S' => {
2930                failed = false;
2931                sock.write_all(&ready()).await?;
2932            }
2933
2934            other => {
2935                sock.write_all(&err_msg(
2936                    "08P01",
2937                    &format!("unexpected frontend message {:?}", other as char),
2938                )).await?;
2939                failed = true;
2940            }
2941        }
2942    }
2943}
2944
2945const READ_ONLY_MSG: &str =
2946    "this endpoint is running read-only (NEDBD_PG_READ_ONLY=1). Writes are \
2947     implemented but disabled on this server — unset the flag to allow them.";
2948
2949fn no_db(db_name: &str) -> Vec<u8> {
2950    err_msg("3D000", &format!(
2951        "database {:?} is not open on this server — create it first \
2952         (POST /v1/databases), or connect with -d <name>", db_name))
2953}
2954
2955/// `pg_catalog.pg_class` → `pg_class`, but `information_schema.tables` keeps
2956/// its qualifier, because `tables` is a plausible collection name and the
2957/// catalogue must never shadow a user's own data.
2958fn catalog_name(n: &str) -> String {
2959    let joined: Vec<&str> = n.split('.').collect();
2960    if joined.len() >= 2 && joined[joined.len() - 2] == "information_schema" {
2961        format!("information_schema.{}", joined[joined.len() - 1])
2962    } else {
2963        joined[joined.len() - 1].to_string()
2964    }
2965}
2966
2967/// Does the SQL evaluator own this statement?
2968///
2969/// Two ways in. The first is obvious: it reads a catalogue relation.
2970///
2971/// The second is a statement with NO relation at all — a select list of
2972/// literals and scalar function calls, which is exactly what this evaluator
2973/// does and which the SQL→NQL path cannot express (NQL is FROM-first). That
2974/// path answers a handful of EXACT spellings from a canned table
2975/// (`SELECT 1`, `SELECT VERSION()`, `SELECT CURRENT_SCHEMA`), and those
2976/// answers are what existing clients already see — so this predicate rescues
2977/// only what it REFUSES, leaving every spelling it does handle alone.
2978///
2979/// That gap was not hypothetical. SQLAlchemy's PostgreSQL dialect opens every
2980/// connection with `select pg_catalog.version()`, which is one character of
2981/// qualification away from the canned `SELECT VERSION()` and therefore missed
2982/// it — so the engine refused the first statement of dialect initialisation
2983/// and NO SQLAlchemy application could connect at all. A canned list of
2984/// spellings is the same brittleness `pgcatalog` exists to avoid; the fix is
2985/// to let the evaluator answer, because it has `version()`,
2986/// `current_setting()` and the rest as real functions.
2987///
2988/// Cheap: one parse, no execution, no storage access.
2989fn sql_engine_owns(sql: &str) -> bool {
2990    let Ok(sel) = crate::sqlselect::parse(sql) else { return false };
2991    let touched = sel.base_relations();
2992    if touched.is_empty() {
2993        return translate(sql).is_err();
2994    }
2995    touched.iter().any(|t| crate::pgcatalog::is_catalog(&catalog_name(t)))
2996}
2997
2998/// Run a `SELECT` through the full SQL engine when it touches the catalogue.
2999///
3000/// The gate is deliberately narrow: a statement goes to `sqlselect` only when
3001/// one of its tables is a catalogue relation. Everything else keeps the
3002/// SQL→NQL path, which has the index pushdown, `AS OF`, `TRACE` and the
3003/// bounded scans — and whose join story is a real planning question rather
3004/// than a nested loop. Routing a large collection through a nested-loop join
3005/// would be a promise this engine cannot keep.
3006///
3007/// `None` means "not mine": the caller falls through to the ordinary path, so
3008/// the error the client sees is the ordinary path's error rather than a
3009/// confusing one from a parser that was never meant to handle the statement.
3010fn try_catalog_select(
3011    sql: &str,
3012    db: Option<&Arc<Db>>,
3013) -> Result<Option<(Executed, crate::sqlplan::Plan)>, Vec<u8>> {
3014    let sel = match crate::sqlselect::parse(sql) {
3015        Ok(sel) => sel,
3016        Err(why) => {
3017            // A statement that plainly reads the catalogue but that this
3018            // engine cannot parse gets the PARSE error, not the NQL path's.
3019            //
3020            // Falling through unconditionally produced an actively false
3021            // message: `\d` and `\dp` were told "JOIN is not supported",
3022            // which stopped being true the moment joins started working — and
3023            // a wrong explanation is worse than a blunt one, because it sends
3024            // the reader to fix the wrong thing.
3025            if mentions_catalog(sql) {
3026                return Err(err_msg("0A000", &format!(
3027                    "this catalogue query uses SQL this endpoint does not \
3028                     implement: {}", why)));
3029            }
3030            return Ok(None);
3031        }
3032    };
3033
3034    // Which relations does it read — at ANY depth? `\dd` names its catalogue
3035    // relations only inside a derived table, and `\dT` only inside two
3036    // subqueries; a walk over the top-level FROM list alone would route both
3037    // to the NQL path, which cannot parse them and would report an error that
3038    // sends the reader to fix the wrong thing.
3039    if !sql_engine_owns(sql) {
3040        return Ok(None);
3041    }
3042
3043    let resolve = |name: &str| -> anyhow::Result<Option<Box<dyn crate::sqlselect::Relation>>> {
3044        let cname = catalog_name(name);
3045        if let Some(rows) = crate::pgcatalog::rows(&cname, db) {
3046            // A synthesised catalogue relation is small and built eagerly;
3047            // wrapping it satisfies the streaming contract without pretending
3048            // it is lazy.
3049            return Ok(Some(crate::sqlselect::from_vec(rows)));
3050        }
3051        // A join between a catalogue relation and a real collection is
3052        // legitimate, so a user table still resolves.
3053        //
3054        // NOTE: `nql::query` materialises the whole collection, so this side
3055        // is eager even though the evaluator no longer requires it to be.
3056        // Making the storage scan itself lazy is the other half of the work
3057        // and is tracked in HANDOFF — stated here so nobody reads the
3058        // streaming interface as a claim that storage is already streaming.
3059        match db {
3060            Some(db) => match crate::nql::query(db, &format!("FROM {}", cname)) {
3061                Ok((rows, _)) => Ok(Some(crate::sqlselect::from_vec(rows))),
3062                Err(_) => Ok(None),
3063            },
3064            None => Ok(None),
3065        }
3066    };
3067
3068    let (cols, rows, plan) = crate::sqlselect::execute_explain(
3069        &sel,
3070        &resolve,
3071        crate::sqljoin::JoinExec::Auto,
3072    )
3073    .map_err(|e| err_msg("42601", &e.to_string()))?;
3074
3075    Ok(Some((
3076        Executed {
3077            rows,
3078            // The KEY is what the row is stored under; the NAME is what the
3079            // client sees. They differ when a select list has duplicate output
3080            // names, which PostgreSQL permits and generated SQL relies on.
3081            project: cols
3082                .iter()
3083                .map(|c| Col::renamed(&c.key, &c.name))
3084                .collect(),
3085            has_rows: true,
3086            tag: "SELECT".into(),
3087            tag_counts_rows: true,
3088        },
3089        plan,
3090    )))
3091}
3092
3093/// Strip a leading `EXPLAIN`, returning the statement it wraps.
3094///
3095/// `ANALYZE` and `VERBOSE` are accepted and ignored: this endpoint always
3096/// executes and always reports actual rows, so `EXPLAIN` and
3097/// `EXPLAIN ANALYZE` genuinely do the same thing here. Accepting the keyword
3098/// and silently doing the honest thing beats refusing a client's spelling.
3099fn strip_explain(sql: &str) -> Option<&str> {
3100    let t = sql.trim().trim_end_matches(';').trim();
3101    let mut rest = t.strip_prefix("EXPLAIN").or_else(|| t.strip_prefix("explain"))?;
3102    // Require a word boundary so `EXPLAINED` is not mistaken for a keyword.
3103    if !rest.starts_with(char::is_whitespace) {
3104        return None;
3105    }
3106    rest = rest.trim_start();
3107    loop {
3108        let low = rest.to_lowercase();
3109        if let Some(r) = low.strip_prefix("analyze").or_else(|| low.strip_prefix("analyse")) {
3110            if r.starts_with(char::is_whitespace) || r.is_empty() {
3111                rest = rest[rest.len() - r.len()..].trim_start();
3112                continue;
3113            }
3114        }
3115        if let Some(r) = low.strip_prefix("verbose") {
3116            if r.starts_with(char::is_whitespace) || r.is_empty() {
3117                rest = rest[rest.len() - r.len()..].trim_start();
3118                continue;
3119            }
3120        }
3121        break;
3122    }
3123    Some(rest)
3124}
3125
3126/// One text column named `QUERY PLAN`, which is exactly the shape PostgreSQL
3127/// returns — so `psql` prints it without special handling.
3128fn plan_result(lines: Vec<String>) -> Executed {
3129    Executed {
3130        rows: lines
3131            .into_iter()
3132            .map(|l| serde_json::json!({ "QUERY PLAN": l }))
3133            .collect(),
3134        project: vec![Col::same("QUERY PLAN")],
3135        has_rows: true,
3136        tag: "EXPLAIN".into(),
3137        tag_counts_rows: false,
3138    }
3139}
3140
3141/// Does the raw SQL plainly read a catalogue relation?
3142///
3143/// A cheap text check, used only to decide WHICH error to report when the
3144/// statement cannot be parsed — never to decide what a parsable statement
3145/// means. `pg_` is the giveaway: every catalogue relation is prefixed, and so
3146/// is the `pg_catalog` schema qualifier.
3147fn mentions_catalog(sql: &str) -> bool {
3148    let low = sql.to_lowercase();
3149    low.contains("pg_catalog.")
3150        || low.contains("information_schema.")
3151        || low.contains("from pg_")
3152        || low.contains("join pg_")
3153}
3154
3155/// The catalogue relation a translated query reads from, if any.
3156///
3157/// Reads the collection straight off the parsed NQL rather than re-parsing the
3158/// SQL, so it cannot disagree with what the executor is about to run.
3159fn catalog_target(nql: &str) -> Option<String> {
3160    let coll = crate::nql::parse(nql).ok()?.coll;
3161    if crate::pgcatalog::is_catalog(&coll) {
3162        Some(coll)
3163    } else {
3164        None
3165    }
3166}
3167
3168/// True when the statement carried a RETURNING clause. Checked against the raw
3169/// SQL because `RETURNING *` yields an EMPTY projection, which is otherwise
3170/// indistinguishable from "no RETURNING at all".
3171fn wants_returning(sql: &str) -> bool {
3172    find_kw(&sql.to_uppercase(), "RETURNING").is_some()
3173}
3174
3175/// A unique key for a server-assigned INSERT id.
3176fn next_row_id() -> String {
3177    use std::sync::atomic::{AtomicU64, Ordering};
3178    static N: AtomicU64 = AtomicU64::new(0);
3179    let n = N.fetch_add(1, Ordering::Relaxed);
3180    let ts = std::time::SystemTime::now()
3181        .duration_since(std::time::UNIX_EPOCH)
3182        .map(|d| d.as_micros())
3183        .unwrap_or(0);
3184    format!("r{}{}", ts, n)
3185}
3186
3187/// One executed statement, held apart from any wire encoding.
3188///
3189/// This type is why the simple and extended protocols share an execution path
3190/// rather than growing two copies of the SQL→NEDB semantics. The simple path
3191/// encodes it immediately; the extended path parks it in a portal and dribbles
3192/// the rows out across successive `Execute` messages. Both get identical
3193/// answers because both call `execute_stmt`.
3194pub struct Executed {
3195    /// The rows the client gets — a SELECT's result, or a write's `RETURNING`.
3196    pub rows: Vec<Value>,
3197    /// How to project them (empty = every key in the row).
3198    pub project: Vec<Col>,
3199    /// Whether the client asked for rows at all. Distinct from `rows.is_empty()`:
3200    /// a `SELECT` matching nothing still owes a `RowDescription`, while an
3201    /// `UPDATE` without `RETURNING` owes `NoData`.
3202    pub has_rows: bool,
3203    /// The command tag, already rendered — except for a SELECT, where the row
3204    /// count is only known once the rows have actually been sent.
3205    pub tag: String,
3206    /// True when `tag` is a SELECT-shaped tag whose count is the rows sent.
3207    pub tag_counts_rows: bool,
3208}
3209
3210impl Executed {
3211    fn nothing(tag: &str) -> Self {
3212        Executed { rows: vec![], project: vec![], has_rows: false, tag: tag.to_string(), tag_counts_rows: false }
3213    }
3214    /// Render the final `CommandComplete` given how many rows went out.
3215    fn tag_for(&self, sent: usize) -> String {
3216        if self.tag_counts_rows { format!("{} {}", self.tag, sent) } else { self.tag.clone() }
3217    }
3218}
3219
3220/// Run ONE statement. `Err` carries an already-encoded `ErrorResponse`.
3221///
3222/// Every SQL→NEDB decision lives here, which is the point: the extended query
3223/// protocol added below is then purely a matter of message framing, and cannot
3224/// drift from the simple path's semantics.
3225fn execute_stmt(
3226    stmt_sql: &str,
3227    db_name: &str,
3228    db: Option<&Arc<Db>>,
3229    read_only: bool,
3230) -> Result<Executed, Vec<u8>> {
3231    // The full SQL engine gets first refusal, but ONLY for statements that
3232    // touch the catalogue — see `try_catalog_select`. It has to run before
3233    // `translate`, because `translate` targets NQL and NQL cannot express a
3234    // join, a CASE or a scalar function at all.
3235    // EXPLAIN reports which engine would run the statement, and a plan only
3236    // when the SQL evaluator is the engine that actually runs it. Describing a
3237    // pipeline the statement would not take is the one thing an EXPLAIN must
3238    // never do.
3239    if let Some(inner) = strip_explain(stmt_sql) {
3240        if let Some((_, plan)) = try_catalog_select(inner, db)? {
3241            return Ok(plan_result(plan.render()));
3242        }
3243        let mut lines = vec![];
3244        match translate(inner) {
3245            Ok(_) => {
3246                lines.push(
3247                    "NQL path — this statement is translated to NQL and \
3248                     executed by the storage engine, not by the SQL evaluator."
3249                        .to_string(),
3250                );
3251                lines.push(
3252                    "No plan is reported, because the SQL evaluator is not \
3253                     what runs it. Reporting one would describe a pipeline \
3254                     that never executed."
3255                        .to_string(),
3256                );
3257                lines.push(
3258                    "The SQL evaluator (joins, CASE, scalar functions, a \
3259                     hash-join planner) currently serves catalogue queries."
3260                        .to_string(),
3261                );
3262            }
3263            Err(why) => lines.push(format!("cannot be executed: {why}")),
3264        }
3265        return Ok(plan_result(lines));
3266    }
3267
3268    if let Some((done, _plan)) = try_catalog_select(stmt_sql, db)? {
3269        return Ok(done);
3270    }
3271
3272    let stmt = translate(stmt_sql).map_err(|why| err_msg("0A000", &why))?;
3273
3274    // Every arm below that touches storage needs a database; resolve the
3275    // "no such database" answer once instead of at each use.
3276    macro_rules! need_db {
3277        () => {
3278            match db {
3279                Some(db) => db,
3280                None => return Err(no_db(db_name)),
3281            }
3282        };
3283    }
3284    macro_rules! need_write {
3285        () => {
3286            if read_only {
3287                return Err(err_msg("25006", READ_ONLY_MSG));
3288            }
3289        };
3290    }
3291
3292    match stmt {
3293        Stmt::Ok(tag) => Ok(Executed::nothing(if tag.is_empty() { "SELECT 0" } else { tag })),
3294
3295        Stmt::Canned { cols, row } => {
3296            // Fold the canned answer into an ordinary row so the encoders,
3297            // the portal machinery and `Describe` all see one shape.
3298            let mut obj = serde_json::Map::new();
3299            for (c, v) in cols.iter().zip(row.iter()) {
3300                obj.insert(c.clone(), Value::String(v.clone()));
3301            }
3302            Ok(Executed {
3303                rows: vec![Value::Object(obj)],
3304                project: cols.iter().map(|c| Col::same(c)).collect(),
3305                has_rows: true,
3306                tag: "SELECT".into(),
3307                tag_counts_rows: true,
3308            })
3309        }
3310
3311        Stmt::Query { nql, project } => {
3312            // A catalogue relation is synthesised from the live database
3313            // rather than read from it — but it is still queried with the
3314            // ORDINARY predicate path, so WHERE / ORDER BY / LIMIT and the
3315            // `~` operators work on it because they are the same operators.
3316            //
3317            // Checked BEFORE `need_db!()`: `SELECT * FROM pg_namespace` has to
3318            // answer even when the client connected without naming a database,
3319            // which is exactly what psql does on startup. Refusing there is
3320            // how "psql cannot connect" starts.
3321            if let Some(coll) = catalog_target(&nql) {
3322                let rows = crate::pgcatalog::rows(&coll, db)
3323                    .expect("catalog_target only returns names pgcatalog serves");
3324                let rows = crate::nql::query_rows(rows, &nql)
3325                    .map_err(|e| err_msg("42601", &e.to_string()))?;
3326                return Ok(Executed {
3327                    rows, project, has_rows: true,
3328                    tag: "SELECT".into(), tag_counts_rows: true,
3329                });
3330            }
3331            let db = need_db!();
3332            let (rows, _) = crate::nql::query(db, &nql).map_err(|e| {
3333                err_msg("42601", &format!("{} (translated to NQL: {})", e, nql))
3334            })?;
3335            Ok(Executed { rows, project, has_rows: true, tag: "SELECT".into(), tag_counts_rows: true })
3336        }
3337
3338        Stmt::Insert { coll, rows, returning } => {
3339            let db = need_db!();
3340            need_write!();
3341            let mut written: Vec<Value> = vec![];
3342            for (i, r) in rows.iter().enumerate() {
3343                // The engine requires an id. When the statement did not supply
3344                // one, mint a unique key rather than silently overwriting a
3345                // shared default.
3346                let id = match &r.id {
3347                    Some(id) => id.clone(),
3348                    None => format!("{}-{}", next_row_id(), i),
3349                };
3350                let node = db
3351                    .put(&coll, &id, Value::Object(r.doc.clone()),
3352                         r.caused_by.clone(), r.valid_from.clone(), r.valid_to.clone())
3353                    .map_err(|e| err_msg("XX000", &format!("INSERT failed: {}", e)))?;
3354                written.push(crate::nql::node_to_json(&node));
3355            }
3356            let n = written.len();
3357            let has_rows = wants_returning(stmt_sql);
3358            Ok(Executed {
3359                rows: if has_rows { written } else { vec![] },
3360                project: returning,
3361                has_rows,
3362                // Postgres reports `INSERT <oid> <rows>`; the oid is always 0.
3363                tag: format!("INSERT 0 {}", n),
3364                tag_counts_rows: false,
3365            })
3366        }
3367
3368        Stmt::Update { coll, set, nql, returning } => {
3369            let db = need_db!();
3370            need_write!();
3371            // Matching rows come from an ordinary NQL read, so the whole
3372            // predicate surface works inside an UPDATE.
3373            let (matched, _) = crate::nql::query(db, &nql).map_err(|e| {
3374                err_msg("42601", &format!("{} (translated to NQL: {})", e, nql))
3375            })?;
3376            let mut written: Vec<Value> = vec![];
3377            for row in &matched {
3378                let id = match row.get("_id").and_then(|v| v.as_str()) {
3379                    Some(id) => id.to_string(),
3380                    None => continue,
3381                };
3382                // Merge onto the CURRENT stored document, not onto the query
3383                // row: a query row carries injected `_`-prefixed metadata that
3384                // must never be written back into the payload.
3385                let mut doc = match db.get(&coll, &id) {
3386                    Some(n) => match n.data {
3387                        Value::Object(m) => m,
3388                        _ => serde_json::Map::new(),
3389                    },
3390                    None => continue,
3391                };
3392                for (k, v) in &set {
3393                    doc.insert(k.clone(), v.clone());
3394                }
3395                // An UPDATE is a NEW VERSION — the prior value stays readable
3396                // with AS OF SYSTEM TIME. That is the whole point.
3397                let node = db
3398                    .put(&coll, &id, Value::Object(doc), vec![], None, None)
3399                    .map_err(|e| err_msg("XX000", &format!("UPDATE failed: {}", e)))?;
3400                written.push(crate::nql::node_to_json(&node));
3401            }
3402            let n = written.len();
3403            let has_rows = wants_returning(stmt_sql);
3404            Ok(Executed {
3405                rows: if has_rows { written } else { vec![] },
3406                project: returning,
3407                has_rows,
3408                tag: format!("UPDATE {}", n),
3409                tag_counts_rows: false,
3410            })
3411        }
3412
3413        Stmt::Delete { coll, nql, returning } => {
3414            let db = need_db!();
3415            need_write!();
3416            let (matched, _) = crate::nql::query(db, &nql).map_err(|e| {
3417                err_msg("42601", &format!("{} (translated to NQL: {})", e, nql))
3418            })?;
3419            // RETURNING must be captured BEFORE the delete: after the tombstone
3420            // the row is no longer readable by id.
3421            let returned = matched.clone();
3422            let mut n = 0usize;
3423            for row in &matched {
3424                if let Some(id) = row.get("_id").and_then(|v| v.as_str()) {
3425                    match db.delete(&coll, id) {
3426                        Ok(true) => n += 1,
3427                        Ok(false) => {}
3428                        Err(e) => return Err(err_msg("XX000", &format!("DELETE failed: {}", e))),
3429                    }
3430                }
3431            }
3432            let has_rows = wants_returning(stmt_sql);
3433            Ok(Executed {
3434                rows: if has_rows { returned } else { vec![] },
3435                project: returning,
3436                has_rows,
3437                tag: format!("DELETE {}", n),
3438                tag_counts_rows: false,
3439            })
3440        }
3441    }
3442}
3443
3444/// Execute a simple-query payload, which may hold several `;`-separated statements.
3445fn run_simple_query(sql: &str, db_name: &str, db: Option<&Arc<Db>>, read_only: bool) -> Vec<u8> {
3446    let mut out = vec![];
3447    let statements = split_statements(sql);
3448    if statements.is_empty() {
3449        // EmptyQueryResponse
3450        return Out::msg(b'I').finish();
3451    }
3452    for stmt_sql in statements {
3453        match execute_stmt(&stmt_sql, db_name, db, read_only) {
3454            // Abandon the rest of the batch on the first error, as Postgres does.
3455            Err(encoded) => {
3456                out.extend_from_slice(&encoded);
3457                return out;
3458            }
3459            Ok(ex) => {
3460                if ex.has_rows {
3461                    out.extend_from_slice(&encode_rows(&ex.rows, &ex.project));
3462                }
3463                out.extend_from_slice(&command_complete(&ex.tag_for(ex.rows.len())));
3464            }
3465        }
3466    }
3467    out
3468}
3469
3470/// Split on `;` at the top level, ignoring separators inside string literals.
3471fn split_statements(sql: &str) -> Vec<String> {
3472    let mut out = vec![];
3473    let mut cur = String::new();
3474    let mut in_s = false;
3475    for c in sql.chars() {
3476        match c {
3477            '\'' => { in_s = !in_s; cur.push(c); }
3478            ';' if !in_s => {
3479                if !cur.trim().is_empty() { out.push(cur.clone()); }
3480                cur.clear();
3481            }
3482            _ => cur.push(c),
3483        }
3484    }
3485    if !cur.trim().is_empty() {
3486        out.push(cur);
3487    }
3488    out
3489}
3490
3491/// Bind and serve the Postgres read endpoint until the process exits.
3492pub async fn run(host: &str, port: u16, resolver: Arc<dyn DbResolver>) -> anyhow::Result<()> {
3493    // Writes are ON by default — that is the parity position. An operator who
3494    // wants the "system of proof beside your database" deployment, where this
3495    // door must never mutate anything, sets NEDBD_PG_READ_ONLY=1.
3496    let read_only = std::env::var("NEDBD_PG_READ_ONLY")
3497        .map(|v| v == "1" || v.eq_ignore_ascii_case("true"))
3498        .unwrap_or(false);
3499    let listener = TcpListener::bind((host, port)).await?;
3500    println!("  pgwire   postgres endpoint on {}:{} — psql / DBeaver / psycopg ({})",
3501             host, port,
3502             if read_only { "SELECT only — read-only mode" } else { "SELECT + INSERT/UPDATE/DELETE" });
3503    loop {
3504        let (sock, _peer) = match listener.accept().await {
3505            Ok(v) => v,
3506            Err(e) => {
3507                eprintln!("  [pgwire] accept failed: {}", e);
3508                continue;
3509            }
3510        };
3511        let r = Arc::clone(&resolver);
3512        tokio::spawn(async move {
3513            let _ = sock.set_nodelay(true);
3514            if let Err(e) = handle(sock, r, read_only).await {
3515                // A client disconnecting mid-message is routine, not an incident.
3516                if e.kind() != std::io::ErrorKind::UnexpectedEof
3517                    && e.kind() != std::io::ErrorKind::ConnectionReset
3518                {
3519                    eprintln!("  [pgwire] connection error: {}", e);
3520                }
3521            }
3522        });
3523    }
3524}
3525
3526// ─────────────────────────────────────────────────────────────────────────────
3527
3528#[cfg(test)]
3529mod explain_tests {
3530    use super::*;
3531
3532    #[test]
3533    fn a_bare_explain_is_stripped() {
3534        assert_eq!(strip_explain("EXPLAIN SELECT 1"), Some("SELECT 1"));
3535        assert_eq!(strip_explain("explain select 1"), Some("select 1"));
3536        assert_eq!(strip_explain("  EXPLAIN   SELECT 1 ;  "), Some("SELECT 1"));
3537    }
3538
3539    #[test]
3540    fn analyze_and_verbose_are_accepted_and_ignored() {
3541        // This endpoint always executes and always reports actual rows, so
3542        // EXPLAIN and EXPLAIN ANALYZE genuinely do the same thing. Accepting
3543        // the client's spelling beats refusing it.
3544        assert_eq!(strip_explain("EXPLAIN ANALYZE SELECT 1"), Some("SELECT 1"));
3545        assert_eq!(strip_explain("EXPLAIN ANALYSE SELECT 1"), Some("SELECT 1"));
3546        assert_eq!(strip_explain("EXPLAIN VERBOSE SELECT 1"), Some("SELECT 1"));
3547        assert_eq!(strip_explain("EXPLAIN ANALYZE VERBOSE SELECT 1"), Some("SELECT 1"));
3548        assert_eq!(strip_explain("explain analyze verbose select 1"), Some("select 1"));
3549    }
3550
3551    #[test]
3552    fn a_word_merely_starting_with_explain_is_not_a_keyword() {
3553        assert_eq!(strip_explain("EXPLAINED SELECT 1"), None);
3554        assert_eq!(strip_explain("SELECT 1"), None);
3555        assert_eq!(strip_explain("SELECT explain FROM t"), None);
3556    }
3557
3558    #[test]
3559    fn a_column_named_analyze_is_not_eaten() {
3560        // `analyzed` merely starts with the keyword; the word boundary check
3561        // is what stops it being consumed as an option.
3562        assert_eq!(strip_explain("EXPLAIN analyzed_view"), Some("analyzed_view"));
3563    }
3564
3565    #[test]
3566    fn the_plan_result_has_postgres_shape() {
3567        let e = plan_result(vec!["Seq Scan on t".into(), "note".into()]);
3568        assert_eq!(e.project.len(), 1);
3569        assert_eq!(e.project[0].out, "QUERY PLAN");
3570        assert_eq!(e.rows.len(), 2);
3571        assert_eq!(e.rows[0]["QUERY PLAN"], "Seq Scan on t");
3572        assert_eq!(e.tag, "EXPLAIN");
3573        // EXPLAIN's tag carries no row count in PostgreSQL.
3574        assert!(!e.tag_counts_rows);
3575    }
3576}
3577
3578#[cfg(test)]
3579mod tests {
3580    use super::*;
3581    use serde_json::json;
3582
3583    fn q(sql: &str) -> String {
3584        match translate(sql) {
3585            Ok(Stmt::Query { nql, .. }) => nql,
3586            other => panic!("expected a query for {:?}, got {:?}", sql, other),
3587        }
3588    }
3589    /// Output column names, in order.
3590    fn proj(sql: &str) -> Vec<String> {
3591        match translate(sql) {
3592            Ok(Stmt::Query { project, .. }) => project.iter().map(|c| c.out.clone()).collect(),
3593            other => panic!("expected a query for {:?}, got {:?}", sql, other),
3594        }
3595    }
3596    /// (source key, output name) pairs, for the aggregate renaming.
3597    fn proj_pairs(sql: &str) -> Vec<(String, String)> {
3598        match translate(sql) {
3599            Ok(Stmt::Query { project, .. }) =>
3600                project.iter().map(|c| (c.src.clone(), c.out.clone())).collect(),
3601            other => panic!("expected a query for {:?}, got {:?}", sql, other),
3602        }
3603    }
3604    fn names(cols: &[Col]) -> Vec<String> { cols.iter().map(|c| c.out.clone()).collect() }
3605
3606    /// The full projection, so a test can assert the SRC and the OUT
3607    /// separately — they are different jobs and conflating them is how an
3608    /// alias got lost.
3609    fn cols_of(sql: &str) -> Vec<Col> {
3610        match translate(sql).unwrap() {
3611            Stmt::Query { project, .. } => project,
3612            other => panic!("{:?}", other),
3613        }
3614    }
3615
3616    #[test]
3617    fn select_star_becomes_bare_from() {
3618        assert_eq!(q("SELECT * FROM orders"), "FROM orders");
3619        assert_eq!(q("select * from orders;"), "FROM orders");
3620        assert_eq!(proj("SELECT * FROM orders"), Vec::<String>::new());
3621    }
3622
3623    #[test]
3624    fn a_column_list_becomes_a_projection_not_a_clause() {
3625        // NQL has no projection, so the column list is carried separately and
3626        // applied to the returned rows.
3627        assert_eq!(q("SELECT status, total FROM orders"), "FROM orders");
3628        assert_eq!(proj("SELECT status, total FROM orders"), vec!["status", "total"]);
3629    }
3630
3631    #[test]
3632    fn a_qualifier_reduces_to_the_field_while_an_ALIAS_is_the_name_the_client_sees() {
3633        // Two different jobs, and they used to be conflated. The SRC is what
3634        // NEDB reads out of the row, so a qualifier must be stripped from it.
3635        // The OUT is the name the CLIENT looks the column up by, so an alias
3636        // must be KEPT in it — `SELECT status AS s` returns a column called
3637        // `s`, and answering with one called `status` hands a client a result
3638        // it cannot find. SQLAlchemy writes `count(*) AS count_1` and then
3639        // reads `count_1`.
3640        let cols = cols_of("SELECT o.status AS s, o.total total, o.region FROM orders o");
3641        assert_eq!(cols.iter().map(|c| c.src.clone()).collect::<Vec<_>>(),
3642                   vec!["status", "total", "region"]);
3643        assert_eq!(cols.iter().map(|c| c.out.clone()).collect::<Vec<_>>(),
3644                   vec!["s", "total", "region"]);
3645        assert_eq!(q("SELECT * FROM public.orders"), "FROM orders");
3646        assert_eq!(q("SELECT * FROM \"orders\""), "FROM orders");
3647    }
3648
3649    #[test]
3650    fn a_select_list_may_MIX_columns_with_an_aggregate() {
3651        // What a GROUP BY query actually looks like. The previous parser
3652        // refused any list containing a parenthesis, so this whole shape was
3653        // unreachable even though NQL expresses it natively — and it is the
3654        // single most common grouped query an ORM emits.
3655        // The aggregate sits IMMEDIATELY AFTER the group key — verified
3656        // against the running engine, which refuses the other order with
3657        // "only one aggregate per query".
3658        assert_eq!(q("SELECT status, count(*) AS count_1 FROM orders GROUP BY status"),
3659                   "FROM orders GROUP BY status COUNT");
3660        // SQL puts GROUP BY before ORDER BY / LIMIT; the aggregate still lands
3661        // on the key, and the rest of the tail follows.
3662        assert_eq!(q("SELECT status, count(*) FROM orders WHERE total > 1 GROUP BY status ORDER BY status LIMIT 5"),
3663                   "FROM orders WHERE total > 1 GROUP BY status COUNT ORDER BY status LIMIT 5");
3664        // A bare aggregate with NO grouping still goes after the collection.
3665        assert_eq!(q("SELECT count(*) FROM orders"), "FROM orders COUNT");
3666        assert_eq!(q("SELECT sum(total) FROM orders"), "FROM orders SUM total");
3667        // More than one group key is refused by name: NQL groups by a single
3668        // field, and using only the first would aggregate over rows the query
3669        // meant to keep apart.
3670        let e = translate("SELECT status, count(*) FROM orders GROUP BY status, region").unwrap_err();
3671        assert!(e.contains("GROUP BY takes one key"), "{}", e);
3672        let cols = cols_of("SELECT status, count(*) AS count_1 FROM orders GROUP BY status");
3673        assert_eq!(cols.iter().map(|c| c.src.clone()).collect::<Vec<_>>(),
3674                   vec!["status", "count"]);
3675        assert_eq!(cols.iter().map(|c| c.out.clone()).collect::<Vec<_>>(),
3676                   vec!["status", "count_1"]);
3677
3678        // A named aggregate rides along with `count`, because an NQL grouped
3679        // row carries both.
3680        let cols = cols_of("SELECT status, count(*), sum(total) FROM orders GROUP BY status");
3681        assert_eq!(cols.iter().map(|c| c.src.clone()).collect::<Vec<_>>(),
3682                   vec!["status", "count", "sum_total"]);
3683        assert_eq!(q("SELECT status, count(*), sum(total) FROM orders GROUP BY status"),
3684                   "FROM orders GROUP BY status SUM total");
3685
3686        // A qualifier on the aggregate's column is stripped like any other.
3687        assert_eq!(q("SELECT o.status, sum(o.total) FROM orders o GROUP BY o.status"),
3688                   "FROM orders GROUP BY status SUM total");
3689
3690        // Two NAMED aggregates cannot both be carried, and that is refused by
3691        // name rather than silently dropping one.
3692        let e = translate("SELECT status, sum(total), avg(total) FROM orders GROUP BY status")
3693            .unwrap_err();
3694        assert!(e.contains("only one of SUM/AVG/MIN/MAX"), "{}", e);
3695
3696        // A column that is neither a key nor an aggregate is still refused.
3697        let e = translate("SELECT status, total, count(*) FROM orders GROUP BY status")
3698            .unwrap_err();
3699        assert!(e.contains("must appear in the GROUP BY clause"), "{}", e);
3700    }
3701
3702    #[test]
3703    fn ORDER_BY_an_ordinal_resolves_to_that_select_list_column() {
3704        // SQL lets a sort key be a POSITION, and clients write it constantly.
3705        // NQL has no ordinals — it read the `1` as a literal and refused with
3706        // "expected field name, got Num(1.0)". node-postgres sent
3707        // `GROUP BY status ORDER BY 1` in the harness's first run.
3708        assert_eq!(q("SELECT status, total FROM orders ORDER BY 1"),
3709                   "FROM orders ORDER BY status");
3710        assert_eq!(q("SELECT status, total FROM orders ORDER BY 2 DESC"),
3711                   "FROM orders ORDER BY total DESC");
3712        // Several keys, mixing ordinals with names, and a direction on each.
3713        assert_eq!(q("SELECT status, total FROM orders ORDER BY 2 DESC, 1"),
3714                   "FROM orders ORDER BY total DESC, status");
3715        assert_eq!(q("SELECT status, total FROM orders ORDER BY 1, total DESC"),
3716                   "FROM orders ORDER BY status, total DESC");
3717        // An ordinal survives the GROUP BY splice, and resolves to the group
3718        // key rather than to the literal 1 — which is the exact shape that
3719        // failed in CI.
3720        assert_eq!(q("SELECT status, count(*) AS n FROM orders GROUP BY status ORDER BY 1"),
3721                   "FROM orders GROUP BY status COUNT ORDER BY status");
3722        // An ordinal may name the AGGREGATE column too.
3723        assert_eq!(q("SELECT status, count(*) AS n FROM orders GROUP BY status ORDER BY 2 DESC"),
3724                   "FROM orders GROUP BY status COUNT ORDER BY count DESC");
3725        // The clause boundary is respected: a following LIMIT is not swallowed
3726        // into the sort list, and `LIMIT 1` is not mistaken for an ordinal.
3727        assert_eq!(q("SELECT status, total FROM orders ORDER BY 2 LIMIT 1"),
3728                   "FROM orders ORDER BY total LIMIT 1");
3729        // A `1` anywhere else stays a literal.
3730        assert_eq!(q("SELECT status FROM orders WHERE total > 1 ORDER BY 1"),
3731                   "FROM orders WHERE total > 1 ORDER BY status");
3732
3733        // Out of range, and `SELECT *` where there is no list to index, are
3734        // both refused with the reason — guessing a column would sort by
3735        // something the query never named.
3736        let e = translate("SELECT status FROM orders ORDER BY 4").unwrap_err();
3737        assert!(e.contains("out of range") && e.contains("1 column"), "{}", e);
3738        let e = translate("SELECT * FROM orders ORDER BY 1").unwrap_err();
3739        assert!(e.contains("no list to index"), "{}", e);
3740    }
3741
3742    #[test]
3743    fn count_of_a_subquery_flattens_only_when_the_two_counts_MUST_agree() {
3744        // `.count()` in every ORM wraps the whole query in a derived table.
3745        // Counting rows that ARE the inner query's rows is counting the inner
3746        // query, so this is an identity, not an approximation.
3747        assert_eq!(
3748            q("SELECT count(*) AS count_1 FROM (SELECT orders._id AS a, orders.status AS b \
3749               FROM orders WHERE orders.status = 'paid') AS anon_1"),
3750            // Verified against the running engine: with no GROUP BY the
3751            // aggregate may sit either side of WHERE and answers identically.
3752            r#"FROM orders COUNT WHERE status = "paid""#);
3753        // No predicate at all.
3754        assert_eq!(q("SELECT count(*) FROM (SELECT orders._id FROM orders) AS anon_1"),
3755                   "FROM orders COUNT");
3756        // ORDER BY cannot change a count, so it is dropped rather than refused.
3757        assert_eq!(q("SELECT count(*) FROM (SELECT _id FROM orders ORDER BY total DESC) AS a"),
3758                   "FROM orders COUNT");
3759        // The outer alias is the name the client reads the column back by.
3760        let cols = cols_of("SELECT count(*) AS count_1 FROM (SELECT _id FROM orders) AS a");
3761        assert_eq!(cols[0].src, "count");
3762        assert_eq!(cols[0].out, "count_1");
3763
3764        // Each guard is a construct that would make the two counts DIFFERENT
3765        // numbers, so each is refused rather than silently flattened.
3766        for sql in [
3767            // LIMIT / OFFSET cap the rows before they are counted
3768            "SELECT count(*) FROM (SELECT _id FROM orders LIMIT 1) AS a",
3769            "SELECT count(*) FROM (SELECT _id FROM orders OFFSET 1) AS a",
3770            // the inner rows ARE the groups
3771            "SELECT count(*) FROM (SELECT status FROM orders GROUP BY status) AS a",
3772            // an inner aggregate already reduced the rows to one
3773            "SELECT count(*) FROM (SELECT count(*) FROM orders) AS a",
3774            "SELECT count(*) FROM (SELECT sum(total) FROM orders) AS a",
3775            // the outer list would need the derived table's own columns
3776            "SELECT count(*), status FROM (SELECT status FROM orders) AS a",
3777            "SELECT status FROM (SELECT status FROM orders) AS a",
3778            // one level is the claim
3779            "SELECT count(*) FROM (SELECT x FROM (SELECT _id AS x FROM orders) AS b) AS a",
3780        ] {
3781            let e = translate(sql).unwrap_err();
3782            assert!(e.contains("subqueries in FROM"), "{} -> {}", sql, e);
3783        }
3784
3785        // DISTINCT and the set operators are caught EARLIER, by their own
3786        // rules, which scan the whole statement before the FROM list is even
3787        // read. Asserted separately so the test records which check owns each
3788        // refusal rather than implying one catch-all does.
3789        for (sql, needle) in [
3790            ("SELECT count(*) FROM (SELECT DISTINCT status FROM orders) AS a", "DISTINCT"),
3791            ("SELECT count(*) FROM (SELECT a FROM t UNION SELECT b FROM u) AS x", "UNION"),
3792        ] {
3793            let e = translate(sql).unwrap_err();
3794            assert!(e.contains(needle), "{} -> {}", sql, e);
3795        }
3796    }
3797
3798    #[test]
3799    fn a_QUALIFIED_column_in_WHERE_finds_its_field_instead_of_ZERO_ROWS() {
3800        // THE silent wrong answer. NQL looks a field up FLAT, so
3801        // `WHERE orders.status = 'paid'` asked for a field literally named
3802        // "orders.status", no document had one, and the query returned ZERO
3803        // ROWS with no error — an empty result that reads exactly like "you
3804        // have no paid orders". Every ORM qualifies its predicates, so every
3805        // filtered SQLAlchemy query answered empty and `.get(pk)` answered
3806        // None.
3807        assert_eq!(q("SELECT _id FROM orders WHERE orders.status = 'paid'"),
3808                   r#"FROM orders WHERE status = "paid""#);
3809        assert_eq!(q("SELECT _id FROM orders WHERE orders.total > 50"),
3810                   "FROM orders WHERE total > 50");
3811        // Every clause in the tail, not just WHERE.
3812        assert_eq!(q("SELECT _id FROM orders ORDER BY orders.total DESC LIMIT 2"),
3813                   "FROM orders ORDER BY total DESC LIMIT 2");
3814        assert_eq!(q("SELECT status, count(*) FROM orders GROUP BY orders.status"),
3815                   "FROM orders GROUP BY status COUNT");
3816
3817        // An alias is a legal qualifier and is accepted as one. It is also
3818        // REMOVED from the tail, because NQL has no alias syntax and reported
3819        // an "unexpected token" on it.
3820        assert_eq!(q("SELECT o.status FROM orders o WHERE o.status = 'paid'"),
3821                   r#"FROM orders WHERE status = "paid""#);
3822        assert_eq!(q("SELECT o.status FROM orders AS o WHERE o.total > 1"),
3823                   "FROM orders WHERE total > 1");
3824
3825        // A qualifier naming NEITHER the collection nor its alias is an
3826        // ERROR, not a strip. Stripping it would answer from the one relation
3827        // that IS present, which is a different wrong answer in the same
3828        // empty-looking clothes.
3829        let e = translate("SELECT _id FROM orders WHERE nosuch.status = 'paid'").unwrap_err();
3830        assert!(e.contains("no table or alias named \"nosuch\""), "{}", e);
3831        let e = translate("SELECT _id FROM orders o WHERE p.status = 'paid'").unwrap_err();
3832        assert!(e.contains("aliased \"o\""), "the message names the alias in scope: {}", e);
3833
3834        // A dot INSIDE a literal is data, not a qualifier.
3835        assert_eq!(q("SELECT _id FROM orders WHERE status = 'pa.id'"),
3836                   r#"FROM orders WHERE status = "pa.id""#);
3837        // ...and a decimal point is not one either.
3838        assert_eq!(q("SELECT _id FROM orders WHERE total > 1.5"),
3839                   "FROM orders WHERE total > 1.5");
3840
3841        // UPDATE and DELETE carry the same tail, and had the same bug.
3842        match translate("UPDATE orders o SET status = 'x' WHERE o.total > 5").unwrap() {
3843            Stmt::Update { coll, nql, .. } => {
3844                assert_eq!(coll, "orders", "the alias is not part of the collection name");
3845                assert_eq!(nql, "FROM orders WHERE total > 5");
3846            }
3847            other => panic!("{:?}", other),
3848        }
3849        match translate("DELETE FROM orders o WHERE o.status = 'paid'").unwrap() {
3850            Stmt::Delete { coll, nql, .. } => {
3851                assert_eq!(coll, "orders");
3852                assert_eq!(nql, r#"FROM orders WHERE status = "paid""#);
3853            }
3854            other => panic!("{:?}", other),
3855        }
3856
3857        // `AS OF SYSTEM TIME` also begins with AS and is NOT an alias.
3858        assert_eq!(q("SELECT _id FROM orders AS OF SYSTEM TIME 3 WHERE orders.total > 1"),
3859                   "FROM orders AS OF 3 WHERE total > 1");
3860    }
3861
3862    #[test]
3863    fn where_clauses_pass_through_with_sql_literals_rewritten() {
3864        assert_eq!(q("SELECT * FROM orders WHERE status = 'paid'"),
3865                   r#"FROM orders WHERE status = "paid""#);
3866        assert_eq!(q("SELECT * FROM orders WHERE status <> 'paid'"),
3867                   r#"FROM orders WHERE status != "paid""#);
3868        assert_eq!(q("SELECT * FROM orders WHERE status IN ('paid','open')"),
3869                   r#"FROM orders WHERE status IN ("paid","open")"#);
3870    }
3871
3872    /// SQL escapes an embedded quote by doubling it. That must become ONE
3873    /// character inside the NQL string, not terminate it.
3874    #[test]
3875    fn a_doubled_sql_quote_is_one_literal_character() {
3876        assert_eq!(q("SELECT * FROM t WHERE name = 'it''s'"),
3877                   r#"FROM t WHERE name = "it's""#);
3878    }
3879
3880    /// A double quote inside a SQL literal has to be escaped for NQL, whose
3881    /// lexer collapses \" — otherwise it would close the string early.
3882    #[test]
3883    fn a_double_quote_inside_a_sql_literal_is_escaped_for_nql() {
3884        assert_eq!(q(r#"SELECT * FROM t WHERE name = 'say "hi"'"#),
3885                   r#"FROM t WHERE name = "say \"hi\"""#);
3886    }
3887
3888    #[test]
3889    fn the_shared_clauses_are_handed_to_nql_unchanged() {
3890        assert_eq!(q("SELECT * FROM orders ORDER BY total DESC LIMIT 10 OFFSET 5"),
3891                   "FROM orders ORDER BY total DESC LIMIT 10 OFFSET 5");
3892        assert_eq!(q("SELECT * FROM orders GROUP BY region"), "FROM orders GROUP BY region");
3893        assert_eq!(q("SELECT * FROM o WHERE total BETWEEN 1 AND 9 ORDER BY a, b DESC"),
3894                   "FROM o WHERE total BETWEEN 1 AND 9 ORDER BY a, b DESC");
3895    }
3896
3897    /// An aggregate must surface as ONE column, named as SQL names it.
3898    ///
3899    /// NQL answers `SUM(total)` with `{count, sum_total, value}` — `value`
3900    /// being a back-compat alias. Passing that straight through gave
3901    /// `SELECT COUNT(*)` two columns (`count`, `value`) where SQL promises
3902    /// one, and leaked an internal key name onto the wire.
3903    #[test]
3904    fn an_aggregate_is_one_column_named_as_sql_names_it() {
3905        assert_eq!(proj_pairs("SELECT COUNT(*) FROM orders"),
3906                   vec![("count".to_string(), "count".to_string())]);
3907        assert_eq!(proj_pairs("SELECT SUM(total) FROM orders"),
3908                   vec![("sum_total".to_string(), "sum".to_string())]);
3909        assert_eq!(proj_pairs("SELECT avg(total) FROM orders"),
3910                   vec![("avg_total".to_string(), "avg".to_string())]);
3911        assert_eq!(proj_pairs("SELECT MIN(total) FROM orders"),
3912                   vec![("min_total".to_string(), "min".to_string())]);
3913        // And the encoded result really is one column with that name.
3914        let rows = vec![json!({"count": 4, "sum_total": 420, "value": 420})];
3915        let p = vec![Col::renamed("sum_total", "sum")];
3916        let cols = columns_for(&rows, &p);
3917        assert_eq!(names(&cols), vec!["sum"], "one column, SQL's name");
3918        assert_eq!(cell(rows[0].get(&cols[0].src)), Some("420".to_string()));
3919    }
3920
3921    /// A grouped NQL row holds the group key, `count` and the aggregate —
3922    /// nothing else. Projecting another column found nothing and rendered
3923    /// NULL, which is a silent wrong answer. Postgres errors; so do we, in
3924    /// Postgres's own words.
3925    #[test]
3926    fn a_bare_column_with_group_by_is_refused_not_nulled() {
3927        let e = translate("SELECT region, total FROM orders GROUP BY region").unwrap_err();
3928        assert!(e.contains("must appear in the GROUP BY clause"), "{}", e);
3929        assert!(e.contains("total"), "the message names the offending column: {}", e);
3930
3931        // The group key itself, and `count`, are both legitimate.
3932        assert!(translate("SELECT region FROM orders GROUP BY region").is_ok());
3933        assert!(translate("SELECT region, count FROM orders GROUP BY region").is_ok());
3934        // As is an aggregate over the grouped set.
3935        assert!(translate("SELECT SUM(total) FROM orders GROUP BY region").is_ok());
3936        // And `*` is unaffected — it returns whatever the grouped row holds.
3937        assert!(translate("SELECT * FROM orders GROUP BY region").is_ok());
3938    }
3939
3940    #[test]
3941    fn count_star_becomes_nql_count() {
3942        assert_eq!(q("SELECT COUNT(*) FROM orders"), "FROM orders COUNT");
3943        assert_eq!(q("SELECT count(*) FROM orders WHERE total > 5"),
3944                   "FROM orders COUNT WHERE total > 5");
3945    }
3946
3947    #[test]
3948    fn aggregates_carry_their_target_column() {
3949        assert_eq!(q("SELECT SUM(total) FROM orders"), "FROM orders SUM total");
3950        assert_eq!(q("SELECT avg(total) FROM orders WHERE region = 'eu'"),
3951                   r#"FROM orders AVG total WHERE region = "eu""#);
3952        assert!(translate("SELECT SUM(*) FROM orders").is_err());
3953    }
3954
3955    /// The bridge worth having: Postgres spells time travel
3956    /// `AS OF SYSTEM TIME`, and NEDB's is sequence-addressed and permanent.
3957    #[test]
3958    fn as_of_system_time_bridges_to_nql_as_of() {
3959        assert_eq!(q("SELECT * FROM orders AS OF SYSTEM TIME 42"),
3960                   "FROM orders AS OF 42");
3961        assert_eq!(q("SELECT * FROM orders AS OF SYSTEM TIME 42 WHERE total > 1"),
3962                   "FROM orders AS OF 42 WHERE total > 1");
3963        // A wall-clock timestamp is refused with the reason, not silently ignored.
3964        let e = translate("SELECT * FROM orders AS OF SYSTEM TIME '2026-01-01'").unwrap_err();
3965        assert!(e.contains("sequence number"), "{}", e);
3966    }
3967
3968    #[test]
3969    fn handshake_queries_are_answered_so_clients_can_connect() {
3970        assert!(matches!(translate("SELECT version()"), Ok(Stmt::Canned { .. })));
3971        assert!(matches!(translate("SHOW transaction_isolation"), Ok(Stmt::Canned { .. })));
3972        assert!(matches!(translate("SELECT current_schema()"), Ok(Stmt::Canned { .. })));
3973        assert!(matches!(translate("SET extra_float_digits = 3"), Ok(Stmt::Ok(_))));
3974        assert!(matches!(translate("BEGIN"), Ok(Stmt::Ok(_))));
3975        assert!(matches!(translate(""), Ok(Stmt::Ok(_))));
3976    }
3977
3978    /// Every refusal has to name the boundary. "Syntax error" would send a
3979    /// developer hunting for a typo that is not there.
3980    #[test]
3981    fn unsupported_sql_is_refused_with_a_reason() {
3982        for (sql, expect) in [
3983            ("INSERT INTO t VALUES (1)", "explicit column list"),
3984            ("CREATE TABLE t (a int)", "DDL"),
3985            ("TRUNCATE t", "append-only"),
3986            ("GRANT ALL ON t TO x", "privilege system"),
3987            ("SELECT * FROM a JOIN b ON a.x = b.x", "JOIN is not supported"),
3988            ("SELECT * FROM a UNION SELECT * FROM b", "UNION"),
3989            ("SELECT DISTINCT region FROM orders", "GROUP BY"),
3990            ("SELECT * FROM (SELECT 1) x", "subqueries in FROM"),
3991            ("SELECT * FROM a, b", "more than one collection"),
3992            ("SELECT lower(status) FROM orders", "expressions in the select list"),
3993            ("VACUUM", "only SELECT"),
3994        ] {
3995            let e = translate(sql).unwrap_err();
3996            assert!(e.contains(expect), "for {:?} expected {:?} in {:?}", sql, expect, e);
3997        }
3998    }
3999
4000    // ── writes ───────────────────────────────────────────────────────────────
4001    //
4002    // SQL's write semantics and NEDB's append-only model line up: INSERT is a
4003    // put, UPDATE is a new version, DELETE is a tombstone. These tests pin the
4004    // parse; tests/test_pgwire.py proves the behaviour against a live server,
4005    // including that the PRIOR value is still readable afterwards.
4006
4007    fn ins(sql: &str) -> (String, Vec<InsertRow>, Vec<Col>) {
4008        match translate(sql) {
4009            Ok(Stmt::Insert { coll, rows, returning }) => (coll, rows, returning),
4010            other => panic!("expected INSERT for {:?}, got {:?}", sql, other),
4011        }
4012    }
4013
4014    #[test]
4015    fn insert_becomes_a_put_per_row() {
4016        let (coll, rows, ret) = ins("INSERT INTO orders (_id, status, total) VALUES ('o1', 'paid', 120)");
4017        assert_eq!(coll, "orders");
4018        assert_eq!(rows.len(), 1);
4019        assert_eq!(rows[0].id.as_deref(), Some("o1"));
4020        assert_eq!(rows[0].doc.get("status"), Some(&json!("paid")));
4021        assert_eq!(rows[0].doc.get("total"), Some(&json!(120)));
4022        // `_id` is the key, not a payload field.
4023        assert!(!rows[0].doc.contains_key("_id"));
4024        assert!(ret.is_empty());
4025    }
4026
4027    #[test]
4028    fn a_multi_row_insert_yields_one_row_each() {
4029        let (_, rows, _) = ins(
4030            "INSERT INTO t (id, n) VALUES ('a', 1), ('b', 2), ('c', 3)");
4031        assert_eq!(rows.len(), 3);
4032        assert_eq!(rows[1].id.as_deref(), Some("b"));
4033        assert_eq!(rows[2].doc.get("n"), Some(&json!(3)));
4034    }
4035
4036    #[test]
4037    fn an_insert_without_an_id_column_lets_the_server_assign_one() {
4038        let (_, rows, _) = ins("INSERT INTO t (n) VALUES (1)");
4039        assert_eq!(rows[0].id, None, "the executor mints a unique key");
4040        assert_eq!(rows[0].doc.get("n"), Some(&json!(1)));
4041    }
4042
4043    /// Provenance is reachable from SQL, not only from the HTTP API — which is
4044    /// the point of having writes here at all.
4045    #[test]
4046    fn insert_lifts_provenance_out_of_reserved_columns() {
4047        let (_, rows, _) = ins(
4048            "INSERT INTO audit (_id, _caused_by, _valid_from, kind) \
4049             VALUES ('e1', 'abc123', '2026-01-01', 'reprice')");
4050        assert_eq!(rows[0].caused_by, vec!["abc123".to_string()]);
4051        assert_eq!(rows[0].valid_from.as_deref(), Some("2026-01-01"));
4052        assert_eq!(rows[0].doc.get("kind"), Some(&json!("reprice")));
4053        // None of the reserved names leak into the stored payload.
4054        for k in ["_id", "_caused_by", "_valid_from"] {
4055            assert!(!rows[0].doc.contains_key(k), "{} leaked into the doc", k);
4056        }
4057    }
4058
4059    #[test]
4060    fn insert_values_cover_the_scalar_types() {
4061        let (_, rows, _) = ins(
4062            "INSERT INTO t (s, i, f, b, n) VALUES ('x', 42, 1.5, TRUE, NULL)");
4063        assert_eq!(rows[0].doc.get("s"), Some(&json!("x")));
4064        assert_eq!(rows[0].doc.get("i"), Some(&json!(42)));
4065        assert_eq!(rows[0].doc.get("f"), Some(&json!(1.5)));
4066        assert_eq!(rows[0].doc.get("b"), Some(&json!(true)));
4067        assert_eq!(rows[0].doc.get("n"), Some(&Value::Null));
4068    }
4069
4070    /// A doubled '' is one literal quote, and a comma inside a string is not a
4071    /// value separator.
4072    #[test]
4073    fn insert_literals_survive_quotes_and_commas() {
4074        let (_, rows, _) = ins("INSERT INTO t (a, b) VALUES ('it''s', 'x,y')");
4075        assert_eq!(rows[0].doc.get("a"), Some(&json!("it's")));
4076        assert_eq!(rows[0].doc.get("b"), Some(&json!("x,y")));
4077    }
4078
4079    #[test]
4080    fn insert_refuses_what_it_cannot_store_faithfully() {
4081        // An unevaluated expression stored as text would be a wrong value.
4082        assert!(translate("INSERT INTO t (a) VALUES (1 + 1)").is_err());
4083        assert!(translate("INSERT INTO t (a) VALUES (now())").is_err());
4084        // Column/value count mismatch.
4085        let e = translate("INSERT INTO t (a, b) VALUES (1)").unwrap_err();
4086        assert!(e.contains("values for"), "{}", e);
4087        // No column list at all.
4088        let e2 = translate("INSERT INTO t VALUES (1)").unwrap_err();
4089        assert!(e2.contains("explicit column list"), "{}", e2);
4090    }
4091
4092    #[test]
4093    fn update_finds_rows_with_the_full_predicate_surface() {
4094        match translate("UPDATE orders SET status = 'void' WHERE total < 50 AND region IN ('eu')") {
4095            Ok(Stmt::Update { coll, set, nql, .. }) => {
4096                assert_eq!(coll, "orders");
4097                assert_eq!(set, vec![("status".to_string(), json!("void"))]);
4098                // The WHERE became ordinary NQL, so IN/BETWEEN/LIKE all work.
4099                assert_eq!(nql, r#"FROM orders WHERE total < 50 AND region IN ("eu")"#);
4100            }
4101            other => panic!("expected UPDATE, got {:?}", other),
4102        }
4103    }
4104
4105    #[test]
4106    fn update_without_where_targets_the_whole_collection() {
4107        // Postgres allows it, so parity allows it.
4108        match translate("UPDATE t SET a = 1") {
4109            Ok(Stmt::Update { nql, .. }) => assert_eq!(nql, "FROM t"),
4110            other => panic!("expected UPDATE, got {:?}", other),
4111        }
4112    }
4113
4114    #[test]
4115    fn update_handles_several_assignments() {
4116        match translate("UPDATE t SET a = 1, b = 'x,y', c = NULL WHERE id = 'k'") {
4117            Ok(Stmt::Update { set, .. }) => {
4118                assert_eq!(set.len(), 3);
4119                assert_eq!(set[1], ("b".to_string(), json!("x,y")));
4120                assert_eq!(set[2], ("c".to_string(), Value::Null));
4121            }
4122            other => panic!("expected UPDATE, got {:?}", other),
4123        }
4124        assert!(translate("UPDATE t SET").is_err());
4125        assert!(translate("UPDATE t SET a").is_err());
4126    }
4127
4128    #[test]
4129    fn delete_becomes_a_predicate_over_the_collection() {
4130        match translate("DELETE FROM orders WHERE status = 'void'") {
4131            Ok(Stmt::Delete { coll, nql, .. }) => {
4132                assert_eq!(coll, "orders");
4133                assert_eq!(nql, r#"FROM orders WHERE status = "void""#);
4134            }
4135            other => panic!("expected DELETE, got {:?}", other),
4136        }
4137        match translate("DELETE FROM t") {
4138            Ok(Stmt::Delete { nql, .. }) => assert_eq!(nql, "FROM t"),
4139            other => panic!("expected DELETE, got {:?}", other),
4140        }
4141    }
4142
4143    #[test]
4144    fn returning_is_parsed_off_every_write() {
4145        let (_, _, ret) = ins("INSERT INTO t (a) VALUES (1) RETURNING a, _id");
4146        assert_eq!(ret.iter().map(|c| c.out.clone()).collect::<Vec<_>>(), vec!["a", "_id"]);
4147        // `RETURNING *` is an empty projection — every column — which is why
4148        // the executor checks the raw SQL for the keyword instead.
4149        let (_, _, star) = ins("INSERT INTO t (a) VALUES (1) RETURNING *");
4150        assert!(star.is_empty());
4151        assert!(wants_returning("INSERT INTO t (a) VALUES (1) RETURNING *"));
4152        assert!(!wants_returning("INSERT INTO t (a) VALUES (1)"));
4153
4154        match translate("UPDATE t SET a = 1 WHERE id = 'k' RETURNING a") {
4155            Ok(Stmt::Update { nql, returning, .. }) => {
4156                assert_eq!(returning.len(), 1);
4157                // RETURNING must NOT leak into the predicate.
4158                assert!(!nql.to_uppercase().contains("RETURNING"), "{}", nql);
4159            }
4160            other => panic!("expected UPDATE, got {:?}", other),
4161        }
4162        match translate("DELETE FROM t WHERE id = 'k' RETURNING *") {
4163            Ok(Stmt::Delete { nql, .. }) =>
4164                assert!(!nql.to_uppercase().contains("RETURNING"), "{}", nql),
4165            other => panic!("expected DELETE, got {:?}", other),
4166        }
4167    }
4168
4169    #[test]
4170    fn a_keyword_inside_a_value_is_not_a_clause() {
4171        match translate("UPDATE t SET note = 'where returning from' WHERE id = 'k'") {
4172            Ok(Stmt::Update { set, nql, .. }) => {
4173                assert_eq!(set[0].1, json!("where returning from"));
4174                assert_eq!(nql, r#"FROM t WHERE id = "k""#);
4175            }
4176            other => panic!("expected UPDATE, got {:?}", other),
4177        }
4178    }
4179
4180    #[test]
4181    fn split_top_respects_quotes_and_nesting() {
4182        assert_eq!(split_top("a, b, c", ',').len(), 3);
4183        assert_eq!(split_top("(1, 2), (3, 4)", ',').len(), 2);
4184        assert_eq!(split_top("'a,b', c", ',').len(), 2);
4185        assert_eq!(split_top("'it''s, fine', c", ',').len(), 2);
4186    }
4187
4188    #[test]
4189    fn comments_and_whitespace_do_not_confuse_the_translator() {
4190        assert_eq!(q("SELECT *\n  FROM orders  -- trailing note\n"), "FROM orders");
4191        assert_eq!(q("SELECT /* inline */ * FROM orders"), "FROM orders");
4192        // A keyword inside a string literal must not be treated as a clause.
4193        assert_eq!(q("SELECT * FROM t WHERE note = 'from here to JOIN'"),
4194                   r#"FROM t WHERE note = "from here to JOIN""#);
4195    }
4196
4197    #[test]
4198    fn find_kw_ignores_quotes_parens_and_substrings() {
4199        assert_eq!(find_kw("SELECT A FROM B", "FROM"), Some(9));
4200        assert_eq!(find_kw("SELECT 'FROM' FROM B", "FROM"), Some(14));
4201        assert_eq!(find_kw("SELECT F(x FROM y) FROM B", "FROM"), Some(19));
4202        assert_eq!(find_kw("SELECT FROMAGE", "FROM"), None);
4203        assert_eq!(find_kw("SELECT X_FROM", "FROM"), None);
4204    }
4205
4206    // ── result encoding ──────────────────────────────────────────────────────
4207
4208    #[test]
4209    fn provenance_columns_sort_after_the_users_own_fields() {
4210        let rows = vec![json!({"_id":"1","_hash":"ab","status":"paid","total":9})];
4211        assert_eq!(names(&columns_for(&rows, &[])),
4212                   vec!["status", "total", "_hash", "_id"]);
4213    }
4214
4215    #[test]
4216    fn an_explicit_projection_sets_the_column_order() {
4217        let rows = vec![json!({"a":1,"b":2})];
4218        let p = vec![Col::same("b"), Col::same("a")];
4219        assert_eq!(names(&columns_for(&rows, &p)), vec!["b", "a"]);
4220    }
4221
4222    #[test]
4223    fn columns_are_the_union_across_sparse_rows() {
4224        // A document store has no schema, so row 2 may carry a field row 1 lacks.
4225        let rows = vec![json!({"a":1}), json!({"b":2})];
4226        assert_eq!(names(&columns_for(&rows, &[])), vec!["a", "b"]);
4227    }
4228
4229    #[test]
4230    fn type_oids_follow_the_first_non_null_value() {
4231        let rows = vec![json!({"i":1,"f":1.5,"b":true,"s":"x","n":null})];
4232        assert_eq!(oid_for(&rows, "i"), OID_INT8);
4233        assert_eq!(oid_for(&rows, "f"), OID_FLOAT8);
4234        assert_eq!(oid_for(&rows, "b"), OID_BOOL);
4235        assert_eq!(oid_for(&rows, "s"), OID_TEXT);
4236        // All-null and absent columns fall back to text rather than guessing.
4237        assert_eq!(oid_for(&rows, "n"), OID_TEXT);
4238        assert_eq!(oid_for(&rows, "absent"), OID_TEXT);
4239    }
4240
4241    #[test]
4242    fn a_column_that_is_null_in_the_first_row_still_gets_its_type() {
4243        let rows = vec![json!({"v": null}), json!({"v": 7})];
4244        assert_eq!(oid_for(&rows, "v"), OID_INT8);
4245    }
4246
4247    #[test]
4248    fn cells_render_in_postgres_text_format() {
4249        assert_eq!(cell(Some(&json!("x"))), Some("x".to_string()));
4250        assert_eq!(cell(Some(&json!(true))), Some("t".to_string()));
4251        assert_eq!(cell(Some(&json!(false))), Some("f".to_string()));
4252        assert_eq!(cell(Some(&json!(42))), Some("42".to_string()));
4253        assert_eq!(cell(Some(&json!(null))), None);
4254        assert_eq!(cell(None), None);
4255        // Nested values render as JSON text rather than being dropped.
4256        assert_eq!(cell(Some(&json!({"a":1}))), Some("{\"a\":1}".to_string()));
4257    }
4258
4259    /// The framing has to be exact or the client desynchronises and hangs.
4260    /// Length covers the length field itself but not the tag byte.
4261    #[test]
4262    fn message_framing_length_excludes_the_tag() {
4263        let mut m = Out::msg(b'Z');
4264        m.bytes(b"I");
4265        let bytes = m.finish();
4266        assert_eq!(bytes[0], b'Z');
4267        assert_eq!(i32::from_be_bytes([bytes[1], bytes[2], bytes[3], bytes[4]]), 5);
4268        assert_eq!(bytes.len(), 6);
4269    }
4270
4271    #[test]
4272    fn a_result_set_encodes_as_description_then_rows_then_complete() {
4273        let rows = vec![json!({"a": 1}), json!({"a": 2})];
4274        let out = encode_result(&rows, &[]);
4275        assert_eq!(out[0], b'T');
4276        let tags: Vec<u8> = {
4277            // Walk the message stream by its own length prefixes.
4278            let mut t = vec![];
4279            let mut i = 0usize;
4280            while i < out.len() {
4281                t.push(out[i]);
4282                let len = i32::from_be_bytes([out[i+1], out[i+2], out[i+3], out[i+4]]) as usize;
4283                i += 1 + len;
4284            }
4285            t
4286        };
4287        assert_eq!(tags, vec![b'T', b'D', b'D', b'C'],
4288                   "one description, one row each, one completion");
4289    }
4290
4291    /// A statement must emit EXACTLY ONE CommandComplete. A write with
4292    /// RETURNING that reused the SELECT encoder sent two, and the visible
4293    /// symptom was RETURNING yielding no rows: the client took the first tag
4294    /// as the end of the statement and threw the description away.
4295    #[test]
4296    fn a_write_with_returning_emits_exactly_one_command_complete() {
4297        let rows = vec![json!({"_id": "o1", "total": 9})];
4298        let mut out = encode_rows(&rows, &[Col::same("_id")]);
4299        out.extend_from_slice(&command_complete("INSERT 0 1"));
4300        let mut tags = vec![];
4301        let mut i = 0usize;
4302        while i < out.len() {
4303            tags.push(out[i]);
4304            let len = i32::from_be_bytes([out[i+1], out[i+2], out[i+3], out[i+4]]) as usize;
4305            i += 1 + len;
4306        }
4307        assert_eq!(tags, vec![b'T', b'D', b'C'], "one description, one row, ONE tag");
4308        assert_eq!(tags.iter().filter(|t| **t == b'C').count(), 1);
4309        // encode_rows alone must not carry a tag at all.
4310        assert!(!encode_rows(&rows, &[]).contains(&b'C')
4311                || encode_rows(&rows, &[]).iter().filter(|b| **b == b'C').count() > 0);
4312        let bare = encode_rows(&rows, &[Col::same("_id")]);
4313        let mut bare_tags = vec![];
4314        let mut j = 0usize;
4315        while j < bare.len() {
4316            bare_tags.push(bare[j]);
4317            let len = i32::from_be_bytes([bare[j+1], bare[j+2], bare[j+3], bare[j+4]]) as usize;
4318            j += 1 + len;
4319        }
4320        assert_eq!(bare_tags, vec![b'T', b'D'], "encode_rows never appends a tag");
4321    }
4322
4323    #[test]
4324    fn an_empty_result_still_sends_a_description() {
4325        let out = encode_result(&[], &[Col::same("a")]);
4326        assert_eq!(out[0], b'T', "clients need the shape even with no rows");
4327    }
4328
4329    #[test]
4330    fn statements_split_on_top_level_semicolons_only() {
4331        assert_eq!(split_statements("SELECT 1; SELECT 2").len(), 2);
4332        assert_eq!(split_statements("SELECT ';'").len(), 1);
4333        assert_eq!(split_statements("SELECT 1;").len(), 1);
4334        assert_eq!(split_statements("   ").len(), 0);
4335    }
4336
4337    #[test]
4338    fn an_error_names_its_sqlstate() {
4339        let e = String::from_utf8_lossy(&err_msg("0A000", "x")).to_string();
4340        assert!(e.contains("ERROR"));
4341        assert!(e.contains("0A000"));
4342    }
4343
4344    // ── the extended query protocol ─────────────────────────────────────────
4345
4346    #[test]
4347    fn placeholders_are_counted_outside_string_literals() {
4348        assert_eq!(param_count("SELECT a FROM t WHERE b = $1 AND c = $2"), 2);
4349        assert_eq!(param_count("SELECT a FROM t"), 0);
4350        // The highest index wins, because a parameter may be reused.
4351        assert_eq!(param_count("WHERE a = $2 OR b = $2 OR c = $1"), 2);
4352        assert_eq!(param_count("SELECT a FROM t WHERE b = '$1'"), 0,
4353                   "a placeholder inside a literal is data, not a parameter");
4354        assert_eq!(param_count("WHERE a = $10 AND b = $1"), 10,
4355                   "two-digit indexes must not be read as $1 followed by 0");
4356    }
4357
4358    #[test]
4359    fn parameters_are_spliced_as_literals() {
4360        let out = substitute_params("WHERE a = $1 AND b = $2 AND c = $3",
4361            &[Some("'x'".into()), Some("42".into()), None]).unwrap();
4362        assert_eq!(out, "WHERE a = 'x' AND b = 42 AND c = NULL");
4363    }
4364
4365    #[test]
4366    fn substitution_leaves_string_literals_alone() {
4367        let out = substitute_params("WHERE a = '$1' AND b = $1", &[Some("9".into())]).unwrap();
4368        assert_eq!(out, "WHERE a = '$1' AND b = 9");
4369    }
4370
4371    #[test]
4372    fn too_few_parameters_is_an_error_not_a_silent_null() {
4373        // The alternative — treating a missing parameter as NULL — turns a
4374        // client bug into a wrong answer with a 200-shaped response.
4375        let e = substitute_params("WHERE a = $2", &[Some("1".into())]).unwrap_err();
4376        assert!(e.contains("$2"), "{}", e);
4377    }
4378
4379    #[test]
4380    fn a_quote_in_a_parameter_cannot_escape_its_literal() {
4381        let lit = decode_param(Some(b"it's"), OID_TEXT, 0).unwrap().unwrap();
4382        assert_eq!(lit, "'it''s'");
4383        // And it survives a round trip through the splice unchanged.
4384        let out = substitute_params("WHERE a = $1", &[Some(lit)]).unwrap();
4385        assert_eq!(out, "WHERE a = 'it''s'");
4386    }
4387
4388    #[test]
4389    fn binary_parameters_decode_in_every_width_psycopg_sends() {
4390        // These are the exact encodings read off a psycopg3 wire transcript:
4391        // a small int arrives as int2, a float as float8, a bool as one byte.
4392        assert_eq!(decode_param(Some(&[0x00, 0x2a]), OID_INT2, 1).unwrap().unwrap(), "42");
4393        assert_eq!(decode_param(Some(&[0, 0, 0, 7]), OID_INT4, 1).unwrap().unwrap(), "7");
4394        assert_eq!(
4395            decode_param(Some(&[0, 0, 0, 0, 0, 0, 0, 9]), OID_INT8, 1).unwrap().unwrap(), "9");
4396        assert_eq!(
4397            decode_param(Some(&0x400c_0000_0000_0000u64.to_be_bytes()), OID_FLOAT8, 1)
4398                .unwrap().unwrap(), "3.5");
4399        assert_eq!(decode_param(Some(&[1]), OID_BOOL, 1).unwrap().unwrap(), "TRUE");
4400        assert_eq!(decode_param(Some(&[0]), OID_BOOL, 1).unwrap().unwrap(), "FALSE");
4401    }
4402
4403    #[test]
4404    fn a_negative_binary_integer_keeps_its_sign() {
4405        assert_eq!(decode_param(Some(&(-5i32).to_be_bytes()), OID_INT4, 1).unwrap().unwrap(), "-5");
4406        assert_eq!(decode_param(Some(&(-5i16).to_be_bytes()), OID_INT2, 1).unwrap().unwrap(), "-5");
4407    }
4408
4409    #[test]
4410    fn a_binary_parameter_of_the_wrong_width_is_refused() {
4411        // Truncating or zero-extending would produce a plausible wrong number,
4412        // which is the failure mode worth engineering against.
4413        let e = decode_param(Some(&[0x2a]), OID_INT4, 1).unwrap_err();
4414        assert!(e.contains("4 bytes"), "{}", e);
4415    }
4416
4417    #[test]
4418    fn an_unspecified_text_parameter_is_treated_as_a_string() {
4419        // psycopg3 declares OID 0 only for `str`; every number it sends carries
4420        // a real numeric OID. So quoting here is grounded, not a guess.
4421        assert_eq!(decode_param(Some(b"hello"), 0, 0).unwrap().unwrap(), "'hello'");
4422    }
4423
4424    #[test]
4425    fn a_null_parameter_decodes_to_none_in_every_format() {
4426        assert_eq!(decode_param(None, OID_TEXT, 0).unwrap(), None);
4427        assert_eq!(decode_param(None, OID_INT8, 1).unwrap(), None);
4428    }
4429
4430    #[test]
4431    fn an_unsupported_binary_type_says_so_by_name() {
4432        let e = decode_param(Some(&[0u8; 8]), 1114, 1).unwrap_err();
4433        assert!(e.contains("1114"), "{}", e);
4434        assert!(e.contains("text"), "the error should point at the way out: {}", e);
4435    }
4436
4437    #[test]
4438    fn a_text_number_that_is_not_a_number_gets_quoted() {
4439        // Splicing it in bare would emit a naked identifier into the NQL text
4440        // and fail somewhere far away from the cause.
4441        assert_eq!(decode_param(Some(b"oops"), OID_INT8, 0).unwrap().unwrap(), "'oops'");
4442    }
4443
4444    #[test]
4445    fn a_client_declared_type_is_believed_over_inference() {
4446        // The client is about to encode its argument that way; overriding it
4447        // would break the decode.
4448        let oids = infer_param_oids("SELECT a FROM t WHERE b = $1 AND c = $2", &[OID_INT4, 0], None);
4449        assert_eq!(oids, vec![OID_INT4, OID_TEXT]);
4450    }
4451
4452    #[test]
4453    fn parameter_arity_is_taken_from_the_sql_when_the_client_declares_none() {
4454        // asyncpg declares nothing and then refuses the call if the count that
4455        // comes back is wrong, so this is the load-bearing path for it.
4456        let oids = infer_param_oids("SELECT a FROM t WHERE b = $1 AND c = $2", &[], None);
4457        assert_eq!(oids.len(), 2);
4458    }
4459
4460    #[test]
4461    fn the_field_behind_each_placeholder_is_identified() {
4462        assert_eq!(
4463            param_fields("SELECT a FROM t WHERE qty > $1 AND status = $2", 2),
4464            vec![Some("qty".to_string()), Some("status".to_string())]);
4465    }
4466
4467    #[test]
4468    fn word_operators_do_not_hide_the_field() {
4469        assert_eq!(param_fields("SELECT a FROM t WHERE name LIKE $1", 1),
4470                   vec![Some("name".to_string())]);
4471        assert_eq!(param_fields("SELECT a FROM t WHERE qty BETWEEN $1 AND $2", 2),
4472                   vec![Some("qty".to_string()), Some("qty".to_string())]);
4473        assert_eq!(param_fields("SELECT a FROM t WHERE region IN ($1, $2)", 2),
4474                   vec![Some("region".to_string()), Some("region".to_string())]);
4475    }
4476
4477    #[test]
4478    fn a_clause_position_types_from_the_grammar_not_from_a_column() {
4479        // `AS OF SYSTEM TIME $1` has no column beside it — the token to its
4480        // left is the word TIME. Typing it text made asyncpg refuse to send
4481        // the sequence number at all.
4482        assert_eq!(
4483            infer_param_oids("SELECT a FROM t AS OF SYSTEM TIME $1 WHERE b = $2", &[], None),
4484            vec![OID_INT8, OID_TEXT]);
4485        assert_eq!(infer_param_oids("SELECT a FROM t AS OF $1", &[], None), vec![OID_INT8]);
4486        // VALID AS OF also ends with "AS OF", but its argument is a DATE
4487        // STRING. Checking the longer clause first is load-bearing.
4488        assert_eq!(
4489            infer_param_oids("SELECT a FROM t VALID AS OF $1", &[], None), vec![OID_TEXT]);
4490        assert_eq!(
4491            infer_param_oids("SELECT a FROM t LIMIT $1 OFFSET $2", &[], None),
4492            vec![OID_INT8, OID_INT8]);
4493    }
4494
4495    #[test]
4496    fn an_aggregate_column_types_from_what_the_aggregate_means() {
4497        // No document holds a field called `count`, so sampling stored data
4498        // finds nothing and falls back to text — which hands a binary client
4499        // the string "2" for COUNT(*).
4500        assert_eq!(aggregate_oid("count", None, "t"), Some(OID_INT8));
4501        assert_eq!(aggregate_oid("avg_fee", None, "t"), Some(OID_FLOAT8),
4502                   "an average is fractional even over integers");
4503        // SUM/MIN/MAX inherit the field's type; with no database to sample,
4504        // that resolves to text, and `_seq` is known from the engine contract.
4505        assert_eq!(aggregate_oid("max__seq", None, "t"), Some(OID_INT8));
4506        assert_eq!(aggregate_oid("total", None, "t"), None, "not an aggregate");
4507    }
4508
4509    #[test]
4510    fn the_parse_probe_uses_a_literal_that_every_clause_accepts() {
4511        // Stubbing with NULL was the obvious choice and the wrong one: clauses
4512        // that validate their argument rejected it, so `AS OF SYSTEM TIME $1`
4513        // failed at Parse before a real sequence was ever bound.
4514        let probe = probe_sql("SELECT a FROM t AS OF SYSTEM TIME $1 WHERE b = $2", 2);
4515        assert!(!probe.contains("NULL"), "{}", probe);
4516        assert!(translate(&probe).is_ok(), "the probe must parse: {}", probe);
4517    }
4518
4519    #[test]
4520    fn a_column_with_mixed_types_across_documents_is_advertised_as_text() {
4521        // Taking the first non-null value's type told the client `int8` and
4522        // then sent it "n/a" — which fails to parse client-side, and on the
4523        // binary path cannot be encoded at all.
4524        let rows = vec![json!({"x": 3}), json!({"x": "n/a"})];
4525        assert_eq!(oid_for(&rows, "x"), OID_TEXT);
4526        // Integers and floats in one column widen rather than conflict.
4527        let rows = vec![json!({"x": 3}), json!({"x": 1.5})];
4528        assert_eq!(oid_for(&rows, "x"), OID_FLOAT8);
4529        // A leading null must not decide the type.
4530        let rows = vec![json!({"x": Value::Null}), json!({"x": 7})];
4531        assert_eq!(oid_for(&rows, "x"), OID_INT8);
4532    }
4533
4534    #[test]
4535    fn binary_output_encodes_each_advertised_type() {
4536        assert_eq!(cell_binary(Some(&json!(true)), OID_BOOL).unwrap().unwrap(), vec![1]);
4537        assert_eq!(cell_binary(Some(&json!(42)), OID_INT8).unwrap().unwrap(),
4538                   42i64.to_be_bytes().to_vec());
4539        assert_eq!(cell_binary(Some(&json!(3.5)), OID_FLOAT8).unwrap().unwrap(),
4540                   3.5f64.to_be_bytes().to_vec());
4541        // For the text family, binary and text are the same bytes.
4542        assert_eq!(cell_binary(Some(&json!("hi")), OID_TEXT).unwrap().unwrap(), b"hi".to_vec());
4543        assert_eq!(cell_binary(Some(&Value::Null), OID_INT8).unwrap(), None);
4544        // A boolean renders as `t`/`f` in text but one byte in binary.
4545        assert_eq!(cell(Some(&json!(true))).unwrap(), "t");
4546    }
4547
4548    #[test]
4549    fn a_value_that_does_not_fit_its_advertised_binary_type_is_refused() {
4550        // Advertised types come from a bounded sample, so a field that only
4551        // turns heterogeneous outside it lands here. Sending a zero, or the
4552        // text bytes under a binary header, would corrupt the value in a way
4553        // the client cannot detect — so it is an error instead.
4554        let e = cell_binary(Some(&json!("nope")), OID_INT8).unwrap_err();
4555        assert!(e.contains("a string"), "{}", e);
4556        assert!(e.contains("more than one type"), "the error should explain WHY: {}", e);
4557    }
4558
4559    #[test]
4560    fn a_row_description_carries_the_requested_format_per_column() {
4561        let cols = [Col::same("a"), Col::same("b")];
4562        let m = row_description_fmt(&cols, &[OID_INT8, OID_TEXT], &[1, 0]);
4563        assert_eq!(m[0], b'T');
4564        // The trailing i16 of each field entry is its format code.
4565        assert_eq!(m[m.len() - 1], 0, "the last column was requested as text");
4566    }
4567
4568    #[test]
4569    fn a_qualified_column_resolves_to_its_bare_name() {
4570        assert_eq!(param_fields("SELECT a FROM t WHERE t.qty = $1", 1),
4571                   vec![Some("qty".to_string())]);
4572    }
4573
4574    #[test]
4575    fn insert_placeholders_map_positionally_to_the_column_list() {
4576        assert_eq!(
4577            param_fields("INSERT INTO t (_id, qty, status) VALUES ($1, $2, $3)", 3),
4578            vec![Some("_id".to_string()), Some("qty".to_string()), Some("status".to_string())]);
4579    }
4580
4581    #[test]
4582    fn a_set_clause_placeholder_finds_its_column() {
4583        assert_eq!(param_fields("UPDATE t SET status = $1 WHERE _id = $2", 2),
4584                   vec![Some("status".to_string()), Some("_id".to_string())]);
4585    }
4586
4587    #[test]
4588    fn the_target_collection_is_found_for_every_statement_kind() {
4589        assert_eq!(stmt_collection("SELECT a FROM inv WHERE b = $1"), "inv");
4590        assert_eq!(stmt_collection("UPDATE inv SET a = $1"), "inv");
4591        assert_eq!(stmt_collection("DELETE FROM inv WHERE a = $1"), "inv");
4592        assert_eq!(stmt_collection("INSERT INTO inv (a) VALUES ($1)"), "inv");
4593        // Clients qualify as schema.table; NEDB has one namespace.
4594        assert_eq!(stmt_collection("SELECT a FROM public.inv"), "inv");
4595        assert_eq!(stmt_collection("INSERT INTO inv(a) VALUES ($1)"), "inv");
4596    }
4597
4598    #[test]
4599    fn engine_metadata_fields_type_without_touching_storage() {
4600        assert_eq!(infer_field_oid(None, "t", "_seq"), OID_INT8);
4601        assert_eq!(infer_field_oid(None, "t", "_id"), OID_TEXT);
4602    }
4603
4604    #[test]
4605    fn the_protocol_acknowledgements_are_single_empty_messages() {
4606        // Each is a tag plus a 4-byte length of exactly 4.
4607        for (m, tag) in [
4608            (parse_complete(), b'1'), (bind_complete(), b'2'),
4609            (close_complete(), b'3'), (no_data(), b'n'), (portal_suspended(), b's'),
4610        ] {
4611            assert_eq!(m.len(), 5, "{:?}", tag as char);
4612            assert_eq!(m[0], tag);
4613            assert_eq!(i32::from_be_bytes([m[1], m[2], m[3], m[4]]), 4);
4614        }
4615    }
4616
4617    #[test]
4618    fn parameter_description_reports_its_arity_and_types() {
4619        let m = parameter_description(&[OID_TEXT, OID_INT8]);
4620        assert_eq!(m[0], b't');
4621        assert_eq!(i16::from_be_bytes([m[5], m[6]]), 2);
4622        assert_eq!(i32::from_be_bytes([m[7], m[8], m[9], m[10]]), OID_TEXT);
4623        assert_eq!(i32::from_be_bytes([m[11], m[12], m[13], m[14]]), OID_INT8);
4624    }
4625
4626    #[test]
4627    fn a_cstring_is_taken_without_its_terminator() {
4628        let body = b"one\0two\0".to_vec();
4629        let mut at = 0usize;
4630        assert_eq!(take_cstr(&body, &mut at), "one");
4631        assert_eq!(take_cstr(&body, &mut at), "two");
4632        assert_eq!(at, body.len());
4633    }
4634
4635    #[test]
4636    fn truncated_integers_are_reported_rather_than_read_past_the_end() {
4637        let body = vec![0u8, 1];
4638        let mut at = 0usize;
4639        assert!(take_i32(&body, &mut at).is_err());
4640        let mut at = 0usize;
4641        assert!(take_i16(&body, &mut at).is_ok());
4642    }
4643
4644    #[test]
4645    fn a_binary_result_format_request_is_refused_rather_than_faked() {
4646        // Sending text under a binary header corrupts every value silently,
4647        // which is far worse than an error naming the limitation.
4648        let out = encode_rows(&[], &[Col::same("a")]);
4649        let desc_format = &out[out.len() - 2..];
4650        assert_eq!(i16::from_be_bytes([desc_format[0], desc_format[1]]), 0,
4651                   "every column is advertised as text format");
4652    }
4653
4654    #[test]
4655    fn a_float_parameter_does_not_render_as_rust_infinity() {
4656        assert_eq!(fmt_float(f64::INFINITY), "'Infinity'");
4657        assert_eq!(fmt_float(f64::NEG_INFINITY), "'-Infinity'");
4658        assert_eq!(fmt_float(f64::NAN), "'NaN'");
4659        assert_eq!(fmt_float(3.0), "3", "a whole float should not gain a .0 tail");
4660        assert_eq!(fmt_float(3.5), "3.5");
4661    }
4662}