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
10 changes: 9 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,15 @@ 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
## 2026-09-05 — mqdb-core 0.7.12, mqdb-agent 0.8.26, mqdb-cluster 0.4.12, mqdb-vault 0.1.4, mqdb-wasm 0.3.6

### Added

- **Client-facing compare-and-set on update and delete.** A write may now carry a reserved `_expected_version` in its payload; the server rejects the operation if the row's current `_version` differs from it. This is the second primitive of on-disconnect hold reclaim (`docs/design/hold-reclaim.md`): a janitor reclaims an abandoned hold with `delete` guarded on the version it last observed, so a holder that renewed in the meantime is never clobbered. The check is **terminal** — it fails with `PreconditionFailed` (transport code 412), which the write-retry loops treat as non-retryable, distinct from an optimistic-concurrency `Conflict` that is retried. Enforced identically in agent mode (`update_with_expected`/`delete_with_expected`) and in both cluster write paths (client-dispatched local-primary and forwarded). Omitting `_expected_version` leaves behavior unchanged.

### Notes

- `_expected_version` is stripped from the payload before validation and storage, so it never persists as a field. Embedded (WASM) mode ignores it (no versioned write path).

### Fixed

Expand Down
8 changes: 4 additions & 4 deletions 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-agent/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "mqdb-agent"
version = "0.8.25"
version = "0.8.26"
edition.workspace = true
license = "Apache-2.0"
authors.workspace = true
Expand Down
188 changes: 177 additions & 11 deletions crates/mqdb-agent/src/database/crud.rs
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,33 @@ impl Database {
fields: Value,
update_constraint_data: Option<(Value, Value)>,
caller: &CallerContext<'_>,
) -> Result<Value> {
self.update_with_expected(
entity_name,
id,
fields,
update_constraint_data,
None,
caller,
)
.await
}

