Skip to content

Feat/storage cassandra in tree - #338

Open
jcshepherd wants to merge 59 commits into
mainfrom
feat/storage-cassandra-in-tree
Open

jcshepherd wants to merge 59 commits into
mainfrom
feat/storage-cassandra-in-tree

Conversation

@jcshepherd

@jcshepherd jcshepherd commented Sep 10, 2026 •

Copy link
Copy Markdown
Collaborator

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

  • I have read CONTRIBUTING.md
  • All tests pass (cargo test --workspace)
  • Code is formatted (cargo fmt --check)
  • Clippy is clean (cargo clippy -- -W clippy::pedantic)
  • I have added or updated tests for new functionality
  • I have updated documentation if behavior changed
  • Breaking changes are noted below (if any)
  • If this changes the wire protocol, Storage trait, auth model, on-disk
    format, 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.

Comment thread crates/storage-cassandra/src/bootstrapper.rs Fixed
Comment thread crates/storage-cassandra/src/bootstrapper.rs Fixed
Comment thread crates/storage-cassandra/src/bootstrapper.rs Fixed
Comment thread crates/storage-cassandra/src/bootstrapper.rs Fixed
Comment thread crates/storage-cassandra/src/bootstrapper.rs Fixed
Comment thread crates/storage-cassandra/src/bootstrapper.rs Fixed

@amrith amrith left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Would you take a look at the failing checks - or are they related to 341?

@LeeroyHannigan LeeroyHannigan left a comment •

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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:

  1. CI on this head is red on clippy (three too_many_arguments and one from_str naming 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.
  2. 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} \

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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>,

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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} \

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 = ?");

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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(

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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"),

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread crates/storage/src/config.rs Outdated
///
/// 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

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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);

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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()],

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread tests/test_transaction_operations.py Outdated
@@ -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):

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

@jcshepherd jcshepherd mentioned this pull request Sep 15, 2026
8 tasks
@jcshepherd
jcshepherd changed the base branch from main to fix/hlc-instance-id-wiring September 15, 2026 19:15
@jcshepherd
jcshepherd force-pushed the feat/storage-cassandra-in-tree branch from f6af48f to 8151972 Compare September 15, 2026 20:37
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.
…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.
@jcshepherd
jcshepherd changed the base branch from fix/hlc-instance-id-wiring to fix/streams-sequence-normalization September 15, 2026 22:46
jcshepherd and others added 8 commits September 15, 2026 23:20
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.
robinnsc and others added 3 commits September 18, 2026 18:17
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.
test: Porting over the Cassandra Repo integration suite
Base automatically changed from fix/streams-sequence-normalization to main September 22, 2026 00:33
…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.
…ts to avoid unpredictable exhaustion of the connection pool
@jcshepherd
jcshepherd added this pull request to stack #374 September 28, 2026 21:40
…, query scan follows Amazon DDB key sorting

rules in both directions (mirroring #362 for PostgreSQL and SQLite).

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants