Skip to content

fix(streams): derive shard ids from the table id, label streams to the millisecond, sweep deleted tables' shards - #383

Open
robinnsc wants to merge 1 commit into
mainfrom
fix/pg-stream-shard-ids
Open

robinnsc wants to merge 1 commit into
mainfrom
fix/pg-stream-shard-ids

Conversation

@robinnsc

@robinnsc robinnsc commented Oct 5, 2026 •

Copy link
Copy Markdown
Collaborator

What

Stream shard ids on PostgreSQL and SQLite are built from the table id instead of the table name, stream labels on all three backends carry milliseconds, and PostgreSQL removes a deleted table's shards once its stream has aged out.

  • crates/storage-postgres/src/stream_engine.rs, crates/storage-sqlite/src/stream.rs: shardId-{table_id}-{i:016} in place of shardId-{table_name}-{i:016}. The table id is a UUID, so new shard ids are 61 characters. MongoDB already builds its shard ids this way.
  • stream_label is now YYYY-MM-DDThh:mm:ss.sss, the shape the service uses. It is formatted in one place, extenddb_storage::util::format_stream_label (UTC, three fractional digits, unit-tested), and bound by every backend: PostgreSQL (stream_engine.rs, update_table.rs) and SQLite (stream.rs, update_table.rs) previously formatted it in SQL with to_char and strftime, MongoDB in a private Rust function; the three produced identical bytes but only tests kept them so.
  • PostgreSQL cleanup_expired_stream_records (the hourly retention sweep) now also removes the shards of a table that is gone from the catalog once none of its records remain. Candidates are shards older than the retention window whose table has no records; a freshly created table is never a candidate because its shards are younger than the window, so the gap between inserting shards and committing the catalog row cannot be mistaken for a deletion, and a live table's unused shards are left alone. SQLite already removed both with the table.
  • Existing shard rows and labels are not rewritten. They keep their ids and ARNs and keep working, because shard ids and labels are looked up as opaque values. Shards created from now on get the new form, including on a pre-existing table that enables a stream for the first time.
  • The PostgreSQL storage-level test harnesses (key_collation.rs, vector_control_plane.rs, and the new stream_shard_retention.rs) build their scratch database by applying every .sql file under migrations/ and data_migrations/ in filename order, instead of each carrying a copied list of the files.
  • Docs: docs/manuals/02-design-guide.md (shard id format, the sweep, the remaining label window), docs/design/13-storage-mongodb.md (label format).

Why

No issue filed. Found in the 1.0 readiness review (P0-5: scope shard ids by table id, and clean up shards on delete or at the retention horizon).

stream_shards is keyed by shard id alone, and the id contained only the table name, so any second stream-enabled table under one name collided:

  • PostgreSQL keeps a deleted table's shard rows. DeleteTable followed by CreateTable with the same name and a stream failed with InternalServerError every time (duplicate key value violates unique constraint "stream_shards_pkey"). Delete-and-recreate under one name is routine in test harnesses and works in DynamoDB.
  • On PostgreSQL and SQLite, a second account creating a stream-enabled table with a name another account already uses failed the same way.
  • A table name longer than about 40 characters gave a ShardId over the 65 characters the AWS SDKs accept, so the SDK rejected DescribeStream and GetShardIterator calls carrying it client-side. [Bug] GetShardIterator/GetRecords rejected client-side by AWS SDKs for tables named 3–6 characters long #247 fixed the lower bound for short names; this fixes the upper one.

With recreate working, a second defect became reachable. Labels had one-second resolution, so a table deleted and recreated within the same second (control_plane_delay_seconds of 0 makes this easy) got the same stream ARN as the old table, and DescribeStream on the old ARN resolved to the new table's stream. On MongoDB, where recreate already worked, the new test hit this in 1 of 3 runs on main.

And PostgreSQL never removed a deleted table's shard rows at all: records were trimmed by the retention sweep, shards were not, so every deleted stream-enabled table left four rows behind for good.

Not closed here

