Skip to content
Merged
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
2 changes: 1 addition & 1 deletion .github/workflows/integration.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
31 changes: 7 additions & 24 deletions crates/storage-mongodb/src/table_engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
};
Expand Down Expand Up @@ -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
Expand All @@ -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);
}
}
Expand Down
91 changes: 81 additions & 10 deletions crates/storage-postgres/src/stream_engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand All @@ -28,21 +28,87 @@ 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<u64, StorageError> {
let candidates: Vec<String> = 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<String> =
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<String> = 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,
account_id: &str,
table_name: &str,
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') \
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()))?;

Expand All @@ -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) \
Expand Down Expand Up @@ -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())
})
}
Expand Down
4 changes: 2 additions & 2 deletions crates/storage-postgres/src/update_table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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()))?;
Expand Down
34 changes: 25 additions & 9 deletions crates/storage-postgres/tests/key_collation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,26 @@ fn base_conn() -> Option<String> {
(!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<String> {
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 \
Expand All @@ -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");
Expand Down
34 changes: 25 additions & 9 deletions crates/storage-postgres/tests/put_create_race.rs
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,26 @@ fn base_conn() -> Option<String> {
(!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<String> {
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 \
Expand All @@ -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");
Expand Down
Loading
Loading