From ed9910435443fd9337dbdee8d03d0658a7a55a65 Mon Sep 17 00:00:00 2001 From: Scott Robinson Date: Sun, 4 Oct 2026 23:00:40 +0000 Subject: [PATCH] fix(postgres): read TransactGetItems from one snapshot TransactGetItems ran its reads in a READ COMMITTED transaction, where each SELECT takes a fresh snapshot. A TransactWriteItems that committed between two of those reads was observed half-applied: with one writer moving two items to the same generation and four readers, 27 of 3,578 reads returned the items at different generations. Begin the read transaction as REPEATABLE READ READ ONLY so every read is served from the snapshot taken at the first one. On a primary, a REPEATABLE READ transaction that only reads cannot fail with a serialization error (PostgreSQL raises those only when such a transaction modifies a row changed since its snapshot), so no retry path is added. tests/test_transact_get_snapshot.py runs the writer-and-readers probe for five seconds and fails on any torn read; it fails on main and passes with this change. Readiness assessment P0-3. Signed-off-by: Scott Robinson --- .../storage-postgres/src/data/transactions.rs | 13 ++- tests/test_transact_get_snapshot.py | 92 +++++++++++++++++++ 2 files changed, 104 insertions(+), 1 deletion(-) create mode 100644 tests/test_transact_get_snapshot.py diff --git a/crates/storage-postgres/src/data/transactions.rs b/crates/storage-postgres/src/data/transactions.rs index d69f90043..2c0f7a4b7 100644 --- a/crates/storage-postgres/src/data/transactions.rs +++ b/crates/storage-postgres/src/data/transactions.rs @@ -53,9 +53,20 @@ impl PostgresEngine { return Err(StorageError::TransactionCanceled(reasons)); } + // TransactGetItems promises that every item comes from one point in + // time. Under the pool's default READ COMMITTED, each SELECT takes its + // own snapshot, so a TransactWriteItems that commits between two + // reads is observed half-applied. REPEATABLE READ takes one snapshot + // at the first read and serves every read from it. On a primary, a + // REPEATABLE READ transaction that only reads cannot fail with a + // serialization error, because those are raised only when it tries to + // modify a row changed since its snapshot, so no retry is added. Every + // write in this backend goes through this same pool, so it cannot be + // a hot standby in a deployment that accepts writes; on a standby a + // recovery conflict could cancel this read, as it could any query. let mut tx = self .data_pool - .begin() + .begin_with("BEGIN ISOLATION LEVEL REPEATABLE READ READ ONLY") .await .map_err(|e| StorageError::Internal(e.to_string()))?; diff --git a/tests/test_transact_get_snapshot.py b/tests/test_transact_get_snapshot.py new file mode 100644 index 000000000..94e0e2ef7 --- /dev/null +++ b/tests/test_transact_get_snapshot.py @@ -0,0 +1,92 @@ +# Copyright 2026 ExtendDB contributors +# SPDX-License-Identifier: Apache-2.0 + +"""TransactGetItems reads every item from one point in time. + +One writer moves two items to the same generation inside a single +TransactWriteItems while readers fetch both with TransactGetItems. Every +successful read must see both items at the same generation; a mismatch is a +torn read (half of a committed transaction observed). + +DynamoDB may cancel a TransactGetItems that conflicts with an in-flight write +(TransactionCanceledException); that is a valid outcome and is not counted. +""" + +from __future__ import annotations + +import threading +import time + +from botocore.exceptions import ClientError + +from conftest import scoped_table +from test_concurrency import _make_client + +DURATION_S = 5.0 +READERS = 4 +KEYS = ("snap-a", "snap-b") + + +def test_transact_get_items_never_returns_a_torn_read(dynamodb_client): + with scoped_table(dynamodb_client) as table: + for k in KEYS: + dynamodb_client.put_item(TableName=table, Item={"pk": {"S": k}, "gen": {"N": "0"}}) + + deadline = time.monotonic() + DURATION_S + torn: list[list[str]] = [] + reads = [0] + writes = [0] + errors: list[BaseException] = [] + lock = threading.Lock() + + def writer() -> None: + client = _make_client() + gen = 0 + try: + while time.monotonic() < deadline: + gen += 1 + client.transact_write_items( + TransactItems=[ + {"Put": {"TableName": table, "Item": {"pk": {"S": k}, "gen": {"N": str(gen)}}}} + for k in KEYS + ] + ) + with lock: + writes[0] += 1 + except BaseException as e: # noqa: BLE001 - surfaced below + errors.append(e) + + def reader() -> None: + client = _make_client() + try: + while time.monotonic() < deadline: + try: + resp = client.transact_get_items( + TransactItems=[ + {"Get": {"TableName": table, "Key": {"pk": {"S": k}}}} for k in KEYS + ] + ) + except ClientError as e: + if e.response["Error"]["Code"] == "TransactionCanceledException": + continue + raise + gens = [r["Item"]["gen"]["N"] for r in resp["Responses"]] + with lock: + reads[0] += 1 + if len(set(gens)) != 1: + torn.append(gens) + except BaseException as e: # noqa: BLE001 - surfaced below + errors.append(e) + + threads = [threading.Thread(target=writer)] + [ + threading.Thread(target=reader) for _ in range(READERS) + ] + for t in threads: + t.start() + for t in threads: + t.join() + + assert not errors, errors[0] + assert writes[0] > 10, f"writer made too little progress ({writes[0]} transactions)" + assert reads[0] > 10, f"readers made too little progress ({reads[0]} reads)" + assert not torn, f"{len(torn)} of {reads[0]} TransactGetItems were torn, e.g. {torn[:5]}"