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
4 changes: 2 additions & 2 deletions Cargo.lock

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

27 changes: 27 additions & 0 deletions crates/stitch-harness/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -88,6 +89,7 @@ pub struct BrokerHarness {
anonymous_set: bool,
users: Vec<(String, String)>,
acl: Vec<AclEntry>,
unique_constraints: Vec<(String, Vec<String>)>,
}

impl Default for BrokerHarness {
Expand All @@ -106,6 +108,7 @@ impl BrokerHarness {
anonymous_set: false,
users: Vec::new(),
acl: Vec::new(),
unique_constraints: Vec::new(),
}
}

Expand Down Expand Up @@ -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<String>,
fields: impl IntoIterator<Item = impl Into<String>>,
) -> 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.
Expand All @@ -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)
Expand Down
11 changes: 11 additions & 0 deletions crates/stitch-wasm/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).

## [0.4.0] - 2026-08-26

### 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
Expand Down
4 changes: 2 additions & 2 deletions crates/stitch-wasm/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand All @@ -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"
Expand Down
32 changes: 32 additions & 0 deletions crates/stitch/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,38 @@ 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).

## [0.5.0] - 2026-08-26

### 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 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 `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

### Added
Expand Down
2 changes: 1 addition & 1 deletion crates/stitch/Cargo.toml
Original file line number Diff line number Diff line change
@@ -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"
Expand Down
76 changes: 43 additions & 33 deletions crates/stitch/src/offline_queue.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::{FlushSummary, MutationSender, OfflineQueue, RejectedInsert};
use crate::rt::{Shared, new_id, now_millis};
use crate::types::{Operation, PendingMutation, Record};
use async_trait::async_trait;
Expand Down Expand Up @@ -172,8 +172,9 @@ async fn flush_consolidated(
sender: &dyn MutationSender,
root_entity: &str,
mut remove_records: impl FnMut(Vec<String>) -> RemoveFuture,
) -> usize {
) -> FlushSummary {
let mut retained = 0usize;
let mut rejected_inserts: Vec<RejectedInsert> = Vec::new();
for mutation in consolidated {
let attempt = match mutation.op {
Operation::Insert => {
Expand Down Expand Up @@ -204,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
Expand Down Expand Up @@ -231,20 +244,6 @@ async fn flush_consolidated(
Err(_) => FlushOutcome::Drop,
}
}
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
}
}
Err(err) if err.is_permanent_mutation() => {
tracing::error!(
entity = %mutation.entity,
Expand All @@ -271,16 +270,28 @@ async fn flush_consolidated(
FlushOutcome::Drop => {
remove_records(mutation.record_ids).await;
}
FlushOutcome::DropRejected => {
rejected_inserts.push(RejectedInsert {
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,
rejected_inserts,
}
}

enum FlushOutcome {
Drop,
DropRejected,
Keep,
}

Expand Down Expand Up @@ -396,10 +407,10 @@ impl OfflineQueue for PersistentOfflineQueue {
Ok(())
}

async fn flush(&self, sender: &dyn MutationSender) -> Result<usize> {
async fn flush(&self, sender: &dyn MutationSender) -> Result<FlushSummary> {
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
}
Expand Down Expand Up @@ -460,9 +471,9 @@ impl OfflineQueue for PersistentOfflineQueue {
}

impl PersistentOfflineQueue {
async fn do_flush(&self, sender: &dyn MutationSender) -> Result<usize> {
async fn do_flush(&self, sender: &dyn MutationSender) -> Result<FlushSummary> {
let Some(user) = self.current_user() else {
return Ok(0);
return Ok(FlushSummary::default());
};
let filters = vec![Filter::new(
"userId".into(),
Expand All @@ -471,7 +482,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);
Expand All @@ -483,8 +494,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)
}
}

Expand Down Expand Up @@ -543,14 +553,14 @@ impl OfflineQueue for InMemoryOfflineQueue {
Ok(())
}

async fn flush(&self, sender: &dyn MutationSender) -> Result<usize> {
async fn flush(&self, sender: &dyn MutationSender) -> Result<FlushSummary> {
let _guard = match FlushGuard::try_acquire(&self.flushing) {
Some(g) => g,
None => return Ok(0),
None => return Ok(FlushSummary::default()),
};
let snapshot: Vec<StoredRow> = self.rows.lock_guard().clone();
if snapshot.is_empty() {
return Ok(0);
return Ok(FlushSummary::default());
}
let consolidated = consolidate(snapshot);
let rows_handle: &Mutex<Vec<StoredRow>> = &self.rows;
Expand All @@ -562,12 +572,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<String> = 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<()> {
Expand Down Expand Up @@ -804,10 +814,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(())
}
Expand All @@ -829,7 +839,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(())
Expand All @@ -856,7 +866,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);
Expand Down
Loading