diff --git a/crates/storage-postgres/src/data/transactions.rs b/crates/storage-postgres/src/data/transactions.rs index d69f9004..2c0f7a4b 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 00000000..94e0e2ef --- /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]}"