Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,12 @@ All notable changes to this project will be documented in this file.

Each entry lists the date and the crate versions that were released.

## 2026-09-01 — mqdb-cluster 0.4.11

### Fixed

- **TTL cleanup now routes through the replicated delete path (cluster mode).** The cluster TTL sweep removed expired rows with a raw in-memory `entities.remove` — it never released the row's unique-constraint guards (so an expired `unique` value became permanently unclaimable), emitted no change event (so watchers and cross-node subscribers never learned of the deletion), and did not replicate (each node independently mutated its own copy). The sweep now reaps each expired row this node is the data-partition **primary** for through the same path as a normal delete (`db_delete_prepare` + `db_commit` + `release_unique_for_deleted_record` + change event, plus resource-grant clearing), so the delete releases guards, emits a change event, and replicates to the replicas. Runs under the controller write lock so the scan and the deletes are atomic with respect to other controller messages. Companion to the agent-mode fix in 0.8.25; part of on-disconnect hold reclaim (`docs/design/hold-reclaim.md`). FK cascade on TTL expiry remains a follow-up (holds are leaf entities).

## 2026-09-01 — mqdb-agent 0.8.25

### Fixed
Expand Down
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion crates/mqdb-cluster/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "mqdb-cluster"
version = "0.4.10"
version = "0.4.11"
publish = false
edition.workspace = true
license = "AGPL-3.0-only"
Expand Down
53 changes: 19 additions & 34 deletions crates/mqdb-cluster/src/cluster/db/data_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -481,40 +481,25 @@ impl DbDataStore {
/// # Panics
/// Panics if the internal lock is poisoned.
#[must_use]
pub fn cleanup_expired_ttl(&self, now_secs: u64) -> Vec<(String, String)> {
let expired: Vec<(String, String, String)> = {
let entities = self.read_entities();
entities
.iter()
.filter_map(|(key, entity)| {
let data: serde_json::Value = serde_json::from_slice(&entity.data).ok()?;
let expires_at = data
.get("_expires_at")
.and_then(serde_json::Value::as_u64)?;
if expires_at <= now_secs {
Some((
key.clone(),
entity.entity_str().to_string(),
entity.id_str().to_string(),
))
} else {
None
}
})
.collect()
};

if expired.is_empty() {
return Vec::new();
}

let mut entities = self.write_entities();
let mut deleted = Vec::new();
for (key, entity_name, entity_id) in expired {
entities.remove(&key);
deleted.push((entity_name, entity_id));
}
deleted
/// Return `(entity, id)` for every row whose `_expires_at` has passed. Read-only:
/// the actual delete is routed through the replicated delete path (so it releases
/// unique guards, emits a change event, and replicates) by the caller.
pub fn scan_expired_ttl(&self, now_secs: u64) -> Vec<(String, String)> {
let entities = self.read_entities();
entities
.values()
.filter_map(|entity| {
let data: serde_json::Value = serde_json::from_slice(&entity.data).ok()?;
let expires_at = data
.get("_expires_at")
.and_then(serde_json::Value::as_u64)?;
if expires_at <= now_secs {
Some((entity.entity_str().to_string(), entity.id_str().to_string()))
} else {
None
}
})
.collect()
}
}

Expand Down
49 changes: 49 additions & 0 deletions crates/mqdb-cluster/src/cluster/db_handler/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -594,6 +594,55 @@ async fn delete_releases_unique_claim() {
);
}

#[tokio::test]
async fn ttl_reap_releases_unique_claim_and_removes_row() {
use super::super::db::ClusterConstraint;

let mut ctrl = setup_controller_all_partitions();
ctrl.constraint_add(&ClusterConstraint::unique("users", "uniq_email", "email"))
.await
.unwrap();

let data_bytes =
serde_json::to_vec(&serde_json::json!({"email": "a@x.com", "_expires_at": 100})).unwrap();
let rec = ctrl
.db_create("users", "u1", &data_bytes, 1000)
.await
.unwrap();
let value = serde_json::to_vec(&serde_json::json!("a@x.com")).unwrap();
ctrl.stores()
.db_unique
.reassert("users", "email", &value, "u1", rec.partition(), 1000);
assert!(ctrl.stores().unique_get("users", "email", &value).is_some());

let reaped = ctrl.reap_expired_ttl(1_000_000).await;
assert_eq!(reaped, 1, "the expired row must be reaped");

assert!(
ctrl.db_get("users", "u1").is_none(),
"the expired row must be removed"
);
assert!(
ctrl.stores().unique_get("users", "email", &value).is_none(),
"the TTL reap must release the committed unique claim so the value is reclaimable"
);
}

