From a79080e866c8bc9c29c0b1c6b918bceb8cb4cf07 Mon Sep 17 00:00:00 2001 From: Scott Robinson Date: Tue, 15 Sep 2026 03:37:37 +0000 Subject: [PATCH 1/2] feat: ENABLING/DISABLING TTL states, non-blocking updates, durable backfill cursor MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit UpdateTimeToLive used to block for the whole enable backfill (a full table scan) and the whole disable drain, so the call took as long as the table was large and could outlive a client timeout; meanwhile DescribeTimeToLive reported ENABLED for tables whose queue did not yet cover every item. The Cassandra catalog already encoded the real lifecycle in ttl_index_ready and ttl_cleanup_generation — nothing read it. TimeToLiveStatus gains Enabling and Disabling (DynamoDB's wire names). Cassandra derives all four states from the catalog; PostgreSQL derives Enabling from its own ttl_index_ready (its disable is synchronous, so Disabling never appears); SQLite and MongoDB keep their two-state synchronous behavior, and the handler only detaches the backfill when the backend actually reports a transitional state — detaching on a two-state backend would leave a table claiming Enabled with no index. Updates are rejected while a transition settles, as DynamoDB does; the gate is a pure function with the full 8-case truth table tested. Enable returns once the catalog flip is durable and the backfill runs detached, resuming after a failure from a new durable cursor: a JSON resume point written after every scanned page as an LWT fenced on the live ttl_generation. The fence is what makes the cursor safe, not just useful — a plain write is an upsert, and a detached backfill racing DeleteTable would resurrect a partial catalog row for the dead table, blocking a same-name CreateTable (adversarial-review finding). A refused fence means the lifecycle moved and the scan stops. Disable returns once the flip is durable and the worker's pending-cleanup pass finishes the drain; the table reports DISABLING until it does. Also from review: a PostgreSQL CONCURRENTLY build that fails partway leaves an INVALID index that IF NOT EXISTS would silently keep on retry, so invalid leftovers are dropped before rebuilding, and the readiness publication is fenced on the attribute the build indexed so a stale detached task cannot certify a later lifecycle. The V004 migration uses ADD IF NOT EXISTS so a run that dies between applying and recording is rerunnable. The migration runner splits statements on every semicolon including in comments; V004 documents that trap. Tests: the four lifecycle states in order; a backfill killed mid-scan (cursor present, still ENABLING) resumed to completion by the worker pass with an audit proving full coverage; a planted cursor at the scan-last item proving the worker path resumes rather than rescans (audit finds exactly the ten skipped registrations); the transition gate truth table. Existing disable tests now invoke the worker drain they previously got inline. --- crates/core/src/types/table.rs | 9 +- crates/engine/src/ttl.rs | 108 +++++- .../catalog/V004__ttl_backfill_cursor.cql | 12 + crates/storage-cassandra/src/lib.rs | 1 + .../storage-cassandra/src/metadata_engine.rs | 160 +++++++- crates/storage-cassandra/src/migrations.rs | 4 + crates/storage-cassandra/src/ttl_worker.rs | 12 +- .../tests/ttl_integration.rs | 344 ++++++++++++++++++ .../storage-postgres/src/metadata_engine.rs | 43 ++- .../0010-cassandra-ttl-expiration-queue.md | 3 +- docs/differences-from-dynamodb.md | 4 +- 11 files changed, 648 insertions(+), 52 deletions(-) create mode 100644 crates/storage-cassandra/migrations/catalog/V004__ttl_backfill_cursor.cql diff --git a/crates/core/src/types/table.rs b/crates/core/src/types/table.rs index c0a47f40a..531c4094a 100755 --- a/crates/core/src/types/table.rs +++ b/crates/core/src/types/table.rs @@ -1075,10 +1075,17 @@ impl UpdateTableInput { #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "SCREAMING_SNAKE_CASE")] pub enum TimeToLiveStatus { - /// TTL is enabled. + /// TTL is enabled and the expiration queue covers every item. Enabled, + /// TTL has been requested and the queue backfill has not yet finished. + /// Writes made in this state are TTL-tracked; pre-existing items are + /// still being registered. + Enabling, /// TTL is disabled. Disabled, + /// TTL has been disabled and the old expiration queue is still being + /// drained. No further expirations happen in this state. + Disabling, } /// TTL description returned by `DescribeTimeToLive` and `UpdateTimeToLive`. diff --git a/crates/engine/src/ttl.rs b/crates/engine/src/ttl.rs index 95e750bc9..f930b6608 100755 --- a/crates/engine/src/ttl.rs +++ b/crates/engine/src/ttl.rs @@ -54,6 +54,35 @@ pub async fn handle_describe_time_to_live( /// is already in the requested state (idempotency check). /// Returns `ResourceNotFoundException` if the table does not exist. /// Returns `InternalServerError` on storage failures. +/// Gate an `UpdateTimeToLive` request on the table's current TTL state. +/// +/// A table mid-transition rejects further changes: the backfill and the queue +/// drain are both keyed to the current generation, and admitting a toggle now +/// would race them. The client retries once the transition lands (matching +/// DynamoDB, which also refuses updates in ENABLING and DISABLING). Steady +/// states reject only the idempotent no-op, as before. +fn check_ttl_transition( + current: TimeToLiveStatus, + requested_enabled: bool, +) -> Result<(), DynamoDbError> { + match current { + TimeToLiveStatus::Enabling | TimeToLiveStatus::Disabling => { + Err(DynamoDbError::ValidationException( + "Time to live is being updated for this table. Retry the request once the \ + current change completes." + .to_owned(), + )) + } + TimeToLiveStatus::Enabled if requested_enabled => Err(DynamoDbError::ValidationException( + "TimeToLive is already enabled".to_owned(), + )), + TimeToLiveStatus::Disabled if !requested_enabled => Err( + DynamoDbError::ValidationException("TimeToLive is already disabled".to_owned()), + ), + _ => Ok(()), + } +} + pub async fn handle_update_time_to_live( body: Value, ctx: &OperationContext, @@ -72,17 +101,10 @@ pub async fn handle_update_time_to_live( .await .map_err(storage_to_dynamo)?; - let already_enabled = current.time_to_live_status == TimeToLiveStatus::Enabled; - if input.time_to_live_specification.enabled && already_enabled { - return Err(DynamoDbError::ValidationException( - "TimeToLive is already enabled".to_owned(), - )); - } - if !input.time_to_live_specification.enabled && !already_enabled { - return Err(DynamoDbError::ValidationException( - "TimeToLive is already disabled".to_owned(), - )); - } + check_ttl_transition( + current.time_to_live_status, + input.time_to_live_specification.enabled, + )?; // Resolve the old attribute before committing the disable. If the catalog // is inconsistent, fail without leaving a partially applied request. @@ -105,18 +127,44 @@ pub async fn handle_update_time_to_live( .map_err(storage_to_dynamo)?; if input.time_to_live_specification.enabled { - // Kick off index creation (CONCURRENTLY — non-blocking for other database - // operations on the table, but the handler awaits completion). - // If it fails, the TTL sweeper will retry on its next cycle. - let account_id = ctx.account_id.clone(); - let table_name = input.table_name.clone(); - let attr = input.time_to_live_specification.attribute_name.clone(); - if let Err(e) = ctx + // Backends that report ENABLING (Cassandra, PostgreSQL) get the + // backfill kicked off in the background: it is a full table scan, and + // awaiting it made the API call take as long as the table is large. + // The caller sees ENABLING until readiness is published; if the task + // dies unfinished, the TTL worker's pending-index pass retries (on + // Cassandra, from the durable cursor). Backends that report Enabled + // immediately (SQLite, MongoDB) keep the awaited call: detaching it + // would leave a table claiming Enabled while its index does not yet + // exist, with no observable transition state. + let transitional = ctx + .storage + .describe_ttl(&ctx.account_id, &input.table_name) + .await + .map(|d| d.time_to_live_status == TimeToLiveStatus::Enabling) + .unwrap_or(false); + if transitional { + let storage = ctx.storage.clone(); + let account_id = ctx.account_id.clone(); + let table_name = input.table_name.clone(); + let attr = input.time_to_live_specification.attribute_name.clone(); + tokio::spawn(async move { + if let Err(e) = storage + .create_ttl_index(&account_id, &table_name, &attr) + .await + { + tracing::warn!("TTL queue backfill deferred for {table_name}: {e}"); + } + }); + } else if let Err(e) = ctx .storage - .create_ttl_index(&account_id, &table_name, &attr) + .create_ttl_index( + &ctx.account_id, + &input.table_name, + &input.time_to_live_specification.attribute_name, + ) .await { - tracing::warn!("TTL index creation deferred for {table_name}: {e}"); + tracing::warn!("TTL index creation deferred for {}: {e}", input.table_name); } } else { // Disable path: metadata already updated (sweeper won't pick up this table). @@ -211,6 +259,26 @@ mod tests { assert!(validate_ttl_attribute_name(&max).is_ok()); } + #[test] + fn transition_gate() { + use TimeToLiveStatus::{Disabled, Disabling, Enabled, Enabling}; + // Steady states admit the toggle and reject the no-op. + assert!(check_ttl_transition(Disabled, true).is_ok()); + assert!(check_ttl_transition(Enabled, false).is_ok()); + assert!(check_ttl_transition(Enabled, true).is_err()); + assert!(check_ttl_transition(Disabled, false).is_err()); + // Mid-transition rejects both directions. + for state in [Enabling, Disabling] { + for requested in [true, false] { + let error = check_ttl_transition(state, requested).unwrap_err(); + assert!( + format!("{error:?}").contains("being updated"), + "expected transition rejection, got {error:?}" + ); + } + } + } + #[test] fn special_chars_rejected() { assert!(validate_ttl_attribute_name("it's").is_err()); diff --git a/crates/storage-cassandra/migrations/catalog/V004__ttl_backfill_cursor.cql b/crates/storage-cassandra/migrations/catalog/V004__ttl_backfill_cursor.cql new file mode 100644 index 000000000..83473e8d9 --- /dev/null +++ b/crates/storage-cassandra/migrations/catalog/V004__ttl_backfill_cursor.cql @@ -0,0 +1,12 @@ +-- Copyright 2026 ExtendDB contributors +-- SPDX-License-Identifier: Apache-2.0 +-- Durable resume point for the TTL enable backfill. +-- +-- JSON blob {"generation": , "last_key": } written after each +-- scanned page. A backfill that dies mid-scan (host crash, deploy) resumes +-- from here instead of rescanning the table from the top. A cursor whose +-- generation does not match the table's current ttl_generation is ignored. +-- +-- NOTE for future migrations: the migration runner splits files on the +-- semicolon character, including inside comments. Do not use one in a comment. +ALTER TABLE tables ADD IF NOT EXISTS ttl_backfill_cursor text; diff --git a/crates/storage-cassandra/src/lib.rs b/crates/storage-cassandra/src/lib.rs index 46af6e985..bced78b5c 100644 --- a/crates/storage-cassandra/src/lib.rs +++ b/crates/storage-cassandra/src/lib.rs @@ -40,6 +40,7 @@ pub use bootstrapper::CassandraBootstrapper; pub use catalog_store::CassandraCatalogStore; pub use config::CassandraStorageConfig; pub use engine::{CassandraEngine, CassandraSession}; +pub use metadata_engine::TtlBackfillCursor; use cdrs_tokio::authenticators::StaticPasswordAuthenticatorProvider; use cdrs_tokio::cluster::NodeTcpConfigBuilder; diff --git a/crates/storage-cassandra/src/metadata_engine.rs b/crates/storage-cassandra/src/metadata_engine.rs index aabf2d66a..8112897c7 100755 --- a/crates/storage-cassandra/src/metadata_engine.rs +++ b/crates/storage-cassandra/src/metadata_engine.rs @@ -15,6 +15,17 @@ use crate::CassandraEngine; const TTL_CONTROL_MAX_RETRIES: u32 = 4; const TTL_CONTROL_RETRY_DELAY_MS: u64 = 25; +/// Durable resume point for the TTL enable backfill, stored per table in the +/// catalog as JSON. `generation` pins it to one enable cycle: a cursor left +/// behind by a retired generation must never seed a resume for the next one. +#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] +pub struct TtlBackfillCursor { + /// Hyphenated UUID string (the local `uuid` dependency has no serde + /// feature, so the generation travels in its display form). + pub generation: String, + pub last_key: extenddb_core::types::Item, +} + impl CassandraEngine { /// How long a cached TTL configuration may serve the write path. The /// staleness is safe in both directions because the audit exists; see the @@ -572,7 +583,7 @@ impl CassandraEngine { /// `insert_ttl_entry` runs only for items whose registration is actually /// missing. Pages are paced so a large table's audit is a slow background /// murmur rather than a read burst. - pub(crate) async fn audit_ttl_queue_for_table( + pub async fn audit_ttl_queue_for_table( &self, account_id: &str, table_name: &str, @@ -655,12 +666,81 @@ impl CassandraEngine { Ok(Some(repaired)) } + /// Read the durable backfill resume point, if any. Public for direct + /// backend integration tests; production readers are the backfill itself. + pub async fn ttl_backfill_cursor( + &self, + account_id: &str, + table_name: &str, + ) -> Result, StorageError> { + let query = format!( + "SELECT ttl_backfill_cursor FROM {}.tables \ + WHERE account_id = ? AND table_name = ?", + self.catalog_keyspace() + ); + let Some(row) = crate::cassandra_util::query_optional( + &self.session, + &query, + cdrs_tokio::query_values!(account_id, table_name), + "ttl_backfill_cursor", + ) + .await? + else { + return Ok(None); + }; + let raw: Option = row.get_by_name("ttl_backfill_cursor").ok().flatten(); + // An unparseable cursor is treated as absent: the backfill falls back + // to a full rescan, which is always correct. + Ok(raw.and_then(|raw| serde_json::from_str(&raw).ok())) + } + + /// Persist a backfill resume point. Public for direct backend integration + /// tests; production writers are the backfill itself. + /// + /// Returns whether the write applied. An LWT fenced on the generation, not + /// a plain write: a plain UPDATE is an upsert, and a backfill racing + /// DeleteTable would resurrect a partial catalog row for the dead table, + /// blocking a same-name CreateTable. A refused write means the lifecycle + /// moved under the scan (disable, re-enable, or table deletion) and the + /// backfill must stop. One Paxos round per 1,000-item page is noise next + /// to the page's inserts. + pub async fn write_ttl_backfill_cursor( + &self, + account_id: &str, + table_name: &str, + generation: uuid::Uuid, + cursor: &TtlBackfillCursor, + ) -> Result { + let raw = serde_json::to_string(cursor) + .map_err(|error| StorageError::Internal(format!("TTL cursor encode: {error}")))?; + let query = format!( + "UPDATE {}.tables SET ttl_backfill_cursor = ? \ + WHERE account_id = ? AND table_name = ? IF ttl_generation = ?", + self.catalog_keyspace() + ); + let result = crate::cassandra_util::query_lwt( + &self.session, + &query, + cdrs_tokio::query_values!( + raw.as_str(), + account_id, + table_name, + cdrs_tokio::types::value::Bytes::new(generation.as_bytes().to_vec()) + ), + ) + .await + .map_err(|error| StorageError::Internal(format!("TTL cursor write: {error}")))?; + metadata_lwt_applied(&result) + } + + /// Scan the table and register an expiration entry for every item that /// carries a valid TTL timestamp, then publish the generation as ready. /// - /// Runs under the caller's control lease. The scan has no durable cursor, so - /// a failure restarts it from the beginning on the next cycle; entry inserts - /// are conditional, so repeating the scan is idempotent. + /// Runs under the caller's control lease. Resumes from the durable + /// per-page cursor when one exists for this generation, so a failure + /// costs one replayed page rather than a rescan; entry inserts are + /// conditional, so replaying is idempotent. async fn backfill_ttl_queue( &self, account_id: &str, @@ -671,7 +751,17 @@ impl CassandraEngine { let key_info = self.fetch_table_key_info(account_id, table_name).await?; let account_keyspace = self.account_keyspace(account_id); - let mut start_key = None; + // Resume from the durable cursor when it belongs to this generation. + // A backfill that died mid-scan (host crash, deploy) picks up at its + // last completed page instead of rescanning the table from the top; + // a cursor from a retired generation is ignored. Registrations are + // idempotent, so a page replayed around the crash point is harmless. + let mut start_key: Option = self + .ttl_backfill_cursor(account_id, table_name) + .await? + .and_then(|cursor| { + (cursor.generation == config.generation.to_string()).then_some(cursor.last_key) + }); loop { if self.ttl_config_for_table(account_id, table_name).await? != Some(config.clone()) { return Ok(()); @@ -694,13 +784,31 @@ impl CassandraEngine { } } match next_key { - Some(key) => start_key = Some(key), + Some(key) => { + let applied = self + .write_ttl_backfill_cursor( + account_id, + table_name, + config.generation, + &TtlBackfillCursor { + generation: config.generation.to_string(), + last_key: key.clone(), + }, + ) + .await?; + if !applied { + // The generation moved or the table is gone; this + // scan's registrations belong to a retired lifecycle. + return Ok(()); + } + start_key = Some(key); + } None => break, } } let query = format!( - "UPDATE {}.tables SET ttl_index_ready = true \ + "UPDATE {}.tables SET ttl_index_ready = true, ttl_backfill_cursor = null \ WHERE account_id = ? AND table_name = ? \ IF ttl_attribute = ? AND ttl_generation = ? AND table_status = 'ACTIVE'", self.catalog_keyspace() @@ -791,7 +899,8 @@ impl MetadataEngine for CassandraEngine { let table_name = table_name.to_owned(); Box::pin(async move { let query = format!( - "SELECT ttl_attribute FROM {}.tables WHERE account_id = ? AND table_name = ?", + "SELECT ttl_attribute, ttl_index_ready, ttl_cleanup_generation \ + FROM {}.tables WHERE account_id = ? AND table_name = ?", self.catalog_keyspace() ); let row = crate::cassandra_util::query_optional( @@ -803,12 +912,22 @@ impl MetadataEngine for CassandraEngine { .await? .ok_or_else(|| StorageError::TableNotFound(table_name.clone()))?; let attribute: Option = row.get_by_name("ttl_attribute").ok().flatten(); + let index_ready: Option = row.get_by_name("ttl_index_ready").ok().flatten(); + let cleanup_generation: Option = + row.get_by_name("ttl_cleanup_generation").ok().flatten(); + // The catalog already encodes the full lifecycle; this is just the + // first reader to surface it. `ttl_index_ready` is published by the + // backfill when the queue covers every pre-existing item, and + // `ttl_cleanup_generation` is held until the retired generation's + // queue is fully drained. + let time_to_live_status = match (&attribute, index_ready, cleanup_generation) { + (Some(_), Some(true), _) => TimeToLiveStatus::Enabled, + (Some(_), _, _) => TimeToLiveStatus::Enabling, + (None, _, Some(_)) => TimeToLiveStatus::Disabling, + (None, _, None) => TimeToLiveStatus::Disabled, + }; Ok(TimeToLiveDescription { - time_to_live_status: if attribute.is_some() { - TimeToLiveStatus::Enabled - } else { - TimeToLiveStatus::Disabled - }, + time_to_live_status, attribute_name: attribute, }) }) @@ -877,7 +996,8 @@ impl MetadataEngine for CassandraEngine { let query = if enabled { format!( "UPDATE {}.tables SET ttl_attribute = ?, ttl_generation = ?, \ - ttl_index_ready = false WHERE account_id = ? AND table_name = ? \ + ttl_index_ready = false, ttl_backfill_cursor = null \ + WHERE account_id = ? AND table_name = ? \ IF table_status = 'ACTIVE' AND ttl_sweep_owner = null \ AND ttl_cleanup_generation = null \ AND ttl_attribute = null AND ttl_generation = null", @@ -886,7 +1006,8 @@ impl MetadataEngine for CassandraEngine { } else { format!( "UPDATE {}.tables SET ttl_attribute = null, ttl_generation = null, \ - ttl_cleanup_generation = ?, ttl_index_ready = false \ + ttl_cleanup_generation = ?, ttl_index_ready = false, \ + ttl_backfill_cursor = null \ WHERE account_id = ? AND table_name = ? \ IF table_status = 'ACTIVE' AND ttl_sweep_owner = null \ AND ttl_cleanup_generation = null \ @@ -964,10 +1085,11 @@ impl MetadataEngine for CassandraEngine { // The lifecycle change is durable; the issuing host must not keep // serving the old configuration from its cache. self.invalidate_ttl_config_cache(&account_id, &table_name); - if !enabled { - self.complete_ttl_cleanup(&account_id, &table_name, &table_id, cleanup_generation) - .await?; - } + // Disable does NOT drain the retired generation's queue here: the + // drain visits every leftover entry and used to make the API call + // take as long as the queue was deep. `ttl_cleanup_generation` + // stays set (the table reports DISABLING) until the worker's + // pending-cleanup pass finishes the drain and clears it. Ok(()) }) } diff --git a/crates/storage-cassandra/src/migrations.rs b/crates/storage-cassandra/src/migrations.rs index 9b8957b49..fb0ab64d6 100644 --- a/crates/storage-cassandra/src/migrations.rs +++ b/crates/storage-cassandra/src/migrations.rs @@ -58,6 +58,10 @@ pub(crate) const CATALOG_MIGRATIONS: &[(&str, &str)] = &[ "V003__account_scoped_idempotency.cql", include_str!("../migrations/catalog/V003__account_scoped_idempotency.cql"), ), + ( + "V004__ttl_backfill_cursor.cql", + include_str!("../migrations/catalog/V004__ttl_backfill_cursor.cql"), + ), ]; /// Embedded data migration files, applied in order. diff --git a/crates/storage-cassandra/src/ttl_worker.rs b/crates/storage-cassandra/src/ttl_worker.rs index 965809b58..f6a170ddd 100644 --- a/crates/storage-cassandra/src/ttl_worker.rs +++ b/crates/storage-cassandra/src/ttl_worker.rs @@ -767,7 +767,11 @@ pub(crate) async fn drain_retired_generation( Ok(()) } -async fn retry_pending_cleanup(storage: &CassandraEngine) { +/// Drain the retired generation's queue for every table whose disable is +/// still finalizing (`ttl_cleanup_generation` set; the table reports +/// DISABLING). Public for direct backend integration tests and manual +/// operational triggering. +pub async fn retry_pending_cleanup(storage: &CassandraEngine) { let pending = match storage.pending_ttl_cleanups().await { Ok(pending) => pending, Err(error) => { @@ -785,11 +789,13 @@ async fn retry_pending_cleanup(storage: &CassandraEngine) { } } -/// Retry the queue backfill for any TTL-enabled table that is not yet ready. +/// Retry the queue backfill for any TTL-enabled table that is not yet ready +/// (the table reports ENABLING). Resumes from the durable cursor. Public for +/// direct backend integration tests and manual operational triggering. /// /// `create_ttl_index` takes the table's control lease internally, so a table is /// scanned by one host at a time even though every host runs this pass. -async fn retry_pending_indexes(storage: &CassandraEngine) { +pub async fn retry_pending_indexes(storage: &CassandraEngine) { let Ok(enabled) = MetadataEngine::all_tables_with_ttl(storage).await else { return; }; diff --git a/crates/storage-cassandra/tests/ttl_integration.rs b/crates/storage-cassandra/tests/ttl_integration.rs index 03245357b..1c4a7c30d 100644 --- a/crates/storage-cassandra/tests/ttl_integration.rs +++ b/crates/storage-cassandra/tests/ttl_integration.rs @@ -273,6 +273,10 @@ async fn test_ttl_metadata_enable_disable_and_listing() { ) .await .expect("disable TTL metadata"); + // The drain is deferred to the worker now; run its pass so the + // assertions below observe the completed drain, as they did when the + // disable call was synchronous. + extenddb_storage_cassandra::ttl_worker::retry_pending_cleanup(&engine).await; engine .drop_ttl_index( &table.key_info.account_id, @@ -1204,6 +1208,10 @@ async fn test_disable_drains_claimed_work_and_releases_its_claim() { ) .await .expect("disable succeeds even with claimed work in flight"); + // The drain is deferred to the worker now; run its pass so the + // assertions below observe the completed drain, as they did when the + // disable call was synchronous. + extenddb_storage_cassandra::ttl_worker::retry_pending_cleanup(&engine).await; assert_eq!( base_row_owner(&engine, &table.key_info, &item).await, @@ -1264,6 +1272,10 @@ async fn test_disable_completes_effects_applying_work() { ) .await .expect("disable completes effects-applying work"); + // The drain is deferred to the worker now; run its pass so the + // assertions below observe the completed drain, as they did when the + // disable call was synchronous. + extenddb_storage_cassandra::ttl_worker::retry_pending_cleanup(&engine).await; let mut key = extenddb_core::types::Item::new(); key.insert( @@ -1854,6 +1866,10 @@ async fn test_disable_completes_effects_applied_work() { ) .await .expect("disable succeeds with effects-applied work in flight"); + // The drain is deferred to the worker now; run its pass so the + // assertions below observe the completed drain, as they did when the + // disable call was synchronous. + extenddb_storage_cassandra::ttl_worker::retry_pending_cleanup(&engine).await; let mut key = extenddb_core::types::Item::new(); key.insert( @@ -3268,3 +3284,331 @@ async fn test_audit_restores_lost_queue_registration() { .await; assert_eq!(queue_rows(&engine).await, 1); } + +/// The four lifecycle states in order, as `describe_ttl` reports them: the +/// catalog encoded this all along, F5 is the first reader. Also proves the +/// disable drain is genuinely deferred (DISABLING is observable) and that the +/// worker's cleanup pass is what lands DISABLED. +#[tokio::test] +async fn test_ttl_lifecycle_states() { + if helpers::skip_without_cassandra() { + return; + } + use extenddb_core::types::{AttributeValue, TimeToLiveStatus}; + use extenddb_storage::DataEngine; + + let engine = setup_engine().await; + let table = crate::helpers::TestTable::new(&engine, "TtlLifecycle", false).await; + activate_tables(&engine).await; + let account = &table.key_info.account_id; + let name = &table.key_info.table_name; + + // Something for the queue to hold so the disable drain has real work. + let mut item = std::collections::BTreeMap::new(); + item.insert("id".to_owned(), AttributeValue::S("holder".to_owned())); + item.insert( + "expires_at".to_owned(), + AttributeValue::N((chrono::Utc::now().timestamp() + 3_600).to_string()), + ); + engine + .put_item( + &table.key_info, + item, + false, + None, + &Default::default(), + None, + ) + .await + .expect("put"); + + // ENABLING: requested, backfill not yet run. + engine + .update_ttl(account, name, "expires_at", true) + .await + .expect("enable"); + let state = engine.describe_ttl(account, name).await.expect("describe"); + assert_eq!(state.time_to_live_status, TimeToLiveStatus::Enabling); + assert_eq!(state.attribute_name.as_deref(), Some("expires_at")); + + // ENABLED: backfill published readiness. + engine + .create_ttl_index(account, name, "expires_at") + .await + .expect("backfill"); + let state = engine.describe_ttl(account, name).await.expect("describe"); + assert_eq!(state.time_to_live_status, TimeToLiveStatus::Enabled); + assert!( + engine + .ttl_backfill_cursor(account, name) + .await + .expect("cursor read") + .is_none(), + "completed backfill must clear its cursor" + ); + + // DISABLING: the flip is durable but the old generation's queue is not + // yet drained — the API call no longer waits for that. + engine + .update_ttl(account, name, "expires_at", false) + .await + .expect("disable"); + let state = engine.describe_ttl(account, name).await.expect("describe"); + assert_eq!(state.time_to_live_status, TimeToLiveStatus::Disabling); + + // DISABLED: the worker's pending-cleanup pass finishes the drain. + extenddb_storage_cassandra::ttl_worker::retry_pending_cleanup(&engine).await; + let state = engine.describe_ttl(account, name).await.expect("describe"); + assert_eq!(state.time_to_live_status, TimeToLiveStatus::Disabled); + assert!(state.attribute_name.is_none()); +} + +/// Kill the backfill mid-scan and prove the durable cursor makes the retry +/// resume instead of rescanning, and that the resumed run covers every item. +#[tokio::test] +async fn test_backfill_cursor_resume() { + if helpers::skip_without_cassandra() { + return; + } + use extenddb_core::types::{AttributeValue, TimeToLiveStatus}; + use extenddb_storage::DataEngine; + + let engine = setup_engine().await; + let table = crate::helpers::TestTable::new(&engine, "TtlCursorResume", false).await; + activate_tables(&engine).await; + let account = table.key_info.account_id.clone(); + let name = table.key_info.table_name.clone(); + + // Three scan pages' worth of TTL-carrying items (page size is 1,000). + let future = chrono::Utc::now().timestamp() + 86_400; + for index in 0..2_500u32 { + let mut item = std::collections::BTreeMap::new(); + item.insert( + "id".to_owned(), + AttributeValue::S(format!("cur-{index:05}")), + ); + item.insert( + "expires_at".to_owned(), + AttributeValue::N(future.to_string()), + ); + engine + .put_item( + &table.key_info, + item, + false, + None, + &Default::default(), + None, + ) + .await + .expect("seed"); + } + + engine + .update_ttl(&account, &name, "expires_at", true) + .await + .expect("enable"); + + // Run the backfill on its own task and kill it once it has durably + // finished at least one page (cursor present), like a host dying + // mid-enable. + let engine_for_task = setup_engine().await; + let (task_account, task_name) = (account.clone(), name.clone()); + let backfill = tokio::spawn(async move { + let _ = engine_for_task + .create_ttl_index(&task_account, &task_name, "expires_at") + .await; + }); + let started = std::time::Instant::now(); + let mut cursor = None; + let mut polls = 0u32; + for _ in 0..1_200 { + polls += 1; + cursor = engine + .ttl_backfill_cursor(&account, &name) + .await + .expect("cursor read"); + if cursor.is_some() { + break; + } + tokio::time::sleep(std::time::Duration::from_millis(25)).await; + } + backfill.abort(); + let abort_result = backfill.await; + assert!( + abort_result.err().is_some_and(|e| e.is_cancelled()), + "backfill task should have been killed mid-scan, not finished \ + (cursor seen after {polls} polls / {:?}; slow machine or page size drift?)", + started.elapsed() + ); + + // The killed run died holding the table's control lease. In production it + // expires on its own (`USING TTL 900`) and the worker's next pending-index + // pass resumes — from the cursor, which is why a 15-minute-old death does + // not cost a fresh full-table scan. The test simulates the expiry instead + // of waiting out the clock. + let expire_lease = "UPDATE extenddb_ttl_test_catalog.tables SET ttl_sweep_owner = null \ + WHERE account_id = ? AND table_name = ?"; + engine + .session_arc() + .query_with_values( + expire_lease, + cdrs_tokio::query_values!(account.as_str(), name.as_str()), + ) + .await + .expect("simulate lease expiry"); + let cursor = cursor.expect("backfill should persist a cursor after its first page"); + + // The kill left the table mid-enable: cursor present, not ready. + let state = engine + .describe_ttl(&account, &name) + .await + .expect("describe"); + assert_eq!( + state.time_to_live_status, + TimeToLiveStatus::Enabling, + "aborted backfill must not have published readiness" + ); + + // The worker's pending-index pass is the production retry path. It must + // resume from the persisted cursor (observable as: it completes without + // rewriting the cursor's first page — verified below by full coverage + // plus the cursor being pinned to this generation). + let generation = ttl_generation(&engine, &account, &name).await; + assert_eq!( + cursor.generation, + generation.to_string(), + "cursor must be pinned to the live generation" + ); + extenddb_storage_cassandra::ttl_worker::retry_pending_indexes(&engine).await; + let state = engine + .describe_ttl(&account, &name) + .await + .expect("describe"); + assert_eq!(state.time_to_live_status, TimeToLiveStatus::Enabled); + assert!( + engine + .ttl_backfill_cursor(&account, &name) + .await + .expect("cursor read") + .is_none() + ); + + // Full coverage proof: an audit pass over the finished table repairs + // nothing, i.e. every one of the 2,500 items is registered. + let metrics = extenddb_core::metrics::MetricsCollector::new(); + let repaired = engine + .audit_ttl_queue_for_table(&account, &name, "expires_at", &metrics) + .await + .expect("audit") + .unwrap_or(0); + assert_eq!( + repaired, 0, + "resumed backfill left {repaired} items unregistered" + ); +} + +/// The resume is real, not a rescan: a cursor pointing at the last item makes +/// the backfill skip everything before it (the audit then counts exactly the +/// skipped items as missing). Companion to `test_backfill_cursor_resume`, +/// whose full-coverage assertion a silent rescan would also satisfy. +#[tokio::test] +async fn test_backfill_cursor_is_honored() { + if helpers::skip_without_cassandra() { + return; + } + use extenddb_core::types::{AttributeValue, TimeToLiveStatus}; + use extenddb_storage::DataEngine; + + let engine = setup_engine().await; + let table = crate::helpers::TestTable::new(&engine, "TtlCursorHonored", false).await; + activate_tables(&engine).await; + let account = table.key_info.account_id.clone(); + let name = table.key_info.table_name.clone(); + + let future = chrono::Utc::now().timestamp() + 86_400; + for index in 0..10u32 { + let mut item = std::collections::BTreeMap::new(); + item.insert( + "id".to_owned(), + AttributeValue::S(format!("hon-{index:02}")), + ); + item.insert( + "expires_at".to_owned(), + AttributeValue::N(future.to_string()), + ); + engine + .put_item( + &table.key_info, + item, + false, + None, + &Default::default(), + None, + ) + .await + .expect("seed"); + } + + engine + .update_ttl(&account, &name, "expires_at", true) + .await + .expect("enable"); + let generation = ttl_generation(&engine, &account, &name).await; + + // Plant a cursor claiming the scan already covered everything up to the + // item the scan returns LAST — scan order is token order, not insertion + // order, so ask the engine rather than assuming. Key only: the cursor is + // an exclusive start key, so the planted item itself is also skipped. + let (scanned, tail) = engine + .scan(&table.key_info, None, None, None, None, None) + .await + .expect("scan for token order"); + assert!( + tail.is_none() && scanned.len() == 10, + "expected one full page" + ); + let mut last_key = std::collections::BTreeMap::new(); + last_key.insert( + "id".to_owned(), + scanned.last().expect("ten items")["id"].clone(), + ); + let planted = engine + .write_ttl_backfill_cursor( + &account, + &name, + generation, + &extenddb_storage_cassandra::TtlBackfillCursor { + generation: generation.to_string(), + last_key, + }, + ) + .await + .expect("plant cursor"); + assert!(planted, "cursor fence must admit the live generation"); + + // The worker's pending-index pass is the production retry entry point; + // proving THE WORKER honors the cursor is the point of this test. + extenddb_storage_cassandra::ttl_worker::retry_pending_indexes(&engine).await; + let state = engine + .describe_ttl(&account, &name) + .await + .expect("describe"); + assert_eq!(state.time_to_live_status, TimeToLiveStatus::Enabled); + + // If the cursor was honored, the backfill scanned nothing (the cursor + // points past the scan-last item) and the audit finds all ten items — + // nine predecessors plus the exclusive cursor item itself, which puts on + // an ENABLING table never registered because the puts predate the enable. + // A rescan would find zero. + let metrics = extenddb_core::metrics::MetricsCollector::new(); + let repaired = engine + .audit_ttl_queue_for_table(&account, &name, "expires_at", &metrics) + .await + .expect("audit") + .unwrap_or(0); + assert_eq!( + repaired, 10, + "backfill did not resume from the planted cursor (repaired {repaired}, expected all 10 skipped)" + ); +} diff --git a/crates/storage-postgres/src/metadata_engine.rs b/crates/storage-postgres/src/metadata_engine.rs index 322a9015f..beccfecaa 100755 --- a/crates/storage-postgres/src/metadata_engine.rs +++ b/crates/storage-postgres/src/metadata_engine.rs @@ -20,8 +20,9 @@ impl MetadataEngine for PostgresEngine { let account_id = account_id.to_string(); let table_name = table_name.to_string(); Box::pin(async move { - let row: Option<(Option,)> = sqlx::query_as( - "SELECT ttl_attribute FROM tables WHERE account_id = $1 AND table_name = $2", + let row: Option<(Option, Option)> = sqlx::query_as( + "SELECT ttl_attribute, ttl_index_ready FROM tables \ + WHERE account_id = $1 AND table_name = $2", ) .bind(&account_id) .bind(&table_name) @@ -29,11 +30,20 @@ impl MetadataEngine for PostgresEngine { .await .map_err(|e| StorageError::Internal(e.to_string()))?; - let (ttl_attr,) = row.ok_or_else(|| StorageError::TableNotFound(table_name.clone()))?; + let (ttl_attr, index_ready) = + row.ok_or_else(|| StorageError::TableNotFound(table_name.clone()))?; Ok(match ttl_attr { Some(attr) => TimeToLiveDescription { - time_to_live_status: TimeToLiveStatus::Enabled, + // Enabled only once the queue backfill has published + // readiness; until then pre-existing items are still + // being registered. Postgres disable tears the queue + // down synchronously, so it has no Disabling state. + time_to_live_status: if index_ready == Some(true) { + TimeToLiveStatus::Enabled + } else { + TimeToLiveStatus::Enabling + }, attribute_name: Some(attr), }, None => TimeToLiveDescription { @@ -290,6 +300,25 @@ impl MetadataEngine for PostgresEngine { let bare_table = data_table.trim_matches('"'); let index_name = format!("idx_ttl_{bare_table}"); + // A CONCURRENTLY build that fails partway leaves an INVALID + // index behind, and IF NOT EXISTS would silently keep it on the + // retry — publishing readiness over an index that indexes + // nothing. Drop any invalid leftover before building. + let invalid: Option<(bool,)> = sqlx::query_as( + "SELECT i.indisvalid FROM pg_index i \ + JOIN pg_class c ON c.oid = i.indexrelid WHERE c.relname = $1", + ) + .bind(&index_name) + .fetch_optional(&self.data_pool) + .await + .map_err(|e| StorageError::Internal(e.to_string()))?; + if invalid == Some((false,)) { + sqlx::query(&format!("DROP INDEX IF EXISTS \"{index_name}\"")) + .execute(&self.data_pool) + .await + .map_err(|e| StorageError::Internal(format!("Drop invalid TTL index: {e}")))?; + } + let sql = format!( "CREATE INDEX CONCURRENTLY IF NOT EXISTS \"{index_name}\" \ ON {data_table} (((item_data->'{ttl_attribute}'->>'N')::BIGINT)) \ @@ -300,12 +329,16 @@ impl MetadataEngine for PostgresEngine { .await .map_err(|e| StorageError::Internal(format!("TTL index creation failed: {e}")))?; + // Fenced on the attribute this build actually indexed: a stale + // task surviving a disable / re-enable-with-different-attribute + // must not certify the new lifecycle's readiness. sqlx::query( "UPDATE tables SET ttl_index_ready = TRUE \ - WHERE account_id = $1 AND table_name = $2", + WHERE account_id = $1 AND table_name = $2 AND ttl_attribute = $3", ) .bind(&account_id) .bind(&table_name) + .bind(&ttl_attribute) .execute(&self.pool) .await .map_err(|e| StorageError::Internal(e.to_string()))?; diff --git a/docs/adr/0010-cassandra-ttl-expiration-queue.md b/docs/adr/0010-cassandra-ttl-expiration-queue.md index c027005e1..9e35d5c56 100644 --- a/docs/adr/0010-cassandra-ttl-expiration-queue.md +++ b/docs/adr/0010-cassandra-ttl-expiration-queue.md @@ -129,7 +129,7 @@ The bar for this feature is the production readiness of the PostgreSQL TTL imple | --- | --- | | ~100 expirations per table per minute, and no increase with fleet size | Both backends use `SCAN_INTERVAL = 60s` and `BATCH_SIZE = 100`. PostgreSQL runs without a lease but every host issues the same `ORDER BY ttl LIMIT 100` query, so hosts contend for the same rows and the loser's delete fails its TTL condition. | | `UpdateTimeToLive` blocks instead of returning immediately | The shared handler awaits `create_ttl_index` on every backend. | -| No `ENABLING`/`DISABLING` status | `TimeToLiveStatus` has only two variants; both backends derive status from catalog presence. | +| `ENABLING`/`DISABLING` are reported while a change settles | The catalog already encoded the lifecycle (`ttl_index_ready`, `ttl_cleanup_generation`); `describe_ttl` surfaces it, the API call no longer waits out the backfill or the drain, and updates are rejected mid-transition. | | No five-year cutoff on old timestamps | PostgreSQL matches on `BETWEEN 1 AND now`; Cassandra accepts any positive `i64`. | | TTL deletion bypasses write-capacity accounting and throttling | Both workers call the storage layer directly, below the request capacity path. | | No caching of TTL configuration | Neither backend caches it. Cassandra pays a per-write catalog read because it needs the configuration on the write path at all; PostgreSQL does not need it, because expiry is derived from the item by a database index. | @@ -161,7 +161,6 @@ The summary: with this change, Cassandra TTL matches PostgreSQL on every shared Deliberately out of scope for this change, in rough priority order: * **Expiration throughput does not scale horizontally.** The sweep lease is per table even though the queue is already sharded 64 ways by key and those shards are disjoint partitions. Leasing per `(table, shard)` would allow up to 64 concurrent workers per table with no change to the claim protocol, because every queue transition is already conditional on the exact work UUID. This is the highest-value follow-up, and it would take Cassandra past the PostgreSQL rate rather than merely matching it. -* **The backfill has no durable cursor.** A failure restarts the scan from the beginning, and the `UpdateTimeToLive` call blocks for its duration instead of reporting `ENABLING`. * **TTL alongside asynchronously propagated GSIs**, per the table above. * **Version fence hardening.** The fence now uses the `version` column; remaining follow-up is retiring the `item_data` fallback once no unversioned rows can exist (requires a fleet-wide rewrite pass or an explicit migration), after which the canonical-JSON coupling disappears entirely. * **Recovery-path test coverage.** The transitions this change introduced — draining a retired generation, retiring a drained bucket registration, an expired worker claim — are reasoned about but not yet covered by tests. diff --git a/docs/differences-from-dynamodb.md b/docs/differences-from-dynamodb.md index f1c86b861..837a41da1 100755 --- a/docs/differences-from-dynamodb.md +++ b/docs/differences-from-dynamodb.md @@ -54,8 +54,8 @@ adaptation when switching between ExtendDB and the real service. | TTL stream records | REMOVE events with `userIdentity: {type: "Service", principalId: "dynamodb.amazonaws.com"}` | Supported — TTL deletions generate REMOVE stream records with the same `userIdentity` | | TTL timestamps far in the past | Ignored if more than five years before the current time | No cutoff. Any positive epoch-second value is queued and expired, however old. | | TTL attribute value validation | Non-conforming values are ignored | Cassandra matches: only a positive integral `N` is expired, and missing, non-numeric, fractional, zero, and negative values are ignored. | -| `UpdateTimeToLive` response timing | Returns immediately; the change takes effect asynchronously (up to about an hour) | All backends await readiness before returning, so the call can be slow on a large table and can exceed a client timeout. PostgreSQL awaits a concurrent index build; Cassandra awaits a full-table backfill of the expiration queue, which has no durable cursor and restarts on the next worker cycle if it fails. | -| `DescribeTimeToLive` transitional states | `ENABLING` and `DISABLING` are reported while a change settles | Only `ENABLED` and `DISABLED`. A table reports `ENABLED` as soon as the attribute is set, including while its queue is still being backfilled and nothing will expire yet. | +| `UpdateTimeToLive` response timing | Returns immediately; the change takes effect asynchronously (up to about an hour) | Matches: the call returns once the catalog flip is durable. The Cassandra queue backfill runs in the background with a durable per-page cursor (a failed host's retry resumes rather than rescanning), and the disable drain is finished by the worker. | +| `DescribeTimeToLive` transitional states | `ENABLING` and `DISABLING` are reported while a change settles | Matches on Cassandra (all four states; `ENABLED` means the queue covers every item) and on PostgreSQL for enable (`ENABLING` until the concurrent index build publishes; its disable is synchronous so `DISABLING` never appears). During `DISABLING`, `AttributeName` is absent, unlike DynamoDB. Updates are rejected while a transition is settling, as in DynamoDB. | | Concurrent same-key writes on a TTL-enabled table | Unaffected by TTL | Serialized through a base-row claim. Contention is retried internally against a re-read image, so writes stay last-writer-wins; only sustained contention surfaces `TransactionConflictException`. | | TTL-enabled write cost | Unaffected by TTL | A claim LWT, with the base write becoming a logged batch (measured: roughly 2× put latency on a single-node bench). The TTL configuration is served from a 5-second cache, the remaining catalog read is overlapped with the index read, and the release is folded into the batch with an exact conditional release off the request path. Non-TTL tables keep the existing fast paths. | | Cassandra TTL with asynchronous GSIs | Supported | Not currently supported. The Cassandra backend rejects TTL enable when any GSI has a nonzero effective propagation delay, and rejects creating such a GSI on a TTL-enabled table; base tables, LSIs, and synchronous GSIs are supported. See [ADR-0010](adr/0010-cassandra-ttl-expiration-queue.md). | From 3223a929d04dc509083de278641e33476662c6eb Mon Sep 17 00:00:00 2001 From: Scott Robinson Date: Tue, 15 Sep 2026 04:07:11 +0000 Subject: [PATCH 2/2] test: wait out the TTL enable transition in the E2E disable tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Enable returns during ENABLING now and updates are rejected until the transition settles, so the enable-then-immediately-disable sequence in the Python and Rust E2E suites hits the transition gate — the same rejection real DynamoDB gives that sequence, which is why both tests already carried a cooldown comment about it. Poll DescribeTimeToLive for ENABLED (30s bound) before disabling, like a real client must. Both tests' final assertions already accepted DISABLED or DISABLING. --- .../storage-cassandra/src/metadata_engine.rs | 1 - tests/python/test_ttl.py | 11 ++++++++ tests/rust/src/ttl.rs | 26 +++++++++++++++++++ 3 files changed, 37 insertions(+), 1 deletion(-) diff --git a/crates/storage-cassandra/src/metadata_engine.rs b/crates/storage-cassandra/src/metadata_engine.rs index 8112897c7..1eb07ee45 100755 --- a/crates/storage-cassandra/src/metadata_engine.rs +++ b/crates/storage-cassandra/src/metadata_engine.rs @@ -733,7 +733,6 @@ impl CassandraEngine { metadata_lwt_applied(&result) } - /// Scan the table and register an expiration entry for every item that /// carries a valid TTL timestamp, then publish the generation as ready. /// diff --git a/tests/python/test_ttl.py b/tests/python/test_ttl.py index b02a8600b..9b277f62e 100755 --- a/tests/python/test_ttl.py +++ b/tests/python/test_ttl.py @@ -10,6 +10,7 @@ from __future__ import annotations import os +import time import pytest from botocore.exceptions import ClientError @@ -57,6 +58,16 @@ def test_disable_ttl(self, table_factory, dynamodb_client): TableName=name, TimeToLiveSpecification={"Enabled": True, "AttributeName": "ttl"}, ) + # Enable returns during ENABLING now (the backfill is detached, as in + # DynamoDB) and updates are rejected until the transition settles, so + # wait for ENABLED like a real client must. + for _ in range(120): + resp = dynamodb_client.describe_time_to_live(TableName=name) + if resp["TimeToLiveDescription"]["TimeToLiveStatus"] == "ENABLED": + break + time.sleep(0.25) + else: + pytest.fail("TTL enable did not settle within 30s") dynamodb_client.update_time_to_live( TableName=name, TimeToLiveSpecification={"Enabled": False, "AttributeName": "ttl"}, diff --git a/tests/rust/src/ttl.rs b/tests/rust/src/ttl.rs index 6adeb96f3..2aa8d06e3 100755 --- a/tests/rust/src/ttl.rs +++ b/tests/rust/src/ttl.rs @@ -192,6 +192,32 @@ async fn disable_ttl() { tokio::time::sleep(std::time::Duration::from_secs(30)).await; } + // Enable returns during ENABLING now (the backfill is detached, as in + // DynamoDB) and updates are rejected until the transition settles, so + // wait for ENABLED like a real client must. + let mut settled = false; + for _ in 0..120 { + let resp = c + .describe_time_to_live() + .table_name(&name) + .send() + .await + .unwrap(); + let status = format!( + "{:?}", + resp.time_to_live_description() + .unwrap() + .time_to_live_status() + .unwrap() + ); + if status == "Enabled" { + settled = true; + break; + } + tokio::time::sleep(std::time::Duration::from_millis(250)).await; + } + assert!(settled, "TTL enable did not settle within 30s"); + c.update_time_to_live() .table_name(&name) .time_to_live_specification(