From f7f7bc7c65fc82f26b679edb46e10083ee2aa6ca Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Fri, 9 Oct 2026 04:30:20 +0800 Subject: [PATCH 1/5] feat(sql): coerce array literals written to VECTOR columns An array literal bound for a column declared VECTOR(dim) becomes an array of floats on every write path and every engine, including schemaless and columnar-family collections that advertise the column as text. The declared type text is now kept on every raw DDL column so the rule applies there. A non-numeric element is refused with SQLSTATE 42804 naming the column. --- nodedb-sql/src/error.rs | 10 + .../src/planner/declared_type_coerce.rs | 96 +++++++- .../src/planner/declared_vector_coerce.rs | 229 ++++++++++++++++++ nodedb-sql/src/planner/mod.rs | 1 + nodedb-sql/src/types/collection.rs | 2 + .../planner/catalog_adapter/type_convert.rs | 9 +- nodedb/src/control/planner/plan_error_map.rs | 5 + .../shared/ddl/neutral/column_default.rs | 4 +- .../wire/cases/sql_vector_array_literal.rs | 186 ++++++++++++++ 9 files changed, 532 insertions(+), 10 deletions(-) create mode 100644 nodedb-sql/src/planner/declared_vector_coerce.rs create mode 100644 nodedb/tests/wire/cases/sql_vector_array_literal.rs diff --git a/nodedb-sql/src/error.rs b/nodedb-sql/src/error.rs index 4b47cf6ec..5ea3fd987 100644 --- a/nodedb-sql/src/error.rs +++ b/nodedb-sql/src/error.rs @@ -82,6 +82,16 @@ pub enum SqlError { #[error("type mismatch: {detail}")] TypeMismatch { detail: String }, + /// An element of an array bound for a declared `VECTOR(dim)` column is + /// not a number. Rendered as SQLSTATE `42804` (datatype_mismatch). + /// `element` names the refused element's kind, and its text for a string. + #[error("column '{column}': expected a number for VECTOR({dim}), got {element}")] + VectorElementNotNumeric { + column: String, + dim: usize, + element: String, + }, + /// A statement lists a different number of targets than expressions, such /// as an `INSERT` whose target column list does not match its `SELECT` /// list. diff --git a/nodedb-sql/src/planner/declared_type_coerce.rs b/nodedb-sql/src/planner/declared_type_coerce.rs index 3dbd953b6..226a68cad 100644 --- a/nodedb-sql/src/planner/declared_type_coerce.rs +++ b/nodedb-sql/src/planner/declared_type_coerce.rs @@ -56,6 +56,9 @@ //! one with more than `p - s` integer digits is refused. A plain `DECIMAL` //! keeps every digit it is given. //! +//! A column declared `VECTOR(dim)` turns an array literal into an array of +//! floats, for every engine. See `declared_vector_coerce` for the element rule. +//! //! Every other declared type has one unambiguous literal form already and //! passes through untouched. //! @@ -79,6 +82,7 @@ use nodedb_types::datetime::NdbDateTime; use rust_decimal::Decimal; use rust_decimal::prelude::ToPrimitive; +use super::declared_vector_coerce::{coerce_to_vector, declared_vector_dim}; use super::dml_helpers::{check_declared_float_ranges, check_declared_int_ranges}; use crate::error::{Result, SqlError}; use crate::types::{ColumnInfo, SqlDataType, SqlExpr, SqlValue}; @@ -98,7 +102,7 @@ pub fn coerce_write_literal(column: &ColumnInfo, value: SqlValue) -> Result, column: &str) -> bool { exempt_column.is_some_and(|exempt| exempt.eq_ignore_ascii_case(column)) } +/// Coerce one literal written to `column`, reported under the name `name`. +/// +/// A column declared `VECTOR(dim)` takes the vector rule even when its +/// advertised [`SqlDataType`] is text, as on a schemaless or columnar-family +/// collection. Every other column takes [`coerce_value`] for its type. +fn coerce_for_column(name: &str, column: &ColumnInfo, value: SqlValue) -> Result { + match declared_vector_dim(column) { + Some(dim) => coerce_to_vector(name, value, dim), + None => coerce_value(name, value, &column.data_type), + } +} + /// Coerce one literal to `declared`, returning it unchanged when the declared /// type imposes no representation of its own. /// @@ -214,12 +230,12 @@ pub(crate) fn coerce_value( coerce_to_instant(column, value, declared) } SqlDataType::Decimal(Some(typmod)) => coerce_to_decimal(column, value, *typmod), + SqlDataType::Vector(dim) => coerce_to_vector(column, value, *dim), SqlDataType::String | SqlDataType::Bool | SqlDataType::Bytes | SqlDataType::Decimal(None) | SqlDataType::Uuid - | SqlDataType::Vector(_) | SqlDataType::Geometry | SqlDataType::Json | SqlDataType::Unknown => Ok(value), @@ -863,4 +879,76 @@ mod tests { decimal("1.5") ); } + + fn fractional_vector() -> SqlValue { + SqlValue::Array(vec![decimal("0.1"), SqlValue::Int(2), decimal("0.3")]) + } + + fn float_vector() -> SqlValue { + SqlValue::Array(vec![ + SqlValue::Float(0.1), + SqlValue::Float(2.0), + SqlValue::Float(0.3), + ]) + } + + /// A strict or KV `VECTOR(3)` column turns every element into a float on + /// the `VALUES`, `SET`, and DEFAULT paths alike. + #[test] + fn a_typed_vector_column_coerces_its_elements_to_floats() { + let columns = [column("embedding", SqlDataType::Vector(3))]; + assert_eq!( + coerced(&columns, "embedding", fractional_vector()).expect("VALUES coerces"), + float_vector() + ); + + let mut assignments = vec![( + "embedding".to_string(), + SqlExpr::Literal(fractional_vector()), + )]; + coerce_assignments_to_declared_types(&columns, &mut assignments, None) + .expect("SET coerces"); + assert!(matches!( + &assignments[0].1, + SqlExpr::Literal(value) if *value == float_vector() + )); + + assert_eq!( + coerce_write_literal(&columns[0], fractional_vector()).expect("DEFAULT coerces"), + float_vector() + ); + } + + /// A schemaless or columnar-family `VECTOR(3)` column advertises text, + /// and still takes the vector rule from its declared type text. + #[test] + fn a_text_advertised_vector_column_coerces_from_its_declared_text() { + let mut embedding = column("embedding", SqlDataType::String); + embedding.raw_type = Some("VECTOR(3)".to_string()); + assert_eq!( + coerced( + std::slice::from_ref(&embedding), + "embedding", + fractional_vector() + ) + .expect("VALUES coerces"), + float_vector() + ); + } + + /// A text element is refused with the typed vector error naming the + /// column, on every engine's write path. + #[test] + fn a_vector_column_refuses_a_text_element() { + let columns = [column("embedding", SqlDataType::Vector(2))]; + let value = SqlValue::Array(vec![decimal("0.1"), SqlValue::String("abc".into())]); + let err = coerced(&columns, "embedding", value).expect_err("text is not a number"); + assert!( + matches!( + err, + SqlError::VectorElementNotNumeric { ref column, dim: 2, .. } if column == "embedding" + ), + "{err}" + ); + } } diff --git a/nodedb-sql/src/planner/declared_vector_coerce.rs b/nodedb-sql/src/planner/declared_vector_coerce.rs new file mode 100644 index 000000000..8d91d6b3a --- /dev/null +++ b/nodedb-sql/src/planner/declared_vector_coerce.rs @@ -0,0 +1,229 @@ +// SPDX-License-Identifier: Apache-2.0 + +//! Coerce an array literal bound for a declared `VECTOR(dim)` column into an +//! array of floats. +//! +//! A fractional literal resolves to [`SqlValue::Decimal`], and the Origin +//! planner serializes a non-integer decimal as a msgpack string. Left as is, +//! `ARRAY[0.1, 0.2]` reaches every engine as an array of strings. The strict +//! and columnar encoders refuse each string element, and the schemaless +//! engine stores the strings. Converting each element here gives every +//! engine and every write path the same array of floats. +//! +//! An integer, float, or decimal element becomes an `f64`. Any other element +//! is refused with [`SqlError::VectorElementNotNumeric`]. A string element is +//! refused even when its text spells a number, because the client wrote text. +//! A number outside the `f32` range a vector element holds is refused with +//! [`SqlError::FloatOutOfRange`]. +//! +//! A value that is not an array passes through. The engines read a vector +//! text literal or raw vector bytes by their own rule. The element count is +//! checked by the engines, which own that error. + +use nodedb_types::columnar::ColumnType; +use rust_decimal::prelude::ToPrimitive; + +use crate::error::{Result, SqlError}; +use crate::types::{ColumnInfo, SqlDataType, SqlValue}; + +/// The declared type name a vector range error reports. +const VECTOR_TYPE_NAME: &str = "VECTOR"; + +/// The dimension of `column` when it is declared `VECTOR(dim)`. +/// +/// A strict or key-value column carries the dimension in its +/// [`SqlDataType::Vector`]. A schemaless or columnar-family column +/// advertises a vector as text, so its dimension comes from the declared +/// type text in [`ColumnInfo::raw_type`]. +pub(crate) fn declared_vector_dim(column: &ColumnInfo) -> Option { + if let SqlDataType::Vector(dim) = column.data_type { + return Some(dim); + } + match column + .raw_type + .as_deref() + .and_then(ColumnType::from_declared_type) + { + Some(ColumnType::Vector(dim)) => Some(dim as usize), + _ => None, + } +} + +/// Coerce `value`, bound for the `VECTOR(dim)` column `column`. +/// +/// An array becomes an array of [`SqlValue::Float`]. Every other value is +/// returned unchanged. +pub(crate) fn coerce_to_vector(column: &str, value: SqlValue, dim: usize) -> Result { + let SqlValue::Array(elements) = value else { + return Ok(value); + }; + elements + .into_iter() + .map(|element| vector_element(column, element, dim).map(SqlValue::Float)) + .collect::>>() + .map(SqlValue::Array) +} + +/// One vector element as the `f64` the column stores. +fn vector_element(column: &str, element: SqlValue, dim: usize) -> Result { + let number = match element { + SqlValue::Int(i) => i as f64, + SqlValue::Float(f) => f, + SqlValue::Decimal(d) => d.to_f64().ok_or_else(|| SqlError::TypeMismatch { + detail: format!("column '{column}': '{d}' is not representable as VECTOR({dim})"), + })?, + other => { + return Err(SqlError::VectorElementNotNumeric { + column: column.to_string(), + dim, + element: element_description(&other), + }); + } + }; + if !number.is_finite() || (number as f32).is_infinite() { + return Err(SqlError::FloatOutOfRange { + column: column.to_string(), + value: number, + declared_type: VECTOR_TYPE_NAME, + }); + } + Ok(number) +} + +/// How a refused element is named: its kind, and its text for a string. +fn element_description(element: &SqlValue) -> String { + match element { + SqlValue::String(s) => format!("text '{s}'"), + SqlValue::Null => "null".to_string(), + SqlValue::Bool(_) => "a boolean".to_string(), + SqlValue::Bytes(_) => "bytes".to_string(), + SqlValue::Array(_) => "a nested array".to_string(), + SqlValue::Timestamp(_) => "a timestamp".to_string(), + SqlValue::Timestamptz(_) => "a timestamptz".to_string(), + SqlValue::Int(_) => "an integer".to_string(), + SqlValue::Float(_) => "a float".to_string(), + SqlValue::Decimal(_) => "a decimal".to_string(), + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn decimal(text: &str) -> SqlValue { + SqlValue::Decimal(text.parse().expect("decimal literal")) + } + + fn column(data_type: SqlDataType, raw_type: Option<&str>) -> ColumnInfo { + ColumnInfo { + name: "embedding".to_string(), + data_type, + nullable: true, + is_primary_key: false, + default: None, + raw_type: raw_type.map(str::to_string), + int_width: None, + float_width: None, + } + } + + #[test] + fn fractional_elements_become_floats() { + let value = SqlValue::Array(vec![decimal("0.1"), decimal("0.2"), decimal("0.3")]); + assert_eq!( + coerce_to_vector("embedding", value, 3).expect("coerces"), + SqlValue::Array(vec![ + SqlValue::Float(0.1), + SqlValue::Float(0.2), + SqlValue::Float(0.3), + ]) + ); + } + + #[test] + fn mixed_integer_and_fractional_elements_become_floats() { + let value = SqlValue::Array(vec![SqlValue::Int(1), decimal("0.5"), SqlValue::Float(2.0)]); + assert_eq!( + coerce_to_vector("embedding", value, 3).expect("coerces"), + SqlValue::Array(vec![ + SqlValue::Float(1.0), + SqlValue::Float(0.5), + SqlValue::Float(2.0), + ]) + ); + } + + #[test] + fn a_string_element_is_refused_naming_the_column() { + for text in ["abc", "0.1"] { + let value = SqlValue::Array(vec![decimal("0.1"), SqlValue::String(text.into())]); + let error = coerce_to_vector("embedding", value, 2).expect_err("refused"); + assert_eq!( + error, + SqlError::VectorElementNotNumeric { + column: "embedding".into(), + dim: 2, + element: format!("text '{text}'"), + } + ); + } + } + + #[test] + fn null_bool_and_nested_elements_are_refused() { + for element in [ + SqlValue::Null, + SqlValue::Bool(true), + SqlValue::Array(vec![SqlValue::Int(1)]), + ] { + let value = SqlValue::Array(vec![element]); + assert!(matches!( + coerce_to_vector("embedding", value, 1), + Err(SqlError::VectorElementNotNumeric { .. }) + )); + } + } + + #[test] + fn an_element_past_the_f32_range_is_refused() { + let value = SqlValue::Array(vec![SqlValue::Float(1e39)]); + assert!(matches!( + coerce_to_vector("embedding", value, 1), + Err(SqlError::FloatOutOfRange { .. }) + )); + } + + #[test] + fn a_non_array_value_passes_through() { + for value in [ + SqlValue::Null, + SqlValue::String("[0.1,0.2]".into()), + SqlValue::Bytes(vec![0; 8]), + ] { + assert_eq!( + coerce_to_vector("embedding", value.clone(), 2).expect("passes"), + value + ); + } + } + + #[test] + fn the_dimension_comes_from_the_type_or_the_declared_text() { + assert_eq!( + declared_vector_dim(&column(SqlDataType::Vector(3), None)), + Some(3) + ); + assert_eq!( + declared_vector_dim(&column(SqlDataType::String, Some("VECTOR(4) NOT NULL"))), + Some(4) + ); + assert_eq!( + declared_vector_dim(&column(SqlDataType::String, Some("TEXT"))), + None + ); + assert_eq!( + declared_vector_dim(&column(SqlDataType::String, None)), + None + ); + } +} diff --git a/nodedb-sql/src/planner/mod.rs b/nodedb-sql/src/planner/mod.rs index e97486d7d..a4d8f4ecc 100644 --- a/nodedb-sql/src/planner/mod.rs +++ b/nodedb-sql/src/planner/mod.rs @@ -19,6 +19,7 @@ pub mod const_fold; pub mod cp_projection; pub mod cte; pub mod declared_type_coerce; +pub mod declared_vector_coerce; pub mod defaults; pub mod dml; pub mod dml_helpers; diff --git a/nodedb-sql/src/types/collection.rs b/nodedb-sql/src/types/collection.rs index 69be19055..ffa202ce2 100644 --- a/nodedb-sql/src/types/collection.rs +++ b/nodedb-sql/src/types/collection.rs @@ -90,6 +90,8 @@ pub struct ColumnInfo { /// `None` for columns synthesized by the planner (e.g. auto-injected `id`). /// Columnar INSERT converters use this to reconstruct the exact `ColumnType` /// so JSON / Geometry / UUID columns are not incorrectly inferred as String. + /// Write coercion reads a `VECTOR(dim)` declaration from it on a column + /// whose `data_type` advertises text. pub raw_type: Option, /// Declared width of an integer column, resolved once from the catalog at /// adapter-construction time. `None` for non-integer columns and for diff --git a/nodedb/src/control/planner/catalog_adapter/type_convert.rs b/nodedb/src/control/planner/catalog_adapter/type_convert.rs index 62382febb..6b8917bfe 100644 --- a/nodedb/src/control/planner/catalog_adapter/type_convert.rs +++ b/nodedb/src/control/planner/catalog_adapter/type_convert.rs @@ -154,9 +154,7 @@ pub(crate) fn convert_collection_type( if !profile.is_timeseries() && name.eq_ignore_ascii_case(pk_name) { continue; } - let mut column = declared_column_info(name, type_str); - column.raw_type = Some(type_str.clone()); - columns.push(column); + columns.push(declared_column_info(name, type_str)); } let pk = if profile.is_timeseries() { None @@ -169,7 +167,8 @@ pub(crate) fn convert_collection_type( } /// The planner-facing column a raw DDL declaration (`name`, `type_str`) -/// resolves to: its SQL type, declared numeric width, and DEFAULT text. +/// resolves to: its SQL type, declared numeric width, DEFAULT text, and the +/// declared text itself. /// /// `type_str` is the text that followed the column name in the DDL, modifiers /// included (`SMALLINT DEFAULT 5`, `TIMESTAMP TIME_KEY`). The schemaless and @@ -184,7 +183,7 @@ pub(crate) fn declared_column_info(name: &str, type_str: &str) -> ColumnInfo { nullable: true, is_primary_key: false, default: declared_default(type_str), - raw_type: None, + raw_type: Some(type_str.to_string()), int_width: IntWidth::from_declared_type(type_str), float_width: FloatWidth::from_declared_type(type_str), } diff --git a/nodedb/src/control/planner/plan_error_map.rs b/nodedb/src/control/planner/plan_error_map.rs index 75168676f..3f3d79302 100644 --- a/nodedb/src/control/planner/plan_error_map.rs +++ b/nodedb/src/control/planner/plan_error_map.rs @@ -54,6 +54,11 @@ pub(crate) fn map_plan_error( // row-scope evaluator raises, so it carries the same code. nodedb_sql::SqlError::DivisionByZero => crate::Error::DivisionByZero, nodedb_sql::SqlError::DataException { detail } => crate::Error::DataException { detail }, + // A vector element that is not a number is the wrong kind for the + // column: `42804`, the code the strict encoder gives the same element. + nodedb_sql::SqlError::VectorElementNotNumeric { .. } => crate::Error::DatatypeMismatch { + detail: error.to_string(), + }, nodedb_sql::SqlError::InvalidLimitValue { clause, value } => { crate::Error::InvalidLimitValue { clause, value } } diff --git a/nodedb/src/control/server/shared/ddl/neutral/column_default.rs b/nodedb/src/control/server/shared/ddl/neutral/column_default.rs index 2431ca28a..9183deb14 100644 --- a/nodedb/src/control/server/shared/ddl/neutral/column_default.rs +++ b/nodedb/src/control/server/shared/ddl/neutral/column_default.rs @@ -116,7 +116,9 @@ pub(super) fn validate_clause_expr(clause: &str, owner: &str, expr: &str) -> Res fn clause_error(clause: &str, owner: &str, error: &SqlError) -> DdlError { let sqlstate = match error { SqlError::UndefinedFunction { .. } => sqlstate::UNDEFINED_FUNCTION, - SqlError::TypeMismatch { .. } => sqlstate::DATATYPE_MISMATCH, + SqlError::TypeMismatch { .. } | SqlError::VectorElementNotNumeric { .. } => { + sqlstate::DATATYPE_MISMATCH + } SqlError::IntegerOutOfRange { .. } | SqlError::FloatOutOfRange { .. } | SqlError::DecimalOutOfRange { .. } diff --git a/nodedb/tests/wire/cases/sql_vector_array_literal.rs b/nodedb/tests/wire/cases/sql_vector_array_literal.rs new file mode 100644 index 000000000..33608de97 --- /dev/null +++ b/nodedb/tests/wire/cases/sql_vector_array_literal.rs @@ -0,0 +1,186 @@ +// SPDX-License-Identifier: BUSL-1.1 + +//! `ARRAY[...]` literals written to a declared `VECTOR(dim)` column. +//! +//! A fractional literal resolves to an exact decimal in the planner. Each +//! element must reach the engine as a float, so every engine stores the +//! vector the client wrote. An element that is not a number is refused with +//! SQLSTATE `42804`, naming the column. + +use crate::harness::TestServer; + +/// The numbers of a rendered vector cell: `{0.1,0.2}` or `[0.1,0.2]`. +fn vector_numbers(cell: &str) -> Vec { + cell.trim_matches(|c| matches!(c, '[' | ']' | '{' | '}')) + .split(',') + .map(|n| { + n.trim() + .parse::() + .unwrap_or_else(|e| panic!("vector element {n:?} of {cell:?}: {e}")) + }) + .collect() +} + +/// Read the `embedding` of row `id` and compare it to `expected` within +/// `f32` precision. +async fn assert_embedding(server: &TestServer, collection: &str, id: &str, expected: &[f64]) { + let rows = server + .query_text(&format!( + "SELECT embedding FROM {collection} WHERE id = '{id}'" + )) + .await + .unwrap_or_else(|e| panic!("SELECT {collection} {id}: {e}")); + assert_eq!( + rows.len(), + 1, + "{collection} {id}: expected one row, got {rows:?}" + ); + let got = vector_numbers(&rows[0]); + assert_eq!( + got.len(), + expected.len(), + "{collection} {id}: element count of {:?}", + rows[0] + ); + for (g, e) in got.iter().zip(expected) { + assert!( + (g - e).abs() < 1e-6, + "{collection} {id}: stored {:?}, expected {expected:?}", + rows[0] + ); + } +} + +async fn create(server: &TestServer, name: &str, engine: &str) { + server + .exec(&format!( + "CREATE COLLECTION {name} (id TEXT PRIMARY KEY, embedding VECTOR(3)) \ + WITH (engine = '{engine}')" + )) + .await + .unwrap_or_else(|e| panic!("CREATE {name}: {e}")); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn strict_vector_accepts_a_fractional_array() { + let server = TestServer::start().await; + create(&server, "vec_strict_frac", "document_strict").await; + + server + .exec("INSERT INTO vec_strict_frac (id, embedding) VALUES ('a1', ARRAY[0.1, 0.2, 0.3])") + .await + .expect("fractional ARRAY into VECTOR(3)"); + assert_embedding(&server, "vec_strict_frac", "a1", &[0.1, 0.2, 0.3]).await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn strict_vector_accepts_a_mixed_integer_and_fractional_array() { + let server = TestServer::start().await; + create(&server, "vec_strict_mixed", "document_strict").await; + + server + .exec("INSERT INTO vec_strict_mixed (id, embedding) VALUES ('m1', ARRAY[1, 0.5, -2])") + .await + .expect("mixed ARRAY into VECTOR(3)"); + assert_embedding(&server, "vec_strict_mixed", "m1", &[1.0, 0.5, -2.0]).await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn strict_vector_update_accepts_a_fractional_array() { + let server = TestServer::start().await; + create(&server, "vec_strict_upd", "document_strict").await; + + server + .exec("INSERT INTO vec_strict_upd (id, embedding) VALUES ('u1', ARRAY[1, 2, 3])") + .await + .expect("integer ARRAY into VECTOR(3)"); + server + .exec("UPDATE vec_strict_upd SET embedding = ARRAY[0.25, 0.5, 0.75] WHERE id = 'u1'") + .await + .expect("SET a fractional ARRAY on VECTOR(3)"); + assert_embedding(&server, "vec_strict_upd", "u1", &[0.25, 0.5, 0.75]).await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn strict_vector_refuses_a_non_numeric_element() { + let server = TestServer::start().await; + create(&server, "vec_strict_bad", "document_strict").await; + + let err = server + .exec("INSERT INTO vec_strict_bad (id, embedding) VALUES ('b1', ARRAY[0.1, 'abc', 0.3])") + .await + .expect_err("a text element is not a vector element"); + assert!(err.contains("42804"), "expected SQLSTATE 42804: {err}"); + assert!( + err.contains("embedding"), + "the error names the column: {err}" + ); + + let err = server + .exec("INSERT INTO vec_strict_bad (id, embedding) VALUES ('b2', ARRAY[0.1, '0.2', 0.3])") + .await + .expect_err("numeric text is still text"); + assert!(err.contains("42804"), "expected SQLSTATE 42804: {err}"); + + let rows = server + .query_text("SELECT id FROM vec_strict_bad") + .await + .expect("SELECT vec_strict_bad"); + assert!(rows.is_empty(), "a refused row is never stored: {rows:?}"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn strict_vector_refuses_a_fractional_array_of_the_wrong_length() { + let server = TestServer::start().await; + create(&server, "vec_strict_len", "document_strict").await; + + let err = server + .exec("INSERT INTO vec_strict_len (id, embedding) VALUES ('l1', ARRAY[0.1, 0.2])") + .await + .expect_err("two elements for VECTOR(3)"); + assert!(err.contains("22000"), "expected SQLSTATE 22000: {err}"); + assert!( + err.contains("got 2 elements"), + "the error counts the elements: {err}" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn schemaless_vector_stores_a_fractional_array_as_numbers() { + let server = TestServer::start().await; + create(&server, "vec_loose_frac", "document_schemaless").await; + + server + .exec("INSERT INTO vec_loose_frac (id, embedding) VALUES ('s1', ARRAY[0.1, 0.2, 0.3])") + .await + .expect("fractional ARRAY into a schemaless VECTOR(3)"); + assert_embedding(&server, "vec_loose_frac", "s1", &[0.1, 0.2, 0.3]).await; + + let rows = server + .query_text("SELECT embedding FROM vec_loose_frac WHERE id = 's1'") + .await + .expect("SELECT vec_loose_frac"); + assert!( + !rows[0].contains('"'), + "the elements are stored as numbers, not text: {:?}", + rows[0] + ); + + let err = server + .exec("INSERT INTO vec_loose_frac (id, embedding) VALUES ('s2', ARRAY[0.1, 'abc', 0.3])") + .await + .expect_err("a text element is not a vector element"); + assert!(err.contains("42804"), "expected SQLSTATE 42804: {err}"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn columnar_vector_accepts_a_fractional_array() { + let server = TestServer::start().await; + create(&server, "vec_col_frac", "columnar").await; + + server + .exec("INSERT INTO vec_col_frac (id, embedding) VALUES ('c1', ARRAY[0.1, 0.2, 0.3])") + .await + .expect("fractional ARRAY into a columnar VECTOR(3)"); + assert_embedding(&server, "vec_col_frac", "c1", &[0.1, 0.2, 0.3]).await; +} From 6c1c4ff1eda27fca52d96e670a25282d00e00e56 Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Fri, 9 Oct 2026 04:30:20 +0800 Subject: [PATCH 2/5] refactor(calvin): move the scheduler event loop into run_loop Narrow the CI nextest exclusion to the one cluster test that still fails there and document its cause. --- .config/nextest.toml | 36 +--- .../calvin/scheduler/driver/core/mod.rs | 1 + .../calvin/scheduler/driver/core/run_loop.rs | 193 ++++++++++++++++++ .../calvin/scheduler/driver/core/scheduler.rs | 178 +--------------- 4 files changed, 204 insertions(+), 204 deletions(-) create mode 100644 nodedb/src/control/cluster/calvin/scheduler/driver/core/run_loop.rs diff --git a/.config/nextest.toml b/.config/nextest.toml index 358ac5719..26689fbb1 100644 --- a/.config/nextest.toml +++ b/.config/nextest.toml @@ -212,35 +212,17 @@ server-process-serial = { max-threads = 1 } retries = { backoff = "fixed", count = 1, delay = "2s" } fail-fast = false -# Four cluster tests are excluded from the CI profile ONLY. They still run under -# the default profile, which is what a local `cargo nextest run` uses, because -# they pass locally — CI is the sole environment they fail in. +# One cluster test is excluded from the CI profile ONLY. It still runs under the +# default profile, which is what a local `cargo nextest run` uses. # -# All four hang identically: the coordinator never receives a Calvin completion -# ack and the statement dies with `timed out waiting for Calvin completion`. -# Raising the internal request deadline from 30s to 120s was tried and did NOT -# help — they timed out at 120s as well, which is what rules out slowness and -# establishes a genuine hang. They reproduce on CI on every attempt and have -# never reproduced locally: not in isolation, not inside the full cluster suite, -# and not pinned to two cores. -# -# This is NOT retry-masking. The `retries` above were already being exhausted -# against them on every run; excluding makes the gap explicit rather than -# leaving it to hide in retry noise. -# -# `learner_caught_up_via_real_install_snapshot` and -# `native_implicit_edge_delete_cleans_reverse_cross_node` were passing on CI -# before this release branch and regressed, so start there. The open suspect is -# the sequencer catch-up armed at scheduler spawn, which replays the retained -# sequencer log and may disturb in-flight Calvin transactions. Re-enable by -# diagnosing the Calvin completion path on a CI run — NOT by widening a timeout. +# `learner_caught_up_via_real_install_snapshot` fails on CI runners about once in +# six runs. A newly elected sequencer leader whose log was compacted refuses to +# propose ("sequencer log was truncated below every retained epoch"), so the +# sequencer stalls. The joining node then never catches up on the sequencer +# group, and its authorization lease is withheld. Re-enable it once a leader can +# seed its epoch from the installed sequencer snapshot. NEVER widen a timeout. default-filter = """ -not ( - test(ollp_implicit_edge_delete_cleans_reverse_cross_node) - or test(ollp_implicit_edge_update_moves_edge_cross_node) - or test(learner_caught_up_via_real_install_snapshot) - or test(native_implicit_edge_delete_cleans_reverse_cross_node) -) +not test(learner_caught_up_via_real_install_snapshot) """ # CI-only cluster headroom. On CI runners (~2× slower per-core than a dev diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/mod.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/mod.rs index aaad71d68..d4895393b 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/mod.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/mod.rs @@ -36,6 +36,7 @@ pub mod read_result; mod redo_window; pub mod request; pub mod routing; +pub mod run_loop; pub mod scheduler; pub mod sequencer_proposer; pub mod staged_vote; diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/run_loop.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/run_loop.rs new file mode 100644 index 000000000..60892742b --- /dev/null +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/run_loop.rs @@ -0,0 +1,193 @@ +// SPDX-License-Identifier: BUSL-1.1 + +//! The scheduler event loop: per-pass sweeps and the `select!` over every input. + +use std::sync::Arc; + +use tracing::info; + +use super::catch_up::CatchUpDrain; +use super::scheduler::Scheduler; +use crate::control::shutdown::ShutdownReceiver; + +impl Scheduler { + /// Run the scheduler event loop until shutdown is signaled. + pub async fn run(mut self, mut shutdown: ShutdownReceiver) { + info!( + vshard_id = self.vshard_id, + fully_applied_epoch = self.applied.fully_applied_epoch(), + rebuild_target_epoch = self.rebuild_target_epoch, + "calvin scheduler starting" + ); + + // Low-frequency liveness timer so the top-of-loop stall/barrier sweeps + // run even on an otherwise-idle vShard. Without it, a dropped verdict + // push plus zero further events for this vShard will leave a parked txn + // never re-probing the durable verdict. A fraction of the stall-warn + // window re-probes well within it. + let mut stall_tick = tokio::time::interval(self.config.verdict_stall_warn() / 4); + stall_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + + // Woken when a routed Data-Plane response frees dispatcher capacity. + let capacity_freed = Arc::clone(&self.capacity_freed); + // Set when a tick left armed catch-up unreplayed. The next open-gate + // pass fires the tick at once to resume it. + let mut catch_up_resume = false; + + loop { + // Register for the capacity wake BEFORE the re-send pass. A + // response routed after a refusal but before this point is + // covered by the pass itself. One routed after it wakes the arm. + let capacity_notified = capacity_freed.notified(); + tokio::pin!(capacity_notified); + capacity_notified.as_mut().enable(); + if self.resends_deferred() { + self.redispatch_deferred(); + } + + self.check_dependent_barrier_timeouts(); + self.check_awaiting_verdict_stalls(); + self.resume_metadata_hold(); + + // Open only once the inputs a closed gate held are processed. A + // gate closed for the backlog still takes awaited parts and + // releases (see `super::parts_lane`). + let intake_open = self.pass_intake_lane(); + if intake_open && catch_up_resume { + catch_up_resume = false; + stall_tick.reset_immediately(); + } + + tokio::select! { + biased; + + _ = shutdown.wait_cancelled() => { + info!(vshard_id = self.vshard_id, "calvin scheduler shutting down"); + break; + } + + // A closed channel below means this node retired the vShard or + // started a new scheduler for it. A closed receiver yields at + // once on every poll, so the loop must leave rather than spin. + maybe_completion = self.completion_rx.recv() => { + let Some((txn_id, request_id, resp_opt)) = maybe_completion else { + self.log_superseded("completion"); + break; + }; + // Awaited in the arm: the loop takes no other input + // until this completion, its durability wait included, + // is fully handled. + self.handle_completion(txn_id, request_id, resp_opt).await; + } + + maybe_verdict = self.verdict_rx.recv() => { + let Some(signal) = maybe_verdict else { + self.log_superseded("verdict"); + break; + }; + // A durable global verdict landed: resume the matching + // parked txn into its flush (commit) or drop (abort). + self.handle_verdict_signal(signal); + } + + maybe_event = self.read_result_rx.recv() => { + let Some(event) = maybe_event else { + self.log_superseded("read result"); + break; + }; + self.handle_read_result(event); + } + + maybe_promoted = self.promotion_rx.recv() => { + let Some(promoted) = maybe_promoted else { + self.log_superseded("promotion"); + break; + }; + // A fast-path write-admission guard released an uncontended + // key that one of this scheduler's txns had queued behind; + // `release` already promoted it to holder. Run the normal + // promotion -> dispatch path so it stops being a stalled + // holder in `blocked` and actually executes. + self.dispatch_promoted(promoted); + } + + _ = tokio::time::sleep(super::metadata_hold::HOLD_POLL), + if self.metadata_hold.is_some() => { + // The next loop pass re-checks the held txn's floor. + } + + _ = &mut capacity_notified, if self.resends_deferred() => { + // Capacity freed: the next loop pass re-sends deferred + // requests in FIFO order. + } + + _ = self.install_gate.released(), if self.install_gate.is_waiting() => { + // A snapshot released a group's gate: the waiting flush + // tries again. + self.pump_flush_turn(); + } + + maybe_txn = self.receiver.recv(), if intake_open => { + match maybe_txn { + Some(input) => self.process_scheduler_input(input), + None => { + info!( + vshard_id = self.vshard_id, + "calvin scheduler: receiver channel closed; exiting" + ); + break; + } + } + } + + _ = stall_tick.tick() => { + // Replay any sequencer-fan-out inputs dropped on this replica + // (channel Full/Closed) so a missed `SchedulerInput` never + // permanently diverges this vShard's lock table from its peers. + // O(1) common case (no pending catch-up). See `drain_catch_up`. + // A closed intake gate skips the drain until it opens. + catch_up_resume = + !intake_open || self.drain_catch_up() == CatchUpDrain::Remaining; + // Propose again every owed sequencer entry not yet applied. + self.retry_owed_sequencer_entries(); + // The top-of-loop check_awaiting_verdict_stalls / + // check_dependent_barrier_timeouts and the deferred re-send + // pass run on every wake; this arm guarantees the loop wakes + // to run them (and the drain) when no other event arrives. + } + } + } + // Every txn still pending stays unapplied on this replica. + self.hold_all_redo_records(); + } + + /// Log that a scheduler channel closed and the loop exits. + fn log_superseded(&self, channel: &str) { + info!( + vshard_id = self.vshard_id, + channel, + "calvin scheduler: a channel closed; the vShard left this node or a new scheduler took it" + ); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + use crate::control::cluster::calvin::scheduler::driver::core::test_support::{ + build_test_scheduler, spawn_scheduler_loop, + }; + + /// Retiring a vShard drops its verdict sender. The scheduler then leaves + /// its loop. A closed receiver is always ready, so a loop that ignored the + /// close would spin and starve the runtime. + #[tokio::test] + async fn a_closed_verdict_channel_ends_the_loop() { + let (scheduler, _dir) = build_test_scheduler(0); + let registry = Arc::clone(&scheduler.registry); + let running = spawn_scheduler_loop(scheduler); + registry.unregister_verdict_signal_sender(0); + assert!(running.exits_unprompted().await); + } +} diff --git a/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs b/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs index 812a7e563..002090c4a 100644 --- a/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs +++ b/nodedb/src/control/cluster/calvin/scheduler/driver/core/scheduler.rs @@ -6,7 +6,6 @@ use std::collections::BTreeMap; use std::sync::{Arc, Mutex}; use tokio::sync::{Notify, mpsc}; -use tracing::info; use nodedb_cluster::MultiRaft; use nodedb_cluster::calvin::types::SchedulerInput; @@ -15,7 +14,6 @@ use nodedb_cluster::calvin::{CalvinCompletionRegistry, SequencerStateMachine, Ve use super::super::barrier::{PendingDependentBarrier, ReadResultEvent}; use super::super::config::SchedulerConfig; use super::super::types::{BlockedTxn, PendingTxn}; -use super::catch_up::CatchUpDrain; use super::deferred::DeferredQueue; use super::halt::HaltLatch; use super::intake::IntakeGate; @@ -25,7 +23,6 @@ use crate::bridge::envelope::Response; use crate::control::cluster::calvin::scheduler::lock_manager::{LockManager, TxnId}; use crate::control::cluster::calvin::scheduler::metrics::SchedulerMetrics; use crate::control::cluster::calvin::scheduler::{AppliedGate, NOT_YET_APPLIED_EPOCH}; -use crate::control::shutdown::ShutdownReceiver; use crate::control::state::SharedState; use crate::types::RequestId; @@ -353,165 +350,6 @@ impl Scheduler { // A marker passes once every epoch delivered before it folded. self.report_passed_cuts(); } - - /// Run the scheduler event loop until shutdown is signaled. - pub async fn run(mut self, mut shutdown: ShutdownReceiver) { - info!( - vshard_id = self.vshard_id, - fully_applied_epoch = self.applied.fully_applied_epoch(), - rebuild_target_epoch = self.rebuild_target_epoch, - "calvin scheduler starting" - ); - - // Low-frequency liveness timer so the top-of-loop stall/barrier sweeps - // run even on an otherwise-idle vShard. Without it, a dropped verdict - // push plus zero further events for this vShard will leave a parked txn - // never re-probing the durable verdict. A fraction of the stall-warn - // window re-probes well within it. - let mut stall_tick = tokio::time::interval(self.config.verdict_stall_warn() / 4); - stall_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); - - // Woken when a routed Data-Plane response frees dispatcher capacity. - let capacity_freed = Arc::clone(&self.capacity_freed); - // Set when a tick left armed catch-up unreplayed. The next open-gate - // pass fires the tick at once to resume it. - let mut catch_up_resume = false; - - loop { - // Register for the capacity wake BEFORE the re-send pass. A - // response routed after a refusal but before this point is - // covered by the pass itself. One routed after it wakes the arm. - let capacity_notified = capacity_freed.notified(); - tokio::pin!(capacity_notified); - capacity_notified.as_mut().enable(); - if self.resends_deferred() { - self.redispatch_deferred(); - } - - self.check_dependent_barrier_timeouts(); - self.check_awaiting_verdict_stalls(); - self.resume_metadata_hold(); - - // Open only once the inputs a closed gate held are processed. A - // gate closed for the backlog still takes awaited parts and - // releases (see `super::parts_lane`). - let intake_open = self.pass_intake_lane(); - if intake_open && catch_up_resume { - catch_up_resume = false; - stall_tick.reset_immediately(); - } - - tokio::select! { - biased; - - _ = shutdown.wait_cancelled() => { - info!(vshard_id = self.vshard_id, "calvin scheduler shutting down"); - break; - } - - // A closed channel below means this node retired the vShard or - // started a new scheduler for it. A closed receiver yields at - // once on every poll, so the loop must leave rather than spin. - maybe_completion = self.completion_rx.recv() => { - let Some((txn_id, request_id, resp_opt)) = maybe_completion else { - self.log_superseded("completion"); - break; - }; - // Awaited in the arm: the loop takes no other input - // until this completion, its durability wait included, - // is fully handled. - self.handle_completion(txn_id, request_id, resp_opt).await; - } - - maybe_verdict = self.verdict_rx.recv() => { - let Some(signal) = maybe_verdict else { - self.log_superseded("verdict"); - break; - }; - // A durable global verdict landed: resume the matching - // parked txn into its flush (commit) or drop (abort). - self.handle_verdict_signal(signal); - } - - maybe_event = self.read_result_rx.recv() => { - let Some(event) = maybe_event else { - self.log_superseded("read result"); - break; - }; - self.handle_read_result(event); - } - - maybe_promoted = self.promotion_rx.recv() => { - let Some(promoted) = maybe_promoted else { - self.log_superseded("promotion"); - break; - }; - // A fast-path write-admission guard released an uncontended - // key that one of this scheduler's txns had queued behind; - // `release` already promoted it to holder. Run the normal - // promotion -> dispatch path so it stops being a stalled - // holder in `blocked` and actually executes. - self.dispatch_promoted(promoted); - } - - _ = tokio::time::sleep(super::metadata_hold::HOLD_POLL), - if self.metadata_hold.is_some() => { - // The next loop pass re-checks the held txn's floor. - } - - _ = &mut capacity_notified, if self.resends_deferred() => { - // Capacity freed: the next loop pass re-sends deferred - // requests in FIFO order. - } - - _ = self.install_gate.released(), if self.install_gate.is_waiting() => { - // A snapshot released a group's gate: the waiting flush - // tries again. - self.pump_flush_turn(); - } - - maybe_txn = self.receiver.recv(), if intake_open => { - match maybe_txn { - Some(input) => self.process_scheduler_input(input), - None => { - info!( - vshard_id = self.vshard_id, - "calvin scheduler: receiver channel closed; exiting" - ); - break; - } - } - } - - _ = stall_tick.tick() => { - // Replay any sequencer-fan-out inputs dropped on this replica - // (channel Full/Closed) so a missed `SchedulerInput` never - // permanently diverges this vShard's lock table from its peers. - // O(1) common case (no pending catch-up). See `drain_catch_up`. - // A closed intake gate skips the drain until it opens. - catch_up_resume = - !intake_open || self.drain_catch_up() == CatchUpDrain::Remaining; - // Propose again every owed sequencer entry not yet applied. - self.retry_owed_sequencer_entries(); - // The top-of-loop check_awaiting_verdict_stalls / - // check_dependent_barrier_timeouts and the deferred re-send - // pass run on every wake; this arm guarantees the loop wakes - // to run them (and the drain) when no other event arrives. - } - } - } - // Every txn still pending stays unapplied on this replica. - self.hold_all_redo_records(); - } - - /// Log that a scheduler channel closed and the loop exits. - fn log_superseded(&self, channel: &str) { - info!( - vshard_id = self.vshard_id, - channel, - "calvin scheduler: a channel closed; the vShard left this node or a new scheduler took it" - ); - } } // ── `is_caught_up` sentinel handling ───────────────────────────────────────── @@ -520,21 +358,7 @@ mod tests { use super::*; use std::collections::BTreeSet; - use crate::control::cluster::calvin::scheduler::driver::core::test_support::{ - build_test_scheduler, spawn_scheduler_loop, - }; - - /// Retiring a vShard drops its verdict sender. The scheduler then leaves - /// its loop. A closed receiver is always ready, so a loop that ignored the - /// close would spin and starve the runtime. - #[tokio::test] - async fn a_closed_verdict_channel_ends_the_loop() { - let (scheduler, _dir) = build_test_scheduler(0); - let registry = Arc::clone(&scheduler.registry); - let running = spawn_scheduler_loop(scheduler); - registry.unregister_verdict_signal_sender(0); - assert!(running.exits_unprompted().await); - } + use crate::control::cluster::calvin::scheduler::driver::core::test_support::build_test_scheduler; /// A freshly-recovered scheduler (`fully_applied_epoch` still the /// `NOT_YET_APPLIED_EPOCH` sentinel) with a REAL, non-zero rebuild target must From d550af4caaa376553c3a95b8a108cbdff626912c Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Fri, 9 Oct 2026 04:30:20 +0800 Subject: [PATCH 3/5] feat(crdt): map CRDT errors to typed verdicts A CRDT error now takes the Data-Plane error code of its cause instead of reading as a server fault. The SQLSTATE, public error class, and cluster wire form follow that code, and only a Loro fault stays internal. Add UndefinedObject and ObjectNotInPrerequisiteState to the exhaustive error-code matches. --- .../src/rpc_codec/data_plane_error.rs | 9 + nodedb-crdt/src/error.rs | 30 ++ nodedb-crdt/src/list_ops.rs | 201 ++++++------ nodedb-crdt/src/state/bitemporal_archive.rs | 65 ++-- nodedb-crdt/src/state/changed_rows.rs | 34 +- nodedb-crdt/src/state/core.rs | 280 ++++++++++++++--- nodedb-crdt/src/state/document_cell.rs | 24 ++ nodedb-crdt/src/state/history.rs | 56 +++- nodedb-crdt/src/state/preview.rs | 19 +- nodedb-crdt/src/state/remove_fields.rs | 21 +- nodedb-crdt/src/state/row_image.rs | 81 ++--- nodedb/src/bridge/envelope/crdt_error_code.rs | 296 ++++++++++++++++++ nodedb/src/bridge/envelope/error_code.rs | 9 + nodedb/src/bridge/envelope/error_code_from.rs | 18 +- nodedb/src/bridge/envelope/mod.rs | 1 + .../control/cluster/data_plane_error_wire.rs | 24 ++ .../control/cluster/execution_error_wire.rs | 6 +- .../error_map/class_parity/code_samples.rs | 11 +- .../error_map/class_parity/error_parity.rs | 3 + .../server/dispatch_utils/write_abort.rs | 8 +- .../control/server/pgwire/types/error_map.rs | 8 +- .../src/control/server/shared/ddl/sqlstate.rs | 10 + .../sync/async_dispatch/delta/compensation.rs | 2 + nodedb/src/control/server/sync/refusal.rs | 4 + .../executor/handlers/control/crdt_doc.rs | 4 +- .../executor/handlers/graph_edge_write/put.rs | 5 +- .../handlers/point/apply_put/types.rs | 2 + .../src/engine/crdt/tenant_state/history.rs | 14 +- .../engine/crdt/tenant_state/snapshot_io.rs | 5 +- nodedb/src/error_classify/public.rs | 6 +- nodedb/src/error_classify/unclassified.rs | 21 +- nodedb/src/error_from.rs | 6 +- nodedb/src/error_from_data_plane.rs | 4 + .../crdt_restore_non_default_database.rs | 88 ++++++ 34 files changed, 1102 insertions(+), 273 deletions(-) create mode 100644 nodedb/src/bridge/envelope/crdt_error_code.rs create mode 100644 nodedb/tests/wire/cases/crdt_restore_non_default_database.rs diff --git a/nodedb-cluster/src/rpc_codec/data_plane_error.rs b/nodedb-cluster/src/rpc_codec/data_plane_error.rs index 72ecc56ec..f67852305 100644 --- a/nodedb-cluster/src/rpc_codec/data_plane_error.rs +++ b/nodedb-cluster/src/rpc_codec/data_plane_error.rs @@ -228,6 +228,15 @@ pub enum DataPlaneErrorCode { DatetimeFieldOverflow { detail: String, }, + /// The request names an object that does not exist (SQLSTATE `42704`). + UndefinedObject { + object: String, + }, + /// The object is not in the state the request needs (SQLSTATE `55000`). + ObjectNotInPrerequisiteState { + object: String, + detail: String, + }, } /// Wire mirror of `nodedb_types::text_search::TextColumnFault`. diff --git a/nodedb-crdt/src/error.rs b/nodedb-crdt/src/error.rs index 66bd26f27..82902708b 100644 --- a/nodedb-crdt/src/error.rs +++ b/nodedb-crdt/src/error.rs @@ -53,6 +53,9 @@ pub enum CrdtError { /// Preview requires a quiescent authoritative state and refuses to commit /// a caller's pending auto-commit transaction as a side effect of forking. + /// + /// Every `CrdtState` mutator commits before it returns. This error + /// therefore means that invariant broke, never that a client erred. #[error("delta preview source has {operations} pending operations")] PreviewSourceTransactionPending { operations: usize }, @@ -111,6 +114,33 @@ pub enum CrdtError { discarded: i32, }, + /// The row did not exist at the requested version, so there is no state + /// to read or restore there. + #[error("document `{row_id}` in collection `{collection}` did not exist at the target version")] + RowAbsentAtVersion { collection: String, row_id: String }, + + /// A block-list operation named a row that does not exist, or a row that + /// is not a map of fields. + #[error("row `{row_id}` in collection `{collection}` is absent or is not a map")] + BlockListRowAbsent { collection: String, row_id: String }, + + /// A block-list path segment is absent, or holds a value of the wrong kind. + #[error("block-list path `{list_path}`: segment `{segment}` is absent or has the wrong type")] + BlockListPathUnresolved { list_path: String, segment: String }, + + /// A block-list path resolves to a plain list. Block lists must be movable + /// lists, so concurrent reorders converge without duplicates. + #[error("block-list path `{list_path}` is a plain list, not a movable list")] + BlockListNotMovable { list_path: String }, + + /// A block-list index is past the end of the list. + #[error("block-list `{list_path}` index {index} is out of bounds (len {len})")] + BlockListIndexOutOfBounds { + list_path: String, + index: usize, + len: usize, + }, + /// Loro internal error. #[error("loro error: {0}")] Loro(String), diff --git a/nodedb-crdt/src/list_ops.rs b/nodedb-crdt/src/list_ops.rs index 6d5df9667..937e5caa4 100644 --- a/nodedb-crdt/src/list_ops.rs +++ b/nodedb-crdt/src/list_ops.rs @@ -57,10 +57,11 @@ pub fn list_delete( ) -> Result<()> { let list = get_movable_list(doc, collection, row_id, list_path)?; if index >= list.len() { - return Err(CrdtError::Loro(format!( - "list delete index {index} out of bounds (len={})", - list.len() - ))); + return Err(CrdtError::BlockListIndexOutOfBounds { + list_path: list_path.to_string(), + index, + len: list.len(), + }); } list.delete(index, 1) .map_err(|e| CrdtError::Loro(format!("list delete at {index}: {e}")))?; @@ -84,15 +85,14 @@ pub fn list_move( } let list = get_movable_list(doc, collection, row_id, list_path)?; let len = list.len(); - if from_index >= len { - return Err(CrdtError::Loro(format!( - "list move from_index {from_index} out of bounds (len={len})" - ))); - } - if to_index >= len { - return Err(CrdtError::Loro(format!( - "list move to_index {to_index} out of bounds (len={len})" - ))); + for index in [from_index, to_index] { + if index >= len { + return Err(CrdtError::BlockListIndexOutOfBounds { + list_path: list_path.to_string(), + index, + len, + }); + } } list.mov(from_index, to_index) .map_err(|e| CrdtError::Loro(format!("list move {from_index}→{to_index}: {e}")))?; @@ -129,6 +129,51 @@ pub fn list_get( }) } +/// The row map `list_path` is resolved in. A block list lives inside an +/// existing row map, so an absent or non-map row is refused. +fn block_list_row(doc: &LoroDoc, collection: &str, row_id: &str) -> Result { + match doc.get_map(collection).get(row_id) { + Some(ValueOrContainer::Container(loro::Container::Map(m))) => Ok(m), + _ => Err(CrdtError::BlockListRowAbsent { + collection: collection.to_string(), + row_id: row_id.to_string(), + }), + } +} + +/// Split `list_path` into its parent map segments and the list segment. +fn split_list_path(list_path: &str) -> (Vec<&str>, &str) { + match list_path.rsplit_once('.') { + Some((parents, last)) => (parents.split('.').collect(), last), + None => (Vec::new(), list_path), + } +} + +/// The movable list held by `value` at the last segment of `list_path`. +fn movable_list_at( + value: Option, + list_path: &str, + segment: &str, +) -> Result> { + match value { + Some(ValueOrContainer::Container(loro::Container::MovableList(l))) => Ok(Some(l)), + Some(ValueOrContainer::Container(loro::Container::List(_))) => { + Err(CrdtError::BlockListNotMovable { + list_path: list_path.to_string(), + }) + } + None => Ok(None), + Some(_) => Err(path_unresolved(list_path, segment)), + } +} + +fn path_unresolved(list_path: &str, segment: &str) -> CrdtError { + CrdtError::BlockListPathUnresolved { + list_path: list_path.to_string(), + segment: segment.to_string(), + } +} + /// Navigate to a LoroMovableList at `collection/row_id/list_path`. /// /// `list_path` is a dot-separated field path within the row LoroMap. @@ -140,46 +185,16 @@ fn get_movable_list( row_id: &str, list_path: &str, ) -> Result { - let coll = doc.get_map(collection); - let row = match coll.get(row_id) { - Some(ValueOrContainer::Container(loro::Container::Map(m))) => m, - _ => { - return Err(CrdtError::Loro(format!( - "row '{row_id}' not found or not a map in '{collection}'" - ))); - } - }; - - let segments: Vec<&str> = list_path.split('.').collect(); - let mut current_map = row; - - for (i, segment) in segments.iter().enumerate() { - let is_last = i == segments.len() - 1; - match current_map.get(segment) { - Some(ValueOrContainer::Container(loro::Container::MovableList(l))) if is_last => { - return Ok(l); - } - // Also accept LoroList for backward compatibility (read-only navigation). - Some(ValueOrContainer::Container(loro::Container::List(_))) if is_last => { - return Err(CrdtError::Loro(format!( - "path '{list_path}' resolved to LoroList, not LoroMovableList. \ - Use LoroMovableList for block containers that support reordering." - ))); - } - Some(ValueOrContainer::Container(loro::Container::Map(m))) if !is_last => { - current_map = m; - } - _ => { - return Err(CrdtError::Loro(format!( - "path '{list_path}' segment '{segment}' not found or wrong type" - ))); - } - } + let (parents, last) = split_list_path(list_path); + let mut current_map = block_list_row(doc, collection, row_id)?; + for segment in parents { + current_map = match current_map.get(segment) { + Some(ValueOrContainer::Container(loro::Container::Map(m))) => m, + _ => return Err(path_unresolved(list_path, segment)), + }; } - - Err(CrdtError::Loro(format!( - "path '{list_path}' did not resolve to a movable list" - ))) + movable_list_at(current_map.get(last), list_path, last)? + .ok_or_else(|| path_unresolved(list_path, last)) } /// Navigate to a LoroMovableList at `collection/row_id/list_path`, creating @@ -204,42 +219,9 @@ fn get_or_create_movable_list( row_id: &str, list_path: &str, ) -> Result { - let coll = doc.get_map(collection); - let row = match coll.get(row_id) { - Some(ValueOrContainer::Container(loro::Container::Map(m))) => m, - _ => { - return Err(CrdtError::Loro(format!( - "row '{row_id}' not found or not a map in '{collection}'" - ))); - } - }; - - let segments: Vec<&str> = list_path.split('.').collect(); - let mut current_map = row; - - for (i, segment) in segments.iter().enumerate() { - let is_last = i == segments.len() - 1; - - if is_last { - return match current_map.get(segment) { - Some(ValueOrContainer::Container(loro::Container::MovableList(l))) => Ok(l), - Some(ValueOrContainer::Container(loro::Container::List(_))) => { - Err(CrdtError::Loro(format!( - "path '{list_path}' resolved to LoroList, not LoroMovableList. \ - Use LoroMovableList for block containers that support reordering." - ))) - } - None => current_map - .insert_container(segment, LoroMovableList::new()) - .map_err(|e| { - CrdtError::Loro(format!("create movable list at '{list_path}': {e}")) - }), - Some(_) => Err(CrdtError::Loro(format!( - "path '{list_path}' segment '{segment}' not found or wrong type" - ))), - }; - } - + let (parents, last) = split_list_path(list_path); + let mut current_map = block_list_row(doc, collection, row_id)?; + for segment in parents { current_map = match current_map.get(segment) { Some(ValueOrContainer::Container(loro::Container::Map(m))) => m, None => current_map @@ -249,17 +231,15 @@ fn get_or_create_movable_list( "create intermediate map at '{list_path}' segment '{segment}': {e}" )) })?, - Some(_) => { - return Err(CrdtError::Loro(format!( - "path '{list_path}' segment '{segment}' not found or wrong type" - ))); - } + Some(_) => return Err(path_unresolved(list_path, segment)), }; } - - Err(CrdtError::Loro(format!( - "path '{list_path}' did not resolve to a movable list" - ))) + match movable_list_at(current_map.get(last), list_path, last)? { + Some(list) => Ok(list), + None => current_map + .insert_container(last, LoroMovableList::new()) + .map_err(|e| CrdtError::Loro(format!("create movable list at '{list_path}': {e}"))), + } } #[cfg(test)] @@ -412,14 +392,22 @@ mod tests { fn list_delete_out_of_bounds() { let state = setup_doc_with_blocks(); let err = list_delete(state.doc(), "pages", "doc-1", "blocks", 99); - assert!(err.is_err()); + assert!(matches!( + err, + Err(CrdtError::BlockListIndexOutOfBounds { index: 99, .. }) + )); } #[test] fn get_list_wrong_path_errors() { let state = setup_doc_with_blocks(); let err = list_length(state.doc(), "pages", "doc-1", "nonexistent"); - assert!(err.is_err()); + assert!(matches!( + err, + Err(CrdtError::BlockListPathUnresolved { ref segment, .. }) if segment == "nonexistent" + )); + let err = list_length(state.doc(), "pages", "missing-row", "blocks"); + assert!(matches!(err, Err(CrdtError::BlockListRowAbsent { .. }))); } /// A row with no list at `list_path` yet — the setup previous agents @@ -502,7 +490,10 @@ mod tests { 0, LoroValue::String("x".into()), ); - assert!(err.is_err()); + assert!(matches!( + err, + Err(CrdtError::BlockListPathUnresolved { ref segment, .. }) if segment == "blocks" + )); // The scalar must survive untouched — no silent replacement. match row.get("blocks") { Some(ValueOrContainer::Value(v)) => { @@ -516,13 +507,19 @@ mod tests { fn list_delete_on_missing_list_still_errors() { let state = setup_bare_row("doc-2"); let err = list_delete(state.doc(), "pages", "doc-2", "blocks", 0); - assert!(err.is_err()); + assert!(matches!( + err, + Err(CrdtError::BlockListPathUnresolved { .. }) + )); } #[test] fn list_move_on_missing_list_still_errors() { let state = setup_bare_row("doc-2"); let err = list_move(state.doc(), "pages", "doc-2", "blocks", 0, 1); - assert!(err.is_err()); + assert!(matches!( + err, + Err(CrdtError::BlockListPathUnresolved { .. }) + )); } } diff --git a/nodedb-crdt/src/state/bitemporal_archive.rs b/nodedb-crdt/src/state/bitemporal_archive.rs index 79b9ba6af..95f3f29fd 100644 --- a/nodedb-crdt/src/state/bitemporal_archive.rs +++ b/nodedb-crdt/src/state/bitemporal_archive.rs @@ -56,24 +56,30 @@ impl CrdtState { /// [`CrdtState::upsert`] — the archive has a non-trivial doc-size /// cost and is only meaningful when the collection participates in /// `AS OF` / audit queries. + /// + /// The archive entry and the new row commit together on return, on + /// success and on error. pub fn upsert_versioned( &self, collection: &str, row_id: &str, fields: &[(&str, LoroValue)], ) -> Result<()> { - if let Some((prior_sys_ms, prior_fields)) = self.prior_system_snapshot(collection, row_id) { - let archive = self.doc.get_map(HISTORY_ROOT); - let key = archive_key(collection, row_id, prior_sys_ms); - let slot = archive - .insert_container(&key, LoroMap::new()) - .map_err(|e| CrdtError::Loro(format!("archive insert: {e}")))?; - for (k, v) in &prior_fields { - slot.insert(k.as_str(), v.clone()) - .map_err(|e| CrdtError::Loro(format!("archive field: {e}")))?; + let prior = self.prior_system_snapshot(collection, row_id); + self.doc.mutate(|doc| { + if let Some((prior_sys_ms, prior_fields)) = prior { + let archive = doc.get_map(HISTORY_ROOT); + let key = archive_key(collection, row_id, prior_sys_ms); + let slot = archive + .insert_container(&key, LoroMap::new()) + .map_err(|e| CrdtError::Loro(format!("archive insert: {e}")))?; + for (k, v) in &prior_fields { + slot.insert(k.as_str(), v.clone()) + .map_err(|e| CrdtError::Loro(format!("archive field: {e}")))?; + } } - } - self.upsert(collection, row_id, fields) + Self::replace_row_fields(doc, collection, row_id, fields) + }) } /// Read the row as it was at `asof_ms` (system-time). Scans the @@ -133,25 +139,26 @@ impl CrdtState { /// collection. Returns the number of archive entries deleted. The /// live row is never touched — retention only reclaims history, so /// the current state of every logical row remains readable even - /// when the entire archive is pruned. + /// when the entire archive is pruned. The deletes commit on return. pub fn purge_history_before(&self, collection: &str, cutoff_ms: i64) -> Result { - let archive = self.doc.get_map(HISTORY_ROOT); - let victims: Vec = archive - .keys() - .filter_map(|k| { - let ks = k.to_string(); - let matches = parse_archive_key(&ks) - .is_some_and(|(c, _, ts)| c == collection && ts < cutoff_ms); - matches.then_some(ks) - }) - .collect(); - let count = victims.len(); - for key in victims { - archive - .delete(&key) - .map_err(|e| CrdtError::Loro(format!("archive delete: {e}")))?; - } - Ok(count) + self.doc.mutate(|doc| { + let archive = doc.get_map(HISTORY_ROOT); + let victims: Vec = archive + .keys() + .filter_map(|k| { + let ks = k.to_string(); + let matches = parse_archive_key(&ks) + .is_some_and(|(c, _, ts)| c == collection && ts < cutoff_ms); + matches.then_some(ks) + }) + .collect(); + for key in &victims { + archive + .delete(key) + .map_err(|e| CrdtError::Loro(format!("archive delete: {e}")))?; + } + Ok(victims.len()) + }) } /// Read the live row's (`_ts_system`, field-map) pair when both the diff --git a/nodedb-crdt/src/state/changed_rows.rs b/nodedb-crdt/src/state/changed_rows.rs index 23b10162a..e20db76c1 100644 --- a/nodedb-crdt/src/state/changed_rows.rs +++ b/nodedb-crdt/src/state/changed_rows.rs @@ -76,6 +76,20 @@ mod tests { .expect("upsert"); } + /// A local write left in Loro's open transaction. Every `CrdtState` + /// mutator commits on return, so only the raw document handle can leave + /// one open. + fn put_uncommitted(state: &CrdtState, row: &str, text: &str) { + state + .doc + .get_map("docs") + .insert_container(row, loro::LoroMap::new()) + .expect("row") + .insert("body", LoroValue::from(text)) + .expect("field"); + assert_ne!(state.doc.get_pending_txn_len(), 0); + } + fn rows(tracked: &super::TrackedImport) -> Vec<&str> { tracked.changed_rows.iter().map(String::as_str).collect() } @@ -226,8 +240,7 @@ mod tests { #[test] fn an_uncommitted_local_write_is_not_reported_as_imported() { let (source, target) = synced_pair(); - put(&target, "local", "mine"); - assert_ne!(target.doc.get_pending_txn_len(), 0); + put_uncommitted(&target, "local", "mine"); let before = source.oplog_version_vector(); put(&source, "c", "changed"); let delta = source.export_updates_since(&before).expect("delta"); @@ -239,11 +252,25 @@ mod tests { } #[test] - fn an_uncommitted_local_write_on_an_empty_doc_is_not_reported() { + fn a_committed_local_write_is_not_reported_as_imported() { let source = CrdtState::new(1).expect("source"); put(&source, "a", "one"); let target = CrdtState::new(2).expect("target"); put(&target, "local", "mine"); + assert_eq!(target.doc.get_pending_txn_len(), 0); + + let tracked = target.import_tracked("docs", &source.export_snapshot().expect("snapshot")); + tracked.outcome.as_ref().expect("import"); + assert_eq!(rows(&tracked), vec!["a"]); + assert!(target.row_exists("docs", "local")); + } + + #[test] + fn an_uncommitted_local_write_on_an_empty_doc_is_not_reported() { + let source = CrdtState::new(1).expect("source"); + put(&source, "a", "one"); + let target = CrdtState::new(2).expect("target"); + put_uncommitted(&target, "local", "mine"); let tracked = target.import_tracked("docs", &source.export_snapshot().expect("snapshot")); tracked.outcome.as_ref().expect("import"); @@ -345,7 +372,6 @@ mod tests { .expect("hub imports peer 1"); hub.import(&withheld).expect("hub imports peer 3"); put(&hub, "z", "dependent"); - hub.doc.commit(); let vv = hub.oplog_version_vector(); let end = |peer: u64| vv.get(&peer).copied().unwrap_or(0); let blob = hub diff --git a/nodedb-crdt/src/state/core.rs b/nodedb-crdt/src/state/core.rs index 3b651fdb6..93ef1a2ba 100644 --- a/nodedb-crdt/src/state/core.rs +++ b/nodedb-crdt/src/state/core.rs @@ -89,8 +89,8 @@ impl CrdtState { /// Fetch a row's existing `LoroMap` container, or create one if absent. /// Shared by `upsert` and `set_fields` — both need the same row handle /// before diverging on prune-vs-preserve semantics. - fn row_container(&self, collection: &str, row_id: &str) -> Result { - let coll = self.doc.get_map(collection); + fn row_container(doc: &LoroDoc, collection: &str, row_id: &str) -> Result { + let coll = doc.get_map(collection); match coll.get(row_id) { Some(ValueOrContainer::Container(loro::Container::Map(m))) => Ok(m), _ => coll @@ -128,23 +128,17 @@ impl CrdtState { Ok(()) } - /// Insert or update a row in a collection. + /// The body of [`CrdtState::upsert`], run inside an open mutation. /// - /// This is a REPLACE for scalar fields — every caller passes the - /// complete scalar projection, and any current scalar key absent from - /// `fields` is deleted. It reuses the row's existing `LoroMap` rather - /// than destroying and recreating it, because container-valued keys - /// (e.g. the Notion-style block list in `list_ops.rs`, stored as a - /// container-valued key inside this same row map) cannot be expressed in - /// `fields: &[(&str, LoroValue)]` at all — they are structurally out of - /// scope for this replace and must survive across every call. - pub fn upsert( - &self, + /// `upsert_versioned` calls it in the same mutation as its archive + /// write, so the archive entry and the new row commit together. + pub(in crate::state) fn replace_row_fields( + doc: &LoroDoc, collection: &str, row_id: &str, fields: &[(&str, LoroValue)], ) -> Result<()> { - let row_container = self.row_container(collection, row_id)?; + let row_container = Self::row_container(doc, collection, row_id)?; let incoming_keys: HashSet<&str> = fields.iter().map(|(field, _)| *field).collect(); @@ -169,30 +163,58 @@ impl CrdtState { Self::write_scalar_fields(&row_container, collection, row_id, fields) } + /// Insert or update a row in a collection. + /// + /// This is a REPLACE for scalar fields — every caller passes the + /// complete scalar projection, and any current scalar key absent from + /// `fields` is deleted. It reuses the row's existing `LoroMap` rather + /// than destroying and recreating it, because container-valued keys + /// (e.g. the Notion-style block list in `list_ops.rs`, stored as a + /// container-valued key inside this same row map) cannot be expressed in + /// `fields: &[(&str, LoroValue)]` at all — they are structurally out of + /// scope for this replace and must survive across every call. + /// + /// The write commits on return, on success and on error. + pub fn upsert( + &self, + collection: &str, + row_id: &str, + fields: &[(&str, LoroValue)], + ) -> Result<()> { + self.doc + .mutate(|doc| Self::replace_row_fields(doc, collection, row_id, fields)) + } + /// Partial-merge write: set exactly the provided scalar `fields` on a row /// (LWW-per-field), creating the row if absent, leaving every untouched /// key intact. This is `upsert` WITHOUT the full-projection prune step — /// the UPDATE-SET semantic for `CrdtOp::DocUpsert { partial: true }`. + /// + /// The write commits on return, on success and on error. pub fn set_fields( &self, collection: &str, row_id: &str, fields: &[(&str, LoroValue)], ) -> Result<()> { - let row_container = self.row_container(collection, row_id)?; - Self::write_scalar_fields(&row_container, collection, row_id, fields) + self.doc.mutate(|doc| { + let row_container = Self::row_container(doc, collection, row_id)?; + Self::write_scalar_fields(&row_container, collection, row_id, fields) + }) } - /// Delete a row from a collection. + /// Delete a row from a collection. The delete commits on return. pub fn delete(&self, collection: &str, row_id: &str) -> Result<()> { - let coll = self.doc.get_map(collection); - coll.delete(row_id) - .map_err(|e| CrdtError::Loro(e.to_string()))?; - Ok(()) + self.doc.mutate(|doc| { + doc.get_map(collection) + .delete(row_id) + .map_err(|e| CrdtError::Loro(e.to_string())) + }) } /// Insert a block-map into one row's movable list and populate scalar fields. - /// The raw document handle never leaves this state object. + /// The raw document handle never leaves this state object. The insert + /// commits on return, on success and on error. pub fn list_insert_fields( &self, collection: &str, @@ -201,18 +223,20 @@ impl CrdtState { index: usize, fields: &[(String, LoroValue)], ) -> Result<()> { - let block = crate::list_ops::list_insert_container( - &self.doc, collection, row_id, list_path, index, - )?; - for (key, value) in fields { - block - .insert(key.as_str(), value.clone()) - .map_err(|error| CrdtError::Loro(error.to_string()))?; - } - Ok(()) + self.doc.mutate(|doc| { + let block = + crate::list_ops::list_insert_container(doc, collection, row_id, list_path, index)?; + for (key, value) in fields { + block + .insert(key.as_str(), value.clone()) + .map_err(|error| CrdtError::Loro(error.to_string()))?; + } + Ok(()) + }) } - /// Delete one block from a row-owned movable list. + /// Delete one block from a row-owned movable list. The delete commits on + /// return. pub fn list_delete( &self, collection: &str, @@ -220,10 +244,12 @@ impl CrdtState { list_path: &str, index: usize, ) -> Result<()> { - crate::list_ops::list_delete(&self.doc, collection, row_id, list_path, index) + self.doc + .mutate(|doc| crate::list_ops::list_delete(doc, collection, row_id, list_path, index)) } - /// Move one block within a row-owned movable list. + /// Move one block within a row-owned movable list. The move commits on + /// return. pub fn list_move( &self, collection: &str, @@ -232,9 +258,9 @@ impl CrdtState { from_index: usize, to_index: usize, ) -> Result<()> { - crate::list_ops::list_move( - &self.doc, collection, row_id, list_path, from_index, to_index, - ) + self.doc.mutate(|doc| { + crate::list_ops::list_move(doc, collection, row_id, list_path, from_index, to_index) + }) } /// Return one row-owned movable list's length. @@ -254,15 +280,17 @@ impl CrdtState { } /// Delete all rows in a collection. Returns the number of rows deleted. + /// The deletes commit on return, on success and on error. pub fn clear_collection(&self, collection: &str) -> Result { - let coll = self.doc.get_map(collection); - let keys: Vec = coll.keys().map(|k| k.to_string()).collect(); - let count = keys.len(); - for key in &keys { - coll.delete(key) - .map_err(|e| CrdtError::Loro(e.to_string()))?; - } - Ok(count) + self.doc.mutate(|doc| { + let coll = doc.get_map(collection); + let keys: Vec = coll.keys().map(|k| k.to_string()).collect(); + for key in &keys { + coll.delete(key) + .map_err(|e| CrdtError::Loro(e.to_string()))?; + } + Ok(keys.len()) + }) } /// Read a single row's fields as a `LoroValue::Map`. @@ -770,9 +798,14 @@ mod tests { #[test] fn oplog_version_counts_uncommitted_operations() { let state = CrdtState::new(1).unwrap(); + // Every mutator commits on return, so the raw handle is the only way + // to hold an operation in an open transaction. state - .upsert("items", "a", &[("value", LoroValue::I64(1))]) + .doc + .get_map("items") + .insert("a", LoroValue::I64(1)) .unwrap(); + assert_ne!(state.doc.get_pending_txn_len(), 0); // The memory estimate is cached against this version vector. That is only // sound because Loro advances it when the operation is written rather than @@ -1016,4 +1049,159 @@ mod tests { the collections, not to this list" ); } + + /// The operation count this state has authored, after checking that the + /// write named `what` left no open Loro transaction and authored at + /// least one operation past `before`. + fn committed(state: &CrdtState, before: i32, what: &str) -> i32 { + assert_eq!( + state.doc.get_pending_txn_len(), + 0, + "{what} left its operations in an open Loro transaction" + ); + let after = state.local_op_counter(); + assert!(after > before, "{what} authored no operation"); + after + } + + #[test] + fn every_mutator_commits_before_it_returns() { + let state = CrdtState::new(1).unwrap(); + let mut ops = state.local_op_counter(); + + state + .upsert( + "docs", + "r", + &[("a", LoroValue::I64(1)), ("b", LoroValue::I64(2))], + ) + .unwrap(); + ops = committed(&state, ops, "upsert"); + state + .set_fields("docs", "r", &[("b", LoroValue::I64(3))]) + .unwrap(); + ops = committed(&state, ops, "set_fields"); + for (index, id) in ["b0", "b1"].into_iter().enumerate() { + state + .list_insert_fields( + "docs", + "r", + "blocks", + index, + &[("id".into(), LoroValue::String(id.into()))], + ) + .unwrap(); + ops = committed(&state, ops, "list_insert_fields"); + } + state.list_move("docs", "r", "blocks", 0, 1).unwrap(); + ops = committed(&state, ops, "list_move"); + state.list_delete("docs", "r", "blocks", 0).unwrap(); + ops = committed(&state, ops, "list_delete"); + assert_eq!(state.remove_fields("docs", "r", &["a"]).unwrap(), 1); + ops = committed(&state, ops, "remove_fields"); + state + .restore_row_image("docs", "r", &crate::state::RowImage::Absent) + .unwrap(); + ops = committed(&state, ops, "restore_row_image"); + + for sys_ms in [10, 20] { + state + .upsert_versioned( + "hist", + "h", + &[ + ("_ts_system", LoroValue::I64(sys_ms)), + ("v", LoroValue::I64(sys_ms)), + ], + ) + .unwrap(); + ops = committed(&state, ops, "upsert_versioned"); + } + assert_eq!(state.archive_version_count("hist", "h"), 1); + assert_eq!(state.purge_history_before("hist", 15).unwrap(), 1); + ops = committed(&state, ops, "purge_history_before"); + + state + .upsert("docs", "r2", &[("v", LoroValue::I64(1))]) + .unwrap(); + ops = committed(&state, ops, "upsert"); + let version = state.oplog_version_vector(); + state + .upsert("docs", "r2", &[("v", LoroValue::I64(2))]) + .unwrap(); + ops = committed(&state, ops, "upsert"); + let delta = state.restore_to_version("docs", "r2", &version).unwrap(); + assert!(!delta.is_empty()); + ops = committed(&state, ops, "restore_to_version"); + state.delete("docs", "r2").unwrap(); + ops = committed(&state, ops, "delete"); + assert_eq!(state.clear_collection("hist").unwrap(), 1); + committed(&state, ops, "clear_collection"); + } + + #[test] + fn a_version_taken_after_a_write_restores_through_preview() { + let state = CrdtState::new(1).unwrap(); + state + .upsert("docs", "doc", &[("title", LoroValue::String("v1".into()))]) + .unwrap(); + // A checkpoint reads the version right after the acknowledged write. + let checkpoint = state.oplog_version_vector(); + assert!( + state + .preview_restore_to_version("docs", "doc", &checkpoint) + .expect("a restore preview right after a write must not be refused") + .is_empty(), + "the document is already at the checkpoint" + ); + + state + .set_fields("docs", "doc", &[("title", LoroValue::String("v2".into()))]) + .unwrap(); + let delta = state + .preview_restore_to_version("docs", "doc", &checkpoint) + .expect("preview after a later write"); + assert!(!delta.is_empty()); + state.import(&delta).unwrap(); + assert_eq!( + state.read_field("docs", "doc", "title"), + Some(LoroValue::String("v1".into())), + "the checkpoint must hold the write acknowledged before it" + ); + } + + #[test] + fn a_scalar_write_refused_part_way_commits_and_its_row_image_undoes_it() { + let state = CrdtState::new(1).unwrap(); + state + .upsert("docs", "r", &[("a", LoroValue::I64(1))]) + .unwrap(); + attach_nested_block_list(&state, "docs", "r"); + let image = state.row_image("docs", "r").unwrap(); + + // `a` is written before `blocks` is refused, for both scalar writes. + let partial = [("a", LoroValue::I64(2)), ("blocks", LoroValue::I64(3))]; + for partial_merge in [true, false] { + let refusal = if partial_merge { + state.set_fields("docs", "r", &partial) + } else { + state.upsert("docs", "r", &partial) + }; + assert!(matches!( + refusal, + Err(CrdtError::ScalarFieldShadowsContainer { ref field, .. }) if field == "blocks" + )); + assert_eq!( + state.doc.get_pending_txn_len(), + 0, + "a refused write must not leave its partial operations open" + ); + assert_eq!(state.read_field("docs", "r", "a"), Some(LoroValue::I64(2))); + + state.restore_row_image("docs", "r", &image).unwrap(); + assert_eq!(state.doc.get_pending_txn_len(), 0); + assert_eq!(state.read_field("docs", "r", "a"), Some(LoroValue::I64(1))); + assert_eq!(state.list_length("docs", "r", "blocks").unwrap(), 1); + } + } } diff --git a/nodedb-crdt/src/state/document_cell.rs b/nodedb-crdt/src/state/document_cell.rs index b5a5a3dc2..bdba108f2 100644 --- a/nodedb-crdt/src/state/document_cell.rs +++ b/nodedb-crdt/src/state/document_cell.rs @@ -7,6 +7,8 @@ use std::ops::Deref; use loro::LoroDoc; +use crate::error::Result; + /// A measured point relating the document's operation count to its encoded /// size, taken from one real snapshot export. struct Calibration { @@ -82,6 +84,28 @@ impl DocumentCell { *self = Self::new(doc); } + /// Run one local mutation, then commit its Loro transaction. + /// + /// Every `CrdtState` method that writes the document goes through here. + /// The authoritative document is then quiescent between calls. A version + /// vector read after the call covers the write, and a fork or preview + /// never meets a pending transaction. + /// + /// The commit runs on success and on error. Loro cannot roll back an + /// operation it has applied. A failed mutation therefore commits the + /// operations it authored before the error. A caller that needs the write + /// to be all-or-nothing captures a `RowImage` first. On error it calls + /// `restore_row_image`, which authors and commits the compensating + /// operations. + pub(in crate::state) fn mutate( + &self, + write: impl FnOnce(&LoroDoc) -> Result, + ) -> Result { + let outcome = write(&self.doc); + self.doc.commit(); + outcome + } + /// Estimated encoded size in bytes, as a proxy for memory footprint. /// /// Loro exposes no direct memory metric, and a snapshot export — the diff --git a/nodedb-crdt/src/state/history.rs b/nodedb-crdt/src/state/history.rs index 238a8b2fc..c5472675f 100644 --- a/nodedb-crdt/src/state/history.rs +++ b/nodedb-crdt/src/state/history.rs @@ -155,6 +155,9 @@ impl CrdtState { /// restoring would not change the live row (e.g. restoring to the /// version the document is already at). /// + /// The restore operations commit before the export, on success and on + /// error. + /// /// Historical fields are inspected on the *live* forked container (via /// `LoroMap::get` → `ValueOrContainer`), not the flattened /// `read_at_version` projection: scalar entries are replaced the same @@ -174,7 +177,8 @@ impl CrdtState { return Ok(Vec::new()); } let vv_before = self.doc.oplog_vv(); - apply_restore_to_document(&self.doc, collection, row_id, &historical)?; + self.doc + .mutate(|doc| apply_restore_to_document(doc, collection, row_id, &historical))?; self.doc .export(loro::ExportMode::updates(&vv_before)) .map_err(|e| CrdtError::Loro(format!("restore delta export: {e}"))) @@ -238,10 +242,20 @@ fn historical_row( let forked = fork_at_version(doc, version)?; match forked.get_map(collection).get(row_id) { Some(ValueOrContainer::Container(loro::Container::Map(row))) => Ok(row), - Some(_) => Err(CrdtError::Loro("historical state is not a map".into())), - None => Err(CrdtError::Loro( - "document did not exist at target version".into(), - )), + Some(other) => Err(CrdtError::NonMapRowValue { + collection: collection.to_string(), + row_id: row_id.to_string(), + value: match other { + ValueOrContainer::Value(value) => format!("the scalar {value:?}"), + ValueOrContainer::Container(container) => { + format!("a {:?} container", container.get_type()) + } + }, + }), + None => Err(CrdtError::RowAbsentAtVersion { + collection: collection.to_string(), + row_id: row_id.to_string(), + }), } } @@ -306,7 +320,6 @@ mod tests { &[("title", LoroValue::String(title.into()))], ) .expect("write"); - state.doc.commit(); versions.push(state.oplog_version_vector()); } History { @@ -541,7 +554,6 @@ mod tests { state .upsert("docs", "doc-1", &[("body", LoroValue::String("v1".into()))]) .expect("write"); - state.doc.commit(); let current = state.oplog_version_vector(); assert!( state @@ -559,14 +571,12 @@ mod tests { state .upsert("docs", "a", &[("v", LoroValue::I64(1))]) .expect("write a"); - state.doc.commit(); let after_a = state.local_op_counter(); assert!(after_a > 0); state .upsert("docs", "b", &[("v", LoroValue::I64(2))]) .expect("write b"); - state.doc.commit(); let after_b = state.local_op_counter(); assert!(after_b > after_a); } @@ -579,13 +589,11 @@ mod tests { state .upsert("docs", "a", &[("v", LoroValue::I64(1))]) .expect("write a"); - state.doc.commit(); let end_a = state.local_op_counter(); state .upsert("docs", "b", &[("v", LoroValue::I64(2))]) .expect("write b"); - state.doc.commit(); let row_a_delta = state .export_local_range(start_a, end_a) @@ -604,7 +612,6 @@ mod tests { state .upsert("docs", "a", &[("v", LoroValue::I64(1))]) .expect("write a"); - state.doc.commit(); let delta = state.export_local_range(5, 5).expect("empty range export"); assert!(delta.is_empty()); @@ -619,4 +626,29 @@ mod tests { ); assert!(!target.row_exists("docs", "a")); } + + #[test] + fn restoring_a_row_absent_at_the_version_is_a_typed_refusal() { + let state = CrdtState::new(1).expect("state"); + let before = state.oplog_version_vector(); + state + .upsert("docs", "late", &[("v", LoroValue::I64(1))]) + .expect("write"); + + for result in [ + state.preview_restore_to_version("docs", "late", &before), + state.restore_to_version("docs", "late", &before), + ] { + match result { + Err(CrdtError::RowAbsentAtVersion { collection, row_id }) => { + assert_eq!((collection.as_str(), row_id.as_str()), ("docs", "late")); + } + other => panic!("expected RowAbsentAtVersion, got {other:?}"), + } + } + assert_eq!( + state.read_field("docs", "late", "v"), + Some(LoroValue::I64(1)) + ); + } } diff --git a/nodedb-crdt/src/state/preview.rs b/nodedb-crdt/src/state/preview.rs index 27d04bb75..3add2be6a 100644 --- a/nodedb-crdt/src/state/preview.rs +++ b/nodedb-crdt/src/state/preview.rs @@ -361,9 +361,6 @@ mod tests { let state = CrdtState::new(1).expect("state"); let delta = delta_for("docs", "one", "new"); - // Make the authoritative baseline explicitly committed before any - // snapshot/frontier observation used by this non-mutation regression. - state.doc.commit(); let snapshot_before = state.export_snapshot().expect("before snapshot"); let frontier_before = state.oplog_version_vector(); let preview = state @@ -389,10 +386,17 @@ mod tests { #[test] fn rejects_preview_when_source_has_pending_transaction_without_mutating_it() { let state = CrdtState::new(1).expect("state"); + // Every `CrdtState` mutator commits on return. A pending transaction + // is only reachable through the raw document handle. state - .upsert("local", "pending", &[("value", LoroValue::I64(7))]) + .doc + .get_map("local") + .insert_container("pending", loro::LoroMap::new()) + .expect("pending row") + .insert("value", LoroValue::I64(7)) .expect("local pending write"); let pending_before = state.doc.get_pending_txn_len(); + assert_ne!(pending_before, 0); let row_before = state.read_row("local", "pending"); let delta = delta_for("docs", "one", "remote"); @@ -432,11 +436,6 @@ mod tests { &[("value", LoroValue::String("base".into()))], ) .expect("base"); - // Loro's auto-commit transaction is finalized only at an explicit - // commit or an import/export boundary. Commit the prerequisite so - // the incremental delta below depends on an operation absent from - // the preview target, rather than exporting both writes together. - source.doc.commit(); let version = source.oplog_version_vector(); source .set_fields( @@ -464,7 +463,6 @@ mod tests { &[("value", LoroValue::String("base".into()))], ) .expect("base"); - source.doc.commit(); let version = source.oplog_version_vector(); source .set_fields( @@ -675,7 +673,6 @@ mod tests { source .upsert("docs", "new", &[("value", LoroValue::I64(2))]) .expect("new row"); - source.doc.commit(); let expected_new_ops = positive_oplog_advance(&receiver_oplog, &source.oplog_version_vector()) .expect("monotonic source oplog"); diff --git a/nodedb-crdt/src/state/remove_fields.rs b/nodedb-crdt/src/state/remove_fields.rs index d6831ac21..ca6135f3d 100644 --- a/nodedb-crdt/src/state/remove_fields.rs +++ b/nodedb-crdt/src/state/remove_fields.rs @@ -18,7 +18,8 @@ impl CrdtState { /// A container-valued key is refused with `ScalarFieldShadowsContainer`, /// as `set_fields` refuses it: deleting it discards its nested CRDT state /// (e.g. a row's block list). Every field is checked before any is - /// deleted, so a refused call deletes nothing. + /// deleted, so a refused call deletes nothing. The deletes commit on + /// return. pub fn remove_fields(&self, collection: &str, row_id: &str, fields: &[&str]) -> Result { let coll = self.doc.get_map(collection); let row: LoroMap = match coll.get(row_id) { @@ -35,15 +36,17 @@ impl CrdtState { field: (*field).to_string(), }); } - let mut removed = 0; - for field in fields { - if row.get(field).is_some() { - row.delete(field) - .map_err(|e| CrdtError::Loro(e.to_string()))?; - removed += 1; + self.doc.mutate(|_| { + let mut removed = 0; + for field in fields { + if row.get(field).is_some() { + row.delete(field) + .map_err(|e| CrdtError::Loro(e.to_string()))?; + removed += 1; + } } - } - Ok(removed) + Ok(removed) + }) } } diff --git a/nodedb-crdt/src/state/row_image.rs b/nodedb-crdt/src/state/row_image.rs index 96b357cd6..b735f77f3 100644 --- a/nodedb-crdt/src/state/row_image.rs +++ b/nodedb-crdt/src/state/row_image.rs @@ -66,56 +66,61 @@ impl CrdtState { /// Container-valued fields keep their state. A `Fields` image needs the /// row to still be a map: a scalar write keeps the row's map, so any /// other row shape is refused with `NonMapRowValue`. + /// + /// The compensating operations commit on return, on success and on + /// error. This is the undo path for a failed `upsert` or `set_fields`. pub fn restore_row_image( &self, collection: &str, row_id: &str, image: &RowImage, ) -> Result<()> { - let coll = self.doc.get_map(collection); - match image { - RowImage::Absent => { - if coll.get(row_id).is_some() { - coll.delete(row_id).map_err(loro_error)?; - } - Ok(()) - } - RowImage::Value(value) => coll.insert(row_id, value.clone()).map_err(loro_error), - RowImage::Fields(fields) => { - let row = match coll.get(row_id) { - Some(ValueOrContainer::Container(loro::Container::Map(row))) => row, - other => { - return Err(CrdtError::NonMapRowValue { - collection: collection.to_string(), - row_id: row_id.to_string(), - value: match other { - None => "nothing".to_string(), - Some(ValueOrContainer::Value(value)) => { - format!("the scalar {value:?}") - } - Some(ValueOrContainer::Container(container)) => { - format!("a {:?} container", container.get_type()) - } - }, - }); - } - }; - for (key, _) in scalar_fields(&row) { - if !fields.iter().any(|(field, _)| *field == key) { - row.delete(&key).map_err(loro_error)?; + self.doc.mutate(|doc| { + let coll = doc.get_map(collection); + match image { + RowImage::Absent => { + if coll.get(row_id).is_some() { + coll.delete(row_id).map_err(loro_error)?; } + Ok(()) } - for (field, value) in fields { - match row.get(field) { - Some(ValueOrContainer::Value(current)) if current == *value => {} - _ => { - row.insert(field, value.clone()).map_err(loro_error)?; + RowImage::Value(value) => coll.insert(row_id, value.clone()).map_err(loro_error), + RowImage::Fields(fields) => { + let row = match coll.get(row_id) { + Some(ValueOrContainer::Container(loro::Container::Map(row))) => row, + other => { + return Err(CrdtError::NonMapRowValue { + collection: collection.to_string(), + row_id: row_id.to_string(), + value: match other { + None => "nothing".to_string(), + Some(ValueOrContainer::Value(value)) => { + format!("the scalar {value:?}") + } + Some(ValueOrContainer::Container(container)) => { + format!("a {:?} container", container.get_type()) + } + }, + }); + } + }; + for (key, _) in scalar_fields(&row) { + if !fields.iter().any(|(field, _)| *field == key) { + row.delete(&key).map_err(loro_error)?; + } + } + for (field, value) in fields { + match row.get(field) { + Some(ValueOrContainer::Value(current)) if current == *value => {} + _ => { + row.insert(field, value.clone()).map_err(loro_error)?; + } } } + Ok(()) } - Ok(()) } - } + }) } } diff --git a/nodedb/src/bridge/envelope/crdt_error_code.rs b/nodedb/src/bridge/envelope/crdt_error_code.rs new file mode 100644 index 000000000..236a13427 --- /dev/null +++ b/nodedb/src/bridge/envelope/crdt_error_code.rs @@ -0,0 +1,296 @@ +// SPDX-License-Identifier: BUSL-1.1 + +//! The Data-Plane verdict for each CRDT engine error. +//! +//! A client request causes most CRDT errors: a malformed delta, a version +//! that no longer exists, a write that names an absent document. Each one +//! keeps a client class here. Only server faults and broken invariants +//! answer `Internal`, which a client sees as `XX000`. + +use nodedb_crdt::CrdtError; + +use super::error_code::ErrorCode; + +/// Exhaustive, so a new CRDT error picks its own class rather than +/// defaulting to `Internal`. +impl From<&CrdtError> for ErrorCode { + fn from(e: &CrdtError) -> Self { + match e { + CrdtError::ConstraintViolation { + constraint, + collection, + detail, + } => Self::RejectedConstraint { + constraint: constraint.clone(), + detail: format!("{detail} (collection `{collection}`)"), + }, + // Delta or snapshot bytes the client sent are malformed, past a + // limit, or of a shape no collection holds. Class `22`. + client @ (CrdtError::ImportTooLarge { .. } + | CrdtError::ImportMalformed { .. } + | CrdtError::ImportInvalidOperationRange + | CrdtError::ImportOperationLimitExceeded { .. } + | CrdtError::PreviewDeltaTooLarge { .. } + | CrdtError::PreviewMalformed { .. } + | CrdtError::PreviewInvalidOperationRange + | CrdtError::PreviewOperationLimitExceeded { .. } + | CrdtError::PreviewWriteSetLimitExceeded { .. } + | CrdtError::PreviewTargetMismatch { .. } + | CrdtError::PreviewPostImageTooLarge { .. } + | CrdtError::NonMapRootContainer { .. } + | CrdtError::NonMapRowValue { .. } + | CrdtError::BlockListPathUnresolved { .. } + | CrdtError::BlockListIndexOutOfBounds { .. }) => Self::DataException { + detail: client.to_string(), + }, + // A value of the wrong kind for the field or path it targets. + // `42804`. + client @ (CrdtError::ScalarFieldShadowsContainer { .. } + | CrdtError::BlockListNotMovable { .. }) => Self::DatatypeMismatch { + detail: client.to_string(), + }, + CrdtError::RowAbsentAtVersion { collection, row_id } => Self::UndefinedObject { + object: format!( + "document \"{row_id}\" in collection \"{collection}\" at the target version" + ), + }, + CrdtError::BlockListRowAbsent { collection, row_id } => Self::UndefinedObject { + object: format!("document \"{row_id}\" in collection \"{collection}\""), + }, + CrdtError::UnknownCollection(collection) => Self::UndefinedObject { + object: format!("collection \"{collection}\""), + }, + client @ CrdtError::VersionBeforeCompactionBoundary { .. } => { + Self::ObjectNotInPrerequisiteState { + object: "CRDT version".into(), + detail: client.to_string(), + } + } + // Nothing applied. The same bytes succeed once the missing + // causal history arrives, or once the dead-letter queue drains. + retry @ (CrdtError::ImportPendingDependencies + | CrdtError::PreviewPendingDependencies + | CrdtError::DlqFull { .. }) => Self::RetryableRefusal { + reason: retry.to_string(), + }, + denied @ (CrdtError::AuthExpired { .. } | CrdtError::InvalidSignature { .. }) => { + Self::RejectedAuthz { + resource: denied.to_string(), + } + } + // A replayed sequence number fails the same way every time. + replay @ CrdtError::ReplayDetected { .. } => Self::RejectedPrevalidation { + reason: replay.to_string(), + }, + // Server faults: a Loro failure, a pending transaction on the + // authoritative document, or an apply this node could not run. + fault @ (CrdtError::DeltaApplyFailed(_) + | CrdtError::PreviewSourceTransactionPending { .. } + | CrdtError::Loro(_)) => Self::Internal { + detail: fault.to_string(), + }, + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn code(e: CrdtError) -> ErrorCode { + ErrorCode::from(&e) + } + + #[test] + fn a_missing_document_is_an_undefined_object() { + let absent = code(CrdtError::RowAbsentAtVersion { + collection: "notes".into(), + row_id: "doc".into(), + }); + assert!( + matches!(&absent, ErrorCode::UndefinedObject { object } if object.contains("\"doc\"")), + "got {absent:?}" + ); + assert!(matches!( + code(CrdtError::BlockListRowAbsent { + collection: "pages".into(), + row_id: "p".into(), + }), + ErrorCode::UndefinedObject { .. } + )); + assert!(matches!( + code(CrdtError::UnknownCollection("c".into())), + ErrorCode::UndefinedObject { .. } + )); + } + + #[test] + fn a_version_below_the_compaction_boundary_is_not_in_prerequisite_state() { + assert!(matches!( + code(CrdtError::VersionBeforeCompactionBoundary { + peer: 1, + requested: 2, + discarded: 3, + }), + ErrorCode::ObjectNotInPrerequisiteState { .. } + )); + } + + #[test] + fn malformed_client_input_is_a_data_exception() { + for e in [ + CrdtError::ImportMalformed { detail: "x".into() }, + CrdtError::ImportTooLarge { + limit: 1, + actual: 2, + }, + CrdtError::ImportInvalidOperationRange, + CrdtError::PreviewMalformed { detail: "x".into() }, + CrdtError::PreviewDeltaTooLarge { + limit: 1, + actual: 2, + }, + CrdtError::PreviewPostImageTooLarge { + limit: 1, + actual: 2, + }, + CrdtError::NonMapRowValue { + collection: "c".into(), + row_id: "r".into(), + value: "a List container".into(), + }, + CrdtError::BlockListIndexOutOfBounds { + list_path: "blocks".into(), + index: 9, + len: 1, + }, + ] { + let mapped = code(e); + assert!( + matches!(mapped, ErrorCode::DataException { .. }), + "got {mapped:?}" + ); + } + } + + #[test] + fn a_wrong_kind_of_value_is_a_datatype_mismatch() { + assert!(matches!( + code(CrdtError::ScalarFieldShadowsContainer { + collection: "c".into(), + row_id: "r".into(), + field: "blocks".into(), + }), + ErrorCode::DatatypeMismatch { .. } + )); + } + + #[test] + fn missing_causal_history_is_retryable() { + for e in [ + CrdtError::ImportPendingDependencies, + CrdtError::PreviewPendingDependencies, + ] { + assert!(matches!(code(e), ErrorCode::RetryableRefusal { .. })); + } + } + + #[test] + fn a_constraint_violation_keeps_its_constraint() { + let mapped = code(CrdtError::ConstraintViolation { + constraint: "users_email_unique".into(), + collection: "users".into(), + detail: "duplicate".into(), + }); + assert!( + matches!(&mapped, ErrorCode::RejectedConstraint { constraint, detail } + if constraint == "users_email_unique" && detail.contains("users")), + "got {mapped:?}" + ); + } + + #[test] + fn only_server_faults_stay_internal() { + for e in [ + CrdtError::PreviewSourceTransactionPending { operations: 1 }, + CrdtError::Loro("boom".into()), + CrdtError::DeltaApplyFailed("boom".into()), + ] { + assert!(matches!(code(e), ErrorCode::Internal { .. })); + } + } + + /// pgwire, the native public code and a node hop all render one class. + #[test] + fn a_crdt_error_renders_its_class_on_every_surface() { + use nodedb_types::error::sqlstate; + + use crate::control::server::pgwire::types::error_map::numeric_code_to_sqlstate; + use crate::control::server::pgwire::types::error_to_sqlstate; + + let cases = [ + ( + CrdtError::RowAbsentAtVersion { + collection: "notes".into(), + row_id: "doc".into(), + }, + sqlstate::UNDEFINED_OBJECT, + ), + ( + CrdtError::VersionBeforeCompactionBoundary { + peer: 1, + requested: 2, + discarded: 3, + }, + sqlstate::OBJECT_NOT_IN_PREREQUISITE_STATE, + ), + ( + CrdtError::ImportMalformed { detail: "x".into() }, + sqlstate::DATA_EXCEPTION, + ), + ( + CrdtError::ScalarFieldShadowsContainer { + collection: "c".into(), + row_id: "r".into(), + field: "blocks".into(), + }, + sqlstate::DATATYPE_MISMATCH, + ), + (CrdtError::Loro("boom".into()), sqlstate::INTERNAL_ERROR), + ]; + for (crdt, expected) in cases { + let label = crdt.to_string(); + let err = crate::Error::Crdt(crdt); + let (_, pg_state, _) = error_to_sqlstate(&err); + assert_eq!(pg_state, expected, "{label}"); + + let public = crate::error_classify::classify(&err); + assert_eq!( + numeric_code_to_sqlstate(public.code()).get(..2), + expected.get(..2), + "{label}: public code {}", + public.code() + ); + + let hopped = + crate::Error::from(nodedb_cluster::rpc_codec::TypedClusterError::from(err)); + assert_eq!( + error_to_sqlstate(&hopped).1, + expected, + "{label} after a hop" + ); + } + } + + #[test] + fn a_crdt_error_crossing_the_bridge_keeps_its_class() { + let mapped = ErrorCode::from(crate::Error::Crdt(CrdtError::RowAbsentAtVersion { + collection: "notes".into(), + row_id: "doc".into(), + })); + assert!( + matches!(mapped, ErrorCode::UndefinedObject { .. }), + "got {mapped:?}" + ); + } +} diff --git a/nodedb/src/bridge/envelope/error_code.rs b/nodedb/src/bridge/envelope/error_code.rs index 9a348a9c8..a63576fcb 100644 --- a/nodedb/src/bridge/envelope/error_code.rs +++ b/nodedb/src/bridge/envelope/error_code.rs @@ -221,4 +221,13 @@ pub enum ErrorCode { /// An instant outside the range a `TIMESTAMP` or `TIMESTAMPTZ` holds. /// SQLSTATE `22008` (datetime_field_overflow). DatetimeFieldOverflow { detail: String }, + /// The request names an object that does not exist, such as a document + /// at a version where it had no state. `object` describes it, and the + /// message reads "`object` does not exist". SQLSTATE `42704` + /// (undefined_object). + UndefinedObject { object: String }, + /// The object exists but is not in the state the request needs, such as + /// a version whose history compaction discarded. SQLSTATE `55000` + /// (object_not_in_prerequisite_state). + ObjectNotInPrerequisiteState { object: String, detail: String }, } diff --git a/nodedb/src/bridge/envelope/error_code_from.rs b/nodedb/src/bridge/envelope/error_code_from.rs index 0f3dd6a2e..c94d78be3 100644 --- a/nodedb/src/bridge/envelope/error_code_from.rs +++ b/nodedb/src/bridge/envelope/error_code_from.rs @@ -247,16 +247,23 @@ impl From for ErrorCode { resource: e.to_string(), } } - // Client errors of class `42`, and client errors whose class - // (`25006`, `55`) no Data-Plane code has. `BadRequest` is the - // class their public code has. + // `42704` and `55000`, as the Control Plane gives them. + crate::Error::UndefinedObject { kind, name } => Self::UndefinedObject { + object: format!("{kind} \"{name}\""), + }, + crate::Error::ObjectNotInPrerequisiteState { object, detail } => { + Self::ObjectNotInPrerequisiteState { object, detail } + } + // Each CRDT error keeps the class of what caused it. + crate::Error::Crdt(crdt) => Self::from(&crdt), + // Client errors of class `42`, and client errors whose dedicated + // SQLSTATE (`25006`, a `55` code of their own) no Data-Plane code + // has. `BadRequest` is the class their public code has. e @ (crate::Error::CrdtAdmissionInvalidPlan { .. } | crate::Error::CrdtAdmissionCallerFence | crate::Error::CrdtApplyRequiresAdmission | crate::Error::CloneWriteRequiresMaterialize { .. } - | crate::Error::ObjectNotInPrerequisiteState { .. } | crate::Error::MirrorReadOnly { .. } - | crate::Error::UndefinedObject { .. } | crate::Error::AmbiguousColumn { .. } | crate::Error::ExecutionLimitExceeded { .. } | crate::Error::LimitExceeded { .. } @@ -298,7 +305,6 @@ impl From for ErrorCode { | crate::Error::Serialization { .. } | crate::Error::Codec { .. } | crate::Error::SegmentCorrupted { .. } - | crate::Error::Crdt(_) | crate::Error::Io(_) | crate::Error::Config { .. } | crate::Error::Encryption { .. } diff --git a/nodedb/src/bridge/envelope/mod.rs b/nodedb/src/bridge/envelope/mod.rs index ca2f6ea70..371b131b2 100644 --- a/nodedb/src/bridge/envelope/mod.rs +++ b/nodedb/src/bridge/envelope/mod.rs @@ -2,6 +2,7 @@ //! Request/response envelopes exchanged over the SPSC bridge. +pub mod crdt_error_code; pub mod error_code; pub mod error_code_from; pub mod payload; diff --git a/nodedb/src/control/cluster/data_plane_error_wire.rs b/nodedb/src/control/cluster/data_plane_error_wire.rs index 1fdaf2a6c..e4d6de18c 100644 --- a/nodedb/src/control/cluster/data_plane_error_wire.rs +++ b/nodedb/src/control/cluster/data_plane_error_wire.rs @@ -158,6 +158,10 @@ impl From for DataPlaneErrorCode { ErrorCode::DatatypeMismatch { detail } => Self::DatatypeMismatch { detail }, ErrorCode::InvalidDatetimeFormat { detail } => Self::InvalidDatetimeFormat { detail }, ErrorCode::DatetimeFieldOverflow { detail } => Self::DatetimeFieldOverflow { detail }, + ErrorCode::UndefinedObject { object } => Self::UndefinedObject { object }, + ErrorCode::ObjectNotInPrerequisiteState { object, detail } => { + Self::ObjectNotInPrerequisiteState { object, detail } + } ErrorCode::TextColumn { collection, column, @@ -325,6 +329,10 @@ impl From for ErrorCode { DataPlaneErrorCode::DatetimeFieldOverflow { detail } => { Self::DatetimeFieldOverflow { detail } } + DataPlaneErrorCode::UndefinedObject { object } => Self::UndefinedObject { object }, + DataPlaneErrorCode::ObjectNotInPrerequisiteState { object, detail } => { + Self::ObjectNotInPrerequisiteState { object, detail } + } DataPlaneErrorCode::TextColumn { collection, column, @@ -371,6 +379,22 @@ mod tests { assert_eq!(ErrorCode::from(wire), ErrorCode::DivisionByZero); } + #[test] + fn object_state_codes_roundtrip_verbatim() { + for original in [ + ErrorCode::UndefinedObject { + object: "document \"doc\" in collection \"notes\"".into(), + }, + ErrorCode::ObjectNotInPrerequisiteState { + object: "CRDT version".into(), + detail: "version predates the compaction boundary".into(), + }, + ] { + let wire = DataPlaneErrorCode::from(original.clone()); + assert_eq!(ErrorCode::from(wire), original); + } + } + #[test] fn payload_bearing_code_roundtrips_verbatim() { let original = ErrorCode::RejectedConstraint { diff --git a/nodedb/src/control/cluster/execution_error_wire.rs b/nodedb/src/control/cluster/execution_error_wire.rs index 7d490a588..89323fa95 100644 --- a/nodedb/src/control/cluster/execution_error_wire.rs +++ b/nodedb/src/control/cluster/execution_error_wire.rs @@ -19,6 +19,11 @@ use nodedb_cluster::rpc_codec::{DataPlaneErrorCode, TypedClusterError}; pub(crate) fn execution_error_to_typed(err: crate::Error) -> TypedClusterError { match err { crate::Error::DataPlane(code) => TypedClusterError::DataPlane { code: code.into() }, + // A CRDT error crosses as its Data-Plane verdict, so the coordinator + // answers the SQLSTATE a local execution answers. + crate::Error::Crdt(crdt) => TypedClusterError::DataPlane { + code: crate::bridge::envelope::ErrorCode::from(&crdt).into(), + }, // A statement that ran out of time keeps the wire's own deadline // variant, which the coordinator rebuilds as `Error::DeadlineExceeded`. // Folding it into `Internal` will report a client's own timeout as an @@ -141,7 +146,6 @@ pub(crate) fn execution_error_to_typed(err: crate::Error) -> TypedClusterError { | crate::Error::SegmentCorrupted { .. } | crate::Error::MemoryExhausted { .. } | crate::Error::Backpressure { .. } - | crate::Error::Crdt(_) | crate::Error::Io(_) | crate::Error::Config { .. } | crate::Error::Encryption { .. } diff --git a/nodedb/src/control/gateway/error_map/class_parity/code_samples.rs b/nodedb/src/control/gateway/error_map/class_parity/code_samples.rs index 017889eaa..5b67f3f64 100644 --- a/nodedb/src/control/gateway/error_map/class_parity/code_samples.rs +++ b/nodedb/src/control/gateway/error_map/class_parity/code_samples.rs @@ -8,7 +8,7 @@ use nodedb_types::sync::wire::SyncProvenance; use crate::bridge::envelope::{CounterFault, ErrorCode, SyncHold}; /// The number of `ErrorCode` variants [`variant_index`] numbers. -pub(super) const VARIANT_COUNT: usize = 50; +pub(super) const VARIANT_COUNT: usize = 52; /// A dense index per variant. Exhaustive, so a new variant fails to compile /// here until it gets an index, and [`every_variant_has_a_sample`] then fails @@ -65,6 +65,8 @@ pub(super) fn variant_index(code: &ErrorCode) -> usize { ErrorCode::DatatypeMismatch { .. } => 47, ErrorCode::InvalidDatetimeFormat { .. } => 48, ErrorCode::DatetimeFieldOverflow { .. } => 49, + ErrorCode::UndefinedObject { .. } => 50, + ErrorCode::ObjectNotInPrerequisiteState { .. } => 51, } } @@ -193,6 +195,13 @@ pub(super) fn samples() -> Vec { ErrorCode::DatatypeMismatch { detail: text() }, ErrorCode::InvalidDatetimeFormat { detail: text() }, ErrorCode::DatetimeFieldOverflow { detail: text() }, + ErrorCode::UndefinedObject { + object: "document \"d\"".into(), + }, + ErrorCode::ObjectNotInPrerequisiteState { + object: "CRDT version".into(), + detail: text(), + }, ErrorCode::DispatchCapacity { reason: text() }, ErrorCode::ExpiredBeforeExecution, ErrorCode::BadRequest { detail: text() }, diff --git a/nodedb/src/control/gateway/error_map/class_parity/error_parity.rs b/nodedb/src/control/gateway/error_map/class_parity/error_parity.rs index 32470284d..a1b21c6b0 100644 --- a/nodedb/src/control/gateway/error_map/class_parity/error_parity.rs +++ b/nodedb/src/control/gateway/error_map/class_parity/error_parity.rs @@ -89,6 +89,9 @@ fn classified_sqlstates() -> Vec<(usize, &'static str)> { (44, sqlstate::QUOTA_OVERCOMMIT), (62, sqlstate::SYNTAX_ERROR), (63, sqlstate::SYNTAX_ERROR), + // A CRDT constraint violation renders its constraint's class, never + // the internal-error default. + (74, sqlstate::UNIQUE_VIOLATION), (87, sqlstate::SYNTAX_ERROR), (88, sqlstate::DEPENDENT_OBJECTS_STILL_EXIST), (90, sqlstate::ACTIVE_SQL_TRANSACTION), diff --git a/nodedb/src/control/server/dispatch_utils/write_abort.rs b/nodedb/src/control/server/dispatch_utils/write_abort.rs index ac9cdd029..0ed1cdba0 100644 --- a/nodedb/src/control/server/dispatch_utils/write_abort.rs +++ b/nodedb/src/control/server/dispatch_utils/write_abort.rs @@ -96,6 +96,8 @@ fn is_transient_verdict(code: &ErrorCode) -> bool { | ErrorCode::DatatypeMismatch { .. } | ErrorCode::InvalidDatetimeFormat { .. } | ErrorCode::DatetimeFieldOverflow { .. } + | ErrorCode::UndefinedObject { .. } + | ErrorCode::ObjectNotInPrerequisiteState { .. } | ErrorCode::BadRequest { .. } => false, } } @@ -174,7 +176,11 @@ pub(crate) fn write_definitely_not_applied(code: &ErrorCode) -> bool { | ErrorCode::DatetimeFieldOverflow { .. } | ErrorCode::UndefinedColumn { .. } // A full-text read refused the field before ranking. - | ErrorCode::TextColumn { .. } => true, + | ErrorCode::TextColumn { .. } + // The named object is absent, or not in the state the request needs. + // Both are checked before any mutation. + | ErrorCode::UndefinedObject { .. } + | ErrorCode::ObjectNotInPrerequisiteState { .. } => true, // NOT established — every one of these can be reported by a request // whose write reached, or can have reached, engine state. Emitting an diff --git a/nodedb/src/control/server/pgwire/types/error_map.rs b/nodedb/src/control/server/pgwire/types/error_map.rs index 97cfa04b3..c393d9e39 100644 --- a/nodedb/src/control/server/pgwire/types/error_map.rs +++ b/nodedb/src/control/server/pgwire/types/error_map.rs @@ -323,6 +323,13 @@ pub fn error_to_sqlstate(err: &crate::Error) -> (&'static str, &'static str, Str crate::Error::DataPlane(code) => { crate::control::server::shared::ddl::sqlstate::error_code_to_sqlstate(code) } + // A CRDT error renders the SQLSTATE of its Data-Plane verdict. Only + // a server fault renders `XX000`. + crate::Error::Crdt(crdt) => { + crate::control::server::shared::ddl::sqlstate::error_code_to_sqlstate( + &crate::bridge::envelope::ErrorCode::from(crdt), + ) + } crate::Error::Shaping(e) => ( "ERROR", numeric_code_to_sqlstate(e.code()), @@ -419,7 +426,6 @@ pub fn error_to_sqlstate(err: &crate::Error) -> (&'static str, &'static str, Str | crate::Error::Serialization { .. } | crate::Error::Codec { .. } | crate::Error::SegmentCorrupted { .. } - | crate::Error::Crdt(_) | crate::Error::Io(_) | crate::Error::Config { .. } | crate::Error::Encryption { .. } diff --git a/nodedb/src/control/server/shared/ddl/sqlstate.rs b/nodedb/src/control/server/shared/ddl/sqlstate.rs index d8499dde2..414c22ad9 100644 --- a/nodedb/src/control/server/shared/ddl/sqlstate.rs +++ b/nodedb/src/control/server/shared/ddl/sqlstate.rs @@ -244,6 +244,16 @@ pub fn error_code_to_sqlstate(code: &ErrorCode) -> (&'static str, &'static str, ErrorCode::DatetimeFieldOverflow { detail } => { ("ERROR", sqlstate::DATETIME_FIELD_OVERFLOW, detail.clone()) } + ErrorCode::UndefinedObject { object } => ( + "ERROR", + sqlstate::UNDEFINED_OBJECT, + format!("{object} does not exist"), + ), + ErrorCode::ObjectNotInPrerequisiteState { detail, .. } => ( + "ERROR", + sqlstate::OBJECT_NOT_IN_PREREQUISITE_STATE, + detail.clone(), + ), // The same SQLSTATE the Control Plane gives `crate::Error::BadRequest`. ErrorCode::BadRequest { detail } => ("ERROR", sqlstate::SYNTAX_ERROR, detail.clone()), ErrorCode::TransactionRollback { detail } => { diff --git a/nodedb/src/control/server/sync/async_dispatch/delta/compensation.rs b/nodedb/src/control/server/sync/async_dispatch/delta/compensation.rs index 99cfbf7b9..7fc11a1b9 100644 --- a/nodedb/src/control/server/sync/async_dispatch/delta/compensation.rs +++ b/nodedb/src/control/server/sync/async_dispatch/delta/compensation.rs @@ -228,6 +228,8 @@ fn compensation_hint_for_code(code: &ErrorCode) -> CompensationHint { | ErrorCode::DatatypeMismatch { .. } | ErrorCode::InvalidDatetimeFormat { .. } | ErrorCode::DatetimeFieldOverflow { .. } + | ErrorCode::UndefinedObject { .. } + | ErrorCode::ObjectNotInPrerequisiteState { .. } | ErrorCode::BadRequest { .. } | ErrorCode::ActiveSqlTransaction { .. } | ErrorCode::DependentObjectsExist { .. } diff --git a/nodedb/src/control/server/sync/refusal.rs b/nodedb/src/control/server/sync/refusal.rs index d50436b29..ecf8e3d96 100644 --- a/nodedb/src/control/server/sync/refusal.rs +++ b/nodedb/src/control/server/sync/refusal.rs @@ -67,6 +67,8 @@ pub(super) fn retryable_refusal_reason(error: &crate::Error) -> Option<&str> { | ErrorCode::DatatypeMismatch { .. } | ErrorCode::InvalidDatetimeFormat { .. } | ErrorCode::DatetimeFieldOverflow { .. } + | ErrorCode::UndefinedObject { .. } + | ErrorCode::ObjectNotInPrerequisiteState { .. } | ErrorCode::DispatchCapacity { .. } | ErrorCode::ExpiredBeforeExecution | ErrorCode::BadRequest { .. } @@ -393,6 +395,8 @@ fn is_indeterminate_code(code: &ErrorCode) -> bool { | ErrorCode::DatatypeMismatch { .. } | ErrorCode::InvalidDatetimeFormat { .. } | ErrorCode::DatetimeFieldOverflow { .. } + | ErrorCode::UndefinedObject { .. } + | ErrorCode::ObjectNotInPrerequisiteState { .. } | ErrorCode::BadRequest { .. } | ErrorCode::ActiveSqlTransaction { .. } | ErrorCode::DependentObjectsExist { .. } diff --git a/nodedb/src/data/executor/handlers/control/crdt_doc.rs b/nodedb/src/data/executor/handlers/control/crdt_doc.rs index df8a3cbec..be93cfe9c 100644 --- a/nodedb/src/data/executor/handlers/control/crdt_doc.rs +++ b/nodedb/src/data/executor/handlers/control/crdt_doc.rs @@ -116,7 +116,9 @@ impl CoreLoop { } Ok(Self::encode_crdt_row(engine, collection, document_id)) }); - // A Loro write that fails part-way leaves some fields written. + // A Loro write that fails part-way commits the fields it wrote before + // the error. Loro cannot roll them back, so the undo writes the + // captured row image over them. let bytes = match mutated { Ok(Some(bytes)) => bytes, Ok(None) => { diff --git a/nodedb/src/data/executor/handlers/graph_edge_write/put.rs b/nodedb/src/data/executor/handlers/graph_edge_write/put.rs index acc04cf07..9fadb68bc 100644 --- a/nodedb/src/data/executor/handlers/graph_edge_write/put.rs +++ b/nodedb/src/data/executor/handlers/graph_edge_write/put.rs @@ -290,7 +290,10 @@ mod tests { nodedb_graph::GraphError::RebuildInProgress, ); - assert!(matches!(code, ErrorCode::BadRequest { .. }), "{code:?}"); + assert!( + matches!(code, ErrorCode::ObjectNotInPrerequisiteState { .. }), + "{code:?}" + ); let stored = h .core .edge_store diff --git a/nodedb/src/data/executor/handlers/point/apply_put/types.rs b/nodedb/src/data/executor/handlers/point/apply_put/types.rs index 866edc4ab..f9082bb4d 100644 --- a/nodedb/src/data/executor/handlers/point/apply_put/types.rs +++ b/nodedb/src/data/executor/handlers/point/apply_put/types.rs @@ -181,6 +181,8 @@ pub(in crate::data::executor) fn map_enforcement_error(e: ErrorCode) -> crate::E | ErrorCode::DatatypeMismatch { .. } | ErrorCode::InvalidDatetimeFormat { .. } | ErrorCode::DatetimeFieldOverflow { .. } + | ErrorCode::UndefinedObject { .. } + | ErrorCode::ObjectNotInPrerequisiteState { .. } | ErrorCode::DispatchCapacity { .. } | ErrorCode::ExpiredBeforeExecution | ErrorCode::BadRequest { .. } diff --git a/nodedb/src/engine/crdt/tenant_state/history.rs b/nodedb/src/engine/crdt/tenant_state/history.rs index b56f5be68..4b41d3da8 100644 --- a/nodedb/src/engine/crdt/tenant_state/history.rs +++ b/nodedb/src/engine/crdt/tenant_state/history.rs @@ -89,9 +89,10 @@ impl TenantCrdtEngine { ) -> crate::Result> { let vv = json_to_vv(target_version_json)?; let state = self.collections.get(collection).ok_or_else(|| { - crate::Error::Crdt(nodedb_crdt::CrdtError::Loro( - "document did not exist at target version".into(), - )) + crate::Error::Crdt(nodedb_crdt::CrdtError::RowAbsentAtVersion { + collection: collection.to_string(), + row_id: document_id.to_string(), + }) })?; state .preview_restore_to_version(collection, document_id, &vv) @@ -107,9 +108,10 @@ impl TenantCrdtEngine { ) -> crate::Result> { let vv = json_to_vv(target_version_json)?; let state = self.collections.get(collection).ok_or_else(|| { - crate::Error::Crdt(nodedb_crdt::CrdtError::Loro( - "document did not exist at target version".into(), - )) + crate::Error::Crdt(nodedb_crdt::CrdtError::RowAbsentAtVersion { + collection: collection.to_string(), + row_id: document_id.to_string(), + }) })?; state .restore_to_version(collection, document_id, &vv) diff --git a/nodedb/src/engine/crdt/tenant_state/snapshot_io.rs b/nodedb/src/engine/crdt/tenant_state/snapshot_io.rs index 6fae172c1..cb50262e4 100644 --- a/nodedb/src/engine/crdt/tenant_state/snapshot_io.rs +++ b/nodedb/src/engine/crdt/tenant_state/snapshot_io.rs @@ -71,8 +71,11 @@ impl TenantCrdtEngine { super::ValidatedApplyOutcome::DeadLetterRefused { error, .. } => { Err(crate::Error::Crdt(error)) } + // The caller's bytes do not decode as a snapshot. super::ValidatedApplyOutcome::Malformed => Err(crate::Error::Crdt( - nodedb_crdt::CrdtError::DeltaApplyFailed("malformed snapshot".into()), + nodedb_crdt::CrdtError::ImportMalformed { + detail: "the bytes do not decode as a snapshot".into(), + }, )), super::ValidatedApplyOutcome::PendingDependencies => Err(crate::Error::Crdt( nodedb_crdt::CrdtError::DeltaApplyFailed( diff --git a/nodedb/src/error_classify/public.rs b/nodedb/src/error_classify/public.rs index 9db05ecda..a3b0753e5 100644 --- a/nodedb/src/error_classify/public.rs +++ b/nodedb/src/error_classify/public.rs @@ -260,7 +260,11 @@ pub(crate) fn classify(e: &Error) -> NodeDbError { Error::SegmentCorrupted { detail } => NodeDbError::segment_corrupted(detail), Error::MemoryExhausted { engine } => NodeDbError::memory_exhausted(engine.clone()), Error::Backpressure { engine } => NodeDbError::memory_exhausted(engine.to_string()), - Error::Crdt(crdt_err) => NodeDbError::internal(crdt_err), + // A CRDT error takes the public code of its Data-Plane verdict. Only + // a server fault is `internal`. + Error::Crdt(crdt_err) => crate::error_from_data_plane::data_plane_code_to_public( + crate::bridge::envelope::ErrorCode::from(crdt_err), + ), Error::Io(io_err) => NodeDbError::storage(io_err), Error::Config { detail } => NodeDbError::config(detail), Error::Encryption { detail } => NodeDbError::encryption(detail), diff --git a/nodedb/src/error_classify/unclassified.rs b/nodedb/src/error_classify/unclassified.rs index 26364b1f6..65ecd6bfd 100644 --- a/nodedb/src/error_classify/unclassified.rs +++ b/nodedb/src/error_classify/unclassified.rs @@ -16,7 +16,6 @@ pub(crate) fn is_unclassified_failure(e: &Error) -> bool { | Error::Serialization { .. } | Error::Codec { .. } | Error::SegmentCorrupted { .. } - | Error::Crdt(_) | Error::Io(_) | Error::Config { .. } | Error::Encryption { .. } @@ -31,6 +30,12 @@ pub(crate) fn is_unclassified_failure(e: &Error) -> bool { | Error::CollectionUnstamped { .. } | Error::MaterializedSumResolutionMissing { .. } | Error::CascadeCycle { .. } => true, + // A CRDT error is classified unless its Data-Plane verdict is a + // server fault. + Error::Crdt(crdt) => matches!( + crate::bridge::envelope::ErrorCode::from(crdt), + crate::bridge::envelope::ErrorCode::Internal { .. } + ), // A client-matchable class or a retry contract a caller matches by // variant. Error::RejectedConstraint { .. } @@ -157,4 +162,18 @@ mod tests { detail: "apply error".to_owned(), })); } + + /// A client-caused CRDT error is a verdict. A Loro fault is machinery. + #[test] + fn a_crdt_error_is_unclassified_only_when_it_is_a_server_fault() { + assert!(!is_unclassified_failure(&Error::Crdt( + nodedb_crdt::CrdtError::RowAbsentAtVersion { + collection: "notes".to_owned(), + row_id: "doc".to_owned(), + } + ))); + assert!(is_unclassified_failure(&Error::Crdt( + nodedb_crdt::CrdtError::Loro("encode".to_owned()) + ))); + } } diff --git a/nodedb/src/error_from.rs b/nodedb/src/error_from.rs index c139fe658..fb0d65b0e 100644 --- a/nodedb/src/error_from.rs +++ b/nodedb/src/error_from.rs @@ -318,6 +318,11 @@ impl From for nodedb_cluster::rpc_codec::TypedClusterError { // Keep the verdict typed across a further hop instead of // degrading it to a numeric class on the second forward. Error::DataPlane(code) => TypedClusterError::DataPlane { code: code.into() }, + // A CRDT error crosses as its Data-Plane verdict, so its class + // survives the hop instead of reading as a server fault. + Error::Crdt(crdt) => TypedClusterError::DataPlane { + code: crate::bridge::envelope::ErrorCode::from(&crdt).into(), + }, // Keep the constraint kind typed across a further hop, same as // a Data-Plane verdict, instead of flattening it to one code. Error::RejectedConstraint { @@ -438,7 +443,6 @@ impl From for nodedb_cluster::rpc_codec::TypedClusterError { | Error::SegmentCorrupted { .. } | Error::MemoryExhausted { .. } | Error::Backpressure { .. } - | Error::Crdt(_) | Error::Io(_) | Error::Config { .. } | Error::Encryption { .. } diff --git a/nodedb/src/error_from_data_plane.rs b/nodedb/src/error_from_data_plane.rs index 308f9d440..e8cd847ed 100644 --- a/nodedb/src/error_from_data_plane.rs +++ b/nodedb/src/error_from_data_plane.rs @@ -174,6 +174,10 @@ pub(crate) fn data_plane_code_to_public(code: ErrorCode) -> NodeDbError { ErrorCode::DependentObjectsExist { object, detail } => { NodeDbError::dependent_objects_exist(object, detail) } + ErrorCode::UndefinedObject { object } => NodeDbError::undefined_object(object), + ErrorCode::ObjectNotInPrerequisiteState { object, detail } => { + NodeDbError::object_not_ready(object, detail) + } // Nothing was enqueued, and the same request succeeds once capacity // frees: the retryable overload class. ErrorCode::DispatchCapacity { reason } => NodeDbError::server_overload(reason), diff --git a/nodedb/tests/wire/cases/crdt_restore_non_default_database.rs b/nodedb/tests/wire/cases/crdt_restore_non_default_database.rs new file mode 100644 index 000000000..c6a6f8532 --- /dev/null +++ b/nodedb/tests/wire/cases/crdt_restore_non_default_database.rs @@ -0,0 +1,88 @@ +// SPDX-License-Identifier: BUSL-1.1 + +//! `RESTORE ... SET VERSION` on a CRDT collection in a non-default database. +//! +//! The restore must address the same CRDT document the ordinary write path +//! maintains, which is keyed by the database-qualified collection name. A +//! same-name collection in `default` is a different document and stays +//! untouched. + +use crate::harness::TestServer; + +const COLLECTION: &str = "crdt_restore_notes"; + +async fn exec_ok(server: &TestServer, sql: &str) { + server + .exec(sql) + .await + .unwrap_or_else(|e| panic!("query failed: {e}\nsql: {sql}")); +} + +async fn title_of(server: &TestServer, id: &str) -> String { + let rows = server + .query_rows(&format!("SELECT title FROM {COLLECTION} WHERE id = '{id}'")) + .await + .unwrap_or_else(|e| panic!("read {COLLECTION}/{id}: {e}")); + rows.into_iter() + .next() + .and_then(|row| row.into_iter().next()) + .unwrap_or_else(|| panic!("{COLLECTION}/{id} not found")) +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn restore_set_version_in_non_default_database_restores_checkpointed_state() { + let server = TestServer::start().await; + let create = + format!("CREATE TABLE {COLLECTION} (id TEXT PRIMARY KEY, title TEXT) WITH (crdt='true')"); + + // Default database: same collection name, own row, never restored. + exec_ok(&server, &create).await; + exec_ok( + &server, + &format!("INSERT INTO {COLLECTION} (id, title) VALUES ('doc', 'default-v1')"), + ) + .await; + exec_ok( + &server, + &format!("UPDATE {COLLECTION} SET title = 'default-v2' WHERE id = 'doc'"), + ) + .await; + + exec_ok(&server, "CREATE DATABASE d2").await; + exec_ok(&server, "USE DATABASE d2").await; + exec_ok(&server, &create).await; + exec_ok( + &server, + &format!("INSERT INTO {COLLECTION} (id, title) VALUES ('doc', 'd2-v1')"), + ) + .await; + exec_ok( + &server, + &format!("CREATE CHECKPOINT 'cp1' ON {COLLECTION} WHERE id = 'doc'"), + ) + .await; + exec_ok( + &server, + &format!("UPDATE {COLLECTION} SET title = 'd2-v2' WHERE id = 'doc'"), + ) + .await; + assert_eq!(title_of(&server, "doc").await, "d2-v2"); + + exec_ok( + &server, + &format!("RESTORE {COLLECTION} SET VERSION = 'cp1' WHERE id = 'doc'"), + ) + .await; + assert_eq!( + title_of(&server, "doc").await, + "d2-v1", + "RESTORE in d2 must show the checkpointed state through the normal read path" + ); + + exec_ok(&server, "USE DATABASE default").await; + assert_eq!( + title_of(&server, "doc").await, + "default-v2", + "RESTORE in d2 must leave the same-name default collection untouched" + ); +} From 3eaf2645b0f5e2fb919972cdb13783771bb22ad5 Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Fri, 9 Oct 2026 04:30:20 +0800 Subject: [PATCH 4/5] test(wire): cover SERIAL in the column list and register new cases --- nodedb/tests/wire/cases/mod.rs | 3 + .../wire/cases/sql_serial_column_list.rs | 90 +++++++++++++++++++ 2 files changed, 93 insertions(+) create mode 100644 nodedb/tests/wire/cases/sql_serial_column_list.rs diff --git a/nodedb/tests/wire/cases/mod.rs b/nodedb/tests/wire/cases/mod.rs index ff24408a8..5696b97e4 100644 --- a/nodedb/tests/wire/cases/mod.rs +++ b/nodedb/tests/wire/cases/mod.rs @@ -50,6 +50,7 @@ mod columnar_read_row_level_security; mod columnar_write_row_level_security; mod command_complete_tag_conformance; mod continuous_aggregate_restart; +mod crdt_restore_non_default_database; mod crdt_write_rls_database_scope; mod database_boundary; mod database_metrics; @@ -271,6 +272,7 @@ mod sql_security_e2e; mod sql_sequence_row_scope; mod sql_sequence_row_scope_refusals; mod sql_sequences; +mod sql_serial_column_list; mod sql_spatial_index_ddl; mod sql_subquery_from; mod sql_subquery_text_spatial_composition; @@ -331,6 +333,7 @@ mod sql_update_expressions; mod sql_update_from; mod sql_update_primary_key_identity; mod sql_utf8_expressions; +mod sql_vector_array_literal; mod sql_vector_index_ddl; mod sql_where_expressions; mod sql_where_instant_literals; diff --git a/nodedb/tests/wire/cases/sql_serial_column_list.rs b/nodedb/tests/wire/cases/sql_serial_column_list.rs new file mode 100644 index 000000000..5faee18cc --- /dev/null +++ b/nodedb/tests/wire/cases/sql_serial_column_list.rs @@ -0,0 +1,90 @@ +// SPDX-License-Identifier: BUSL-1.1 + +//! `SERIAL` and `BIGSERIAL` in the parenthesised column list. +//! +//! `CREATE COLLECTION c (n SERIAL, v TEXT)` and +//! `CREATE COLLECTION c FIELDS (n SERIAL, v TEXT)` expand the type +//! identically. An insert that omits the column takes 1, 2, 3 in order. + +use crate::harness::TestServer; + +/// Insert three rows that omit `n`, then assert `n` reads back 1, 2, 3. +async fn assert_serial_allocates_in_order(server: &TestServer, collection: &str) { + for value in ["a", "b", "c"] { + server + .exec(&format!("INSERT INTO {collection} (v) VALUES ('{value}')")) + .await + .unwrap_or_else(|e| panic!("insert into {collection}: {e}")); + } + let rows = server + .query_text(&format!("SELECT n FROM {collection} ORDER BY n")) + .await + .unwrap(); + let keys: Vec<&str> = rows.iter().map(|row| row.trim()).collect(); + assert_eq!( + keys, + vec!["1", "2", "3"], + "{collection}: n must allocate 1, 2, 3 in order" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn serial_in_column_list_schemaless() { + let server = TestServer::start().await; + server + .exec("CREATE COLLECTION serial_cl_schemaless (n SERIAL, v TEXT)") + .await + .unwrap(); + assert_serial_allocates_in_order(&server, "serial_cl_schemaless").await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn serial_in_column_list_document_strict() { + let server = TestServer::start().await; + server + .exec( + "CREATE COLLECTION serial_cl_strict (n SERIAL PRIMARY KEY, v TEXT) \ + WITH (engine='document_strict')", + ) + .await + .unwrap(); + assert_serial_allocates_in_order(&server, "serial_cl_strict").await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn bigserial_in_column_list_schemaless() { + let server = TestServer::start().await; + server + .exec("CREATE COLLECTION bigserial_cl_schemaless (n BIGSERIAL, v TEXT)") + .await + .unwrap(); + assert_serial_allocates_in_order(&server, "bigserial_cl_schemaless").await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn bigserial_in_column_list_document_strict() { + let server = TestServer::start().await; + server + .exec( + "CREATE COLLECTION bigserial_cl_strict (n BIGSERIAL PRIMARY KEY, v TEXT) \ + WITH (engine='document_strict')", + ) + .await + .unwrap(); + assert_serial_allocates_in_order(&server, "bigserial_cl_strict").await; +} + +/// The column-list spelling creates the same implicit sequence as `FIELDS`. +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn serial_in_column_list_creates_implicit_sequence() { + let server = TestServer::start().await; + server + .exec("CREATE COLLECTION serial_cl_seq (n SERIAL, v TEXT)") + .await + .unwrap(); + let rows = server.query_text("SHOW SEQUENCES").await.unwrap(); + assert!( + !rows.is_empty(), + "SERIAL in the column list must create an implicit sequence" + ); +} From 66bcdb2734aa20fc8bb0b35be710ed86c36b16f2 Mon Sep 17 00:00:00 2001 From: Farhan Syah Date: Fri, 9 Oct 2026 04:37:06 +0800 Subject: [PATCH 5/5] fix(deps): bump rustls to 0.23.45 in fuzz lockfile rustls before 0.23.45 accepts TLS 1.3 handshake messages across encryption level boundaries (GHSA-2mjx-qc3c-rqvc). Update the fuzz lockfile to the patched release, along with rustls-webpki and the other transitive updates the resolver pulled in. --- fuzz/Cargo.lock | 26 ++++++++++++++++++-------- 1 file changed, 18 insertions(+), 8 deletions(-) diff --git a/fuzz/Cargo.lock b/fuzz/Cargo.lock index 7112da464..287774457 100644 --- a/fuzz/Cargo.lock +++ b/fuzz/Cargo.lock @@ -337,9 +337,9 @@ checksum = "f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53" [[package]] name = "aws-lc-rs" -version = "1.17.3" +version = "1.18.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "00bdb5da18dac48ca2cc7cd4a98e533e8635a58e2361d13a1a4ee3888e0d72f1" +checksum = "b281d307588d634de920874890732659e2e7672f72b5e10e81badc1a8a83621e" dependencies = [ "aws-lc-sys", "zeroize", @@ -347,9 +347,9 @@ dependencies = [ [[package]] name = "aws-lc-sys" -version = "0.43.0" +version = "0.45.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "43103168cc76fe62678a375e722fc9cb3a0146159ac5828bc4f0dfd755c2224c" +checksum = "9bff6c3b54fad79a2e60b8102caf565819711497c1f5f092f49508e2f5c31b27" dependencies = [ "cc", "cmake", @@ -1217,6 +1217,12 @@ version = "0.5.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "fc0fef456e4baa96da950455cd02c081ca953b141298e41db3fc7e36b1da849c" +[[package]] +name = "hex" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "7f24254aa9a54b5c858eaee2f5bccdb46aaf0e486a595ed5fd8f86ba55232a70" + [[package]] name = "hkdf" version = "0.12.4" @@ -2070,6 +2076,7 @@ dependencies = [ "serde_json", "sonic-rs", "thiserror", + "ulid", "uuid", "zerompk", ] @@ -2125,6 +2132,7 @@ dependencies = [ name = "nodedb-query" version = "0.5.0" dependencies = [ + "hex", "nodedb-fts", "nodedb-spatial", "nodedb-types", @@ -2159,6 +2167,7 @@ name = "nodedb-sql" version = "0.5.0" dependencies = [ "chrono", + "hex", "nodedb-query", "nodedb-spatial", "nodedb-types", @@ -2192,6 +2201,7 @@ dependencies = [ "bytemuck", "crc32c", "getrandom 0.4.3", + "hex", "nanoid", "nodedb-codec", "rand 0.10.2", @@ -2912,9 +2922,9 @@ dependencies = [ [[package]] name = "rustls" -version = "0.23.42" +version = "0.23.45" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "3c54fcab019b409d04215d3a17cb438fd7fbf192ee61461f20f4fe18704bc138" +checksum = "0d41d731c7d2f962d1ccc364cec258de3c0e93b38c2fb3ba97ac74513048d634" dependencies = [ "aws-lc-rs", "once_cell", @@ -2975,9 +2985,9 @@ checksum = "f87165f0995f63a9fbeea62b64d10b4d9d8e78ec6d7d51fb2125fda7bb36788f" [[package]] name = "rustls-webpki" -version = "0.103.13" +version = "0.103.15" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "61c429a8649f110dddef65e2a5ad240f747e85f7758a6bccc7e5777bd33f756e" +checksum = "f3c3cf1d8b1e7d4927e2d154c3fcb02979afb9939629c62cd9048d4f07b60ac2" dependencies = [ "aws-lc-rs", "ring",