diff --git a/CHANGELOG.md b/CHANGELOG.md index 3a637db0..a07b91c8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/Cargo.lock b/Cargo.lock index 98548805..044eca36 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1420,7 +1420,7 @@ dependencies = [ [[package]] name = "mqdb-agent" -version = "0.8.25" +version = "0.8.26" dependencies = [ "arc-swap", "argon2", @@ -1482,7 +1482,7 @@ dependencies = [ [[package]] name = "mqdb-cluster" -version = "0.4.11" +version = "0.4.12" dependencies = [ "arc-swap", "bebytes", @@ -1515,7 +1515,7 @@ dependencies = [ [[package]] name = "mqdb-core" -version = "0.7.11" +version = "0.7.12" dependencies = [ "arc-swap", "bebytes", @@ -1537,7 +1537,7 @@ dependencies = [ [[package]] name = "mqdb-vault" -version = "0.1.3" +version = "0.1.4" dependencies = [ "base64", "mqdb-agent", diff --git a/crates/mqdb-agent/Cargo.toml b/crates/mqdb-agent/Cargo.toml index eccf287e..e920097b 100644 --- a/crates/mqdb-agent/Cargo.toml +++ b/crates/mqdb-agent/Cargo.toml @@ -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 diff --git a/crates/mqdb-agent/src/database/crud.rs b/crates/mqdb-agent/src/database/crud.rs index 3c52bd30..dca0e8a1 100644 --- a/crates/mqdb-agent/src/database/crud.rs +++ b/crates/mqdb-agent/src/database/crud.rs @@ -176,6 +176,33 @@ impl Database { fields: Value, update_constraint_data: Option<(Value, Value)>, caller: &CallerContext<'_>, + ) -> Result { + 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, + caller: &CallerContext<'_>, ) -> Result { let retryable = update_constraint_data.is_none(); let mut attempt = 1; @@ -186,6 +213,7 @@ impl Database { &id, fields.clone(), update_constraint_data.clone(), + expected_version, caller, ) .await; @@ -204,6 +232,7 @@ impl Database { id: &str, fields: Value, update_constraint_data: Option<(Value, Value)>, + expected_version: Option, caller: &CallerContext<'_>, ) -> Result { let key = keys::encode_data_key(entity_name, id); @@ -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(), @@ -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, ) -> 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 => { @@ -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, ) -> 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(), @@ -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 { @@ -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; diff --git a/crates/mqdb-agent/src/transport_execute.rs b/crates/mqdb-agent/src/transport_execute.rs index 230ed45d..8ebb32fd 100644 --- a/crates/mqdb-agent/src/transport_execute.rs +++ b/crates/mqdb-agent/src/transport_execute.rs @@ -138,6 +138,7 @@ impl Database { entity, id, mut fields, + expected_version, } => { match ownership.evaluate(&entity, sender) { OwnershipDecision::Check { @@ -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, @@ -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(()) => { diff --git a/crates/mqdb-agent/tests/integration_test.rs b/crates/mqdb-agent/tests/integration_test.rs index e1cc6f16..efc3b31c 100644 --- a/crates/mqdb-agent/tests/integration_test.rs +++ b/crates/mqdb-agent/tests/integration_test.rs @@ -1747,6 +1747,7 @@ async fn test_share_grant_gates_read_update_delete() { entity: "diagrams".into(), id: diagram_id.clone(), fields: json!({ "title": title }), + expected_version: None, }, Some(sender), None, @@ -1806,6 +1807,7 @@ async fn test_share_grant_gates_read_update_delete() { Request::Delete { entity: "diagrams".into(), id: diagram_id.clone(), + expected_version: None, }, Some("bob"), None, @@ -1824,6 +1826,7 @@ async fn test_share_grant_gates_read_update_delete() { Request::Delete { entity: "diagrams".into(), id: diagram_id.clone(), + expected_version: None, }, Some("alice"), None, @@ -2296,6 +2299,7 @@ async fn test_delete_clears_shares_no_inheritance() { Request::Delete { entity: "diagrams".into(), id: id.clone(), + expected_version: None, }, Some("alice"), None, @@ -2389,6 +2393,7 @@ async fn test_child_enforcement_derives_from_parent() { entity: "nodes".into(), id: node_id.clone(), fields: json!({"label": "edited"}), + expected_version: None, }, Some(sender), None, @@ -2403,6 +2408,7 @@ async fn test_child_enforcement_derives_from_parent() { Request::Delete { entity: "nodes".into(), id: node_id.clone(), + expected_version: None, }, Some(sender), None, @@ -2561,6 +2567,7 @@ async fn test_child_reparent_blocked() { entity: "nodes".into(), id: node_id.clone(), fields: json!({"diagramId": d2, "label": "moved"}), + expected_version: None, }, Some("bob"), None, @@ -2650,6 +2657,7 @@ async fn test_unconfigured_child_unrestricted() { entity: "widgets".into(), id: widget_id.clone(), fields: json!({"k": "v2"}), + expected_version: None, }, Some("bob"), None, @@ -3461,6 +3469,7 @@ async fn test_reshare_legacy_grant_no_duplicate_and_demotes() { entity: "diagrams".into(), id: "d1".into(), fields: json!({"title": "hijack"}), + expected_version: None, }, Some("cid-bob"), None, @@ -3694,6 +3703,7 @@ async fn test_cascade_delete_events_carry_recipients() { Request::Delete { entity: "diagrams".into(), id: diagram_id.clone(), + expected_version: None, }, Some("alice"), None, diff --git a/crates/mqdb-cluster/Cargo.toml b/crates/mqdb-cluster/Cargo.toml index 99a11758..0e81dcfe 100644 --- a/crates/mqdb-cluster/Cargo.toml +++ b/crates/mqdb-cluster/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-cluster" -version = "0.4.11" +version = "0.4.12" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/crates/mqdb-cluster/src/cluster/db_handler/json_ops.rs b/crates/mqdb-cluster/src/cluster/db_handler/json_ops.rs index 503100bd..91d30fa4 100644 --- a/crates/mqdb-cluster/src/cluster/db_handler/json_ops.rs +++ b/crates/mqdb-cluster/src/cluster/db_handler/json_ops.rs @@ -28,6 +28,18 @@ pub(super) enum JsonOpResult { PendingFkDelete(Box), } +/// The reserved `_expected_version` compare-and-set precondition carried in a delete +/// request body (deletes otherwise ignore their payload). Returns the client-facing error +/// message when the field is present but malformed, so the CAS can never be silently skipped. +fn expected_version_from_delete_payload(payload: &[u8]) -> Result, String> { + if payload.is_empty() { + return Ok(None); + } + let mut value: Value = + serde_json::from_slice(payload).map_err(|e| format!("invalid JSON payload: {e}"))?; + mqdb_core::protocol::take_expected_version(&mut value).map_err(|e| e.to_string()) +} + /// Outcome of resolving an identity-mode (email) grantee before a share/unshare /// is routed to the resource primary. #[cfg(feature = "http-api")] @@ -264,6 +276,10 @@ impl DbRequestHandler { { return JsonOpResult::Response(err); } + let expected_version = match expected_version_from_delete_payload(payload) { + Ok(v) => v, + Err(msg) => return JsonOpResult::Response(Self::json_error(400, &msg)), + }; let partition = data_partition(entity, id); if !controller.is_primary_for_partition(partition) { let forwarded = controller @@ -272,7 +288,7 @@ impl DbRequestHandler { JsonDbOp::Delete, entity, Some(id), - &[], + payload, response_topic, correlation_data, sender, @@ -287,7 +303,7 @@ impl DbRequestHandler { )) }; } - self.handle_json_delete(controller, entity, id, mqtt_ctx) + self.handle_json_delete(controller, entity, id, expected_version, mqtt_ctx) .await } DbTopicOperation::JsonList { entity } => { @@ -960,6 +976,10 @@ impl DbRequestHandler { if let Some(obj) = updates.as_object_mut() { obj.remove("__mqdb_fk_expected"); } + let expected_version = match mqdb_core::protocol::take_expected_version(&mut updates) { + Ok(v) => v, + Err(e) => return JsonOpResult::Response(Self::json_error(400, &e.to_string())), + }; let vault_crypto = self.resolve_vault_crypto(entity, sender); let (old_data, merged_data) = match self.vault_merge_with_existing( controller, @@ -972,6 +992,19 @@ impl DbRequestHandler { Err(response) => return JsonOpResult::Response(response), }; + if let Some(expected) = expected_version { + let current = old_data + .get("_version") + .and_then(Value::as_u64) + .unwrap_or(0); + if current != expected { + return JsonOpResult::Response(Self::json_error( + 412, + &format!("{entity}/{id}: expected _version {expected}, found {current}"), + )); + } + } + if let Some(err) = Self::validate_against_schema(controller, entity, &merged_data) { return JsonOpResult::Response(err); } @@ -1201,6 +1234,7 @@ impl DbRequestHandler { controller: &mut NodeController, entity: &str, id: &str, + expected_version: Option, mqtt_ctx: &super::MqttRequestContext<'_>, ) -> JsonOpResult { let sender = mqtt_ctx.sender; @@ -1208,6 +1242,21 @@ impl DbRequestHandler { let response_topic = mqtt_ctx.response_topic.unwrap_or(""); let correlation_data = mqtt_ctx.correlation_data; + if let Some(expected) = expected_version + && let Some(existing) = controller.db_get(entity, id) + { + let current = serde_json::from_slice::(&existing.data) + .ok() + .and_then(|d| d.get("_version").and_then(Value::as_u64)) + .unwrap_or(0); + if current != expected { + return JsonOpResult::Response(Self::json_error( + 412, + &format!("{entity}/{id}: expected _version {expected}, found {current}"), + )); + } + } + let (local_results, pending_remote) = match controller.start_fk_reverse_lookup(entity, id, sender).await { Ok(pair) => pair, diff --git a/crates/mqdb-cluster/src/cluster/db_handler/tests.rs b/crates/mqdb-cluster/src/cluster/db_handler/tests.rs index db2a7855..824dc4a4 100644 --- a/crates/mqdb-cluster/src/cluster/db_handler/tests.rs +++ b/crates/mqdb-cluster/src/cluster/db_handler/tests.rs @@ -643,6 +643,84 @@ async fn ttl_reap_leaves_unexpired_row() { assert!(ctrl.db_get("widgets", "w1").is_some()); } +#[tokio::test] +async fn cluster_update_delete_expected_version_cas() { + let node1 = NodeId::validated(1).unwrap(); + let handler = DbRequestHandler::new(node1); + let mut ctrl = setup_controller_all_partitions(); + + let data = serde_json::to_vec(&serde_json::json!({"n": 0, "_version": 1})).unwrap(); + ctrl.db_create("docs", "d", &data, 1000).await.unwrap(); + + let ctx = MqttRequestContext { + response_topic: Some("resp/t"), + correlation_data: None, + sender: None, + client_id: None, + }; + + let stale = serde_json::to_vec(&serde_json::json!({"n": 1, "_expected_version": 99})).unwrap(); + let r = handler + .handle_publish(&mut ctrl, "$DB/docs/d/update", &stale, &ctx) + .await + .unwrap(); + assert_eq!(parse_json_response(&r.payload)["code"], 412); + + let ok = serde_json::to_vec(&serde_json::json!({"n": 1, "_expected_version": 1})).unwrap(); + let r = handler + .handle_publish(&mut ctrl, "$DB/docs/d/update", &ok, &ctx) + .await + .unwrap(); + assert_eq!(parse_json_response(&r.payload)["status"], "ok"); + + let r = handler + .handle_publish(&mut ctrl, "$DB/docs/d/update", &ok, &ctx) + .await + .unwrap(); + assert_eq!( + parse_json_response(&r.payload)["code"], + 412, + "the version bumped, so the old expected version is now stale" + ); + + let bad_update = + serde_json::to_vec(&serde_json::json!({"n": 2, "_expected_version": "2"})).unwrap(); + let r = handler + .handle_publish(&mut ctrl, "$DB/docs/d/update", &bad_update, &ctx) + .await + .unwrap(); + assert_eq!( + parse_json_response(&r.payload)["code"], + 400, + "a malformed _expected_version must be rejected, not silently applied" + ); + + let bad_delete = serde_json::to_vec(&serde_json::json!({"_expected_version": "2"})).unwrap(); + let r = handler + .handle_publish(&mut ctrl, "$DB/docs/d/delete", &bad_delete, &ctx) + .await + .unwrap(); + assert_eq!( + parse_json_response(&r.payload)["code"], + 400, + "a malformed _expected_version must not silently disable the delete CAS" + ); + + let dstale = serde_json::to_vec(&serde_json::json!({"_expected_version": 1})).unwrap(); + let r = handler + .handle_publish(&mut ctrl, "$DB/docs/d/delete", &dstale, &ctx) + .await + .unwrap(); + assert_eq!(parse_json_response(&r.payload)["code"], 412); + + let dok = serde_json::to_vec(&serde_json::json!({"_expected_version": 2})).unwrap(); + let r = handler + .handle_publish(&mut ctrl, "$DB/docs/d/delete", &dok, &ctx) + .await + .unwrap(); + assert_eq!(parse_json_response(&r.payload)["data"]["deleted"], true); +} + #[tokio::test] async fn json_update_bypasses_ownership_when_no_sender() { let node1 = NodeId::validated(1).unwrap(); diff --git a/crates/mqdb-cluster/src/cluster/node_controller/db_ops.rs b/crates/mqdb-cluster/src/cluster/node_controller/db_ops.rs index a1086e2a..d8fe05ba 100644 --- a/crates/mqdb-cluster/src/cluster/node_controller/db_ops.rs +++ b/crates/mqdb-cluster/src/cluster/node_controller/db_ops.rs @@ -23,6 +23,19 @@ const CASCADE_ACK_TIMEOUT_SECS: u64 = 5; /// the next sweep. const TTL_REAP_MAX_PER_PASS: usize = 1024; +/// Parse the reserved `_expected_version` compare-and-set precondition from a delete +/// request payload (the update path removes it from the parsed object instead). Returns the +/// client-facing error message when the field is present but malformed, so the CAS can never +/// be silently skipped. +fn expected_version_from_payload(payload: &[u8]) -> Result, String> { + if payload.is_empty() { + return Ok(None); + } + let mut value: Value = + serde_json::from_slice(payload).map_err(|e| format!("invalid JSON payload: {e}"))?; + mqdb_core::protocol::take_expected_version(&mut value).map_err(|e| e.to_string()) +} + pub(crate) fn spawn_cascade_ack_waiter( outbox: Option, operation_id: String, @@ -676,10 +689,33 @@ impl NodeController { .and_then(|obj| obj.remove("__mqdb_fk_expected")) .and_then(|v| v.as_str().map(String::from)); + let expected_version = match mqdb_core::protocol::take_expected_version(&mut data) { + Ok(v) => v, + Err(e) => return (Self::json_error(400, &e.to_string()), None), + }; + let (old_data, merged_data) = if let Some(existing) = self.db_get(entity, id) { let existing_data: Value = serde_json::from_slice(&existing.data) .unwrap_or(Value::Object(serde_json::Map::new())); + if let Some(expected) = expected_version { + let current = existing_data + .get("_version") + .and_then(Value::as_u64) + .unwrap_or(0); + if current != expected { + return ( + Self::json_error( + 412, + &format!( + "{entity}/{id}: expected _version {expected}, found {current}" + ), + ), + None, + ); + } + } + if let Some(ref expected) = fk_expected && let Some((field, value)) = expected.split_once('=') { @@ -879,6 +915,28 @@ impl NodeController { request: &JsonDbRequest, from: NodeId, ) -> (Vec, Option) { + let expected_version = match expected_version_from_payload(&request.payload) { + Ok(v) => v, + Err(msg) => return (Self::json_error(400, &msg), None), + }; + if let Some(expected) = expected_version + && let Some(existing) = self.db_get(entity, id) + { + let current = serde_json::from_slice::(&existing.data) + .ok() + .and_then(|d| d.get("_version").and_then(Value::as_u64)) + .unwrap_or(0); + if current != expected { + return ( + Self::json_error( + 412, + &format!("{entity}/{id}: expected _version {expected}, found {current}"), + ), + None, + ); + } + } + let (local_results, pending_remote) = match self .start_fk_reverse_lookup(entity, id, request.sender.as_deref()) .await diff --git a/crates/mqdb-core/Cargo.toml b/crates/mqdb-core/Cargo.toml index 518d32e9..c6ba93a7 100644 --- a/crates/mqdb-core/Cargo.toml +++ b/crates/mqdb-core/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-core" -version = "0.7.11" +version = "0.7.12" edition.workspace = true license = "Apache-2.0" authors.workspace = true diff --git a/crates/mqdb-core/src/error.rs b/crates/mqdb-core/src/error.rs index 6d1a6756..99294ae4 100644 --- a/crates/mqdb-core/src/error.rs +++ b/crates/mqdb-core/src/error.rs @@ -87,6 +87,9 @@ pub enum Error { #[error("forbidden: {0}")] Forbidden(String), + #[error("precondition failed: {0}")] + PreconditionFailed(String), + #[error("cascade blocked: cannot delete {0} - cross-owned entity has non-nullable FK field")] CascadeBlocked(Box), } diff --git a/crates/mqdb-core/src/protocol/mod.rs b/crates/mqdb-core/src/protocol/mod.rs index a6faf5a7..5ec7fc4f 100644 --- a/crates/mqdb-core/src/protocol/mod.rs +++ b/crates/mqdb-core/src/protocol/mod.rs @@ -15,6 +15,8 @@ pub enum ProtocolError { InvalidPayload(#[from] serde_json::Error), #[error("payload too large ({0} bytes, max {MAX_PAYLOAD_SIZE})")] PayloadTooLarge(usize), + #[error("invalid _expected_version (must be a non-negative integer): {0}")] + InvalidExpectedVersion(String), } #[derive(Debug, Clone, Copy, PartialEq, Eq)] @@ -280,10 +282,31 @@ pub fn parse_db_topic(topic: &str) -> Option { } } +/// Remove and return the reserved `_expected_version` compare-and-set precondition from a +/// request payload, so it is never merged into the record as a data field. +/// +/// # Errors +/// Returns [`ProtocolError::InvalidExpectedVersion`] if the field is present but is not a +/// non-negative integer — a malformed precondition must fail loudly, not silently disable +/// the compare-and-set. A missing field or an explicit `null` yields `Ok(None)`. +pub fn take_expected_version(data: &mut Value) -> Result, ProtocolError> { + let Value::Object(obj) = data else { + return Ok(None); + }; + match obj.remove("_expected_version") { + None | Some(Value::Null) => Ok(None), + Some(v) => v + .as_u64() + .map(Some) + .ok_or_else(|| ProtocolError::InvalidExpectedVersion(v.to_string())), + } +} + /// Builds a database request from an operation descriptor and payload. /// /// # Errors -/// Returns an error if JSON deserialization fails or a required ID is missing. +/// Returns an error if JSON deserialization fails, a required ID is missing, or the +/// `_expected_version` precondition is malformed. pub fn build_request(op: DbOperation, payload: &[u8]) -> Result { if payload.len() > MAX_PAYLOAD_SIZE { return Err(ProtocolError::PayloadTooLarge(payload.len())); @@ -311,17 +334,23 @@ pub fn build_request(op: DbOperation, payload: &[u8]) -> Result { let id = op.id.ok_or(ProtocolError::MissingId(DbOp::Update))?; + let mut fields = data; + let expected_version = take_expected_version(&mut fields)?; Ok(Request::Update { entity: op.entity, id, - fields: data, + fields, + expected_version, }) } DbOp::Delete => { let id = op.id.ok_or(ProtocolError::MissingId(DbOp::Delete))?; + let mut data = data; + let expected_version = take_expected_version(&mut data)?; Ok(Request::Delete { entity: op.entity, id, + expected_version, }) } DbOp::List => { @@ -466,6 +495,52 @@ fn extract_list_options(data: &Value) -> ListOptions { mod tests { use super::*; + #[test] + fn take_expected_version_extracts_and_strips() { + let mut data = serde_json::json!({"n": 1, "_expected_version": 7}); + let got = take_expected_version(&mut data).unwrap(); + assert_eq!(got, Some(7)); + assert!(data.get("_expected_version").is_none()); + assert_eq!(data.get("n").and_then(Value::as_u64), Some(1)); + } + + #[test] + fn take_expected_version_absent_or_null_is_none() { + let mut absent = serde_json::json!({"n": 1}); + assert_eq!(take_expected_version(&mut absent).unwrap(), None); + let mut null = serde_json::json!({"_expected_version": null}); + assert_eq!(take_expected_version(&mut null).unwrap(), None); + } + + #[test] + fn take_expected_version_rejects_malformed() { + for bad in [ + serde_json::json!({"_expected_version": "7"}), + serde_json::json!({"_expected_version": 7.5}), + serde_json::json!({"_expected_version": -1}), + serde_json::json!({"_expected_version": true}), + ] { + let mut data = bad.clone(); + let err = take_expected_version(&mut data); + assert!( + matches!(err, Err(ProtocolError::InvalidExpectedVersion(_))), + "a present-but-non-integer precondition must error, not silently skip: {bad:?}" + ); + assert!( + data.get("_expected_version").is_none(), + "the field is still stripped even on rejection" + ); + } + } + + #[test] + fn build_request_delete_rejects_malformed_expected_version() { + let op = parse_db_topic("$DB/users/123/delete").unwrap(); + let payload = serde_json::to_vec(&serde_json::json!({"_expected_version": "5"})).unwrap(); + let err = build_request(op, &payload); + assert!(matches!(err, Err(ProtocolError::InvalidExpectedVersion(_)))); + } + #[test] fn test_parse_db_topic_create() { let op = parse_db_topic("$DB/users/create").unwrap(); diff --git a/crates/mqdb-core/src/transport.rs b/crates/mqdb-core/src/transport.rs index 5d1095dd..0283741d 100644 --- a/crates/mqdb-core/src/transport.rs +++ b/crates/mqdb-core/src/transport.rs @@ -30,10 +30,14 @@ pub enum Request { entity: String, id: String, fields: Value, + #[serde(default)] + expected_version: Option, }, Delete { entity: String, id: String, + #[serde(default)] + expected_version: Option, }, List { entity: String, @@ -120,6 +124,7 @@ pub enum ErrorCode { Forbidden = 403, NotFound = 404, Conflict = 409, + PreconditionFailed = 412, RateLimited = 429, Internal = 500, } @@ -138,6 +143,7 @@ impl ErrorCode { ErrorCode::Forbidden => 7, ErrorCode::NotFound => 5, ErrorCode::Conflict => 6, + ErrorCode::PreconditionFailed => 9, ErrorCode::RateLimited => 8, ErrorCode::Internal => 13, } @@ -218,6 +224,10 @@ impl From for Response { ), ), Error::Conflict(msg) => (ErrorCode::Conflict, format!("conflict: {msg}")), + Error::PreconditionFailed(msg) => ( + ErrorCode::PreconditionFailed, + format!("precondition failed: {msg}"), + ), _ => { tracing::error!(error = %e, "internal error in client request"); (ErrorCode::Internal, "internal error".to_string()) diff --git a/crates/mqdb-vault/Cargo.toml b/crates/mqdb-vault/Cargo.toml index a8461592..2c4cc760 100644 --- a/crates/mqdb-vault/Cargo.toml +++ b/crates/mqdb-vault/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-vault" -version = "0.1.3" +version = "0.1.4" publish = false edition.workspace = true license = "AGPL-3.0-only" diff --git a/crates/mqdb-vault/src/backend.rs b/crates/mqdb-vault/src/backend.rs index b26b3113..6db7960e 100644 --- a/crates/mqdb-vault/src/backend.rs +++ b/crates/mqdb-vault/src/backend.rs @@ -71,6 +71,7 @@ async fn pre_update_encrypt( ownership: &OwnershipConfig, sender_uid: Option<&str>, delta: Value, + expected_version: Option, skip: &[String], ) -> Result<(Request, Option<(Value, Value)>), Response> { if let OwnershipDecision::Check { @@ -91,6 +92,7 @@ async fn pre_update_encrypt( entity: entity.to_string(), id: id.to_string(), fields: delta, + expected_version, }, None, )); @@ -125,6 +127,7 @@ async fn pre_update_encrypt( entity: entity.to_string(), id: id.to_string(), fields: merged, + expected_version, }, Some((plaintext_merged, plaintext_existing)), )) @@ -185,9 +188,18 @@ impl VaultBackend for VaultBackendImpl { entity: ent, id, fields: delta, + expected_version, } => { match pre_update_encrypt( - db, &crypto, &ent, &id, ownership, sender_uid, delta, &skip, + db, + &crypto, + &ent, + &id, + ownership, + sender_uid, + delta, + expected_version, + &skip, ) .await { diff --git a/crates/mqdb-wasm/Cargo.lock b/crates/mqdb-wasm/Cargo.lock index 18b17ba3..c7d1c8c9 100644 --- a/crates/mqdb-wasm/Cargo.lock +++ b/crates/mqdb-wasm/Cargo.lock @@ -342,7 +342,7 @@ dependencies = [ [[package]] name = "mqdb-core" -version = "0.7.11" +version = "0.7.12" dependencies = [ "arc-swap", "bebytes", @@ -360,7 +360,7 @@ dependencies = [ [[package]] name = "mqdb-wasm" -version = "0.3.5" +version = "0.3.6" dependencies = [ "console_error_panic_hook", "futures-channel", diff --git a/crates/mqdb-wasm/Cargo.toml b/crates/mqdb-wasm/Cargo.toml index 415d93b0..9f6331ca 100644 --- a/crates/mqdb-wasm/Cargo.toml +++ b/crates/mqdb-wasm/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-wasm" -version = "0.3.5" +version = "0.3.6" edition = "2024" license = "Apache-2.0" authors = ["LabOverWire"] diff --git a/crates/mqdb-wasm/src/execute.rs b/crates/mqdb-wasm/src/execute.rs index c88f8c3d..3cff2281 100644 --- a/crates/mqdb-wasm/src/execute.rs +++ b/crates/mqdb-wasm/src/execute.rs @@ -36,10 +36,10 @@ impl WasmDatabase { match request { Request::Create { entity, data } => self.create(entity, serialize_js(&data)?).await, Request::Read { entity, id, .. } => self.read(entity, id).await, - Request::Update { entity, id, fields } => { - self.update(entity, id, serialize_js(&fields)?).await - } - Request::Delete { entity, id } => { + Request::Update { + entity, id, fields, .. + } => self.update(entity, id, serialize_js(&fields)?).await, + Request::Delete { entity, id, .. } => { self.delete(entity, id).await?; Ok(JsValue::NULL) } diff --git a/crates/mqdb-wasm/src/indexeddb.rs b/crates/mqdb-wasm/src/indexeddb.rs index 796bdcf3..648933c8 100644 --- a/crates/mqdb-wasm/src/indexeddb.rs +++ b/crates/mqdb-wasm/src/indexeddb.rs @@ -429,8 +429,8 @@ impl AsyncStorageBackend for IndexedDbBackend { } } - async fn flush(&self) -> Result<()> { - Ok(()) + fn flush(&self) -> impl std::future::Future> { + std::future::ready(Ok(())) } } diff --git a/docs/design/hold-reclaim.md b/docs/design/hold-reclaim.md index dbcf6a7e..4643d378 100644 --- a/docs/design/hold-reclaim.md +++ b/docs/design/hold-reclaim.md @@ -138,9 +138,9 @@ Cluster sweep (PR 1b) additional semantics: | PR | Scope | Notes | |----|-------|-------| -| 1a | TTL backstop fix — **agent** | `background.rs`: guard release + `expect_value` on scanned bytes + per-entity batches. Standalone data-loss + seat-lockout bug. | -| 1b | TTL backstop fix — **cluster** | Rewire `handle_ttl_cleanup` through the replicated delete path (primary-gated, guard release, change event, version precondition) — a larger change than 1a, split out to isolate risk. | -| 2 | Client CAS — both paths + new error | Add "retry-loop-doesn't-defeat-CAS" counter-tests. | +| 1a ✅ | TTL backstop fix — **agent** | `background.rs`: guard release + `expect_value` on scanned bytes + per-entity batches. Standalone data-loss + seat-lockout bug. | +| 1b ✅ | TTL backstop fix — **cluster** | Rewire `handle_ttl_cleanup` through the replicated delete path (primary-gated, guard release, change event, version precondition) — a larger change than 1a, split out to isolate risk. | +| 2 ✅ | Client CAS — both paths + new error | `_expected_version` reserved payload key, terminal `PreconditionFailed` (412); agent `update_with_expected`/`delete_with_expected` + both cluster write paths; retry-loop-doesn't-defeat-CAS counter-tests. | | 3 | Presence feed — both modes + topic rule | Gated on decisions 1–2 below. | | app | Janitor + reassert-on-reconnect contract + short keepalive | Out of mqdb (application). |