Skip to content

fix(storage-cassandra): consistently apply LWT to affected table-columns - #389

Open
jcshepherd wants to merge 7 commits into
ci/cassandra-integration-workflowfrom
fix/cassandra-lwt-paxos-corruption
Open

jcshepherd wants to merge 7 commits into
ci/cassandra-integration-workflowfrom
fix/cassandra-lwt-paxos-corruption

Conversation

@jcshepherd

Copy link
Copy Markdown
Collaborator

Mixing plain and LWT operations on the same table, column and partition can corrupt Cassandra's internal Paxos state, causing LWT operations to spin for the full Paxos timeout, potentially destabilizing and crashing the node.

All writes to tables that use IF NOT EXISTS on insert now use IF EXISTS (or conditional IF col = val) on update/delete. delete_user's logged batch replaced with sequential apply_lwt calls (atomicity gap documented in ADR-0021). Test keyspace accumulation causing JVM OOM fixed via PENDING_DROPS queue drained at test setup.

What

  1. Table-columns mutated via LWTs are consistently mutated by LWTs (not non-LWT operations which can corrupt LWT Paxos).
  2. In some cases, this meant backing out logged batch (atomic batch) operations in favor of consistent LWT behavior. This is tech debt, captured in ADR-0021.
  3. More effective clean-up between tests to eliminate warnings about table and keyspace counts.

Why

Under load, Paxos corruption will destabilize and eventually crash a node.

Testing done

Full suite of Cassandra tests (no changes outside storage-cassandra):
devtools/run-cassandra-tests -- --pytest --rust --comprehensive --rust-integration

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

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.

Mixing plain and LWT operations on the same table, column and partition
can corrupt Cassandra's internal Paxos state, causing LWT operations to spin
for the full Paxos timeout, potentially destabilizing and crashing the node.

All writes to tables that use IF NOT EXISTS on insert now use IF EXISTS
(or conditional IF col = val) on update/delete. delete_user's logged batch
replaced with sequential apply_lwt calls (atomicity gap documented in
ADR-0021). Test keyspace accumulation causing JVM OOM fixed via PENDING_DROPS
queue drained at test setup.

@robinnsc robinnsc 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.

A few general notes that don't anchor to changed lines:

delete_table.rs:142, worker_store.rs:158

Both of these paths still issue a plain DELETE FROM tables. Given that the row is created with IF NOT EXISTS and this PR moves the remaining tables writes over to LWT, these two deletes may want the same treatment, as otherwise the plain/LWT mixing on the tables row persists on what is probably its most frequently exercised path.

Scope of the "every IF NOT EXISTS table" claim

While working through the catalog I found a handful of other tables that appear to still mix plain and conditional writes on the same rows. backups_by_arn is inserted with IF NOT EXISTS and then updated or deleted plainly at backup_engine.rs:476, 496, 559, 670, 686; iam_users has a plain password update at users.rs:416; admin_users has plain delete and update paths at admin_store.rs:120, 145; and the bootstrapper inserts accounts and admin_users without a condition. On the data path, transaction_ledger.rs, the idempotency token delete at transactions.rs:931, ttl_expirations in ttl.rs, and the unconditional item deletes in delete_item.rs all sit alongside LWTs on the same rows. If the data-path tables are intentionally deferred, a note to that effect in the description or the ADR would help set expectations, and otherwise these could potentially be picked up in a follow-up.

let session = std::sync::Arc::clone(session);
crate::cassandra_util::apply_lwt(
&session,
&format!("DELETE FROM {ks}.iam_users WHERE account_id = ? AND user_name = ? IF EXISTS"),

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.

With the logged batch removed, delete_user now deletes the iam_users row first, and then proceeds through the tag, access key, policy, and group membership deletes as separate calls, each of which returns early on error. If any of the later steps fails, the user row is already gone while its access keys remain in the access_keys table. The credential lookup in credential_store.rs:117 resolves access keys by key ID alone and doesn't verify that the parent user still exists, and the credential cache invalidation in the server layer only runs once delete_user returns Ok, which means a partially failed delete could leave a removed user's credentials continuing to authenticate. It might be worth reordering the cascade so that access_keys is deleted first and iam_users last, or alternatively having the credential lookup confirm the parent user row exists, so that the "orphaned rows" described in the ADR are limited to inert catalog entries rather than live credentials.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Good call, especially on deleting the access key first. That's done, and the user delete is last which will enable the entire operation to make progress if retried.

for key_id in &key_ids {
batch = batch.add_query(
format!("DELETE FROM {ks}.access_keys WHERE access_key_id = ?"),
// access_keys uses plain INSERT — plain DELETE is safe

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.

Minor, but the "access_keys uses plain INSERT" comment may not be accurate. The import path at access_keys.rs:199 inserts with IF NOT EXISTS, and the existing delete_access_key at access_keys.rs:93 is already conditional, so under the rule this PR introduces, this delete would likely want to go through apply_lwt as well.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Thanks: create_access_key now uses IF NOT EXISTS, delete_user now uses IF EXISTS for the
access_keys delete, and the comment is corrected.

})?;
// IF EXISTS returns false only if the row was deleted between our
// earlier SELECT and this UPDATE — treat as TableNotFound.
if !crate::cassandra_util::lwt_applied(&result).unwrap_or(true) {

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.

Two observations here. First, the tables LWT now commits before the indexes batch runs, so a failure in the batch would leave the table row updated with no corresponding index rows. The recovery worker at workers.rs:489 only scans index rows that already exist in CREATING, so it wouldn't pick this case up. Since conditional batches can't span tables there may not be a clean way to restore atomicity, but it could be worth a comment or an ADR note. Second, the [applied] == false branch here returns TableNotFound without releasing the propagation holds taken earlier in the function, and since holds have no TTL, the affected indexes would remain skipped by workers indefinitely. This branch probably wants the same release loop as the error path. Relatedly, lwt_applied(...).unwrap_or(true) treats a result that fails to parse as applied, and surfacing that as an error may be the safer default.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Updated to correctly release the propagation holds, not silently swallow parse errors, left commetn on the (lack of) atomicity concern.

"UPDATE {catalog_keyspace}.tables SET stream_label = '{label}' \
WHERE account_id = '{account_id}' AND table_name = '{table_name}'"
));
// Update stream_label via LWT to avoid mixing plain/LWT writes on the

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 has a similar shape to the update_table change, in that the stream_label LWT and the shard insert batch are now independent, and the [applied] result of the LWT is not checked. If the shard batch fails after the label is written, ListStreams would report a stream from the catalog label while DescribeStream finds no shards for it. Could consider checking the applied result and, on batch failure, clearing the label so the two stay consistent.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Addressed.