The stream ARN still resolves by table name and label, so a table deleted and recreated within the same millisecond would still give the old ARN the new stream. Milliseconds narrow the window from the second-resolution one; closing it needs the table id in the ARN or label, which changes the ARN shape and should be its own change.

Testing done

New tests/test_stream_table_reuse.py:

  • recreate under the same name: the new stream ARN differs, the new stream carries only the new table's records, and the old ARN returns ResourceNotFoundException;
  • two accounts, one table name: each stream carries only its own account's records;
  • a 221-character table name: every ShardId is within 28..65 and GetShardIterator and GetRecords work through boto3;
  • disable then re-enable a stream on one table: records written after re-enabling are delivered.

On main, the first three fail on PostgreSQL and the second and third fail on SQLite. On this branch all four pass on PostgreSQL (control_plane_delay_seconds 0, 0.05, and 0.25), SQLite, and MongoDB, together with test_streams.py (20 passed, 1 xfailed on each).

New crates/storage-postgres/tests/stream_shard_retention.rs (added to the PostgreSQL storage-level CI step): shards are kept after DeleteTable and while a record of the table remains inside the window, removed once the records are gone, and a live table's shards are untouched even with no records.

The file is ExtendDB-only (module skip when EXTENDDB_TEST_ENDPOINT is unset), like test_streams.py: it signs Streams requests against the ExtendDB endpoint, and the two-account case needs the management API.

Full PostgreSQL pytest run (tests/, import/export excluded as in CI) and the comprehensive suite (tests/python), against servers built from this branch and from main: no test fails on the branch that passes on main, and the comprehensive suite passes 331 of 331 on both.

cargo fmt --all -- --check
cargo +1.97.0 clippy --workspace --all-targets -- -D warnings   # the CI toolchain
cargo test --workspace
EXTENDDB_TEST_PG_CONNECTION_STRING=... cargo test -p extenddb-storage-postgres --test stream_shard_retention
cargo +1.88.0 check --workspace --locked

Checklist

  • I have read CONTRIBUTING.md
  • All tests pass (cargo test --workspace)
  • Code is formatted (cargo fmt --check)
  • Clippy is clean (cargo clippy -- -W clippy::pedantic)
  • I have added or updated tests for new functionality
  • I have updated documentation if behavior changed
  • Breaking changes are noted below (if any)
  • If this changes the wire protocol, Storage trait, auth model, on-disk format, or public CLI surface, an RFC has been accepted or is linked below. Otherwise, an ADR captures the decision (link below).

ADR / RFC: n/a. Shard ids and stream labels are opaque server-issued values; existing rows are left as they are, and no trait, schema, or CLI change.

Breaking changes

New streams get ShardIds and StreamLabels in a different shape. Both are opaque in the DynamoDB API, and the new label shape is the service's own. A client that parsed the table name out of a ShardId, or expected a label with no fractional seconds, would see the difference. Existing streams are unchanged. On PostgreSQL, a deleted table's shards are now removed once its records have aged out of the 24-hour window.


By submitting this pull request, I confirm that my contribution is made under the terms of the Apache License 2.0 and I agree to the Developer Certificate of Origin (DCO). See CONTRIBUTING.md for details.

