From 80d2257a8239eab361c4dc1164a9f9d3d724ab1d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Wed, 26 Aug 2026 15:24:49 -0300 Subject: [PATCH 1/4] roll back optimistic create on remote conflict --- crates/stitch-harness/src/lib.rs | 27 ++++++++++ crates/stitch-wasm/CHANGELOG.md | 11 ++++ crates/stitch/CHANGELOG.md | 33 ++++++++++++ crates/stitch/src/error.rs | 1 + crates/stitch/src/offline_queue.rs | 69 +++++++++++++------------ crates/stitch/src/queue.rs | 23 ++++++++- crates/stitch/src/store.rs | 71 ++++++++++++++++++++------ crates/stitch/tests/broker_harness.rs | 73 +++++++++++++++++++++++++++ crates/stitch/tests/offline_queue.rs | 23 +++++++-- 9 files changed, 280 insertions(+), 51 deletions(-) diff --git a/crates/stitch-harness/src/lib.rs b/crates/stitch-harness/src/lib.rs index 5913e7b..0ec357c 100644 --- a/crates/stitch-harness/src/lib.rs +++ b/crates/stitch-harness/src/lib.rs @@ -26,6 +26,7 @@ use std::net::SocketAddr; use std::sync::atomic::{AtomicU16, Ordering}; use mqdb_agent::{Database, MqdbAgent}; +use mqdb_core::schema::{FieldDefinition, FieldType, Schema}; use mqdb_core::types::ScopeConfig; use mqtt5::broker::PasswordAuthProvider; use tempfile::TempDir; @@ -88,6 +89,7 @@ pub struct BrokerHarness { anonymous_set: bool, users: Vec<(String, String)>, acl: Vec, + unique_constraints: Vec<(String, Vec)>, } impl Default for BrokerHarness { @@ -106,6 +108,7 @@ impl BrokerHarness { anonymous_set: false, users: Vec::new(), acl: Vec::new(), + unique_constraints: Vec::new(), } } @@ -147,6 +150,21 @@ impl BrokerHarness { self } + /// Register a UNIQUE constraint on `fields` of `entity` on the broker before + /// it starts. A create whose `fields` value already exists under a different + /// row id is rejected with a 409, letting tests exercise the client's + /// conflict handling against real broker enforcement. + #[must_use] + pub fn unique_constraint( + mut self, + entity: impl Into, + fields: impl IntoIterator>, + ) -> Self { + self.unique_constraints + .push((entity.into(), fields.into_iter().map(Into::into).collect())); + self + } + /// Bind an ephemeral loopback port, write the password/ACL files, and run the /// agent. The broker stays up until the returned [`RunningBroker`] is dropped /// or [`RunningBroker::shutdown`] is called. @@ -160,6 +178,15 @@ impl BrokerHarness { let dir = TempDir::new()?; let db = Database::open_without_background_tasks(dir.path().join("agent")).await?; + for (entity, fields) in &self.unique_constraints { + let mut schema = Schema::new(entity.clone()); + for field in fields { + schema = schema.add_field(FieldDefinition::new(field.clone(), FieldType::String)); + } + db.add_schema(schema).await?; + db.add_unique_constraint(entity.clone(), fields.clone()) + .await?; + } let mut agent = MqdbAgent::new(db) .with_bind_address(addr) .with_anonymous(self.anonymous) diff --git a/crates/stitch-wasm/CHANGELOG.md b/crates/stitch-wasm/CHANGELOG.md index b004b94..7091ce7 100644 --- a/crates/stitch-wasm/CHANGELOG.md +++ b/crates/stitch-wasm/CHANGELOG.md @@ -5,6 +5,17 @@ All notable changes to the `stitch-wasm` crate are documented here. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +### Fixed + +- `Store.create` (the `create` binding) now rejects with the broker error when an + insert loses a unique-key race (409 `Conflict`) or is denied (403 `Ownership`), + rolling back the optimistic local row, instead of resolving with an id and + leaving a phantom row. A browser client can now claim an exclusive key — e.g. a + seat hold keyed by a UNIQUE constraint — and observe the loss via a rejected + promise. Inherited from `stitch-sync`. + ## [0.3.0] - 2026-08-11 ### Added diff --git a/crates/stitch/CHANGELOG.md b/crates/stitch/CHANGELOG.md index 4faeec0..89cff49 100644 --- a/crates/stitch/CHANGELOG.md +++ b/crates/stitch/CHANGELOG.md @@ -5,6 +5,39 @@ All notable changes to the `stitch-sync` crate are documented here. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +### Fixed + +- `Store::create` now surfaces a remote rejection instead of hiding it. When a + remote is connected and the broker rejects the insert with `Conflict` (409, + e.g. a unique-constraint collision) or `Ownership` (403), the optimistic local + write is rolled back from memory, persistence, and the offline queue, and the + error is returned. Previously the rejection was logged at `warn` and `create` + returned `Ok(id)`, leaving a phantom local row — so a client racing for an + exclusive key (e.g. a seat/hold) believed it had won. This makes `Store::create` + usable for exclusive-key claims. Note the broker must enforce the uniqueness for + a 409 to arise: `mqdb-agent` has no primary-key collision on `id` (a duplicate + `id` is a silent last-writer-wins overwrite), so an exclusive key must be a + UNIQUE constraint on a non-`id` field. +- The offline-queue flush no longer overwrites the winner on a losing insert. A + queued insert that flushes into a unique `Conflict` is now dropped (and reported + in the flush result) instead of being converted into a blind `sync_update`. When + a `create` returned `Ok(id)` after a transient error and then lost the race, the + online flush-retry now rolls the losing local row back (memory + persistence + + queue), so the client converges to not-holding while continuously connected + rather than only on the next reconnect. +- `Error::is_permanent_mutation` now classifies the mqdb `UniqueViolation` variant + as permanent, so a locally-raised unique violation is no longer left as neither + transient nor permanent. + +### Changed + +- **Breaking:** `OfflineQueue::flush` returns a `FlushSummary` (the `retained` + count plus the `ConflictedInsert`s that were dropped on a unique conflict) + instead of a bare `usize`, so the store can roll back the local rows of losing + inserts. External implementors of the trait must update their signature. + ## [0.4.0] - 2026-08-11 ### Added diff --git a/crates/stitch/src/error.rs b/crates/stitch/src/error.rs index 3ffdc61..343276d 100644 --- a/crates/stitch/src/error.rs +++ b/crates/stitch/src/error.rs @@ -111,6 +111,7 @@ impl Error { source.as_ref(), mqdb_core::error::Error::Validation(_) | mqdb_core::error::Error::ConstraintViolation(_) + | mqdb_core::error::Error::UniqueViolation { .. } | mqdb_core::error::Error::ForeignKeyViolation { .. } | mqdb_core::error::Error::ForeignKeyRestrict { .. } | mqdb_core::error::Error::NotNullViolation { .. } diff --git a/crates/stitch/src/offline_queue.rs b/crates/stitch/src/offline_queue.rs index e44cfe8..25f028f 100644 --- a/crates/stitch/src/offline_queue.rs +++ b/crates/stitch/src/offline_queue.rs @@ -2,7 +2,7 @@ use crate::error::Result; use crate::lock::MutexExt; use crate::origin::Origin; use crate::persistence::PersistenceLayer; -pub use crate::queue::{MutationSender, OfflineQueue}; +pub use crate::queue::{ConflictedInsert, FlushSummary, MutationSender, OfflineQueue}; use crate::rt::{Shared, new_id, now_millis}; use crate::types::{Operation, PendingMutation, Record}; use async_trait::async_trait; @@ -172,8 +172,9 @@ async fn flush_consolidated( sender: &dyn MutationSender, root_entity: &str, mut remove_records: impl FnMut(Vec) -> RemoveFuture, -) -> usize { +) -> FlushSummary { let mut retained = 0usize; + let mut conflicted_inserts: Vec = Vec::new(); for mutation in consolidated { let attempt = match mutation.op { Operation::Insert => { @@ -232,18 +233,13 @@ async fn flush_consolidated( } } Err(err) if err.is_conflict() && mutation.op == Operation::Insert => { - if let Some(data) = mutation.data.clone() { - match sender - .sync_update(&mutation.entity, &mutation.scope_id, &mutation.id, data) - .await - { - Ok(()) => FlushOutcome::Drop, - Err(e) if e.is_transient() => FlushOutcome::Keep, - Err(_) => FlushOutcome::Drop, - } - } else { - FlushOutcome::Drop - } + tracing::warn!( + entity = %mutation.entity, + id = %mutation.id, + error = %err, + "dropping queued insert: a unique key is already held by another row" + ); + FlushOutcome::DropConflict } Err(err) if err.is_permanent_mutation() => { tracing::error!( @@ -271,16 +267,28 @@ async fn flush_consolidated( FlushOutcome::Drop => { remove_records(mutation.record_ids).await; } + FlushOutcome::DropConflict => { + conflicted_inserts.push(ConflictedInsert { + entity: mutation.entity.clone(), + id: mutation.id.clone(), + scope_id: mutation.scope_id.clone(), + }); + remove_records(mutation.record_ids).await; + } FlushOutcome::Keep => { retained += 1; } } } - retained + FlushSummary { + retained, + conflicted_inserts, + } } enum FlushOutcome { Drop, + DropConflict, Keep, } @@ -396,10 +404,10 @@ impl OfflineQueue for PersistentOfflineQueue { Ok(()) } - async fn flush(&self, sender: &dyn MutationSender) -> Result { + async fn flush(&self, sender: &dyn MutationSender) -> Result { let _guard = match FlushGuard::try_acquire(&self.flushing) { Some(g) => g, - None => return Ok(0), + None => return Ok(FlushSummary::default()), }; self.do_flush(sender).await } @@ -460,9 +468,9 @@ impl OfflineQueue for PersistentOfflineQueue { } impl PersistentOfflineQueue { - async fn do_flush(&self, sender: &dyn MutationSender) -> Result { + async fn do_flush(&self, sender: &dyn MutationSender) -> Result { let Some(user) = self.current_user() else { - return Ok(0); + return Ok(FlushSummary::default()); }; let filters = vec![Filter::new( "userId".into(), @@ -471,7 +479,7 @@ impl PersistentOfflineQueue { )]; let rows = self.list_rows(filters).await?; if rows.is_empty() { - return Ok(0); + return Ok(FlushSummary::default()); } let consolidated = consolidate(rows); let persistence = Shared::clone(&self.persistence); @@ -483,8 +491,7 @@ impl PersistentOfflineQueue { } }) }; - let retained = flush_consolidated(consolidated, sender, &self.root_entity, remove).await; - Ok(retained) + Ok(flush_consolidated(consolidated, sender, &self.root_entity, remove).await) } } @@ -543,14 +550,14 @@ impl OfflineQueue for InMemoryOfflineQueue { Ok(()) } - async fn flush(&self, sender: &dyn MutationSender) -> Result { + async fn flush(&self, sender: &dyn MutationSender) -> Result { let _guard = match FlushGuard::try_acquire(&self.flushing) { Some(g) => g, - None => return Ok(0), + None => return Ok(FlushSummary::default()), }; let snapshot: Vec = self.rows.lock_guard().clone(); if snapshot.is_empty() { - return Ok(0); + return Ok(FlushSummary::default()); } let consolidated = consolidate(snapshot); let rows_handle: &Mutex> = &self.rows; @@ -562,12 +569,12 @@ impl OfflineQueue for InMemoryOfflineQueue { set.lock_guard().extend(ids); }) }; - let retained = flush_consolidated(consolidated, sender, &self.root_entity, remove).await; + let summary = flush_consolidated(consolidated, sender, &self.root_entity, remove).await; let flushed: Vec = remove_set.lock_guard().clone(); rows_handle .lock_guard() .retain(|r| !flushed.contains(&r.record_id)); - Ok(retained) + Ok(summary) } async fn clear(&self) -> Result<()> { @@ -804,10 +811,10 @@ mod tests { fail_create_transient: true, ..Default::default() }; - assert_eq!(queue.flush(&failing).await?, 1); + assert_eq!(queue.flush(&failing).await?.retained, 1); let recovered = FakeSender::default(); - assert_eq!(queue.flush(&recovered).await?, 0); + assert_eq!(queue.flush(&recovered).await?.retained, 0); assert_eq!(recovered.creates.lock_guard().len(), 1); Ok(()) } @@ -829,7 +836,7 @@ mod tests { read_result: Mutex::new(Some(record(&[("id", "t1"), ("status", "running")]))), ..Default::default() }; - assert_eq!(queue.flush(&sender).await?, 0); + assert_eq!(queue.flush(&sender).await?.retained, 0); assert_eq!(sender.updates.load(Ordering::SeqCst), 1); assert_eq!(sender.creates.lock_guard().len(), 1); Ok(()) @@ -856,7 +863,7 @@ mod tests { .await?; let sender = FakeSender::default(); - assert_eq!(queue.flush(&sender).await?, 0); + assert_eq!(queue.flush(&sender).await?.retained, 0); assert_eq!(sender.updates.load(Ordering::SeqCst), 0); let creates = sender.creates.lock_guard(); assert_eq!(creates.len(), 1); diff --git a/crates/stitch/src/queue.rs b/crates/stitch/src/queue.rs index 00393f7..7463f54 100644 --- a/crates/stitch/src/queue.rs +++ b/crates/stitch/src/queue.rs @@ -2,6 +2,25 @@ use crate::error::Result; use crate::types::{Operation, PendingMutation, Record}; use async_trait::async_trait; +/// A queued insert dropped during flush because the broker reported the row's +/// unique key as already held by another row (409). The store rolls the +/// optimistic local copy back so a losing racer converges to not-holding. +#[derive(Debug, Clone)] +pub struct ConflictedInsert { + pub entity: String, + pub id: String, + pub scope_id: String, +} + +/// Outcome of an offline-queue flush: how many mutations were retained (still +/// pending, needing another pass) and which inserts were dropped on a unique +/// conflict so the caller can roll their local rows back. +#[derive(Debug, Default)] +pub struct FlushSummary { + pub retained: usize, + pub conflicted_inserts: Vec, +} + /// Sends a queued mutation to the remote during an offline-queue flush, and /// reads/deletes local rows while reconciling. Implemented by the remote sync /// layer. @@ -42,7 +61,7 @@ pub trait OfflineQueue: Send + Sync { scope_id: &str, op: Operation, ) -> Result<()>; - async fn flush(&self, sender: &dyn MutationSender) -> Result; + async fn flush(&self, sender: &dyn MutationSender) -> Result; async fn clear(&self) -> Result<()>; async fn pending_for_scope(&self, scope_id: &str) -> Result>; async fn has_pending_insert(&self, entity: &str, entity_id: &str) -> Result; @@ -60,7 +79,7 @@ pub trait OfflineQueue { scope_id: &str, op: Operation, ) -> Result<()>; - async fn flush(&self, sender: &dyn MutationSender) -> Result; + async fn flush(&self, sender: &dyn MutationSender) -> Result; async fn clear(&self) -> Result<()>; async fn pending_for_scope(&self, scope_id: &str) -> Result>; async fn has_pending_insert(&self, entity: &str, entity_id: &str) -> Result; diff --git a/crates/stitch/src/store.rs b/crates/stitch/src/store.rs index 519c2a7..1f5ee28 100644 --- a/crates/stitch/src/store.rs +++ b/crates/stitch/src/store.rs @@ -299,6 +299,14 @@ impl Store { /// owning scope; for the root entity it's ignored and the row's own id /// becomes the scope. Returns the new row's id (either taken from the /// `data.id` field or freshly generated). + /// + /// When a remote is connected and the broker rejects the insert as a + /// [`Error::Conflict`] (409, e.g. a unique-constraint collision) or an + /// [`Error::Ownership`] (403), the optimistic local write is rolled back + /// from memory, persistence, and the offline queue, and the error is + /// returned — so a caller racing for an exclusive key observes the loss + /// instead of a phantom local row. Transient failures keep the row queued + /// for retry and still return the id. pub async fn create( &self, entity: &str, @@ -375,16 +383,15 @@ impl Store { .await; } } - Err(err) if err.is_ownership() => { - if let Some(queue) = &inner.queue { - let _ = queue - .remove(entity, &id, &effective_scope, Operation::Insert) - .await; - } - } Err(err) if err.is_transient() => { inner.flush_notify.notify_one(); } + Err(err) if err.is_conflict() || err.is_ownership() => { + inner + .rollback_local_create(entity, &id, &effective_scope, origin) + .await; + return Err(err); + } Err(err) => { tracing::warn!( entity = %entity, @@ -1266,6 +1273,30 @@ impl StoreInner { inner: Shared::clone(self), } } + + async fn rollback_local_create(&self, entity: &str, id: &str, scope_id: &str, origin: Origin) { + let _ = self.memory.delete(entity, id, origin).await; + if let Some(persistence) = &self.persistence + && !origin.skips_persistence() + { + let _ = persistence.delete(entity, id, origin).await; + } + if let Some(queue) = &self.queue { + let _ = queue.remove(entity, id, scope_id, Operation::Insert).await; + } + } + + async fn rollback_conflicted_inserts(&self, conflicts: &[crate::queue::ConflictedInsert]) { + for conflicted in conflicts { + self.rollback_local_create( + &conflicted.entity, + &conflicted.id, + &conflicted.scope_id, + Origin::Local, + ) + .await; + } + } } struct InnerLocalAccessor { @@ -1451,14 +1482,17 @@ async fn flush_loop(inner: Shared) { break; }; let sender: &dyn crate::queue::MutationSender = remote.as_ref(); - let retained = match queue.flush(sender).await { - Ok(retained) => retained, + let summary = match queue.flush(sender).await { + Ok(summary) => summary, Err(err) => { tracing::warn!(error = %err, "offline queue flush failed"); break; } }; - if retained == 0 { + inner + .rollback_conflicted_inserts(&summary.conflicted_inserts) + .await; + if summary.retained == 0 { break; } rt::sleep(RETAIN_BACKOFF).await; @@ -1679,11 +1713,18 @@ async fn on_connected(inner: Shared) { if let (Some(queue), Some(remote)) = (queue_ref, remote.as_ref()) { let sender: &dyn crate::queue::MutationSender = remote.as_ref(); - let _ = queue.flush(sender).await; - if let Ok(retained) = queue.flush(sender).await - && retained > 0 - { - inner.flush_notify.notify_one(); + if let Ok(summary) = queue.flush(sender).await { + inner + .rollback_conflicted_inserts(&summary.conflicted_inserts) + .await; + } + if let Ok(summary) = queue.flush(sender).await { + inner + .rollback_conflicted_inserts(&summary.conflicted_inserts) + .await; + if summary.retained > 0 { + inner.flush_notify.notify_one(); + } } } diff --git a/crates/stitch/tests/broker_harness.rs b/crates/stitch/tests/broker_harness.rs index 2455483..abe33cd 100644 --- a/crates/stitch/tests/broker_harness.rs +++ b/crates/stitch/tests/broker_harness.rs @@ -110,3 +110,76 @@ async fn wrong_password_never_reaches_connected() { broker.shutdown().await; } + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn losing_create_on_unique_conflict_returns_err_and_leaves_no_phantom() { + init_tracing(); + let broker = BrokerHarness::new() + .scope("project", "projectId") + .unique_constraint("task", ["title"]) + .start() + .await + .expect("start broker"); + + let store = Store::with_client_id( + fixture_config(), + StoreOptions { + persistence: None, + remote: Some(RemoteConfig::new(broker.tcp_url())), + }, + "app-client".into(), + ); + store.initialize().await.expect("initialize"); + assert!(reaches_connected(&store).await, "must reach Connected"); + + store + .create( + "project", + "", + make_record(&[("id", json!("p1")), ("name", json!("Alpha"))]), + Origin::Local, + ) + .await + .expect("create project"); + + store + .create( + "task", + "p1", + make_record(&[ + ("id", json!("t1")), + ("title", json!("seatA")), + ("projectId", json!("p1")), + ]), + Origin::Local, + ) + .await + .expect("first claim wins"); + + let err = store + .create( + "task", + "p1", + make_record(&[ + ("id", json!("t2")), + ("title", json!("seatA")), + ("projectId", json!("p1")), + ]), + Origin::Local, + ) + .await + .expect_err("a second claim on the same unique key must be rejected"); + assert!(err.is_conflict(), "expected Conflict, got {err:?}"); + + assert!( + store.read("task", "t2").await.expect("read t2").is_none(), + "the losing create must leave no phantom local row" + ); + assert!( + store.read("task", "t1").await.expect("read t1").is_some(), + "the winning row must remain" + ); + + store.shutdown().await.expect("shutdown store"); + broker.shutdown().await; +} diff --git a/crates/stitch/tests/offline_queue.rs b/crates/stitch/tests/offline_queue.rs index 037ccd0..0457630 100644 --- a/crates/stitch/tests/offline_queue.rs +++ b/crates/stitch/tests/offline_queue.rs @@ -423,7 +423,7 @@ async fn update_on_root_with_not_found_triggers_hard_delete() { } #[tokio::test] -async fn conflict_on_insert_switches_to_update() { +async fn conflict_on_insert_drops_without_stealing_the_unique_key() { let queue = InMemoryOfflineQueue::new("project".to_string()); queue .queue(pending( @@ -438,10 +438,27 @@ async fn conflict_on_insert_switches_to_update() { let sender = MockSender::new(Behavior::Ok); sender.set("create", Behavior::Conflict); - queue.flush(sender.as_ref()).await.unwrap(); + let summary = queue.flush(sender.as_ref()).await.unwrap(); let log = sender.log(); assert!(log.iter().any(|e| e.starts_with("create:")), "log: {log:?}"); - assert!(log.iter().any(|e| e == "update:task:p1:t1"), "log: {log:?}"); + assert!( + !log.iter().any(|e| e.starts_with("update:")), + "a losing insert must not be converted into an update: {log:?}" + ); + assert_eq!( + summary.retained, 0, + "the conflicting insert is dropped, not retained" + ); + assert!(queue.pending_for_scope("p1").await.unwrap().is_empty()); + assert_eq!( + summary.conflicted_inserts.len(), + 1, + "the dropped insert must be reported so the store can roll its local row back" + ); + let reported = &summary.conflicted_inserts[0]; + assert_eq!(reported.entity, "task"); + assert_eq!(reported.id, "t1"); + assert_eq!(reported.scope_id, "p1"); } #[tokio::test] From c82d394cfc653d8e507649bde3d3f1704d5175bf Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Wed, 26 Aug 2026 16:51:14 -0300 Subject: [PATCH 2/4] bump stitch-sync 0.5.0 and stitch-wasm 0.4.0 --- Cargo.lock | 4 ++-- crates/stitch-wasm/CHANGELOG.md | 2 +- crates/stitch-wasm/Cargo.toml | 4 ++-- crates/stitch/CHANGELOG.md | 2 +- crates/stitch/Cargo.toml | 2 +- 5 files changed, 7 insertions(+), 7 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 6583dbc..5fcfffb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2404,7 +2404,7 @@ dependencies = [ [[package]] name = "stitch-sync" -version = "0.4.0" +version = "0.5.0" dependencies = [ "async-trait", "flume", @@ -2444,7 +2444,7 @@ dependencies = [ [[package]] name = "stitch-wasm" -version = "0.3.0" +version = "0.4.0" dependencies = [ "futures", "js-sys", diff --git a/crates/stitch-wasm/CHANGELOG.md b/crates/stitch-wasm/CHANGELOG.md index 7091ce7..d269270 100644 --- a/crates/stitch-wasm/CHANGELOG.md +++ b/crates/stitch-wasm/CHANGELOG.md @@ -5,7 +5,7 @@ All notable changes to the `stitch-wasm` crate are documented here. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). -## [Unreleased] +## [0.4.0] - 2026-08-26 ### Fixed diff --git a/crates/stitch-wasm/Cargo.toml b/crates/stitch-wasm/Cargo.toml index c415e91..2c01bbc 100644 --- a/crates/stitch-wasm/Cargo.toml +++ b/crates/stitch-wasm/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "stitch-wasm" -version = "0.3.0" +version = "0.4.0" edition = "2024" rust-version = "1.88" license = "Apache-2.0" @@ -20,7 +20,7 @@ futures = "0.3.32" js-sys = "0.3" serde = { version = "1", features = ["derive"] } serde_json = "1" -stitch = { package = "stitch-sync", version = "0.4.0", path = "../stitch" } +stitch = { package = "stitch-sync", version = "0.5.0", path = "../stitch" } tokio = { version = "1.52.3", default-features = false, features = ["sync"] } wasm-bindgen = "0.2" wasm-bindgen-futures = "0.4" diff --git a/crates/stitch/CHANGELOG.md b/crates/stitch/CHANGELOG.md index 89cff49..f935b87 100644 --- a/crates/stitch/CHANGELOG.md +++ b/crates/stitch/CHANGELOG.md @@ -5,7 +5,7 @@ All notable changes to the `stitch-sync` crate are documented here. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). -## [Unreleased] +## [0.5.0] - 2026-08-26 ### Fixed diff --git a/crates/stitch/Cargo.toml b/crates/stitch/Cargo.toml index 28ff1ae..9d18190 100644 --- a/crates/stitch/Cargo.toml +++ b/crates/stitch/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "stitch-sync" -version = "0.4.0" +version = "0.5.0" edition = "2024" rust-version = "1.90" license = "Apache-2.0" From 8c356f296814b7035855afb5f670357390728d37 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Fri, 28 Aug 2026 10:35:41 -0300 Subject: [PATCH 3/4] remove dead unique-violation classification and simplify reconnect flush --- crates/stitch/CHANGELOG.md | 3 --- crates/stitch/src/error.rs | 1 - crates/stitch/src/store.rs | 11 ++++++----- 3 files changed, 6 insertions(+), 9 deletions(-) diff --git a/crates/stitch/CHANGELOG.md b/crates/stitch/CHANGELOG.md index f935b87..725c30e 100644 --- a/crates/stitch/CHANGELOG.md +++ b/crates/stitch/CHANGELOG.md @@ -27,9 +27,6 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 online flush-retry now rolls the losing local row back (memory + persistence + queue), so the client converges to not-holding while continuously connected rather than only on the next reconnect. -- `Error::is_permanent_mutation` now classifies the mqdb `UniqueViolation` variant - as permanent, so a locally-raised unique violation is no longer left as neither - transient nor permanent. ### Changed diff --git a/crates/stitch/src/error.rs b/crates/stitch/src/error.rs index 343276d..3ffdc61 100644 --- a/crates/stitch/src/error.rs +++ b/crates/stitch/src/error.rs @@ -111,7 +111,6 @@ impl Error { source.as_ref(), mqdb_core::error::Error::Validation(_) | mqdb_core::error::Error::ConstraintViolation(_) - | mqdb_core::error::Error::UniqueViolation { .. } | mqdb_core::error::Error::ForeignKeyViolation { .. } | mqdb_core::error::Error::ForeignKeyRestrict { .. } | mqdb_core::error::Error::NotNullViolation { .. } diff --git a/crates/stitch/src/store.rs b/crates/stitch/src/store.rs index 1f5ee28..def1991 100644 --- a/crates/stitch/src/store.rs +++ b/crates/stitch/src/store.rs @@ -307,6 +307,12 @@ impl Store { /// returned — so a caller racing for an exclusive key observes the loss /// instead of a phantom local row. Transient failures keep the row queued /// for retry and still return the id. + /// + /// `create` upserts by `id`: calling it with an `id` that already exists + /// locally overwrites that row, and a subsequent remote rejection rolls the + /// `id` back to empty rather than restoring the prior row. Use a fresh `id` + /// for an exclusive-key claim, and [`Store::update`] to modify an existing + /// row. pub async fn create( &self, entity: &str, @@ -1713,11 +1719,6 @@ async fn on_connected(inner: Shared) { if let (Some(queue), Some(remote)) = (queue_ref, remote.as_ref()) { let sender: &dyn crate::queue::MutationSender = remote.as_ref(); - if let Ok(summary) = queue.flush(sender).await { - inner - .rollback_conflicted_inserts(&summary.conflicted_inserts) - .await; - } if let Ok(summary) = queue.flush(sender).await { inner .rollback_conflicted_inserts(&summary.conflicted_inserts) From 81c3aa7ed9dac88a2f73f397b8b568528c26e64e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Fri, 28 Aug 2026 13:50:33 -0300 Subject: [PATCH 4/4] roll back ownership-rejected inserts on flush and dedupe queue removal --- crates/stitch/CHANGELOG.md | 20 ++++++++------- crates/stitch/src/offline_queue.rs | 33 +++++++++++++----------- crates/stitch/src/queue.rs | 15 ++++++----- crates/stitch/src/store.rs | 23 ++++++++--------- crates/stitch/tests/offline_queue.rs | 38 ++++++++++++++++++++++++++-- 5 files changed, 84 insertions(+), 45 deletions(-) diff --git a/crates/stitch/CHANGELOG.md b/crates/stitch/CHANGELOG.md index 725c30e..8141643 100644 --- a/crates/stitch/CHANGELOG.md +++ b/crates/stitch/CHANGELOG.md @@ -21,19 +21,21 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 `id` is a silent last-writer-wins overwrite), so an exclusive key must be a UNIQUE constraint on a non-`id` field. - The offline-queue flush no longer overwrites the winner on a losing insert. A - queued insert that flushes into a unique `Conflict` is now dropped (and reported - in the flush result) instead of being converted into a blind `sync_update`. When - a `create` returned `Ok(id)` after a transient error and then lost the race, the - online flush-retry now rolls the losing local row back (memory + persistence + - queue), so the client converges to not-holding while continuously connected - rather than only on the next reconnect. + queued insert the broker rejects with `Conflict` (409) or `Ownership` (403) is + now dropped and reported in the flush result — instead of a unique conflict + being converted into a blind `sync_update`, or an ownership denial being dropped + while its local row lingered. When a `create` returned `Ok(id)` after a transient + error and then lost the race, the online flush-retry rolls the losing local row + back (memory + persistence), matching the synchronous create path, so the client + converges to not-holding while continuously connected rather than only on the + next reconnect. ### Changed - **Breaking:** `OfflineQueue::flush` returns a `FlushSummary` (the `retained` - count plus the `ConflictedInsert`s that were dropped on a unique conflict) - instead of a bare `usize`, so the store can roll back the local rows of losing - inserts. External implementors of the trait must update their signature. + count plus the `RejectedInsert`s the broker refused) instead of a bare `usize`, + so the store can roll back the local rows of rejected inserts. External + implementors of the trait must update their signature. ## [0.4.0] - 2026-08-11 diff --git a/crates/stitch/src/offline_queue.rs b/crates/stitch/src/offline_queue.rs index 25f028f..f164cd8 100644 --- a/crates/stitch/src/offline_queue.rs +++ b/crates/stitch/src/offline_queue.rs @@ -2,7 +2,7 @@ use crate::error::Result; use crate::lock::MutexExt; use crate::origin::Origin; use crate::persistence::PersistenceLayer; -pub use crate::queue::{ConflictedInsert, FlushSummary, MutationSender, OfflineQueue}; +pub use crate::queue::{FlushSummary, MutationSender, OfflineQueue, RejectedInsert}; use crate::rt::{Shared, new_id, now_millis}; use crate::types::{Operation, PendingMutation, Record}; use async_trait::async_trait; @@ -174,7 +174,7 @@ async fn flush_consolidated( mut remove_records: impl FnMut(Vec) -> RemoveFuture, ) -> FlushSummary { let mut retained = 0usize; - let mut conflicted_inserts: Vec = Vec::new(); + let mut rejected_inserts: Vec = Vec::new(); for mutation in consolidated { let attempt = match mutation.op { Operation::Insert => { @@ -205,6 +205,18 @@ async fn flush_consolidated( let outcome = match attempt { Ok(()) => FlushOutcome::Drop, Err(err) if err.is_transient() => FlushOutcome::Keep, + Err(err) + if (err.is_conflict() || err.is_ownership()) + && mutation.op == Operation::Insert => + { + tracing::warn!( + entity = %mutation.entity, + id = %mutation.id, + error = %err, + "dropping rejected queued insert; rolling back the optimistic local row" + ); + FlushOutcome::DropRejected + } Err(err) if err.is_ownership() => FlushOutcome::Drop, Err(err) if err.is_not_found() && mutation.op == Operation::Delete => { FlushOutcome::Drop @@ -232,15 +244,6 @@ async fn flush_consolidated( Err(_) => FlushOutcome::Drop, } } - Err(err) if err.is_conflict() && mutation.op == Operation::Insert => { - tracing::warn!( - entity = %mutation.entity, - id = %mutation.id, - error = %err, - "dropping queued insert: a unique key is already held by another row" - ); - FlushOutcome::DropConflict - } Err(err) if err.is_permanent_mutation() => { tracing::error!( entity = %mutation.entity, @@ -267,8 +270,8 @@ async fn flush_consolidated( FlushOutcome::Drop => { remove_records(mutation.record_ids).await; } - FlushOutcome::DropConflict => { - conflicted_inserts.push(ConflictedInsert { + FlushOutcome::DropRejected => { + rejected_inserts.push(RejectedInsert { entity: mutation.entity.clone(), id: mutation.id.clone(), scope_id: mutation.scope_id.clone(), @@ -282,13 +285,13 @@ async fn flush_consolidated( } FlushSummary { retained, - conflicted_inserts, + rejected_inserts, } } enum FlushOutcome { Drop, - DropConflict, + DropRejected, Keep, } diff --git a/crates/stitch/src/queue.rs b/crates/stitch/src/queue.rs index 7463f54..1d71a6e 100644 --- a/crates/stitch/src/queue.rs +++ b/crates/stitch/src/queue.rs @@ -2,23 +2,24 @@ use crate::error::Result; use crate::types::{Operation, PendingMutation, Record}; use async_trait::async_trait; -/// A queued insert dropped during flush because the broker reported the row's -/// unique key as already held by another row (409). The store rolls the -/// optimistic local copy back so a losing racer converges to not-holding. +/// A queued insert the broker rejected during flush — either a `Conflict` (409, +/// its unique key is already held by another row) or an `Ownership` denial (403). +/// Such an insert can never land, so the store rolls the optimistic local copy +/// back and a losing racer converges to not-holding. #[derive(Debug, Clone)] -pub struct ConflictedInsert { +pub struct RejectedInsert { pub entity: String, pub id: String, pub scope_id: String, } /// Outcome of an offline-queue flush: how many mutations were retained (still -/// pending, needing another pass) and which inserts were dropped on a unique -/// conflict so the caller can roll their local rows back. +/// pending, needing another pass) and which inserts the broker rejected so the +/// caller can roll their local rows back. #[derive(Debug, Default)] pub struct FlushSummary { pub retained: usize, - pub conflicted_inserts: Vec, + pub rejected_inserts: Vec, } /// Sends a queued mutation to the remote during an offline-queue flush, and diff --git a/crates/stitch/src/store.rs b/crates/stitch/src/store.rs index def1991..c3350fe 100644 --- a/crates/stitch/src/store.rs +++ b/crates/stitch/src/store.rs @@ -1280,27 +1280,26 @@ impl StoreInner { } } - async fn rollback_local_create(&self, entity: &str, id: &str, scope_id: &str, origin: Origin) { + async fn rollback_local_row(&self, entity: &str, id: &str, origin: Origin) { let _ = self.memory.delete(entity, id, origin).await; if let Some(persistence) = &self.persistence && !origin.skips_persistence() { let _ = persistence.delete(entity, id, origin).await; } + } + + async fn rollback_local_create(&self, entity: &str, id: &str, scope_id: &str, origin: Origin) { + self.rollback_local_row(entity, id, origin).await; if let Some(queue) = &self.queue { let _ = queue.remove(entity, id, scope_id, Operation::Insert).await; } } - async fn rollback_conflicted_inserts(&self, conflicts: &[crate::queue::ConflictedInsert]) { - for conflicted in conflicts { - self.rollback_local_create( - &conflicted.entity, - &conflicted.id, - &conflicted.scope_id, - Origin::Local, - ) - .await; + async fn rollback_rejected_inserts(&self, rejected: &[crate::queue::RejectedInsert]) { + for insert in rejected { + self.rollback_local_row(&insert.entity, &insert.id, Origin::Local) + .await; } } } @@ -1496,7 +1495,7 @@ async fn flush_loop(inner: Shared) { } }; inner - .rollback_conflicted_inserts(&summary.conflicted_inserts) + .rollback_rejected_inserts(&summary.rejected_inserts) .await; if summary.retained == 0 { break; @@ -1721,7 +1720,7 @@ async fn on_connected(inner: Shared) { let sender: &dyn crate::queue::MutationSender = remote.as_ref(); if let Ok(summary) = queue.flush(sender).await { inner - .rollback_conflicted_inserts(&summary.conflicted_inserts) + .rollback_rejected_inserts(&summary.rejected_inserts) .await; if summary.retained > 0 { inner.flush_notify.notify_one(); diff --git a/crates/stitch/tests/offline_queue.rs b/crates/stitch/tests/offline_queue.rs index 0457630..7af8b2c 100644 --- a/crates/stitch/tests/offline_queue.rs +++ b/crates/stitch/tests/offline_queue.rs @@ -451,11 +451,45 @@ async fn conflict_on_insert_drops_without_stealing_the_unique_key() { ); assert!(queue.pending_for_scope("p1").await.unwrap().is_empty()); assert_eq!( - summary.conflicted_inserts.len(), + summary.rejected_inserts.len(), 1, "the dropped insert must be reported so the store can roll its local row back" ); - let reported = &summary.conflicted_inserts[0]; + let reported = &summary.rejected_inserts[0]; + assert_eq!(reported.entity, "task"); + assert_eq!(reported.id, "t1"); + assert_eq!(reported.scope_id, "p1"); +} + +#[tokio::test] +async fn ownership_denied_insert_is_dropped_and_reported_for_rollback() { + let queue = InMemoryOfflineQueue::new("project".to_string()); + queue + .queue(pending( + Operation::Insert, + "task", + "t1", + "p1", + make_record(&[("id", json!("t1")), ("title", json!("x"))]), + )) + .await + .unwrap(); + + let sender = MockSender::new(Behavior::Ok); + sender.set("create", Behavior::Ownership); + let summary = queue.flush(sender.as_ref()).await.unwrap(); + + assert_eq!( + summary.retained, 0, + "an ownership-denied insert is not retained" + ); + assert!(queue.pending_for_scope("p1").await.unwrap().is_empty()); + assert_eq!( + summary.rejected_inserts.len(), + 1, + "an ownership-denied insert must be reported so the store rolls its local row back" + ); + let reported = &summary.rejected_inserts[0]; assert_eq!(reported.entity, "task"); assert_eq!(reported.id, "t1"); assert_eq!(reported.scope_id, "p1");