diff --git a/.github/workflows/integration.yml b/.github/workflows/integration.yml index 717a84f6d..cfdfdf2ee 100644 --- a/.github/workflows/integration.yml +++ b/.github/workflows/integration.yml @@ -305,7 +305,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 --test put_create_race --test twi_conflict + run: cargo test --release -p extenddb-storage-postgres --test vector_control_plane --test put_create_race --test twi_conflict --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/put_create_race.rs b/crates/storage-postgres/tests/put_create_race.rs index 8b38f012a..98b02757a 100644 --- a/crates/storage-postgres/tests/put_create_race.rs +++ b/crates/storage-postgres/tests/put_create_race.rs @@ -72,6 +72,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 \ @@ -97,15 +117,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/twi_conflict.rs b/crates/storage-postgres/tests/twi_conflict.rs index 6752c2e62..5d7c6fa14 100644 --- a/crates/storage-postgres/tests/twi_conflict.rs +++ b/crates/storage-postgres/tests/twi_conflict.rs @@ -70,6 +70,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 \ @@ -95,15 +115,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/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/