#[tokio::test]
async fn ttl_reap_leaves_unexpired_row() {
let mut ctrl = setup_controller_all_partitions();

let data_bytes =
serde_json::to_vec(&serde_json::json!({"v": 1, "_expires_at": 9_000_000_000u64})).unwrap();
ctrl.db_create("widgets", "w1", &data_bytes, 1000)
.await
.unwrap();

let reaped = ctrl.reap_expired_ttl(1000).await;
assert_eq!(reaped, 0, "a not-yet-expired row must not be reaped");
assert!(ctrl.db_get("widgets", "w1").is_some());
}

#[tokio::test]
async fn json_update_bypasses_ownership_when_no_sender() {
let node1 = NodeId::validated(1).unwrap();
Expand Down
102 changes: 76 additions & 26 deletions crates/mqdb-cluster/src/cluster/node_controller/db_ops.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,11 @@ use crate::cluster::db_handler::helpers::{parse_projection, validate_projection_

const CASCADE_ACK_TIMEOUT_SECS: u64 = 5;

/// Upper bound on rows reaped in a single TTL sweep pass, so a mass expiry cannot hold the
/// controller write lock long enough to stall heartbeats/Raft. The remainder is reaped on
/// the next sweep.
const TTL_REAP_MAX_PER_PASS: usize = 1024;

pub(crate) fn spawn_cascade_ack_waiter(
outbox: Option<crate::cluster::store_manager::outbox::ClusterOutbox>,
operation_id: String,
Expand Down Expand Up @@ -932,32 +937,11 @@ impl<T: ClusterTransport> NodeController<T> {
cascade: Option<&CascadeOutboxPayload>,
ack_receivers: Vec<oneshot::Receiver<bool>>,
) -> Vec<u8> {
match self.db_delete_prepare(entity, id) {
Ok((db_entity, write)) => {
let data: Value = serde_json::from_slice(&db_entity.data).unwrap_or(Value::Null);
let event = ChangeEvent::delete(entity.to_string(), id.to_string(), data.clone());
let outbox = build_change_event_outbox(&event);
if let Some(cas) = cascade {
self.db_commit_with_cascade(write, outbox.clone(), cas)
.await;
} else {
self.db_commit(write, outbox.clone()).await;
}
self.release_unique_for_deleted_record(entity, id, &data)
.await;
self.publish_and_deliver_change_event(event, &outbox.operation_id)
.await;
if self.ownership.owner_field(entity).is_some() {
self.clear_all_resource_grants(entity, id).await;
}
if let Some(cas) = cascade {
spawn_cascade_ack_waiter(
self.stores.cluster_outbox().cloned(),
cas.operation_id.clone(),
ack_receivers,
true,
);
}
match self
.delete_record_replicated(entity, id, cascade, ack_receivers)
.await
{
Ok(()) => {
let result = serde_json::json!({
"status": "ok",
"data": {"id": id, "deleted": true}
Expand All @@ -971,6 +955,72 @@ impl<T: ClusterTransport> NodeController<T> {
}
}

/// The effectful core of a delete, shared by the client-facing delete and the TTL
/// reap: prepare + replicate the delete, release the record's unique guards, emit the
/// change event, and clear resource grants. Returns `NotFound` if the row is gone.
async fn delete_record_replicated(
&mut self,
entity: &str,
id: &str,
cascade: Option<&CascadeOutboxPayload>,
ack_receivers: Vec<oneshot::Receiver<bool>>,
) -> Result<(), super::db::DbDataStoreError> {
let (db_entity, write) = self.db_delete_prepare(entity, id)?;
let data: Value = serde_json::from_slice(&db_entity.data).unwrap_or(Value::Null);
let event = ChangeEvent::delete(entity.to_string(), id.to_string(), data.clone());
let outbox = build_change_event_outbox(&event);
if let Some(cas) = cascade {
self.db_commit_with_cascade(write, outbox.clone(), cas)
.await;
} else {
self.db_commit(write, outbox.clone()).await;
}
self.release_unique_for_deleted_record(entity, id, &data)
.await;
self.publish_and_deliver_change_event(event, &outbox.operation_id)
.await;
if self.ownership.owner_field(entity).is_some() {
self.clear_all_resource_grants(entity, id).await;
}
if let Some(cas) = cascade {
spawn_cascade_ack_waiter(
self.stores.cluster_outbox().cloned(),
cas.operation_id.clone(),
ack_receivers,
true,
);
}
Ok(())
}

/// Reap TTL-expired rows this node is the data-partition primary for, routing each
/// through the replicated delete core so it releases the row's unique guards, emits a
/// change event, and replicates to the replicas. Non-primary nodes skip; the primary's
/// replicated delete reaches them. Runs under the controller write lock, so the scan
/// and the deletes are atomic w.r.t. other controller messages (a renewal cannot
/// interleave). At most `TTL_REAP_MAX_PER_PASS` rows are reaped per call so a mass
/// expiry cannot hold the lock long enough to stall heartbeats/Raft; the remainder is
/// reaped on the next sweep. Returns the number reaped.
pub(crate) async fn reap_expired_ttl(&mut self, now_secs: u64) -> usize {
let mut expired = self.stores.db_data.scan_expired_ttl(now_secs);
expired.truncate(TTL_REAP_MAX_PER_PASS);
let mut reaped = 0usize;
for (entity, id) in expired {
let partition = crate::cluster::db::data_partition(&entity, &id);
if !self.is_primary_for_partition(partition) {
continue;
}
if self
.delete_record_replicated(&entity, &id, None, Vec::new())
.await
.is_ok()
{
reaped += 1;
}
}
reaped
}

pub(crate) fn prepare_fk_side_effects(
&self,
results: &[super::fk::FkReverseLookupResult],
Expand Down
8 changes: 4 additions & 4 deletions crates/mqdb-cluster/src/cluster_agent/event_loop.rs
Original file line number Diff line number Diff line change
Expand Up @@ -485,10 +485,10 @@ impl ClusteredAgent {
let now_secs = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_secs());
let ctrl = self.controller.read().await;
let expired_ttl = ctrl.stores().db_data.cleanup_expired_ttl(now_secs);
if !expired_ttl.is_empty() {
info!(count = expired_ttl.len(), "cleaned up TTL-expired entities");
let mut ctrl = self.controller.write().await;
let reaped = ctrl.reap_expired_ttl(now_secs).await;
if reaped > 0 {
info!(count = reaped, "cleaned up TTL-expired entities");
}
}

Expand Down
14 changes: 14 additions & 0 deletions docs/design/hold-reclaim.md
Original file line number Diff line number Diff line change
Expand Up @@ -87,6 +87,20 @@ Known limitations of the agent sweep (PR 1a), deliberate / pre-existing:
guard (skipping is the safer of the two); the row is retried on the next sweep. A
permanent failure would retry (and re-log) each interval.

Cluster sweep (PR 1b) additional semantics:
- **Generic reap, no FK cascade.** `reap_expired_ttl` reaps any entity carrying
`_expires_at`, not only leaf `holds`, and (like the agent) does not run FK
cascade/set-null/restrict. TTL on a non-leaf FK parent therefore breaks referential
integrity silently (pre-existing; the old raw-remove sweep also skipped cascade). Use
TTL only on leaf entities until the cascade follow-up.
- **Primary-gated + capped.** Only the data-partition primary reaps; the delete replicates
to replicas (so replicas must not also reap — that would be N redundant deletes). A
replica that misses the fire-and-forget replicated delete keeps the expired row until it
is re-replicated or the partition fails over to it and it sweeps — the same reliability
class as any async-replicated write. At most `TTL_REAP_MAX_PER_PASS` (1024) rows are
reaped per pass so a mass expiry can't hold the controller write lock long enough to
stall heartbeats/Raft; the remainder drains on subsequent sweeps.

### B. Client CAS (`expected_version`, terminal, both paths)

- Shared: add `expected_version: Option<u64>` to `Request::Update`/`Request::Delete`
Expand Down
Loading