/// Like [`update`](Self::update) but rejects with [`Error::PreconditionFailed`] (a
/// terminal, non-retried error) if the row's current `_version` differs from
/// `expected_version` — a client-facing compare-and-set.
///
/// # Errors
/// Returns an error if the entity is not found, the version precondition fails,
/// validation fails, or storage fails.
pub async fn update_with_expected(
&self,
entity_name: String,
id: String,
fields: Value,
update_constraint_data: Option<(Value, Value)>,
expected_version: Option<u64>,
caller: &CallerContext<'_>,
) -> Result<Value> {
let retryable = update_constraint_data.is_none();
let mut attempt = 1;
Expand All @@ -186,6 +213,7 @@ impl Database {
&id,
fields.clone(),
update_constraint_data.clone(),
expected_version,
caller,
)
.await;
Expand All @@ -204,6 +232,7 @@ impl Database {
id: &str,
fields: Value,
update_constraint_data: Option<(Value, Value)>,
expected_version: Option<u64>,
caller: &CallerContext<'_>,
) -> Result<Value> {
let key = keys::encode_data_key(entity_name, id);
Expand All @@ -227,6 +256,13 @@ impl Database {
.get("_version")
.and_then(Value::as_u64)
.unwrap_or(0);
if let Some(expected) = expected_version
&& existing_version != expected
{
return Err(Error::PreconditionFailed(format!(
"{entity_name}/{id}: expected _version {expected}, found {existing_version}"
)));
}
if let Value::Object(ref mut obj) = updated_data {
obj.insert(
"_version".to_string(),
Expand Down Expand Up @@ -308,18 +344,35 @@ impl Database {
client_id: Option<&str>,
scope_config: &ScopeConfig,
ownership: &OwnershipConfig,
) -> Result<()> {
let caller = CallerContext {
sender,
client_id,
scope_config,
};
self.delete_with_expected(entity_name, id, &caller, ownership, None)
.await
}

/// Like [`delete`](Self::delete) but rejects with [`Error::PreconditionFailed`] (a
/// terminal, non-retried error) if the row's current `_version` differs from
/// `expected_version` — a client-facing compare-and-set.
///
/// # Errors
/// Returns an error if the entity is not found, the version precondition fails, or
/// constraint validation fails.
pub async fn delete_with_expected(
&self,
entity_name: String,
id: String,
caller: &CallerContext<'_>,
ownership: &OwnershipConfig,
expected_version: Option<u64>,
) -> Result<()> {
let mut attempt = 1;
loop {
let result = self
.try_delete_once(
&entity_name,
&id,
sender,
client_id,
scope_config,
ownership,
)
.try_delete_once(&entity_name, &id, caller, ownership, expected_version)
.await;
match result {
Err(Error::Conflict(_)) if attempt < MAX_WRITE_ATTEMPTS => {
Expand All @@ -335,13 +388,16 @@ impl Database {
&self,
entity_name: &str,
id: &str,
sender: Option<&str>,
client_id: Option<&str>,
scope_config: &ScopeConfig,
caller: &CallerContext<'_>,
ownership: &OwnershipConfig,
expected_version: Option<u64>,
) -> Result<()> {
use mqdb_core::constraint::{DeleteOperation, OwnershipContext};

let sender = caller.sender;
let client_id = caller.client_id;
let scope_config = caller.scope_config;

let key = keys::encode_data_key(entity_name, id);
let existing_data = self.storage.get(&key)?.ok_or_else(|| Error::NotFound {
entity: entity_name.to_string(),
Expand All @@ -351,6 +407,19 @@ impl Database {
let existing_entity =
Entity::deserialize(entity_name.to_string(), id.to_string(), &existing_data)?;

if let Some(expected) = expected_version {
let existing_version = existing_entity
.data
.get("_version")
.and_then(Value::as_u64)
.unwrap_or(0);
if existing_version != expected {
return Err(Error::PreconditionFailed(format!(
"{entity_name}/{id}: expected _version {expected}, found {existing_version}"
)));
}
}

let ownership_ctx = sender
.filter(|_| !ownership.is_empty())
.map(|s| OwnershipContext {
Expand Down Expand Up @@ -1049,6 +1118,103 @@ mod concurrency_tests {
);
}

#[tokio::test(flavor = "multi_thread")]
async fn update_and_delete_with_expected_version_are_terminal_cas() {
let (_tmp, db) = test_db().await;
let scope = ScopeConfig::default();
let caller = CallerContext {
sender: None,
client_id: None,
scope_config: &scope,
};

db.create(
"doc".to_string(),
json!({ "id": "d", "n": 0 }),
None,
None,
None,
&scope,
)
.await
.unwrap();

let stale = db
.update_with_expected(
"doc".to_string(),
"d".to_string(),
json!({ "n": 1 }),
None,
Some(99),
&caller,
)
.await;
assert!(
matches!(stale, Err(mqdb_core::error::Error::PreconditionFailed(_))),
"a wrong expected version must be a terminal precondition failure, got {stale:?}"
);

let ok = db
.update_with_expected(
"doc".to_string(),
"d".to_string(),
json!({ "n": 1 }),
None,
Some(1),
&caller,
)
.await
.expect("matching expected version must succeed");
assert_eq!(ok["_version"].as_u64(), Some(2));

let stale2 = db
.update_with_expected(
"doc".to_string(),
"d".to_string(),
json!({ "n": 2 }),
None,
Some(1),
&caller,
)
.await;
assert!(
matches!(stale2, Err(mqdb_core::error::Error::PreconditionFailed(_))),
"the version bumped, so the old expected version is now stale"
);

let ownership = OwnershipConfig::default();
let caller = CallerContext {
sender: None,
client_id: None,
scope_config: &scope,
};
let del_stale = db
.delete_with_expected(
"doc".to_string(),
"d".to_string(),
&caller,
&ownership,
Some(1),
)
.await;
assert!(
matches!(
del_stale,
Err(mqdb_core::error::Error::PreconditionFailed(_))
),
"delete with a wrong expected version must be rejected, got {del_stale:?}"
);
db.delete_with_expected(
"doc".to_string(),
"d".to_string(),
&caller,
&ownership,
Some(2),
)
.await
.expect("delete with the matching version must succeed");
}

#[tokio::test(flavor = "multi_thread")]
async fn add_unique_constraint_rejects_existing_duplicates() {
let (_tmp, db) = test_db().await;
Expand Down
23 changes: 20 additions & 3 deletions crates/mqdb-agent/src/transport_execute.rs
Original file line number Diff line number Diff line change
Expand Up @@ -138,6 +138,7 @@ impl Database {
entity,
id,
mut fields,
expected_version,
} => {
match ownership.evaluate(&entity, sender) {
OwnershipDecision::Check {
Expand Down Expand Up @@ -191,14 +192,25 @@ impl Database {
scope_config,
};
match self
.update(entity, id, fields, update_constraint, &caller)
.update_with_expected(
entity,
id,
fields,
update_constraint,
expected_version,
&caller,
)
.await
{
Ok(v) => Response::ok(v),
Err(e) => e.into(),
}
}
Request::Delete { entity, id } => {
Request::Delete {
entity,
id,
expected_version,
} => {
match ownership.evaluate(&entity, sender) {
OwnershipDecision::Check {
owner_field,
Expand Down Expand Up @@ -233,8 +245,13 @@ impl Database {
let id_clone = id.clone();
let shareable = ownership.owner_field(&entity).is_some();
let entity_clone = entity.clone();
let caller = CallerContext {
sender,
client_id,
scope_config,
};
match self
.delete(entity, id, sender, client_id, scope_config, ownership)
.delete_with_expected(entity, id, &caller, ownership, expected_version)
.await
{
Ok(()) => {
Expand Down
Loading
Loading