Feat/storage cassandra in tree - #338
jcshepherd wants to merge 59 commits into
Conversation
amrith
left a comment
There was a problem hiding this comment.
Would you take a look at the failing checks - or are they related to 341?
There was a problem hiding this comment.
Requesting changes. Twelve blocking items are commented inline; the full analysis with the probes, the should-fix list, and a proposed acceptance path comes separately.
Four items are not code and have no line to sit on:
- CI on this head is red on clippy (three
too_many_argumentsand onefrom_strnaming in the crate), test, notices-current (Cargo.lock is stale under--locked), launcher-tests, CodeQL (account_id written to a log), and both MongoDB legs. The 1.88.0 leg was cancelled, so MSRV never completed. - No workflow in this PR runs anything against Cassandra, so none of the crate tests, the pytest suite, or the Rust integration suite cover this backend in CI. RFC-0002 requires conformance results tracked in CI, and this backend implements every trait, so that obligation is total. The workflow in #339 targets this branch rather than main.
For context on what holds up: 111 differential write probes and every read probe matched the postgres backend exactly, the shared pytest suite passes 860 tests against Cassandra, the tenant boundary held on every cross-account probe, and throughput at 32 clients is within roughly 10 percent of postgres on writes. The defects are concentrated in the write paths that depend on what they read. UpdateItem already has a sound protocol; applying it to conditional PutItem, conditional DeleteItem, PutItem versioning, and the transaction token would close five of the twelve.
| if indexes.is_empty() && stream_stmt.is_none() && ttl_config.is_none() { | ||
| // Fast path: no batch needed. | ||
| let insert_query = format!( | ||
| "INSERT INTO {data_keyspace}.{ddb_table} \ |
There was a problem hiding this comment.
Conditional PutItem is not atomic. The row is read at 148 and the condition evaluated in Rust at 189, then this INSERT carries no IF and no version check, so two writers on the same key both pass and both write. Probe with attribute_not_exists(guard) on a non-key attribute: both succeeded in 30 of 30 trials, postgres 0 of 60. Needs the OCC treatment UpdateItem has (read version, write with IF version = ?, retry).
| /// This is the only condition we map directly to a null-aware LWT without a | ||
| /// prior read. | ||
| pub(crate) fn is_attribute_not_exists_key( | ||
| condition: Option<&Expr>, |
There was a problem hiding this comment.
This never receives the expression name map, so attribute_not_exists(#p) with {"#p": "pk"} compares the literal #p against the key names, misses the LWT fast path, and falls into the read-then-write path above. Same request, two spellings, different atomicity: alias form both writers succeeded 20 of 20, plain form 0 of 20. The doc comment states this resolution happens; resolve maps before the comparison.
| if indexes.is_empty() && stream_stmt.is_none() && ttl_config.is_none() { | ||
| // Fast path: no batch needed. | ||
| let insert_query = format!( | ||
| "INSERT INTO {data_keyspace}.{ddb_table} \ |
There was a problem hiding this comment.
This write leaves the version cell untouched, so a concurrent UpdateItem's IF version = n still applies and overwrites the put with a post-image computed from the pre-put row. Probe with PutItem racing SET u = 1: the put was lost in 17 of 20 trials (pre-put attributes present, put attributes gone) with both calls reporting success. Every write path needs to advance or clear version under the same protocol.
| let mutation_timestamp = chrono::Utc::now().timestamp_micros(); | ||
|
|
||
| // Delete the item (with index updates if needed). | ||
| let delete_cql = format!("DELETE FROM {data_keyspace}.{ddb_table} WHERE pk = ?"); |
There was a problem hiding this comment.
The condition is evaluated at 469, then this DELETE has no IF and no version check. Probe: seed {pk, v=1}, DeleteItem IF v = 1 racing UpdateItem SET v = 2: both calls succeeded in 30 of 30 trials and the item was gone in 16 of 30, which no serial order allows. Needs DELETE ... IF version = ? against the version read, with retry.
| "Idempotency account does not match transaction account".to_owned(), | ||
| )); | ||
| } | ||
| self.check_idempotency_token( |
There was a problem hiding this comment.
The token is reserved durably at 861 before the ledger entry exists, and nothing deletes it when the transaction is cancelled. The retry then matches the stored fingerprint at 893 and returns IdempotentReplay, which the engine reports as a successful response with no writes. Executed: transaction cancelled by a failing condition, blocking item removed, identical call with the same token returned ok and wrote nothing. The SDKs generate this token and retry TransactWriteItems automatically. Reserve it inside the ledger commit, or release it on every cancellation path.
| gsi_default_delay_ms: Arc::new(std::sync::atomic::AtomicU64::new(10)), // Default 10ms | ||
| gsi_queue: crate::gsi_queue::GsiQueue::new(), | ||
| hlc: crate::stream_util::new_shared_hlc( | ||
| config.instance_id.as_deref().unwrap_or("default"), |
There was a problem hiding this comment.
instance_id is always None here, so every node falls back to "default" and derives the same HLC node id (startup logs node_id=137 instance_id="default"). set_instance_id has no caller anywhere in the tree. stream_records keys on (shard_id, sequence_number), so two nodes in the same millisecond with the same counter and shard overwrite each other's record. ADR 0018's uniqueness argument depends on this being wired.
| /// | ||
| /// Called by the server before constructing the storage engine, passing | ||
| /// a stable unique identifier for this server process (typically the | ||
| /// bind address, e.g. "192.168.1.1:18443"). Backends that need per-instance |
There was a problem hiding this comment.
Added to the shared trait for one backend's need and never called: grep -rn set_instance_id crates/ finds only the definitions, this default, the Cassandra impl, and the wrapper in extenddb-config. Wire it in the serve path or derive the node id from something already unique.
| }); | ||
| let resolved_admin_user = cassandra_user | ||
| .unwrap_or_else(|| std::env::var("USER").unwrap_or_else(|_| "cassandra".to_owned())); | ||
| let resolved_keyspace_prefix = keyspace_prefix.unwrap_or(prefix); |
There was a problem hiding this comment.
Re-running init against a cluster where the extenddb role already exists produces a config the server cannot start with. Executed with --keyspace-prefix r_probe --cassandra-user cassandra: init created r_probe_catalog and r_probe_account_<id>, printed "User 'extenddb' already exists", then wrote keyspace_prefix = "extenddb" and no username or password; serve failed with "Connection timeout after 10s". Persist the resolved prefix, and either fail loudly or rotate the password when the role exists.
| /// Returns a standard test configuration for Cassandra. | ||
| pub fn test_config() -> CassandraStorageConfig { | ||
| let mut config = CassandraStorageConfig { | ||
| contact_points: vec!["127.0.0.1:9042".to_string()], |
There was a problem hiding this comment.
Hardcoded contact point with no environment override and no skip gate, so the 25 tests in this crate's integration files fail cargo test --workspace unless a Cassandra node happens to be on 9042. That is the red test job and the first checklist item in the PR template. The other backends keep live service tests out of the workspace suite or gate them.
| @@ -326,6 +326,26 @@ def test_transact_write_conditional_put_fail(dynamodb_client, hash_table): | |||
| resp = dynamodb_client.get_item(TableName=hash_table, Key={"pk": {"S": "cp-2"}}) | |||
| assert resp["Item"]["v"]["S"] == "old" | |||
|
|
|||
| def test_transact_write_conditional_put_fail_diag(dynamodb_client, hash_table): | |||
There was a problem hiding this comment.
This test prints and asserts nothing, so it passes on every backend regardless of behavior. Its printed branch ("NO EXCEPTION - transaction succeeded (wrong)") describes a conditional Put inside TransactWriteItems succeeding when it must fail, which is the same defect class as the put and delete comments above. Drop it from this PR; the asserting test directly above covers the case once the backend is fixed.
f6af48f to
8151972
Compare
Backends that need per-node identity (e.g. to construct unique hybrid logical clock (HLC) values for stream sequence ids), can receive an identity via StorageConfig during extenddb serve. Default implementations to avoid breaking backends that don't need this.
Migrates extenddb-cassandra-plugin into crates/storage-cassandra as a first-class backend, selected via --features cassandra.
…must_use, unwrap, casing
…es on the same key SSE and table class attributes were defaulting to None: now they are plumbed through. Concurrent writes to the same key are now linearized by IF NOT EXISTS for inserts, and item-level version numbers for updates. This uses Cassandra's lightweight transactions (LWTs) which are Paxos-based and require two round-trips per write, adding roughly ~2-3ms/write without contention. With contention, short backoffs may add up to 200ms before either writing successfully, or failing the write back to the caller.
The branch reconstruction (845f933 and the rebuilt series after it) re-applied the backend squash and the PR #338 OCC changes but not merged PR #339, so the follow-up work reviewed and approved there vanished from the branch. This restores it, adapted to the reconstructed base rather than re-applied verbatim: #338's OCC layer (occ_write / occ_delete_lwt, always-non-null version) and the rebuilt claim protocol supersede #339's put-path version bumps and claim-fence changes, and those stay superseded — the base's newer per-page repair cursors and table-not-found tolerance are also kept where #339's equivalents were older. What comes back is everything orthogonal: * The self-healing audit: a leaseless, paced rescan of every ready TTL table that re-registers any item whose queue entry is missing (a lost registration was otherwise permanent for an item never written again, with nothing alarming). One quorum read per healthy item, deferring to claimed work, standing aside on lifecycle movement. Its repair count (TtlAuditRepairedEntryCount) is the drift signal. * Transaction recovery reconciles the pre-commit image: PREPARE persists each row's prior image into the ledger (serde-defaulted, old ledgers stay insert-only) and COMMITTING recovery retires the old queue entry as the live path does, DELETE ops included. The insert-only reconcile helpers lose their last caller and are removed. * Sweep throughput: concurrent waves (16 in flight) over the queue's disjoint shard partitions, batch 100 -> 1,000, lease renewed per 64-row wave, per-row pre-effects gate as a quorum config read rather than a lease renewal (per-row renewals serialize on one catalog partition's Paxos), row failures confined to the row and counted per table (TtlSweepRowErrorCount). * The 5-second epoched TTL-config cache on the write path, invalidated on update_ttl; update_table's GSI-fence decision reads quorum and uncached. The idle outbox pass visits only keyspaces of accounts with TTL tables, widening every tenth cycle. * The repair worker emits TtlRepairMarkerCount, the unresolved-marker gauge ADR-0010 asks operators to watch, from an (observed, resolved) tuple merged into the base's newer cursor-paging implementation. * The Cassandra CI workflow (deleted outright by the rebase), the skip-without-Cassandra test gating that keeps the serviceless workspace test job green, the benchmark harness, the audit integration test, the workspace clippy fixes, and the ADR / differences-doc updates. Verified against a fresh Cassandra 4.1: 25/25 ttl_integration serial (88s), 1 metadata, 1,145 workspace lib tests, clippy -D warnings clean, fmt clean.
…mark The ported direct integration suite caught this: the OCC fast path added in the #338 rework (plain delete, no TTL claim, no transaction owner — the most common delete shape) returned after its fence LWT without writing partition_max_delete_timestamp. The transaction prepare path consults that high-water mark to order new-item PUTs against deletes in the same partition, so the common delete silently disarmed the check the claimed and transactional delete paths still fed. Both fast-path arms (sort-key and pk-only) now write the mark before the fence, matching the claimed path's write-before-delete ordering: a mark advanced for a delete the fence then refuses only widens a conservative check, while the reverse order can lose the mark on a crash between the two.
The original extenddb-cassandra-plugin repository carried 139 direct integration tests against the storage traits — real node, no HTTP server in between — of which the in-tree port had picked up only the four tag tests. The other 135 covered CRUD, Query/Scan semantics and pagination, transactions (TransactGet, TransactWrite rollback and idempotency, ledger operations), GSI/LSI physical schema and the async propagation workers, streams, access keys, accounts, users, groups, roles, policies, admin and settings stores, and backup/restore state. None of that had automated coverage in-tree. One test binary (direct_integration), one module per area, sharing tests/common/mod.rs — the in-tree descendant of the plug-in's helpers.rs, which is why the port is mostly mechanical: keyspace-prefix adaptation, the two CreateTableInput fields added since the fork, and the skip-without-Cassandra guard on all 139 tests so the serviceless workspace test job stays green. Four tests the inventory flagged as weak were strengthened rather than copied: the table lifecycle and settings tests printed outcomes instead of asserting them (the settings test also now restores the live control-plane knob it mutates, which used to leak past the test run); the pk-only delete-timestamp test carried dead bindings and now verifies the delete; and the sync-GSI test's TODO was hiding a real semantic difference — the plug-in engine wrote GSIs synchronously by default, in-tree the default is 10ms async propagation via a queue no test worker drains, so its "verify the write path executes without crashing" would never have caught a lost index write. It now pins the index to synchronous propagation (the same knob production routing reads) and asserts the index row is visible through the index-scan path. The suite runs in parallel: every test provisions its own account and keyspace. Wired into the Cassandra CI workflow after the existing serial suites.
…ckfill cursor 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.
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.
Review suggestion from jcshepherd. The standalone metadata_integration binary held one test — the in-tree consolidation of the plug-in repo's four tag tests — while the ported suite carried those four originals, so merging the binaries also deduplicates: the consolidated phased lifecycle test replaces the four in direct/metadata_engine.rs (same assertions, one engine connection instead of four catalog-migration setups). It is parallel-safe like the rest of the suite. The TTL suite stays its own serial binary: its tests drive global sweep, repair, and audit passes that would see every concurrent test's tables, so the two binaries have deliberately opposite concurrency contracts. One fewer test binary to link, one fewer CI step.
…, durable backfill cursor
test: Porting over the Cassandra Repo integration suite
…havior and GSI initialization. (PR338, B12) - extenddb init now honors a keyspace prefix provided to init, instead of defaulting to 'extenddb'. If extenddb init is run on a partially initialized Cassandra instance with an extenddb application user already configured, requires application user password to re-run. (PR338, B6) - New GSIs go into creating state, and don't enter active state until backfill is complete.
…ole account lockout, Cargo.lock update
…ts to avoid unpredictable exhaustion of the connection pool
…d enable script to regenerate it
…, query scan follows Amazon DDB key sorting rules in both directions (mirroring #362 for PostgreSQL and SQLite).
What
This is a full-featured Cassandra backend for ExtendDB. The only current ExtendDB feature that it does not support is vector indexes.
Why
Refer to docs/rfcs/draft-cassandra-backend.md (in the PR) for background and design notes. Detailed design notes on specific aspects are available in the ADRs including in this PR.
Closes #
Testing done
pytests, rust-integration tests, unit tests all pass. Clippy is a work in progress.
Checklist
cargo test --workspace)cargo fmt --check)cargo clippy -- -W clippy::pedantic)Storagetrait, auth model, on-diskformat, or public CLI surface, an RFC has been accepted or is linked
below. Otherwise, an ADR captures the decision (link below).
ADR / RFC: docs/rfcs/draft-cassandra-backend.md
Breaking changes
By submitting this pull request, I confirm that my contribution is made under
the terms of the Apache License 2.0 and I agree to the Developer Certificate of
Origin (DCO). See CONTRIBUTING.md for details.