diff --git a/crates/storage-mongodb/src/table_engine.rs b/crates/storage-mongodb/src/table_engine.rs index 87a5df20..a579a239 100644 --- a/crates/storage-mongodb/src/table_engine.rs +++ b/crates/storage-mongodb/src/table_engine.rs @@ -26,23 +26,29 @@ use crate::MongoEngine; use crate::data::data_collection_name; /// Format a timestamp as a DynamoDB-style stream label: -/// `YYYY-MM-DDThh:mm:ss` (second precision, no timezone). +/// `YYYY-MM-DDThh:mm:ss.sss` (millisecond precision, no timezone), the +/// shape the service uses. /// /// Matches the postgres backend's -/// `to_char(NOW(), 'YYYY-MM-DD"T"HH24:MI:SS')` output byte-for-byte -/// so a stream ARN issued by one backend is parseable by tooling that -/// only ever saw the other. The `time` crate's `Iso8601::DEFAULT` -/// emits nanoseconds with a trailing `Z` — pushing that through AWS- -/// SDK parsers or postgres-shaped tests failed unpredictably. D-m8. +/// `to_char(clock_timestamp(), 'YYYY-MM-DD"T"HH24:MI:SS.MS')` output +/// byte-for-byte so a stream ARN issued by one backend is parseable by +/// tooling that only ever saw the other. The `time` crate's +/// `Iso8601::DEFAULT` emits nanoseconds with a trailing `Z`; pushing that +/// through AWS-SDK parsers or postgres-shaped tests failed unpredictably. +/// D-m8. Milliseconds rather than seconds because the label is what tells +/// one stream on a table name from the next: a table deleted and recreated +/// within the same second would otherwise get the same stream ARN, and the +/// old ARN would resolve to the new table's stream. fn format_stream_label(now: time::OffsetDateTime) -> String { format!( - "{:04}-{:02}-{:02}T{:02}:{:02}:{:02}", + "{:04}-{:02}-{:02}T{:02}:{:02}:{:02}.{:03}", now.year(), u8::from(now.month()), now.day(), now.hour(), now.minute(), now.second(), + now.millisecond(), ) } diff --git a/crates/storage-postgres/src/stream_engine.rs b/crates/storage-postgres/src/stream_engine.rs index 151a9018..d159a462 100755 --- a/crates/storage-postgres/src/stream_engine.rs +++ b/crates/storage-postgres/src/stream_engine.rs @@ -36,7 +36,7 @@ impl PostgresEngine { table_id: &str, ) -> Result { let label: String = sqlx::query_scalar( - "UPDATE tables SET stream_label = to_char(NOW(), 'YYYY-MM-DD\"T\"HH24:MI:SS') \ + "UPDATE tables SET stream_label = to_char(clock_timestamp(), 'YYYY-MM-DD\"T\"HH24:MI:SS.MS') \ WHERE account_id = $1 AND table_name = $2 \ RETURNING stream_label", ) @@ -55,10 +55,14 @@ impl PostgresEngine { .map_err(|e| StorageError::Internal(e.to_string()))?; for i in 0..SHARDS_PER_STREAM { - // Zero-padded to 16 digits so the shard ID is always at - // least 28 characters (minimum length the AWS SDKs enforce for ShardId) - // even for the shortest legal table name. - let shard_id = format!("shardId-{table_name}-{i:016}"); + // Derived from the table id, not its name. A name is reused by + // delete-and-recreate and by every account that picks it, while + // `stream_shards.shard_id` is unique across the data database + // and PostgreSQL keeps a deleted table's shards. A name of + // up to 255 bytes also overflowed the 65-character ShardId limit + // the AWS SDKs enforce. The table id is a UUID, giving a + // 61-character id. + let shard_id = format!("shardId-{table_id}-{i:016}"); let start_seq = format!("{:021}", 0); sqlx::query( "INSERT INTO stream_shards (shard_id, table_id, starting_sequence_number) \ diff --git a/crates/storage-postgres/src/update_table.rs b/crates/storage-postgres/src/update_table.rs index 3cbdada3..c71ab763 100755 --- a/crates/storage-postgres/src/update_table.rs +++ b/crates/storage-postgres/src/update_table.rs @@ -363,7 +363,7 @@ impl PostgresEngine { if current_label.is_none() { sqlx::query( "UPDATE tables SET stream_label = \ - to_char(NOW(), 'YYYY-MM-DD\"T\"HH24:MI:SS') \ + to_char(clock_timestamp(), 'YYYY-MM-DD\"T\"HH24:MI:SS.MS') \ WHERE account_id = $1 AND table_name = $2", ) .bind(account_id) diff --git a/crates/storage-sqlite/src/stream.rs b/crates/storage-sqlite/src/stream.rs index 3d8850e8..c2befe73 100644 --- a/crates/storage-sqlite/src/stream.rs +++ b/crates/storage-sqlite/src/stream.rs @@ -34,7 +34,7 @@ impl SqliteEngine { table_id: &str, ) -> Result { let label: String = sqlx::query_scalar( - "UPDATE tables SET stream_label = strftime('%Y-%m-%dT%H:%M:%S','now') \ + "UPDATE tables SET stream_label = strftime('%Y-%m-%dT%H:%M:%f','now') \ WHERE account_id = ? AND table_name = ? RETURNING stream_label", ) .bind(account_id) @@ -44,10 +44,13 @@ impl SqliteEngine { .map_err(|e| StorageError::Internal(e.to_string()))?; for i in 0..SHARDS_PER_STREAM { - // Zero-padded to 16 digits so the shard ID is always at - // least 28 characters (minimum length the AWS SDKs enforce for ShardId) - // even for the shortest legal table name. - let shard_id = format!("shardId-{table_name}-{i:016}"); + // Derived from the table id, not its name. A name is shared by + // every account that picks it, while `stream_shards.shard_id` is + // unique across the database. A name of up to 255 bytes also + // overflowed the 65-character ShardId limit + // the AWS SDKs enforce. The table id is a UUID, giving a + // 61-character id. + let shard_id = format!("shardId-{table_id}-{i:016}"); let start_seq = format!("{:021}", 0); sqlx::query( "INSERT INTO stream_shards (shard_id, table_id, starting_sequence_number) \ diff --git a/crates/storage-sqlite/src/update_table.rs b/crates/storage-sqlite/src/update_table.rs index 821bda1d..654dd50c 100644 --- a/crates/storage-sqlite/src/update_table.rs +++ b/crates/storage-sqlite/src/update_table.rs @@ -293,7 +293,7 @@ impl SqliteEngine { .map_err(|e| StorageError::Internal(e.to_string()))?; if label.is_none() { sqlx::query( - "UPDATE tables SET stream_label = strftime('%Y-%m-%dT%H:%M:%S','now') \ + "UPDATE tables SET stream_label = strftime('%Y-%m-%dT%H:%M:%f','now') \ WHERE account_id = ? AND table_name = ?", ) .bind(account_id) diff --git a/docs/design/13-storage-mongodb.md b/docs/design/13-storage-mongodb.md index 3176ae08..2f55cfa2 100644 --- a/docs/design/13-storage-mongodb.md +++ b/docs/design/13-storage-mongodb.md @@ -720,10 +720,13 @@ only a first-time enable rotates it. A repeat `UpdateTable` `{ StreamEnabled: true }` would otherwise duplicate the shard set and invalidate stream ARNs previously handed out to consumers. -**`stream_label` format.** `YYYY-MM-DDThh:mm:ss` (second precision, no -timezone). Byte-for-byte compatible with the PostgreSQL backend so an -ARN issued by one backend is parseable by tooling that only ever saw -the other. See `format_stream_label` in `table_engine.rs`. +**`stream_label` format.** `YYYY-MM-DDThh:mm:ss.sss` (millisecond +precision, no timezone), the shape the service uses. Byte-for-byte +compatible with the PostgreSQL and SQLite backends so an ARN issued by +one backend is parseable by tooling that only ever saw another. Second +precision gave a table deleted and recreated within one second the same +stream ARN. Labels written before the change keep their second-precision +form. See `format_stream_label` in `table_engine.rs`. ### 5.6 Write conflict handling diff --git a/docs/manuals/02-design-guide.md b/docs/manuals/02-design-guide.md index 13cd6c9d..c1904261 100755 --- a/docs/manuals/02-design-guide.md +++ b/docs/manuals/02-design-guide.md @@ -138,7 +138,7 @@ For UpdateItem, the `new_image` is not known until after `apply_update` runs ins ### Shard Model -Each stream has a fixed set of shards (currently 4 shards per stream). Shard IDs are deterministic (`shardId--0000000000000000` through `shardId-
-0000000000000003`), zero-padded to 16 digits so every shard ID meets the AWS SDKs' 28-character minimum for `ShardId`, regardless of table name length. Sequence numbers are monotonically increasing integers. +Each stream has a fixed set of shards (currently 4 shards per stream). On PostgreSQL and SQLite, shard IDs embed the table's id, a UUID assigned at creation (`shardId--0000000000000000` through `shardId--0000000000000003`), so a new shard ID is 61 characters, inside the AWS SDKs' 28-to-65-character bounds for `ShardId`. The table name is not used because it is reused by delete-and-recreate and by every account that picks it, while shard IDs must be unique across the data database. Shards created by earlier versions keep their `shardId-
-` IDs. Sequence numbers are monotonically increasing integers. ### Iterator Types diff --git a/tests/test_stream_table_reuse.py b/tests/test_stream_table_reuse.py new file mode 100644 index 00000000..2563a894 --- /dev/null +++ b/tests/test_stream_table_reuse.py @@ -0,0 +1,164 @@ +# Copyright 2026 ExtendDB contributors +# SPDX-License-Identifier: Apache-2.0 + +"""A stream-enabled table name can be reused. + +Shard ids used to be derived from the table name, while the shard store is +keyed by shard id alone (and on PostgreSQL keeps a deleted table's shards). +A second stream-enabled table under the same name then failed CreateTable +with InternalServerError: a delete-and-recreate in one account on +PostgreSQL, or the same name chosen by a second account on PostgreSQL and +SQLite. A long table name also produced +a ShardId over the 65-character limit the AWS SDKs enforce. + +ExtendDB only: Amazon DynamoDB serves streams on a separate endpoint, and the +two-account case needs the management API. +""" + +from __future__ import annotations + +import os +import uuid + +import boto3 +import pytest +from botocore.exceptions import ClientError + +from conftest import wait_for_active, wait_for_deleted +from test_idempotency_account_scope import two_account_clients # noqa: F401 - fixture + +pytestmark = pytest.mark.skipif( + not os.environ.get("EXTENDDB_TEST_ENDPOINT", "").strip(), + reason="streams are served on the same endpoint by ExtendDB only", +) + +STREAM_SPEC = {"StreamEnabled": True, "StreamViewType": "NEW_IMAGE"} +# The AWS SDK model bounds ShardId to 28..65 characters. +SHARD_ID_MIN, SHARD_ID_MAX = 28, 65 + + +def _streams_client_for(ddb_client): + """A dynamodbstreams client sharing the given client's endpoint and keys.""" + creds = ddb_client._request_signer._credentials # noqa: SLF001 - test helper + endpoint = ddb_client.meta.endpoint_url + kwargs: dict = dict( + service_name="dynamodbstreams", + region_name=ddb_client.meta.region_name, + endpoint_url=endpoint, + aws_access_key_id=creds.access_key, + aws_secret_access_key=creds.secret_key, + aws_session_token=creds.token, + ) + if endpoint.startswith("https://"): + kwargs["verify"] = False + return boto3.client(**kwargs) + + +def _create_stream_table(client, name: str) -> str: + client.create_table( + TableName=name, + AttributeDefinitions=[{"AttributeName": "pk", "AttributeType": "S"}], + KeySchema=[{"AttributeName": "pk", "KeyType": "HASH"}], + BillingMode="PAY_PER_REQUEST", + StreamSpecification=STREAM_SPEC, + ) + wait_for_active(client, name) + return client.describe_table(TableName=name)["Table"]["LatestStreamArn"] + + +def _drop(client, name: str) -> None: + try: + client.delete_table(TableName=name) + except client.exceptions.ResourceNotFoundException: + return + wait_for_deleted(client, name) + + +def _stream_pks(streams, stream_arn: str) -> list[str]: + """Partition keys of every record in the stream, all shards, from the start.""" + desc = streams.describe_stream(StreamArn=stream_arn)["StreamDescription"] + pks: list[str] = [] + for shard in desc["Shards"]: + it = streams.get_shard_iterator( + StreamArn=stream_arn, ShardId=shard["ShardId"], ShardIteratorType="TRIM_HORIZON" + )["ShardIterator"] + for _ in range(20): + resp = streams.get_records(ShardIterator=it) + pks += [r["dynamodb"]["Keys"]["pk"]["S"] for r in resp["Records"]] + it = resp.get("NextShardIterator") + if not it or not resp["Records"]: + break + return pks + + +def test_recreate_stream_table_under_same_name(dynamodb_client): + streams = _streams_client_for(dynamodb_client) + name = f"stream-reuse-{uuid.uuid4().hex[:8]}" + try: + first_arn = _create_stream_table(dynamodb_client, name) + dynamodb_client.put_item(TableName=name, Item={"pk": {"S": "first"}}) + _drop(dynamodb_client, name) + + second_arn = _create_stream_table(dynamodb_client, name) + assert second_arn != first_arn + dynamodb_client.put_item(TableName=name, Item={"pk": {"S": "second"}}) + + # The new stream carries only the new table's records, and the old + # ARN does not resolve to it. + assert _stream_pks(streams, second_arn) == ["second"] + with pytest.raises(ClientError) as exc: + streams.describe_stream(StreamArn=first_arn) + assert exc.value.response["Error"]["Code"] == "ResourceNotFoundException" + finally: + _drop(dynamodb_client, name) + + +def test_two_accounts_stream_tables_with_one_name(two_account_clients): # noqa: F811 + client_a, client_b = two_account_clients + name = f"stream-shared-{uuid.uuid4().hex[:8]}" + try: + arn_a = _create_stream_table(client_a, name) + arn_b = _create_stream_table(client_b, name) + client_a.put_item(TableName=name, Item={"pk": {"S": "from-a"}}) + client_b.put_item(TableName=name, Item={"pk": {"S": "from-b"}}) + + assert _stream_pks(_streams_client_for(client_a), arn_a) == ["from-a"] + assert _stream_pks(_streams_client_for(client_b), arn_b) == ["from-b"] + finally: + _drop(client_a, name) + _drop(client_b, name) + + +def test_shard_ids_fit_sdk_bounds_for_a_long_table_name(dynamodb_client): + streams = _streams_client_for(dynamodb_client) + name = f"stream-long-{uuid.uuid4().hex[:8]}-" + "x" * 200 + try: + arn = _create_stream_table(dynamodb_client, name) + shards = streams.describe_stream(StreamArn=arn)["StreamDescription"]["Shards"] + assert shards + for shard in shards: + assert SHARD_ID_MIN <= len(shard["ShardId"]) <= SHARD_ID_MAX, shard["ShardId"] + # boto3 validates ShardId length client-side, so this call is the + # end-to-end check that the id is usable. + dynamodb_client.put_item(TableName=name, Item={"pk": {"S": "k"}}) + assert _stream_pks(streams, arn) == ["k"] + finally: + _drop(dynamodb_client, name) + + +def test_stream_disable_and_reenable(dynamodb_client): + streams = _streams_client_for(dynamodb_client) + name = f"stream-toggle-{uuid.uuid4().hex[:8]}" + try: + _create_stream_table(dynamodb_client, name) + dynamodb_client.update_table( + TableName=name, StreamSpecification={"StreamEnabled": False} + ) + wait_for_active(dynamodb_client, name) + dynamodb_client.update_table(TableName=name, StreamSpecification=STREAM_SPEC) + wait_for_active(dynamodb_client, name) + arn = dynamodb_client.describe_table(TableName=name)["Table"]["LatestStreamArn"] + dynamodb_client.put_item(TableName=name, Item={"pk": {"S": "after"}}) + assert "after" in _stream_pks(streams, arn) + finally: + _drop(dynamodb_client, name)