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 {
277 coll: String,
278 set: Vec<(String, Value)>,
279 /// The SQL `WHERE …` as written (column qualifiers stripped), which is
280 /// what actually selects the rows. See `rows_for_write`.
281 where_sql: String,
282 /// The same predicate rendered as NQL. No longer used to SELECT
283 /// anything — kept because it is the translation the `translate_*`
284 /// tests pin, and because an operator reading a 42601 wants to see it.
285 nql: String,
286 returning: Vec<Col>,
287 },
288 /// `DELETE FROM coll [WHERE …] [RETURNING …]` — a tombstone per match.
289 Delete { coll: String, where_sql: String, nql: String, returning: Vec<Col> },
290 /// Answer from a fixed table — the handshake queries clients send on connect.
291 Canned { cols: Vec<String>, row: Vec<String> },
292 /// Nothing to do (empty statement, or a SET the client does not need honoured).
293 Ok(&'static str),
294}
295
296/// One row of an `INSERT`: an explicit id when the statement supplied one, the
297/// document body, and optional provenance lifted out of reserved columns.
298#[derive(Debug, PartialEq, Clone)]
299pub struct InsertRow {
300 /// From an `_id` or `id` column. `None` means the server assigns one.
301 pub id: Option<String>,
302 pub doc: serde_json::Map<String, Value>,
303 /// From a `_caused_by` column — the causal parents, so provenance is
304 /// reachable from SQL rather than only from the HTTP API.
305 pub caused_by: Vec<String>,
306 pub valid_from: Option<String>,
307 pub valid_to: Option<String>,
308}
309
310/// Strip SQL comments and collapse whitespace, so the matchers below can be
311/// simple without being fragile about formatting.
312fn normalise(sql: &str) -> String {
313 let mut out = String::with_capacity(sql.len());
314 let mut chars = sql.chars().peekable();
315 let mut in_s = false;
316 while let Some(c) = chars.next() {
317 if in_s {
318 out.push(c);
319 if c == '\'' { in_s = false; }
320 continue;
321 }
322 match c {
323 '\'' => { in_s = true; out.push(c); }
324 '-' if chars.peek() == Some(&'-') => {
325 // line comment
326 for n in chars.by_ref() { if n == '\n' { break; } }
327 out.push(' ');
328 }
329 '/' if chars.peek() == Some(&'*') => {
330 chars.next();
331 let mut prev = ' ';
332 while let Some(n) = chars.next() {
333 if prev == '*' && n == '/' { break; }
334 prev = n;
335 }
336 out.push(' ');
337 }
338 _ => out.push(c),
339 }
340 }
341 out.split_whitespace().collect::<Vec<_>>().join(" ")
342}
343
344/// Rewrite SQL literal/operator spellings into NQL's.
345///
346/// Only `'…'` → `"…"` and `<>` → `!=`. Done with an explicit scan rather than a
347/// regex so a quote inside a string cannot be mistaken for a delimiter: SQL
348/// escapes an embedded quote by doubling it (`'it''s'`), and that has to become
349/// a single character inside the NQL string rather than terminating it.
350/// Drop the table qualifier from every column reference in a clause tail.
351///
352/// # The silent wrong answer this removes
353///
354/// NQL has no notion of a qualifier: `field_value` looks a field up FLAT, in
355/// one map. So `WHERE orders.status = 'paid'` asked for a field literally
356/// named `orders.status`, no document had one, and the query returned ZERO
357/// ROWS — with no error and no warning, an empty result that reads exactly
358/// like "you have no paid orders".
359///
360/// Every ORM qualifies its predicates. SQLAlchemy emits
361/// `SELECT orders._id FROM orders WHERE orders.status = 'paid'` for the most
362/// ordinary filter there is, so EVERY filtered query answered empty, `.get(pk)`
363/// answered `None`, and `filter_by` answered nothing. The select list had
364/// always stripped qualifiers; the tail was "handed to the NQL parser
365/// unchanged", which is right for the clause GRAMMAR and wrong for a name NQL
366/// cannot interpret.
367///
368/// # Why a mismatched qualifier is an ERROR, not a strip
369///
370/// A qualifier naming something other than this statement's own collection
371/// means the query referenced a relation that is not in its FROM clause.
372/// Stripping it would answer with rows from the one relation that IS there —
373/// a different wrong answer wearing the same empty-looking clothes. Aliases
374/// are refused on this path already, so the collection's own name is the only
375/// qualifier that can be correct.
376///
377/// Runs BEFORE `sql_literals_to_nql`, so only SQL's single-quoted strings have
378/// to be skipped — the rewrite to NQL's double-quoted form has not happened
379/// yet, and a qualifier can never appear inside a literal.
380fn strip_column_qualifiers(
381 tail: &str,
382 coll: &str,
383 alias: Option<&str>,
384) -> Result<String, String> {
385 let bare = coll.rsplit('.').next().unwrap_or(coll);
386 let b: Vec<char> = tail.chars().collect();
387 let mut out = String::with_capacity(tail.len());
388 let mut i = 0usize;
389 let ident_start = |c: char| c.is_alphabetic() || c == '_';
390 let ident_char = |c: char| c.is_alphanumeric() || c == '_';
391
392 while i < b.len() {
393 // A single-quoted literal is copied through untouched.
394 if b[i] == '\'' {
395 out.push(b[i]);
396 i += 1;
397 while i < b.len() {
398 out.push(b[i]);
399 if b[i] == '\'' {
400 // A doubled '' is one literal quote, not a close.
401 if b.get(i + 1) == Some(&'\'') {
402 out.push('\'');
403 i += 2;
404 continue;
405 }
406 i += 1;
407 break;
408 }
409 i += 1;
410 }
411 continue;
412 }
413 // A double-quoted run is copied through too. NQL reads double quotes
414 // as a STRING delimiter rather than an identifier one, so a SQL
415 // delimited identifier is a genuine divergence — but it already fails
416 // LOUDLY in the NQL parser ("expected field name, got Str"), and a
417 // loud failure is not this function's problem to solve quietly.
418 if b[i] == '"' {
419 out.push(b[i]);
420 i += 1;
421 while i < b.len() {
422 out.push(b[i]);
423 if b[i] == '"' { i += 1; break; }
424 i += 1;
425 }
426 continue;
427 }
428 if !ident_start(b[i]) {
429 // A number like `1.5` starts with a digit, so it never enters the
430 // identifier branch and its dot is never touched.
431 out.push(b[i]);
432 i += 1;
433 continue;
434 }
435
436 let start = i;
437 while i < b.len() && ident_char(b[i]) {
438 i += 1;
439 }
440 let word: String = b[start..i].iter().collect();
441
442 // `qual.field` — a dot followed immediately by another identifier.
443 if b.get(i) == Some(&'.') && b.get(i + 1).is_some_and(|c| ident_start(*c)) {
444 let fstart = i + 1;
445 let mut j = fstart;
446 while j < b.len() && ident_char(b[j]) {
447 j += 1;
448 }
449 let field: String = b[fstart..j].iter().collect();
450 // A qualified FUNCTION call (`pg_catalog.something(`) is left
451 // exactly as written: this path does not implement functions at
452 // all, and NQL's own refusal names the function, which is more use
453 // to the reader than a claim about relations.
454 let is_call = b[j..].iter().find(|c| !c.is_whitespace()) == Some(&'(');
455 if is_call {
456 out.push_str(&word);
457 out.push('.');
458 out.push_str(&field);
459 i = j;
460 continue;
461 }
462 let matches_alias = alias.is_some_and(|a| word.eq_ignore_ascii_case(a));
463 if matches_alias || word.eq_ignore_ascii_case(bare) || word.eq_ignore_ascii_case(coll) {
464 out.push_str(&field);
465 i = j;
466 continue;
467 }
468 return Err(format!(
469 "no table or alias named {:?} in this query — this statement reads \
470 {:?}{}, and a qualifier naming anything else would have to be \
471 answered from a relation that is not in its FROM clause",
472 word,
473 bare,
474 alias.map(|a| format!(" (aliased {:?})", a)).unwrap_or_default()
475 ));
476 }
477 out.push_str(&word);
478 }
479 Ok(out)
480}
481
482/// Rewrite `SELECT count(*) FROM (<inner>) [AS] alias` into a flat count over
483/// the inner query's own collection and predicate — or `None` when the shapes
484/// do not permit it.
485///
486/// `None` is a REFUSAL, never a fallback: every caller reports the boundary
487/// rather than trying something else, because the alternative to an exact
488/// count is a wrong one.
489fn flatten_count_of_subquery(projection: &str, rest: &str) -> Option<String> {
490 // The outer select list must be nothing but `count(*)`, optionally
491 // aliased. Any other column would have to come from the derived table's
492 // output, which a flat count does not produce.
493 let (outer_expr, outer_alias) = split_output_alias(projection.trim());
494 let ou = outer_expr.to_uppercase().replace(' ', "");
495 if ou != "COUNT(*)" {
496 return None;
497 }
498
499 // Take the balanced parenthesised span, honouring literals so a `)` inside
500 // a string cannot close it early.
501 let b: Vec<char> = rest.chars().collect();
502 let mut depth = 0i32;
503 let mut in_s = false;
504 let mut end = None;
505 for (i, &c) in b.iter().enumerate() {
506 match c {
507 '\'' => in_s = !in_s,
508 '(' if !in_s => depth += 1,
509 ')' if !in_s => {
510 depth -= 1;
511 if depth == 0 {
512 end = Some(i);
513 break;
514 }
515 }
516 _ => {}
517 }
518 }
519 let end = end?;
520 let inner = b[1..end].iter().collect::<String>().trim().to_string();
521
522 // Nothing may follow the derived table but its alias — a join or a second
523 // FROM item changes what is being counted.
524 let trailing = b[end + 1..].iter().collect::<String>();
525 let (_alias, after) = split_table_alias(trailing.trim());
526 if !after.trim().is_empty() {
527 return None;
528 }
529
530 let iu = inner.to_uppercase();
531 if !iu.starts_with("SELECT") {
532 return None;
533 }
534 // Each of these would make the inner row count differ from the flat one.
535 for kw in ["LIMIT", "OFFSET", "GROUP BY", "HAVING", "UNION", "INTERSECT", "EXCEPT", "JOIN"] {
536 if find_kw(&iu, kw).is_some() {
537 return None;
538 }
539 }
540 if find_kw(&iu, "DISTINCT").is_some() {
541 return None;
542 }
543 // An inner aggregate already reduced the rows to one.
544 let inner_from = find_kw(&iu, "FROM")?;
545 let inner_list = inner[..inner_from].to_uppercase();
546 for agg in ["COUNT(", "SUM(", "AVG(", "MIN(", "MAX(", "ARRAY_AGG(", "STRING_AGG("] {
547 if inner_list.contains(agg) {
548 return None;
549 }
550 }
551 // A nested derived table is not walked — one level is the claim.
552 let inner_rest = inner[inner_from + 4..].trim();
553 if inner_rest.starts_with('(') {
554 return None;
555 }
556
557 // `ORDER BY` cannot change a count, so it is dropped rather than refused.
558 let mut tail = inner_rest.to_string();
559 let tu = tail.to_uppercase();
560 if let Some(ob) = find_kw(&tu, "ORDER BY") {
561 tail = tail[..ob].trim_end().to_string();
562 }
563 Some(format!(
564 "SELECT count(*){} FROM {}",
565 outer_alias.map(|a| format!(" AS {}", a)).unwrap_or_default(),
566 tail
567 ))
568}
569
570/// Words that begin a clause and can therefore never be a bare table alias.
571///
572/// `AS` is absent on purpose: it introduces an alias, and `AS OF` is
573/// disambiguated by looking at the word after it.
574const CLAUSE_WORDS: &[&str] = &[
575 "WHERE", "GROUP", "ORDER", "LIMIT", "OFFSET", "HAVING", "FOR", "VALID",
576 "TRACE", "TRAVERSE", "SEARCH", "RETURNING", "UNION", "INTERSECT", "EXCEPT",
577 "JOIN", "LEFT", "RIGHT", "INNER", "FULL", "CROSS", "ON", "USING", "SET",
578];
579
580/// Take a table alias off the front of a clause tail: `FROM orders o WHERE …`.
581///
582/// Returns the alias and the rest of the tail. The alias is REMOVED because
583/// NQL has no table-alias syntax and would report an "unexpected token" on it
584/// — which is how `FROM orders o` used to fail. Removing it here and teaching
585/// `strip_column_qualifiers` to accept it is what makes `SELECT o.status FROM
586/// orders o` work at all.
587///
588/// `AS OF SYSTEM TIME` also starts with `AS`, so the word AFTER `AS` decides:
589/// `AS OF` is a time-travel clause, anything else is an alias.
590fn split_table_alias(tail: &str) -> (Option<String>, &str) {
591 let t = tail.trim_start();
592 let first_end = t.find(char::is_whitespace).unwrap_or(t.len());
593 let first = &t[..first_end];
594 let fu = first.to_uppercase();
595
596 if fu == "AS" {
597 let rest = t[first_end..].trim_start();
598 let end = rest.find(char::is_whitespace).unwrap_or(rest.len());
599 let word = &rest[..end];
600 if word.eq_ignore_ascii_case("OF") {
601 return (None, t); // `AS OF …`, not an alias
602 }
603 if word.is_empty() {
604 return (None, t);
605 }
606 return (Some(word.trim_matches('"').to_string()), rest[end..].trim_start());
607 }
608 if first.is_empty() || CLAUSE_WORDS.contains(&fu.as_str()) {
609 return (None, t);
610 }
611 // A bare identifier here can only be an alias — the collection name was
612 // already consumed by the caller.
613 if first.chars().next().is_some_and(|c| c.is_alphabetic() || c == '_' || c == '"') {
614 return (Some(first.trim_matches('"').to_string()), t[first_end..].trim_start());
615 }
616 (None, t)
617}
618
619/// Split on a delimiter that is at PAREN DEPTH ZERO and outside a literal.
620///
621/// `projection.split(',')` cuts `SUM(a, b)` in half; a select list is not a
622/// flat comma list once it can contain calls.
623fn split_top_level(s: &str, delim: char) -> Vec<String> {
624 let mut out = vec![];
625 let mut cur = String::new();
626 let mut depth = 0i32;
627 let mut in_s = false;
628 let mut in_d = false;
629 for c in s.chars() {
630 match c {
631 '\'' if !in_d => { in_s = !in_s; cur.push(c); }
632 '"' if !in_s => { in_d = !in_d; cur.push(c); }
633 '(' if !in_s && !in_d => { depth += 1; cur.push(c); }
634 ')' if !in_s && !in_d => { depth -= 1; cur.push(c); }
635 c if c == delim && depth == 0 && !in_s && !in_d => {
636 out.push(std::mem::take(&mut cur));
637 }
638 _ => cur.push(c),
639 }
640 }
641 out.push(cur);
642 out
643}
644
645/// Split `expr AS name` / `expr name` into the expression and its output name.
646///
647/// The alias is the name the CLIENT will look the column up by — SQLAlchemy
648/// reads `count(*) AS count_1` back as `count_1`, so dropping the alias and
649/// returning a column called `count` hands it a result it cannot find.
650fn split_output_alias(p: &str) -> (&str, Option<&str>) {
651 let pu = p.to_uppercase();
652 if let Some(at) = find_kw(&pu, "AS") {
653 let alias = p[at + 2..].trim().trim_matches('"');
654 if !alias.is_empty() {
655 return (p[..at].trim(), Some(alias));
656 }
657 }
658 // A bare alias: `count(*) count_1`. Only after a closing paren or a plain
659 // identifier, and never when the tail is itself part of the expression —
660 // so the split point is the LAST whitespace outside any parenthesis.
661 let b: Vec<char> = p.chars().collect();
662 let mut depth = 0i32;
663 let mut in_s = false;
664 let mut cut = None;
665 for (i, &c) in b.iter().enumerate() {
666 match c {
667 '\'' => in_s = !in_s,
668 '(' if !in_s => depth += 1,
669 ')' if !in_s => depth -= 1,
670 c if c.is_whitespace() && depth == 0 && !in_s => cut = Some(i),
671 _ => {}
672 }
673 }
674 match cut {
675 Some(i) => {
676 let alias = p[i..].trim().trim_matches('"');
677 if alias.is_empty() { (p, None) } else { (p[..i].trim(), Some(alias)) }
678 }
679 None => (p, None),
680 }
681}
682
683/// One string, in NQL's spelling — double-quoted, inner quotes escaped.
684///
685/// These values arrive already UNQUOTED from the SQL parser, so they cannot be
686/// pasted into an NQL query as-is: a value containing `"` would close the
687/// literal early and the rest of it would be parsed as grammar. Which is the
688/// shape of an injection, not merely a syntax error.
689fn nql_string(s: &str) -> String {
690 format!("\"{}\"", s.replace('\\', "\\\\").replace('"', "\\\""))
691}
692
693fn sql_literals_to_nql(s: &str) -> String {
694 let mut out = String::with_capacity(s.len());
695 let mut it = s.chars().peekable();
696 while let Some(c) = it.next() {
697 match c {
698 '\'' => {
699 out.push('"');
700 while let Some(ch) = it.next() {
701 if ch == '\'' {
702 if it.peek() == Some(&'\'') {
703 it.next();
704 out.push('\''); // doubled '' is one literal quote
705 } else {
706 break;
707 }
708 } else if ch == '"' {
709 // A double quote inside a SQL literal must be escaped
710 // for NQL, whose lexer collapses \" to a literal quote.
711 out.push('\\');
712 out.push('"');
713 } else {
714 out.push(ch);
715 }
716 }
717 out.push('"');
718 }
719 '<' if it.peek() == Some(&'>') => { it.next(); out.push_str("!="); }
720 _ => out.push(c),
721 }
722 }
723 out
724}
725
726fn strip_prefix_ci(s: &str, prefix: &str) -> Option<String> {
727 if s.len() >= prefix.len() && s[..prefix.len()].eq_ignore_ascii_case(prefix) {
728 Some(s[prefix.len()..].trim_start().to_string())
729 } else {
730 None
731 }
732}
733
734/// Find a top-level keyword (not inside quotes or parentheses), returning its
735/// byte offset. Case-insensitive, and only matches on word boundaries.
736fn find_kw(s: &str, kw: &str) -> Option<usize> {
737 let bytes = s.as_bytes();
738 let k = kw.as_bytes();
739 let mut depth = 0i32;
740 let mut in_s = false;
741 let mut in_d = false;
742 let mut i = 0usize;
743 while i < bytes.len() {
744 let c = bytes[i];
745 if in_s { if c == b'\'' { in_s = false; } i += 1; continue; }
746 if in_d { if c == b'"' { in_d = false; } i += 1; continue; }
747 match c {
748 b'\'' => { in_s = true; i += 1; continue; }
749 b'"' => { in_d = true; i += 1; continue; }
750 b'(' => { depth += 1; i += 1; continue; }
751 b')' => { depth -= 1; i += 1; continue; }
752 _ => {}
753 }
754 if depth == 0 && i + k.len() <= bytes.len()
755 && bytes[i..i + k.len()].eq_ignore_ascii_case(k)
756 {
757 let before_ok = i == 0 || !(bytes[i - 1] as char).is_alphanumeric() && bytes[i - 1] != b'_';
758 let after = i + k.len();
759 let after_ok = after >= bytes.len()
760 || !(bytes[after] as char).is_alphanumeric() && bytes[after] != b'_';
761 if before_ok && after_ok {
762 return Some(i);
763 }
764 }
765 i += 1;
766 }
767 None
768}
769
770/// Split a comma-separated list at the TOP level, ignoring commas inside
771/// quotes or parentheses — so `VALUES (1, 'a,b'), (2, 'c')` splits into two
772/// groups and not four.
773fn split_top(s: &str, sep: char) -> Vec<String> {
774 let mut out = vec![];
775 let mut cur = String::new();
776 let mut depth = 0i32;
777 let mut in_s = false;
778 let mut it = s.chars().peekable();
779 while let Some(c) = it.next() {
780 if in_s {
781 cur.push(c);
782 if c == '\'' {
783 // A doubled '' is an escaped quote, not the end of the literal.
784 if it.peek() == Some(&'\'') { cur.push(it.next().unwrap()); } else { in_s = false; }
785 }
786 continue;
787 }
788 match c {
789 '\'' => { in_s = true; cur.push(c); }
790 '(' => { depth += 1; cur.push(c); }
791 ')' => { depth -= 1; cur.push(c); }
792 x if x == sep && depth == 0 => { out.push(cur.trim().to_string()); cur.clear(); }
793 _ => cur.push(c),
794 }
795 }
796 if !cur.trim().is_empty() { out.push(cur.trim().to_string()); }
797 out
798}
799
800/// Parse one SQL scalar literal into JSON.
801///
802/// Deliberately narrow: a string, a number, a boolean, or NULL. Anything else
803/// — a function call, an expression, a cast — is refused by name rather than
804/// coerced into a string that would silently store the wrong value.
805fn sql_value(raw: &str) -> Result<Value, String> {
806 let t = raw.trim();
807 if t.is_empty() {
808 return Err("empty value".into());
809 }
810 let up = t.to_uppercase();
811 if up == "NULL" { return Ok(Value::Null); }
812 if up == "TRUE" { return Ok(Value::Bool(true)); }
813 if up == "FALSE" { return Ok(Value::Bool(false)); }
814 if t.starts_with('\'') && t.ends_with('\'') && t.len() >= 2 {
815 // Unwrap, collapsing the SQL '' escape to one quote.
816 let inner = &t[1..t.len() - 1];
817 return Ok(Value::String(inner.replace("''", "'")));
818 }
819 if let Ok(i) = t.parse::<i64>() { return Ok(Value::from(i)); }
820 if let Ok(f) = t.parse::<f64>() { return Ok(Value::from(f)); }
821 Err(format!(
822 "cannot use {:?} as a value — this endpoint accepts string literals, \
823 numbers, TRUE/FALSE and NULL. Expressions, casts and function calls \
824 are not evaluated, because storing an unevaluated expression as text \
825 would be worse than refusing it", t))
826}
827
828/// Pull a trailing `RETURNING …` off a statement, returning (head, columns).
829fn split_returning(tail: &str) -> (String, Vec<Col>) {
830 let tu = tail.to_uppercase();
831 match find_kw(&tu, "RETURNING") {
832 None => (tail.to_string(), vec![]),
833 Some(at) => {
834 let head = tail[..at].trim().to_string();
835 let list = tail[at + "RETURNING".len()..].trim();
836 if list == "*" {
837 return (head, vec![]); // empty projection = every column
838 }
839 let cols = split_top(list, ',')
840 .into_iter()
841 .map(|p| {
842 let raw = p.split_whitespace().next().unwrap_or(&p).to_string();
843 let name = raw.rsplit('.').next().unwrap_or(&raw).trim_matches('"').to_string();
844 Col::same(&name)
845 })
846 .collect();
847 (head, cols)
848 }
849 }
850}
851
852/// Columns whose names are reserved: they carry provenance rather than data.
853fn take_reserved(doc: &mut serde_json::Map<String, Value>) -> (Option<String>, Vec<String>, Option<String>, Option<String>) {
854 let id = doc.remove("_id").or_else(|| doc.remove("id"))
855 .and_then(|v| match v {
856 Value::String(s) => Some(s),
857 Value::Null => None,
858 other => Some(other.to_string()), // a numeric key is a fine id
859 });
860 let caused_by = match doc.remove("_caused_by") {
861 Some(Value::String(s)) => vec![s],
862 Some(Value::Array(a)) => a.into_iter()
863 .filter_map(|v| v.as_str().map(str::to_string)).collect(),
864 _ => vec![],
865 };
866 let vf = doc.remove("_valid_from").and_then(|v| v.as_str().map(str::to_string));
867 let vt = doc.remove("_valid_to").and_then(|v| v.as_str().map(str::to_string));
868 (id, caused_by, vf, vt)
869}
870
871/// `INSERT INTO coll (c1, c2) VALUES (v1, v2), (…) [RETURNING …]`
872fn translate_insert(sql: &str) -> Result<Stmt, String> {
873 let rest = strip_prefix_ci(sql, "INSERT")
874 .and_then(|r| strip_prefix_ci(&r, "INTO"))
875 .ok_or("expected INSERT INTO")?;
876 // Locate VALUES first. Everything before it is `coll (col, …)`; searching
877 // for `(` without that bound finds the VALUES parenthesis instead and
878 // swallows the keyword into the collection name.
879 let ru = rest.to_uppercase();
880 let values_at = find_kw(&ru, "VALUES").ok_or(
881 "expected VALUES — `INSERT … SELECT` is not supported on this endpoint")?;
882 let head = rest[..values_at].trim().to_string();
883 let open = head.find('(').ok_or(
884 "INSERT needs an explicit column list — `INSERT INTO t (a, b) VALUES (…)`. \
885 NEDB is schemaless, so there is no declared column order to infer from")?;
886 let coll = head[..open].trim().trim_matches('"');
887 let coll = coll.rsplit('.').next().unwrap_or(coll).to_string();
888 if coll.is_empty() {
889 return Err("expected a collection name after INSERT INTO".into());
890 }
891 let close = head.rfind(')').ok_or("unterminated column list")?;
892 if close < open {
893 return Err("malformed column list".into());
894 }
895 let tail_from_values = rest[values_at..].to_string();
896 let cols: Vec<String> = split_top(&head[open + 1..close], ',')
897 .into_iter()
898 .map(|c| c.trim().trim_matches('"').to_string())
899 .collect();
900 if cols.is_empty() {
901 return Err("the column list is empty".into());
902 }
903
904 let after = strip_prefix_ci(&tail_from_values, "VALUES")
905 .ok_or("expected VALUES after the column list")?;
906 let (values_part, returning) = split_returning(&after);
907
908 let mut rows = vec![];
909 for group in split_top(&values_part, ',') {
910 let g = group.trim();
911 if !(g.starts_with('(') && g.ends_with(')')) {
912 return Err(format!("expected a parenthesised row of values, got {:?}", g));
913 }
914 let vals = split_top(&g[1..g.len() - 1], ',');
915 if vals.len() != cols.len() {
916 return Err(format!(
917 "{} values for {} columns — every row must match the column list",
918 vals.len(), cols.len()));
919 }
920 let mut doc = serde_json::Map::new();
921 for (c, v) in cols.iter().zip(vals.iter()) {
922 doc.insert(c.clone(), sql_value(v)?);
923 }
924 let (id, caused_by, valid_from, valid_to) = take_reserved(&mut doc);
925 rows.push(InsertRow { id, doc, caused_by, valid_from, valid_to });
926 }
927 if rows.is_empty() {
928 return Err("INSERT with no rows".into());
929 }
930 Ok(Stmt::Insert { coll, rows, returning })
931}
932
933/// `UPDATE coll SET a = 1, b = 'x' [WHERE …] [RETURNING …]`
934fn translate_update(sql: &str) -> Result<Stmt, String> {
935 let rest = strip_prefix_ci(sql, "UPDATE").ok_or("expected UPDATE")?;
936 let ru = rest.to_uppercase();
937 let set_at = find_kw(&ru, "SET").ok_or("expected SET in UPDATE")?;
938 // `UPDATE orders o SET …` — Postgres allows an alias here, and taking the
939 // whole span as the collection name made it part of the name ("orders o").
940 let target = rest[..set_at].trim();
941 let mut parts = target.split_whitespace();
942 let coll = parts.next().unwrap_or("").trim_matches('"');
943 let coll = coll.rsplit('.').next().unwrap_or(coll).to_string();
944 let upd_alias: Option<String> = match parts.next() {
945 Some(w) if w.eq_ignore_ascii_case("AS") => {
946 parts.next().map(|a| a.trim_matches('"').to_string())
947 }
948 Some(w) => Some(w.trim_matches('"').to_string()),
949 None => None,
950 };
951 if coll.is_empty() {
952 return Err("expected a collection name after UPDATE".into());
953 }
954 let after_set = rest[set_at + 3..].trim().to_string();
955 let (after_set, returning) = split_returning(&after_set);
956
957 // WHERE ends the assignment list; everything after it is a NQL predicate.
958 let au = after_set.to_uppercase();
959 let (assigns_raw, where_raw) = match find_kw(&au, "WHERE") {
960 Some(at) => (after_set[..at].to_string(), after_set[at..].to_string()),
961 None => (after_set.clone(), String::new()),
962 };
963
964 let mut set = vec![];
965 for a in split_top(&assigns_raw, ',') {
966 let eq = a.find('=').ok_or(format!("expected `col = value` in SET, got {:?}", a))?;
967 let col = a[..eq].trim().trim_matches('"').to_string();
968 if col.is_empty() {
969 return Err("empty column name in SET".into());
970 }
971 set.push((col, sql_value(&a[eq + 1..])?));
972 }
973 if set.is_empty() {
974 return Err("UPDATE with no assignments".into());
975 }
976 // The matching rows are found with an ordinary NQL read, so the whole
977 // predicate surface (IN, BETWEEN, LIKE, OR, …) works in an UPDATE too.
978 let where_raw = strip_column_qualifiers(where_raw.trim(), &coll, upd_alias.as_deref())?;
979 let nql = format!("FROM {} {}", coll, sql_literals_to_nql(&where_raw))
980 .trim().to_string();
981 Ok(Stmt::Update { coll, set, where_sql: where_raw, nql, returning })
982}
983
984/// `DELETE FROM coll [WHERE …] [RETURNING …]`
985fn translate_delete(sql: &str) -> Result<Stmt, String> {
986 let rest = strip_prefix_ci(sql, "DELETE")
987 .and_then(|r| strip_prefix_ci(&r, "FROM"))
988 .ok_or("expected DELETE FROM")?;
989 let (rest, returning) = split_returning(&rest);
990 let end = rest.find(' ').unwrap_or(rest.len());
991 let coll = rest[..end].trim().trim_matches('"');
992 let coll = coll.rsplit('.').next().unwrap_or(coll).to_string();
993 if coll.is_empty() {
994 return Err("expected a collection name after DELETE FROM".into());
995 }
996 let (del_alias, where_raw) = split_table_alias(rest[end..].trim());
997 let where_raw = strip_column_qualifiers(where_raw, &coll, del_alias.as_deref())?;
998 let nql = format!("FROM {} {}", coll, sql_literals_to_nql(&where_raw))
999 .trim().to_string();
1000 Ok(Stmt::Delete { coll, where_sql: where_raw, nql, returning })
1001}
1002
1003/// Translate one SQL statement into something executable, or explain why not.
1004pub fn translate(sql_raw: &str) -> Result<Stmt, String> {
1005 let sql = normalise(sql_raw);
1006 let sql = sql.trim().trim_end_matches(';').trim();
1007 if sql.is_empty() {
1008 return Ok(Stmt::Ok(""));
1009 }
1010 let upper = sql.to_uppercase();
1011
1012 // ── the handshake. Clients issue these before anything useful; answering
1013 // them with plausible values is the difference between "connects" and
1014 // "hangs on startup". They are canned on purpose — NEDB has no pg_catalog
1015 // and pretending otherwise would be worse than a clear boundary.
1016 if upper.starts_with("SET ") || upper.starts_with("BEGIN") || upper.starts_with("COMMIT")
1017 || upper.starts_with("ROLLBACK") || upper.starts_with("DISCARD")
1018 || upper.starts_with("LISTEN ") || upper.starts_with("UNLISTEN ")
1019 {
1020 // Accepted and ignored: there is one implicit read-only transaction.
1021 return Ok(Stmt::Ok(if upper.starts_with("SET") { "SET" } else { "OK" }));
1022 }
1023 if upper.starts_with("SHOW ") {
1024 let name = sql[5..].trim().to_lowercase();
1025 let val = match name.as_str() {
1026 "transaction_isolation" | "default_transaction_isolation" => "read committed",
1027 "server_version" => SERVER_VERSION,
1028 "server_encoding" | "client_encoding" => "UTF8",
1029 "standard_conforming_strings" => "on",
1030 "is_superuser" => "off",
1031 _ => "",
1032 };
1033 return Ok(Stmt::Canned { cols: vec![name], row: vec![val.to_string()] });
1034 }
1035 if upper == "SELECT VERSION()" {
1036 return Ok(Stmt::Canned {
1037 cols: vec!["version".into()],
1038 row: vec![full_version_string()],
1039 });
1040 }
1041 if upper == "SELECT 1" || upper == "SELECT 1;" {
1042 return Ok(Stmt::Canned { cols: vec!["?column?".into()], row: vec!["1".into()] });
1043 }
1044 if upper.starts_with("SELECT CURRENT_SCHEMA") {
1045 return Ok(Stmt::Canned { cols: vec!["current_schema".into()], row: vec!["public".into()] });
1046 }
1047 if upper.starts_with("SELECT CURRENT_DATABASE") {
1048 return Ok(Stmt::Canned { cols: vec!["current_database".into()], row: vec!["nedb".into()] });
1049 }
1050 if upper.starts_with("SELECT CURRENT_USER") || upper.starts_with("SELECT USER") {
1051 return Ok(Stmt::Canned { cols: vec!["current_user".into()], row: vec!["nedb".into()] });
1052 }
1053
1054 // ── writes ───────────────────────────────────────────────────────────────
1055 // SQL's write semantics and NEDB's append-only model line up, so these are
1056 // first-class rather than refused. See the `Stmt` doc comment.
1057 if upper.starts_with("INSERT") { return translate_insert(sql); }
1058 if upper.starts_with("UPDATE") { return translate_update(sql); }
1059 if upper.starts_with("DELETE") { return translate_delete(sql); }
1060
1061 // ── the refusals that remain, each naming the boundary ──────────────────
1062 for (kw, why) in [
1063 ("CREATE", "DDL is not supported — collections are created implicitly by the first write to them, because NEDB is schemaless"),
1064 ("ALTER", "DDL is not supported — there is no schema to alter"),
1065 ("DROP", "DDL is not supported; drop a database with DELETE /v1/databases/<db>"),
1066 ("TRUNCATE", "not supported, and not an oversight: NEDB is append-only so that history cannot be discarded. That is the product"),
1067 ("COPY", "not supported; use GET /v1/databases/<db>/since for bulk export"),
1068 ("GRANT", "there is no SQL-level privilege system; auth is the bearer token"),
1069 ("REVOKE", "there is no SQL-level privilege system; auth is the bearer token"),
1070 ] {
1071 if upper.starts_with(kw) {
1072 return Err(format!("{} is not supported — {}", kw, why));
1073 }
1074 }
1075 if !upper.starts_with("SELECT") {
1076 return Err(format!(
1077 "only SELECT, INSERT, UPDATE and DELETE are supported on the Postgres \
1078 endpoint (got {:?})",
1079 sql.split_whitespace().next().unwrap_or("")
1080 ));
1081 }
1082 for (kw, why) in [
1083 (" JOIN ", "JOIN is not supported — NQL is single-collection; join in your client or model the relation with LINK/TRAVERSE"),
1084 (" UNION ", "UNION is not supported"),
1085 (" INTERSECT ", "INTERSECT is not supported"),
1086 (" EXCEPT ", "EXCEPT is not supported"),
1087 (" OVER (", "window functions are not supported"),
1088 ("DISTINCT ", "DISTINCT is not supported — GROUP BY <col> gives the distinct values with counts"),
1089 ] {
1090 if upper.contains(kw) {
1091 return Err(why.to_string());
1092 }
1093 }
1094 if find_kw(&upper, "FROM").is_none() {
1095 return Err("SELECT without FROM is not supported on this endpoint".into());
1096 }
1097
1098 // ── SELECT <projection> FROM <rest> ──────────────────────────────────────
1099 let after_select = strip_prefix_ci(sql, "SELECT").ok_or("expected SELECT")?;
1100 let from_at = find_kw(&after_select.to_uppercase(), "FROM")
1101 .ok_or("expected FROM after the select list")?;
1102 let projection = after_select[..from_at].trim().to_string();
1103 let rest = after_select[from_at + 4..].trim().to_string();
1104 if rest.is_empty() {
1105 return Err("expected a collection name after FROM".into());
1106 }
1107 // ── the one derived table with a provable flat equivalent ───────────────
1108 //
1109 // `SELECT count(*) FROM (SELECT … FROM coll WHERE …) AS anon` is what
1110 // EVERY ORM emits for `.count()` — SQLAlchemy's `Query.count()` wraps the
1111 // whole query in a subquery unconditionally. Refusing it means "SQLAlchemy
1112 // works, except counting", which is not a boundary anyone would accept.
1113 //
1114 // Counting a derived table whose rows are exactly the inner query's rows
1115 // is counting the inner query, so the rewrite is an IDENTITY rather than
1116 // an approximation. Each guard below names a construct that would break
1117 // that identity, and anything carrying one is still refused:
1118 //
1119 // * `LIMIT` / `OFFSET` — caps the row count before it is counted
1120 // * `DISTINCT` — collapses duplicates, so the counts differ
1121 // * `GROUP BY` — the inner rows ARE the groups
1122 // * an inner aggregate — already one row, counting it answers 1
1123 // * anything but `count(*)` outside — the outer list would need the
1124 // inner columns, which a flat count cannot supply
1125 if rest.starts_with('(') {
1126 if let Some(flat) = flatten_count_of_subquery(&projection, &rest) {
1127 // Recurses ONCE at most: the rewrite is only produced when the
1128 // inner FROM names a real collection, so the flat statement can
1129 // never re-enter this branch.
1130 return translate(&flat);
1131 }
1132 return Err("subqueries in FROM are not supported — except \
1133 `SELECT count(*) FROM (…)`, which is rewritten to a flat \
1134 count when the inner query has no LIMIT, OFFSET, DISTINCT, \
1135 GROUP BY or aggregate of its own (any of those would make the \
1136 two counts different numbers)".into());
1137 }
1138 let coll_end = rest.find(' ').unwrap_or(rest.len());
1139 let coll = &rest[..coll_end];
1140 if coll.contains(',') {
1141 return Err("selecting from more than one collection is not supported (no JOIN)".into());
1142 }
1143 // Postgres clients often qualify as schema.table; NEDB has one namespace,
1144 // so the schema is dropped — EXCEPT for `information_schema`, whose table
1145 // names (`tables`, `columns`) are words a user could plausibly name a
1146 // collection. Keeping the qualifier there is what stops
1147 // `SELECT * FROM information_schema.tables` and a real collection called
1148 // `tables` from resolving to the same thing.
1149 let bare = coll.rsplit('.').next().unwrap_or(coll).trim_matches('"');
1150 let qualified = coll
1151 .split('.')
1152 .map(|p| p.trim_matches('"'))
1153 .collect::<Vec<_>>()
1154 .join(".");
1155 let coll = if qualified.starts_with("information_schema.") {
1156 qualified.as_str()
1157 } else {
1158 bare
1159 };
1160 let tail = rest[coll_end..].trim();
1161
1162 // ── the select list ──────────────────────────────────────────────────────
1163 //
1164 // Parsed ITEM BY ITEM, which is what lets a list MIX plain columns with an
1165 // aggregate — and that mixture is exactly what a `GROUP BY` query is.
1166 // SQLAlchemy writes `SELECT orders.status, count(*) AS count_1 FROM orders
1167 // GROUP BY orders.status` for the most ordinary grouped query there is,
1168 // and the previous check refused any list containing a parenthesis at all,
1169 // so the whole shape was unreachable even though NQL expresses it
1170 // natively.
1171 //
1172 // NQL's grouped row carries the group key, `count`, and at most one NAMED
1173 // aggregate — so `count(*)` is always available and one of SUM/AVG/MIN/MAX
1174 // may join it. A second named aggregate is refused by name rather than
1175 // silently dropped.
1176 let mut agg_clause = String::new();
1177 let mut agg_srcs: Vec<String> = vec![];
1178 let mut project: Vec<Col> = vec![];
1179
1180 if projection == "*" {
1181 // everything
1182 } else {
1183 for part in split_top_level(&projection, ',') {
1184 let p = part.trim();
1185 if p.is_empty() {
1186 return Err("empty column in the select list".into());
1187 }
1188 let (expr, alias) = split_output_alias(p);
1189 let eu = expr.to_uppercase();
1190
1191 // COUNT(*) and COUNT(col) both become NQL's bare COUNT: NQL counts
1192 // the group, and a per-column non-null count is not expressible.
1193 if eu.starts_with("COUNT(") {
1194 if agg_clause.is_empty() {
1195 agg_clause = " COUNT".to_string();
1196 }
1197 agg_srcs.push("count".to_string());
1198 project.push(Col::renamed("count", alias.unwrap_or("count")));
1199 continue;
1200 }
1201 if let Some(agg) = ["SUM", "AVG", "MIN", "MAX"]
1202 .iter()
1203 .find(|a| eu.starts_with(&format!("{}(", a)))
1204 {
1205 let inner = expr[agg.len() + 1..].trim_end_matches(')').trim();
1206 if inner.is_empty() || inner == "*" {
1207 return Err(format!("{}() needs a column", agg));
1208 }
1209 let inner = inner.rsplit('.').next().unwrap_or(inner).trim_matches('"');
1210 let named = format!("{} {}", agg, inner);
1211 if !agg_clause.is_empty() && agg_clause.trim() != "COUNT" && agg_clause.trim() != named {
1212 return Err(format!(
1213 "only one of SUM/AVG/MIN/MAX is supported per statement \
1214 (already have {:?}, then {:?}) — NQL's grouped row carries \
1215 the group key, `count`, and ONE named aggregate",
1216 agg_clause.trim(), named));
1217 }
1218 agg_clause = format!(" {}", named);
1219 // NQL emits `<agg>_<field>`; SQL names the column after the
1220 // function unless the query aliased it.
1221 let src = format!("{}_{}", agg.to_lowercase(), inner);
1222 project.push(Col::renamed(&src, alias.unwrap_or(&agg.to_lowercase())));
1223 agg_srcs.push(src);
1224 continue;
1225 }
1226 // A paren used to be the whole test for "is this an expression",
1227 // and it let every paren-free one through: `total * 2` became a
1228 // FIELD NAME, no document had a field called "total * 2", and the
1229 // column came back blank for every row with no error. Same silent
1230 // class as the qualified-WHERE bug -- a wrong answer that looks
1231 // like data. So the test is now the positive one: what survives
1232 // has to BE a column reference.
1233 let bare = expr.rsplit('.').next().unwrap_or(expr).trim_matches('"');
1234 let is_column = !bare.is_empty()
1235 && !bare.starts_with(|c: char| c.is_ascii_digit())
1236 && bare.chars().all(|c| c.is_alphanumeric() || c == '_' || c == '$');
1237 if !is_column {
1238 return Err(format!(
1239 "expressions in the select list are not supported ({:?}) — \
1240 supported: *, a column list, COUNT(*), or SUM/AVG/MIN/MAX(col). \
1241 Compute it in your client, or read the column and map it there",
1242 p));
1243 }
1244 let name = bare;
1245 project.push(Col::renamed(name, alias.unwrap_or(name)));
1246 }
1247 }
1248
1249 // ── clause tail: AS OF SYSTEM TIME → AS OF, then pass the rest through ──
1250 //
1251 // The clause keywords NQL shares with SQL (WHERE, GROUP BY, HAVING,
1252 // ORDER BY, LIMIT, OFFSET) are deliberately handed to the NQL parser
1253 // unchanged rather than re-parsed here. NQL is the authority on what is
1254 // valid; re-implementing its grammar would give two parsers to disagree.
1255 // `FROM orders o WHERE …` — the alias is taken off the tail (NQL has no
1256 // alias syntax) and then ACCEPTED as a qualifier on the columns.
1257 let (alias, tail) = split_table_alias(tail);
1258 let mut tail = strip_column_qualifiers(tail, coll, alias.as_deref())?;
1259 let tu = tail.to_uppercase();
1260 if let Some(at) = find_kw(&tu, "AS OF SYSTEM TIME") {
1261 let before = tail[..at].to_string();
1262 let after = tail[at + "AS OF SYSTEM TIME".len()..].trim_start().to_string();
1263 // Take the sequence token; the rest of the tail follows it.
1264 let end = after.find(' ').unwrap_or(after.len());
1265 let seq = after[..end].trim().trim_matches('\'').trim_matches('"').to_string();
1266 if seq.parse::<u64>().is_err() {
1267 return Err(format!(
1268 "AS OF SYSTEM TIME takes a NEDB sequence number here, not a timestamp (got {:?}). \
1269 NEDB's history is sequence-addressed and never garbage-collected, so a seq is \
1270 exact where a wall-clock time would be approximate", seq));
1271 }
1272 tail = format!("{} AS OF {} {}", before.trim(), seq, after[end..].trim())
1273 .trim()
1274 .to_string();
1275 }
1276
1277 // ── ORDER BY <ordinal> → ORDER BY <that select-list column> ─────────────
1278 //
1279 // SQL lets a sort key be a POSITION in the select list, and clients write
1280 // it constantly — `ORDER BY 1, 2` is how psql's own catalogue queries sort,
1281 // and node-postgres sent `GROUP BY status ORDER BY 1` in the very first
1282 // run of the driver harness. NQL has no ordinals: it read the `1` as a
1283 // literal and refused with "expected field name, got Num(1.0)".
1284 //
1285 // The projection is already parsed here, so the position resolves to a
1286 // real field name. An ordinal past the end of the select list, or one used
1287 // with `SELECT *` where there is no list to index, is refused with the
1288 // reason — guessing a column would sort by something the query never named.
1289 let tu_ord = tail.to_uppercase();
1290 if let Some(ob_at) = find_kw(&tu_ord, "ORDER BY") {
1291 let start = ob_at + "ORDER BY".len();
1292 // The clause runs to the next one, or to the end of the tail.
1293 let end = ["LIMIT", "OFFSET", "GROUP BY", "TRACE", "TRAVERSE", "SEARCH"]
1294 .iter()
1295 .filter_map(|k| find_kw(&tu_ord[start..], k).map(|at| start + at))
1296 .min()
1297 .unwrap_or(tail.len());
1298 let mut keys = vec![];
1299 for item in split_top_level(&tail[start..end], ',') {
1300 let item = item.trim();
1301 if item.is_empty() {
1302 continue;
1303 }
1304 let mut parts = item.split_whitespace();
1305 let first = parts.next().unwrap_or("");
1306 let rest: Vec<&str> = parts.collect();
1307 match first.parse::<usize>() {
1308 Ok(n) if n >= 1 => {
1309 let col = project.get(n - 1).ok_or_else(|| {
1310 if project.is_empty() {
1311 format!(
1312 "ORDER BY {} is a select-list POSITION, and `SELECT *` \
1313 has no list to index — name the column instead", n)
1314 } else {
1315 format!(
1316 "ORDER BY {} is out of range: the select list has {} \
1317 column(s)", n, project.len())
1318 }
1319 })?;
1320 keys.push(
1321 std::iter::once(col.src.as_str())
1322 .chain(rest.iter().copied())
1323 .collect::<Vec<_>>()
1324 .join(" "),
1325 );
1326 }
1327 // Not an ordinal — a named column, or `1 + 1`, which NQL will
1328 // judge for itself.
1329 _ => keys.push(item.to_string()),
1330 }
1331 }
1332 tail = format!("{} ORDER BY {} {}", &tail[..ob_at], keys.join(", "), &tail[end..])
1333 .split_whitespace()
1334 .collect::<Vec<_>>()
1335 .join(" ");
1336 }
1337
1338 // ── GROUP BY: refuse a bare column that SQL would refuse ─────────────────
1339 //
1340 // A grouped NQL row holds only the group key, `count` and the aggregate —
1341 // so projecting `total` from `GROUP BY region` found nothing and rendered
1342 // NULL. Silently answering NULL for a column the query cannot produce is
1343 // the exact failure shape this engine keeps getting bitten by, so it is an
1344 // error, using Postgres's own wording so the message is already familiar.
1345 let mut gkey: Option<String> = None;
1346 let tu_all = tail.to_uppercase();
1347 if let Some(gb_at) = find_kw(&tu_all, "GROUP BY") {
1348 let head = tail[..gb_at].trim_end().to_string();
1349 let after = tail[gb_at + "GROUP BY".len()..].trim_start();
1350 let key_end = after.find(|c: char| c == ' ' || c == ',').unwrap_or(after.len());
1351 let group_key = after[..key_end].trim().trim_matches('"').to_string();
1352 let after_key = after[key_end..].trim_start();
1353 gkey = Some(group_key.clone());
1354
1355 // NQL groups by ONE field. Taking the first key and leaving the rest
1356 // in the tail would group by something narrower than the query asked
1357 // for — more rows than Postgres returns, each aggregating too much.
1358 if after_key.starts_with(',') {
1359 return Err(format!(
1360 "GROUP BY takes one key here (got {:?} and more) — NQL groups by a \
1361 single field, and grouping by only the first would aggregate over \
1362 rows the query meant to keep apart",
1363 group_key));
1364 }
1365
1366 for c in &project {
1367 let ok = c.src == group_key
1368 || c.src == "count"
1369 || agg_srcs.contains(&c.src);
1370 if !ok {
1371 return Err(format!(
1372 "column {:?} must appear in the GROUP BY clause or be used in an \
1373 aggregate function — a grouped row carries the group key, `count`, \
1374 and the aggregate, nothing else",
1375 c.src));
1376 }
1377 }
1378
1379 // NQL's aggregate belongs IMMEDIATELY AFTER the group key
1380 // (`GROUP BY status COUNT`), not after the collection name. Emitting
1381 // `FROM orders COUNT GROUP BY status` is refused by the NQL parser
1382 // with "only one aggregate per query" — which is how the most
1383 // ordinary grouped query an ORM writes still failed even once its
1384 // select list parsed.
1385 //
1386 // `count` rides along free with a named aggregate — an NQL grouped row
1387 // carries the key, `count` AND the aggregate — so only the named one
1388 // is emitted when both were asked for.
1389 tail = format!("{} GROUP BY {}{} {}", head, group_key, agg_clause, after_key)
1390 .split_whitespace()
1391 .collect::<Vec<_>>()
1392 .join(" ");
1393 agg_clause.clear();
1394 }
1395
1396 // ── HAVING <agg> → the spelling NQL's grouped row actually carries ──────
1397 //
1398 // NQL's grouped row has fields named `count` and `<agg>_<field>`, and its
1399 // HAVING matches on those. Every SQL client writes something else:
1400 //
1401 // HAVING count(*) > 1 -> NQL parse error (loud, fine)
1402 // HAVING COUNT > 1 -> ZERO ROWS, no error
1403 // HAVING n > 1 -> ZERO ROWS, no error (`n` being the SQL alias)
1404 //
1405 // The last two are the dangerous ones: HAVING is advertised as supported,
1406 // and a filter that silently matches nothing reads as "no groups qualified"
1407 // rather than "your predicate named a field that does not exist". So the
1408 // aggregate spellings are translated, and anything left that is not a
1409 // group-key or aggregate field is refused BY NAME.
1410 let tu_hav = tail.to_uppercase();
1411 if let Some(h_at) = find_kw(&tu_hav, "HAVING") {
1412 let start = h_at + "HAVING".len();
1413 let end = ["ORDER BY", "LIMIT", "OFFSET"]
1414 .iter()
1415 .filter_map(|k| find_kw(&tu_hav[start..], k).map(|at| start + at))
1416 .min()
1417 .unwrap_or(tail.len());
1418 let clause = tail[start..end].to_string();
1419 // The left-hand side of the first comparison is the key being filtered.
1420 let lhs_end = clause
1421 .find(|c: char| "<>=!".contains(c))
1422 .unwrap_or(clause.len());
1423 let lhs = clause[..lhs_end].trim();
1424 if !lhs.is_empty() {
1425 let lu = lhs.to_uppercase();
1426 // `count(*)`, `COUNT(*)`, `count`, or the alias the query gave the
1427 // count -- all mean NQL's `count`.
1428 // The alias test has to tie THIS column to the count. Asking only
1429 // "is there a count anywhere in the projection" matched the GROUP
1430 // BY key too, so `HAVING status > 'a'` -- a perfectly legitimate
1431 // filter on the group key -- was rewritten into `count > 'a'`.
1432 let is_count = lu == "COUNT" || lu.replace(' ', "") == "COUNT(*)"
1433 || project.iter().any(|c| c.out.eq_ignore_ascii_case(lhs) && c.src == "count");
1434 let mapped = if is_count {
1435 Some("count".to_string())
1436 } else {
1437 // A named aggregate, by its NQL source name or by its alias.
1438 agg_srcs.iter().find(|s| s.eq_ignore_ascii_case(lhs)).cloned().or_else(|| {
1439 project.iter()
1440 .find(|c| c.out.eq_ignore_ascii_case(lhs) && agg_srcs.contains(&c.src))
1441 .map(|c| c.src.clone())
1442 })
1443 };
1444 match mapped {
1445 Some(m) => {
1446 // The space matters: `count> 1` happens to parse today, but
1447 // relying on the tokenizer being forgiving is how a rewrite
1448 // breaks the next time the grammar tightens.
1449 let rewritten = format!("{} {}", m, clause[lhs_end..].trim());
1450 tail = format!("{} HAVING {} {}",
1451 tail[..h_at].trim(), rewritten.trim(), tail[end..].trim())
1452 .trim().to_string();
1453 }
1454 None if gkey.as_deref().map(|g| g.eq_ignore_ascii_case(lhs)) == Some(true) => {}
1455 None => {
1456 return Err(format!(
1457 "HAVING names {:?}, which this grouped row does not carry. \
1458 It has the group key{}{}. Filtering on anything else would \
1459 answer zero rows rather than report a mistake",
1460 lhs,
1461 gkey.as_deref().map(|g| format!(" ({:?})", g)).unwrap_or_default(),
1462 if agg_srcs.is_empty() { String::new() }
1463 else { format!(", plus {}", agg_srcs.join(", ")) }));
1464 }
1465 }
1466 }
1467 }
1468
1469
1470 let tail = sql_literals_to_nql(&tail);
1471 let nql = format!("FROM {}{}{}", coll,
1472 if agg_clause.is_empty() { String::new() } else { agg_clause },
1473 if tail.is_empty() { String::new() } else { format!(" {}", tail) });
1474
1475 Ok(Stmt::Query { nql: nql.trim().to_string(), project })
1476}
1477
1478const SERVER_VERSION: &str = "15.0";
1479
1480/// The `version()` string, for the SQL engine's `version()` function.
1481pub fn version_string() -> String {
1482 full_version_string()
1483}
1484
1485fn full_version_string() -> String {
1486 format!(
1487 "PostgreSQL {} (NEDB {}) — tamper-evident, append-only, permanent \
1488 history. SELECT + INSERT/UPDATE/DELETE; an UPDATE is a new version, \
1489 so prior values stay readable with AS OF SYSTEM TIME.",
1490 SERVER_VERSION,
1491 env!("CARGO_PKG_VERSION")
1492 )
1493}
1494
1495// ── result shaping ──────────────────────────────────────────────────────────
1496
1497/// Pick the column order for a result set.
1498///
1499/// With an explicit projection, that order. Otherwise the union of keys across
1500/// the returned rows — `_`-prefixed provenance columns last, so `psql` shows
1501/// the user's own fields first and `_hash` does not push `status` off screen.
1502fn columns_for(rows: &[Value], project: &[Col]) -> Vec<Col> {
1503 if !project.is_empty() {
1504 return project.to_vec();
1505 }
1506 let mut plain: Vec<String> = vec![];
1507 let mut meta: Vec<String> = vec![];
1508 for r in rows {
1509 if let Value::Object(m) = r {
1510 for k in m.keys() {
1511 let target = if k.starts_with('_') { &mut meta } else { &mut plain };
1512 if !target.contains(k) {
1513 target.push(k.clone());
1514 }
1515 }
1516 }
1517 }
1518 // The user's own fields keep the DOCUMENT'S order -- `serde_json`'s
1519 // `preserve_order` is on crate-wide precisely so they can, and Postgres
1520 // orders `*` by column definition rather than alphabetically. Sorting them
1521 // here made `SELECT *` answer in a different column order than the SQL
1522 // evaluator did, so a client reading by POSITION got different columns
1523 // depending on a deployment flag. Only the provenance block is sorted.
1524 meta.sort();
1525 plain.extend(meta);
1526 plain.into_iter().map(|k| Col::same(&k)).collect()
1527}
1528
1529/// The Postgres type of one JSON value.
1530fn oid_of_value(v: &Value) -> Option<i32> {
1531 match v {
1532 Value::Null => None,
1533 Value::Bool(_) => Some(OID_BOOL),
1534 Value::Number(n) => Some(if n.is_i64() || n.is_u64() { OID_INT8 } else { OID_FLOAT8 }),
1535 Value::String(_) => Some(OID_TEXT),
1536 // Arrays and objects render as their JSON text.
1537 _ => Some(OID_TEXT),
1538 }
1539}
1540
1541/// Reconcile two observed types for the same column.
1542///
1543/// A relational column has one type by construction. A NEDB collection does
1544/// not: document 1 may hold `qty: 3` and document 2 `qty: "three"`. Widening
1545/// to `text` on a conflict is the only answer that can carry both, and mixed
1546/// integers and floats widen to float8 for the same reason.
1547fn unify_oid(a: i32, b: i32) -> i32 {
1548 if a == b {
1549 return a;
1550 }
1551 match (a, b) {
1552 (OID_INT8, OID_FLOAT8) | (OID_FLOAT8, OID_INT8) => OID_FLOAT8,
1553 _ => OID_TEXT,
1554 }
1555}
1556
1557/// The type of `col` across EVERY row in the result, not just the first.
1558///
1559/// Taking the first non-null value's type was a latent wrong answer: a column
1560/// holding `3` in row one and `"n/a"` in row two was advertised as `int8`, and
1561/// a client that believes the description then fails parsing `"n/a"` as an
1562/// integer — or, on the binary path, cannot be sent the value at all.
1563/// Public alias so `pgcatalog` types a column EXACTLY as the wire does.
1564///
1565/// The catalogue reporting `bigint` for a column the protocol then sends as
1566/// text would be a self-contradiction a client is entitled to trust, so both
1567/// go through this one function rather than two that agree today.
1568pub fn oid_for_column(rows: &[Value], col: &str) -> i32 {
1569 oid_for(rows, col)
1570}
1571
1572/// Did any row actually carry a non-null value for this column?
1573///
1574/// `oid_for` cannot answer this: it folds "no evidence" and "evidence, all
1575/// text" into the same `OID_TEXT`. The difference matters, because one of
1576/// those is a measurement and the other is a default standing in for one.
1577fn has_evidence(rows: &[Value], col: &str) -> bool {
1578 rows.iter().any(|r| matches!(r.get(col), Some(v) if !v.is_null()))
1579}
1580
1581fn oid_for(rows: &[Value], col: &str) -> i32 {
1582 let mut acc: Option<i32> = None;
1583 for r in rows {
1584 if let Some(o) = r.get(col).and_then(oid_of_value) {
1585 acc = Some(match acc {
1586 None => o,
1587 Some(prev) => unify_oid(prev, o),
1588 });
1589 if acc == Some(OID_TEXT) {
1590 break; // text absorbs everything; no need to look further
1591 }
1592 }
1593 }
1594 acc.unwrap_or(OID_TEXT)
1595}
1596
1597/// Render one cell in the text format Postgres clients expect for format 0.
1598fn cell(v: Option<&Value>) -> Option<String> {
1599 match v {
1600 None | Some(Value::Null) => None, // NULL on the wire
1601 Some(Value::String(s)) => Some(s.clone()),
1602 Some(Value::Bool(b)) => Some(if *b { "t".into() } else { "f".into() }),
1603 Some(other) => Some(other.to_string()),
1604 }
1605}
1606
1607/// Render one cell in binary format for the type the column was advertised as.
1608///
1609/// Needed because asyncpg asks for binary results — it is not an optimisation
1610/// there, it is the only format it requests, so without this it cannot read a
1611/// single row. Text-format clients never reach this path.
1612///
1613/// A value that does not fit the advertised type is an error rather than a
1614/// coercion. The advertised type comes from sampling stored documents, so a
1615/// mismatch means the field is genuinely heterogeneous beyond the sample, and
1616/// quietly sending a zero (or the text bytes under a binary header) would
1617/// corrupt the value in a way the client cannot detect.
1618fn cell_binary(v: Option<&Value>, oid: i32) -> Result<Option<Vec<u8>>, String> {
1619 let v = match v {
1620 None | Some(Value::Null) => return Ok(None),
1621 Some(v) => v,
1622 };
1623 let as_f64 = |n: &serde_json::Number| n.as_f64()
1624 .ok_or_else(|| "a number too large to send as float8".to_string());
1625 Ok(Some(match (oid, v) {
1626 (OID_BOOL, Value::Bool(b)) => vec![u8::from(*b)],
1627 (OID_INT2, Value::Number(n)) => {
1628 let i = n.as_i64().ok_or("not an integer")?;
1629 i16::try_from(i).map_err(|_| format!("{} does not fit in int2", i))?
1630 .to_be_bytes().to_vec()
1631 }
1632 (OID_INT4, Value::Number(n)) => {
1633 let i = n.as_i64().ok_or("not an integer")?;
1634 i32::try_from(i).map_err(|_| format!("{} does not fit in int4", i))?
1635 .to_be_bytes().to_vec()
1636 }
1637 (OID_INT8, Value::Number(n)) => {
1638 n.as_i64().ok_or("not an integer")?.to_be_bytes().to_vec()
1639 }
1640 (OID_FLOAT4, Value::Number(n)) => (as_f64(n)? as f32).to_be_bytes().to_vec(),
1641 (OID_FLOAT8, Value::Number(n)) => as_f64(n)?.to_be_bytes().to_vec(),
1642 // For the text family, binary and text are the same bytes.
1643 (OID_TEXT | OID_VARCHAR | OID_NAME | OID_UNKNOWN | OID_JSON, _) => {
1644 cell(Some(v)).unwrap_or_default().into_bytes()
1645 }
1646 // jsonb is a one-byte version header then the JSON text.
1647 (OID_JSONB, _) => {
1648 let mut b = vec![1u8];
1649 b.extend_from_slice(cell(Some(v)).unwrap_or_default().as_bytes());
1650 b
1651 }
1652 (oid, val) => {
1653 let kind = match val {
1654 Value::Bool(_) => "a boolean",
1655 Value::Number(_) => "a number",
1656 Value::String(_) => "a string",
1657 Value::Array(_) => "an array",
1658 _ => "an object",
1659 };
1660 return Err(format!(
1661 "cannot send {} in binary format as type OID {} — the field holds \
1662 more than one type across documents, so it cannot be described \
1663 by a single Postgres type. Select it with a text cast, or use a \
1664 text-format client",
1665 kind, oid
1666 ));
1667 }
1668 }))
1669}
1670
1671/// A `RowDescription`, with a per-column wire format code.
1672fn row_description_fmt(cols: &[Col], oids: &[i32], fmts: &[i16]) -> Vec<u8> {
1673 let mut m = Out::msg(b'T');
1674 m.i16(cols.len() as i16);
1675 for (i, c) in cols.iter().enumerate() {
1676 m.cstr(&c.out);
1677 m.i32(0); // table OID — unknown
1678 m.i16((i + 1) as i16); // column attribute number
1679 m.i32(oids.get(i).copied().unwrap_or(OID_TEXT));
1680 m.i16(-1); // variable length
1681 m.i32(-1); // no type modifier
1682 m.i16(fmts.get(i).copied().unwrap_or(0));
1683 }
1684 m.finish()
1685}
1686
1687fn row_description(cols: &[Col], oids: &[i32]) -> Vec<u8> {
1688 row_description_fmt(cols, oids, &[])
1689}
1690
1691fn data_row_bytes(vals: &[Option<Vec<u8>>]) -> Vec<u8> {
1692 let mut m = Out::msg(b'D');
1693 m.i16(vals.len() as i16);
1694 for v in vals {
1695 match v {
1696 None => m.i32(-1),
1697 Some(b) => {
1698 m.i32(b.len() as i32);
1699 m.bytes(b);
1700 }
1701 }
1702 }
1703 m.finish()
1704}
1705
1706fn data_row(vals: &[Option<String>]) -> Vec<u8> {
1707 let owned: Vec<Option<Vec<u8>>> =
1708 vals.iter().map(|v| v.as_ref().map(|s| s.as_bytes().to_vec())).collect();
1709 data_row_bytes(&owned)
1710}
1711
1712/// Encode just the rows: `T` followed by one `D` per row, and NO
1713/// `CommandComplete`.
1714///
1715/// Split out because a write with `RETURNING` must emit `T`/`D`* and then its
1716/// OWN tag (`INSERT 0 3`, `UPDATE 1`). The first cut called `encode_result`
1717/// there, which appends `CommandComplete("SELECT n")` — so one statement sent
1718/// TWO CommandComplete messages. That is a protocol violation, and the visible
1719/// symptom was `RETURNING` silently yielding no rows at all: the client took
1720/// the first tag as the end of the statement and discarded the description.
1721pub fn encode_rows(rows: &[Value], project: &[Col]) -> Vec<u8> {
1722 let cols = columns_for(rows, project);
1723 let oids: Vec<i32> = cols.iter().map(|c| oid_for(rows, &c.src)).collect();
1724 let mut out = row_description(&cols, &oids);
1725 for r in rows {
1726 let vals: Vec<Option<String>> = cols.iter().map(|c| cell(r.get(&c.src))).collect();
1727 out.extend_from_slice(&data_row(&vals));
1728 }
1729 out
1730}
1731
1732/// A complete SELECT response: rows plus `CommandComplete("SELECT n")`.
1733pub fn encode_result(rows: &[Value], project: &[Col]) -> Vec<u8> {
1734 let mut out = encode_rows(rows, project);
1735 out.extend_from_slice(&command_complete(&format!("SELECT {}", rows.len())));
1736 out
1737}
1738
1739// ── the extended query protocol: Parse / Bind / Describe / Execute ──────────
1740//
1741// Why this exists at all: psycopg3, asyncpg and the JDBC driver do not speak
1742// the simple query protocol for parameterised statements. Without these six
1743// messages they cannot run a single query — psycopg3 hangs waiting for a
1744// `ParseComplete`, and asyncpg refuses before it ever sends a `Bind`. "psql
1745// works" is not the same as "the drivers your evaluators use work".
1746//
1747// Two facts about real drivers shaped everything below, and both were read off
1748// a wire transcript rather than assumed:
1749//
1750// 1. psycopg3 sends parameters in a MIXED format — a `str` as OID 0 in text
1751// format, but an `int` as int2/int4/int8 in BINARY, a float as float8
1752// binary, a bool as a single binary byte. A text-only decoder gets `\x00*`
1753// where it expected `42`.
1754//
1755// 2. asyncpg declares NO parameter types in `Parse` and then asks
1756// `Describe(statement)`, encoding its arguments from whatever OIDs come
1757// back. Answering "text" for all of them does not degrade gracefully — it
1758// makes asyncpg REFUSE the call client-side ("expected str, got int").
1759//
1760// (2) is the reason `infer_param_oids` exists. NEDB is schemaless, so there is
1761// no catalogue to read a column's type out of — the only honest source of truth
1762// is the data already stored, so the type is sampled from it.
1763
1764/// Parameter/result type OIDs handled on the binary path.
1765const OID_INT2: i32 = 21;
1766const OID_INT4: i32 = 23;
1767const OID_OID: i32 = 26;
1768const OID_FLOAT4: i32 = 700;
1769const OID_VARCHAR: i32 = 1043;
1770const OID_NAME: i32 = 19;
1771const OID_UNKNOWN: i32 = 705;
1772const OID_JSON: i32 = 114;
1773const OID_JSONB: i32 = 3802;
1774
1775/// How many `$n` placeholders a statement carries, and the highest index used.
1776///
1777/// Scans outside string literals so a `'$1'` inside a value is not mistaken for
1778/// a placeholder. Dollar-quoted bodies (`$tag$…$tag$`) are not recognised —
1779/// they need a procedural language NEDB does not have.
1780fn param_count(sql: &str) -> usize {
1781 let b = sql.as_bytes();
1782 let mut i = 0usize;
1783 let mut in_s = false;
1784 let mut max = 0usize;
1785 while i < b.len() {
1786 let c = b[i];
1787 if in_s {
1788 if c == b'\'' {
1789 in_s = false;
1790 }
1791 i += 1;
1792 continue;
1793 }
1794 if c == b'\'' {
1795 in_s = true;
1796 i += 1;
1797 continue;
1798 }
1799 if c == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
1800 let mut j = i + 1;
1801 let mut n = 0usize;
1802 while j < b.len() && b[j].is_ascii_digit() {
1803 n = n * 10 + (b[j] - b'0') as usize;
1804 j += 1;
1805 }
1806 max = max.max(n);
1807 i = j;
1808 continue;
1809 }
1810 i += 1;
1811 }
1812 max
1813}
1814
1815/// The JSON-shaped type of `field` as it is actually stored, sampled from the
1816/// collection, mapped onto the nearest Postgres OID.
1817///
1818/// This is the schemaless answer to "what type is this column?". A relational
1819/// server reads its catalogue; NEDB has none, so it reads the data. Sampling a
1820/// bounded number of rows keeps a `Describe` cheap, and the first row that
1821/// actually carries the field decides — a field missing from row one but
1822/// present in row nine still types correctly.
1823fn infer_field_oid(db: Option<&Arc<Db>>, coll: &str, field: &str) -> i32 {
1824 // `_`-prefixed names are engine metadata, not stored document fields, so
1825 // they type from the engine's own contract — no sampling, and no database
1826 // handle needed.
1827 match field {
1828 "_seq" => return OID_INT8,
1829 "_id" | "_hash" | "_prev" | "_collection" | "_valid_from" | "_valid_to" => return OID_TEXT,
1830 _ => {}
1831 }
1832 // A catalogue relation types its own columns. Sampling a USER collection
1833 // named `pg_type` finds nothing and falls back to text — and asyncpg,
1834 // which declares parameter types client-side and refuses the call when
1835 // the server's answer is wrong, then rejected `WHERE oid = $1` with
1836 // "expected str, got int" before a single byte was sent.
1837 if !field.is_empty() && crate::pgcatalog::is_catalog(coll) {
1838 if let Some(rows) = crate::pgcatalog::rows(coll, db) {
1839 return oid_for(&rows, field);
1840 }
1841 }
1842 let db = match db {
1843 Some(db) => db,
1844 None => return OID_TEXT,
1845 };
1846 if coll.is_empty() || field.is_empty() {
1847 return OID_TEXT;
1848 }
1849 let rows = match crate::nql::query(db, &format!("FROM {} LIMIT {}", coll, TYPE_SAMPLE)) {
1850 Ok((rows, _)) => rows,
1851 Err(_) => return OID_TEXT,
1852 };
1853 // Unified over the sample, not taken from the first hit: a field that is a
1854 // number in one document and a string in another has to be advertised as
1855 // text or a client cannot decode every row of it.
1856 oid_for(&rows, field)
1857}
1858
1859/// The type of an aggregate output column, which no document holds.
1860///
1861/// Sampling stored documents cannot type these: `COUNT(*)` produces a column
1862/// called `count` that exists in no document, so the sampler finds nothing and
1863/// falls back to text. A text-format client papers over that, but a binary
1864/// client is then handed the digits of a number under a text header and
1865/// `COUNT(*)` comes back as the string `"2"` instead of the integer `2`.
1866///
1867/// So aggregates are typed from what the aggregate MEANS: a count is always an
1868/// integer, an average is always fractional, and min/max/sum inherit the type
1869/// of the field they were computed over.
1870/// Column names and wire types for a statement the EVALUATOR will answer.
1871///
1872/// `describe_shape` derived both by calling `translate()`, which means it
1873/// described the TRANSLATOR's output. That was right while the translator
1874/// answered; once the evaluator did, the two disagreed about the one thing
1875/// `Describe` exists to report.
1876///
1877/// They disagree on naming. `SELECT sum(total)` is column `sum_total` to the
1878/// translator and `sum` to the evaluator, so `aggregate_oid("sum", ..)` found
1879/// no `sum_` prefix, fell through to `infer_field_oid(db, coll, "sum")`, found
1880/// no stored field called `sum`, and answered `OID_TEXT`.
1881///
1882/// A text OID is not a cosmetic defect in the BINARY protocol. `Describe`
1883/// happens before `Execute`, so the client is told the column is text and
1884/// decodes the bytes that way: asyncpg received the string `'420'` where
1885/// `420` was meant, and `AS OF SYSTEM TIME $1` came back `total='66'`. The
1886/// text protocol was unaffected — it re-derives types from the rows it
1887/// actually has — which is why psycopg2's suite stayed green while asyncpg's
1888/// did not.
1889///
1890/// Typed from the PARSED SELECT rather than from a sample of the output,
1891/// because `Describe` has no rows yet. That is also why this cannot simply
1892/// reuse the row-sniffing path.
1893fn evaluator_shape(
1894 sql: &str,
1895 db: Option<&Arc<Db>>,
1896 coll: &str,
1897) -> Option<(Vec<Col>, Vec<i32>)> {
1898 let sel = crate::sqlselect::parse(sql).ok()?;
1899 // `*` expands from the rows, which Describe does not have. Declining is
1900 // honest; the caller falls back and the text path types it from the rows.
1901 if sel.items.iter().any(|i| matches!(i.expr, crate::sqlselect::Expr::Star
1902 | crate::sqlselect::Expr::QualifiedStar(_)))
1903 {
1904 return None;
1905 }
1906
1907 let mut cols: Vec<Col> = Vec::new();
1908 let mut oids: Vec<i32> = Vec::new();
1909 for item in &sel.items {
1910 let name = match &item.alias {
1911 Some(a) => a.clone(),
1912 None => match &item.expr {
1913 crate::sqlselect::Expr::Column { name, .. } => name.clone(),
1914 crate::sqlselect::Expr::Agg { name, .. } => name.to_ascii_lowercase(),
1915 crate::sqlselect::Expr::Func { name, .. } => name.to_ascii_lowercase(),
1916 // Anything else is named by a rule this function should not
1917 // try to reproduce from memory. Declining beats guessing a
1918 // name the evaluator will not use.
1919 _ => return None,
1920 },
1921 };
1922 oids.push(expr_oid(&item.expr, db, coll)?);
1923 cols.push(Col::renamed(&name, &name));
1924 }
1925 if cols.is_empty() {
1926 return None;
1927 }
1928 Some((cols, oids))
1929}
1930
1931/// The wire type of one select-list expression.
1932fn expr_oid(e: &crate::sqlselect::Expr, db: Option<&Arc<Db>>, coll: &str) -> Option<i32> {
1933 use crate::sqlselect::Expr;
1934 match e {
1935 Expr::Column { name, .. } => Some(infer_field_oid(db, coll, name)),
1936 Expr::Literal(v) => Some(oid_of_value(v).unwrap_or(OID_TEXT)),
1937 // Aggregates are `Agg`, NOT `Func`. Matching only `Func` here is what
1938 // made this whole fallback inert: `expr_oid` answered None for every
1939 // aggregate, `evaluator_shape` propagated the None, and the caller's
1940 // `unwrap_or(OID_TEXT)` shipped `sum` as text. The unit tests did not
1941 // catch it because they exercised `aggregate_oid`, which types from a
1942 // NAME; nothing typed from a parsed expression until this existed.
1943 Expr::Agg { name, args, .. } | Expr::Func { name, args } => {
1944 let f = name.to_ascii_lowercase();
1945 match f.as_str() {
1946 // COUNT is a count whatever it counts.
1947 "count" => Some(OID_INT8),
1948 // An average is fractional even over integers — the case the
1949 // translator also special-cased.
1950 "avg" => Some(OID_FLOAT8),
1951 // SUM/MIN/MAX inherit the type they range over, so the
1952 // argument has to be resolved rather than assumed numeric.
1953 "sum" | "min" | "max" => match args.first() {
1954 Some(Expr::Column { name, .. }) => match infer_field_oid(db, coll, name) {
1955 OID_INT8 => Some(OID_INT8),
1956 OID_FLOAT8 => Some(OID_FLOAT8),
1957 other => Some(other),
1958 },
1959 _ => None,
1960 },
1961 _ => None,
1962 }
1963 }
1964 _ => None,
1965 }
1966}
1967
1968fn aggregate_oid(src: &str, db: Option<&Arc<Db>>, coll: &str) -> Option<i32> {
1969 if src == "count" {
1970 return Some(OID_INT8);
1971 }
1972 for (prefix, fixed) in [
1973 ("count_", Some(OID_INT8)),
1974 ("avg_", Some(OID_FLOAT8)),
1975 ("sum_", None),
1976 ("min_", None),
1977 ("max_", None),
1978 ] {
1979 if let Some(field) = src.strip_prefix(prefix) {
1980 return Some(match fixed {
1981 Some(oid) => oid,
1982 // SUM/MIN/MAX of an integer field is an integer; of a
1983 // fractional field, fractional.
1984 None => match infer_field_oid(db, coll, field) {
1985 OID_INT8 => OID_INT8,
1986 OID_FLOAT8 => OID_FLOAT8,
1987 // Summing or ordering a non-numeric field is not
1988 // meaningful; let the row-derived type answer.
1989 other => other,
1990 },
1991 });
1992 }
1993 }
1994 None
1995}
1996
1997/// How many documents to sample when typing a column.
1998///
1999/// Bounded so a `Describe` stays cheap. It is a sample, so a field that only
2000/// turns heterogeneous outside it can still surprise us — which is exactly why
2001/// `cell_binary` refuses a mismatch loudly instead of coercing.
2002const TYPE_SAMPLE: usize = 200;
2003
2004/// The collection a statement reads from or writes to, for type sampling.
2005fn stmt_collection(sql: &str) -> String {
2006 let s = normalise(sql);
2007 let up = s.to_uppercase();
2008 let after = if let Some(at) = find_kw(&up, "FROM") {
2009 &s[at + 4..]
2010 } else if let Some(rest) = strip_prefix_ci(&s, "UPDATE") {
2011 return rest
2012 .split_whitespace()
2013 .next()
2014 .unwrap_or("")
2015 .rsplit('.')
2016 .next()
2017 .unwrap_or("")
2018 .trim_matches('"')
2019 .to_string();
2020 } else if let Some(rest) = strip_prefix_ci(&s, "INSERT INTO") {
2021 return rest
2022 .split(|c: char| c.is_whitespace() || c == '(')
2023 .find(|t| !t.is_empty())
2024 .unwrap_or("")
2025 .rsplit('.')
2026 .next()
2027 .unwrap_or("")
2028 .trim_matches('"')
2029 .to_string();
2030 } else {
2031 return String::new();
2032 };
2033 after
2034 .trim()
2035 .split(|c: char| c.is_whitespace())
2036 .find(|t| !t.is_empty())
2037 .unwrap_or("")
2038 .rsplit('.')
2039 .next()
2040 .unwrap_or("")
2041 .trim_matches('"')
2042 .to_string()
2043}
2044
2045/// Which document field each `$n` is being compared against.
2046///
2047/// Three shapes cover essentially all driver-generated SQL:
2048/// `WHERE qty > $1` → the identifier immediately left of the operator
2049/// `SET status = $1` → same shape, inside the SET list
2050/// `INSERT INTO t (a,b) VALUES ($1,$2)` → positional against the column list
2051///
2052/// Anything it cannot read returns `None`, which types as `text`. Guessing
2053/// wrong here would make a driver encode a value the engine then fails to
2054/// match, so an unknown is left unknown on purpose.
2055fn param_fields(sql: &str, n_params: usize) -> Vec<Option<String>> {
2056 let s = normalise(sql);
2057 let mut out = vec![None; n_params];
2058
2059 // The INSERT column list maps positionally, which is more reliable than
2060 // scanning leftwards through a VALUES tuple.
2061 let up = s.to_uppercase();
2062 if up.starts_with("INSERT") {
2063 if let (Some(open), Some(vals_at)) = (s.find('('), find_kw(&up, "VALUES")) {
2064 if open < vals_at {
2065 if let Some(close) = s[open..vals_at].rfind(')') {
2066 let cols: Vec<String> = split_top(&s[open + 1..open + close], ',')
2067 .into_iter()
2068 .map(|c| c.trim().trim_matches('"').to_string())
2069 .collect();
2070 // `$1` is the first placeholder in the first tuple, and so on.
2071 let tail = &s[vals_at..];
2072 let mut seen = 0usize;
2073 let b = tail.as_bytes();
2074 let mut i = 0usize;
2075 let mut in_s = false;
2076 while i < b.len() {
2077 if in_s {
2078 if b[i] == b'\'' { in_s = false; }
2079 i += 1;
2080 continue;
2081 }
2082 if b[i] == b'\'' { in_s = true; i += 1; continue; }
2083 if b[i] == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
2084 let mut j = i + 1;
2085 let mut num = 0usize;
2086 while j < b.len() && b[j].is_ascii_digit() {
2087 num = num * 10 + (b[j] - b'0') as usize;
2088 j += 1;
2089 }
2090 if num >= 1 && num <= n_params {
2091 if let Some(c) = cols.get(seen % cols.len().max(1)) {
2092 out[num - 1] = Some(c.clone());
2093 }
2094 }
2095 seen += 1;
2096 i = j;
2097 continue;
2098 }
2099 i += 1;
2100 }
2101 return out;
2102 }
2103 }
2104 }
2105 }
2106
2107 // Otherwise: for each `$n`, walk left past the operator to the identifier.
2108 let b = s.as_bytes();
2109 let mut i = 0usize;
2110 let mut in_s = false;
2111 while i < b.len() {
2112 if in_s {
2113 if b[i] == b'\'' { in_s = false; }
2114 i += 1;
2115 continue;
2116 }
2117 if b[i] == b'\'' { in_s = true; i += 1; continue; }
2118 if b[i] == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
2119 let mut j = i + 1;
2120 let mut num = 0usize;
2121 while j < b.len() && b[j].is_ascii_digit() {
2122 num = num * 10 + (b[j] - b'0') as usize;
2123 j += 1;
2124 }
2125 if num >= 1 && num <= n_params {
2126 let left = &s[..i];
2127 // Skip the operator characters and whitespace sitting between
2128 // the identifier and the placeholder.
2129 let trimmed = left.trim_end_matches(|c: char| {
2130 c.is_whitespace() || "=<>!+-*/%(,".contains(c)
2131 });
2132 // A word operator (`LIKE`, `IN`, `BETWEEN`, `AND`) also sits
2133 // between them; step over it to reach the real identifier.
2134 let mut tok = trimmed
2135 .rsplit(|c: char| c.is_whitespace() || c == '(' || c == ',')
2136 .find(|t| !t.is_empty())
2137 .unwrap_or("")
2138 .trim_matches('"');
2139 let mut before = trimmed;
2140 for _ in 0..4 {
2141 let upper_tok = tok.to_uppercase();
2142 // `BETWEEN $1 AND $2` puts BOTH a word operator and an
2143 // earlier placeholder between `$2` and the column it
2144 // constrains, so a placeholder has to be stepped over too —
2145 // otherwise the upper bound of every range query types as
2146 // text while the lower bound types correctly.
2147 if upper_tok.starts_with('$')
2148 || matches!(upper_tok.as_str(),
2149 "LIKE" | "ILIKE" | "IN" | "BETWEEN" | "AND" | "OR" | "NOT" | "IS") {
2150 before = before[..before.len() - tok.len()].trim_end_matches(|c: char| {
2151 c.is_whitespace() || "=<>!(,".contains(c)
2152 });
2153 tok = before
2154 .rsplit(|c: char| c.is_whitespace() || c == '(' || c == ',')
2155 .find(|t| !t.is_empty())
2156 .unwrap_or("")
2157 .trim_matches('"');
2158 } else {
2159 break;
2160 }
2161 }
2162 if !tok.is_empty()
2163 && tok.chars().all(|c| c.is_alphanumeric() || c == '_' || c == '.')
2164 && !tok.chars().next().map(|c| c.is_ascii_digit()).unwrap_or(true)
2165 {
2166 out[num - 1] = Some(tok.rsplit('.').next().unwrap_or(tok).to_string());
2167 }
2168 }
2169 i = j;
2170 continue;
2171 }
2172 i += 1;
2173 }
2174 out
2175}
2176
2177/// The type of a placeholder sitting in a CLAUSE position rather than beside a
2178/// column.
2179///
2180/// `AS OF SYSTEM TIME $1` has no column to sample — the token to its left is
2181/// the word `TIME`. Its type comes from the grammar instead, which is both
2182/// cheaper and more certain than any inference: a system-time bound is a
2183/// sequence number, a valid-time bound is a date string, and a page bound is an
2184/// integer. Without this, a parameterised time-travel query typed as text and
2185/// asyncpg refused to send the integer at all.
2186fn clause_param_oids(sql: &str, n_params: usize) -> Vec<Option<i32>> {
2187 let s = normalise(sql);
2188 let mut out = vec![None; n_params];
2189 let b = s.as_bytes();
2190 let mut i = 0usize;
2191 let mut in_s = false;
2192 while i < b.len() {
2193 if in_s {
2194 if b[i] == b'\'' { in_s = false; }
2195 i += 1;
2196 continue;
2197 }
2198 if b[i] == b'\'' { in_s = true; i += 1; continue; }
2199 if b[i] == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
2200 let mut j = i + 1;
2201 let mut num = 0usize;
2202 while j < b.len() && b[j].is_ascii_digit() {
2203 num = num * 10 + (b[j] - b'0') as usize;
2204 j += 1;
2205 }
2206 if num >= 1 && num <= n_params {
2207 let left = s[..i].trim_end().to_uppercase();
2208 // VALID AS OF is checked FIRST: it ends with "AS OF" too, and
2209 // its argument is a DATE STRING, not a sequence number.
2210 out[num - 1] = if left.ends_with("VALID AS OF") {
2211 Some(OID_TEXT)
2212 } else if left.ends_with("AS OF SYSTEM TIME")
2213 || left.ends_with("FOR SYSTEM_TIME AS OF")
2214 || left.ends_with("AS OF")
2215 || left.ends_with("LIMIT")
2216 || left.ends_with("OFFSET")
2217 {
2218 Some(OID_INT8)
2219 } else {
2220 None
2221 };
2222 }
2223 i = j;
2224 continue;
2225 }
2226 i += 1;
2227 }
2228 out
2229}
2230
2231/// The OIDs to advertise for `$1..$n`, sampled from stored data.
2232///
2233/// `declared` is what the client itself put in `Parse`. A client that states a
2234/// type is believed — it is about to encode its arguments that way, and second
2235///-guessing it would break the decode. Only the unspecified slots are inferred.
2236fn infer_param_oids(sql: &str, declared: &[i32], db: Option<&Arc<Db>>) -> Vec<i32> {
2237 let n = param_count(sql).max(declared.len());
2238 if n == 0 {
2239 return vec![];
2240 }
2241 let coll = stmt_collection(sql);
2242 let fields = param_fields(sql, n);
2243 let clauses = clause_param_oids(sql, n);
2244 (0..n)
2245 .map(|i| match declared.get(i) {
2246 Some(&oid) if oid != 0 => oid,
2247 // A clause position knows its own type from the grammar, so it
2248 // outranks sampling a column that is not even there.
2249 _ => match clauses[i] {
2250 Some(oid) => oid,
2251 None => match &fields[i] {
2252 Some(f) => infer_field_oid(db, &coll, f),
2253 None => OID_TEXT,
2254 },
2255 },
2256 })
2257 .collect()
2258}
2259
2260/// Decode one bound parameter into the SQL literal text to splice into the
2261/// statement.
2262///
2263/// `None` means SQL NULL. Format 1 is binary — see the module note on psycopg3
2264/// sending small integers as int2.
2265fn decode_param(raw: Option<&[u8]>, oid: i32, format: i16) -> Result<Option<String>, String> {
2266 let bytes = match raw {
2267 None => return Ok(None),
2268 Some(b) => b,
2269 };
2270 let quote = |s: &str| format!("'{}'", s.replace('\'', "''"));
2271
2272 if format == 0 {
2273 let s = String::from_utf8_lossy(bytes).to_string();
2274 return Ok(Some(match oid {
2275 OID_BOOL => {
2276 let t = matches!(s.as_str(), "t" | "true" | "TRUE" | "1" | "yes" | "on");
2277 if t { "TRUE".into() } else { "FALSE".into() }
2278 }
2279 OID_INT2 | OID_INT4 | OID_INT8 | OID_OID | OID_FLOAT4 | OID_FLOAT8 => {
2280 // Validate rather than trust: an unparseable "number" spliced
2281 // in bare would become a bare identifier in the NQL text and
2282 // produce a baffling error far from its cause.
2283 if s.parse::<f64>().is_ok() { s } else { quote(&s) }
2284 }
2285 // OID 0 with text format is psycopg3's `str`. Confirmed on the
2286 // wire: it declares a real numeric OID whenever the value is a
2287 // number, so an unspecified text parameter is genuinely a string
2288 // and quoting it is right rather than a guess.
2289 _ => quote(&s),
2290 }));
2291 }
2292 if format != 1 {
2293 return Err(format!("unsupported parameter format code {}", format));
2294 }
2295
2296 // ── binary ──────────────────────────────────────────────────────────────
2297 let need = |n: usize| -> Result<(), String> {
2298 if bytes.len() == n {
2299 Ok(())
2300 } else {
2301 Err(format!(
2302 "binary parameter of type OID {} should be {} bytes, got {}",
2303 oid, n, bytes.len()
2304 ))
2305 }
2306 };
2307 Ok(Some(match oid {
2308 OID_BOOL => {
2309 need(1)?;
2310 if bytes[0] != 0 { "TRUE".into() } else { "FALSE".into() }
2311 }
2312 OID_INT2 => {
2313 need(2)?;
2314 i16::from_be_bytes([bytes[0], bytes[1]]).to_string()
2315 }
2316 OID_INT4 => {
2317 need(4)?;
2318 i32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]).to_string()
2319 }
2320 OID_OID => {
2321 need(4)?;
2322 u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]).to_string()
2323 }
2324 OID_INT8 => {
2325 need(8)?;
2326 i64::from_be_bytes(bytes[..8].try_into().unwrap()).to_string()
2327 }
2328 OID_FLOAT4 => {
2329 need(4)?;
2330 let f = f32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]);
2331 fmt_float(f as f64)
2332 }
2333 OID_FLOAT8 => {
2334 need(8)?;
2335 fmt_float(f64::from_be_bytes(bytes[..8].try_into().unwrap()))
2336 }
2337 OID_TEXT | OID_VARCHAR | OID_NAME | OID_UNKNOWN | OID_JSON | 0 => {
2338 quote(&String::from_utf8_lossy(bytes))
2339 }
2340 OID_JSONB => {
2341 // jsonb binary is a 1-byte version header followed by the JSON text.
2342 let body = if bytes.first() == Some(&1) { &bytes[1..] } else { bytes };
2343 quote(&String::from_utf8_lossy(body))
2344 }
2345 other => {
2346 return Err(format!(
2347 "parameter type OID {} is not supported in binary format — \
2348 the supported set is bool, int2/int4/int8, float4/float8, \
2349 text/varchar/json/jsonb. Send it as text, or cast it in the \
2350 statement",
2351 other
2352 ))
2353 }
2354 }))
2355}
2356
2357/// Render a float without Rust's `inf`/`NaN` spellings leaking into SQL text.
2358fn fmt_float(f: f64) -> String {
2359 if f.is_nan() {
2360 "'NaN'".into()
2361 } else if f.is_infinite() {
2362 if f > 0.0 { "'Infinity'".into() } else { "'-Infinity'".into() }
2363 } else if f.fract() == 0.0 && f.abs() < 1e15 {
2364 format!("{:.0}", f)
2365 } else {
2366 f.to_string()
2367 }
2368}
2369
2370/// Splice decoded parameters into the statement text.
2371///
2372/// Textual substitution, deliberately: the whole SQL surface is already a text
2373/// translation into NQL, so one representation is simpler and cannot disagree
2374/// with itself. Every value arrives already rendered as a SQL literal by
2375/// `decode_param`, with embedded quotes doubled, so a parameter cannot break
2376/// out of its literal and alter the statement's shape.
2377fn substitute_params(sql: &str, params: &[Option<String>]) -> Result<String, String> {
2378 let b = sql.as_bytes();
2379 let mut out = String::with_capacity(sql.len() + 16);
2380 let mut i = 0usize;
2381 let mut in_s = false;
2382 while i < b.len() {
2383 let c = b[i];
2384 if in_s {
2385 out.push(c as char);
2386 if c == b'\'' { in_s = false; }
2387 i += 1;
2388 continue;
2389 }
2390 if c == b'\'' {
2391 in_s = true;
2392 out.push('\'');
2393 i += 1;
2394 continue;
2395 }
2396 if c == b'$' && i + 1 < b.len() && b[i + 1].is_ascii_digit() {
2397 let mut j = i + 1;
2398 let mut n = 0usize;
2399 while j < b.len() && b[j].is_ascii_digit() {
2400 n = n * 10 + (b[j] - b'0') as usize;
2401 j += 1;
2402 }
2403 match params.get(n.wrapping_sub(1)) {
2404 Some(Some(lit)) => out.push_str(lit),
2405 Some(None) => out.push_str("NULL"),
2406 None => {
2407 return Err(format!(
2408 "bind message supplies {} parameter(s) but the statement uses ${}",
2409 params.len(), n
2410 ))
2411 }
2412 }
2413 i = j;
2414 continue;
2415 }
2416 out.push(c as char);
2417 i += 1;
2418 }
2419 Ok(out)
2420}
2421
2422/// A parsed statement, held for the life of the connection (or until `Close`).
2423struct Prepared {
2424 sql: String,
2425 /// OIDs advertised for `$1..$n` — what `ParameterDescription` reports and
2426 /// what `Bind` values are decoded as.
2427 param_oids: Vec<i32>,
2428 /// The advertised output shape, computed on demand and then reused.
2429 ///
2430 /// Lazy because working it out samples stored documents, and a text-format
2431 /// client that never sends `Describe(statement)` should not pay for a scan
2432 /// on every `Parse` — psycopg3 parses once per query.
2433 ///
2434 /// `Some(None)` means "computed, and this statement returns no rows".
2435 out_shape: Option<Option<(Vec<Col>, Vec<i32>)>>,
2436}
2437
2438/// The output columns and types a statement advertises, computed once.
2439fn prepared_shape<'a>(
2440 p: &'a mut Prepared,
2441 db: Option<&Arc<Db>>,
2442) -> &'a Option<(Vec<Col>, Vec<i32>)> {
2443 if p.out_shape.is_none() {
2444 p.out_shape = Some(describe_shape(&p.sql, db, p.param_oids.len()));
2445 }
2446 p.out_shape.as_ref().expect("just filled")
2447}
2448
2449/// A bound statement: fully substituted SQL plus, once run, its result.
2450struct Portal {
2451 sql: String,
2452 /// Filled by the first `Describe` or `Execute` and reused afterwards.
2453 ///
2454 /// Executing once and streaming from the buffer is what makes a suspended
2455 /// portal safe: a second `Execute` on a partially-drained `INSERT` must
2456 /// continue the row stream, not perform the insert again.
2457 result: Option<PortalResult>,
2458 /// The output shape, frozen at the first `Describe`/`Execute`.
2459 ///
2460 /// A schemaless store derives `SELECT *`'s columns from the rows it found,
2461 /// which would let a `Describe` and a later `Execute` disagree about the
2462 /// column count — and a driver that was told three fields and handed two
2463 /// mis-decodes the row rather than failing loudly. Freezing the shape and
2464 /// projecting every row onto it makes the result set rectangular, as SQL
2465 /// promises. The simple protocol keeps the dynamic behaviour, where there
2466 /// is no `Describe` to contradict.
2467 frozen: Option<Vec<Col>>,
2468 /// Result-column format codes requested by `Bind`. Empty = all text.
2469 formats: Vec<i16>,
2470 /// The shape this portal's statement advertised, carried over from the
2471 /// prepared statement when any column is to be sent in BINARY.
2472 ///
2473 /// It has to be the ADVERTISED shape rather than one derived from the rows
2474 /// in hand: asyncpg built its decoders from `Describe`, so re-deriving a
2475 /// different type here would hand it bytes it cannot read.
2476 declared: Option<(Vec<Col>, Vec<i32>)>,
2477}
2478
2479impl Portal {
2480 /// The format code for column `i`, following the protocol's shorthands:
2481 /// no codes means all-text, one code applies to every column.
2482 fn format_of(&self, i: usize) -> i16 {
2483 match self.formats.len() {
2484 0 => 0,
2485 1 => self.formats[0],
2486 _ => self.formats.get(i).copied().unwrap_or(0),
2487 }
2488 }
2489 /// The columns and types to advertise and encode with.
2490 fn shape(&self, r: &PortalResult) -> (Vec<Col>, Vec<i32>) {
2491 match &self.declared {
2492 Some((cols, oids)) if self.formats.iter().any(|f| *f == 1) => {
2493 (cols.clone(), oids.clone())
2494 }
2495 _ => {
2496 let cols = columns_for(&r.rows, &r.project);
2497 let oids = cols.iter().map(|c| oid_for(&r.rows, &c.src)).collect();
2498 (cols, oids)
2499 }
2500 }
2501 }
2502}
2503
2504struct PortalResult {
2505 rows: Vec<Value>,
2506 project: Vec<Col>,
2507 has_rows: bool,
2508 tag: String,
2509 tag_counts_rows: bool,
2510 /// How many rows have gone out across all `Execute`s on this portal.
2511 sent: usize,
2512}
2513
2514fn parse_complete() -> Vec<u8> { Out::msg(b'1').finish() }
2515fn bind_complete() -> Vec<u8> { Out::msg(b'2').finish() }
2516fn close_complete() -> Vec<u8> { Out::msg(b'3').finish() }
2517fn no_data() -> Vec<u8> { Out::msg(b'n').finish() }
2518fn portal_suspended() -> Vec<u8> { Out::msg(b's').finish() }
2519
2520fn parameter_description(oids: &[i32]) -> Vec<u8> {
2521 let mut m = Out::msg(b't');
2522 m.i16(oids.len() as i16);
2523 for o in oids {
2524 m.i32(*o);
2525 }
2526 m.finish()
2527}
2528
2529/// Split a NUL-terminated string off the front of a message body.
2530fn take_cstr(body: &[u8], at: &mut usize) -> String {
2531 let start = *at;
2532 while *at < body.len() && body[*at] != 0 {
2533 *at += 1;
2534 }
2535 let s = String::from_utf8_lossy(&body[start..*at]).to_string();
2536 if *at < body.len() {
2537 *at += 1; // step over the NUL
2538 }
2539 s
2540}
2541
2542fn take_i16(body: &[u8], at: &mut usize) -> Result<i16, String> {
2543 if *at + 2 > body.len() {
2544 return Err("truncated message".into());
2545 }
2546 let v = i16::from_be_bytes([body[*at], body[*at + 1]]);
2547 *at += 2;
2548 Ok(v)
2549}
2550
2551fn take_i32(body: &[u8], at: &mut usize) -> Result<i32, String> {
2552 if *at + 4 > body.len() {
2553 return Err("truncated message".into());
2554 }
2555 let v = i32::from_be_bytes([body[*at], body[*at + 1], body[*at + 2], body[*at + 3]]);
2556 *at += 4;
2557 Ok(v)
2558}
2559
2560/// The field names a collection actually holds, sampled from stored documents.
2561///
2562/// The answer to `SELECT *` on a store with no schema. Sorted, because
2563/// `serde_json`'s map is ordered and both this and the row encoder must agree
2564/// on column order or the values land under the wrong headings.
2565fn sample_columns(db: Option<&Arc<Db>>, coll: &str) -> Vec<Col> {
2566 let db = match db {
2567 Some(db) => db,
2568 None => return vec![],
2569 };
2570 let rows = match crate::nql::query(db, &format!("FROM {} LIMIT 25", coll)) {
2571 Ok((rows, _)) => rows,
2572 Err(_) => return vec![],
2573 };
2574 let mut names: Vec<String> = vec![];
2575 for r in &rows {
2576 if let Value::Object(m) = r {
2577 for k in m.keys() {
2578 if !names.iter().any(|n| n == k) {
2579 names.push(k.clone());
2580 }
2581 }
2582 }
2583 }
2584 names.sort();
2585 names.iter().map(|n| Col::same(n)).collect()
2586}
2587
2588/// The result shape of a statement, worked out WITHOUT running it.
2589///
2590/// Needed for `Describe(statement)`, which arrives before any `Bind` — asyncpg
2591/// builds its row decoders from the answer. Only the select list is read off
2592/// the result; nothing touches storage except the type sampling.
2593///
2594/// Returns `None` when the statement returns no rows at all (`NoData`).
2595fn describe_shape(
2596 sql: &str,
2597 db: Option<&Arc<Db>>,
2598 n_params: usize,
2599) -> Option<(Vec<Col>, Vec<i32>)> {
2600 let probe = probe_sql(sql, n_params);
2601
2602 // The SQL evaluator describes its own output. It has to: `translate`
2603 // cannot parse a catalogue join at all, so without this a `Describe`
2604 // answered `NoData` — and a client told a SELECT has no output never
2605 // reads its rows.
2606 //
2607 // The probe is EXECUTED here, which is affordable precisely because this
2608 // path only serves catalogue relations and relation-free select lists.
2609 // Column types come from the values it actually produced, unified across
2610 // the rows by the same `oid_for` every other path uses — so a column
2611 // advertised `int8` is one the wire really encodes as int8.
2612 let coll = stmt_collection(sql);
2613
2614 if sql_engine_owns(&probe) {
2615 if let Ok(Some((done, _))) = try_catalog_select(&probe, db) {
2616 if done.project.is_empty() {
2617 return None;
2618 }
2619 // Sniffing the probe's OUTPUT is only sound when the probe
2620 // produced output. It frequently does not, and the reason is
2621 // structural rather than unlucky: `probe_sql` substitutes `0` for
2622 // every parameter, so `... WHERE region = $1` becomes
2623 // `... WHERE region = 0`, matches nothing, and hands this line an
2624 // empty `rows`. `oid_for` then finds no evidence and returns its
2625 // `unwrap_or(OID_TEXT)` default.
2626 //
2627 // In the BINARY protocol that default is not a shrug, it is a
2628 // wrong answer the client cannot recover from: `Describe`
2629 // precedes `Execute`, so asyncpg was told `sum` was text and
2630 // decoded 420 as the string "420". The text protocol re-derives
2631 // types from the rows it really got, which is why psycopg2's
2632 // suite stayed green throughout and only asyncpg's went red.
2633 //
2634 // So: evidence where there is evidence, and static inference from
2635 // the STORED data where there is none — which is what the
2636 // translator's `infer_field_oid` was doing all along.
2637 let fallback = evaluator_shape(&probe, db, &coll);
2638 let oids: Vec<i32> = done
2639 .project
2640 .iter()
2641 .enumerate()
2642 .map(|(i, c)| {
2643 let seen = has_evidence(&done.rows, &c.src);
2644 if seen {
2645 oid_for(&done.rows, &c.src)
2646 } else {
2647 fallback
2648 .as_ref()
2649 .and_then(|(_, o)| o.get(i).copied())
2650 .unwrap_or(OID_TEXT)
2651 }
2652 })
2653 .collect();
2654 return Some((done.project, oids));
2655 }
2656 }
2657
2658
2659 // Ask the engine that will actually answer. Falls through when the
2660 // evaluator declines to describe itself — `SELECT *` expands from rows
2661 // Describe has not read — and the translator's shape is then the better
2662 // of the two available answers rather than the right one.
2663 if sql_engine_owns(&probe) {
2664 if let Some(shape) = evaluator_shape(&probe, db, &coll) {
2665 return Some(shape);
2666 }
2667 }
2668
2669 let stmt = translate(&probe).ok()?;
2670
2671 let cols = match stmt {
2672 Stmt::Ok(_) => return None,
2673 Stmt::Canned { cols, .. } => cols.iter().map(|c| Col::same(c)).collect(),
2674 Stmt::Query { project, .. } => {
2675 if project.is_empty() { sample_columns(db, &coll) } else { project }
2676 }
2677 Stmt::Insert { returning, .. } | Stmt::Update { returning, .. } | Stmt::Delete { returning, .. } => {
2678 if !wants_returning(sql) {
2679 return None;
2680 }
2681 if returning.is_empty() { sample_columns(db, &coll) } else { returning }
2682 }
2683 };
2684 if cols.is_empty() {
2685 // Nothing could be determined. `NoData` is a lie for a SELECT, but a
2686 // RowDescription with zero columns is a worse one — it tells the client
2687 // the query definitively has no output.
2688 return None;
2689 }
2690 let oids = cols
2691 .iter()
2692 .map(|c| {
2693 aggregate_oid(&c.src, db, &coll)
2694 .unwrap_or_else(|| infer_field_oid(db, &coll, &c.src))
2695 })
2696 .collect();
2697 Some((cols, oids))
2698}
2699
2700/// A parse-only stand-in for a parameterised statement.
2701///
2702/// Substituting `NULL` was the obvious choice and the wrong one: a clause that
2703/// validates its argument rejects it, so `AS OF SYSTEM TIME $1` failed at
2704/// `Parse` — before the client ever bound a real sequence number. `0` parses
2705/// everywhere a literal can appear, and since only the SELECT list is read back
2706/// out, the stub's value never reaches an answer.
2707fn probe_sql(sql: &str, n_params: usize) -> String {
2708 let stub: Vec<Option<String>> = vec![Some("0".to_string()); n_params];
2709 substitute_params(sql, &stub).unwrap_or_else(|_| sql.to_string())
2710}
2711
2712/// Run a portal's statement if it has not run yet, then report its shape.
2713fn ensure_executed(
2714 portal: &mut Portal,
2715 db_name: &str,
2716 db: Option<&Arc<Db>>,
2717 read_only: bool,
2718) -> Result<(), Vec<u8>> {
2719 if portal.result.is_some() {
2720 return Ok(());
2721 }
2722 let ex = execute_stmt(&portal.sql, db_name, db, read_only)?;
2723 // Freeze the output shape on first sight so `Describe` and every later
2724 // `Execute` describe the same rectangle.
2725 let project = if let Some(f) = &portal.frozen {
2726 f.clone()
2727 } else {
2728 let p = if ex.project.is_empty() {
2729 columns_for(&ex.rows, &[])
2730 } else {
2731 ex.project.clone()
2732 };
2733 portal.frozen = Some(p.clone());
2734 p
2735 };
2736 portal.result = Some(PortalResult {
2737 rows: ex.rows,
2738 project,
2739 has_rows: ex.has_rows,
2740 tag: ex.tag,
2741 tag_counts_rows: ex.tag_counts_rows,
2742 sent: 0,
2743 });
2744 Ok(())
2745}
2746
2747// ── connection handling ─────────────────────────────────────────────────────
2748
2749async fn read_exact(sock: &mut TcpStream, n: usize) -> std::io::Result<Vec<u8>> {
2750 let mut buf = vec![0u8; n];
2751 sock.read_exact(&mut buf).await?;
2752 Ok(buf)
2753}
2754
2755async fn read_i32(sock: &mut TcpStream) -> std::io::Result<i32> {
2756 let b = read_exact(sock, 4).await?;
2757 Ok(i32::from_be_bytes([b[0], b[1], b[2], b[3]]))
2758}
2759
2760fn parse_startup_params(body: &[u8]) -> HashMap<String, String> {
2761 let mut out = HashMap::new();
2762 let mut parts = body.split(|b| *b == 0).map(|s| String::from_utf8_lossy(s).to_string());
2763 while let (Some(k), Some(v)) = (parts.next(), parts.next()) {
2764 if k.is_empty() {
2765 break;
2766 }
2767 out.insert(k, v);
2768 }
2769 out
2770}
2771
2772/// Serve one client connection to completion.
2773async fn handle(mut sock: TcpStream, resolver: Arc<dyn DbResolver>, read_only: bool) -> std::io::Result<()> {
2774 // ── startup, including the SSL negotiation clients try first ────────────
2775 let params = loop {
2776 let len = read_i32(&mut sock).await?;
2777 if len < 8 || len > 1 << 20 {
2778 return Ok(()); // nonsense framing — drop the connection
2779 }
2780 let code = read_i32(&mut sock).await?;
2781 let body = read_exact(&mut sock, (len - 8) as usize).await?;
2782 match code {
2783 SSL_REQUEST | GSS_REQUEST => {
2784 // Decline and let the client retry in the clear.
2785 sock.write_all(b"N").await?;
2786 continue;
2787 }
2788 CANCEL_REQUEST => return Ok(()), // nothing cancellable: reads are synchronous
2789 PROTO_V3 => break parse_startup_params(&body),
2790 other => {
2791 let major = other >> 16;
2792 sock.write_all(&err_msg(
2793 "0A000",
2794 &format!("unsupported frontend protocol {}.{} — this endpoint speaks 3.0",
2795 major, other & 0xffff),
2796 )).await?;
2797 return Ok(());
2798 }
2799 }
2800 };
2801
2802 let db_name = params.get("database").cloned().unwrap_or_default();
2803
2804 // Resolve the database ONCE, here, on a blocking thread.
2805 //
2806 // A Postgres connection is bound to one database for its whole life, so
2807 // per-connection resolution is both correct and simpler than resolving per
2808 // statement — and it keeps the lock acquisition off the async worker.
2809 let resolved: Option<Arc<Db>> = {
2810 let r = Arc::clone(&resolver);
2811 let name = db_name.clone();
2812 tokio::task::spawn_blocking(move || r.resolve(&name))
2813 .await
2814 .unwrap_or(None)
2815 };
2816
2817 // ── auth: mirror the HTTP surface ───────────────────────────────────────
2818 if let Some(expected) = resolver.token() {
2819 // AuthenticationCleartextPassword (3)
2820 let mut m = Out::msg(b'R');
2821 m.i32(3);
2822 sock.write_all(&m.finish()).await?;
2823
2824 let tag = read_exact(&mut sock, 1).await?;
2825 if tag[0] != b'p' {
2826 sock.write_all(&err_msg("28000", "expected a password message")).await?;
2827 return Ok(());
2828 }
2829 let len = read_i32(&mut sock).await?;
2830 if len < 4 || len > 1 << 16 {
2831 return Ok(());
2832 }
2833 let body = read_exact(&mut sock, (len - 4) as usize).await?;
2834 let supplied = String::from_utf8_lossy(&body).trim_end_matches('\0').to_string();
2835 // Constant-time-ish: compare lengths and bytes without early return.
2836 let ok = supplied.len() == expected.len()
2837 && supplied.bytes().zip(expected.bytes()).fold(0u8, |a, (x, y)| a | (x ^ y)) == 0;
2838 if !ok {
2839 sock.write_all(&err_msg("28P01", "password authentication failed")).await?;
2840 return Ok(());
2841 }
2842 }
2843
2844 let mut m = Out::msg(b'R');
2845 m.i32(0); // AuthenticationOk
2846 sock.write_all(&m.finish()).await?;
2847
2848 for (k, v) in [
2849 ("server_version", SERVER_VERSION),
2850 ("server_encoding", "UTF8"),
2851 ("client_encoding", "UTF8"),
2852 ("DateStyle", "ISO, MDY"),
2853 ("integer_datetimes", "on"),
2854 ("standard_conforming_strings", "on"),
2855 ("application_name", "nedbd"),
2856 ] {
2857 let mut p = Out::msg(b'S');
2858 p.cstr(k);
2859 p.cstr(v);
2860 sock.write_all(&p.finish()).await?;
2861 }
2862 let mut k = Out::msg(b'K');
2863 k.i32(std::process::id() as i32);
2864 k.i32(0);
2865 sock.write_all(&k.finish()).await?;
2866 sock.write_all(&ready()).await?;
2867
2868 // ── message loop ────────────────────────────────────────────────────────
2869 //
2870 // Prepared statements and portals live for the connection. `""` is the
2871 // unnamed statement/portal, which every driver reuses constantly — it is an
2872 // ordinary entry in the map rather than a special case.
2873 let mut prepared: HashMap<String, Prepared> = HashMap::new();
2874 let mut portals: HashMap<String, Portal> = HashMap::new();
2875 // After an error inside an extended-protocol sequence, everything up to the
2876 // next `Sync` is discarded. Skipping this is how a server ends up answering
2877 // a Bind the client has already abandoned, and the stream desynchronises.
2878 let mut failed = false;
2879
2880 loop {
2881 let mut tag = [0u8; 1];
2882 if sock.read_exact(&mut tag).await.is_err() {
2883 return Ok(()); // client hung up
2884 }
2885 let len = read_i32(&mut sock).await?;
2886 if len < 4 || len > 64 << 20 {
2887 return Ok(());
2888 }
2889 let body = read_exact(&mut sock, (len - 4) as usize).await?;
2890
2891 // `Sync` always clears the error state; `Terminate` always applies.
2892 if failed && tag[0] != b'S' && tag[0] != b'X' {
2893 continue;
2894 }
2895
2896 match tag[0] {
2897 b'X' => return Ok(()), // Terminate
2898
2899 b'Q' => {
2900 let sql = String::from_utf8_lossy(&body).trim_end_matches('\0').to_string();
2901 let out = run_simple_query(&sql, &db_name, resolved.as_ref(), read_only);
2902 sock.write_all(&out).await?;
2903 sock.write_all(&ready()).await?;
2904 // A simple query closes the unnamed portal, per the protocol.
2905 portals.remove("");
2906 }
2907
2908 // ── Parse: name, SQL, declared parameter type OIDs ─────────────
2909 b'P' => {
2910 let mut at = 0usize;
2911 let name = take_cstr(&body, &mut at);
2912 let sql = take_cstr(&body, &mut at);
2913 let n = take_i16(&body, &mut at).unwrap_or(0).max(0) as usize;
2914 let mut declared = Vec::with_capacity(n);
2915 let mut bad = false;
2916 for _ in 0..n {
2917 match take_i32(&body, &mut at) {
2918 Ok(o) => declared.push(o),
2919 Err(_) => { bad = true; break; }
2920 }
2921 }
2922 if bad {
2923 sock.write_all(&err_msg("08P01", "malformed Parse message")).await?;
2924 failed = true;
2925 continue;
2926 }
2927 // Reject unsupported SQL here rather than at Execute, so the
2928 // client learns at the point it asked — which is also where
2929 // Postgres reports it.
2930 //
2931 // The SQL evaluator gets asked first, or a catalogue query
2932 // would be refused at `Parse` by the NQL path that was never
2933 // going to run it — and the extended protocol is where every
2934 // ORM and async driver lives, so refusing here refuses them
2935 // all.
2936 let probe = probe_sql(&sql, param_count(&sql));
2937 if !sql_engine_owns(&probe) {
2938 if let Err(why) = translate(&probe) {
2939 sock.write_all(&err_msg("0A000", &why)).await?;
2940 failed = true;
2941 continue;
2942 }
2943 }
2944 let param_oids = infer_param_oids(&sql, &declared, resolved.as_ref());
2945 prepared.insert(name, Prepared { sql, param_oids, out_shape: None });
2946 sock.write_all(&parse_complete()).await?;
2947 }
2948
2949 // ── Bind: portal, statement, formats, values, result formats ───
2950 b'B' => {
2951 let mut at = 0usize;
2952 let portal_name = take_cstr(&body, &mut at);
2953 let stmt_name = take_cstr(&body, &mut at);
2954 if !prepared.contains_key(&stmt_name) {
2955 sock.write_all(&err_msg("26000", &format!(
2956 "prepared statement {:?} does not exist", stmt_name))).await?;
2957 failed = true;
2958 continue;
2959 }
2960 let p = &prepared[&stmt_name];
2961 let mut want_formats: Vec<i16> = vec![];
2962 let res: Result<String, String> = (|| {
2963 let nfmt = take_i16(&body, &mut at)? .max(0) as usize;
2964 let mut fmts = Vec::with_capacity(nfmt);
2965 for _ in 0..nfmt {
2966 fmts.push(take_i16(&body, &mut at)?);
2967 }
2968 let nparam = take_i16(&body, &mut at)?.max(0) as usize;
2969 let mut vals: Vec<Option<String>> = Vec::with_capacity(nparam);
2970 for i in 0..nparam {
2971 let l = take_i32(&body, &mut at)?;
2972 let raw: Option<Vec<u8>> = if l < 0 {
2973 None
2974 } else {
2975 let l = l as usize;
2976 if at + l > body.len() {
2977 return Err("truncated Bind parameter".into());
2978 }
2979 let v = body[at..at + l].to_vec();
2980 at += l;
2981 Some(v)
2982 };
2983 // Zero format codes means "all text"; one means "this
2984 // format for every parameter"; otherwise one per value.
2985 let f = match fmts.len() {
2986 0 => 0,
2987 1 => fmts[0],
2988 _ => *fmts.get(i).unwrap_or(&0),
2989 };
2990 let oid = *p.param_oids.get(i).unwrap_or(&OID_TEXT);
2991 vals.push(decode_param(raw.as_deref(), oid, f)?);
2992 }
2993 // Result format codes. asyncpg asks for binary on every
2994 // column, so honouring these is not an optimisation — it
2995 // is the difference between asyncpg reading rows and
2996 // refusing the result outright.
2997 let nres = take_i16(&body, &mut at)?.max(0) as usize;
2998 for _ in 0..nres {
2999 let f = take_i16(&body, &mut at)?;
3000 if f != 0 && f != 1 {
3001 return Err(format!("unknown result format code {}", f));
3002 }
3003 want_formats.push(f);
3004 }
3005 substitute_params(&p.sql, &vals)
3006 })();
3007 match res {
3008 Ok(sql) => {
3009 // Binary encoding must use the types the client was
3010 // TOLD about, so pull the advertised shape across.
3011 let declared = if want_formats.iter().any(|f| *f == 1) {
3012 let p = prepared.get_mut(&stmt_name).expect("checked above");
3013 prepared_shape(p, resolved.as_ref()).clone()
3014 } else {
3015 None
3016 };
3017 portals.insert(portal_name, Portal {
3018 sql, result: None, frozen: None,
3019 formats: want_formats, declared,
3020 });
3021 sock.write_all(&bind_complete()).await?;
3022 }
3023 Err(why) => {
3024 sock.write_all(&err_msg("08P01", &why)).await?;
3025 failed = true;
3026 }
3027 }
3028 }
3029
3030 // ── Describe: 'S' statement, or 'P' portal ─────────────────────
3031 b'D' => {
3032 let kind = body.first().copied().unwrap_or(b'S');
3033 let mut at = 1usize;
3034 let name = take_cstr(&body, &mut at);
3035 if kind == b'S' {
3036 if !prepared.contains_key(&name) {
3037 sock.write_all(&err_msg("26000", &format!(
3038 "prepared statement {:?} does not exist", name))).await?;
3039 failed = true;
3040 continue;
3041 }
3042 let p = prepared.get_mut(&name).expect("checked above");
3043 let oids = p.param_oids.clone();
3044 // asyncpg encodes its arguments from this, so the count has
3045 // to be right or it refuses the call before sending a Bind.
3046 sock.write_all(¶meter_description(&oids)).await?;
3047 // Describe(statement) happens before Bind, so the requested
3048 // result format is not known yet; Postgres reports text
3049 // here too and the client's own Bind decides the encoding.
3050 let out = match prepared_shape(p, resolved.as_ref()) {
3051 Some((cols, col_oids)) => row_description(cols, col_oids),
3052 None => no_data(),
3053 };
3054 sock.write_all(&out).await?;
3055 } else {
3056 let portal = match portals.get_mut(&name) {
3057 Some(p) => p,
3058 None => {
3059 sock.write_all(&err_msg("34000", &format!(
3060 "portal {:?} does not exist", name))).await?;
3061 failed = true;
3062 continue;
3063 }
3064 };
3065 // A bound portal can be run: doing it here means the
3066 // RowDescription reports the columns and types actually
3067 // present, which is strictly better than a guess. psycopg3
3068 // takes this path on every query.
3069 match ensure_executed(portal, &db_name, resolved.as_ref(), read_only) {
3070 Err(encoded) => {
3071 sock.write_all(&encoded).await?;
3072 failed = true;
3073 }
3074 Ok(()) => {
3075 let r = portal.result.as_ref().expect("just executed");
3076 if !r.has_rows {
3077 sock.write_all(&no_data()).await?;
3078 } else {
3079 let (cols, oids) = portal.shape(r);
3080 let fmts: Vec<i16> =
3081 (0..cols.len()).map(|i| portal.format_of(i)).collect();
3082 sock.write_all(&row_description_fmt(&cols, &oids, &fmts)).await?;
3083 }
3084 }
3085 }
3086 }
3087 }
3088
3089 // ── Execute: portal, maximum rows (0 = all) ────────────────────
3090 b'E' => {
3091 let mut at = 0usize;
3092 let name = take_cstr(&body, &mut at);
3093 let max_rows = take_i32(&body, &mut at).unwrap_or(0);
3094 let portal = match portals.get_mut(&name) {
3095 Some(p) => p,
3096 None => {
3097 sock.write_all(&err_msg("34000", &format!(
3098 "portal {:?} does not exist", name))).await?;
3099 failed = true;
3100 continue;
3101 }
3102 };
3103 if let Err(encoded) = ensure_executed(portal, &db_name, resolved.as_ref(), read_only) {
3104 sock.write_all(&encoded).await?;
3105 failed = true;
3106 continue;
3107 }
3108 let r = portal.result.as_ref().expect("just executed");
3109 if !r.has_rows {
3110 let tag = r.tag.clone();
3111 sock.write_all(&command_complete(&tag)).await?;
3112 continue;
3113 }
3114 let (cols, oids) = portal.shape(r);
3115 let limit = if max_rows > 0 {
3116 (r.sent + max_rows as usize).min(r.rows.len())
3117 } else {
3118 r.rows.len()
3119 };
3120 // Encode the whole batch BEFORE writing any of it. A value that
3121 // cannot be sent in the advertised binary type has to become an
3122 // error instead of a truncated row stream — half a result set
3123 // followed by an error is far harder to diagnose than an error.
3124 let mut encoded: Vec<Vec<u8>> = Vec::with_capacity(limit - r.sent);
3125 let mut fail: Option<String> = None;
3126 for row in &r.rows[r.sent..limit] {
3127 let mut vals: Vec<Option<Vec<u8>>> = Vec::with_capacity(cols.len());
3128 for (i, c) in cols.iter().enumerate() {
3129 let v = row.get(&c.src);
3130 let got = if portal.format_of(i) == 1 {
3131 cell_binary(v, oids.get(i).copied().unwrap_or(OID_TEXT))
3132 .map_err(|e| format!("column {:?}: {}", c.out, e))
3133 } else {
3134 Ok(cell(v).map(|s| s.into_bytes()))
3135 };
3136 match got {
3137 Ok(b) => vals.push(b),
3138 Err(e) => { fail = Some(e); break; }
3139 }
3140 }
3141 if fail.is_some() {
3142 break;
3143 }
3144 encoded.push(data_row_bytes(&vals));
3145 }
3146 if let Some(why) = fail {
3147 sock.write_all(&err_msg("22P03", &why)).await?;
3148 failed = true;
3149 continue;
3150 }
3151 let mut out = vec![];
3152 for e in &encoded {
3153 out.extend_from_slice(e);
3154 }
3155 let r = portal.result.as_mut().expect("just executed");
3156 r.sent = limit;
3157 // More rows left and the client capped the batch: suspend the
3158 // portal instead of completing it. This is what a JDBC
3159 // `setFetchSize` and a psycopg3 server-side cursor rely on.
3160 if max_rows > 0 && r.sent < r.rows.len() {
3161 out.extend_from_slice(&portal_suspended());
3162 } else {
3163 let tag = if r.tag_counts_rows {
3164 format!("{} {}", r.tag, r.sent)
3165 } else {
3166 r.tag.clone()
3167 };
3168 out.extend_from_slice(&command_complete(&tag));
3169 }
3170 sock.write_all(&out).await?;
3171 }
3172
3173 // ── Close: 'S' statement, or 'P' portal ───────────────────────
3174 b'C' => {
3175 let kind = body.first().copied().unwrap_or(b'S');
3176 let mut at = 1usize;
3177 let name = take_cstr(&body, &mut at);
3178 if kind == b'S' {
3179 prepared.remove(&name);
3180 } else {
3181 portals.remove(&name);
3182 }
3183 // Closing something that was never open is explicitly not an
3184 // error in the protocol.
3185 sock.write_all(&close_complete()).await?;
3186 }
3187
3188 // Flush: everything is written unbuffered already, so this is a
3189 // no-op — but it must NOT produce a ReadyForQuery, or a client that
3190 // flushes mid-sequence (asyncpg does, after Describe) loses sync.
3191 b'H' => {}
3192
3193 b'S' => {
3194 failed = false;
3195 sock.write_all(&ready()).await?;
3196 }
3197
3198 other => {
3199 sock.write_all(&err_msg(
3200 "08P01",
3201 &format!("unexpected frontend message {:?}", other as char),
3202 )).await?;
3203 failed = true;
3204 }
3205 }
3206 }
3207}
3208
3209const READ_ONLY_MSG: &str =
3210 "this endpoint is running read-only (NEDBD_PG_READ_ONLY=1). Writes are \
3211 implemented but disabled on this server — unset the flag to allow them.";
3212
3213fn no_db(db_name: &str) -> Vec<u8> {
3214 err_msg("3D000", &format!(
3215 "database {:?} is not open on this server — create it first \
3216 (POST /v1/databases), or connect with -d <name>", db_name))
3217}
3218
3219/// `pg_catalog.pg_class` → `pg_class`, but `information_schema.tables` keeps
3220/// its qualifier, because `tables` is a plausible collection name and the
3221/// catalogue must never shadow a user's own data.
3222fn catalog_name(n: &str) -> String {
3223 let joined: Vec<&str> = n.split('.').collect();
3224 if joined.len() >= 2 && joined[joined.len() - 2] == "information_schema" {
3225 format!("information_schema.{}", joined[joined.len() - 1])
3226 } else {
3227 joined[joined.len() - 1].to_string()
3228 }
3229}
3230
3231/// Does the SQL evaluator own this statement?
3232///
3233/// Two ways in. The first is obvious: it reads a catalogue relation.
3234///
3235/// The second is a statement with NO relation at all — a select list of
3236/// literals and scalar function calls, which is exactly what this evaluator
3237/// does and which the SQL→NQL path cannot express (NQL is FROM-first). That
3238/// path answers a handful of EXACT spellings from a canned table
3239/// (`SELECT 1`, `SELECT VERSION()`, `SELECT CURRENT_SCHEMA`), and those
3240/// answers are what existing clients already see — so this predicate rescues
3241/// only what it REFUSES, leaving every spelling it does handle alone.
3242///
3243/// That gap was not hypothetical. SQLAlchemy's PostgreSQL dialect opens every
3244/// connection with `select pg_catalog.version()`, which is one character of
3245/// qualification away from the canned `SELECT VERSION()` and therefore missed
3246/// it — so the engine refused the first statement of dialect initialisation
3247/// and NO SQLAlchemy application could connect at all. A canned list of
3248/// spellings is the same brittleness `pgcatalog` exists to avoid; the fix is
3249/// to let the evaluator answer, because it has `version()`,
3250/// `current_setting()` and the rest as real functions.
3251///
3252/// Cheap: one parse, no execution, no storage access.
3253/// Opt-in: route USER-collection `SELECT`s through the SQL evaluator too.
3254///
3255/// `NEDBD_SQL_ENGINE=1`. Default OFF, and the default is the point — this
3256/// changes which engine answers ordinary queries, and the two engines have to
3257/// be shown to agree before anyone's production reads move. Flipping it is a
3258/// deployment decision, not a build one, so it is read from the environment
3259/// once rather than compiled in.
3260///
3261/// What it unlocks is everything the translator refuses because NQL cannot
3262/// express it: joins, subqueries, `EXISTS`, `UNION`/`INTERSECT`/`EXCEPT`,
3263/// several named aggregates in one grouped row, `array_agg(x ORDER BY y)`.
3264/// What it must not lose is what only the translator has — and a statement the
3265/// evaluator's grammar cannot parse (`TRACE`, `SEARCH`, `VALID AS OF`,
3266/// `TRAVERSE`, every write) still falls through to the translator on its own,
3267/// because `parse` fails and this function is never consulted.
3268/// NQL's table-level verbs, gathered per relation name.
3269///
3270/// A struct rather than the tuple this started as. It held
3271/// `(valid_as_of, search)`; adding `TRACE` and `TRAVERSE` would have made it a
3272/// four-tuple indexed by `.0` through `.3`, and the resolver reads these in a
3273/// different order than it builds them — which is precisely how a positional
3274/// tuple turns into `SEARCH` being rendered where `VALID AS OF` was meant.
3275#[derive(Default, Clone)]
3276struct TableVerbs {
3277 valid_as_of: Option<String>,
3278 search: Option<String>,
3279 /// The edge type for `TRACE <edge>`.
3280 trace: Option<String>,
3281 /// `REVERSE` — walk effects rather than causes.
3282 trace_reverse: bool,
3283 /// The relation name for `TRAVERSE <rel>`.
3284 traverse: Option<String>,
3285}
3286
3287impl TableVerbs {
3288 /// Does this relation carry any verb the catalogue cannot answer?
3289 fn first_unsupported_on_catalogue(&self) -> Option<&'static str> {
3290 if self.valid_as_of.is_some() {
3291 Some("VALID AS OF")
3292 } else if self.search.is_some() {
3293 Some("SEARCH")
3294 } else if self.trace.is_some() {
3295 Some("TRACE")
3296 } else if self.traverse.is_some() {
3297 Some("TRAVERSE")
3298 } else {
3299 None
3300 }
3301 }
3302
3303}
3304
3305/// Whether `NEDBD_SQL_ENGINE` is still set in someone's environment.
3306///
3307/// The flag no longer selects anything — the evaluator answers every SELECT it
3308/// can parse. It is read only so a deployment that still exports it is TOLD
3309/// the variable is now inert, rather than left believing it is holding a
3310/// switch that no longer exists. Silence here is how an operator ends up
3311/// certain their reads are on the old path.
3312fn stale_sql_engine_flag() -> bool {
3313 use std::sync::OnceLock;
3314 static ON: OnceLock<bool> = OnceLock::new();
3315 *ON.get_or_init(|| {
3316 let set = std::env::var("NEDBD_SQL_ENGINE").is_ok();
3317 if set {
3318 eprintln!(
3319 "[nedbd] NEDBD_SQL_ENGINE is set but no longer does anything. The SQL \
3320 evaluator now answers every SELECT it can parse; statements it cannot \
3321 parse still fall through to the translator. You can remove the variable."
3322 );
3323 }
3324 set
3325 })
3326}
3327
3328/// The pre-filtered scan, still spelled in NQL.
3329///
3330/// The LAST place a relation is expressed as text, and it survives for a
3331/// reason that does not apply to the others: the pre-filter is an
3332/// OPTIMISATION. `sqlpush` renders the part of the `WHERE` that NQL evaluates
3333/// identically, so pushing it saves reading rows — and the evaluator's real
3334/// `WHERE` runs above regardless, so getting it wrong costs a wasted row and
3335/// never an answer. Everything else about the scan is a MEANING, and meanings
3336/// now travel as a `relation::Scan` that cannot drop a field.
3337///
3338/// Derived FROM that same struct rather than from the original clauses, so the
3339/// two cannot disagree about what is being read. When the index scan learns to
3340/// take a predicate directly, this function and NQL's parser go together.
3341fn compose_prefiltered(cname: &str, scan: &crate::relation::Scan, pre: &str) -> String {
3342 let mut q = format!("FROM {}", cname);
3343 if let Some(seq) = scan.as_of {
3344 q.push_str(&format!(" AS OF {}", seq));
3345 }
3346 if let Some(d) = &scan.valid_as_of {
3347 q.push_str(&format!(" VALID AS OF {}", nql_string(d)));
3348 }
3349 q.push_str(&format!(" WHERE {}", pre));
3350 if let Some(t) = &scan.search {
3351 q.push_str(&format!(" SEARCH {}", nql_string(t)));
3352 }
3353 if let Some(edge) = &scan.trace {
3354 q.push_str(&format!(" TRACE {}", edge));
3355 if scan.trace_reverse {
3356 q.push_str(" REVERSE");
3357 }
3358 }
3359 if let Some(rel) = &scan.traverse {
3360 q.push_str(&format!(" TRAVERSE {}", rel));
3361 }
3362 q
3363}
3364
3365fn sql_engine_owns(sql: &str) -> bool {
3366 let Ok(sel) = crate::sqlselect::parse(sql) else { return false };
3367 let touched = sel.base_relations();
3368 if touched.is_empty() {
3369 return translate(sql).is_err();
3370 }
3371 if touched.iter().any(|t| crate::pgcatalog::is_catalog(&catalog_name(t))) {
3372 return true;
3373 }
3374 // A user collection reaches the evaluator too, unconditionally. There is
3375 // ONE evaluator now.
3376 //
3377 // This used to return `sql_engine_for_collections()` — an env flag,
3378 // default OFF, on the argument that a collection "has a working answer on
3379 // both paths, so the choice between them is a judgement about parity".
3380 // That argument stopped being true. The translator's answer is not a
3381 // second correct answer, it is a worse one:
3382 //
3383 // SELECT who FROM orders translator -> who, total, _id, _hash,
3384 // _seq, _coll
3385 // evaluator -> who
3386 //
3387 // The projection list was ignored entirely, because NQL has no projection
3388 // to translate it into. `sum(total), avg(total)` in one grouped row is not
3389 // slow on the translator, it is unrepresentable. A flag whose two
3390 // positions give different answers to the same correct SQL is not a
3391 // parity switch, it is a bug with a toggle.
3392 //
3393 // What made this safe to flip is that the fallthrough was never the flag.
3394 // A statement this evaluator cannot PARSE never reaches here — `parse`
3395 // fails at the top of this function and the translator takes it, which is
3396 // still how every write, and anything outside the SELECT grammar, is
3397 // served. Removing the flag narrows nothing; it stops answering parseable
3398 // SQL with a translation of it.
3399 //
3400 // Called here only for its one-shot warning: this is the first point at
3401 // which a deployment still exporting the variable is demonstrably running
3402 // the evaluator, which is exactly when saying so is useful.
3403 let _ = stale_sql_engine_flag();
3404
3405 // The translator has not gone anywhere. It still answers every write and
3406 // every statement this evaluator cannot parse, so the two paths still
3407 // coexist and still have to agree where both can answer. That agreement is
3408 // proven by tests/test_pgwire_parity.py, which spawns two daemons and
3409 // compares them — and which became a TAUTOLOGY the moment the flag it used
3410 // to tell them apart stopped selecting anything. Its own header warned
3411 // about exactly this failure, from the environment side; this is the same
3412 // failure from the code side.
3413 //
3414 // So the lever survives for the harness, under a name no one will mistake
3415 // for a product switch, and pointed the other way: it forces the
3416 // TRANSLATOR rather than enabling the evaluator. Nothing in the product
3417 // reads it, the default path has no flag in it at all, and a parity run
3418 // that forgets to set it compares the evaluator with itself and is
3419 // supposed to look wrong.
3420 !force_translator_for_parity()
3421}
3422
3423/// TEST-ONLY. Forces user collections back onto the translator.
3424///
3425/// Not a supported configuration and not a fallback: it exists so
3426/// `test_pgwire_parity.py` can still put a translator daemon next to an
3427/// evaluator daemon now that `NEDBD_SQL_ENGINE` selects nothing. Setting it in
3428/// production gives you the projection-dropping answers this change removed.
3429fn force_translator_for_parity() -> bool {
3430 use std::sync::OnceLock;
3431 static ON: OnceLock<bool> = OnceLock::new();
3432 *ON.get_or_init(|| {
3433 let on = matches!(
3434 std::env::var("NEDB_PARITY_FORCE_TRANSLATOR").as_deref(),
3435 Ok("1") | Ok("true") | Ok("on")
3436 );
3437 if on {
3438 eprintln!(
3439 "[nedbd] NEDB_PARITY_FORCE_TRANSLATOR is set — user collections are being \
3440 answered by the TRANSLATOR. This is a test lever for the parity harness, \
3441 not a supported configuration: projections are dropped on this path."
3442 );
3443 }
3444 on
3445 })
3446}
3447
3448/// Run a `SELECT` through the full SQL engine when it touches the catalogue.
3449///
3450/// The gate is deliberately narrow: a statement goes to `sqlselect` only when
3451/// one of its tables is a catalogue relation. Everything else keeps the
3452/// SQL→NQL path, which has the index pushdown, `AS OF`, `TRACE` and the
3453/// bounded scans — and whose join story is a real planning question rather
3454/// than a nested loop. Routing a large collection through a nested-loop join
3455/// would be a promise this engine cannot keep.
3456///
3457/// `None` means "not mine": the caller falls through to the ordinary path, so
3458/// the error the client sees is the ordinary path's error rather than a
3459/// confusing one from a parser that was never meant to handle the statement.
3460fn try_catalog_select(
3461 sql: &str,
3462 db: Option<&Arc<Db>>,
3463) -> Result<Option<(Executed, crate::sqlplan::Plan)>, Vec<u8>> {
3464 let sel = match crate::sqlselect::parse(sql) {
3465 Ok(sel) => sel,
3466 Err(why) => {
3467 // A statement that plainly reads the catalogue but that this
3468 // engine cannot parse gets the PARSE error, not the NQL path's.
3469 //
3470 // Falling through unconditionally produced an actively false
3471 // message: `\d` and `\dp` were told "JOIN is not supported",
3472 // which stopped being true the moment joins started working — and
3473 // a wrong explanation is worse than a blunt one, because it sends
3474 // the reader to fix the wrong thing.
3475 if mentions_catalog(sql) {
3476 return Err(err_msg("0A000", &format!(
3477 "this catalogue query uses SQL this endpoint does not \
3478 implement: {}", why)));
3479 }
3480 return Ok(None);
3481 }
3482 };
3483
3484 // Which relations does it read — at ANY depth? `\dd` names its catalogue
3485 // relations only inside a derived table, and `\dT` only inside two
3486 // subqueries; a walk over the top-level FROM list alone would route both
3487 // to the NQL path, which cannot parse them and would report an error that
3488 // sends the reader to fix the wrong thing.
3489 if !sql_engine_owns(sql) {
3490 return Ok(None);
3491 }
3492
3493 // The storage pre-filter, resolved per relation NAME and computed once.
3494 //
3495 // The resolver is handed a name (`orders`) but the WHERE clause qualifies
3496 // by BINDING (`o.status` for `FROM orders o`), so the predicate has to be
3497 // looked up by name and rendered against that relation's binding. Getting
3498 // this wrong is silent: the pre-filter simply never matches and the scan
3499 // quietly reads the whole collection, which is exactly what EXPLAIN caught
3500 // the first time round — `Seq Scan on orders o (actual rows=3)` when the
3501 // query wanted two.
3502 //
3503 // A name appearing TWICE (a self-join, `FROM t a JOIN t b`) maps to two
3504 // different bindings with different predicates, and one scan cannot serve
3505 // both. Those are dropped rather than guessed at.
3506 // `AS OF SYSTEM TIME <seq>`, per relation name.
3507 //
3508 // The resolver is keyed by NAME, so one collection named twice gets ONE
3509 // scan. `FROM orders AS OF 1 o JOIN orders n` asks for that collection at
3510 // two different sequences at once, and a single scan cannot serve both.
3511 //
3512 // This is REFUSED rather than resolved to one of them, and the reason is
3513 // worth keeping: the first version dropped the qualifier when a name was
3514 // ambiguous — the same "don't guess" instinct that is right for a
3515 // pre-filter. It is wrong here. Dropping a pre-filter costs a wasted row;
3516 // dropping an AS OF answers a question about the past with data from the
3517 // present, and it does it silently. The query `... orders AS OF 1 o JOIN
3518 // orders n ...` returned the CURRENT value for both sides and looked fine.
3519 let temporal: std::collections::HashMap<String, u64> = {
3520 // Gather every sequence each name is read at first, INCLUDING the
3521 // absent one, then judge. Deciding as we walk got this wrong: the
3522 // first arm of a self-join was judged before it had been recorded, so
3523 // a legitimate pair reported the wrong reason.
3524 let mut seen: std::collections::HashMap<String, Vec<Option<u64>>> =
3525 std::collections::HashMap::new();
3526 for t in sel.from.iter().chain(sel.joins.iter().map(|j| &j.table)) {
3527 seen.entry(catalog_name(&t.name).to_ascii_lowercase())
3528 .or_default()
3529 .push(t.as_of);
3530 }
3531 let mut out: std::collections::HashMap<String, u64> = std::collections::HashMap::new();
3532 for (key, ats) in &seen {
3533 let mut distinct: Vec<Option<u64>> = ats.clone();
3534 distinct.sort();
3535 distinct.dedup();
3536 match distinct.as_slice() {
3537 // One sequence for this name, however many times it appears.
3538 [Some(seq)] => {
3539 out.insert(key.clone(), *seq);
3540 }
3541 [None] => {}
3542 // More than one. Say WHICH disagreement it is, because the two
3543 // read very differently to whoever wrote the query.
3544 _ => {
3545 let mixed_tip = distinct.contains(&None);
3546 let seqs: Vec<String> =
3547 distinct.iter().flatten().map(|s| s.to_string()).collect();
3548 let detail = if mixed_tip {
3549 format!(
3550 "at the tip and AS OF {}",
3551 seqs.join(" and "))
3552 } else {
3553 format!("AS OF {}", seqs.join(" and "))
3554 };
3555 return Err(err_msg("0A000", &format!(
3556 "{:?} is read {} in one statement. This endpoint reads each \
3557 collection once per statement, so it cannot serve both — and \
3558 answering from either one would silently return the same rows for \
3559 both arms, which is the comparison failing to be a comparison. Ask \
3560 the two questions separately.",
3561 key, detail)));
3562 }
3563 }
3564 }
3565 out
3566 };
3567
3568 // NQL's own verbs, per relation name: `(VALID AS OF, SEARCH)`.
3569 //
3570 // Same one-scan-per-name constraint as the temporal map, and the same
3571 // verdict for the same reason: two different values for one scan is
3572 // REFUSED, because silently picking one would answer a different question
3573 // than the one asked and look like it worked.
3574 let nql_verbs: std::collections::HashMap<String, TableVerbs> = {
3575 let mut out: std::collections::HashMap<String, TableVerbs> =
3576 std::collections::HashMap::new();
3577 for t in sel.from.iter().chain(sel.joins.iter().map(|j| &j.table)) {
3578 let k = catalog_name(&t.name).to_ascii_lowercase();
3579 let e = out.entry(k.clone()).or_default();
3580 // REVERSE rides with the edge type rather than being reconciled on
3581 // its own: `TRACE caused_by` and `TRACE caused_by REVERSE` are two
3582 // different questions about the same edge, and reconciling the
3583 // direction separately would let them merge into one scan.
3584 if t.trace.is_some() {
3585 e.trace_reverse = t.trace_reverse;
3586 }
3587 for (slot, incoming, verb) in [
3588 (&mut e.valid_as_of, &t.valid_as_of, "VALID AS OF"),
3589 (&mut e.search, &t.search, "SEARCH"),
3590 (&mut e.trace, &t.trace, "TRACE"),
3591 (&mut e.traverse, &t.traverse, "TRAVERSE"),
3592 ] {
3593 match (slot.as_deref(), incoming.as_deref()) {
3594 (Some(a), Some(b)) if a != b => {
3595 return Err(err_msg("0A000", &format!(
3596 "{:?} is read with two different {} arguments in one statement \
3597 ({:?} and {:?}). This endpoint reads each collection once, so \
3598 it cannot serve both. Ask the two questions separately.",
3599 k, verb, a, b)));
3600 }
3601 (None, Some(b)) => *slot = Some(b.to_string()),
3602 _ => {}
3603 }
3604 }
3605 }
3606 out
3607 };
3608
3609 let pushdown_prefilters: std::collections::HashMap<String, String> = {
3610 let refs: Vec<&crate::sqlselect::TableRef> = sel
3611 .from
3612 .iter()
3613 .chain(sel.joins.iter().map(|j| &j.table))
3614 .collect();
3615 let bindings: Vec<String> = refs.iter().map(|t| t.binding()).collect();
3616 let nullable = crate::sqlpush::nullable_bindings(&sel);
3617 let mut out = std::collections::HashMap::new();
3618 let mut ambiguous: Vec<String> = vec![];
3619 for t in &refs {
3620 let key = catalog_name(&t.name).to_ascii_lowercase();
3621 if out.contains_key(&key) || ambiguous.contains(&key) {
3622 out.remove(&key);
3623 ambiguous.push(key);
3624 continue;
3625 }
3626 if let Some(p) = crate::sqlpush::nql_prefilter(
3627 sel.where_.as_ref(), &t.binding(), &bindings, &nullable) {
3628 out.insert(key, p);
3629 }
3630 }
3631 out
3632 };
3633
3634 let resolve = |name: &str| -> anyhow::Result<Option<Box<dyn crate::sqlselect::Relation>>> {
3635 let cname = catalog_name(name);
3636 // A catalogue relation is SYNTHESISED from the current shape of the
3637 // store: it has no log, so it has no history, and there is nothing for
3638 // a temporal or full-text qualifier to mean.
3639 //
3640 // Refused rather than ignored, and the difference is the entire point.
3641 // Ignoring `AS OF SYSTEM TIME 0` answers a question about the past with
3642 // present-day rows and looks like it worked — and that is exactly what
3643 // started happening here the moment the SQL parser learned `AS OF`:
3644 // before, the statement failed to parse and fell through to the
3645 // translator, which refused it properly. Teaching one layer a clause
3646 // silently un-taught another layer's refusal, and a test written long
3647 // before this change is what caught it.
3648 {
3649 let k = cname.to_ascii_lowercase();
3650 let bad = if temporal.contains_key(&k) {
3651 Some("AS OF SYSTEM TIME")
3652 } else {
3653 nql_verbs.get(&k).and_then(|v| v.first_unsupported_on_catalogue())
3654 };
3655 if let Some(clause) = bad {
3656 if crate::pgcatalog::is_catalog(&cname) {
3657 anyhow::bail!(
3658 "{} is not supported on the catalogue relation {:?} — a catalogue is \
3659 synthesised from the store's current shape rather than read from the \
3660 log, so it has no history to reach and no document text to search. \
3661 Ignoring the clause would answer your question with present-day rows \
3662 and look like it worked",
3663 clause, cname);
3664 }
3665 }
3666 }
3667 if let Some(rows) = crate::pgcatalog::rows(&cname, db) {
3668 // A synthesised catalogue relation is small and built eagerly;
3669 // wrapping it satisfies the streaming contract without pretending
3670 // it is lazy.
3671 return Ok(Some(crate::sqlselect::from_vec(rows)));
3672 }
3673 // A join between a catalogue relation and a real collection is
3674 // legitimate, so a user table still resolves.
3675 //
3676 // `nql::query` materialises whatever it is asked for, so what it is
3677 // ASKED for is the whole cost of this line. It used to be
3678 // `FROM <collection>` — every document, unconditionally, before a
3679 // single predicate ran. Free on a catalogue relation of a few dozen
3680 // synthesised rows; on a user collection it is the difference between
3681 // reading one document and reading all of them.
3682 //
3683 // `sqlpush::nql_prefilter` renders the part of the WHERE that NQL is
3684 // known to evaluate identically, and the full WHERE still runs above
3685 // this — so the pre-filter can only ever cost a wasted row, never an
3686 // answer. See the module note in `sqlpush` for why each refused
3687 // construct is refused.
3688 //
3689 // Still eager, and deliberately not claimed otherwise: this narrows
3690 // WHAT is materialised, not WHETHER it is. A lazy storage scan is the
3691 // other half and is tracked in HANDOFF.
3692 let key = cname.to_ascii_lowercase();
3693 let pre = pushdown_prefilters.get(&key);
3694 // Composed in NQL'S OWN CLAUSE ORDER, which its grammar fixes as
3695 //
3696 // FROM coll [AS OF seq] [VALID AS OF "date"] [WHERE p] [SEARCH "t"]
3697 //
3698 // and which is not negotiable: emit `AS OF` after `WHERE` and the NQL
3699 // parser reads it as part of the predicate expression. This is the
3700 // whole mechanism behind "NQL folded into neSQL" — the SQL side parses
3701 // the verbs and composes joins and subqueries around them, while the
3702 // NQL engine remains the one implementation that executes them.
3703 // Built once, parameterised by whether the pre-filter is included, so
3704 // the retry below cannot diverge from the real query by forgetting a
3705 // clause.
3706 //
3707 // It previously did. The retry was hand-rolled as
3708 // FROM <coll> [AS OF <seq>]
3709 // on the stated grounds that "the fallback drops the PRE-FILTER, which
3710 // is free". Dropping the pre-filter IS free -- the full WHERE runs
3711 // above. But that string also dropped VALID AS OF and SEARCH, which
3712 // are not free and have no equivalent up there: the retry answered
3713 // with rows nobody asked about and looked like it worked. The AS OF
3714 // case had already been found and special-cased; the other two were
3715 // the same bug standing next to it.
3716 let Some(db) = db else { return Ok(None) };
3717
3718 // The scan as DATA. No string is built and none is parsed: the
3719 // qualifiers go to the store as fields.
3720 //
3721 // This replaced `crate::nql::query(db, &compose(true))`, which
3722 // rendered `FROM coll AS OF n VALID AS OF '...' WHERE ... SEARCH '...'`
3723 // into text and handed it back to the NQL parser. That was a
3724 // translation living inside the thing built to stop translating, and
3725 // it failed the same way translations do: the retry path composed its
3726 // own shorter string and dropped two clauses, and `SEARCH 'o''brien'`
3727 // was a quoting question rather than a value.
3728 let verbs = nql_verbs.get(&key);
3729 let scan = crate::relation::Scan {
3730 coll: cname.to_string(),
3731 as_of: temporal.get(&key).copied(),
3732 valid_as_of: verbs.and_then(|v| v.valid_as_of.clone()),
3733 search: verbs.and_then(|v| v.search.clone()),
3734 trace: verbs.and_then(|v| v.trace.clone()),
3735 trace_reverse: verbs.map(|v| v.trace_reverse).unwrap_or(false),
3736 traverse: verbs.and_then(|v| v.traverse.clone()),
3737 trace_limit: crate::relation::DEFAULT_TRACE_LIMIT,
3738 };
3739
3740 // The pre-filter is the one part still expressed in NQL, because it is
3741 // the one part that is an OPTIMISATION rather than a meaning: the full
3742 // `WHERE` runs in the evaluator above regardless, so a pre-filter can
3743 // only ever save a row, never change an answer. When NQL declines it,
3744 // the scan simply happens unfiltered — which is what the query would
3745 // have done anyway, and no clause is lost with it because the scan is
3746 // a struct and the struct does not change.
3747 if let Some(p) = pre {
3748 let filtered = compose_prefiltered(&cname, &scan, p);
3749 if let Ok((rows, _)) = crate::nql::query(db, &filtered) {
3750 return Ok(Some(crate::sqlselect::from_vec(rows)));
3751 }
3752 }
3753 // A collection that does not exist is NOT an empty one.
3754 //
3755 // `nql::query` used to error on an unknown collection, and the `Err`
3756 // arm returned `Ok(None)` — which the evaluator reports as
3757 // `relation "x" does not exist`. Reading the store directly lost that
3758 // for free, because `relation::read` on a name nothing was ever
3759 // written under returns an empty Vec, indistinguishable from a
3760 // collection that exists and is empty.
3761 //
3762 // The cost of getting this wrong is a typo answering successfully:
3763 // `SELECT * FROM orders JOIN x ON true` returned `[]` rather than
3764 // naming `x`, and an empty join result looks exactly like a correct
3765 // answer about data that isn't there.
3766 //
3767 // `list_ids_including_deleted` rather than `collections`, so a
3768 // collection whose rows have all been deleted still EXISTS. Its
3769 // tombstones are the evidence it did.
3770 // A CATALOGUE relation is exempt, and the distinction is deliberate.
3771 // `pg_db_role_setting` and friends are things NEDB has nothing for;
3772 // the documented behaviour is that they are EMPTY rather than an
3773 // error, because a client introspecting the catalogue is asking "is
3774 // there anything here" and "no" is a valid answer. `psql \drds` walks
3775 // exactly such a relation, and my first version of this check broke
3776 // it. A user collection is the opposite case: nobody types a
3777 // collection name hoping it does not exist.
3778 // Membership is by SCHEMA, not by a list of names we happen to
3779 // implement. `is_catalog` alone was not enough: `pg_db_role_setting`
3780 // is in neither its match arm nor EMPTY_CATALOG, so `psql \drds`
3781 // started reporting `relation "pg_catalog.pg_db_role_setting" does
3782 // not exist` — a regression against the documented stance that what
3783 // NEDB has nothing for is EMPTY rather than an error. Enumerating
3784 // catalogue relations means the next introspection command psql
3785 // grows breaks the same way.
3786 let catalogue = crate::pgcatalog::is_catalog(&cname)
3787 || cname.starts_with("pg_")
3788 || cname.starts_with("information_schema.");
3789 let known = catalogue
3790 || db.collections().iter().any(|c| c == &cname)
3791 || !db.list_ids_including_deleted(&cname).is_empty();
3792
3793 // TWO CONTEXTS, TWO RIGHT ANSWERS — and they used to be distinguished
3794 // for free, because the evaluator only ever served catalogue
3795 // relations. Now that it serves user collections too, the distinction
3796 // has to be made on purpose or one of the two answers is lost.
3797 //
3798 // SINGLE RELATION -> EMPTY. NEDB is schemaless and a collection is
3799 // created by its first write, so "does not exist" and "is empty"
3800 // are the same observable state. Erroring makes it impossible to
3801 // read a collection before writing to it.
3802 //
3803 // A JOIN -> ERROR. Nobody joins against a relation they believe is
3804 // absent; there the name is a typo or a bug, and an empty join
3805 // result is indistinguishable from a correct answer about data that
3806 // is not there. `SELECT * FROM orders JOIN x ON true` returning []
3807 // is the failure this guards.
3808 //
3809 // I flattened both into "error" first, which broke `psql \drds` and
3810 // the documented schemaless read. The rule is the one the test for it
3811 // already spelled out.
3812 if !known && !sel.joins.is_empty() {
3813 return Ok(None);
3814 }
3815 Ok(Some(crate::sqlselect::from_vec(crate::relation::read_json(db, &scan))))
3816 };
3817
3818 let (cols, rows, plan) = crate::sqlselect::execute_explain(
3819 &sel,
3820 &resolve,
3821 crate::sqljoin::JoinExec::Auto,
3822 )
3823 .map_err(|e| err_msg("42601", &e.to_string()))?;
3824
3825 Ok(Some((
3826 Executed {
3827 rows,
3828 // The KEY is what the row is stored under; the NAME is what the
3829 // client sees. They differ when a select list has duplicate output
3830 // names, which PostgreSQL permits and generated SQL relies on.
3831 project: cols
3832 .iter()
3833 .map(|c| Col::renamed(&c.key, &c.name))
3834 .collect(),
3835 has_rows: true,
3836 tag: "SELECT".into(),
3837 tag_counts_rows: true,
3838 },
3839 plan,
3840 )))
3841}
3842
3843/// Strip a leading `EXPLAIN`, returning the statement it wraps.
3844///
3845/// `ANALYZE` and `VERBOSE` are accepted and ignored: this endpoint always
3846/// executes and always reports actual rows, so `EXPLAIN` and
3847/// `EXPLAIN ANALYZE` genuinely do the same thing here. Accepting the keyword
3848/// and silently doing the honest thing beats refusing a client's spelling.
3849fn strip_explain(sql: &str) -> Option<&str> {
3850 let t = sql.trim().trim_end_matches(';').trim();
3851 let mut rest = t.strip_prefix("EXPLAIN").or_else(|| t.strip_prefix("explain"))?;
3852 // Require a word boundary so `EXPLAINED` is not mistaken for a keyword.
3853 if !rest.starts_with(char::is_whitespace) {
3854 return None;
3855 }
3856 rest = rest.trim_start();
3857 loop {
3858 let low = rest.to_lowercase();
3859 if let Some(r) = low.strip_prefix("analyze").or_else(|| low.strip_prefix("analyse")) {
3860 if r.starts_with(char::is_whitespace) || r.is_empty() {
3861 rest = rest[rest.len() - r.len()..].trim_start();
3862 continue;
3863 }
3864 }
3865 if let Some(r) = low.strip_prefix("verbose") {
3866 if r.starts_with(char::is_whitespace) || r.is_empty() {
3867 rest = rest[rest.len() - r.len()..].trim_start();
3868 continue;
3869 }
3870 }
3871 break;
3872 }
3873 Some(rest)
3874}
3875
3876/// One text column named `QUERY PLAN`, which is exactly the shape PostgreSQL
3877/// returns — so `psql` prints it without special handling.
3878fn plan_result(lines: Vec<String>) -> Executed {
3879 Executed {
3880 rows: lines
3881 .into_iter()
3882 .map(|l| serde_json::json!({ "QUERY PLAN": l }))
3883 .collect(),
3884 project: vec![Col::same("QUERY PLAN")],
3885 has_rows: true,
3886 tag: "EXPLAIN".into(),
3887 tag_counts_rows: false,
3888 }
3889}
3890
3891/// Does the raw SQL plainly read a catalogue relation?
3892///
3893/// A cheap text check, used only to decide WHICH error to report when the
3894/// statement cannot be parsed — never to decide what a parsable statement
3895/// means. `pg_` is the giveaway: every catalogue relation is prefixed, and so
3896/// is the `pg_catalog` schema qualifier.
3897fn mentions_catalog(sql: &str) -> bool {
3898 let low = sql.to_lowercase();
3899 low.contains("pg_catalog.")
3900 || low.contains("information_schema.")
3901 || low.contains("from pg_")
3902 || low.contains("join pg_")
3903}
3904
3905/// The catalogue relation a translated query reads from, if any.
3906///
3907/// Reads the collection straight off the parsed NQL rather than re-parsing the
3908/// SQL, so it cannot disagree with what the executor is about to run.
3909fn catalog_target(nql: &str) -> Option<String> {
3910 let coll = crate::nql::parse(nql).ok()?.coll;
3911 if crate::pgcatalog::is_catalog(&coll) {
3912 Some(coll)
3913 } else {
3914 None
3915 }
3916}
3917
3918/// True when the statement carried a RETURNING clause. Checked against the raw
3919/// SQL because `RETURNING *` yields an EMPTY projection, which is otherwise
3920/// indistinguishable from "no RETURNING at all".
3921fn wants_returning(sql: &str) -> bool {
3922 find_kw(&sql.to_uppercase(), "RETURNING").is_some()
3923}
3924
3925/// A unique key for a server-assigned INSERT id.
3926fn next_row_id() -> String {
3927 use std::sync::atomic::{AtomicU64, Ordering};
3928 static N: AtomicU64 = AtomicU64::new(0);
3929 let n = N.fetch_add(1, Ordering::Relaxed);
3930 let ts = std::time::SystemTime::now()
3931 .duration_since(std::time::UNIX_EPOCH)
3932 .map(|d| d.as_micros())
3933 .unwrap_or(0);
3934 format!("r{}{}", ts, n)
3935}
3936
3937/// One executed statement, held apart from any wire encoding.
3938///
3939/// This type is why the simple and extended protocols share an execution path
3940/// rather than growing two copies of the SQL→NEDB semantics. The simple path
3941/// encodes it immediately; the extended path parks it in a portal and dribbles
3942/// the rows out across successive `Execute` messages. Both get identical
3943/// answers because both call `execute_stmt`.
3944pub struct Executed {
3945 /// The rows the client gets — a SELECT's result, or a write's `RETURNING`.
3946 pub rows: Vec<Value>,
3947 /// How to project them (empty = every key in the row).
3948 pub project: Vec<Col>,
3949 /// Whether the client asked for rows at all. Distinct from `rows.is_empty()`:
3950 /// a `SELECT` matching nothing still owes a `RowDescription`, while an
3951 /// `UPDATE` without `RETURNING` owes `NoData`.
3952 pub has_rows: bool,
3953 /// The command tag, already rendered — except for a SELECT, where the row
3954 /// count is only known once the rows have actually been sent.
3955 pub tag: String,
3956 /// True when `tag` is a SELECT-shaped tag whose count is the rows sent.
3957 pub tag_counts_rows: bool,
3958}
3959
3960impl Executed {
3961 fn nothing(tag: &str) -> Self {
3962 Executed { rows: vec![], project: vec![], has_rows: false, tag: tag.to_string(), tag_counts_rows: false }
3963 }
3964 /// Render the final `CommandComplete` given how many rows went out.
3965 fn tag_for(&self, sent: usize) -> String {
3966 if self.tag_counts_rows { format!("{} {}", self.tag, sent) } else { self.tag.clone() }
3967 }
3968}
3969
3970/// Run ONE statement. `Err` carries an already-encoded `ErrorResponse`.
3971///
3972/// Every SQL→NEDB decision lives here, which is the point: the extended query
3973/// protocol added below is then purely a matter of message framing, and cannot
3974/// drift from the simple path's semantics.
3975/// Run one neSQL statement against a database, in process.
3976///
3977/// # Why this exists
3978///
3979/// Until this, the engine had exactly one SQL execution path and it was welded
3980/// to the wire protocol: `execute_stmt` is private, takes the connection's
3981/// read-only flag, and reports failure as ALREADY-ENCODED Postgres error bytes.
3982/// Nothing outside a pgwire session could run SQL against a `Db`.
3983///
3984/// That was survivable while the only SQL client was a socket. It stopped being
3985/// survivable when neSQL — which owns the language — needed to run the language
3986/// from a CLI, because the alternatives were a CLI that opens a TCP connection
3987/// to its own process, or a second SQL front end living in the CLI. The second
3988/// one is worse than it sounds: it makes the CLI a quieter second authority on
3989/// what the language accepts, and the first divergence between them would be
3990/// discovered by a user, not by us.
3991///
3992/// So the path the wire already takes is exposed, with the error decoded into
3993/// text. Same parser, same translator, same evaluator, same decision about
3994/// which engine runs a statement — one authority.
3995/// The rows an `UPDATE` or `DELETE` will act on — chosen by the SQL evaluator.
3996///
3997/// This used to render the predicate as NQL and run `nql::query`, which meant
3998/// a write could only match what the NQL parser understood, even though the
3999/// statement arrived as SQL and the read path had long since stopped needing a
4000/// translation. `UPDATE … WHERE _id IN (SELECT …)` was unreachable for exactly
4001/// that reason: the subquery translated into NQL text the NQL parser cannot
4002/// parse. Selecting with a real `SELECT` closes that gap by not having a second
4003/// predicate implementation to fall short of the first.
4004///
4005/// Whole rows, not just `_id`: `DELETE … RETURNING` has to capture the row
4006/// BEFORE the tombstone, so the selection is what it returns.
4007///
4008/// # The unknown-collection guard is not incidental
4009///
4010/// `nql::query` ERRORS on a collection that does not exist; the evaluator's
4011/// scan returns no rows, because a schemaless read of an absent collection is
4012/// legitimately empty. Swapping one for the other without this check would
4013/// turn `UPDATE nowhere SET x = 1` from a loud 42P01 into a silent
4014/// `UPDATE 0` — a write that reports success having done nothing, which is the
4015/// worst available outcome and the reason this function refuses first.
4016fn rows_for_write(db: &Arc<Db>, coll: &str, where_sql: &str, nql: &str)
4017 -> std::result::Result<Vec<Value>, Vec<u8>>
4018{
4019 let known = db.collections().iter().any(|c| c == coll)
4020 || !db.list_ids_including_deleted(coll).is_empty();
4021 if !known {
4022 return Err(err_msg("42P01", &format!("relation \"{}\" does not exist", coll)));
4023 }
4024 let sel = format!("SELECT * FROM {} {}", coll, where_sql).trim().to_string();
4025 // read_only: this is the SELECT half of the write, and nothing it does
4026 // should be able to write. The caller already passed `need_write!()`.
4027 execute_sql(db, &sel, true)
4028 .map(|done| done.rows)
4029 .map_err(|e| err_msg("42601", &format!(
4030 "{} (selecting rows with: {}; the NQL rendering of this predicate \
4031 would have been: {})", e, sel, nql)))
4032}
4033
4034pub fn execute_sql(db: &Arc<Db>, sql: &str, read_only: bool)
4035 -> std::result::Result<Executed, String>
4036{
4037 execute_stmt(sql, "", Some(db), read_only).map_err(|wire| decode_wire_error(&wire))
4038}
4039
4040/// Pull the human-readable message out of an encoded ErrorResponse.
4041///
4042/// The wire format is a sequence of NUL-terminated `field-code || text` runs
4043/// terminated by an empty field. `M` is the primary message and `C` the
4044/// SQLSTATE; both are reported, because a caller who loses the SQLSTATE loses
4045/// the only machine-stable part of the error.
4046fn decode_wire_error(buf: &[u8]) -> String {
4047 let mut code: Option<String> = None;
4048 let mut msg: Option<String> = None;
4049 // Skip the 1-byte tag and 4-byte length when they are present.
4050 let body = if buf.len() > 5 { &buf[5..] } else { buf };
4051 let mut i = 0usize;
4052 while i < body.len() && body[i] != 0 {
4053 let field = body[i];
4054 i += 1;
4055 let start = i;
4056 while i < body.len() && body[i] != 0 { i += 1; }
4057 let text = String::from_utf8_lossy(&body[start..i]).into_owned();
4058 i += 1; // the NUL
4059 match field {
4060 b'C' => code = Some(text),
4061 b'M' => msg = Some(text),
4062 _ => {}
4063 }
4064 }
4065 match (code, msg) {
4066 (Some(c), Some(m)) => format!("{} ({})", m, c),
4067 (None, Some(m)) => m,
4068 // Never silently produce an empty error. A failure we cannot read is
4069 // still a failure, and saying so beats returning "".
4070 _ => format!(
4071 "the engine refused the statement and the error could not be decoded ({} bytes of wire response)", buf.len()
4072 ),
4073 }
4074}
4075
4076fn execute_stmt(
4077 stmt_sql: &str,
4078 db_name: &str,
4079 db: Option<&Arc<Db>>,
4080 read_only: bool,
4081) -> Result<Executed, Vec<u8>> {
4082 // The full SQL engine gets first refusal, but ONLY for statements that
4083 // touch the catalogue — see `try_catalog_select`. It has to run before
4084 // `translate`, because `translate` targets NQL and NQL cannot express a
4085 // join, a CASE or a scalar function at all.
4086 // EXPLAIN reports which engine would run the statement, and a plan only
4087 // when the SQL evaluator is the engine that actually runs it. Describing a
4088 // pipeline the statement would not take is the one thing an EXPLAIN must
4089 // never do.
4090 if let Some(inner) = strip_explain(stmt_sql) {
4091 if let Some((_, plan)) = try_catalog_select(inner, db)? {
4092 return Ok(plan_result(plan.render()));
4093 }
4094 let mut lines = vec![];
4095 match translate(inner) {
4096 Ok(_) => {
4097 lines.push(
4098 "NQL path — this statement is translated to NQL and \
4099 executed by the storage engine, not by the SQL evaluator."
4100 .to_string(),
4101 );
4102 lines.push(
4103 "No plan is reported, because the SQL evaluator is not \
4104 what runs it. Reporting one would describe a pipeline \
4105 that never executed."
4106 .to_string(),
4107 );
4108 lines.push(
4109 "The SQL evaluator (joins, CASE, scalar functions, a \
4110 hash-join planner) currently serves catalogue queries."
4111 .to_string(),
4112 );
4113 }
4114 Err(why) => lines.push(format!("cannot be executed: {why}")),
4115 }
4116 return Ok(plan_result(lines));
4117 }
4118
4119 if let Some((done, _plan)) = try_catalog_select(stmt_sql, db)? {
4120 return Ok(done);
4121 }
4122
4123 let stmt = translate(stmt_sql).map_err(|why| err_msg("0A000", &why))?;
4124
4125 // Every arm below that touches storage needs a database; resolve the
4126 // "no such database" answer once instead of at each use.
4127 macro_rules! need_db {
4128 () => {
4129 match db {
4130 Some(db) => db,
4131 None => return Err(no_db(db_name)),
4132 }
4133 };
4134 }
4135 macro_rules! need_write {
4136 () => {
4137 if read_only {
4138 return Err(err_msg("25006", READ_ONLY_MSG));
4139 }
4140 };
4141 }
4142
4143 match stmt {
4144 Stmt::Ok(tag) => Ok(Executed::nothing(if tag.is_empty() { "SELECT 0" } else { tag })),
4145
4146 Stmt::Canned { cols, row } => {
4147 // Fold the canned answer into an ordinary row so the encoders,
4148 // the portal machinery and `Describe` all see one shape.
4149 let mut obj = serde_json::Map::new();
4150 for (c, v) in cols.iter().zip(row.iter()) {
4151 obj.insert(c.clone(), Value::String(v.clone()));
4152 }
4153 Ok(Executed {
4154 rows: vec![Value::Object(obj)],
4155 project: cols.iter().map(|c| Col::same(c)).collect(),
4156 has_rows: true,
4157 tag: "SELECT".into(),
4158 tag_counts_rows: true,
4159 })
4160 }
4161
4162 Stmt::Query { nql, project } => {
4163 // A catalogue relation is synthesised from the live database
4164 // rather than read from it — but it is still queried with the
4165 // ORDINARY predicate path, so WHERE / ORDER BY / LIMIT and the
4166 // `~` operators work on it because they are the same operators.
4167 //
4168 // Checked BEFORE `need_db!()`: `SELECT * FROM pg_namespace` has to
4169 // answer even when the client connected without naming a database,
4170 // which is exactly what psql does on startup. Refusing there is
4171 // how "psql cannot connect" starts.
4172 if let Some(coll) = catalog_target(&nql) {
4173 let rows = crate::pgcatalog::rows(&coll, db)
4174 .expect("catalog_target only returns names pgcatalog serves");
4175 let rows = crate::nql::query_rows(rows, &nql)
4176 .map_err(|e| err_msg("42601", &e.to_string()))?;
4177 return Ok(Executed {
4178 rows, project, has_rows: true,
4179 tag: "SELECT".into(), tag_counts_rows: true,
4180 });
4181 }
4182 let db = need_db!();
4183 let (rows, _) = crate::nql::query(db, &nql).map_err(|e| {
4184 err_msg("42601", &format!("{} (translated to NQL: {})", e, nql))
4185 })?;
4186 Ok(Executed { rows, project, has_rows: true, tag: "SELECT".into(), tag_counts_rows: true })
4187 }
4188
4189 Stmt::Insert { coll, rows, returning } => {
4190 let db = need_db!();
4191 need_write!();
4192 let mut written: Vec<Value> = vec![];
4193 for (i, r) in rows.iter().enumerate() {
4194 // The engine requires an id. When the statement did not supply
4195 // one, mint a unique key rather than silently overwriting a
4196 // shared default.
4197 let id = match &r.id {
4198 Some(id) => id.clone(),
4199 None => format!("{}-{}", next_row_id(), i),
4200 };
4201 let node = db
4202 .put(&coll, &id, Value::Object(r.doc.clone()),
4203 r.caused_by.clone(), r.valid_from.clone(), r.valid_to.clone())
4204 .map_err(|e| err_msg("XX000", &format!("INSERT failed: {}", e)))?;
4205 written.push(crate::nql::node_to_json(&node));
4206 }
4207 let n = written.len();
4208 let has_rows = wants_returning(stmt_sql);
4209 Ok(Executed {
4210 rows: if has_rows { written } else { vec![] },
4211 project: returning,
4212 has_rows,
4213 // Postgres reports `INSERT <oid> <rows>`; the oid is always 0.
4214 tag: format!("INSERT 0 {}", n),
4215 tag_counts_rows: false,
4216 })
4217 }
4218
4219 Stmt::Update { coll, set, where_sql, nql, returning } => {
4220 let db = need_db!();
4221 need_write!();
4222 // Rows come from the SQL evaluator, so an UPDATE matches exactly
4223 // what a SELECT with the same WHERE matches — one predicate
4224 // implementation, not two.
4225 let matched = rows_for_write(db, &coll, &where_sql, &nql)?;
4226 let mut written: Vec<Value> = vec![];
4227 for row in &matched {
4228 let id = match row.get("_id").and_then(|v| v.as_str()) {
4229 Some(id) => id.to_string(),
4230 None => continue,
4231 };
4232 // Merge onto the CURRENT stored document, not onto the query
4233 // row: a query row carries injected `_`-prefixed metadata that
4234 // must never be written back into the payload.
4235 let mut doc = match db.get(&coll, &id) {
4236 Some(n) => match n.data {
4237 Value::Object(m) => m,
4238 _ => serde_json::Map::new(),
4239 },
4240 None => continue,
4241 };
4242 for (k, v) in &set {
4243 doc.insert(k.clone(), v.clone());
4244 }
4245 // An UPDATE is a NEW VERSION — the prior value stays readable
4246 // with AS OF SYSTEM TIME. That is the whole point.
4247 let node = db
4248 .put(&coll, &id, Value::Object(doc), vec![], None, None)
4249 .map_err(|e| err_msg("XX000", &format!("UPDATE failed: {}", e)))?;
4250 written.push(crate::nql::node_to_json(&node));
4251 }
4252 let n = written.len();
4253 let has_rows = wants_returning(stmt_sql);
4254 Ok(Executed {
4255 rows: if has_rows { written } else { vec![] },
4256 project: returning,
4257 has_rows,
4258 tag: format!("UPDATE {}", n),
4259 tag_counts_rows: false,
4260 })
4261 }
4262
4263 Stmt::Delete { coll, where_sql, nql, returning } => {
4264 let db = need_db!();
4265 need_write!();
4266 let matched = rows_for_write(db, &coll, &where_sql, &nql)?;
4267 // RETURNING must be captured BEFORE the delete: after the tombstone
4268 // the row is no longer readable by id.
4269 let returned = matched.clone();
4270 let mut n = 0usize;
4271 for row in &matched {
4272 if let Some(id) = row.get("_id").and_then(|v| v.as_str()) {
4273 match db.delete(&coll, id) {
4274 Ok(true) => n += 1,
4275 Ok(false) => {}
4276 Err(e) => return Err(err_msg("XX000", &format!("DELETE failed: {}", e))),
4277 }
4278 }
4279 }
4280 let has_rows = wants_returning(stmt_sql);
4281 Ok(Executed {
4282 rows: if has_rows { returned } else { vec![] },
4283 project: returning,
4284 has_rows,
4285 tag: format!("DELETE {}", n),
4286 tag_counts_rows: false,
4287 })
4288 }
4289 }
4290}
4291
4292/// Execute a simple-query payload, which may hold several `;`-separated statements.
4293fn run_simple_query(sql: &str, db_name: &str, db: Option<&Arc<Db>>, read_only: bool) -> Vec<u8> {
4294 let mut out = vec![];
4295 let statements = split_statements(sql);
4296 if statements.is_empty() {
4297 // EmptyQueryResponse
4298 return Out::msg(b'I').finish();
4299 }
4300 for stmt_sql in statements {
4301 match execute_stmt(&stmt_sql, db_name, db, read_only) {
4302 // Abandon the rest of the batch on the first error, as Postgres does.
4303 Err(encoded) => {
4304 out.extend_from_slice(&encoded);
4305 return out;
4306 }
4307 Ok(ex) => {
4308 if ex.has_rows {
4309 out.extend_from_slice(&encode_rows(&ex.rows, &ex.project));
4310 }
4311 out.extend_from_slice(&command_complete(&ex.tag_for(ex.rows.len())));
4312 }
4313 }
4314 }
4315 out
4316}
4317
4318/// Split on `;` at the top level, ignoring separators inside string literals.
4319fn split_statements(sql: &str) -> Vec<String> {
4320 let mut out = vec![];
4321 let mut cur = String::new();
4322 let mut in_s = false;
4323 for c in sql.chars() {
4324 match c {
4325 '\'' => { in_s = !in_s; cur.push(c); }
4326 ';' if !in_s => {
4327 if !cur.trim().is_empty() { out.push(cur.clone()); }
4328 cur.clear();
4329 }
4330 _ => cur.push(c),
4331 }
4332 }
4333 if !cur.trim().is_empty() {
4334 out.push(cur);
4335 }
4336 out
4337}
4338
4339/// Bind and serve the Postgres read endpoint until the process exits.
4340pub async fn run(host: &str, port: u16, resolver: Arc<dyn DbResolver>) -> anyhow::Result<()> {
4341 // Writes are ON by default — that is the parity position. An operator who
4342 // wants the "system of proof beside your database" deployment, where this
4343 // door must never mutate anything, sets NEDBD_PG_READ_ONLY=1.
4344 let read_only = std::env::var("NEDBD_PG_READ_ONLY")
4345 .map(|v| v == "1" || v.eq_ignore_ascii_case("true"))
4346 .unwrap_or(false);
4347 let listener = TcpListener::bind((host, port)).await?;
4348 println!(" pgwire postgres endpoint on {}:{} — psql / DBeaver / psycopg ({})",
4349 host, port,
4350 if read_only { "SELECT only — read-only mode" } else { "SELECT + INSERT/UPDATE/DELETE" });
4351 loop {
4352 let (sock, _peer) = match listener.accept().await {
4353 Ok(v) => v,
4354 Err(e) => {
4355 eprintln!(" [pgwire] accept failed: {}", e);
4356 continue;
4357 }
4358 };
4359 let r = Arc::clone(&resolver);
4360 tokio::spawn(async move {
4361 let _ = sock.set_nodelay(true);
4362 if let Err(e) = handle(sock, r, read_only).await {
4363 // A client disconnecting mid-message is routine, not an incident.
4364 if e.kind() != std::io::ErrorKind::UnexpectedEof
4365 && e.kind() != std::io::ErrorKind::ConnectionReset
4366 {
4367 eprintln!(" [pgwire] connection error: {}", e);
4368 }
4369 }
4370 });
4371 }
4372}
4373
4374// ─────────────────────────────────────────────────────────────────────────────
4375
4376#[cfg(test)]
4377mod explain_tests {
4378 use super::*;
4379
4380 #[test]
4381 fn a_bare_explain_is_stripped() {
4382 assert_eq!(strip_explain("EXPLAIN SELECT 1"), Some("SELECT 1"));
4383 assert_eq!(strip_explain("explain select 1"), Some("select 1"));
4384 assert_eq!(strip_explain(" EXPLAIN SELECT 1 ; "), Some("SELECT 1"));
4385 }
4386
4387 #[test]
4388 fn analyze_and_verbose_are_accepted_and_ignored() {
4389 // This endpoint always executes and always reports actual rows, so
4390 // EXPLAIN and EXPLAIN ANALYZE genuinely do the same thing. Accepting
4391 // the client's spelling beats refusing it.
4392 assert_eq!(strip_explain("EXPLAIN ANALYZE SELECT 1"), Some("SELECT 1"));
4393 assert_eq!(strip_explain("EXPLAIN ANALYSE SELECT 1"), Some("SELECT 1"));
4394 assert_eq!(strip_explain("EXPLAIN VERBOSE SELECT 1"), Some("SELECT 1"));
4395 assert_eq!(strip_explain("EXPLAIN ANALYZE VERBOSE SELECT 1"), Some("SELECT 1"));
4396 assert_eq!(strip_explain("explain analyze verbose select 1"), Some("select 1"));
4397 }
4398
4399 #[test]
4400 fn a_word_merely_starting_with_explain_is_not_a_keyword() {
4401 assert_eq!(strip_explain("EXPLAINED SELECT 1"), None);
4402 assert_eq!(strip_explain("SELECT 1"), None);
4403 assert_eq!(strip_explain("SELECT explain FROM t"), None);
4404 }
4405
4406 #[test]
4407 fn a_column_named_analyze_is_not_eaten() {
4408 // `analyzed` merely starts with the keyword; the word boundary check
4409 // is what stops it being consumed as an option.
4410 assert_eq!(strip_explain("EXPLAIN analyzed_view"), Some("analyzed_view"));
4411 }
4412
4413 #[test]
4414 fn the_plan_result_has_postgres_shape() {
4415 let e = plan_result(vec!["Seq Scan on t".into(), "note".into()]);
4416 assert_eq!(e.project.len(), 1);
4417 assert_eq!(e.project[0].out, "QUERY PLAN");
4418 assert_eq!(e.rows.len(), 2);
4419 assert_eq!(e.rows[0]["QUERY PLAN"], "Seq Scan on t");
4420 assert_eq!(e.tag, "EXPLAIN");
4421 // EXPLAIN's tag carries no row count in PostgreSQL.
4422 assert!(!e.tag_counts_rows);
4423 }
4424}
4425
4426#[cfg(test)]
4427mod tests {
4428 use super::*;
4429 use serde_json::json;
4430
4431 fn q(sql: &str) -> String {
4432 match translate(sql) {
4433 Ok(Stmt::Query { nql, .. }) => nql,
4434 other => panic!("expected a query for {:?}, got {:?}", sql, other),
4435 }
4436 }
4437 /// Output column names, in order.
4438 fn proj(sql: &str) -> Vec<String> {
4439 match translate(sql) {
4440 Ok(Stmt::Query { project, .. }) => project.iter().map(|c| c.out.clone()).collect(),
4441 other => panic!("expected a query for {:?}, got {:?}", sql, other),
4442 }
4443 }
4444 /// (source key, output name) pairs, for the aggregate renaming.
4445 fn proj_pairs(sql: &str) -> Vec<(String, String)> {
4446 match translate(sql) {
4447 Ok(Stmt::Query { project, .. }) =>
4448 project.iter().map(|c| (c.src.clone(), c.out.clone())).collect(),
4449 other => panic!("expected a query for {:?}, got {:?}", sql, other),
4450 }
4451 }
4452 fn names(cols: &[Col]) -> Vec<String> { cols.iter().map(|c| c.out.clone()).collect() }
4453
4454 /// The full projection, so a test can assert the SRC and the OUT
4455 /// separately — they are different jobs and conflating them is how an
4456 /// alias got lost.
4457 fn cols_of(sql: &str) -> Vec<Col> {
4458 match translate(sql).unwrap() {
4459 Stmt::Query { project, .. } => project,
4460 other => panic!("{:?}", other),
4461 }
4462 }
4463
4464 #[test]
4465 fn select_star_becomes_bare_from() {
4466 assert_eq!(q("SELECT * FROM orders"), "FROM orders");
4467 assert_eq!(q("select * from orders;"), "FROM orders");
4468 assert_eq!(proj("SELECT * FROM orders"), Vec::<String>::new());
4469 }
4470
4471 #[test]
4472 fn a_column_list_becomes_a_projection_not_a_clause() {
4473 // NQL has no projection, so the column list is carried separately and
4474 // applied to the returned rows.
4475 assert_eq!(q("SELECT status, total FROM orders"), "FROM orders");
4476 assert_eq!(proj("SELECT status, total FROM orders"), vec!["status", "total"]);
4477 }
4478
4479 #[test]
4480 fn a_qualifier_reduces_to_the_field_while_an_ALIAS_is_the_name_the_client_sees() {
4481 // Two different jobs, and they used to be conflated. The SRC is what
4482 // NEDB reads out of the row, so a qualifier must be stripped from it.
4483 // The OUT is the name the CLIENT looks the column up by, so an alias
4484 // must be KEPT in it — `SELECT status AS s` returns a column called
4485 // `s`, and answering with one called `status` hands a client a result
4486 // it cannot find. SQLAlchemy writes `count(*) AS count_1` and then
4487 // reads `count_1`.
4488 let cols = cols_of("SELECT o.status AS s, o.total total, o.region FROM orders o");
4489 assert_eq!(cols.iter().map(|c| c.src.clone()).collect::<Vec<_>>(),
4490 vec!["status", "total", "region"]);
4491 assert_eq!(cols.iter().map(|c| c.out.clone()).collect::<Vec<_>>(),
4492 vec!["s", "total", "region"]);
4493 assert_eq!(q("SELECT * FROM public.orders"), "FROM orders");
4494 assert_eq!(q("SELECT * FROM \"orders\""), "FROM orders");
4495 }
4496
4497 #[test]
4498 fn a_select_list_may_MIX_columns_with_an_aggregate() {
4499 // What a GROUP BY query actually looks like. The previous parser
4500 // refused any list containing a parenthesis, so this whole shape was
4501 // unreachable even though NQL expresses it natively — and it is the
4502 // single most common grouped query an ORM emits.
4503 // The aggregate sits IMMEDIATELY AFTER the group key — verified
4504 // against the running engine, which refuses the other order with
4505 // "only one aggregate per query".
4506 assert_eq!(q("SELECT status, count(*) AS count_1 FROM orders GROUP BY status"),
4507 "FROM orders GROUP BY status COUNT");
4508 // SQL puts GROUP BY before ORDER BY / LIMIT; the aggregate still lands
4509 // on the key, and the rest of the tail follows.
4510 assert_eq!(q("SELECT status, count(*) FROM orders WHERE total > 1 GROUP BY status ORDER BY status LIMIT 5"),
4511 "FROM orders WHERE total > 1 GROUP BY status COUNT ORDER BY status LIMIT 5");
4512 // A bare aggregate with NO grouping still goes after the collection.
4513 assert_eq!(q("SELECT count(*) FROM orders"), "FROM orders COUNT");
4514 assert_eq!(q("SELECT sum(total) FROM orders"), "FROM orders SUM total");
4515 // More than one group key is refused by name: NQL groups by a single
4516 // field, and using only the first would aggregate over rows the query
4517 // meant to keep apart.
4518 let e = translate("SELECT status, count(*) FROM orders GROUP BY status, region").unwrap_err();
4519 assert!(e.contains("GROUP BY takes one key"), "{}", e);
4520 let cols = cols_of("SELECT status, count(*) AS count_1 FROM orders GROUP BY status");
4521 assert_eq!(cols.iter().map(|c| c.src.clone()).collect::<Vec<_>>(),
4522 vec!["status", "count"]);
4523 assert_eq!(cols.iter().map(|c| c.out.clone()).collect::<Vec<_>>(),
4524 vec!["status", "count_1"]);
4525
4526 // A named aggregate rides along with `count`, because an NQL grouped
4527 // row carries both.
4528 let cols = cols_of("SELECT status, count(*), sum(total) FROM orders GROUP BY status");
4529 assert_eq!(cols.iter().map(|c| c.src.clone()).collect::<Vec<_>>(),
4530 vec!["status", "count", "sum_total"]);
4531 assert_eq!(q("SELECT status, count(*), sum(total) FROM orders GROUP BY status"),
4532 "FROM orders GROUP BY status SUM total");
4533
4534 // A qualifier on the aggregate's column is stripped like any other.
4535 assert_eq!(q("SELECT o.status, sum(o.total) FROM orders o GROUP BY o.status"),
4536 "FROM orders GROUP BY status SUM total");
4537
4538 // Two NAMED aggregates cannot both be carried, and that is refused by
4539 // name rather than silently dropping one.
4540 let e = translate("SELECT status, sum(total), avg(total) FROM orders GROUP BY status")
4541 .unwrap_err();
4542 assert!(e.contains("only one of SUM/AVG/MIN/MAX"), "{}", e);
4543
4544 // A column that is neither a key nor an aggregate is still refused.
4545 let e = translate("SELECT status, total, count(*) FROM orders GROUP BY status")
4546 .unwrap_err();
4547 assert!(e.contains("must appear in the GROUP BY clause"), "{}", e);
4548 }
4549
4550 #[test]
4551 fn ORDER_BY_an_ordinal_resolves_to_that_select_list_column() {
4552 // SQL lets a sort key be a POSITION, and clients write it constantly.
4553 // NQL has no ordinals — it read the `1` as a literal and refused with
4554 // "expected field name, got Num(1.0)". node-postgres sent
4555 // `GROUP BY status ORDER BY 1` in the harness's first run.
4556 assert_eq!(q("SELECT status, total FROM orders ORDER BY 1"),
4557 "FROM orders ORDER BY status");
4558 assert_eq!(q("SELECT status, total FROM orders ORDER BY 2 DESC"),
4559 "FROM orders ORDER BY total DESC");
4560 // Several keys, mixing ordinals with names, and a direction on each.
4561 assert_eq!(q("SELECT status, total FROM orders ORDER BY 2 DESC, 1"),
4562 "FROM orders ORDER BY total DESC, status");
4563 assert_eq!(q("SELECT status, total FROM orders ORDER BY 1, total DESC"),
4564 "FROM orders ORDER BY status, total DESC");
4565 // An ordinal survives the GROUP BY splice, and resolves to the group
4566 // key rather than to the literal 1 — which is the exact shape that
4567 // failed in CI.
4568 assert_eq!(q("SELECT status, count(*) AS n FROM orders GROUP BY status ORDER BY 1"),
4569 "FROM orders GROUP BY status COUNT ORDER BY status");
4570 // An ordinal may name the AGGREGATE column too.
4571 assert_eq!(q("SELECT status, count(*) AS n FROM orders GROUP BY status ORDER BY 2 DESC"),
4572 "FROM orders GROUP BY status COUNT ORDER BY count DESC");
4573 // The clause boundary is respected: a following LIMIT is not swallowed
4574 // into the sort list, and `LIMIT 1` is not mistaken for an ordinal.
4575 assert_eq!(q("SELECT status, total FROM orders ORDER BY 2 LIMIT 1"),
4576 "FROM orders ORDER BY total LIMIT 1");
4577 // A `1` anywhere else stays a literal.
4578 assert_eq!(q("SELECT status FROM orders WHERE total > 1 ORDER BY 1"),
4579 "FROM orders WHERE total > 1 ORDER BY status");
4580
4581 // Out of range, and `SELECT *` where there is no list to index, are
4582 // both refused with the reason — guessing a column would sort by
4583 // something the query never named.
4584 let e = translate("SELECT status FROM orders ORDER BY 4").unwrap_err();
4585 assert!(e.contains("out of range") && e.contains("1 column"), "{}", e);
4586 let e = translate("SELECT * FROM orders ORDER BY 1").unwrap_err();
4587 assert!(e.contains("no list to index"), "{}", e);
4588 }
4589
4590 #[test]
4591 fn count_of_a_subquery_flattens_only_when_the_two_counts_MUST_agree() {
4592 // `.count()` in every ORM wraps the whole query in a derived table.
4593 // Counting rows that ARE the inner query's rows is counting the inner
4594 // query, so this is an identity, not an approximation.
4595 assert_eq!(
4596 q("SELECT count(*) AS count_1 FROM (SELECT orders._id AS a, orders.status AS b \
4597 FROM orders WHERE orders.status = 'paid') AS anon_1"),
4598 // Verified against the running engine: with no GROUP BY the
4599 // aggregate may sit either side of WHERE and answers identically.
4600 r#"FROM orders COUNT WHERE status = "paid""#);
4601 // No predicate at all.
4602 assert_eq!(q("SELECT count(*) FROM (SELECT orders._id FROM orders) AS anon_1"),
4603 "FROM orders COUNT");
4604 // ORDER BY cannot change a count, so it is dropped rather than refused.
4605 assert_eq!(q("SELECT count(*) FROM (SELECT _id FROM orders ORDER BY total DESC) AS a"),
4606 "FROM orders COUNT");
4607 // The outer alias is the name the client reads the column back by.
4608 let cols = cols_of("SELECT count(*) AS count_1 FROM (SELECT _id FROM orders) AS a");
4609 assert_eq!(cols[0].src, "count");
4610 assert_eq!(cols[0].out, "count_1");
4611
4612 // Each guard is a construct that would make the two counts DIFFERENT
4613 // numbers, so each is refused rather than silently flattened.
4614 for sql in [
4615 // LIMIT / OFFSET cap the rows before they are counted
4616 "SELECT count(*) FROM (SELECT _id FROM orders LIMIT 1) AS a",
4617 "SELECT count(*) FROM (SELECT _id FROM orders OFFSET 1) AS a",
4618 // the inner rows ARE the groups
4619 "SELECT count(*) FROM (SELECT status FROM orders GROUP BY status) AS a",
4620 // an inner aggregate already reduced the rows to one
4621 "SELECT count(*) FROM (SELECT count(*) FROM orders) AS a",
4622 "SELECT count(*) FROM (SELECT sum(total) FROM orders) AS a",
4623 // the outer list would need the derived table's own columns
4624 "SELECT count(*), status FROM (SELECT status FROM orders) AS a",
4625 "SELECT status FROM (SELECT status FROM orders) AS a",
4626 // one level is the claim
4627 "SELECT count(*) FROM (SELECT x FROM (SELECT _id AS x FROM orders) AS b) AS a",
4628 ] {
4629 let e = translate(sql).unwrap_err();
4630 assert!(e.contains("subqueries in FROM"), "{} -> {}", sql, e);
4631 }
4632
4633 // DISTINCT and the set operators are caught EARLIER, by their own
4634 // rules, which scan the whole statement before the FROM list is even
4635 // read. Asserted separately so the test records which check owns each
4636 // refusal rather than implying one catch-all does.
4637 for (sql, needle) in [
4638 ("SELECT count(*) FROM (SELECT DISTINCT status FROM orders) AS a", "DISTINCT"),
4639 ("SELECT count(*) FROM (SELECT a FROM t UNION SELECT b FROM u) AS x", "UNION"),
4640 ] {
4641 let e = translate(sql).unwrap_err();
4642 assert!(e.contains(needle), "{} -> {}", sql, e);
4643 }
4644 }
4645
4646 #[test]
4647 fn a_QUALIFIED_column_in_WHERE_finds_its_field_instead_of_ZERO_ROWS() {
4648 // THE silent wrong answer. NQL looks a field up FLAT, so
4649 // `WHERE orders.status = 'paid'` asked for a field literally named
4650 // "orders.status", no document had one, and the query returned ZERO
4651 // ROWS with no error — an empty result that reads exactly like "you
4652 // have no paid orders". Every ORM qualifies its predicates, so every
4653 // filtered SQLAlchemy query answered empty and `.get(pk)` answered
4654 // None.
4655 assert_eq!(q("SELECT _id FROM orders WHERE orders.status = 'paid'"),
4656 r#"FROM orders WHERE status = "paid""#);
4657 assert_eq!(q("SELECT _id FROM orders WHERE orders.total > 50"),
4658 "FROM orders WHERE total > 50");
4659 // Every clause in the tail, not just WHERE.
4660 assert_eq!(q("SELECT _id FROM orders ORDER BY orders.total DESC LIMIT 2"),
4661 "FROM orders ORDER BY total DESC LIMIT 2");
4662 assert_eq!(q("SELECT status, count(*) FROM orders GROUP BY orders.status"),
4663 "FROM orders GROUP BY status COUNT");
4664
4665 // An alias is a legal qualifier and is accepted as one. It is also
4666 // REMOVED from the tail, because NQL has no alias syntax and reported
4667 // an "unexpected token" on it.
4668 assert_eq!(q("SELECT o.status FROM orders o WHERE o.status = 'paid'"),
4669 r#"FROM orders WHERE status = "paid""#);
4670 assert_eq!(q("SELECT o.status FROM orders AS o WHERE o.total > 1"),
4671 "FROM orders WHERE total > 1");
4672
4673 // A qualifier naming NEITHER the collection nor its alias is an
4674 // ERROR, not a strip. Stripping it would answer from the one relation
4675 // that IS present, which is a different wrong answer in the same
4676 // empty-looking clothes.
4677 let e = translate("SELECT _id FROM orders WHERE nosuch.status = 'paid'").unwrap_err();
4678 assert!(e.contains("no table or alias named \"nosuch\""), "{}", e);
4679 let e = translate("SELECT _id FROM orders o WHERE p.status = 'paid'").unwrap_err();
4680 assert!(e.contains("aliased \"o\""), "the message names the alias in scope: {}", e);
4681
4682 // A dot INSIDE a literal is data, not a qualifier.
4683 assert_eq!(q("SELECT _id FROM orders WHERE status = 'pa.id'"),
4684 r#"FROM orders WHERE status = "pa.id""#);
4685 // ...and a decimal point is not one either.
4686 assert_eq!(q("SELECT _id FROM orders WHERE total > 1.5"),
4687 "FROM orders WHERE total > 1.5");
4688
4689 // UPDATE and DELETE carry the same tail, and had the same bug.
4690 match translate("UPDATE orders o SET status = 'x' WHERE o.total > 5").unwrap() {
4691 Stmt::Update { coll, nql, .. } => {
4692 assert_eq!(coll, "orders", "the alias is not part of the collection name");
4693 assert_eq!(nql, "FROM orders WHERE total > 5");
4694 }
4695 other => panic!("{:?}", other),
4696 }
4697 match translate("DELETE FROM orders o WHERE o.status = 'paid'").unwrap() {
4698 Stmt::Delete { coll, nql, .. } => {
4699 assert_eq!(coll, "orders");
4700 assert_eq!(nql, r#"FROM orders WHERE status = "paid""#);
4701 }
4702 other => panic!("{:?}", other),
4703 }
4704
4705 // `AS OF SYSTEM TIME` also begins with AS and is NOT an alias.
4706 assert_eq!(q("SELECT _id FROM orders AS OF SYSTEM TIME 3 WHERE orders.total > 1"),
4707 "FROM orders AS OF 3 WHERE total > 1");
4708 }
4709
4710 #[test]
4711 fn where_clauses_pass_through_with_sql_literals_rewritten() {
4712 assert_eq!(q("SELECT * FROM orders WHERE status = 'paid'"),
4713 r#"FROM orders WHERE status = "paid""#);
4714 assert_eq!(q("SELECT * FROM orders WHERE status <> 'paid'"),
4715 r#"FROM orders WHERE status != "paid""#);
4716 assert_eq!(q("SELECT * FROM orders WHERE status IN ('paid','open')"),
4717 r#"FROM orders WHERE status IN ("paid","open")"#);
4718 }
4719
4720 /// SQL escapes an embedded quote by doubling it. That must become ONE
4721 /// character inside the NQL string, not terminate it.
4722 #[test]
4723 fn a_doubled_sql_quote_is_one_literal_character() {
4724 assert_eq!(q("SELECT * FROM t WHERE name = 'it''s'"),
4725 r#"FROM t WHERE name = "it's""#);
4726 }
4727
4728 /// A double quote inside a SQL literal has to be escaped for NQL, whose
4729 /// lexer collapses \" — otherwise it would close the string early.
4730 #[test]
4731 fn a_double_quote_inside_a_sql_literal_is_escaped_for_nql() {
4732 assert_eq!(q(r#"SELECT * FROM t WHERE name = 'say "hi"'"#),
4733 r#"FROM t WHERE name = "say \"hi\"""#);
4734 }
4735
4736 #[test]
4737 fn the_shared_clauses_are_handed_to_nql_unchanged() {
4738 assert_eq!(q("SELECT * FROM orders ORDER BY total DESC LIMIT 10 OFFSET 5"),
4739 "FROM orders ORDER BY total DESC LIMIT 10 OFFSET 5");
4740 assert_eq!(q("SELECT * FROM orders GROUP BY region"), "FROM orders GROUP BY region");
4741 assert_eq!(q("SELECT * FROM o WHERE total BETWEEN 1 AND 9 ORDER BY a, b DESC"),
4742 "FROM o WHERE total BETWEEN 1 AND 9 ORDER BY a, b DESC");
4743 }
4744
4745 /// An aggregate must surface as ONE column, named as SQL names it.
4746 ///
4747 /// NQL answers `SUM(total)` with `{count, sum_total, value}` — `value`
4748 /// being a back-compat alias. Passing that straight through gave
4749 /// `SELECT COUNT(*)` two columns (`count`, `value`) where SQL promises
4750 /// one, and leaked an internal key name onto the wire.
4751 #[test]
4752 fn an_aggregate_is_one_column_named_as_sql_names_it() {
4753 assert_eq!(proj_pairs("SELECT COUNT(*) FROM orders"),
4754 vec![("count".to_string(), "count".to_string())]);
4755 assert_eq!(proj_pairs("SELECT SUM(total) FROM orders"),
4756 vec![("sum_total".to_string(), "sum".to_string())]);
4757 assert_eq!(proj_pairs("SELECT avg(total) FROM orders"),
4758 vec![("avg_total".to_string(), "avg".to_string())]);
4759 assert_eq!(proj_pairs("SELECT MIN(total) FROM orders"),
4760 vec![("min_total".to_string(), "min".to_string())]);
4761 // And the encoded result really is one column with that name.
4762 let rows = vec![json!({"count": 4, "sum_total": 420, "value": 420})];
4763 let p = vec![Col::renamed("sum_total", "sum")];
4764 let cols = columns_for(&rows, &p);
4765 assert_eq!(names(&cols), vec!["sum"], "one column, SQL's name");
4766 assert_eq!(cell(rows[0].get(&cols[0].src)), Some("420".to_string()));
4767 }
4768
4769 /// A grouped NQL row holds the group key, `count` and the aggregate —
4770 /// nothing else. Projecting another column found nothing and rendered
4771 /// NULL, which is a silent wrong answer. Postgres errors; so do we, in
4772 /// Postgres's own words.
4773 #[test]
4774 fn a_bare_column_with_group_by_is_refused_not_nulled() {
4775 let e = translate("SELECT region, total FROM orders GROUP BY region").unwrap_err();
4776 assert!(e.contains("must appear in the GROUP BY clause"), "{}", e);
4777 assert!(e.contains("total"), "the message names the offending column: {}", e);
4778
4779 // The group key itself, and `count`, are both legitimate.
4780 assert!(translate("SELECT region FROM orders GROUP BY region").is_ok());
4781 assert!(translate("SELECT region, count FROM orders GROUP BY region").is_ok());
4782 // As is an aggregate over the grouped set.
4783 assert!(translate("SELECT SUM(total) FROM orders GROUP BY region").is_ok());
4784 // And `*` is unaffected — it returns whatever the grouped row holds.
4785 assert!(translate("SELECT * FROM orders GROUP BY region").is_ok());
4786 }
4787
4788 #[test]
4789 fn count_star_becomes_nql_count() {
4790 assert_eq!(q("SELECT COUNT(*) FROM orders"), "FROM orders COUNT");
4791 assert_eq!(q("SELECT count(*) FROM orders WHERE total > 5"),
4792 "FROM orders COUNT WHERE total > 5");
4793 }
4794
4795 #[test]
4796 fn aggregates_carry_their_target_column() {
4797 assert_eq!(q("SELECT SUM(total) FROM orders"), "FROM orders SUM total");
4798 assert_eq!(q("SELECT avg(total) FROM orders WHERE region = 'eu'"),
4799 r#"FROM orders AVG total WHERE region = "eu""#);
4800 assert!(translate("SELECT SUM(*) FROM orders").is_err());
4801 }
4802
4803 /// The bridge worth having: Postgres spells time travel
4804 /// `AS OF SYSTEM TIME`, and NEDB's is sequence-addressed and permanent.
4805 #[test]
4806 fn as_of_system_time_bridges_to_nql_as_of() {
4807 assert_eq!(q("SELECT * FROM orders AS OF SYSTEM TIME 42"),
4808 "FROM orders AS OF 42");
4809 assert_eq!(q("SELECT * FROM orders AS OF SYSTEM TIME 42 WHERE total > 1"),
4810 "FROM orders AS OF 42 WHERE total > 1");
4811 // A wall-clock timestamp is refused with the reason, not silently ignored.
4812 let e = translate("SELECT * FROM orders AS OF SYSTEM TIME '2026-01-01'").unwrap_err();
4813 assert!(e.contains("sequence number"), "{}", e);
4814 }
4815
4816 /// A select-list item that is not a column reference must be REFUSED, not
4817 /// turned into a field name.
4818 ///
4819 /// The guard used to be `expr.contains('(')`, which only catches expressions
4820 /// that happen to have a paren. `total * 2` sailed through, became the field
4821 /// name "total * 2", matched no document, and the column came back EMPTY for
4822 /// every row with no error. Same silent class as the qualified-WHERE bug: a
4823 /// wrong answer wearing the shape of data.
4824 #[test]
4825 fn a_select_list_expression_is_refused_rather_than_answered_blank() {
4826 for sql in [
4827 "SELECT total * 2 FROM orders",
4828 "SELECT total, total*2 AS doubled FROM orders",
4829 "SELECT total + 1 FROM orders",
4830 "SELECT status || 'x' FROM orders",
4831 "SELECT -total FROM orders",
4832 "SELECT lower(status) FROM orders",
4833 ] {
4834 let e = translate(sql).unwrap_err();
4835 assert!(e.contains("expressions in the select list"), "{} -> {}", sql, e);
4836 }
4837 // ...and the things that ARE column references still pass, or the fix
4838 // would have bought correctness by refusing everything.
4839 assert_eq!(q("SELECT _id, status FROM orders"), "FROM orders");
4840 assert_eq!(q("SELECT \"status\" FROM orders"), "FROM orders");
4841 assert_eq!(q("SELECT orders.status FROM orders"), "FROM orders");
4842 assert_eq!(q("SELECT o.status FROM orders o"), "FROM orders");
4843 assert_eq!(q("SELECT total AS t FROM orders"), "FROM orders");
4844 assert!(translate("SELECT count(*) FROM orders").is_ok());
4845 assert!(translate("SELECT sum(total) FROM orders").is_ok());
4846 }
4847
4848 /// HAVING has to reach NQL in the spelling NQL's grouped row actually uses.
4849 ///
4850 /// An NQL grouped row carries `count` and `<agg>_<field>`. SQL clients write
4851 /// `count(*)`, or the alias they gave it. `count(*)` failed LOUDLY (fine),
4852 /// but `COUNT` and an alias both passed through verbatim and answered ZERO
4853 /// ROWS — which reads as "no groups qualified" rather than "your predicate
4854 /// named a field that does not exist".
4855 #[test]
4856 fn having_is_translated_to_nqls_spelling_and_refuses_an_unknown_key() {
4857 // Every spelling a client might send for the count.
4858 for sql in [
4859 "SELECT status, count(*) AS n FROM orders GROUP BY status HAVING count(*) > 1",
4860 "SELECT status, count(*) AS n FROM orders GROUP BY status HAVING n > 1",
4861 "SELECT status, count(*) FROM orders GROUP BY status HAVING COUNT > 1",
4862 "SELECT status, count(*) FROM orders GROUP BY status HAVING count > 1",
4863 ] {
4864 let got = q(sql);
4865 assert_eq!(got, "FROM orders GROUP BY status COUNT HAVING count > 1",
4866 "{} -> {}", sql, got);
4867 }
4868 // A named aggregate, by its alias -- NQL calls the field `sum_total`.
4869 assert_eq!(q("SELECT status, sum(total) AS s FROM orders GROUP BY status HAVING s > 100"),
4870 "FROM orders GROUP BY status SUM total HAVING sum_total > 100");
4871 // ...and by NQL's own name for it, which must not be rewritten twice.
4872 assert_eq!(q("SELECT status, sum(total) FROM orders GROUP BY status HAVING sum_total > 100"),
4873 "FROM orders GROUP BY status SUM total HAVING sum_total > 100");
4874 // Filtering on the group key itself is legitimate and passes through
4875 // untouched -- the SQL literal becomes an NQL one, as everywhere else.
4876 assert_eq!(q("SELECT status, count(*) FROM orders GROUP BY status HAVING status > 'a'"),
4877 "FROM orders GROUP BY status COUNT HAVING status > \"a\"");
4878 // A key the grouped row cannot carry is an ERROR, not zero rows.
4879 let e = translate(
4880 "SELECT status, count(*) FROM orders GROUP BY status HAVING nosuch > 1").unwrap_err();
4881 assert!(e.contains("HAVING names") && e.contains("nosuch"), "{}", e);
4882 assert!(e.contains("zero rows"), "the message must say what it prevented: {}", e);
4883 }
4884
4885 #[test]
4886 fn handshake_queries_are_answered_so_clients_can_connect() {
4887 assert!(matches!(translate("SELECT version()"), Ok(Stmt::Canned { .. })));
4888 assert!(matches!(translate("SHOW transaction_isolation"), Ok(Stmt::Canned { .. })));
4889 assert!(matches!(translate("SELECT current_schema()"), Ok(Stmt::Canned { .. })));
4890 assert!(matches!(translate("SET extra_float_digits = 3"), Ok(Stmt::Ok(_))));
4891 assert!(matches!(translate("BEGIN"), Ok(Stmt::Ok(_))));
4892 assert!(matches!(translate(""), Ok(Stmt::Ok(_))));
4893 }
4894
4895 /// Every refusal has to name the boundary. "Syntax error" would send a
4896 /// developer hunting for a typo that is not there.
4897 #[test]
4898 fn unsupported_sql_is_refused_with_a_reason() {
4899 for (sql, expect) in [
4900 ("INSERT INTO t VALUES (1)", "explicit column list"),
4901 ("CREATE TABLE t (a int)", "DDL"),
4902 ("TRUNCATE t", "append-only"),
4903 ("GRANT ALL ON t TO x", "privilege system"),
4904 ("SELECT * FROM a JOIN b ON a.x = b.x", "JOIN is not supported"),
4905 ("SELECT * FROM a UNION SELECT * FROM b", "UNION"),
4906 ("SELECT DISTINCT region FROM orders", "GROUP BY"),
4907 ("SELECT * FROM (SELECT 1) x", "subqueries in FROM"),
4908 ("SELECT * FROM a, b", "more than one collection"),
4909 ("SELECT lower(status) FROM orders", "expressions in the select list"),
4910 ("VACUUM", "only SELECT"),
4911 ] {
4912 let e = translate(sql).unwrap_err();
4913 assert!(e.contains(expect), "for {:?} expected {:?} in {:?}", sql, expect, e);
4914 }
4915 }
4916
4917 // ── writes ───────────────────────────────────────────────────────────────
4918 //
4919 // SQL's write semantics and NEDB's append-only model line up: INSERT is a
4920 // put, UPDATE is a new version, DELETE is a tombstone. These tests pin the
4921 // parse; tests/test_pgwire.py proves the behaviour against a live server,
4922 // including that the PRIOR value is still readable afterwards.
4923
4924 fn ins(sql: &str) -> (String, Vec<InsertRow>, Vec<Col>) {
4925 match translate(sql) {
4926 Ok(Stmt::Insert { coll, rows, returning }) => (coll, rows, returning),
4927 other => panic!("expected INSERT for {:?}, got {:?}", sql, other),
4928 }
4929 }
4930
4931 #[test]
4932 fn insert_becomes_a_put_per_row() {
4933 let (coll, rows, ret) = ins("INSERT INTO orders (_id, status, total) VALUES ('o1', 'paid', 120)");
4934 assert_eq!(coll, "orders");
4935 assert_eq!(rows.len(), 1);
4936 assert_eq!(rows[0].id.as_deref(), Some("o1"));
4937 assert_eq!(rows[0].doc.get("status"), Some(&json!("paid")));
4938 assert_eq!(rows[0].doc.get("total"), Some(&json!(120)));
4939 // `_id` is the key, not a payload field.
4940 assert!(!rows[0].doc.contains_key("_id"));
4941 assert!(ret.is_empty());
4942 }
4943
4944 #[test]
4945 fn a_multi_row_insert_yields_one_row_each() {
4946 let (_, rows, _) = ins(
4947 "INSERT INTO t (id, n) VALUES ('a', 1), ('b', 2), ('c', 3)");
4948 assert_eq!(rows.len(), 3);
4949 assert_eq!(rows[1].id.as_deref(), Some("b"));
4950 assert_eq!(rows[2].doc.get("n"), Some(&json!(3)));
4951 }
4952
4953 #[test]
4954 fn an_insert_without_an_id_column_lets_the_server_assign_one() {
4955 let (_, rows, _) = ins("INSERT INTO t (n) VALUES (1)");
4956 assert_eq!(rows[0].id, None, "the executor mints a unique key");
4957 assert_eq!(rows[0].doc.get("n"), Some(&json!(1)));
4958 }
4959
4960 /// Provenance is reachable from SQL, not only from the HTTP API — which is
4961 /// the point of having writes here at all.
4962 #[test]
4963 fn insert_lifts_provenance_out_of_reserved_columns() {
4964 let (_, rows, _) = ins(
4965 "INSERT INTO audit (_id, _caused_by, _valid_from, kind) \
4966 VALUES ('e1', 'abc123', '2026-01-01', 'reprice')");
4967 assert_eq!(rows[0].caused_by, vec!["abc123".to_string()]);
4968 assert_eq!(rows[0].valid_from.as_deref(), Some("2026-01-01"));
4969 assert_eq!(rows[0].doc.get("kind"), Some(&json!("reprice")));
4970 // None of the reserved names leak into the stored payload.
4971 for k in ["_id", "_caused_by", "_valid_from"] {
4972 assert!(!rows[0].doc.contains_key(k), "{} leaked into the doc", k);
4973 }
4974 }
4975
4976 #[test]
4977 fn insert_values_cover_the_scalar_types() {
4978 let (_, rows, _) = ins(
4979 "INSERT INTO t (s, i, f, b, n) VALUES ('x', 42, 1.5, TRUE, NULL)");
4980 assert_eq!(rows[0].doc.get("s"), Some(&json!("x")));
4981 assert_eq!(rows[0].doc.get("i"), Some(&json!(42)));
4982 assert_eq!(rows[0].doc.get("f"), Some(&json!(1.5)));
4983 assert_eq!(rows[0].doc.get("b"), Some(&json!(true)));
4984 assert_eq!(rows[0].doc.get("n"), Some(&Value::Null));
4985 }
4986
4987 /// A doubled '' is one literal quote, and a comma inside a string is not a
4988 /// value separator.
4989 #[test]
4990 fn insert_literals_survive_quotes_and_commas() {
4991 let (_, rows, _) = ins("INSERT INTO t (a, b) VALUES ('it''s', 'x,y')");
4992 assert_eq!(rows[0].doc.get("a"), Some(&json!("it's")));
4993 assert_eq!(rows[0].doc.get("b"), Some(&json!("x,y")));
4994 }
4995
4996 #[test]
4997 fn insert_refuses_what_it_cannot_store_faithfully() {
4998 // An unevaluated expression stored as text would be a wrong value.
4999 assert!(translate("INSERT INTO t (a) VALUES (1 + 1)").is_err());
5000 assert!(translate("INSERT INTO t (a) VALUES (now())").is_err());
5001 // Column/value count mismatch.
5002 let e = translate("INSERT INTO t (a, b) VALUES (1)").unwrap_err();
5003 assert!(e.contains("values for"), "{}", e);
5004 // No column list at all.
5005 let e2 = translate("INSERT INTO t VALUES (1)").unwrap_err();
5006 assert!(e2.contains("explicit column list"), "{}", e2);
5007 }
5008
5009 #[test]
5010 fn update_finds_rows_with_the_full_predicate_surface() {
5011 match translate("UPDATE orders SET status = 'void' WHERE total < 50 AND region IN ('eu')") {
5012 Ok(Stmt::Update { coll, set, nql, .. }) => {
5013 assert_eq!(coll, "orders");
5014 assert_eq!(set, vec![("status".to_string(), json!("void"))]);
5015 // The WHERE became ordinary NQL, so IN/BETWEEN/LIKE all work.
5016 assert_eq!(nql, r#"FROM orders WHERE total < 50 AND region IN ("eu")"#);
5017 }
5018 other => panic!("expected UPDATE, got {:?}", other),
5019 }
5020 }
5021
5022 #[test]
5023 fn update_without_where_targets_the_whole_collection() {
5024 // Postgres allows it, so parity allows it.
5025 match translate("UPDATE t SET a = 1") {
5026 Ok(Stmt::Update { nql, .. }) => assert_eq!(nql, "FROM t"),
5027 other => panic!("expected UPDATE, got {:?}", other),
5028 }
5029 }
5030
5031 #[test]
5032 fn update_handles_several_assignments() {
5033 match translate("UPDATE t SET a = 1, b = 'x,y', c = NULL WHERE id = 'k'") {
5034 Ok(Stmt::Update { set, .. }) => {
5035 assert_eq!(set.len(), 3);
5036 assert_eq!(set[1], ("b".to_string(), json!("x,y")));
5037 assert_eq!(set[2], ("c".to_string(), Value::Null));
5038 }
5039 other => panic!("expected UPDATE, got {:?}", other),
5040 }
5041 assert!(translate("UPDATE t SET").is_err());
5042 assert!(translate("UPDATE t SET a").is_err());
5043 }
5044
5045 #[test]
5046 fn delete_becomes_a_predicate_over_the_collection() {
5047 match translate("DELETE FROM orders WHERE status = 'void'") {
5048 Ok(Stmt::Delete { coll, nql, .. }) => {
5049 assert_eq!(coll, "orders");
5050 assert_eq!(nql, r#"FROM orders WHERE status = "void""#);
5051 }
5052 other => panic!("expected DELETE, got {:?}", other),
5053 }
5054 match translate("DELETE FROM t") {
5055 Ok(Stmt::Delete { nql, .. }) => assert_eq!(nql, "FROM t"),
5056 other => panic!("expected DELETE, got {:?}", other),
5057 }
5058 }
5059
5060 #[test]
5061 fn returning_is_parsed_off_every_write() {
5062 let (_, _, ret) = ins("INSERT INTO t (a) VALUES (1) RETURNING a, _id");
5063 assert_eq!(ret.iter().map(|c| c.out.clone()).collect::<Vec<_>>(), vec!["a", "_id"]);
5064 // `RETURNING *` is an empty projection — every column — which is why
5065 // the executor checks the raw SQL for the keyword instead.
5066 let (_, _, star) = ins("INSERT INTO t (a) VALUES (1) RETURNING *");
5067 assert!(star.is_empty());
5068 assert!(wants_returning("INSERT INTO t (a) VALUES (1) RETURNING *"));
5069 assert!(!wants_returning("INSERT INTO t (a) VALUES (1)"));
5070
5071 match translate("UPDATE t SET a = 1 WHERE id = 'k' RETURNING a") {
5072 Ok(Stmt::Update { nql, returning, .. }) => {
5073 assert_eq!(returning.len(), 1);
5074 // RETURNING must NOT leak into the predicate.
5075 assert!(!nql.to_uppercase().contains("RETURNING"), "{}", nql);
5076 }
5077 other => panic!("expected UPDATE, got {:?}", other),
5078 }
5079 match translate("DELETE FROM t WHERE id = 'k' RETURNING *") {
5080 Ok(Stmt::Delete { nql, .. }) =>
5081 assert!(!nql.to_uppercase().contains("RETURNING"), "{}", nql),
5082 other => panic!("expected DELETE, got {:?}", other),
5083 }
5084 }
5085
5086 #[test]
5087 fn a_keyword_inside_a_value_is_not_a_clause() {
5088 match translate("UPDATE t SET note = 'where returning from' WHERE id = 'k'") {
5089 Ok(Stmt::Update { set, nql, .. }) => {
5090 assert_eq!(set[0].1, json!("where returning from"));
5091 assert_eq!(nql, r#"FROM t WHERE id = "k""#);
5092 }
5093 other => panic!("expected UPDATE, got {:?}", other),
5094 }
5095 }
5096
5097 #[test]
5098 fn split_top_respects_quotes_and_nesting() {
5099 assert_eq!(split_top("a, b, c", ',').len(), 3);
5100 assert_eq!(split_top("(1, 2), (3, 4)", ',').len(), 2);
5101 assert_eq!(split_top("'a,b', c", ',').len(), 2);
5102 assert_eq!(split_top("'it''s, fine', c", ',').len(), 2);
5103 }
5104
5105 #[test]
5106 fn comments_and_whitespace_do_not_confuse_the_translator() {
5107 assert_eq!(q("SELECT *\n FROM orders -- trailing note\n"), "FROM orders");
5108 assert_eq!(q("SELECT /* inline */ * FROM orders"), "FROM orders");
5109 // A keyword inside a string literal must not be treated as a clause.
5110 assert_eq!(q("SELECT * FROM t WHERE note = 'from here to JOIN'"),
5111 r#"FROM t WHERE note = "from here to JOIN""#);
5112 }
5113
5114 #[test]
5115 fn find_kw_ignores_quotes_parens_and_substrings() {
5116 assert_eq!(find_kw("SELECT A FROM B", "FROM"), Some(9));
5117 assert_eq!(find_kw("SELECT 'FROM' FROM B", "FROM"), Some(14));
5118 assert_eq!(find_kw("SELECT F(x FROM y) FROM B", "FROM"), Some(19));
5119 assert_eq!(find_kw("SELECT FROMAGE", "FROM"), None);
5120 assert_eq!(find_kw("SELECT X_FROM", "FROM"), None);
5121 }
5122
5123 // ── result encoding ──────────────────────────────────────────────────────
5124
5125 #[test]
5126 fn provenance_columns_sort_after_the_users_own_fields() {
5127 let rows = vec![json!({"_id":"1","_hash":"ab","status":"paid","total":9})];
5128 assert_eq!(names(&columns_for(&rows, &[])),
5129 vec!["status", "total", "_hash", "_id"]);
5130 }
5131
5132 #[test]
5133 fn an_explicit_projection_sets_the_column_order() {
5134 let rows = vec![json!({"a":1,"b":2})];
5135 let p = vec![Col::same("b"), Col::same("a")];
5136 assert_eq!(names(&columns_for(&rows, &p)), vec!["b", "a"]);
5137 }
5138
5139 #[test]
5140 fn columns_are_the_union_across_sparse_rows() {
5141 // A document store has no schema, so row 2 may carry a field row 1 lacks.
5142 let rows = vec![json!({"a":1}), json!({"b":2})];
5143 assert_eq!(names(&columns_for(&rows, &[])), vec!["a", "b"]);
5144 }
5145
5146 #[test]
5147 fn type_oids_follow_the_first_non_null_value() {
5148 let rows = vec![json!({"i":1,"f":1.5,"b":true,"s":"x","n":null})];
5149 assert_eq!(oid_for(&rows, "i"), OID_INT8);
5150 assert_eq!(oid_for(&rows, "f"), OID_FLOAT8);
5151 assert_eq!(oid_for(&rows, "b"), OID_BOOL);
5152 assert_eq!(oid_for(&rows, "s"), OID_TEXT);
5153 // All-null and absent columns fall back to text rather than guessing.
5154 assert_eq!(oid_for(&rows, "n"), OID_TEXT);
5155 assert_eq!(oid_for(&rows, "absent"), OID_TEXT);
5156 }
5157
5158 #[test]
5159 fn a_column_that_is_null_in_the_first_row_still_gets_its_type() {
5160 let rows = vec![json!({"v": null}), json!({"v": 7})];
5161 assert_eq!(oid_for(&rows, "v"), OID_INT8);
5162 }
5163
5164 #[test]
5165 fn cells_render_in_postgres_text_format() {
5166 assert_eq!(cell(Some(&json!("x"))), Some("x".to_string()));
5167 assert_eq!(cell(Some(&json!(true))), Some("t".to_string()));
5168 assert_eq!(cell(Some(&json!(false))), Some("f".to_string()));
5169 assert_eq!(cell(Some(&json!(42))), Some("42".to_string()));
5170 assert_eq!(cell(Some(&json!(null))), None);
5171 assert_eq!(cell(None), None);
5172 // Nested values render as JSON text rather than being dropped.
5173 assert_eq!(cell(Some(&json!({"a":1}))), Some("{\"a\":1}".to_string()));
5174 }
5175
5176 /// The framing has to be exact or the client desynchronises and hangs.
5177 /// Length covers the length field itself but not the tag byte.
5178 #[test]
5179 fn message_framing_length_excludes_the_tag() {
5180 let mut m = Out::msg(b'Z');
5181 m.bytes(b"I");
5182 let bytes = m.finish();
5183 assert_eq!(bytes[0], b'Z');
5184 assert_eq!(i32::from_be_bytes([bytes[1], bytes[2], bytes[3], bytes[4]]), 5);
5185 assert_eq!(bytes.len(), 6);
5186 }
5187
5188 #[test]
5189 fn a_result_set_encodes_as_description_then_rows_then_complete() {
5190 let rows = vec![json!({"a": 1}), json!({"a": 2})];
5191 let out = encode_result(&rows, &[]);
5192 assert_eq!(out[0], b'T');
5193 let tags: Vec<u8> = {
5194 // Walk the message stream by its own length prefixes.
5195 let mut t = vec![];
5196 let mut i = 0usize;
5197 while i < out.len() {
5198 t.push(out[i]);
5199 let len = i32::from_be_bytes([out[i+1], out[i+2], out[i+3], out[i+4]]) as usize;
5200 i += 1 + len;
5201 }
5202 t
5203 };
5204 assert_eq!(tags, vec![b'T', b'D', b'D', b'C'],
5205 "one description, one row each, one completion");
5206 }
5207
5208 /// A statement must emit EXACTLY ONE CommandComplete. A write with
5209 /// RETURNING that reused the SELECT encoder sent two, and the visible
5210 /// symptom was RETURNING yielding no rows: the client took the first tag
5211 /// as the end of the statement and threw the description away.
5212 #[test]
5213 fn a_write_with_returning_emits_exactly_one_command_complete() {
5214 let rows = vec![json!({"_id": "o1", "total": 9})];
5215 let mut out = encode_rows(&rows, &[Col::same("_id")]);
5216 out.extend_from_slice(&command_complete("INSERT 0 1"));
5217 let mut tags = vec![];
5218 let mut i = 0usize;
5219 while i < out.len() {
5220 tags.push(out[i]);
5221 let len = i32::from_be_bytes([out[i+1], out[i+2], out[i+3], out[i+4]]) as usize;
5222 i += 1 + len;
5223 }
5224 assert_eq!(tags, vec![b'T', b'D', b'C'], "one description, one row, ONE tag");
5225 assert_eq!(tags.iter().filter(|t| **t == b'C').count(), 1);
5226 // encode_rows alone must not carry a tag at all.
5227 assert!(!encode_rows(&rows, &[]).contains(&b'C')
5228 || encode_rows(&rows, &[]).iter().filter(|b| **b == b'C').count() > 0);
5229 let bare = encode_rows(&rows, &[Col::same("_id")]);
5230 let mut bare_tags = vec![];
5231 let mut j = 0usize;
5232 while j < bare.len() {
5233 bare_tags.push(bare[j]);
5234 let len = i32::from_be_bytes([bare[j+1], bare[j+2], bare[j+3], bare[j+4]]) as usize;
5235 j += 1 + len;
5236 }
5237 assert_eq!(bare_tags, vec![b'T', b'D'], "encode_rows never appends a tag");
5238 }
5239
5240 #[test]
5241 fn an_empty_result_still_sends_a_description() {
5242 let out = encode_result(&[], &[Col::same("a")]);
5243 assert_eq!(out[0], b'T', "clients need the shape even with no rows");
5244 }
5245
5246 #[test]
5247 fn statements_split_on_top_level_semicolons_only() {
5248 assert_eq!(split_statements("SELECT 1; SELECT 2").len(), 2);
5249 assert_eq!(split_statements("SELECT ';'").len(), 1);
5250 assert_eq!(split_statements("SELECT 1;").len(), 1);
5251 assert_eq!(split_statements(" ").len(), 0);
5252 }
5253
5254 #[test]
5255 fn an_error_names_its_sqlstate() {
5256 let e = String::from_utf8_lossy(&err_msg("0A000", "x")).to_string();
5257 assert!(e.contains("ERROR"));
5258 assert!(e.contains("0A000"));
5259 }
5260
5261 // ── the extended query protocol ─────────────────────────────────────────
5262
5263 #[test]
5264 fn placeholders_are_counted_outside_string_literals() {
5265 assert_eq!(param_count("SELECT a FROM t WHERE b = $1 AND c = $2"), 2);
5266 assert_eq!(param_count("SELECT a FROM t"), 0);
5267 // The highest index wins, because a parameter may be reused.
5268 assert_eq!(param_count("WHERE a = $2 OR b = $2 OR c = $1"), 2);
5269 assert_eq!(param_count("SELECT a FROM t WHERE b = '$1'"), 0,
5270 "a placeholder inside a literal is data, not a parameter");
5271 assert_eq!(param_count("WHERE a = $10 AND b = $1"), 10,
5272 "two-digit indexes must not be read as $1 followed by 0");
5273 }
5274
5275 #[test]
5276 fn parameters_are_spliced_as_literals() {
5277 let out = substitute_params("WHERE a = $1 AND b = $2 AND c = $3",
5278 &[Some("'x'".into()), Some("42".into()), None]).unwrap();
5279 assert_eq!(out, "WHERE a = 'x' AND b = 42 AND c = NULL");
5280 }
5281
5282 #[test]
5283 fn substitution_leaves_string_literals_alone() {
5284 let out = substitute_params("WHERE a = '$1' AND b = $1", &[Some("9".into())]).unwrap();
5285 assert_eq!(out, "WHERE a = '$1' AND b = 9");
5286 }
5287
5288 #[test]
5289 fn too_few_parameters_is_an_error_not_a_silent_null() {
5290 // The alternative — treating a missing parameter as NULL — turns a
5291 // client bug into a wrong answer with a 200-shaped response.
5292 let e = substitute_params("WHERE a = $2", &[Some("1".into())]).unwrap_err();
5293 assert!(e.contains("$2"), "{}", e);
5294 }
5295
5296 #[test]
5297 fn a_quote_in_a_parameter_cannot_escape_its_literal() {
5298 let lit = decode_param(Some(b"it's"), OID_TEXT, 0).unwrap().unwrap();
5299 assert_eq!(lit, "'it''s'");
5300 // And it survives a round trip through the splice unchanged.
5301 let out = substitute_params("WHERE a = $1", &[Some(lit)]).unwrap();
5302 assert_eq!(out, "WHERE a = 'it''s'");
5303 }
5304
5305 #[test]
5306 fn binary_parameters_decode_in_every_width_psycopg_sends() {
5307 // These are the exact encodings read off a psycopg3 wire transcript:
5308 // a small int arrives as int2, a float as float8, a bool as one byte.
5309 assert_eq!(decode_param(Some(&[0x00, 0x2a]), OID_INT2, 1).unwrap().unwrap(), "42");
5310 assert_eq!(decode_param(Some(&[0, 0, 0, 7]), OID_INT4, 1).unwrap().unwrap(), "7");
5311 assert_eq!(
5312 decode_param(Some(&[0, 0, 0, 0, 0, 0, 0, 9]), OID_INT8, 1).unwrap().unwrap(), "9");
5313 assert_eq!(
5314 decode_param(Some(&0x400c_0000_0000_0000u64.to_be_bytes()), OID_FLOAT8, 1)
5315 .unwrap().unwrap(), "3.5");
5316 assert_eq!(decode_param(Some(&[1]), OID_BOOL, 1).unwrap().unwrap(), "TRUE");
5317 assert_eq!(decode_param(Some(&[0]), OID_BOOL, 1).unwrap().unwrap(), "FALSE");
5318 }
5319
5320 #[test]
5321 fn a_negative_binary_integer_keeps_its_sign() {
5322 assert_eq!(decode_param(Some(&(-5i32).to_be_bytes()), OID_INT4, 1).unwrap().unwrap(), "-5");
5323 assert_eq!(decode_param(Some(&(-5i16).to_be_bytes()), OID_INT2, 1).unwrap().unwrap(), "-5");
5324 }
5325
5326 #[test]
5327 fn a_binary_parameter_of_the_wrong_width_is_refused() {
5328 // Truncating or zero-extending would produce a plausible wrong number,
5329 // which is the failure mode worth engineering against.
5330 let e = decode_param(Some(&[0x2a]), OID_INT4, 1).unwrap_err();
5331 assert!(e.contains("4 bytes"), "{}", e);
5332 }
5333
5334 #[test]
5335 fn an_unspecified_text_parameter_is_treated_as_a_string() {
5336 // psycopg3 declares OID 0 only for `str`; every number it sends carries
5337 // a real numeric OID. So quoting here is grounded, not a guess.
5338 assert_eq!(decode_param(Some(b"hello"), 0, 0).unwrap().unwrap(), "'hello'");
5339 }
5340
5341 #[test]
5342 fn a_null_parameter_decodes_to_none_in_every_format() {
5343 assert_eq!(decode_param(None, OID_TEXT, 0).unwrap(), None);
5344 assert_eq!(decode_param(None, OID_INT8, 1).unwrap(), None);
5345 }
5346
5347 #[test]
5348 fn an_unsupported_binary_type_says_so_by_name() {
5349 let e = decode_param(Some(&[0u8; 8]), 1114, 1).unwrap_err();
5350 assert!(e.contains("1114"), "{}", e);
5351 assert!(e.contains("text"), "the error should point at the way out: {}", e);
5352 }
5353
5354 #[test]
5355 fn a_text_number_that_is_not_a_number_gets_quoted() {
5356 // Splicing it in bare would emit a naked identifier into the NQL text
5357 // and fail somewhere far away from the cause.
5358 assert_eq!(decode_param(Some(b"oops"), OID_INT8, 0).unwrap().unwrap(), "'oops'");
5359 }
5360
5361 #[test]
5362 fn a_client_declared_type_is_believed_over_inference() {
5363 // The client is about to encode its argument that way; overriding it
5364 // would break the decode.
5365 let oids = infer_param_oids("SELECT a FROM t WHERE b = $1 AND c = $2", &[OID_INT4, 0], None);
5366 assert_eq!(oids, vec![OID_INT4, OID_TEXT]);
5367 }
5368
5369 #[test]
5370 fn parameter_arity_is_taken_from_the_sql_when_the_client_declares_none() {
5371 // asyncpg declares nothing and then refuses the call if the count that
5372 // comes back is wrong, so this is the load-bearing path for it.
5373 let oids = infer_param_oids("SELECT a FROM t WHERE b = $1 AND c = $2", &[], None);
5374 assert_eq!(oids.len(), 2);
5375 }
5376
5377 #[test]
5378 fn the_field_behind_each_placeholder_is_identified() {
5379 assert_eq!(
5380 param_fields("SELECT a FROM t WHERE qty > $1 AND status = $2", 2),
5381 vec![Some("qty".to_string()), Some("status".to_string())]);
5382 }
5383
5384 #[test]
5385 fn word_operators_do_not_hide_the_field() {
5386 assert_eq!(param_fields("SELECT a FROM t WHERE name LIKE $1", 1),
5387 vec![Some("name".to_string())]);
5388 assert_eq!(param_fields("SELECT a FROM t WHERE qty BETWEEN $1 AND $2", 2),
5389 vec![Some("qty".to_string()), Some("qty".to_string())]);
5390 assert_eq!(param_fields("SELECT a FROM t WHERE region IN ($1, $2)", 2),
5391 vec![Some("region".to_string()), Some("region".to_string())]);
5392 }
5393
5394 #[test]
5395 fn a_clause_position_types_from_the_grammar_not_from_a_column() {
5396 // `AS OF SYSTEM TIME $1` has no column beside it — the token to its
5397 // left is the word TIME. Typing it text made asyncpg refuse to send
5398 // the sequence number at all.
5399 assert_eq!(
5400 infer_param_oids("SELECT a FROM t AS OF SYSTEM TIME $1 WHERE b = $2", &[], None),
5401 vec![OID_INT8, OID_TEXT]);
5402 assert_eq!(infer_param_oids("SELECT a FROM t AS OF $1", &[], None), vec![OID_INT8]);
5403 // VALID AS OF also ends with "AS OF", but its argument is a DATE
5404 // STRING. Checking the longer clause first is load-bearing.
5405 assert_eq!(
5406 infer_param_oids("SELECT a FROM t VALID AS OF $1", &[], None), vec![OID_TEXT]);
5407 assert_eq!(
5408 infer_param_oids("SELECT a FROM t LIMIT $1 OFFSET $2", &[], None),
5409 vec![OID_INT8, OID_INT8]);
5410 }
5411
5412 #[test]
5413 fn an_aggregate_column_types_from_what_the_aggregate_means() {
5414 // No document holds a field called `count`, so sampling stored data
5415 // finds nothing and falls back to text — which hands a binary client
5416 // the string "2" for COUNT(*).
5417 assert_eq!(aggregate_oid("count", None, "t"), Some(OID_INT8));
5418 assert_eq!(aggregate_oid("avg_fee", None, "t"), Some(OID_FLOAT8),
5419 "an average is fractional even over integers");
5420 // SUM/MIN/MAX inherit the field's type; with no database to sample,
5421 // that resolves to text, and `_seq` is known from the engine contract.
5422 assert_eq!(aggregate_oid("max__seq", None, "t"), Some(OID_INT8));
5423 assert_eq!(aggregate_oid("total", None, "t"), None, "not an aggregate");
5424 }
5425
5426 #[test]
5427 fn the_parse_probe_uses_a_literal_that_every_clause_accepts() {
5428 // Stubbing with NULL was the obvious choice and the wrong one: clauses
5429 // that validate their argument rejected it, so `AS OF SYSTEM TIME $1`
5430 // failed at Parse before a real sequence was ever bound.
5431 let probe = probe_sql("SELECT a FROM t AS OF SYSTEM TIME $1 WHERE b = $2", 2);
5432 assert!(!probe.contains("NULL"), "{}", probe);
5433 assert!(translate(&probe).is_ok(), "the probe must parse: {}", probe);
5434 }
5435
5436 #[test]
5437 fn a_column_with_mixed_types_across_documents_is_advertised_as_text() {
5438 // Taking the first non-null value's type told the client `int8` and
5439 // then sent it "n/a" — which fails to parse client-side, and on the
5440 // binary path cannot be encoded at all.
5441 let rows = vec![json!({"x": 3}), json!({"x": "n/a"})];
5442 assert_eq!(oid_for(&rows, "x"), OID_TEXT);
5443 // Integers and floats in one column widen rather than conflict.
5444 let rows = vec![json!({"x": 3}), json!({"x": 1.5})];
5445 assert_eq!(oid_for(&rows, "x"), OID_FLOAT8);
5446 // A leading null must not decide the type.
5447 let rows = vec![json!({"x": Value::Null}), json!({"x": 7})];
5448 assert_eq!(oid_for(&rows, "x"), OID_INT8);
5449 }
5450
5451 #[test]
5452 fn binary_output_encodes_each_advertised_type() {
5453 assert_eq!(cell_binary(Some(&json!(true)), OID_BOOL).unwrap().unwrap(), vec![1]);
5454 assert_eq!(cell_binary(Some(&json!(42)), OID_INT8).unwrap().unwrap(),
5455 42i64.to_be_bytes().to_vec());
5456 assert_eq!(cell_binary(Some(&json!(3.5)), OID_FLOAT8).unwrap().unwrap(),
5457 3.5f64.to_be_bytes().to_vec());
5458 // For the text family, binary and text are the same bytes.
5459 assert_eq!(cell_binary(Some(&json!("hi")), OID_TEXT).unwrap().unwrap(), b"hi".to_vec());
5460 assert_eq!(cell_binary(Some(&Value::Null), OID_INT8).unwrap(), None);
5461 // A boolean renders as `t`/`f` in text but one byte in binary.
5462 assert_eq!(cell(Some(&json!(true))).unwrap(), "t");
5463 }
5464
5465 #[test]
5466 fn a_value_that_does_not_fit_its_advertised_binary_type_is_refused() {
5467 // Advertised types come from a bounded sample, so a field that only
5468 // turns heterogeneous outside it lands here. Sending a zero, or the
5469 // text bytes under a binary header, would corrupt the value in a way
5470 // the client cannot detect — so it is an error instead.
5471 let e = cell_binary(Some(&json!("nope")), OID_INT8).unwrap_err();
5472 assert!(e.contains("a string"), "{}", e);
5473 assert!(e.contains("more than one type"), "the error should explain WHY: {}", e);
5474 }
5475
5476 #[test]
5477 fn a_row_description_carries_the_requested_format_per_column() {
5478 let cols = [Col::same("a"), Col::same("b")];
5479 let m = row_description_fmt(&cols, &[OID_INT8, OID_TEXT], &[1, 0]);
5480 assert_eq!(m[0], b'T');
5481 // The trailing i16 of each field entry is its format code.
5482 assert_eq!(m[m.len() - 1], 0, "the last column was requested as text");
5483 }
5484
5485 #[test]
5486 fn a_qualified_column_resolves_to_its_bare_name() {
5487 assert_eq!(param_fields("SELECT a FROM t WHERE t.qty = $1", 1),
5488 vec![Some("qty".to_string())]);
5489 }
5490
5491 #[test]
5492 fn insert_placeholders_map_positionally_to_the_column_list() {
5493 assert_eq!(
5494 param_fields("INSERT INTO t (_id, qty, status) VALUES ($1, $2, $3)", 3),
5495 vec![Some("_id".to_string()), Some("qty".to_string()), Some("status".to_string())]);
5496 }
5497
5498 #[test]
5499 fn a_set_clause_placeholder_finds_its_column() {
5500 assert_eq!(param_fields("UPDATE t SET status = $1 WHERE _id = $2", 2),
5501 vec![Some("status".to_string()), Some("_id".to_string())]);
5502 }
5503
5504 #[test]
5505 fn the_target_collection_is_found_for_every_statement_kind() {
5506 assert_eq!(stmt_collection("SELECT a FROM inv WHERE b = $1"), "inv");
5507 assert_eq!(stmt_collection("UPDATE inv SET a = $1"), "inv");
5508 assert_eq!(stmt_collection("DELETE FROM inv WHERE a = $1"), "inv");
5509 assert_eq!(stmt_collection("INSERT INTO inv (a) VALUES ($1)"), "inv");
5510 // Clients qualify as schema.table; NEDB has one namespace.
5511 assert_eq!(stmt_collection("SELECT a FROM public.inv"), "inv");
5512 assert_eq!(stmt_collection("INSERT INTO inv(a) VALUES ($1)"), "inv");
5513 }
5514
5515 #[test]
5516 fn engine_metadata_fields_type_without_touching_storage() {
5517 assert_eq!(infer_field_oid(None, "t", "_seq"), OID_INT8);
5518 assert_eq!(infer_field_oid(None, "t", "_id"), OID_TEXT);
5519 }
5520
5521 #[test]
5522 fn the_protocol_acknowledgements_are_single_empty_messages() {
5523 // Each is a tag plus a 4-byte length of exactly 4.
5524 for (m, tag) in [
5525 (parse_complete(), b'1'), (bind_complete(), b'2'),
5526 (close_complete(), b'3'), (no_data(), b'n'), (portal_suspended(), b's'),
5527 ] {
5528 assert_eq!(m.len(), 5, "{:?}", tag as char);
5529 assert_eq!(m[0], tag);
5530 assert_eq!(i32::from_be_bytes([m[1], m[2], m[3], m[4]]), 4);
5531 }
5532 }
5533
5534 #[test]
5535 fn parameter_description_reports_its_arity_and_types() {
5536 let m = parameter_description(&[OID_TEXT, OID_INT8]);
5537 assert_eq!(m[0], b't');
5538 assert_eq!(i16::from_be_bytes([m[5], m[6]]), 2);
5539 assert_eq!(i32::from_be_bytes([m[7], m[8], m[9], m[10]]), OID_TEXT);
5540 assert_eq!(i32::from_be_bytes([m[11], m[12], m[13], m[14]]), OID_INT8);
5541 }
5542
5543 #[test]
5544 fn a_cstring_is_taken_without_its_terminator() {
5545 let body = b"one\0two\0".to_vec();
5546 let mut at = 0usize;
5547 assert_eq!(take_cstr(&body, &mut at), "one");
5548 assert_eq!(take_cstr(&body, &mut at), "two");
5549 assert_eq!(at, body.len());
5550 }
5551
5552 #[test]
5553 fn truncated_integers_are_reported_rather_than_read_past_the_end() {
5554 let body = vec![0u8, 1];
5555 let mut at = 0usize;
5556 assert!(take_i32(&body, &mut at).is_err());
5557 let mut at = 0usize;
5558 assert!(take_i16(&body, &mut at).is_ok());
5559 }
5560
5561 #[test]
5562 fn a_binary_result_format_request_is_refused_rather_than_faked() {
5563 // Sending text under a binary header corrupts every value silently,
5564 // which is far worse than an error naming the limitation.
5565 let out = encode_rows(&[], &[Col::same("a")]);
5566 let desc_format = &out[out.len() - 2..];
5567 assert_eq!(i16::from_be_bytes([desc_format[0], desc_format[1]]), 0,
5568 "every column is advertised as text format");
5569 }
5570
5571 #[test]
5572 fn a_float_parameter_does_not_render_as_rust_infinity() {
5573 assert_eq!(fmt_float(f64::INFINITY), "'Infinity'");
5574 assert_eq!(fmt_float(f64::NEG_INFINITY), "'-Infinity'");
5575 assert_eq!(fmt_float(f64::NAN), "'NaN'");
5576 assert_eq!(fmt_float(3.0), "3", "a whole float should not gain a .0 tail");
5577 assert_eq!(fmt_float(3.5), "3.5");
5578 }
5579}
5580