From 73686b8987e9ffc52554438d2f6df5d632bd41ee Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Thu, 1 Oct 2026 18:26:41 +0000 Subject: [PATCH 1/6] fix(postgres): cancel deadlocked write transactions instead of failing them Two TransactWriteItems that name the same items in different orders could deadlock in PostgreSQL. The victim came back as HTTP 500, and each deadlock held its row locks for deadlock_timeout first. Run the ops in table and key order, so every write transaction takes its row locks in one order and two of them cannot deadlock on each other. The results keep their request positions. When PostgreSQL still aborts the transaction with deadlock_detected (40P01) or serialization_failure (40001), cancel it with a TransactionConflict reason on the item that hit the abort, the shape Amazon DynamoDB returns. An abort outside the per-item work names every item. To read the SQLSTATE, the transaction helpers now map database errors through db_error. Assisted-by: pi claude-opus-5-5 --- crates/core/src/types/transaction.rs | 10 + crates/storage-postgres/src/data/mod.rs | 4 +- .../storage-postgres/src/data/transactions.rs | 233 +++++++++++-- .../storage-postgres/src/data/tx_helpers.rs | 19 +- crates/storage-postgres/src/pg_util.rs | 87 +++++ crates/storage-postgres/tests/twi_conflict.rs | 327 ++++++++++++++++++ tests/test_transact_write_conflict.py | 208 +++++++++++ 7 files changed, 843 insertions(+), 45 deletions(-) create mode 100644 crates/storage-postgres/tests/twi_conflict.rs create mode 100644 tests/test_transact_write_conflict.py diff --git a/crates/core/src/types/transaction.rs b/crates/core/src/types/transaction.rs index d15e608ef..5010cbb47 100755 --- a/crates/core/src/types/transaction.rs +++ b/crates/core/src/types/transaction.rs @@ -296,6 +296,16 @@ impl CancellationReason { } } + /// Create a reason for an item that another transaction holds. + #[must_use] + pub fn transaction_conflict() -> Self { + Self { + code: "TransactionConflict".to_owned(), + message: Some("Transaction is ongoing for the item".to_owned()), + item: None, + } + } + /// Create a reason for a validation error. #[must_use] pub fn validation_error(msg: impl Into) -> Self { diff --git a/crates/storage-postgres/src/data/mod.rs b/crates/storage-postgres/src/data/mod.rs index d64ceb539..e16ee3b65 100755 --- a/crates/storage-postgres/src/data/mod.rs +++ b/crates/storage-postgres/src/data/mod.rs @@ -92,7 +92,7 @@ macro_rules! bind_sk_fetch_optional { .await } } - .map_err(|e| extenddb_storage::error::StorageError::Internal(e.to_string())) + .map_err($crate::data::index::db_error) }; } @@ -124,7 +124,7 @@ macro_rules! bind_sk_execute { .await } } - .map_err(|e| extenddb_storage::error::StorageError::Internal(e.to_string())) + .map_err($crate::data::index::db_error) }; } diff --git a/crates/storage-postgres/src/data/transactions.rs b/crates/storage-postgres/src/data/transactions.rs index d69f90043..7c633d0f8 100644 --- a/crates/storage-postgres/src/data/transactions.rs +++ b/crates/storage-postgres/src/data/transactions.rs @@ -3,6 +3,7 @@ //! Transactional read/write implementations for the `PostgreSQL` backend. +use std::borrow::Cow; use std::collections::HashMap; use extenddb_core::expression::{self, ExpressionMaps}; @@ -11,14 +12,18 @@ use extenddb_core::types::{ }; use extenddb_core::validation; use extenddb_storage::error::StorageError; +use extenddb_storage::util::pk_to_text; use extenddb_storage::{IdempotencyKey, TransactGetOp, TransactWriteOp}; -use super::index::{IndexMeta, enqueue_async_indexes, fetch_write_path_indexes, sync_indexes}; +use super::index::{ + IndexMeta, db_error, enqueue_async_indexes, fetch_write_path_indexes, sync_indexes, +}; use super::tx_helpers::{ check_idempotency_token_in_tx, delete_item_in_tx, fetch_item_for_update, fetch_item_in_tx, insert_item_if_absent_in_tx, upsert_item_in_tx, write_stream_record_in_tx, }; use crate::PostgresEngine; +use crate::pg_util::is_conflict_abort; /// Bound on insert retries when a transactional write to a nonexistent item /// keeps losing the create race to writers that then roll back. Mirrors the @@ -109,19 +114,29 @@ impl PostgresEngine { .await .map_err(|e| StorageError::Internal(e.to_string()))?; + // A conflict abort outside the per-op loop cannot be tied to one item. + let cancel_all = |e: StorageError| conflict_cancels_all(e, ops.len()); + // Check the idempotency token within the transaction so token storage // and data writes commit together. The token is scoped to its account. if let Some(key) = idempotency { check_idempotency_token_in_tx(&mut tx, key.account_id, key.token, key.fingerprint) - .await?; + .await + .map_err(cancel_all)?; } - let mut reasons: Vec = Vec::with_capacity(ops.len()); + let mut reasons: Vec = vec![CancellationReason::none(); ops.len()]; // M-3: Collect old/new items from each op for async GSI enqueue after commit. - let mut op_items: Vec<(Option, Option)> = Vec::with_capacity(ops.len()); + let mut op_items: Vec<(Option, Option)> = vec![(None, None); ops.len()]; let mut any_failed = false; + let mut first_invalid: Option<(usize, String)> = None; + let mut failed_op: Option<(usize, StorageError)> = None; - for op in ops { + // Run the ops in key order, not request order, so every transaction + // locks its items in one global order and two cannot deadlock on each + // other. Results keep their request positions. + for i in execution_order(ops) { + let op = &ops[i]; let indexes = &table_indexes[transact_op_table_name(op)]; let reason = execute_transact_write_op( &mut tx, @@ -132,29 +147,43 @@ impl PostgresEngine { ) .await; match reason { - Ok(items) => { - op_items.push(items); - reasons.push(CancellationReason::none()); - } + Ok(items) => op_items[i] = items, Err(TxnOpError::Cancel(r)) => { - op_items.push((None, None)); any_failed = true; - reasons.push(r); + reasons[i] = r; } Err(TxnOpError::Validation(msg)) => { - // Up-front input validation (e.g. empty secondary-index key): - // abort the whole transaction with a top-level - // ValidationException, not a per-item cancellation reason. - return Err(StorageError::Validation(msg)); + // Up-front input validation (e.g. empty secondary-index key) + // fails the request with a top-level ValidationException. + // Report the earliest invalid op in request order. + if first_invalid.as_ref().is_none_or(|(j, _)| i < *j) { + first_invalid = Some((i, msg)); + } } Err(TxnOpError::Storage(e)) => { - // Infrastructure error — abort the entire transaction - // without leaking internal details into cancellation reasons. - return Err(StorageError::Internal(e.to_string())); + // Infrastructure error, or PostgreSQL aborted the + // transaction: run no further ops. + failed_op = Some((i, e)); + break; } } } + if let Some((_, msg)) = first_invalid { + return Err(StorageError::Validation(msg)); + } + if let Some((i, e)) = failed_op { + if is_conflict_abort(&e) { + // PostgreSQL broke a lock conflict by aborting this transaction. + // Cancel it with the contended item named, as the service does. + reasons[i] = CancellationReason::transaction_conflict(); + return Err(StorageError::TransactionCanceled(reasons)); + } + // Infrastructure error: abort without leaking internal details + // into cancellation reasons. + return Err(StorageError::Internal(e.to_string())); + } + if any_failed { return Err(StorageError::TransactionCanceled(reasons)); } @@ -180,7 +209,8 @@ impl PostgresEngine { old_item.as_ref(), new_item.as_ref(), ) - .await?; + .await + .map_err(cancel_all)?; } } @@ -207,7 +237,8 @@ impl PostgresEngine { new_item.as_ref(), sys_delay, ) - .await?; + .await + .map_err(cancel_all)?; // Vector maintenance for all three write kinds in one place, rather // than in each branch above: this loop already visits exactly the ops @@ -224,15 +255,14 @@ impl PostgresEngine { new_item.as_ref(), sys_delay, ) - .await?; + .await + .map_err(cancel_all)?; if n > 0 || vector_n > 0 { needs_notify = true; } } - tx.commit() - .await - .map_err(|e| StorageError::Internal(e.to_string()))?; + tx.commit().await.map_err(db_error).map_err(cancel_all)?; if needs_notify && let Some(ref q) = self.gsi_queue { q.notify_workers(); @@ -260,6 +290,43 @@ impl PostgresEngine { } } +/// Map a conflict abort to a cancellation that names every item. Other +/// errors pass through unchanged. +fn conflict_cancels_all(e: StorageError, n_ops: usize) -> StorageError { + if is_conflict_abort(&e) { + StorageError::TransactionCanceled(vec![CancellationReason::transaction_conflict(); n_ops]) + } else { + e + } +} + +/// Request positions of `ops`, sorted by table and primary key. +fn execution_order(ops: &[TransactWriteOp<'_>]) -> Vec { + let mut order: Vec = (0..ops.len()).collect(); + order.sort_by_cached_key(|&i| lock_key(&ops[i])); + order +} + +/// The table and the primary key values of the item an op touches. +fn lock_key<'a>(op: &'a TransactWriteOp<'_>) -> (&'a str, Vec>) { + let (key_info, key) = match op { + TransactWriteOp::Put { key_info, item, .. } => (key_info, *item), + TransactWriteOp::Delete { key_info, key, .. } + | TransactWriteOp::Update { key_info, key, .. } + | TransactWriteOp::ConditionCheck { key_info, key, .. } => (key_info, *key), + }; + let values = key_info + .key_schema + .iter() + .map(|k| { + key.get(&k.attribute_name) + .and_then(|v| pk_to_text(v).ok()) + .unwrap_or_default() + }) + .collect(); + (&key_info.table_id, values) +} + /// Extract the table name from a transactional write operation. fn transact_op_table_name<'a>(op: &'a TransactWriteOp<'_>) -> &'a str { match op { @@ -423,11 +490,9 @@ async fn execute_transact_write_op( // canceled transaction with a TransactionConflict // reason, the DDB-canonical contention shape, never a // 500 (matches the MongoDB backend's exhaustion path). - return Err(TxnOpError::Cancel(CancellationReason { - code: "TransactionConflict".to_owned(), - message: Some("Transaction is ongoing for the item".to_owned()), - item: None, - })); + return Err(TxnOpError::Cancel( + CancellationReason::transaction_conflict(), + )); } } } @@ -612,11 +677,9 @@ async fn execute_transact_write_op( if attempt >= MAX_CREATE_RACE_ATTEMPTS { // Sustained create-then-delete churn. Same // TransactionConflict cancellation as the Put arm. - return Err(TxnOpError::Cancel(CancellationReason { - code: "TransactionConflict".to_owned(), - message: Some("Transaction is ongoing for the item".to_owned()), - item: None, - })); + return Err(TxnOpError::Cancel( + CancellationReason::transaction_conflict(), + )); } } } @@ -693,3 +756,105 @@ fn eval_condition( } Ok(()) } + +#[cfg(test)] +mod tests { + use extenddb_core::types::{KeySchemaElement, KeyType, TableKeyInfo}; + + use super::*; + + fn table(id: &str) -> TableKeyInfo { + TableKeyInfo { + table_id: id.to_owned(), + key_schema: vec![ + KeySchemaElement { + attribute_name: "pk".to_owned(), + key_type: KeyType::Hash, + }, + KeySchemaElement { + attribute_name: "sk".to_owned(), + key_type: KeyType::Range, + }, + ], + ..TableKeyInfo::default() + } + } + + fn key(pk: &str, sk: &str) -> Item { + Item::from([ + ("pk".to_owned(), AttributeValue::S(pk.to_owned())), + ("sk".to_owned(), AttributeValue::S(sk.to_owned())), + ]) + } + + fn delete<'a>( + key_info: &'a TableKeyInfo, + key: &'a Item, + maps: &'a ExpressionMaps, + ) -> TransactWriteOp<'a> { + TransactWriteOp::Delete { + key_info, + key, + condition: None, + maps, + return_values_on_ccf: ReturnValuesOnConditionCheckFailure::None, + stream: None, + } + } + + /// The items each request locks, in the order it locks them. + fn locked(ops: &[TransactWriteOp<'_>]) -> Vec<(String, Vec)> { + execution_order(ops) + .into_iter() + .map(|i| { + let (t, k) = lock_key(&ops[i]); + (t.to_owned(), k.into_iter().map(Cow::into_owned).collect()) + }) + .collect() + } + + #[test] + fn requests_lock_shared_items_in_the_same_order() { + let (t1, t2) = (table("t1"), table("t2")); + let maps = ExpressionMaps::default(); + let (a, b, c) = (key("a", "1"), key("a", "2"), key("b", "1")); + let forward = [ + delete(&t2, &a, &maps), + delete(&t1, &c, &maps), + delete(&t1, &b, &maps), + delete(&t1, &a, &maps), + ]; + let backward = [ + delete(&t1, &a, &maps), + delete(&t1, &b, &maps), + delete(&t1, &c, &maps), + delete(&t2, &a, &maps), + ]; + assert_eq!(locked(&forward), locked(&backward)); + assert_eq!(execution_order(&forward), vec![3, 2, 1, 0]); + assert_eq!(execution_order(&backward), vec![0, 1, 2, 3]); + } + + #[test] + fn a_conflict_abort_outside_the_ops_cancels_every_item() { + let deadlock = StorageError::Internal("SQLSTATE 40P01: deadlock detected".to_owned()); + match conflict_cancels_all(deadlock, 3) { + StorageError::TransactionCanceled(reasons) => { + assert_eq!(reasons.len(), 3); + for r in reasons { + assert_eq!(r.code, "TransactionConflict"); + assert_eq!( + r.message.as_deref(), + Some("Transaction is ongoing for the item") + ); + } + } + other => panic!("expected a cancellation, got {other:?}"), + } + let other = StorageError::Internal("SQLSTATE 23505: duplicate key".to_owned()); + assert!(matches!( + conflict_cancels_all(other, 3), + StorageError::Internal(_) + )); + } +} diff --git a/crates/storage-postgres/src/data/tx_helpers.rs b/crates/storage-postgres/src/data/tx_helpers.rs index ed30a1bf3..11fcc7953 100755 --- a/crates/storage-postgres/src/data/tx_helpers.rs +++ b/crates/storage-postgres/src/data/tx_helpers.rs @@ -14,6 +14,7 @@ use extenddb_storage::StreamCapture; use extenddb_storage::error::StorageError; use extenddb_storage::util::{SortKeyValue, parse_sk, pk_to_text, sk_column, sk_info}; +use super::index::db_error; use super::{data_table_name, json_to_item}; /// Fetch a single item within an existing transaction. @@ -86,7 +87,7 @@ pub(super) async fn fetch_item_for_update( .bind(pk_text.as_ref()) .fetch_optional(&mut **tx) .await - .map_err(|e| StorageError::Internal(e.to_string()))?; + .map_err(db_error)?; row.map(|(v,)| v) }; @@ -130,7 +131,7 @@ pub(super) async fn upsert_item_in_tx( .bind(&item_json) .execute(&mut **tx) .await - .map_err(|e| StorageError::Internal(e.to_string()))?; + .map_err(db_error)?; } Ok(()) } @@ -183,7 +184,7 @@ pub(super) async fn insert_item_if_absent_in_tx( .bind(&item_json) .execute(&mut **tx) .await - .map_err(|e| StorageError::Internal(e.to_string()))? + .map_err(db_error)? .rows_affected() }; Ok(rows_affected == 1) @@ -233,14 +234,14 @@ pub(super) async fn delete_item_in_tx( .await } } - .map_err(|e| StorageError::Internal(e.to_string()))?; + .map_err(db_error)?; } else { let sql = format!("DELETE FROM {ddb_table} WHERE pk = $1"); sqlx::query(&sql) .bind(pk_text.as_ref()) .execute(&mut **tx) .await - .map_err(|e| StorageError::Internal(e.to_string()))?; + .map_err(db_error)?; } Ok(()) } @@ -323,7 +324,7 @@ pub(super) async fn write_stream_record_in_tx( .bind(&key_info.table_id) .fetch_all(&mut **tx) .await - .map_err(|e| StorageError::Internal(e.to_string()))?; + .map_err(db_error)?; if shards.is_empty() { // No shards — streams may not be fully set up yet. Skip silently. @@ -339,7 +340,7 @@ pub(super) async fn write_stream_record_in_tx( let (seq_val,): (i64,) = sqlx::query_as("SELECT nextval('stream_seq')") .fetch_one(&mut **tx) .await - .map_err(|e| StorageError::Internal(e.to_string()))?; + .map_err(db_error)?; let seq = format!("{seq_val:021}"); let record = StreamRecord { @@ -380,7 +381,7 @@ pub(super) async fn write_stream_record_in_tx( .bind(&record_json) .execute(&mut **tx) .await - .map_err(|e| StorageError::Internal(e.to_string()))?; + .map_err(db_error)?; Ok(()) } @@ -421,7 +422,7 @@ pub(super) async fn check_idempotency_token_in_tx( .bind(fingerprint) .fetch_optional(&mut **tx) .await - .map_err(|e| StorageError::Internal(e.to_string()))?; + .map_err(db_error)?; match row { Some((_, true)) | None => Ok(()), diff --git a/crates/storage-postgres/src/pg_util.rs b/crates/storage-postgres/src/pg_util.rs index f87a22590..30e1273bd 100755 --- a/crates/storage-postgres/src/pg_util.rs +++ b/crates/storage-postgres/src/pg_util.rs @@ -18,3 +18,90 @@ pub(crate) fn is_fk_violation(e: &sqlx::Error) -> bool { } false } + +/// Check if a mapped error is PostgreSQL aborting the transaction to break a +/// lock conflict: deadlock_detected (40P01) or serialization_failure (40001). +/// Matches the `SQLSTATE` prefix that `data::index::db_error` writes. +pub(crate) fn is_conflict_abort(e: &extenddb_storage::error::StorageError) -> bool { + match e { + extenddb_storage::error::StorageError::Internal(msg) => { + msg.starts_with("SQLSTATE 40P01:") || msg.starts_with("SQLSTATE 40001:") + } + _ => false, + } +} + +#[cfg(test)] +mod tests { + use std::borrow::Cow; + use std::fmt; + + use extenddb_storage::error::StorageError; + + use super::is_conflict_abort; + use crate::data::index::db_error; + + /// A database error with a chosen SQLSTATE, standing in for the server. + #[derive(Debug)] + struct FakeDbError(&'static str, &'static str); + + impl fmt::Display for FakeDbError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.write_str(self.1) + } + } + + impl std::error::Error for FakeDbError {} + + impl sqlx::error::DatabaseError for FakeDbError { + fn message(&self) -> &str { + self.1 + } + fn code(&self) -> Option> { + Some(Cow::Borrowed(self.0)) + } + fn as_error(&self) -> &(dyn std::error::Error + Send + Sync + 'static) { + self + } + fn as_error_mut(&mut self) -> &mut (dyn std::error::Error + Send + Sync + 'static) { + self + } + fn into_error(self: Box) -> Box { + self + } + fn kind(&self) -> sqlx::error::ErrorKind { + sqlx::error::ErrorKind::Other + } + } + + fn mapped(code: &'static str, msg: &'static str) -> StorageError { + db_error(sqlx::Error::Database(Box::new(FakeDbError(code, msg)))) + } + + #[test] + fn deadlock_and_serialization_failure_are_conflict_aborts() { + assert!(is_conflict_abort(&mapped("40P01", "deadlock detected"))); + assert!(is_conflict_abort(&mapped( + "40001", + "could not serialize access due to concurrent update" + ))); + } + + #[test] + fn other_errors_are_not_conflict_aborts() { + // A unique violation, a lock timeout, and the same text without the + // code all stay internal errors. + assert!(!is_conflict_abort(&mapped("23505", "duplicate key value"))); + assert!(!is_conflict_abort(&mapped( + "55P03", + "could not obtain lock" + ))); + assert!(!is_conflict_abort(&StorageError::Internal( + "deadlock detected".to_owned() + ))); + assert!(!is_conflict_abort(&db_error(sqlx::Error::PoolTimedOut))); + assert!(!is_conflict_abort(&StorageError::TransactionConflict( + "SQLSTATE 40P01: deadlock detected".to_owned() + ))); + } +} diff --git a/crates/storage-postgres/tests/twi_conflict.rs b/crates/storage-postgres/tests/twi_conflict.rs new file mode 100644 index 000000000..fb0b39995 --- /dev/null +++ b/crates/storage-postgres/tests/twi_conflict.rs @@ -0,0 +1,327 @@ +// Copyright 2026 ExtendDB contributors +// SPDX-License-Identifier: Apache-2.0 +//! Storage-level tests for lock conflicts inside `TransactWriteItems`. +//! +//! Amazon DynamoDB cancels a write transaction that loses a conflict with a +//! `TransactionConflict` reason on the contended item. These tests force the +//! PostgreSQL side of that: a real deadlock with an outside lock holder, and +//! many transactions that take the same items in opposite request orders. +//! +//! Each test builds its own throwaway database, applies the shipped migrations +//! to it, and drops it when it passes. A failing test leaves its database behind +//! on purpose, named `eddb_twic_*`, so the state that failed can be inspected. +//! +//! Requires `EXTENDDB_TEST_PG_CONNECTION_STRING`, a base URL with no database +//! component (for example `postgresql://postgres@127.0.0.1:5432`), pointing at a +//! server whose role may create and drop databases. Without it every test here +//! reports a skip and passes, the same convention the wire suites use. + +use std::collections::BTreeMap; +use std::time::Duration; + +use extenddb_core::expression::ExpressionMaps; +use extenddb_core::types::{ + AttributeDefinition, AttributeValue, BillingMode, CreateTableInput, Item, KeySchemaElement, + KeyType, ReturnValuesOnConditionCheckFailure, ScalarAttributeType, TableKeyInfo, +}; +use extenddb_storage::error::StorageError; +use extenddb_storage::{DataEngine, TableEngine, TransactWriteOp}; +use extenddb_storage_postgres::{PostgresConfig, PostgresEngine}; +use sqlx::PgPool; +use sqlx::postgres::PgPoolOptions; + +const ACCOUNT: &str = "123456789012"; +const REGION: &str = "us-east-1"; +const TABLE: &str = "t_twi_conflict"; + +struct Scratch { + engine: PostgresEngine, + db: PgPool, + admin: PgPool, + db_name: String, +} + +impl Scratch { + async fn cleanup(self) { + let Scratch { + engine, + db, + admin, + db_name, + } = self; + drop(engine); + db.close().await; + sqlx::query(&format!( + "DROP DATABASE IF EXISTS \"{db_name}\" WITH (FORCE)" + )) + .execute(&admin) + .await + .expect("drop the scratch database"); + admin.close().await; + } +} + +fn base_conn() -> Option { + let conn = std::env::var("EXTENDDB_TEST_PG_CONNECTION_STRING").ok()?; + (!conn.trim().is_empty()).then(|| conn.trim_end_matches('/').to_owned()) +} + +fn skip(test: &str) { + eprintln!( + "SKIP {test}: EXTENDDB_TEST_PG_CONNECTION_STRING is not set, so there is no PostgreSQL \ + to build a scratch catalog in." + ); +} + +async fn scratch() -> Scratch { + let base = base_conn().expect("caller checks base_conn() first"); + let db_name = format!("eddb_twic_{}", uuid::Uuid::new_v4().simple())[..24].to_owned(); + let admin = PgPoolOptions::new() + .max_connections(1) + .connect(&format!("{base}/postgres")) + .await + .expect("connect to the postgres maintenance database"); + sqlx::query(&format!("CREATE DATABASE \"{db_name}\"")) + .execute(&admin) + .await + .expect("create the scratch database"); + let url = format!("{base}/{db_name}"); + let db = PgPoolOptions::new() + .max_connections(2) + .connect(&url) + .await + .expect("connect to the scratch database"); + for sql in [ + include_str!("../migrations/001_schema.sql"), + include_str!("../migrations/002_vector_indexes.sql"), + include_str!("../data_migrations/001_data_schema.sql"), + include_str!("../data_migrations/002_gsi_pending.sql"), + include_str!("../data_migrations/003_idempotency_account_scope.sql"), + include_str!("../data_migrations/004_vector_index_state.sql"), + ] { + sqlx::raw_sql(sql) + .execute(&db) + .await + .expect("apply a shipped migration"); + } + sqlx::query("UPDATE settings SET value = '0' WHERE key = 'control_plane_delay_seconds'") + .execute(&db) + .await + .expect("pin the control-plane delay to zero"); + sqlx::query("INSERT INTO accounts (account_id, account_name) VALUES ($1, $2)") + .bind(ACCOUNT) + .bind(format!("acct-{db_name}")) + .execute(&db) + .await + .expect("seed the account row"); + let engine = PostgresEngine::new( + &PostgresConfig { + connection_string: url, + pool_size: 10, + max_item_size_bytes: 400_000, + }, + REGION, + ) + .await + .expect("open a PostgresEngine on the scratch database"); + Scratch { + engine, + db, + admin, + db_name, + } +} + +/// Create a hash-key table holding items `a` and `b`, and return its key info. +async fn seeded_table(s: &Scratch, maps: &ExpressionMaps) -> TableKeyInfo { + s.engine + .create_table( + ACCOUNT, + CreateTableInput { + table_name: TABLE.to_owned(), + key_schema: vec![KeySchemaElement { + attribute_name: "pk".to_owned(), + key_type: KeyType::Hash, + }], + attribute_definitions: vec![AttributeDefinition { + attribute_name: "pk".to_owned(), + attribute_type: ScalarAttributeType::S, + }], + billing_mode: Some(BillingMode::PayPerRequest), + ..Default::default() + }, + ) + .await + .expect("create the table"); + let key_info = s + .engine + .table_key_info(ACCOUNT, TABLE) + .await + .expect("read the key info"); + for pk in ["a", "b"] { + s.engine + .put_item(&key_info, item(pk, "seed"), false, None, maps, None) + .await + .expect("seed an item"); + } + key_info +} + +fn item(pk: &str, v: &str) -> Item { + BTreeMap::from([ + ("pk".to_owned(), AttributeValue::S(pk.to_owned())), + ("v".to_owned(), AttributeValue::S(v.to_owned())), + ]) +} + +fn put<'a>( + key_info: &'a TableKeyInfo, + item: &'a Item, + maps: &'a ExpressionMaps, +) -> TransactWriteOp<'a> { + TransactWriteOp::Put { + key_info, + item, + condition: None, + maps, + return_values_on_ccf: ReturnValuesOnConditionCheckFailure::None, + stream: None, + } +} + +async fn data_table(db: &PgPool) -> String { + let id: String = + sqlx::query_scalar("SELECT table_id FROM tables WHERE account_id = $1 AND table_name = $2") + .bind(ACCOUNT) + .bind(TABLE) + .fetch_one(db) + .await + .expect("look up the table id"); + format!("\"_ddb_{id}\"") +} + +/// Wait until some backend other than `own_pid` waits on a row lock. +async fn wait_for_lock_waiter(db: &PgPool, own_pid: i32) { + for _ in 0..500 { + let waiting: i64 = sqlx::query_scalar( + "SELECT count(*) FROM pg_stat_activity \ + WHERE datname = current_database() AND wait_event_type = 'Lock' AND pid <> $1", + ) + .bind(own_pid) + .fetch_one(db) + .await + .expect("read pg_stat_activity"); + if waiting > 0 { + return; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } + panic!("the transaction never waited on the held lock"); +} + +#[tokio::test] +async fn a_deadlock_cancels_with_transaction_conflict_on_the_waiting_item() { + let test = "a_deadlock_cancels_with_transaction_conflict_on_the_waiting_item"; + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let maps = ExpressionMaps::default(); + let key_info = seeded_table(&s, &maps).await; + let table = data_table(&s.db).await; + + // An outside transaction holds `b`. The write transaction locks `a`, then + // waits on `b`. The outside transaction then asks for `a`, which closes the + // cycle. The write transaction waited first, so its deadlock check runs + // first and PostgreSQL aborts it. + let mut holder = s.db.begin().await.expect("begin the outside transaction"); + let holder_pid: i32 = sqlx::query_scalar("SELECT pg_backend_pid()") + .fetch_one(&mut *holder) + .await + .expect("read the backend pid"); + sqlx::query(&format!("SELECT 1 FROM {table} WHERE pk = 'b' FOR UPDATE")) + .execute(&mut *holder) + .await + .expect("lock b"); + + let (new_a, new_b) = (item("a", "twi"), item("b", "twi")); + // Request order is b, a: the reason must land on b's request position. + let ops = [put(&key_info, &new_b, &maps), put(&key_info, &new_a, &maps)]; + let twi = s.engine.transact_write_items(&ops, None); + let close_cycle = async { + wait_for_lock_waiter(&s.db, holder_pid).await; + sqlx::query(&format!("SELECT 1 FROM {table} WHERE pk = 'a' FOR UPDATE")) + .execute(&mut *holder) + .await + .expect("the outside transaction gets a once the write transaction aborts"); + }; + // Bounded, so a transaction that never closes the cycle fails the test + // instead of hanging it. + let (result, ()) = tokio::time::timeout(Duration::from_secs(30), async { + tokio::join!(twi, close_cycle) + }) + .await + .expect("the write transaction finished"); + holder.rollback().await.expect("release the outside locks"); + + match result { + Err(StorageError::TransactionCanceled(reasons)) => { + let codes: Vec<&str> = reasons.iter().map(|r| r.code.as_str()).collect(); + assert_eq!(codes, ["TransactionConflict", "None"], "{reasons:?}"); + assert_eq!( + reasons[0].message.as_deref(), + Some("Transaction is ongoing for the item") + ); + assert_eq!(reasons[1].message, None); + } + other => panic!("expected a TransactionConflict cancellation, got {other:?}"), + } + // The canceled transaction applied nothing. + for pk in ["a", "b"] { + let v: String = sqlx::query_scalar(&format!( + "SELECT item_data->'v'->>'S' FROM {table} WHERE pk = $1" + )) + .bind(pk) + .fetch_one(&s.db) + .await + .expect("read the item"); + assert_eq!(v, "seed", "item {pk}"); + } + + s.cleanup().await; +} + +#[tokio::test] +async fn opposite_request_orders_do_not_deadlock() { + let test = "opposite_request_orders_do_not_deadlock"; + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let maps = ExpressionMaps::default(); + let key_info = seeded_table(&s, &maps).await; + + // Eight writers, half naming the items as (a, b) and half as (b, a). Taken + // in request order these deadlock within a few rounds. + let writers = (0..8).map(|w| { + let (engine, key_info, maps) = (&s.engine, &key_info, &maps); + async move { + for round in 0..50 { + let tag = format!("w{w}-r{round}"); + let (new_a, new_b) = (item("a", &tag), item("b", &tag)); + let ops = if w % 2 == 0 { + [put(key_info, &new_a, maps), put(key_info, &new_b, maps)] + } else { + [put(key_info, &new_b, maps), put(key_info, &new_a, maps)] + }; + engine + .transact_write_items(&ops, None) + .await + .unwrap_or_else(|e| panic!("writer {w} round {round}: {e:?}")); + } + } + }); + futures::future::join_all(writers).await; + + s.cleanup().await; +} diff --git a/tests/test_transact_write_conflict.py b/tests/test_transact_write_conflict.py new file mode 100644 index 000000000..d9a36c20a --- /dev/null +++ b/tests/test_transact_write_conflict.py @@ -0,0 +1,208 @@ +# Copyright 2026 ExtendDB contributors +# SPDX-License-Identifier: Apache-2.0 + +"""Concurrent TransactWriteItems that contend for the same items. + +Amazon DynamoDB answers a write transaction that loses a conflict with +TransactionCanceledException (HTTP 400). The contended items carry a +TransactionConflict reason and the other items carry None. Under this load +Amazon DynamoDB also answers a few transactions (about 1 in 1000) with +InternalServerError, so a 5xx is tolerated only as a small share of the +cancellations. A server that answers conflicts with a 5xx instead of a +cancellation fails. Every committed transaction must apply in full. +""" + +from __future__ import annotations + +import os +import threading +import time +import uuid +from concurrent.futures import ThreadPoolExecutor + +import boto3 +import pytest +from botocore.config import Config +from botocore.exceptions import ClientError + +from conftest import scoped_table + +WORKERS = 4 +TXNS_PER_WORKER = 25 +# Bounds the run on a server that stalls on each conflict. +TIME_BUDGET_S = 20.0 +CONFLICT_MESSAGE = "Transaction is ongoing for the item" +CANCEL_PREFIX = ( + "Transaction cancelled, please refer cancellation reasons for specific reasons [" +) + + +@pytest.fixture(scope="module") +def raw_client(endpoint_url): + """A client that never retries, so every 5xx and every cancellation is seen.""" + kwargs: dict = { + "service_name": "dynamodb", + "region_name": os.environ.get("AWS_DEFAULT_REGION", "us-east-1"), + "config": Config( + retries={"total_max_attempts": 1, "mode": "standard"}, + max_pool_connections=WORKERS * 2, + ), + } + if endpoint_url: + kwargs["endpoint_url"] = endpoint_url + if endpoint_url.startswith("https://"): + kwargs["verify"] = False + return boto3.client(**kwargs) + + +@pytest.fixture(scope="module") +def table(dynamodb_client): + with scoped_table(dynamodb_client) as name: + yield name + + +def _add_one(table: str, pk: str) -> dict: + return { + "Update": { + "TableName": table, + "Key": {"pk": {"S": pk}}, + "UpdateExpression": "ADD n :one", + "ExpressionAttributeValues": {":one": {"N": "1"}}, + } + } + + +def _run(client, builders) -> tuple[int, int, list[dict]]: + """Run each builder on WORKERS / len(builders) threads at once. + + Returns the attempts, the committed transactions, and one record per failure. + """ + lock = threading.Lock() + attempts = 0 + committed = 0 + failures: list[dict] = [] + start = threading.Barrier(WORKERS) + deadline = time.monotonic() + TIME_BUDGET_S + + def worker(which: int): + nonlocal attempts, committed + start.wait() + for _ in range(TXNS_PER_WORKER): + if time.monotonic() > deadline: + return + with lock: + attempts += 1 + try: + client.transact_write_items(TransactItems=builders[which]()) + with lock: + committed += 1 + except ClientError as e: + r = e.response + with lock: + failures.append( + { + "builder": which, + "status": r["ResponseMetadata"]["HTTPStatusCode"], + "code": r["Error"]["Code"], + "message": r["Error"]["Message"], + "reasons": r.get("CancellationReasons"), + } + ) + + with ThreadPoolExecutor(max_workers=WORKERS) as pool: + futures = [ + pool.submit(worker, i % len(builders)) for i in range(WORKERS) + ] + for f in futures: + f.result() + return attempts, committed, failures + + +def _assert_conflict_shape(failures: list[dict], n_items: int) -> tuple[list[list[str]], int]: + """Assert the failures are conflict cancellations. + + Returns the reason codes of each cancellation and the number of 5xx. + """ + server_errors = [f for f in failures if f["status"] >= 500] + cancellations = [f for f in failures if f["status"] < 500] + for f in server_errors: + assert (f["status"], f["code"]) == (500, "InternalServerError"), f + assert len(server_errors) * 5 <= len(cancellations), ( + f"{len(server_errors)} 5xx against {len(cancellations)} cancellations, " + f"first: {server_errors[0]}" + ) + shapes = [] + for f in cancellations: + assert f["status"] == 400, f + assert f["code"] == "TransactionCanceledException", f + reasons = f["reasons"] + assert reasons is not None and len(reasons) == n_items, f + codes = [r["Code"] for r in reasons] + assert set(codes) <= {"None", "TransactionConflict"}, f + assert "TransactionConflict" in codes, f + for r in reasons: + if r["Code"] == "TransactionConflict": + assert r.get("Message") == CONFLICT_MESSAGE, f + else: + assert "Message" not in r, f + assert f["message"] == CANCEL_PREFIX + ", ".join(codes) + "]", f + shapes.append(codes) + return shapes, len(server_errors) + + +def _counter(client, table: str, pk: str) -> int: + item = client.get_item(TableName=table, Key={"pk": {"S": pk}}, ConsistentRead=True) + return int(item["Item"]["n"]["N"]) + + +def test_opposite_order_transactions_cancel_instead_of_failing( + dynamodb_client, raw_client, table +): + """Two items updated in opposite orders by many clients at once.""" + a, b = f"hot-a-{uuid.uuid4().hex[:8]}", f"hot-b-{uuid.uuid4().hex[:8]}" + for pk in (a, b): + dynamodb_client.put_item(TableName=table, Item={"pk": {"S": pk}, "n": {"N": "0"}}) + + attempts, committed, failures = _run( + raw_client, + [ + lambda: [_add_one(table, a), _add_one(table, b)], + lambda: [_add_one(table, b), _add_one(table, a)], + ], + ) + + _, n_5xx = _assert_conflict_shape(failures, 2) + assert committed + len(failures) == attempts + # Every committed transaction applied both updates and no canceled one + # applied any. A 5xx leaves the outcome unknown, so it may count either way. + n_a, n_b = _counter(dynamodb_client, table, a), _counter(dynamodb_client, table, b) + assert n_a == n_b, (n_a, n_b) + assert committed <= n_a <= committed + n_5xx, (n_a, committed, n_5xx) + + +def test_conflict_reason_names_only_the_contended_item( + dynamodb_client, raw_client, table +): + """One shared item plus one private item per transaction, in either position.""" + hot = f"hot-{uuid.uuid4().hex[:8]}" + dynamodb_client.put_item(TableName=table, Item={"pk": {"S": hot}, "n": {"N": "0"}}) + + def private() -> str: + return f"own-{uuid.uuid4().hex}" + + attempts, committed, failures = _run( + raw_client, + [ + lambda: [_add_one(table, hot), _add_one(table, private())], + lambda: [_add_one(table, private()), _add_one(table, hot)], + ], + ) + + shapes, n_5xx = _assert_conflict_shape(failures, 2) + cancellations = [f for f in failures if f["status"] < 500] + for codes, f in zip(shapes, cancellations): + # The private item never conflicts, so only the shared item is named. + expected = ["TransactionConflict", "None"] + assert codes == (expected if f["builder"] == 0 else expected[::-1]), f + assert committed + len(failures) == attempts + assert committed <= _counter(dynamodb_client, table, hot) <= committed + n_5xx From e6edd276441f30f7f9670a60c98dd548467b047d Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Thu, 1 Oct 2026 19:09:06 +0000 Subject: [PATCH 2/6] fix(postgres): stop a write transaction once its validation answer is known With the ops running in key order, a request whose earliest invalid op is known no longer runs and locks the rest. A storage test pins that the earliest invalid op in request order is still the one reported, and CI now runs the conflict storage tests against its PostgreSQL. Assisted-by: pi claude-opus-5-5 --- .github/workflows/integration.yml | 8 +- .../storage-postgres/src/data/transactions.rs | 14 +++- crates/storage-postgres/tests/twi_conflict.rs | 80 ++++++++++++++++++- 3 files changed, 96 insertions(+), 6 deletions(-) diff --git a/.github/workflows/integration.yml b/.github/workflows/integration.yml index 5d1ab88ad..1230b1ada 100644 --- a/.github/workflows/integration.yml +++ b/.github/workflows/integration.yml @@ -297,12 +297,14 @@ jobs: # The control plane for vector indexes is not reachable over the wire while # this backend declares no vector search capability, so its tests drive the - # storage layer directly against this job's PostgreSQL. They build their own - # throwaway databases; the connection string is the server, not a database. + # storage layer directly against this job's PostgreSQL. The write-transaction + # conflict tests force a deadlock, which also needs direct access. They build + # their own throwaway databases; the connection string is the server, not a + # database. - name: Run PostgreSQL storage-level tests env: EXTENDDB_TEST_PG_CONNECTION_STRING: postgresql://postgres:devpass@127.0.0.1:5432 - run: cargo test --release -p extenddb-storage-postgres --test vector_control_plane + run: cargo test --release -p extenddb-storage-postgres --test vector_control_plane --test twi_conflict # The daemonized server logs to syslog; dump it so server-side failures # are diagnosable from the job log. diff --git a/crates/storage-postgres/src/data/transactions.rs b/crates/storage-postgres/src/data/transactions.rs index 7c633d0f8..8343bd4ea 100644 --- a/crates/storage-postgres/src/data/transactions.rs +++ b/crates/storage-postgres/src/data/transactions.rs @@ -115,6 +115,8 @@ impl PostgresEngine { .map_err(|e| StorageError::Internal(e.to_string()))?; // A conflict abort outside the per-op loop cannot be tied to one item. + // Unreachable at READ COMMITTED, where the ops already hold every row + // lock; a stricter operator isolation can raise 40001 at commit. let cancel_all = |e: StorageError| conflict_cancels_all(e, ops.len()); // Check the idempotency token within the transaction so token storage @@ -131,6 +133,7 @@ impl PostgresEngine { let mut any_failed = false; let mut first_invalid: Option<(usize, String)> = None; let mut failed_op: Option<(usize, StorageError)> = None; + let mut ran = vec![false; ops.len()]; // Run the ops in key order, not request order, so every transaction // locks its items in one global order and two cannot deadlock on each @@ -162,11 +165,20 @@ impl PostgresEngine { } Err(TxnOpError::Storage(e)) => { // Infrastructure error, or PostgreSQL aborted the - // transaction: run no further ops. + // transaction: run no further ops. A request-earlier op + // that sorts later is then not validated. failed_op = Some((i, e)); break; } } + ran[i] = true; + // Once every op before the earliest invalid one has run, the + // answer is fixed: stop instead of locking the rest. + if let Some((j, _)) = &first_invalid + && ran[..*j].iter().all(|r| *r) + { + break; + } } if let Some((_, msg)) = first_invalid { diff --git a/crates/storage-postgres/tests/twi_conflict.rs b/crates/storage-postgres/tests/twi_conflict.rs index fb0b39995..406c5e2ac 100644 --- a/crates/storage-postgres/tests/twi_conflict.rs +++ b/crates/storage-postgres/tests/twi_conflict.rs @@ -21,8 +21,9 @@ use std::time::Duration; use extenddb_core::expression::ExpressionMaps; use extenddb_core::types::{ - AttributeDefinition, AttributeValue, BillingMode, CreateTableInput, Item, KeySchemaElement, - KeyType, ReturnValuesOnConditionCheckFailure, ScalarAttributeType, TableKeyInfo, + AttributeDefinition, AttributeValue, BillingMode, CreateTableInput, GsiInput, Item, + KeySchemaElement, KeyType, Projection, ProjectionType, ReturnValuesOnConditionCheckFailure, + ScalarAttributeType, TableKeyInfo, }; use extenddb_storage::error::StorageError; use extenddb_storage::{DataEngine, TableEngine, TransactWriteOp}; @@ -325,3 +326,78 @@ async fn opposite_request_orders_do_not_deadlock() { s.cleanup().await; } + +#[tokio::test] +async fn the_earliest_invalid_op_in_request_order_is_reported() { + // The ops run in key order, but Amazon DynamoDB names the first invalid + // item of the request: [z, a] names z's index, [a, z] names a's. + let test = "the_earliest_invalid_op_in_request_order_is_reported"; + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let maps = ExpressionMaps::default(); + let s_attr = |name: &str| AttributeDefinition { + attribute_name: name.to_owned(), + attribute_type: ScalarAttributeType::S, + }; + let hash = |name: &str| KeySchemaElement { + attribute_name: name.to_owned(), + key_type: KeyType::Hash, + }; + let gsi = |index: &str, attr: &str| GsiInput { + index_name: index.to_owned(), + key_schema: vec![hash(attr)], + projection: Projection { + projection_type: ProjectionType::All, + non_key_attributes: None, + }, + provisioned_throughput: None, + }; + s.engine + .create_table( + ACCOUNT, + CreateTableInput { + table_name: TABLE.to_owned(), + key_schema: vec![hash("pk")], + attribute_definitions: vec![s_attr("pk"), s_attr("a1"), s_attr("a2")], + billing_mode: Some(BillingMode::PayPerRequest), + global_secondary_indexes: Some(vec![gsi("gi1", "a1"), gsi("gi2", "a2")]), + ..Default::default() + }, + ) + .await + .expect("create the table"); + let key_info = s + .engine + .table_key_info(ACCOUNT, TABLE) + .await + .expect("read the key info"); + let empty_key = |pk: &str, attr: &str| -> Item { + BTreeMap::from([ + ("pk".to_owned(), AttributeValue::S(pk.to_owned())), + (attr.to_owned(), AttributeValue::S(String::new())), + ]) + }; + let (z, a) = (empty_key("z", "a2"), empty_key("a", "a1")); + + for (ops, expected) in [ + ( + [put(&key_info, &z, &maps), put(&key_info, &a, &maps)], + "IndexName: gi2, IndexKey: a2", + ), + ( + [put(&key_info, &a, &maps), put(&key_info, &z, &maps)], + "IndexName: gi1, IndexKey: a1", + ), + ] { + match s.engine.transact_write_items(&ops, None).await { + Err(StorageError::Validation(msg)) => { + assert!(msg.ends_with(expected), "{msg}"); + } + other => panic!("expected a ValidationException naming {expected}, got {other:?}"), + } + } + + s.cleanup().await; +} From 1d47fd6870dad717da07f54288ba90729c4a25e2 Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Thu, 1 Oct 2026 19:09:06 +0000 Subject: [PATCH 3/6] test: cap the 5xx a contended transaction run may return Bound the tolerated InternalServerError count to a fixed 5 per run instead of a share of the cancellations, so a server that cancels a lot cannot hide a 5xx rate far above the service's. Assisted-by: pi claude-opus-5-5 --- tests/test_transact_write_conflict.py | 20 ++++++++++++++------ 1 file changed, 14 insertions(+), 6 deletions(-) diff --git a/tests/test_transact_write_conflict.py b/tests/test_transact_write_conflict.py index d9a36c20a..63a012788 100644 --- a/tests/test_transact_write_conflict.py +++ b/tests/test_transact_write_conflict.py @@ -7,9 +7,14 @@ TransactionCanceledException (HTTP 400). The contended items carry a TransactionConflict reason and the other items carry None. Under this load Amazon DynamoDB also answers a few transactions (about 1 in 1000) with -InternalServerError, so a 5xx is tolerated only as a small share of the -cancellations. A server that answers conflicts with a 5xx instead of a -cancellation fails. Every committed transaction must apply in full. +InternalServerError, so a small fixed number of 5xx is tolerated. Every +committed transaction must apply in full. + +ExtendDB on PostgreSQL queues contending transactions instead of canceling +them, so there the cancellation shape checks have nothing to check and these +tests prove only that contention never surfaces as a 5xx. The mapping of a +database-detected deadlock to TransactionConflict is pinned by the storage +tests in crates/storage-postgres/tests/twi_conflict.rs. """ from __future__ import annotations @@ -31,6 +36,9 @@ TXNS_PER_WORKER = 25 # Bounds the run on a server that stalls on each conflict. TIME_BUDGET_S = 20.0 +# Amazon DynamoDB returned at most 3 InternalServerError in one run of 100 +# contended transactions; they come in bursts. +MAX_SERVER_ERRORS = 5 CONFLICT_MESSAGE = "Transaction is ongoing for the item" CANCEL_PREFIX = ( "Transaction cancelled, please refer cancellation reasons for specific reasons [" @@ -127,9 +135,8 @@ def _assert_conflict_shape(failures: list[dict], n_items: int) -> tuple[list[lis cancellations = [f for f in failures if f["status"] < 500] for f in server_errors: assert (f["status"], f["code"]) == (500, "InternalServerError"), f - assert len(server_errors) * 5 <= len(cancellations), ( - f"{len(server_errors)} 5xx against {len(cancellations)} cancellations, " - f"first: {server_errors[0]}" + assert len(server_errors) <= MAX_SERVER_ERRORS, ( + f"{len(server_errors)} 5xx in {len(failures)} failures, first: {server_errors[0]}" ) shapes = [] for f in cancellations: @@ -206,3 +213,4 @@ def private() -> str: assert codes == (expected if f["builder"] == 0 else expected[::-1]), f assert committed + len(failures) == attempts assert committed <= _counter(dynamodb_client, table, hot) <= committed + n_5xx + From f708177997703623f1a04718e5d52365e9e8d9ce Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Thu, 1 Oct 2026 19:09:06 +0000 Subject: [PATCH 4/6] docs: describe write transaction contention and lock order Add the contention difference to differences-from-dynamodb.md and the lock order rule for blocking-lock backends to the storage extension guide. Assisted-by: pi claude-opus-5-5 --- docs/differences-from-dynamodb.md | 1 + docs/manuals/12-extending-extenddb-storage.md | 1 + 2 files changed, 2 insertions(+) diff --git a/docs/differences-from-dynamodb.md b/docs/differences-from-dynamodb.md index da6a00f7f..eb6af34e1 100755 --- a/docs/differences-from-dynamodb.md +++ b/docs/differences-from-dynamodb.md @@ -15,6 +15,7 @@ adaptation when switching between ExtendDB and the real service. | Numeric precision on partition/sort keys (MongoDB backend only) | 38 significant digits | 34 significant digits (BSON Decimal128). Values that exceed this precision are rejected at write and query time with a ValidationException rather than silently downcast. PostgreSQL backend supports the full 38 digits. | | Inverted numeric `BETWEEN` on a sort key (MongoDB backend only) | ValidationException ("The BETWEEN operator requires upper bound to be greater than or equal to lower bound") | Same error in all practical cases. The inversion guard compares bounds via `f64`, so a `KeyConditionExpression` `BETWEEN` whose bounds are inverted only beyond f64's ~15–17 significant digits (e.g. `BETWEEN 10000000000000002 AND 10000000000000001`) is not rejected and returns an empty result set instead. Valid ranges are never wrongly rejected. | | Transaction read concern (MongoDB backend only) | No user-configurable equivalent | `snapshot` is the default, fidelity-preserving mode. With `majority` or `local`, `TransactGetItems` is not guaranteed a single point-in-time snapshot and transaction condition reads have weaker isolation. With `local`, a condition can be evaluated against data that is later rolled back after failover. | +| Contended `TransactWriteItems` | A transaction that conflicts with another in-flight transaction on the same item is canceled with a `TransactionConflict` reason ("Transaction is ongoing for the item") | **PostgreSQL** queues the second transaction on the first one's row locks and **SQLite** runs one writer at a time, so both normally commit and no `TransactionConflict` is returned. PostgreSQL locks each transaction's items in table and key order, so two `TransactWriteItems` cannot deadlock on each other. When PostgreSQL still aborts a transaction (a deadlock with another lock holder, or a serialization failure under a stricter `default_transaction_isolation`), the answer is a `TransactionCanceledException` with a `TransactionConflict` reason, not a 500. **MongoDB** retries write conflicts internally and cancels with `TransactionConflict` on every item only after sustained contention. A client that relies on `TransactionConflict` to back off does not see it under ordinary contention. | ## Authentication and Authorization (AWS IAM/STS auth surface used by DynamoDB) diff --git a/docs/manuals/12-extending-extenddb-storage.md b/docs/manuals/12-extending-extenddb-storage.md index d8f7a9b09..c59f59c74 100755 --- a/docs/manuals/12-extending-extenddb-storage.md +++ b/docs/manuals/12-extending-extenddb-storage.md @@ -93,6 +93,7 @@ Key design decisions: - **Condition expressions** are evaluated inside the storage transaction. The engine parses and compiles expressions; the storage layer receives an AST (`Expr`) and evaluates it against the existing item within the same transaction that performs the write. This is critical for correctness — condition checks and writes must be atomic. - **Stream capture** is passed as `Option<&StreamCapture>`. When present, the stream record must be written in the same transaction as the data write. - **Idempotency tokens** for `TransactWriteItems` must be checked and stored atomically with the writes, and must be scoped per account: a `ClientRequestToken` is unique per account in DynamoDB, not globally, so the store must key on `(account_id, token)`. Keying on the token alone lets the same token value from two accounts collide (one account's transaction wrongly replayed or rejected against another's). +- **Lock order in `transact_write_items`**: a backend that takes blocking row locks must lock a transaction's items in one global order, so two write transactions cannot deadlock on each other. The PostgreSQL backend runs the ops sorted by table ID and key and keeps the results in request order. When the database still aborts a transaction to break a conflict (a deadlock or a serialization failure), return `StorageError::TransactionCanceled` with a `TransactionConflict` reason for the affected item, never an internal error. - **Items** are `BTreeMap`. A new backend must handle the full `AttributeValue` type (S, N, B, SS, NS, BS, L, M, BOOL, NULL). - **Query** must support forward/reverse sort order, exclusive start key pagination, and routing to secondary index storage. - **Parallel scan** uses `segment` and `total_segments` to partition the keyspace. From 18447cbea7399950e2c62630d629d85419848b89 Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Thu, 1 Oct 2026 19:21:04 +0000 Subject: [PATCH 5/6] test: pin that an invalid first op answers without waiting on later locks A storage test holds the next op's row from outside the transaction and expects the ValidationException within 5 s. The pytest docstring now says which test can deadlock, and the file ends with one newline. Assisted-by: pi claude-opus-5-5 --- crates/storage-postgres/tests/twi_conflict.rs | 86 +++++++++++++++---- tests/test_transact_write_conflict.py | 5 +- 2 files changed, 70 insertions(+), 21 deletions(-) diff --git a/crates/storage-postgres/tests/twi_conflict.rs b/crates/storage-postgres/tests/twi_conflict.rs index 406c5e2ac..7087939dc 100644 --- a/crates/storage-postgres/tests/twi_conflict.rs +++ b/crates/storage-postgres/tests/twi_conflict.rs @@ -327,16 +327,8 @@ async fn opposite_request_orders_do_not_deadlock() { s.cleanup().await; } -#[tokio::test] -async fn the_earliest_invalid_op_in_request_order_is_reported() { - // The ops run in key order, but Amazon DynamoDB names the first invalid - // item of the request: [z, a] names z's index, [a, z] names a's. - let test = "the_earliest_invalid_op_in_request_order_is_reported"; - if base_conn().is_none() { - return skip(test); - } - let s = scratch().await; - let maps = ExpressionMaps::default(); +/// Create a hash-key table with GSIs `gi1` on `a1` and `gi2` on `a2`. +async fn gsi_table(s: &Scratch) -> TableKeyInfo { let s_attr = |name: &str| AttributeDefinition { attribute_name: name.to_owned(), attribute_type: ScalarAttributeType::S, @@ -368,17 +360,31 @@ async fn the_earliest_invalid_op_in_request_order_is_reported() { ) .await .expect("create the table"); - let key_info = s - .engine + s.engine .table_key_info(ACCOUNT, TABLE) .await - .expect("read the key info"); - let empty_key = |pk: &str, attr: &str| -> Item { - BTreeMap::from([ - ("pk".to_owned(), AttributeValue::S(pk.to_owned())), - (attr.to_owned(), AttributeValue::S(String::new())), - ]) - }; + .expect("read the key info") +} + +/// An item whose `attr` index key is the empty string. +fn empty_key(pk: &str, attr: &str) -> Item { + BTreeMap::from([ + ("pk".to_owned(), AttributeValue::S(pk.to_owned())), + (attr.to_owned(), AttributeValue::S(String::new())), + ]) +} + +#[tokio::test] +async fn the_earliest_invalid_op_in_request_order_is_reported() { + // The ops run in key order, but Amazon DynamoDB names the first invalid + // item of the request: [z, a] names z's index, [a, z] names a's. + let test = "the_earliest_invalid_op_in_request_order_is_reported"; + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let maps = ExpressionMaps::default(); + let key_info = gsi_table(&s).await; let (z, a) = (empty_key("z", "a2"), empty_key("a", "a1")); for (ops, expected) in [ @@ -401,3 +407,45 @@ async fn the_earliest_invalid_op_in_request_order_is_reported() { s.cleanup().await; } + +#[tokio::test] +async fn an_invalid_first_op_answers_without_waiting_on_later_locks() { + // Once the earliest invalid op is known, the rest must not run: here the + // next op's row is held by an outside transaction. + let test = "an_invalid_first_op_answers_without_waiting_on_later_locks"; + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let maps = ExpressionMaps::default(); + let key_info = gsi_table(&s).await; + let b = item("b", "seed"); + s.engine + .put_item(&key_info, b.clone(), false, None, &maps, None) + .await + .expect("seed b"); + let table = data_table(&s.db).await; + let mut holder = s.db.begin().await.expect("begin the outside transaction"); + sqlx::query(&format!("SELECT 1 FROM {table} WHERE pk = 'b' FOR UPDATE")) + .execute(&mut *holder) + .await + .expect("lock b"); + + let a = empty_key("a", "a1"); + let ops = [put(&key_info, &a, &maps), put(&key_info, &b, &maps)]; + let result = tokio::time::timeout( + Duration::from_secs(5), + s.engine.transact_write_items(&ops, None), + ) + .await + .expect("the transaction answered without waiting on b"); + holder.rollback().await.expect("release b"); + match result { + Err(StorageError::Validation(msg)) => { + assert!(msg.ends_with("IndexName: gi1, IndexKey: a1"), "{msg}"); + } + other => panic!("expected a ValidationException, got {other:?}"), + } + + s.cleanup().await; +} diff --git a/tests/test_transact_write_conflict.py b/tests/test_transact_write_conflict.py index 63a012788..76404aaa3 100644 --- a/tests/test_transact_write_conflict.py +++ b/tests/test_transact_write_conflict.py @@ -12,7 +12,9 @@ ExtendDB on PostgreSQL queues contending transactions instead of canceling them, so there the cancellation shape checks have nothing to check and these -tests prove only that contention never surfaces as a 5xx. The mapping of a +tests prove only that contention never surfaces as a 5xx. Only the +opposite-order test can deadlock; the shared-plus-private test checks the +reason positions on Amazon DynamoDB. The mapping of a database-detected deadlock to TransactionConflict is pinned by the storage tests in crates/storage-postgres/tests/twi_conflict.rs. """ @@ -213,4 +215,3 @@ def private() -> str: assert codes == (expected if f["builder"] == 0 else expected[::-1]), f assert committed + len(failures) == attempts assert committed <= _counter(dynamodb_client, table, hot) <= committed + n_5xx - From 0e242c59ca6d1b0d2206fdb8ef612171b5b116cf Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Thu, 1 Oct 2026 19:25:36 +0000 Subject: [PATCH 6/6] style: backtick the SQLSTATE names in the conflict classifier doc Assisted-by: pi claude-opus-5-5 --- crates/storage-postgres/src/pg_util.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/crates/storage-postgres/src/pg_util.rs b/crates/storage-postgres/src/pg_util.rs index 30e1273bd..9c5b8e2db 100755 --- a/crates/storage-postgres/src/pg_util.rs +++ b/crates/storage-postgres/src/pg_util.rs @@ -19,8 +19,8 @@ pub(crate) fn is_fk_violation(e: &sqlx::Error) -> bool { false } -/// Check if a mapped error is PostgreSQL aborting the transaction to break a -/// lock conflict: deadlock_detected (40P01) or serialization_failure (40001). +/// Check if a mapped error is `PostgreSQL` aborting the transaction to break a +/// lock conflict: `deadlock_detected` (40P01) or `serialization_failure` (40001). /// Matches the `SQLSTATE` prefix that `data::index::db_error` writes. pub(crate) fn is_conflict_abort(e: &extenddb_storage::error::StorageError) -> bool { match e {