From a62ba52862dd781ac595a65181a709e89eaa5a52 Mon Sep 17 00:00:00 2001 From: Scott Robinson Date: Tue, 6 Oct 2026 07:40:52 +0000 Subject: [PATCH] fix(streams): derive shard ids from the table id, label streams to the millisecond, sweep deleted tables' shards Shard ids were `shardId--`, 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 --- .github/workflows/integration.yml | 2 +- crates/storage-mongodb/src/table_engine.rs | 31 +-- crates/storage-postgres/src/stream_engine.rs | 91 +++++- crates/storage-postgres/src/update_table.rs | 4 +- .../storage-postgres/tests/key_collation.rs | 34 ++- .../tests/stream_shard_retention.rs | 262 ++++++++++++++++++ .../tests/vector_control_plane.rs | 34 ++- crates/storage-sqlite/src/stream.rs | 30 +- crates/storage-sqlite/src/update_table.rs | 3 +- crates/storage/src/util/arn.rs | 76 +++++ crates/storage/src/util/mod.rs | 4 +- docs/design/13-storage-mongodb.md | 12 +- docs/manuals/02-design-guide.md | 2 +- tests/test_stream_table_reuse.py | 164 +++++++++++ 14 files changed, 674 insertions(+), 75 deletions(-) create mode 100644 crates/storage-postgres/tests/stream_shard_retention.rs create mode 100644 tests/test_stream_table_reuse.py diff --git a/.github/workflows/integration.yml b/.github/workflows/integration.yml index 5d1ab88ad..619398b29 100644 --- a/.github/workflows/integration.yml +++ b/.github/workflows/integration.yml @@ -302,7 +302,7 @@ jobs: - name: Run PostgreSQL storage-level tests env: EXTENDDB_TEST_PG_CONNECTION_STRING: postgresql://postgres:devpass@127.0.0.1:5432 - run: cargo test --release -p extenddb-storage-postgres --test vector_control_plane + run: cargo test --release -p extenddb-storage-postgres --test vector_control_plane --test stream_shard_retention # The daemonized server logs to syslog; dump it so server-side failures # are diagnosable from the job log. diff --git a/crates/storage-mongodb/src/table_engine.rs b/crates/storage-mongodb/src/table_engine.rs index 87a5df208..951b28b9a 100644 --- a/crates/storage-mongodb/src/table_engine.rs +++ b/crates/storage-mongodb/src/table_engine.rs @@ -25,27 +25,6 @@ use extenddb_storage::util::{ 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). -/// -/// 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. -fn format_stream_label(now: time::OffsetDateTime) -> String { - format!( - "{:04}-{:02}-{:02}T{:02}:{:02}:{:02}", - now.year(), - u8::from(now.month()), - now.day(), - now.hour(), - now.minute(), - now.second(), - ) -} - impl TableEngine for MongoEngine { fn create_table( &self, @@ -184,7 +163,7 @@ impl MongoEngine { .as_ref() .is_some_and(|ss| ss.stream_enabled) { - Some(format_stream_label(now)) + Some(extenddb_storage::util::format_stream_label(now)) } else { None }; @@ -793,7 +772,9 @@ impl MongoEngine { .await .map_err(|e| StorageError::Internal(e.to_string()))?; if existing_shard.is_none() { - let label = format_stream_label(time::OffsetDateTime::now_utc()); + let label = extenddb_storage::util::format_stream_label( + time::OffsetDateTime::now_utc(), + ); update_doc.insert("stream_label", &label); self.init_stream_shards(table_id).await?; } else if table_doc @@ -805,7 +786,9 @@ impl MongoEngine { // Shards exist but the label was cleared by a // previous disable — restore a fresh label so the // ARN resolves again. - let label = format_stream_label(time::OffsetDateTime::now_utc()); + let label = extenddb_storage::util::format_stream_label( + time::OffsetDateTime::now_utc(), + ); update_doc.insert("stream_label", &label); } } diff --git a/crates/storage-postgres/src/stream_engine.rs b/crates/storage-postgres/src/stream_engine.rs index 151a90180..6eee09658 100755 --- a/crates/storage-postgres/src/stream_engine.rs +++ b/crates/storage-postgres/src/stream_engine.rs @@ -9,7 +9,7 @@ use extenddb_core::types::{ }; use extenddb_storage::StreamEngine; use extenddb_storage::error::StorageError; -use extenddb_storage::util::{parse_stream_arn, stream_arn}; +use extenddb_storage::util::{new_stream_label, parse_stream_arn, stream_arn}; use futures::future::BoxFuture; use sqlx::PgPool; @@ -28,6 +28,70 @@ impl PostgresEngine { /// # Errors /// /// Returns [`StorageError::Internal`] if any query fails. + /// Remove the shards of deleted tables once their stream has aged out. + /// + /// DeleteTable leaves a table's shards and records in place, as the service + /// keeps a deleted table's stream readable for 24 hours, and the records + /// are trimmed by the retention sweep above. The shards were never + /// trimmed, so every deleted stream-enabled table left four rows behind + /// for good. `stream_shards` is in the data database and `tables` in the + /// catalog, so this is two steps: candidates are shards older than the + /// retention window belonging to a table with no records left, and only those whose table id + /// is absent from the catalog are removed. A table created moments ago 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. A live table's unused shards (after a + /// stream was disabled and re-enabled) are left alone. + async fn cleanup_orphaned_stream_shards( + &self, + retention_hours: i64, + ) -> Result { + let candidates: Vec = sqlx::query_scalar( + "SELECT DISTINCT s.table_id FROM stream_shards s \ + WHERE s.created_at < NOW() - make_interval(hours => $1::integer) \ + AND NOT EXISTS (SELECT 1 FROM stream_records r WHERE r.table_id = s.table_id)", + ) + .bind(retention_hours) + .fetch_all(&self.data_pool) + .await + .map_err(|e| StorageError::Internal(e.to_string()))?; + if candidates.is_empty() { + return Ok(0); + } + let live: Vec = + sqlx::query_scalar("SELECT table_id FROM tables WHERE table_id = ANY($1)") + .bind(&candidates) + .fetch_all(&self.pool) + .await + .map_err(|e| StorageError::Internal(e.to_string()))?; + let gone: Vec = candidates + .into_iter() + .filter(|id| !live.contains(id)) + .collect(); + if gone.is_empty() { + return Ok(0); + } + // Re-check for records: a table deleted since the first query may + // still have records inside the window. A table's shards go together, + // once none of its records remain. + let result = sqlx::query( + "DELETE FROM stream_shards s WHERE s.table_id = ANY($1) \ + AND NOT EXISTS (SELECT 1 FROM stream_records r WHERE r.table_id = s.table_id)", + ) + .bind(&gone) + .execute(&self.data_pool) + .await + .map_err(|e| StorageError::Internal(e.to_string()))?; + if result.rows_affected() > 0 { + tracing::debug!( + shards = result.rows_affected(), + tables = gone.len(), + "removed stream shards of deleted tables past the retention window" + ); + } + Ok(result.rows_affected()) + } + pub(crate) async fn init_stream_shards( tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, data_pool: &PgPool, @@ -35,14 +99,16 @@ impl PostgresEngine { table_name: &str, 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') \ - WHERE account_id = $1 AND table_name = $2 \ - RETURNING stream_label", + // Formatted here, not in SQL, so every backend issues the same label + // shape from one function (`extenddb_storage::util::format_stream_label`). + let label = new_stream_label(); + sqlx::query( + "UPDATE tables SET stream_label = $3 WHERE account_id = $1 AND table_name = $2", ) .bind(account_id) .bind(table_name) - .fetch_one(&mut **tx) + .bind(&label) + .execute(&mut **tx) .await .map_err(|e| StorageError::Internal(e.to_string()))?; @@ -55,10 +121,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) \ @@ -424,6 +494,7 @@ impl StreamEngine for PostgresEngine { .execute(&self.data_pool) .await .map_err(|e| StorageError::Internal(e.to_string()))?; + self.cleanup_orphaned_stream_shards(retention_hours).await?; Ok(result.rows_affected()) }) } diff --git a/crates/storage-postgres/src/update_table.rs b/crates/storage-postgres/src/update_table.rs index 3cbdada33..6e9aa58b9 100755 --- a/crates/storage-postgres/src/update_table.rs +++ b/crates/storage-postgres/src/update_table.rs @@ -362,12 +362,12 @@ impl PostgresEngine { if current_label.is_none() { sqlx::query( - "UPDATE tables SET stream_label = \ - to_char(NOW(), 'YYYY-MM-DD\"T\"HH24:MI:SS') \ + "UPDATE tables SET stream_label = $3 \ WHERE account_id = $1 AND table_name = $2", ) .bind(account_id) .bind(&input.table_name) + .bind(extenddb_storage::util::new_stream_label()) .execute(&mut *tx) .await .map_err(|e| StorageError::Internal(e.to_string()))?; diff --git a/crates/storage-postgres/tests/key_collation.rs b/crates/storage-postgres/tests/key_collation.rs index 263972dc9..4988d17eb 100644 --- a/crates/storage-postgres/tests/key_collation.rs +++ b/crates/storage-postgres/tests/key_collation.rs @@ -65,6 +65,26 @@ fn base_conn() -> Option { (!conn.trim().is_empty()).then(|| conn.trim_end_matches('/').to_owned()) } +/// Every `.sql` file under `migrations/` and `data_migrations/`, each +/// directory in filename order, so a new migration is picked up here without +/// editing this file. +fn shipped_migrations() -> Vec { + let root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")); + let mut out = Vec::new(); + for dir in ["migrations", "data_migrations"] { + let mut files: Vec<_> = std::fs::read_dir(root.join(dir)) + .expect("read the migrations directory") + .map(|e| e.expect("directory entry").path()) + .filter(|p| p.extension().is_some_and(|e| e == "sql")) + .collect(); + files.sort(); + for f in files { + out.push(std::fs::read_to_string(&f).expect("read a migration file")); + } + } + out +} + fn skip(test: &str) { eprintln!( "SKIP {test}: EXTENDDB_TEST_PG_CONNECTION_STRING is not set, so there is no PostgreSQL \ @@ -90,15 +110,11 @@ async fn scratch() -> Scratch { .connect(&url) .await .expect("connect to the scratch database"); - for sql in [ - include_str!("../migrations/001_schema.sql"), - include_str!("../migrations/002_vector_indexes.sql"), - include_str!("../data_migrations/001_data_schema.sql"), - include_str!("../data_migrations/002_gsi_pending.sql"), - include_str!("../data_migrations/003_idempotency_account_scope.sql"), - include_str!("../data_migrations/004_vector_index_state.sql"), - ] { - sqlx::raw_sql(sql) + // One scratch database stands in for both the catalog and the data + // database, so every shipped migration file from both directories is + // applied to it, in filename order. + for sql in shipped_migrations() { + sqlx::raw_sql(&sql) .execute(&db) .await .expect("apply a shipped migration"); diff --git a/crates/storage-postgres/tests/stream_shard_retention.rs b/crates/storage-postgres/tests/stream_shard_retention.rs new file mode 100644 index 000000000..6c07a0cb0 --- /dev/null +++ b/crates/storage-postgres/tests/stream_shard_retention.rs @@ -0,0 +1,262 @@ +// Copyright 2026 ExtendDB contributors +// SPDX-License-Identifier: Apache-2.0 + +//! A deleted table's stream shards are removed once its stream has aged out. +//! +//! DeleteTable keeps a table's shards and records so the stream stays readable +//! for the retention window, as on the service. The records were always +//! trimmed by the retention sweep; the shards were not, and every deleted +//! stream-enabled table left its four shard rows behind for good. The sweep +//! now removes the shards of a table that is gone from the catalog once none +//! of its records remain, and leaves every live table's shards alone. +//! +//! Needs `EXTENDDB_TEST_PG_CONNECTION_STRING` (host-only, e.g. +//! `postgres://user:pass@127.0.0.1:5432`); skips otherwise. Each test builds +//! its own scratch database and drops it afterwards. + +use extenddb_core::types::{ + AttributeDefinition, BillingMode, CreateTableInput, DeleteTableInput, KeySchemaElement, + KeyType, ScalarAttributeType, StreamSpecification, StreamViewType, +}; +use extenddb_storage::{StreamEngine, TableEngine}; +use extenddb_storage_postgres::{PostgresConfig, PostgresEngine}; +use sqlx::PgPool; +use sqlx::postgres::PgPoolOptions; + +const ACCOUNT: &str = "123456789012"; +const REGION: &str = "us-east-1"; + +struct Scratch { + engine: PostgresEngine, + db: PgPool, + admin: PgPool, + db_name: String, +} + +impl Scratch { + async fn cleanup(self) { + let Scratch { + engine, + db, + admin, + db_name, + } = self; + drop(engine); + db.close().await; + sqlx::query(&format!( + "DROP DATABASE IF EXISTS \"{db_name}\" WITH (FORCE)" + )) + .execute(&admin) + .await + .expect("drop the scratch database"); + admin.close().await; + } +} + +fn base_conn() -> Option { + let conn = std::env::var("EXTENDDB_TEST_PG_CONNECTION_STRING").ok()?; + (!conn.trim().is_empty()).then(|| conn.trim_end_matches('/').to_owned()) +} + +/// Every `.sql` file under `migrations/` and `data_migrations/`, each +/// directory in filename order, so a new migration is picked up here without +/// editing this file. +fn shipped_migrations() -> Vec { + let root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")); + let mut out = Vec::new(); + for dir in ["migrations", "data_migrations"] { + let mut files: Vec<_> = std::fs::read_dir(root.join(dir)) + .expect("read the migrations directory") + .map(|e| e.expect("directory entry").path()) + .filter(|p| p.extension().is_some_and(|e| e == "sql")) + .collect(); + files.sort(); + for f in files { + out.push(std::fs::read_to_string(&f).expect("read a migration file")); + } + } + out +} + +async fn scratch() -> Scratch { + let base = base_conn().expect("caller checks base_conn() first"); + let db_name = format!("eddb_shrd_{}", uuid::Uuid::new_v4().simple())[..24].to_owned(); + let admin = PgPoolOptions::new() + .max_connections(1) + .connect(&format!("{base}/postgres")) + .await + .expect("connect to the postgres maintenance database"); + sqlx::query(&format!("CREATE DATABASE \"{db_name}\"")) + .execute(&admin) + .await + .expect("create the scratch database"); + let url = format!("{base}/{db_name}"); + let db = PgPoolOptions::new() + .max_connections(2) + .connect(&url) + .await + .expect("connect to the scratch database"); + // One scratch database stands in for both the catalog and the data + // database, so every shipped migration file from both directories is + // applied to it, in filename order. + for sql in shipped_migrations() { + sqlx::raw_sql(&sql) + .execute(&db) + .await + .expect("apply a shipped migration"); + } + // One scratch database stands in for both the catalog and the data + // database. The catalog schema's `stream_shards` carries a foreign key to + // `tables` that cannot exist in a real deployment (the two live in + // different databases) and that the engine's insert order violates; drop + // it so the scratch layout matches production. + sqlx::query("ALTER TABLE stream_shards DROP CONSTRAINT IF EXISTS stream_shards_table_id_fkey") + .execute(&db) + .await + .expect("match the two-database layout"); + sqlx::query("UPDATE settings SET value = '0' WHERE key = 'control_plane_delay_seconds'") + .execute(&db) + .await + .expect("pin the control-plane delay to zero"); + sqlx::query("INSERT INTO accounts (account_id, account_name) VALUES ($1, $2)") + .bind(ACCOUNT) + .bind(format!("acct-{db_name}")) + .execute(&db) + .await + .expect("seed the account row"); + let engine = PostgresEngine::new( + &PostgresConfig { + connection_string: url, + pool_size: 10, + max_item_size_bytes: 400_000, + }, + REGION, + ) + .await + .expect("open a PostgresEngine on the scratch database"); + Scratch { + engine, + db, + admin, + db_name, + } +} + +fn stream_table(name: &str) -> CreateTableInput { + CreateTableInput { + table_name: name.to_owned(), + key_schema: vec![KeySchemaElement { + attribute_name: "pk".to_owned(), + key_type: KeyType::Hash, + }], + attribute_definitions: vec![AttributeDefinition { + attribute_name: "pk".to_owned(), + attribute_type: ScalarAttributeType::S, + }], + billing_mode: Some(BillingMode::PayPerRequest), + stream_specification: Some(StreamSpecification { + stream_enabled: true, + stream_view_type: Some(StreamViewType::NewImage), + }), + ..Default::default() + } +} + +async fn shard_count(db: &PgPool, table_id: &str) -> i64 { + sqlx::query_scalar("SELECT COUNT(*) FROM stream_shards WHERE table_id = $1") + .bind(table_id) + .fetch_one(db) + .await + .expect("count shards") +} + +#[tokio::test] +async fn shards_of_a_deleted_table_go_once_its_records_are_gone() { + if base_conn().is_none() { + eprintln!("SKIP: EXTENDDB_TEST_PG_CONNECTION_STRING is not set"); + return; + } + let s = scratch().await; + + let gone = s + .engine + .create_table(ACCOUNT, stream_table("t_gone")) + .await + .expect("create the table to delete"); + let live = s + .engine + .create_table(ACCOUNT, stream_table("t_live")) + .await + .expect("create the table that stays"); + let gone_id = gone.table_id.clone(); + let live_id = live.table_id.clone(); + assert_eq!(shard_count(&s.db, &gone_id).await, 4); + assert_eq!(shard_count(&s.db, &live_id).await, 4); + + s.engine + .delete_table( + ACCOUNT, + DeleteTableInput { + table_name: "t_gone".to_owned(), + }, + ) + .await + .expect("delete the table"); + assert_eq!( + shard_count(&s.db, &gone_id).await, + 4, + "DeleteTable keeps the shards so the stream stays readable" + ); + + // One record still inside the retention window keeps the shards. + let shard: String = sqlx::query_scalar( + "SELECT shard_id FROM stream_shards WHERE table_id = $1 ORDER BY shard_id LIMIT 1", + ) + .bind(&gone_id) + .fetch_one(&s.db) + .await + .expect("one shard id"); + sqlx::query( + "INSERT INTO stream_records (shard_id, sequence_number, table_id, event_name, record_data, created_at) \ + VALUES ($1, '000000000000000000001', $2, 'INSERT', '{}'::jsonb, NOW() + interval '1 hour')", + ) + .bind(&shard) + .bind(&gone_id) + .execute(&s.db) + .await + .expect("seed a record that is still inside the window"); + + // A retention of zero hours makes everything already written eligible. + s.engine + .cleanup_expired_stream_records(0) + .await + .expect("sweep"); + assert_eq!( + shard_count(&s.db, &gone_id).await, + 4, + "shards stay while a record of theirs remains" + ); + assert_eq!(shard_count(&s.db, &live_id).await, 4); + + sqlx::query("DELETE FROM stream_records WHERE table_id = $1") + .bind(&gone_id) + .execute(&s.db) + .await + .expect("age the record out"); + s.engine + .cleanup_expired_stream_records(0) + .await + .expect("sweep"); + assert_eq!( + shard_count(&s.db, &gone_id).await, + 0, + "the deleted table's shards are removed" + ); + assert_eq!( + shard_count(&s.db, &live_id).await, + 4, + "a live table's shards are never swept, even with no records" + ); + + s.cleanup().await; +} diff --git a/crates/storage-postgres/tests/vector_control_plane.rs b/crates/storage-postgres/tests/vector_control_plane.rs index 336011702..4c077fa1d 100644 --- a/crates/storage-postgres/tests/vector_control_plane.rs +++ b/crates/storage-postgres/tests/vector_control_plane.rs @@ -115,6 +115,26 @@ fn base_conn() -> Option { (!conn.trim().is_empty()).then(|| conn.trim_end_matches('/').to_owned()) } +/// Every `.sql` file under `migrations/` and `data_migrations/`, each +/// directory in filename order, so a new migration is picked up here without +/// editing this file. +fn shipped_migrations() -> Vec { + let root = std::path::Path::new(env!("CARGO_MANIFEST_DIR")); + let mut out = Vec::new(); + for dir in ["migrations", "data_migrations"] { + let mut files: Vec<_> = std::fs::read_dir(root.join(dir)) + .expect("read the migrations directory") + .map(|e| e.expect("directory entry").path()) + .filter(|p| p.extension().is_some_and(|e| e == "sql")) + .collect(); + files.sort(); + for f in files { + out.push(std::fs::read_to_string(&f).expect("read a migration file")); + } + } + out +} + /// Report the reason a test did nothing, loudly enough to notice in a log. fn skip(test: &str) { eprintln!( @@ -157,15 +177,11 @@ async fn scratch(pgvector: Pgvector) -> Scratch { .expect("create the pgvector extension"); } - for sql in [ - include_str!("../migrations/001_schema.sql"), - include_str!("../migrations/002_vector_indexes.sql"), - include_str!("../data_migrations/001_data_schema.sql"), - include_str!("../data_migrations/002_gsi_pending.sql"), - include_str!("../data_migrations/003_idempotency_account_scope.sql"), - include_str!("../data_migrations/004_vector_index_state.sql"), - ] { - sqlx::raw_sql(sql) + // One scratch database stands in for both the catalog and the data + // database, so every shipped migration file from both directories is + // applied to it, in filename order. + for sql in shipped_migrations() { + sqlx::raw_sql(&sql) .execute(&catalog) .await .expect("apply a shipped migration"); diff --git a/crates/storage-sqlite/src/stream.rs b/crates/storage-sqlite/src/stream.rs index 3d8850e8e..62943f678 100644 --- a/crates/storage-sqlite/src/stream.rs +++ b/crates/storage-sqlite/src/stream.rs @@ -33,21 +33,25 @@ impl SqliteEngine { table_name: &str, table_id: &str, ) -> Result { - let label: String = sqlx::query_scalar( - "UPDATE tables SET stream_label = strftime('%Y-%m-%dT%H:%M:%S','now') \ - WHERE account_id = ? AND table_name = ? RETURNING stream_label", - ) - .bind(account_id) - .bind(table_name) - .fetch_one(&mut **tx) - .await - .map_err(|e| StorageError::Internal(e.to_string()))?; + // Formatted here, not in SQL, so every backend issues the same label + // shape from one function (`extenddb_storage::util::format_stream_label`). + let label = extenddb_storage::util::new_stream_label(); + sqlx::query("UPDATE tables SET stream_label = ? WHERE account_id = ? AND table_name = ?") + .bind(&label) + .bind(account_id) + .bind(table_name) + .execute(&mut **tx) + .await + .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 821bda1d9..d527c3e6a 100644 --- a/crates/storage-sqlite/src/update_table.rs +++ b/crates/storage-sqlite/src/update_table.rs @@ -293,9 +293,10 @@ 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 = ? \ WHERE account_id = ? AND table_name = ?", ) + .bind(extenddb_storage::util::new_stream_label()) .bind(account_id) .bind(&input.table_name) .execute(&mut *tx) diff --git a/crates/storage/src/util/arn.rs b/crates/storage/src/util/arn.rs index b5d266dca..813c217bc 100755 --- a/crates/storage/src/util/arn.rs +++ b/crates/storage/src/util/arn.rs @@ -11,6 +11,40 @@ pub fn index_arn(region: &str, account_id: &str, table_name: &str, index_name: & format!("arn:aws:dynamodb:{region}:{account_id}:table/{table_name}/index/{index_name}") } +/// The label that tells one stream on a table name from the next, in the +/// shape the service uses: `YYYY-MM-DDThh:mm:ss.sss`, UTC, millisecond +/// precision, no zone suffix. It is part of the stream ARN +/// (`.../table//stream/
-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. On PostgreSQL, DeleteTable leaves a table's shards and records in place; the hourly retention sweep trims records older than 24 hours and then removes the shards of a table that no longer exists once none of its records remain (SQLite removes both with the table). A stream ARN carries the table name and a millisecond label, and resolves to the table currently holding that name and label, so a table deleted and recreated within one millisecond would give the old ARN the new stream; that window is tracked, not closed. ### Iterator Types diff --git a/tests/test_stream_table_reuse.py b/tests/test_stream_table_reuse.py new file mode 100644 index 000000000..2563a8946 --- /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)