Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 13 additions & 7 deletions crates/storage-mongodb/src/table_engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
)
}

Expand Down
14 changes: 9 additions & 5 deletions crates/storage-postgres/src/stream_engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,7 @@ impl PostgresEngine {
table_id: &str,
) -> Result<String, StorageError> {
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",
)
Expand All @@ -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) \
Expand Down
2 changes: 1 addition & 1 deletion crates/storage-postgres/src/update_table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
13 changes: 8 additions & 5 deletions crates/storage-sqlite/src/stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,7 +34,7 @@ impl SqliteEngine {
table_id: &str,
) -> Result<String, StorageError> {
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)
Expand All @@ -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) \
Expand Down
2 changes: 1 addition & 1 deletion crates/storage-sqlite/src/update_table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
11 changes: 7 additions & 4 deletions docs/design/13-storage-mongodb.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion docs/manuals/02-design-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -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-<table>-0000000000000000` through `shardId-<table>-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-<table_id>-0000000000000000` through `shardId-<table_id>-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-<table>-<n>` IDs. Sequence numbers are monotonically increasing integers.

### Iterator Types

Expand Down
164 changes: 164 additions & 0 deletions tests/test_stream_table_reuse.py
Original file line number Diff line number Diff line change
@@ -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)
Loading