use std::sync::Arc;

/// Keyspaces queued for async drop between tests.
static PENDING_DROPS: std::sync::Mutex<Vec<String>> = std::sync::Mutex::new(Vec::new());

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.

A few cleanup gaps worth being aware of. The queue is only drained at the start of the next setup_engine, so whatever the last tests in each binary queue up is never dropped; TestTable with owns_keyspace=false no longer drops anything at all; and the catalog tables and indexes rows are no longer cleaned up, since dropping the account keyspace doesn't touch the shared catalog keyspace. Over repeated runs this may bring back the table-count warnings this change was intended to address. One option would be draining the queue from a #[ctor::dtor] or at the end of the heavier tests, and keeping the catalog row deletes in place as IF EXISTS.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Acknowledged but these gaps will have minimal effect on accumulating test resources during a CI run: e.g., at most one keyspace per suite, and there are only three test cases that use own_keyspace=false (to allow an account to create multiple tables during a single test case). At this point, I don't think it's worth the effort to mitigate: the new clean-up logic solves >95% of the integration test resource cleanup problem.

cdrs_tokio::query_values!(changed_json.as_str(), "drain"),
&format!(
"UPDATE {keyspace}.{data_table} SET item_data = ? WHERE pk = ? \
IF prepared_txn_id = ?"

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 and the GSI variant below describe a stale writer whose batch was "pinned before the seal", and the GSI variant previously used USING TIMESTAMP to set that up. Both now write with IF prepared_txn_id = ?, which is necessary for an LWT, but it means the tests are now exercising the post-seal owner-aware path rather than the mixed-timestamp race the doc comments describe. It might be worth either updating the comments to match, or finding another way to retain the stale-timestamp coverage, for instance a plain write against a separate throwaway row.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Yeah, that was a bad change. Corrected.

…tion tests until unpredictable failures can be diagnosed
--rust ran cargo test workspace-wide (including ttl_integration) in the
cassandra-rust job; workspace unit tests are already covered by test.yml.
Drop it, leaving only --rust-integration (the server-facing test crate) where
ttl_integration is covered by continue-on-error until we can stabilize those
tests
@jcshepherd

Copy link
Copy Markdown
Collaborator Author

A few general notes that don't anchor to changed lines:

delete_table.rs:142, worker_store.rs:158

Both of these paths still issue a plain DELETE FROM tables. Given that the row is created with IF NOT EXISTS and this PR moves the remaining tables writes over to LWT, these two deletes may want the same treatment, as otherwise the plain/LWT mixing on the tables row persists on what is probably its most frequently exercised path.

Yup, good catch: will address.

Scope of the "every IF NOT EXISTS table" claim

While working through the catalog I found a handful of other tables that appear to still mix plain and conditional writes on the same rows. backups_by_arn is inserted with IF NOT EXISTS and then updated or deleted plainly at backup_engine.rs:476, 496, 559, 670, 686; iam_users has a plain password update at users.rs:416; admin_users has plain delete and update paths at admin_store.rs:120, 145;

Corrected.

and the bootstrapper inserts accounts and admin_users without a condition.

Corrected.

On the data path, transaction_ledger.rs, the idempotency token delete at transactions.rs:931, ttl_expirations in ttl.rs, and the unconditional item deletes in delete_item.rs all sit alongside LWTs on the same rows. If the data-path tables are intentionally deferred, a note to that effect in the description or the ADR would help set expectations, and otherwise these could potentially be picked up in a follow-up.

I think that's the plan. I'll update the PR description.

@robinnsc
robinnsc self-requested a review October 7, 2026 18:45
- delete_table.rs, worker_store.rs: converted plain DELETE FROM tables to
  LWT IF EXISTS
- backup_engine.rs: convert backups_by_arn deletes and status transitions
  to LWT; error-path deletes use IF EXISTS, state machine transitions use
  conditional IF backup_status = ?
- users.rs: convert password update to IF EXISTS; reorder delete_user
  operations so access_keys are deleted first (disables auth for the user
  immediately) and iam_users last (user deleted only after dependents are
  cleaned up, enabling graceful retries/recovery)
- access_keys.rs: convert create_access_key insert and delete_user's
  access_keys delete to LWTs
- bootstrapper.rs: convert accounts and admin_users inserts to LWTs
  eliminating the read-then-write TOCTOU race
- create_table.rs: check [applied] on stream_label LWT; compensate on
  shard batch failure by clearing the label
- update_table.rs: release propagation holds if [applied]=false; treat
  LWT parse failure as error rather than applied; document atomicity gap
- ttl_integration.rs: restore original plain/USING TIMESTAMP stale-writer
  semantics in changed-image and GSI recovery tests (IF prepared_txn_id
  changed what was being tested)

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.

2 participants