From bc6a39239d0c93c70a52dc1f3b2b34657673c333 Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Mon, 5 Oct 2026 17:37:31 +0000 Subject: [PATCH 1/8] test: race PutItem calls that create the same item Racing puts to a new item must all succeed unless a put's own condition fails against the item before it, and their ReturnValues ALL_OLD images must form one chain, as on Amazon DynamoDB. The tests cover hash and range tables, a condition, a GSI table, and a stream table. Assisted-by: pi claude-opus-5-5 --- tests/test_put_item_create_race.py | 252 +++++++++++++++++++++++++++++ 1 file changed, 252 insertions(+) create mode 100644 tests/test_put_item_create_race.py diff --git a/tests/test_put_item_create_race.py b/tests/test_put_item_create_race.py new file mode 100644 index 00000000..8158cdef --- /dev/null +++ b/tests/test_put_item_create_race.py @@ -0,0 +1,252 @@ +# Copyright 2026 ExtendDB contributors +# SPDX-License-Identifier: Apache-2.0 + +"""Concurrent PutItem calls that create the same new item. + +Amazon DynamoDB applies the puts one after another. Each one succeeds unless +its own condition fails against the item before it: a put that finds the item +already created overwrites it, and with ReturnValues ALL_OLD it returns the +item it replaced. So the old images of one round form a single chain, from no +item to the final one. The tests cover the request shapes that make a put +read the current item first: a condition, ALL_OLD, a GSI, and a stream. +""" + +from __future__ import annotations + +import os +import threading +import time +import uuid + +import boto3 +import pytest +from botocore.config import Config +from botocore.exceptions import ClientError + +from conftest import scoped_table + +WRITERS = 8 +ROUNDS = 15 + + +@pytest.fixture(scope="module") +def raw_client(endpoint_url): + """A client that never retries, so every failure 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=WRITERS * 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 streams_client(endpoint_url): + kwargs: dict = { + "service_name": "dynamodbstreams", + "region_name": os.environ.get("AWS_DEFAULT_REGION", "us-east-1"), + } + if endpoint_url: + kwargs["endpoint_url"] = endpoint_url + if endpoint_url.startswith("https://"): + kwargs["verify"] = False + return boto3.client(**kwargs) + + +def _s(name: str) -> dict: + return {"AttributeName": name, "AttributeType": "S"} + + +def _k(name: str, kind: str) -> dict: + return {"AttributeName": name, "KeyType": kind} + + +@pytest.fixture(scope="module") +def hash_table(dynamodb_client): + with scoped_table(dynamodb_client) as name: + yield name + + +@pytest.fixture(scope="module") +def range_table(dynamodb_client): + with scoped_table( + dynamodb_client, + attribute_definitions=[_s("pk"), _s("sk")], + key_schema=[_k("pk", "HASH"), _k("sk", "RANGE")], + ) as name: + yield name + + +@pytest.fixture(scope="module") +def gsi_table(dynamodb_client): + with scoped_table( + dynamodb_client, + attribute_definitions=[_s("pk"), _s("gk")], + GlobalSecondaryIndexes=[ + { + "IndexName": "gk-index", + "KeySchema": [_k("gk", "HASH")], + "Projection": {"ProjectionType": "ALL"}, + } + ], + ) as name: + yield name + + +@pytest.fixture(scope="module") +def stream_table(dynamodb_client): + with scoped_table( + dynamodb_client, + StreamSpecification={"StreamEnabled": True, "StreamViewType": "NEW_AND_OLD_IMAGES"}, + ) as name: + yield name + + +def _key(table_kind: str, pk: str) -> dict: + key = {"pk": {"S": pk}} + if table_kind == "range": + key["sk"] = {"S": "s"} + return key + + +def _race(client, table: str, table_kind: str, extra: dict) -> tuple[str, list]: + """WRITERS clients put the same new key at once, each with its own `w`. + + Returns the key and, per writer, the old `w` it replaced (None when it + created the item), or the ClientError response when the put failed. + """ + pk = f"k-{uuid.uuid4().hex}" + start = threading.Barrier(WRITERS, timeout=30) + out: list = [None] * WRITERS + errors: list[BaseException] = [] + + def run(i: int): + item = {**_key(table_kind, pk), "w": {"S": f"w{i}"}, "gk": {"S": f"g{i}"}} + try: + start.wait() + r = client.put_item(TableName=table, Item=item, **extra) + out[i] = ("ok", r.get("Attributes", {}).get("w", {}).get("S")) + except ClientError as e: + out[i] = ("error", e.response["Error"]["Code"]) + except BaseException as e: # noqa: BLE001 - surfaced below + errors.append(e) + + threads = [threading.Thread(target=run, args=(i,)) for i in range(WRITERS)] + for t in threads: + t.start() + for t in threads: + t.join() + if errors: + raise errors[0] + return pk, out + + +def _final_w(client, table: str, table_kind: str, pk: str) -> str: + item = client.get_item(TableName=table, Key=_key(table_kind, pk), ConsistentRead=True) + return item["Item"]["w"]["S"] + + +def _assert_all_succeed(out: list): + failed = [o for o in out if o[0] != "ok"] + assert not failed, f"{len(failed)} of {WRITERS} puts failed: {failed}" + + +def _assert_one_chain(out: list, final: str): + """The old images lead from the final item back to no item, once each.""" + prev = {f"w{i}": old for i, (_, old) in enumerate(out)} + seen, cur = 0, final + while cur is not None and seen <= WRITERS: + seen += 1 + cur = prev[cur] + assert seen == WRITERS, f"old images do not form one chain: {prev}, final {final}" + + +@pytest.mark.parametrize("table_kind", ["hash", "range"]) +def test_racing_puts_with_all_old_all_succeed( + request, dynamodb_client, raw_client, table_kind +): + table = request.getfixturevalue(f"{table_kind}_table") + for _ in range(ROUNDS): + pk, out = _race(raw_client, table, table_kind, {"ReturnValues": "ALL_OLD"}) + _assert_all_succeed(out) + _assert_one_chain(out, _final_w(dynamodb_client, table, table_kind, pk)) + + +def test_racing_puts_check_their_condition_against_the_winner( + dynamodb_client, raw_client, range_table +): + """A condition that every item satisfies never fails, whoever created the item.""" + extra = {"ConditionExpression": "attribute_not_exists(zz)", "ReturnValues": "ALL_OLD"} + for _ in range(ROUNDS): + pk, out = _race(raw_client, range_table, "range", extra) + _assert_all_succeed(out) + _assert_one_chain(out, _final_w(dynamodb_client, range_table, "range", pk)) + + +def test_racing_conditional_creates_have_one_winner(raw_client, hash_table): + """Control: with attribute_not_exists(pk), exactly one put creates the item.""" + for _ in range(ROUNDS): + _, out = _race( + raw_client, hash_table, "hash", {"ConditionExpression": "attribute_not_exists(pk)"} + ) + codes = sorted(o[1] if o[0] == "error" else "ok" for o in out) + assert codes == ["ConditionalCheckFailedException"] * (WRITERS - 1) + ["ok"], codes + + +def test_racing_puts_on_a_gsi_table_all_succeed(dynamodb_client, raw_client, gsi_table): + for _ in range(ROUNDS): + pk, out = _race(raw_client, gsi_table, "hash", {}) + _assert_all_succeed(out) + final = dynamodb_client.get_item( + TableName=gsi_table, Key=_key("hash", pk), ConsistentRead=True + )["Item"] + assert final["gk"]["S"] == "g" + final["w"]["S"][1:], final + + +def _stream_events(streams_client, stream_arn: str, pk: str, want: int) -> list[dict]: + """Poll the stream until it holds `want` records for `pk` (or 30 s pass).""" + deadline = time.monotonic() + 30 + records: list[dict] = [] + while time.monotonic() < deadline: + records = [] + shards = streams_client.describe_stream(StreamArn=stream_arn)["StreamDescription"]["Shards"] + for shard in shards: + it = streams_client.get_shard_iterator( + StreamArn=stream_arn, ShardId=shard["ShardId"], ShardIteratorType="TRIM_HORIZON" + )["ShardIterator"] + for _ in range(20): + resp = streams_client.get_records(ShardIterator=it, Limit=1000) + records += [ + r for r in resp.get("Records", []) if r["dynamodb"]["Keys"]["pk"]["S"] == pk + ] + it = resp.get("NextShardIterator") + if not it or not resp.get("Records"): + break + if len(records) >= want: + break + time.sleep(1) + return records + + +def test_racing_puts_on_a_stream_table_all_succeed( + dynamodb_client, raw_client, streams_client, stream_table +): + for _ in range(ROUNDS): + pk, out = _race(raw_client, stream_table, "hash", {}) + _assert_all_succeed(out) + # The last round's stream: one INSERT, then a MODIFY per later put, each + # with the item it replaced as its old image. + arn = dynamodb_client.describe_table(TableName=stream_table)["Table"]["LatestStreamArn"] + records = _stream_events(streams_client, arn, pk, WRITERS) + names = sorted(r["eventName"] for r in records) + assert names == ["INSERT"] + ["MODIFY"] * (WRITERS - 1), names + olds = [r["dynamodb"]["OldImage"]["w"]["S"] for r in records if r["eventName"] == "MODIFY"] + news = [r["dynamodb"]["NewImage"]["w"]["S"] for r in records] + assert len(set(olds)) == WRITERS - 1 and set(olds) <= set(news), (olds, news) From 6f0808114e872e144d690586746f7a1989d66032 Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Mon, 5 Oct 2026 17:37:31 +0000 Subject: [PATCH 2/8] fix(postgres): order a put that loses a create race after the winner A PutItem takes the transactional path when it has a condition, ReturnValues ALL_OLD, ReturnConsumedCapacity, an index, a vector index, or a stream. BatchWriteItem puts take the same path on the same tables, and with ReturnConsumedCapacity INDEXES. On that path a put inserts a missing item with ON CONFLICT DO NOTHING. When a concurrent writer created the item first, the put returned ConditionalCheckFailed, even with no condition. BatchWriteItem returned it too, which Amazon DynamoDB never does. The put now re-reads the winner's row, checks its condition against it, and overwrites it, as UpdateItem already does. ALL_OLD, the index updates, and the stream record get the winner as the old image. If the winner was deleted before the re-read, the put retries the insert, up to 5 inserts. Storage tests hold the competing create open against a live PostgreSQL, and CI runs them. Assisted-by: pi claude-opus-5-5 --- .github/workflows/integration.yml | 8 +- crates/storage-postgres/src/data/put_item.rs | 134 ++++--- .../storage-postgres/tests/put_create_race.rs | 360 ++++++++++++++++++ 3 files changed, 444 insertions(+), 58 deletions(-) create mode 100644 crates/storage-postgres/tests/put_create_race.rs diff --git a/.github/workflows/integration.yml b/.github/workflows/integration.yml index 5d1ab88a..d8d45e6d 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 PutItem create + # race tests hold a create open in an outside transaction, 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 put_create_race # 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/put_item.rs b/crates/storage-postgres/src/data/put_item.rs index e307df66..f3a826bf 100755 --- a/crates/storage-postgres/src/data/put_item.rs +++ b/crates/storage-postgres/src/data/put_item.rs @@ -112,25 +112,28 @@ impl PostgresEngine { .await .map_err(|e| StorageError::Internal(e.to_string()))?; - let old: Option<(serde_json::Value,)> = + let mut old: Option<(serde_json::Value,)> = bind_sk_fetch_optional!(&select_sql, pk_text.as_str(), &sk, &mut *tx)?; - if let Some((ref old_json,)) = old { - let old_item: Item = json_to_item(old_json.clone())?; - match check_condition(condition, &old_item, maps) { - Ok(()) => {} - Err(StorageError::ConditionFailed(_)) => { - return Err(StorageError::ConditionFailed(Some(old_item))); + let mut attempt: u32 = 0; + loop { + if let Some((ref old_json,)) = old { + let old_item: Item = json_to_item(old_json.clone())?; + match check_condition(condition, &old_item, maps) { + Ok(()) => {} + Err(StorageError::ConditionFailed(_)) => { + return Err(StorageError::ConditionFailed(Some(old_item))); + } + Err(e) => return Err(e), } - Err(e) => return Err(e), + // Row exists, condition passed: update in place. + let update_sql = format!( + "UPDATE {ddb_table} SET item_data = $3 WHERE pk = $1 AND {sk_col} = $2" + ); + bind_sk_execute!(&update_sql, pk_text.as_str(), &sk, &item_json, &mut *tx)?; + break; } - // Row exists, condition passed — update in place. - let update_sql = format!( - "UPDATE {ddb_table} SET item_data = $3 WHERE pk = $1 AND {sk_col} = $2" - ); - bind_sk_execute!(&update_sql, pk_text.as_str(), &sk, &item_json, &mut *tx)?; - } else { - // No existing item — condition checks against empty item + // No existing item: the condition checks against an empty item let empty = std::collections::BTreeMap::new(); match check_condition(condition, &empty, maps) { Ok(()) => {} @@ -139,21 +142,23 @@ impl PostgresEngine { } Err(e) => return Err(e), } - // Condition passed against empty — atomic insert, fail if someone beat us. let insert_sql = format!( "INSERT INTO {ddb_table} (pk, {sk_col}, item_data) VALUES ($1, $2, $3) \ ON CONFLICT (pk, {sk_col}) DO NOTHING" ); let result = bind_sk_execute!(&insert_sql, pk_text.as_str(), &sk, &item_json, &mut *tx)?; - if result.rows_affected() == 0 { - // Another transaction inserted between our SELECT and INSERT. - // Fetch the winner to return with ConditionFailed. - let winner: Option<(serde_json::Value,)> = - bind_sk_fetch_optional!(&select_sql, pk_text.as_str(), &sk, &mut *tx)?; - let winner_item = winner.map(|(v,)| json_to_item(v)).transpose()?; - return Err(StorageError::ConditionFailed(winner_item)); + if result.rows_affected() == 1 { + break; + } + // Lost the create race. The insert waited for the winner to + // commit, so the locking read now returns its row: this put + // overwrites it, after its condition is checked against it. + attempt += 1; + if attempt >= MAX_CREATE_RACE_ATTEMPTS { + return Err(create_race_exhausted(attempt)); } + old = bind_sk_fetch_optional!(&select_sql, pk_text.as_str(), &sk, &mut *tx)?; } // Sync GSI/LSI update within transaction (D-4). @@ -269,30 +274,37 @@ impl PostgresEngine { .await .map_err(|e| StorageError::Internal(e.to_string()))?; - let old: Option<(serde_json::Value,)> = sqlx::query_as(&select_sql) - .bind(pk_text.as_str()) - .fetch_optional(&mut *tx) - .await - .map_err(|e| StorageError::Internal(e.to_string()))?; - - if let Some((ref old_json,)) = old { - let old_item: Item = json_to_item(old_json.clone())?; - match check_condition(condition, &old_item, maps) { - Ok(()) => {} - Err(StorageError::ConditionFailed(_)) => { - return Err(StorageError::ConditionFailed(Some(old_item))); - } - Err(e) => return Err(e), - } - // Row exists, condition passed — update in place. - let update_sql = format!("UPDATE {ddb_table} SET item_data = $2 WHERE pk = $1"); - sqlx::query(&update_sql) + let fetch_old = async |tx: &mut sqlx::PgConnection| { + sqlx::query_as::<_, (serde_json::Value,)>(&select_sql) .bind(pk_text.as_str()) - .bind(&item_json) - .execute(&mut *tx) + .fetch_optional(tx) .await - .map_err(|e| StorageError::Internal(e.to_string()))?; - } else { + .map_err(|e| StorageError::Internal(e.to_string())) + }; + let mut old: Option<(serde_json::Value,)> = fetch_old(&mut tx).await?; + + let mut attempt: u32 = 0; + loop { + if let Some((ref old_json,)) = old { + let old_item: Item = json_to_item(old_json.clone())?; + match check_condition(condition, &old_item, maps) { + Ok(()) => {} + Err(StorageError::ConditionFailed(_)) => { + return Err(StorageError::ConditionFailed(Some(old_item))); + } + Err(e) => return Err(e), + } + // Row exists, condition passed: update in place. + let update_sql = + format!("UPDATE {ddb_table} SET item_data = $2 WHERE pk = $1"); + sqlx::query(&update_sql) + .bind(pk_text.as_str()) + .bind(&item_json) + .execute(&mut *tx) + .await + .map_err(|e| StorageError::Internal(e.to_string()))?; + break; + } let empty = std::collections::BTreeMap::new(); match check_condition(condition, &empty, maps) { Ok(()) => {} @@ -301,7 +313,6 @@ impl PostgresEngine { } Err(e) => return Err(e), } - // Condition passed against empty — atomic insert, fail if someone beat us. let insert_sql = format!( "INSERT INTO {ddb_table} (pk, item_data) VALUES ($1, $2) \ ON CONFLICT (pk) DO NOTHING" @@ -312,16 +323,15 @@ impl PostgresEngine { .execute(&mut *tx) .await .map_err(|e| StorageError::Internal(e.to_string()))?; - if result.rows_affected() == 0 { - // Another transaction inserted between our SELECT and INSERT. - let winner: Option<(serde_json::Value,)> = sqlx::query_as(&select_sql) - .bind(pk_text.as_str()) - .fetch_optional(&mut *tx) - .await - .map_err(|e| StorageError::Internal(e.to_string()))?; - let winner_item = winner.map(|(v,)| json_to_item(v)).transpose()?; - return Err(StorageError::ConditionFailed(winner_item)); + if result.rows_affected() == 1 { + break; } + // Lost the create race: overwrite the winner, as above. + attempt += 1; + if attempt >= MAX_CREATE_RACE_ATTEMPTS { + return Err(create_race_exhausted(attempt)); + } + old = fetch_old(&mut tx).await?; } // Sync GSI/LSI update within transaction (D-4). @@ -466,3 +476,17 @@ impl PostgresEngine { json_opt.map(json_to_item).transpose() } } + +/// Bound on insert retries when a put to a missing item keeps losing the +/// create race to a winner that is deleted again before the re-read. Same +/// bound as `UpdateItem`. +const MAX_CREATE_RACE_ATTEMPTS: u32 = 5; + +/// The error after `MAX_CREATE_RACE_ATTEMPTS` lost create races, as `UpdateItem` +/// reports it. +fn create_race_exhausted(attempt: u32) -> StorageError { + StorageError::Internal(format!( + "PutItem could not create the item after {attempt} attempts: a concurrent writer \ + repeatedly created and deleted the row" + )) +} diff --git a/crates/storage-postgres/tests/put_create_race.rs b/crates/storage-postgres/tests/put_create_race.rs new file mode 100644 index 00000000..286a7857 --- /dev/null +++ b/crates/storage-postgres/tests/put_create_race.rs @@ -0,0 +1,360 @@ +// Copyright 2026 ExtendDB contributors +// SPDX-License-Identifier: Apache-2.0 +//! Storage-level tests for a `PutItem` that races another writer to create +//! the same item. +//! +//! Amazon DynamoDB applies such puts one after another: the later put +//! overwrites the item, after its own condition is checked against it. These +//! tests hold the competing create open in an outside transaction, so the put +//! deterministically loses the insert and has to order itself after the winner. +//! +//! 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_putr_*`, 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::{self, Expr, ExpressionMaps}; +use extenddb_core::types::{ + AttributeDefinition, AttributeValue, BillingMode, CreateTableInput, Item, KeySchemaElement, + KeyType, ScalarAttributeType, TableKeyInfo, +}; +use extenddb_storage::error::StorageError; +use extenddb_storage::{DataEngine, TableEngine}; +use extenddb_storage_postgres::{PostgresConfig, PostgresEngine}; +use sqlx::postgres::PgPoolOptions; +use sqlx::{PgPool, Postgres, Transaction}; + +const ACCOUNT: &str = "123456789012"; +const REGION: &str = "us-east-1"; +const TABLE: &str = "t_put_create_race"; + +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_putr_{}", 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, + } +} + +fn s_key(name: &str, key_type: KeyType) -> (KeySchemaElement, AttributeDefinition) { + ( + KeySchemaElement { + attribute_name: name.to_owned(), + key_type, + }, + AttributeDefinition { + attribute_name: name.to_owned(), + attribute_type: ScalarAttributeType::S, + }, + ) +} + +/// Create the test table, keyed on `pk` and, when `range` is set, also on `sk`. +async fn table(s: &Scratch, range: bool) -> TableKeyInfo { + let mut keys = vec![s_key("pk", KeyType::Hash)]; + if range { + keys.push(s_key("sk", KeyType::Range)); + } + let (key_schema, attribute_definitions) = keys.into_iter().unzip(); + s.engine + .create_table( + ACCOUNT, + CreateTableInput { + table_name: TABLE.to_owned(), + key_schema, + attribute_definitions, + billing_mode: Some(BillingMode::PayPerRequest), + ..Default::default() + }, + ) + .await + .expect("create the table"); + s.engine + .table_key_info(ACCOUNT, TABLE) + .await + .expect("read the key info") +} + +fn item(range: bool, v: &str) -> Item { + let mut item = BTreeMap::from([ + ("pk".to_owned(), AttributeValue::S("c".to_owned())), + ("v".to_owned(), AttributeValue::S(v.to_owned())), + ]); + if range { + item.insert("sk".to_owned(), AttributeValue::S("1".to_owned())); + } + item +} + +fn condition(text: &str) -> Expr { + let tokens = expression::tokenize(text).expect("tokenize"); + expression::parse_condition(&tokens).expect("parse") +} + +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}\"") +} + +/// Begin an outside transaction that has created `winner` but not committed: a +/// create in flight. Returns it and its backend pid. +async fn create_in_flight( + db: &PgPool, + table: &str, + winner: &Item, +) -> (Transaction<'static, Postgres>, i32) { + let mut creator = db.begin().await.expect("begin the outside transaction"); + let pid: i32 = sqlx::query_scalar("SELECT pg_backend_pid()") + .fetch_one(&mut *creator) + .await + .expect("read the backend pid"); + let data = serde_json::to_value(winner).expect("serialize the winner"); + let sql = if winner.contains_key("sk") { + format!("INSERT INTO {table} (pk, sk_s, item_data) VALUES ('c', '1', $1)") + } else { + format!("INSERT INTO {table} (pk, item_data) VALUES ('c', $1)") + }; + sqlx::query(&sql) + .bind(data) + .execute(&mut *creator) + .await + .expect("insert the winner"); + (creator, pid) +} + +/// Wait until a backend other than `own_pid` waits on a 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 put never waited on the create in flight"); +} + +/// Run `put` (ReturnValues ALL_OLD) while the winner's create is in flight, +/// commit the create once the put waits on it, and return the put's result. +async fn put_after_winner( + s: &Scratch, + key_info: &TableKeyInfo, + range: bool, + condition: Option<&Expr>, +) -> Result, StorageError> { + let table = data_table(&s.db).await; + let (creator, creator_pid) = create_in_flight(&s.db, &table, &item(range, "winner")).await; + let maps = ExpressionMaps::default(); + let put = s + .engine + .put_item(key_info, item(range, "loser"), true, condition, &maps, None); + let commit = async { + wait_for_lock_waiter(&s.db, creator_pid).await; + creator.commit().await.expect("commit the create"); + }; + let (result, ()) = + tokio::time::timeout(Duration::from_secs(30), async { tokio::join!(put, commit) }) + .await + .expect("the put finished"); + result +} + +async fn stored_v(s: &Scratch) -> String { + let table = data_table(&s.db).await; + sqlx::query_scalar(&format!( + "SELECT item_data->'v'->>'S' FROM {table} WHERE pk = 'c'" + )) + .fetch_one(&s.db) + .await + .expect("read the stored item") +} + +async fn overwrites_the_winner(test: &str, range: bool) { + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let key_info = table(&s, range).await; + let old = put_after_winner(&s, &key_info, range, None) + .await + .expect("the put succeeds"); + assert_eq!( + old, + Some(item(range, "winner")), + "ALL_OLD returns the winner" + ); + assert_eq!( + stored_v(&s).await, + "loser", + "the later put overwrote the winner" + ); + s.cleanup().await; +} + +#[tokio::test] +async fn a_put_that_loses_the_create_race_overwrites_the_winner() { + overwrites_the_winner( + "a_put_that_loses_the_create_race_overwrites_the_winner", + false, + ) + .await; +} + +#[tokio::test] +async fn a_put_that_loses_the_create_race_overwrites_the_winner_on_a_range_table() { + overwrites_the_winner( + "a_put_that_loses_the_create_race_overwrites_the_winner_on_a_range_table", + true, + ) + .await; +} + +#[tokio::test] +async fn a_lost_create_race_checks_the_condition_against_the_winner() { + // The condition holds for both the missing item and the winner, so the + // put succeeds after the winner instead of failing on the lost insert. + let test = "a_lost_create_race_checks_the_condition_against_the_winner"; + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let key_info = table(&s, false).await; + let cond = condition("attribute_not_exists(gone)"); + let old = put_after_winner(&s, &key_info, false, Some(&cond)) + .await + .expect("the condition holds for the winner too"); + assert_eq!(old, Some(item(false, "winner"))); + assert_eq!(stored_v(&s).await, "loser"); + s.cleanup().await; +} + +#[tokio::test] +async fn a_lost_create_race_fails_a_condition_the_winner_breaks() { + // `attribute_not_exists(pk)` holds for the missing item but not for the + // winner: the put fails its condition and returns the winner it saw. + let test = "a_lost_create_race_fails_a_condition_the_winner_breaks"; + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let key_info = table(&s, false).await; + let cond = condition("attribute_not_exists(pk)"); + match put_after_winner(&s, &key_info, false, Some(&cond)).await { + Err(StorageError::ConditionFailed(old)) => { + assert_eq!(old, Some(item(false, "winner"))); + } + other => panic!("expected ConditionFailed, got {other:?}"), + } + assert_eq!(stored_v(&s).await, "winner", "the failed put wrote nothing"); + s.cleanup().await; +} From d1d883f33e84db700a1e3a61069452898e5c54af Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Mon, 5 Oct 2026 19:24:38 +0000 Subject: [PATCH 3/8] test: check index and stream contents after racing puts The GSI and LSI tests give each put its own index key and check that only the final item's entry stays in the index. The stream tests check every round in sequence order: each record's old image is the record before it. New tests race BatchWriteItem puts on the stream table, and PutItem calls with ReturnConsumedCapacity. Assisted-by: pi claude-opus-5-5 --- tests/test_put_item_create_race.py | 197 +++++++++++++++++++++++------ 1 file changed, 157 insertions(+), 40 deletions(-) diff --git a/tests/test_put_item_create_race.py b/tests/test_put_item_create_race.py index 8158cdef..cb573661 100644 --- a/tests/test_put_item_create_race.py +++ b/tests/test_put_item_create_race.py @@ -8,7 +8,10 @@ already created overwrites it, and with ReturnValues ALL_OLD it returns the item it replaced. So the old images of one round form a single chain, from no item to the final one. The tests cover the request shapes that make a put -read the current item first: a condition, ALL_OLD, a GSI, and a stream. +read the current item first: a condition, ALL_OLD, ReturnConsumedCapacity, a +GSI, an LSI, and a stream, and BatchWriteItem puts to a stream table. On index +and stream tables they also check that only the final item stays in the index, +and that each stream record's old image is the record before it. """ from __future__ import annotations @@ -100,6 +103,23 @@ def gsi_table(dynamodb_client): yield name +@pytest.fixture(scope="module") +def lsi_table(dynamodb_client): + with scoped_table( + dynamodb_client, + attribute_definitions=[_s("pk"), _s("sk"), _s("lk")], + key_schema=[_k("pk", "HASH"), _k("sk", "RANGE")], + LocalSecondaryIndexes=[ + { + "IndexName": "lk-index", + "KeySchema": [_k("pk", "HASH"), _k("lk", "RANGE")], + "Projection": {"ProjectionType": "ALL"}, + } + ], + ) as name: + yield name + + @pytest.fixture(scope="module") def stream_table(dynamodb_client): with scoped_table( @@ -116,11 +136,17 @@ def _key(table_kind: str, pk: str) -> dict: return key -def _race(client, table: str, table_kind: str, extra: dict) -> tuple[str, list]: +def _race( + client, table: str, table_kind: str, extra: dict, batch: bool = False +) -> tuple[str, list]: """WRITERS clients put the same new key at once, each with its own `w`. + Each writer's item also carries index keys unique to it: `gk` across the + whole table, and `lk` within the key. With `batch`, each writer sends a + BatchWriteItem with one PutRequest instead of a PutItem. + Returns the key and, per writer, the old `w` it replaced (None when it - created the item), or the ClientError response when the put failed. + created the item or used BatchWriteItem), or the error when the put failed. """ pk = f"k-{uuid.uuid4().hex}" start = threading.Barrier(WRITERS, timeout=30) @@ -128,9 +154,19 @@ def _race(client, table: str, table_kind: str, extra: dict) -> tuple[str, list]: errors: list[BaseException] = [] def run(i: int): - item = {**_key(table_kind, pk), "w": {"S": f"w{i}"}, "gk": {"S": f"g{i}"}} + item = { + **_key(table_kind, pk), + "w": {"S": f"w{i}"}, + "gk": {"S": f"g-{pk}-{i}"}, + "lk": {"S": f"l{i}"}, + } try: start.wait() + if batch: + r = client.batch_write_item(RequestItems={table: [{"PutRequest": {"Item": item}}]}) + left = r.get("UnprocessedItems") + out[i] = ("unprocessed", left) if left else ("ok", None) + return r = client.put_item(TableName=table, Item=item, **extra) out[i] = ("ok", r.get("Attributes", {}).get("w", {}).get("S")) except ClientError as e: @@ -148,9 +184,14 @@ def run(i: int): return pk, out +def _final(client, table: str, table_kind: str, pk: str) -> dict: + return client.get_item(TableName=table, Key=_key(table_kind, pk), ConsistentRead=True)[ + "Item" + ] + + def _final_w(client, table: str, table_kind: str, pk: str) -> str: - item = client.get_item(TableName=table, Key=_key(table_kind, pk), ConsistentRead=True) - return item["Item"]["w"]["S"] + return _final(client, table, table_kind, pk)["w"]["S"] def _assert_all_succeed(out: list): @@ -179,6 +220,18 @@ def test_racing_puts_with_all_old_all_succeed( _assert_one_chain(out, _final_w(dynamodb_client, table, table_kind, pk)) +def test_racing_puts_that_return_consumed_capacity_all_succeed( + dynamodb_client, raw_client, hash_table +): + extra = {"ReturnConsumedCapacity": "TOTAL"} + for _ in range(ROUNDS): + pk, out = _race(raw_client, hash_table, "hash", extra) + _assert_all_succeed(out) + assert _final_w(dynamodb_client, hash_table, "hash", pk) in { + f"w{i}" for i in range(WRITERS) + } + + def test_racing_puts_check_their_condition_against_the_winner( dynamodb_client, raw_client, range_table ): @@ -200,53 +253,117 @@ def test_racing_conditional_creates_have_one_winner(raw_client, hash_table): assert codes == ["ConditionalCheckFailedException"] * (WRITERS - 1) + ["ok"], codes -def test_racing_puts_on_a_gsi_table_all_succeed(dynamodb_client, raw_client, gsi_table): +def _gsi_keys(client, table: str) -> dict[str, list[str]]: + """Every `gk` in the GSI, grouped by the base key it points at.""" + by_pk: dict[str, list[str]] = {} + kwargs = {"TableName": table, "IndexName": "gk-index"} + while True: + resp = client.scan(**kwargs) + for item in resp["Items"]: + by_pk.setdefault(item["pk"]["S"], []).append(item["gk"]["S"]) + if "LastEvaluatedKey" not in resp: + return {pk: sorted(gks) for pk, gks in by_pk.items()} + kwargs["ExclusiveStartKey"] = resp["LastEvaluatedKey"] + + +def test_racing_puts_on_a_gsi_table_leave_only_the_final_entry( + dynamodb_client, raw_client, gsi_table +): + """The GSI holds the final item's entry and none of the items it replaced.""" + want: dict[str, list[str]] = {} for _ in range(ROUNDS): pk, out = _race(raw_client, gsi_table, "hash", {}) _assert_all_succeed(out) - final = dynamodb_client.get_item( - TableName=gsi_table, Key=_key("hash", pk), ConsistentRead=True - )["Item"] - assert final["gk"]["S"] == "g" + final["w"]["S"][1:], final + want[pk] = [_final(dynamodb_client, gsi_table, "hash", pk)["gk"]["S"]] + # The GSI is eventually consistent on Amazon DynamoDB, so poll. + deadline = time.monotonic() + 30 + while (got := _gsi_keys(dynamodb_client, gsi_table)) != want and time.monotonic() < deadline: + time.sleep(0.5) + wrong = {pk: (want[pk], got.get(pk)) for pk in want if got.get(pk) != want[pk]} + assert not wrong and got.keys() == want.keys(), f"stale GSI entries (want, got): {wrong}" -def _stream_events(streams_client, stream_arn: str, pk: str, want: int) -> list[dict]: - """Poll the stream until it holds `want` records for `pk` (or 30 s pass).""" +def test_racing_puts_on_an_lsi_table_leave_only_the_final_entry( + dynamodb_client, raw_client, lsi_table +): + """A strongly consistent LSI query returns the final item's entry only.""" + for _ in range(ROUNDS): + pk, out = _race(raw_client, lsi_table, "range", {"ReturnValues": "ALL_OLD"}) + _assert_all_succeed(out) + final = _final(dynamodb_client, lsi_table, "range", pk) + _assert_one_chain(out, final["w"]["S"]) + items = dynamodb_client.query( + TableName=lsi_table, + IndexName="lk-index", + ConsistentRead=True, + KeyConditionExpression="pk = :pk", + ExpressionAttributeValues={":pk": {"S": pk}}, + )["Items"] + assert [i["lk"]["S"] for i in items] == [final["lk"]["S"]], (final, items) + + +def _stream_records(streams_client, stream_arn: str) -> dict[str, list[dict]]: + """Every record in the stream, grouped by key, in sequence order.""" + by_pk: dict[str, list[dict]] = {} + shards = streams_client.describe_stream(StreamArn=stream_arn)["StreamDescription"]["Shards"] + for shard in shards: + it = streams_client.get_shard_iterator( + StreamArn=stream_arn, ShardId=shard["ShardId"], ShardIteratorType="TRIM_HORIZON" + )["ShardIterator"] + empty = 0 + for _ in range(50): + resp = streams_client.get_records(ShardIterator=it, Limit=1000) + for r in resp.get("Records", []): + by_pk.setdefault(r["dynamodb"]["Keys"]["pk"]["S"], []).append(r) + it = resp.get("NextShardIterator") + empty = 0 if resp.get("Records") else empty + 1 + if not it or empty >= 3: + break + for records in by_pk.values(): + records.sort(key=lambda r: int(r["dynamodb"]["SequenceNumber"])) + return by_pk + + +def _assert_stream_chains(dynamodb_client, streams_client, table: str, pks: list[str]): + """Per key: one INSERT, then a MODIFY per later put, each with the item + before it as its old image.""" + arn = dynamodb_client.describe_table(TableName=table)["Table"]["LatestStreamArn"] deadline = time.monotonic() + 30 - records: list[dict] = [] - while time.monotonic() < deadline: - records = [] - shards = streams_client.describe_stream(StreamArn=stream_arn)["StreamDescription"]["Shards"] - for shard in shards: - it = streams_client.get_shard_iterator( - StreamArn=stream_arn, ShardId=shard["ShardId"], ShardIteratorType="TRIM_HORIZON" - )["ShardIterator"] - for _ in range(20): - resp = streams_client.get_records(ShardIterator=it, Limit=1000) - records += [ - r for r in resp.get("Records", []) if r["dynamodb"]["Keys"]["pk"]["S"] == pk - ] - it = resp.get("NextShardIterator") - if not it or not resp.get("Records"): - break - if len(records) >= want: + while True: + by_pk = _stream_records(streams_client, arn) + if all(len(by_pk.get(pk, [])) >= WRITERS for pk in pks) or time.monotonic() > deadline: break time.sleep(1) - return records + for pk in pks: + records = by_pk.get(pk, []) + names = [r["eventName"] for r in records] + assert names == ["INSERT"] + ["MODIFY"] * (WRITERS - 1), (pk, names) + images = [ + (r["dynamodb"].get("OldImage", {}).get("w"), r["dynamodb"]["NewImage"]["w"]) + for r in records + ] + for (_, before), (old, _) in zip(images, images[1:]): + assert old == before, f"old image is not the record before it: {pk} {images}" def test_racing_puts_on_a_stream_table_all_succeed( dynamodb_client, raw_client, streams_client, stream_table ): + pks = [] for _ in range(ROUNDS): pk, out = _race(raw_client, stream_table, "hash", {}) _assert_all_succeed(out) - # The last round's stream: one INSERT, then a MODIFY per later put, each - # with the item it replaced as its old image. - arn = dynamodb_client.describe_table(TableName=stream_table)["Table"]["LatestStreamArn"] - records = _stream_events(streams_client, arn, pk, WRITERS) - names = sorted(r["eventName"] for r in records) - assert names == ["INSERT"] + ["MODIFY"] * (WRITERS - 1), names - olds = [r["dynamodb"]["OldImage"]["w"]["S"] for r in records if r["eventName"] == "MODIFY"] - news = [r["dynamodb"]["NewImage"]["w"]["S"] for r in records] - assert len(set(olds)) == WRITERS - 1 and set(olds) <= set(news), (olds, news) + pks.append(pk) + _assert_stream_chains(dynamodb_client, streams_client, stream_table, pks) + + +def test_racing_batch_puts_on_a_stream_table_all_succeed( + dynamodb_client, raw_client, streams_client, stream_table +): + """BatchWriteItem puts read the current item for the stream record too.""" + pks = [] + for _ in range(ROUNDS): + pk, out = _race(raw_client, stream_table, "hash", {}, batch=True) + _assert_all_succeed(out) + pks.append(pk) + _assert_stream_chains(dynamodb_client, streams_client, stream_table, pks) From 3b739539a88527b6b3f6e1fac81209fede9034e2 Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Mon, 5 Oct 2026 19:24:38 +0000 Subject: [PATCH 4/8] test: cover a put whose create race winner is deleted Two storage tests park the put's insert behind a trigger while an outside transaction commits a winner and a second one locks it. The winner is then deleted before the put re-reads it. In one test the put retries its insert and creates the item. In the other it loses five inserts and returns an internal error, with nothing written. The create race tests now wait for a backend that the creator blocks, not for any lock waiter. Assisted-by: pi claude-opus-5-5 --- .../storage-postgres/tests/put_create_race.rs | 213 +++++++++++++++++- 1 file changed, 206 insertions(+), 7 deletions(-) diff --git a/crates/storage-postgres/tests/put_create_race.rs b/crates/storage-postgres/tests/put_create_race.rs index 286a7857..7229470a 100644 --- a/crates/storage-postgres/tests/put_create_race.rs +++ b/crates/storage-postgres/tests/put_create_race.rs @@ -7,6 +7,8 @@ //! overwrites the item, after its own condition is checked against it. These //! tests hold the competing create open in an outside transaction, so the put //! deterministically loses the insert and has to order itself after the winner. +//! The last two tests also delete the winner before the put re-reads it, so the +//! put retries its insert, and they bound those retries. //! //! 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 @@ -29,7 +31,7 @@ use extenddb_storage::error::StorageError; use extenddb_storage::{DataEngine, TableEngine}; use extenddb_storage_postgres::{PostgresConfig, PostgresEngine}; use sqlx::postgres::PgPoolOptions; -use sqlx::{PgPool, Postgres, Transaction}; +use sqlx::{Connection, PgConnection, PgPool, Postgres, Transaction}; const ACCOUNT: &str = "123456789012"; const REGION: &str = "us-east-1"; @@ -225,14 +227,17 @@ async fn create_in_flight( (creator, pid) } -/// Wait until a backend other than `own_pid` waits on a lock. -async fn wait_for_lock_waiter(db: &PgPool, own_pid: i32) { +/// Wait until a backend waits for a lock of kind `event` (`transactionid` or +/// `advisory`) that the backend `blocker` holds. +async fn wait_until_blocked_by(db: &PgPool, blocker: i32, event: &str) { 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", + WHERE datname = current_database() AND wait_event_type = 'Lock' \ + AND wait_event = $2 AND $1 = ANY(pg_blocking_pids(pid))", ) - .bind(own_pid) + .bind(blocker) + .bind(event) .fetch_one(db) .await .expect("read pg_stat_activity"); @@ -241,7 +246,7 @@ async fn wait_for_lock_waiter(db: &PgPool, own_pid: i32) { } tokio::time::sleep(Duration::from_millis(10)).await; } - panic!("the put never waited on the create in flight"); + panic!("no backend waited on the {event} lock of backend {blocker}"); } /// Run `put` (ReturnValues ALL_OLD) while the winner's create is in flight, @@ -259,7 +264,7 @@ async fn put_after_winner( .engine .put_item(key_info, item(range, "loser"), true, condition, &maps, None); let commit = async { - wait_for_lock_waiter(&s.db, creator_pid).await; + wait_until_blocked_by(&s.db, creator_pid, "transactionid").await; creator.commit().await.expect("commit the create"); }; let (result, ()) = @@ -358,3 +363,197 @@ async fn a_lost_create_race_fails_a_condition_the_winner_breaks() { assert_eq!(stored_v(&s).await, "winner", "the failed put wrote nothing"); s.cleanup().await; } + +// The tests below take the arm where the winner is gone again when the put +// re-reads it. A BEFORE INSERT trigger parks the put's insert on an advisory +// lock, the gate. While the insert is parked, an outside transaction commits a +// winner, and a second one, the locker, locks it FOR UPDATE. When the gate +// opens, the insert loses at once: the winner is committed, and a row lock does +// not make an insert wait. The put's locking re-read then waits on the locker, +// which deletes the winner and commits, so the re-read returns no row. + +/// The advisory lock key of the gate. +const GATE: i64 = 0x5075_7452; + +/// Open a connection to the scratch database outside its pools. Returns it and +/// its backend pid. +async fn connect(s: &Scratch) -> (PgConnection, i32) { + let base = base_conn().expect("caller checks base_conn() first"); + let mut conn = PgConnection::connect(&format!("{base}/{}", s.db_name)) + .await + .expect("connect to the scratch database"); + let pid: i32 = sqlx::query_scalar("SELECT pg_backend_pid()") + .fetch_one(&mut conn) + .await + .expect("read the backend pid"); + (conn, pid) +} + +/// Make every insert into `table` pass the gate first, unless its transaction +/// sets `test.bypass_gate`. Passing takes the gate shared and releases it, so +/// an insert parks while a session holds the gate exclusively. +async fn install_gate(s: &Scratch, table: &str) { + sqlx::raw_sql(&format!( + "CREATE FUNCTION pass_gate() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN \ + IF current_setting('test.bypass_gate', true) = 'on' THEN RETURN NEW; END IF; \ + PERFORM pg_advisory_lock_shared({GATE}); \ + PERFORM pg_advisory_unlock_shared({GATE}); \ + RETURN NEW; \ + END $$; \ + CREATE TRIGGER pass_gate BEFORE INSERT ON {table} \ + FOR EACH ROW EXECUTE FUNCTION pass_gate();" + )) + .execute(&s.db) + .await + .expect("install the gate trigger"); +} + +/// Commit the item with `v` as its value, past the gate. +async fn commit_winner(s: &Scratch, table: &str, v: &str) { + let mut tx = s.db.begin().await.expect("begin the winner"); + sqlx::query("SET LOCAL test.bypass_gate = 'on'") + .execute(&mut *tx) + .await + .expect("bypass the gate"); + sqlx::query(&format!( + "INSERT INTO {table} (pk, item_data) VALUES ('c', $1)" + )) + .bind(serde_json::to_value(item(false, v)).expect("serialize the winner")) + .execute(&mut *tx) + .await + .expect("insert the winner"); + tx.commit().await.expect("commit the winner"); +} + +/// With the put's insert parked at the gate, make it lose to a winner `v` that +/// is deleted before the put re-reads it. With `rearm`, the gate closes again +/// behind the insert, so the put's next insert parks too. +async fn lose_to_a_vanishing_winner( + s: &Scratch, + table: &str, + gate: &mut PgConnection, + v: &str, + rearm: bool, +) { + commit_winner(s, table, v).await; + let (mut locker, locker_pid) = connect(s).await; + sqlx::query("BEGIN") + .execute(&mut locker) + .await + .expect("begin the locker"); + let locked: Option<(serde_json::Value,)> = sqlx::query_as(&format!( + "SELECT item_data FROM {table} WHERE pk = 'c' FOR UPDATE" + )) + .fetch_optional(&mut locker) + .await + .expect("lock the winner"); + assert!(locked.is_some(), "the locker holds the committed winner"); + sqlx::query("SELECT pg_advisory_unlock($1)") + .bind(GATE) + .execute(&mut *gate) + .await + .expect("open the gate"); + if rearm { + // Granted only after the parked insert has passed the gate. + sqlx::query("SELECT pg_advisory_lock($1)") + .bind(GATE) + .execute(&mut *gate) + .await + .expect("close the gate again"); + } + wait_until_blocked_by(&s.db, locker_pid, "transactionid").await; + sqlx::query(&format!("DELETE FROM {table} WHERE pk = 'c'")) + .execute(&mut locker) + .await + .expect("delete the winner"); + sqlx::query("COMMIT") + .execute(&mut locker) + .await + .expect("commit the delete"); +} + +/// Install the gate on a fresh hash table and close it. Returns the table's +/// key info, its data table, and the session that holds the gate. +async fn gated_table(s: &Scratch) -> (TableKeyInfo, String, PgConnection, i32) { + let key_info = table(s, false).await; + let data = data_table(&s.db).await; + install_gate(s, &data).await; + let (mut gate, gate_pid) = connect(s).await; + sqlx::query("SELECT pg_advisory_lock($1)") + .bind(GATE) + .execute(&mut gate) + .await + .expect("close the gate"); + (key_info, data, gate, gate_pid) +} + +#[tokio::test] +async fn a_put_whose_create_race_winner_is_deleted_creates_the_item() { + let test = "a_put_whose_create_race_winner_is_deleted_creates_the_item"; + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let (key_info, data, mut gate, gate_pid) = gated_table(&s).await; + let maps = ExpressionMaps::default(); + let put = s + .engine + .put_item(&key_info, item(false, "loser"), true, None, &maps, None); + let driver = async { + wait_until_blocked_by(&s.db, gate_pid, "advisory").await; + lose_to_a_vanishing_winner(&s, &data, &mut gate, "winner", false).await; + }; + let (result, ()) = + tokio::time::timeout(Duration::from_secs(30), async { tokio::join!(put, driver) }) + .await + .expect("the put finished"); + assert_eq!( + result.expect("the retried insert succeeds"), + None, + "the winner is gone, so there is no old image" + ); + assert_eq!(stored_v(&s).await, "loser"); + s.cleanup().await; +} + +#[tokio::test] +async fn a_put_that_keeps_losing_the_create_race_gives_up_after_five_inserts() { + let test = "a_put_that_keeps_losing_the_create_race_gives_up_after_five_inserts"; + if base_conn().is_none() { + return skip(test); + } + let s = scratch().await; + let (key_info, data, mut gate, gate_pid) = gated_table(&s).await; + let maps = ExpressionMaps::default(); + let put = s + .engine + .put_item(&key_info, item(false, "loser"), true, None, &maps, None); + let driver = async { + for k in 1..=4 { + wait_until_blocked_by(&s.db, gate_pid, "advisory").await; + lose_to_a_vanishing_winner(&s, &data, &mut gate, &format!("w{k}"), true).await; + } + // The fifth insert loses to a winner that stays: the put gives up. + wait_until_blocked_by(&s.db, gate_pid, "advisory").await; + commit_winner(&s, &data, "w5").await; + sqlx::query("SELECT pg_advisory_unlock($1)") + .bind(GATE) + .execute(&mut gate) + .await + .expect("open the gate"); + }; + let (result, ()) = + tokio::time::timeout(Duration::from_secs(30), async { tokio::join!(put, driver) }) + .await + .expect("the put finished"); + match result { + Err(StorageError::Internal(m)) => assert!(m.contains("after 5 attempts"), "{m}"), + other => panic!("expected Internal after 5 inserts, got {other:?}"), + } + assert_eq!( + stored_v(&s).await, + "w5", + "the put that gave up wrote nothing" + ); + s.cleanup().await; +} From cebbba079201fb74c5ef3daa74f54dd509177cfe Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Mon, 5 Oct 2026 19:24:38 +0000 Subject: [PATCH 5/8] docs: say how a write to a missing item resolves a create race The high-level design says that a missing item has no row to lock, and how a PutItem or UpdateItem that loses the insert orders itself after the winner. Two comments in the put code now match what the code does. Assisted-by: pi claude-opus-5-5 --- crates/storage-postgres/src/data/put_item.rs | 12 +++++++----- docs/design/02-high-level-design.md | 1 + 2 files changed, 8 insertions(+), 5 deletions(-) diff --git a/crates/storage-postgres/src/data/put_item.rs b/crates/storage-postgres/src/data/put_item.rs index f3a826bf..9ca34d8f 100755 --- a/crates/storage-postgres/src/data/put_item.rs +++ b/crates/storage-postgres/src/data/put_item.rs @@ -151,9 +151,11 @@ impl PostgresEngine { if result.rows_affected() == 1 { break; } - // Lost the create race. The insert waited for the winner to - // commit, so the locking read now returns its row: this put - // overwrites it, after its condition is checked against it. + // Lost the create race. The winner has committed (the insert + // waited for it if it was in flight). The locking read returns + // its row, and this put overwrites it after checking its + // condition against it. If a delete committed since, the read + // returns no row and the insert is retried. attempt += 1; if attempt >= MAX_CREATE_RACE_ATTEMPTS { return Err(create_race_exhausted(attempt)); @@ -482,8 +484,8 @@ impl PostgresEngine { /// bound as `UpdateItem`. const MAX_CREATE_RACE_ATTEMPTS: u32 = 5; -/// The error after `MAX_CREATE_RACE_ATTEMPTS` lost create races, as `UpdateItem` -/// reports it. +/// The error after `MAX_CREATE_RACE_ATTEMPTS` lost create races. Like the one +/// `UpdateItem` returns, it is an internal error (HTTP 500). fn create_race_exhausted(attempt: u32) -> StorageError { StorageError::Internal(format!( "PutItem could not create the item after {attempt} attempts: a concurrent writer \ diff --git a/docs/design/02-high-level-design.md b/docs/design/02-high-level-design.md index a08cf641..3a680e61 100755 --- a/docs/design/02-high-level-design.md +++ b/docs/design/02-high-level-design.md @@ -390,6 +390,7 @@ Read-modify-write operations (UpdateItem, PutItem with conditions, DeleteItem wi - **Atomicity:** The condition check and the write happen against the same snapshot. - **Serialization:** Concurrent updates to the same item are serialized by PostgreSQL's row lock, not by any in-memory mutex. - **No TOCTOU races:** Another request cannot modify the item between the condition check and the write. +- **Missing items:** A missing item has no row to lock, so `INSERT ... ON CONFLICT DO NOTHING` decides which writer creates it. A PutItem or UpdateItem that loses re-reads the winner `FOR UPDATE`, checks its condition against it, and writes after it. If the winner was deleted in the meantime, it retries the insert, up to 5 inserts in total. There is no in-memory locking (no `Mutex`, `RwLock`, or similar) on the data path. All contention is managed by PostgreSQL. From 1d380372cfddc4ec26f7b0d2b326aeeb12b04eb5 Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Mon, 5 Oct 2026 20:31:27 +0000 Subject: [PATCH 6/8] test: run the create race gate tests on a range table too The deleted-winner and give-up tests now run on a hash table and on a hash and range table, so the retry arm and the bound of the range branch have tests too. When a driver step fails, the driver opens the gate, and the test reports what the put returned instead of only that no backend waited. Assisted-by: pi claude-opus-5-5 --- .../storage-postgres/tests/put_create_race.rs | 157 +++++++++++++----- 1 file changed, 112 insertions(+), 45 deletions(-) diff --git a/crates/storage-postgres/tests/put_create_race.rs b/crates/storage-postgres/tests/put_create_race.rs index 7229470a..e7fc773e 100644 --- a/crates/storage-postgres/tests/put_create_race.rs +++ b/crates/storage-postgres/tests/put_create_race.rs @@ -228,8 +228,9 @@ async fn create_in_flight( } /// Wait until a backend waits for a lock of kind `event` (`transactionid` or -/// `advisory`) that the backend `blocker` holds. -async fn wait_until_blocked_by(db: &PgPool, blocker: i32, event: &str) { +/// `advisory`) that the backend `blocker` holds. Fails after 5 s, so the caller +/// can report what the put did instead. +async fn wait_until_blocked_by(db: &PgPool, blocker: i32, event: &str) -> Result<(), String> { for _ in 0..500 { let waiting: i64 = sqlx::query_scalar( "SELECT count(*) FROM pg_stat_activity \ @@ -242,11 +243,13 @@ async fn wait_until_blocked_by(db: &PgPool, blocker: i32, event: &str) { .await .expect("read pg_stat_activity"); if waiting > 0 { - return; + return Ok(()); } tokio::time::sleep(Duration::from_millis(10)).await; } - panic!("no backend waited on the {event} lock of backend {blocker}"); + Err(format!( + "no backend waited on the {event} lock of backend {blocker}" + )) } /// Run `put` (ReturnValues ALL_OLD) while the winner's create is in flight, @@ -264,13 +267,17 @@ async fn put_after_winner( .engine .put_item(key_info, item(range, "loser"), true, condition, &maps, None); let commit = async { - wait_until_blocked_by(&s.db, creator_pid, "transactionid").await; + let waited = wait_until_blocked_by(&s.db, creator_pid, "transactionid").await; creator.commit().await.expect("commit the create"); + waited }; - let (result, ()) = + let (result, waited) = tokio::time::timeout(Duration::from_secs(30), async { tokio::join!(put, commit) }) .await .expect("the put finished"); + if let Err(e) = waited { + panic!("{e}; the put returned {result:?}"); + } result } @@ -409,19 +416,22 @@ async fn install_gate(s: &Scratch, table: &str) { } /// Commit the item with `v` as its value, past the gate. -async fn commit_winner(s: &Scratch, table: &str, v: &str) { +async fn commit_winner(s: &Scratch, table: &str, range: bool, v: &str) { let mut tx = s.db.begin().await.expect("begin the winner"); sqlx::query("SET LOCAL test.bypass_gate = 'on'") .execute(&mut *tx) .await .expect("bypass the gate"); - sqlx::query(&format!( - "INSERT INTO {table} (pk, item_data) VALUES ('c', $1)" - )) - .bind(serde_json::to_value(item(false, v)).expect("serialize the winner")) - .execute(&mut *tx) - .await - .expect("insert the winner"); + let sql = if range { + format!("INSERT INTO {table} (pk, sk_s, item_data) VALUES ('c', '1', $1)") + } else { + format!("INSERT INTO {table} (pk, item_data) VALUES ('c', $1)") + }; + sqlx::query(&sql) + .bind(serde_json::to_value(item(range, v)).expect("serialize the winner")) + .execute(&mut *tx) + .await + .expect("insert the winner"); tx.commit().await.expect("commit the winner"); } @@ -431,11 +441,12 @@ async fn commit_winner(s: &Scratch, table: &str, v: &str) { async fn lose_to_a_vanishing_winner( s: &Scratch, table: &str, + range: bool, gate: &mut PgConnection, v: &str, rearm: bool, -) { - commit_winner(s, table, v).await; +) -> Result<(), String> { + commit_winner(s, table, range, v).await; let (mut locker, locker_pid) = connect(s).await; sqlx::query("BEGIN") .execute(&mut locker) @@ -461,7 +472,7 @@ async fn lose_to_a_vanishing_winner( .await .expect("close the gate again"); } - wait_until_blocked_by(&s.db, locker_pid, "transactionid").await; + wait_until_blocked_by(&s.db, locker_pid, "transactionid").await?; sqlx::query(&format!("DELETE FROM {table} WHERE pk = 'c'")) .execute(&mut locker) .await @@ -470,12 +481,14 @@ async fn lose_to_a_vanishing_winner( .execute(&mut locker) .await .expect("commit the delete"); + Ok(()) } -/// Install the gate on a fresh hash table and close it. Returns the table's -/// key info, its data table, and the session that holds the gate. -async fn gated_table(s: &Scratch) -> (TableKeyInfo, String, PgConnection, i32) { - let key_info = table(s, false).await; +/// Install the gate on a fresh table, keyed on `pk` and, when `range` is set, +/// also on `sk`, and close the gate. Returns the table's key info, its data +/// table, and the session that holds the gate. +async fn gated_table(s: &Scratch, range: bool) -> (TableKeyInfo, String, PgConnection, i32) { + let key_info = table(s, range).await; let data = data_table(&s.db).await; install_gate(s, &data).await; let (mut gate, gate_pid) = connect(s).await; @@ -487,26 +500,41 @@ async fn gated_table(s: &Scratch) -> (TableKeyInfo, String, PgConnection, i32) { (key_info, data, gate, gate_pid) } -#[tokio::test] -async fn a_put_whose_create_race_winner_is_deleted_creates_the_item() { - let test = "a_put_whose_create_race_winner_is_deleted_creates_the_item"; +/// Release every advisory lock `gate` holds. A driver that stops early calls +/// this, so that a parked insert goes on and the put can report its result. +async fn open_gate_fully(gate: &mut PgConnection) { + sqlx::query("SELECT pg_advisory_unlock_all()") + .execute(gate) + .await + .expect("open the gate"); +} + +async fn deleted_winner_is_retried(test: &str, range: bool) { if base_conn().is_none() { return skip(test); } let s = scratch().await; - let (key_info, data, mut gate, gate_pid) = gated_table(&s).await; + let (key_info, data, mut gate, gate_pid) = gated_table(&s, range).await; let maps = ExpressionMaps::default(); let put = s .engine - .put_item(&key_info, item(false, "loser"), true, None, &maps, None); + .put_item(&key_info, item(range, "loser"), true, None, &maps, None); let driver = async { - wait_until_blocked_by(&s.db, gate_pid, "advisory").await; - lose_to_a_vanishing_winner(&s, &data, &mut gate, "winner", false).await; + let steps = async { + wait_until_blocked_by(&s.db, gate_pid, "advisory").await?; + lose_to_a_vanishing_winner(&s, &data, range, &mut gate, "winner", false).await + } + .await; + open_gate_fully(&mut gate).await; + steps }; - let (result, ()) = + let (result, driven) = tokio::time::timeout(Duration::from_secs(30), async { tokio::join!(put, driver) }) .await .expect("the put finished"); + if let Err(e) = driven { + panic!("{e}; the put returned {result:?}"); + } assert_eq!( result.expect("the retried insert succeeds"), None, @@ -517,35 +545,56 @@ async fn a_put_whose_create_race_winner_is_deleted_creates_the_item() { } #[tokio::test] -async fn a_put_that_keeps_losing_the_create_race_gives_up_after_five_inserts() { - let test = "a_put_that_keeps_losing_the_create_race_gives_up_after_five_inserts"; +async fn a_put_whose_create_race_winner_is_deleted_creates_the_item() { + deleted_winner_is_retried( + "a_put_whose_create_race_winner_is_deleted_creates_the_item", + false, + ) + .await; +} + +#[tokio::test] +async fn a_put_whose_create_race_winner_is_deleted_creates_the_item_on_a_range_table() { + deleted_winner_is_retried( + "a_put_whose_create_race_winner_is_deleted_creates_the_item_on_a_range_table", + true, + ) + .await; +} + +async fn create_race_churn_gives_up(test: &str, range: bool) { if base_conn().is_none() { return skip(test); } let s = scratch().await; - let (key_info, data, mut gate, gate_pid) = gated_table(&s).await; + let (key_info, data, mut gate, gate_pid) = gated_table(&s, range).await; let maps = ExpressionMaps::default(); let put = s .engine - .put_item(&key_info, item(false, "loser"), true, None, &maps, None); + .put_item(&key_info, item(range, "loser"), true, None, &maps, None); let driver = async { - for k in 1..=4 { - wait_until_blocked_by(&s.db, gate_pid, "advisory").await; - lose_to_a_vanishing_winner(&s, &data, &mut gate, &format!("w{k}"), true).await; + let steps = async { + for k in 1..=4 { + wait_until_blocked_by(&s.db, gate_pid, "advisory").await?; + lose_to_a_vanishing_winner(&s, &data, range, &mut gate, &format!("w{k}"), true) + .await?; + } + // The fifth insert loses to a winner that stays: the put gives up. + wait_until_blocked_by(&s.db, gate_pid, "advisory").await?; + commit_winner(&s, &data, range, "w5").await; + Ok::<(), String>(()) } - // The fifth insert loses to a winner that stays: the put gives up. - wait_until_blocked_by(&s.db, gate_pid, "advisory").await; - commit_winner(&s, &data, "w5").await; - sqlx::query("SELECT pg_advisory_unlock($1)") - .bind(GATE) - .execute(&mut gate) - .await - .expect("open the gate"); + .await; + open_gate_fully(&mut gate).await; + steps }; - let (result, ()) = + let (result, driven) = tokio::time::timeout(Duration::from_secs(30), async { tokio::join!(put, driver) }) .await .expect("the put finished"); + if let Err(e) = driven { + panic!("{e}; the put returned {result:?}"); + } match result { Err(StorageError::Internal(m)) => assert!(m.contains("after 5 attempts"), "{m}"), other => panic!("expected Internal after 5 inserts, got {other:?}"), @@ -557,3 +606,21 @@ async fn a_put_that_keeps_losing_the_create_race_gives_up_after_five_inserts() { ); s.cleanup().await; } + +#[tokio::test] +async fn a_put_that_keeps_losing_the_create_race_gives_up_after_five_inserts() { + create_race_churn_gives_up( + "a_put_that_keeps_losing_the_create_race_gives_up_after_five_inserts", + false, + ) + .await; +} + +#[tokio::test] +async fn a_put_that_keeps_losing_the_create_race_gives_up_after_five_inserts_on_a_range_table() { + create_race_churn_gives_up( + "a_put_that_keeps_losing_the_create_race_gives_up_after_five_inserts_on_a_range_table", + true, + ) + .await; +} From a31108ec3fb1cf546225e901f6d8187440ea3077 Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Mon, 5 Oct 2026 20:39:03 +0000 Subject: [PATCH 7/8] test: say that four create race tests delete the winner The module doc of the storage tests still said two. Assisted-by: pi claude-opus-5-5 --- crates/storage-postgres/tests/put_create_race.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/crates/storage-postgres/tests/put_create_race.rs b/crates/storage-postgres/tests/put_create_race.rs index e7fc773e..b6915eee 100644 --- a/crates/storage-postgres/tests/put_create_race.rs +++ b/crates/storage-postgres/tests/put_create_race.rs @@ -7,8 +7,9 @@ //! overwrites the item, after its own condition is checked against it. These //! tests hold the competing create open in an outside transaction, so the put //! deterministically loses the insert and has to order itself after the winner. -//! The last two tests also delete the winner before the put re-reads it, so the -//! put retries its insert, and they bound those retries. +//! The last four tests, on hash and on hash and range tables, also delete the +//! winner before the put re-reads it, so the put retries its insert, and they +//! bound those retries. //! //! 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 From f5515a369be89468d77ffa898f34908ccbe079e8 Mon Sep 17 00:00:00 2001 From: Anandh Somasundaram Date: Mon, 5 Oct 2026 21:47:12 +0000 Subject: [PATCH 8/8] test: expect the plain-table create races to fail on MongoDB for now On MongoDB, a put with no condition on a table with no index and no stream returns HTTP 500 when it loses the create race, and the MongoDB write-race fix repairs it. That fix lands as its own PR. Until it is on main, the three affected tests are marked xfail when the MongoDB test runner runs them. The marker is not strict, so the tests also pass once the fix is in. Remove the marker then. Assisted-by: pi claude-opus-5-5 --- tests/test_put_item_create_race.py | 13 +++++++++++++ 1 file changed, 13 insertions(+) diff --git a/tests/test_put_item_create_race.py b/tests/test_put_item_create_race.py index cb573661..cfe22c30 100644 --- a/tests/test_put_item_create_race.py +++ b/tests/test_put_item_create_race.py @@ -32,6 +32,17 @@ ROUNDS = 15 +# TEMPORARY: on MongoDB, a put without a condition on a table with no index and +# no stream returns HTTP 500 when it loses the create race. The MongoDB +# write-race fix repairs it. That fix lands separately; remove this marker +# once it is on main. The MongoDB test runner sets EXTENDDB_TEST_MONGODB_CONTAINER. +XFAIL_UNTIL_MONGODB_FIX = pytest.mark.xfail( + bool(os.environ.get("EXTENDDB_TEST_MONGODB_CONTAINER", "").strip()), + reason="MongoDB returns HTTP 500 for a lost create race, fixed by the MongoDB write-race fix", + strict=False, +) + + @pytest.fixture(scope="module") def raw_client(endpoint_url): """A client that never retries, so every failure is seen.""" @@ -209,6 +220,7 @@ def _assert_one_chain(out: list, final: str): assert seen == WRITERS, f"old images do not form one chain: {prev}, final {final}" +@XFAIL_UNTIL_MONGODB_FIX @pytest.mark.parametrize("table_kind", ["hash", "range"]) def test_racing_puts_with_all_old_all_succeed( request, dynamodb_client, raw_client, table_kind @@ -220,6 +232,7 @@ def test_racing_puts_with_all_old_all_succeed( _assert_one_chain(out, _final_w(dynamodb_client, table, table_kind, pk)) +@XFAIL_UNTIL_MONGODB_FIX def test_racing_puts_that_return_consumed_capacity_all_succeed( dynamodb_client, raw_client, hash_table ):