@robinnsc
robinnsc force-pushed the fix/pg-stream-shard-ids branch from 82362ed to 57e1438 Compare October 6, 2026 07:43
@robinnsc robinnsc changed the title fix(streams): derive shard ids from the table id and label streams to the millisecond fix(streams): derive shard ids from the table id, label streams to the millisecond, sweep deleted tables' shards Oct 6, 2026
fn format_stream_label(now: time::OffsetDateTime) -> String {
format!(
"{:04}-{:02}-{:02}T{:02}:{:02}:{:02}",
"{:04}-{:02}-{:02}T{:02}:{:02}:{:02}.{:03}",

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Not blocking, but is this something that should/could be pulled up out of the backend layers? It'd be good to avoid this incompatibility in one place only.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yeah thats a good call, label is now formatted in one place (extenddb_storage::util::format_stream_label)

.await
.expect("connect to the scratch database");
for sql in [
include_str!("../migrations/001_schema.sql"),

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This could be a bit painful to keep up to date. Is it worth considering applying any *.sql files in either migrations/ or data_migrations/ , rather than have to remember to update this list every time a new migration file is added?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Done

Comment thread crates/storage-sqlite/src/stream.rs Outdated
) -> 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') \

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Doesn't this result in a different stream_label format than the one produced for MongoDB?

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

strftime('%f') gives SS.SSS, to_char('.MS') gives three digits, and Mongo padded to three, so all three produced identical bytes. But I agree, was unclear, and with the formatter change there's now just one implementation

@jcshepherd jcshepherd left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

A couple questions/comments below ...

@robinnsc
robinnsc force-pushed the fix/pg-stream-shard-ids branch from 57e1438 to 0bdeba0 Compare October 7, 2026 02:35
…e millisecond, sweep deleted tables' shards

Shard ids were `shardId-<table name>-<n>`, while `stream_shards` is keyed by
shard id alone. Any second stream-enabled table under one name collided:

- PostgreSQL keeps a deleted table's shards, so DeleteTable followed by
  CreateTable with the same name and a stream failed with
  InternalServerError every time.
- On PostgreSQL and SQLite, a second account creating a stream-enabled
  table with a name another account already uses failed the same way.
- A table name over about 40 characters produced a ShardId longer than the
  65 characters the AWS SDKs accept, so DescribeStream and GetShardIterator
  calls carrying it were rejected client-side.

Use the table id (a fresh UUID per table) in place of the name on
PostgreSQL and SQLite, as the MongoDB backend already does. New shard ids
are 61 characters. Shard rows that already exist keep their ids and keep
working; shards created from now on, including on a pre-existing table
that enables a stream for the first time, get the new form.

With recreate working, a second defect becomes reachable: stream labels had
one-second resolution on all three backends, so a table deleted and
recreated within the same second got the same stream ARN, and the old ARN
resolved to the new table's stream. Labels now carry milliseconds
(`2026-10-05T01:54:50.312`), the shape the service uses. Existing labels
are unchanged. On MongoDB, where recreate already worked, this reproduced
in one of three runs of the new test. The window is narrowed, not closed:
the ARN still resolves by name and label, so a recreate within one
millisecond would still alias; that is documented and tracked.

The label is now formatted in one place, `extenddb_storage::util::
format_stream_label`, and bound by every backend (PostgreSQL and SQLite
previously formatted it in SQL with `to_char` and `strftime`, MongoDB in its
own Rust function), so the three cannot drift apart and the shape has unit
tests of its own, including that it is always UTC.

PostgreSQL also never removed a deleted table's shard rows: DeleteTable
leaves shards and records in place so the stream can age out, the hourly
retention sweep trimmed the records, and the four shard rows per stream
stayed for good. The sweep now also removes the shards of a table that is
gone from the catalog once none of its records remain. Candidates are
shards older than the retention window whose table has no records; a
freshly created table is never a candidate because its shards are younger
than the window, so the gap between inserting shards and committing the
catalog row cannot be mistaken for a deletion, and a live table's unused
shards are left alone. SQLite already removed both with the table.

The PostgreSQL storage-level test harnesses (key_collation,
vector_control_plane, and the new stream_shard_retention) build their
scratch database by applying every .sql file under migrations/ and
data_migrations/ in filename order, instead of each carrying a copied
list of the files, so a new migration does not need to be added in three
more places.

Tests: tests/test_stream_table_reuse.py (recreate under one name, two
accounts with one name, a 221-character name within the SDK's ShardId
bounds, disable and re-enable), and
crates/storage-postgres/tests/stream_shard_retention.rs (shards kept after
DeleteTable and while a record remains, removed once the records are gone,
a live table's shards untouched), added to the PostgreSQL storage-level
CI step. format_stream_label unit tests for the three-digit shape and UTC.

Signed-off-by: Scott Robinson <robinnsc@amazon.com>
@robinnsc
robinnsc force-pushed the fix/pg-stream-shard-ids branch from 0bdeba0 to a62ba52 Compare October 7, 2026 03:09
@robinnsc
robinnsc requested a review from jcshepherd October 7, 2026 03:20

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants