diff --git a/CHANGELOG.md b/CHANGELOG.md index 81bee42a..63565f02 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,7 +4,11 @@ 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-08-29 — mqdb-cluster 0.4.10 +## 2026-09-01 — mqdb-agent 0.8.25 + +### Fixed + +- **TTL cleanup now releases unique guards and is version-guarded (agent mode).** The TTL sweep deleted expired rows with a bare `batch.remove` — it never released the row's unique-constraint guards, so a row with a `unique` field that expired left its guard behind and the value became **permanently unclaimable** (a new create for that value kept hitting the stale guard); and it carried no precondition, so a row renewed between the sweep's scan and its commit was deleted on the stale snapshot (**silent data loss**). The sweep now mirrors the normal delete: `expect_value` on the exact scanned bytes plus `release_unique_guards`, per row (each expired row is reaped in its own batch, so one concurrently-renewed row no longer aborts the rest). Groundwork for on-disconnect hold reclaim (`docs/design/hold-reclaim.md`); the race and the reclaim invariant are model-checked in `specs/AbandonedHoldReclaim.tla`. ### Added diff --git a/Cargo.lock b/Cargo.lock index dccc8466..e3aba79d 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1420,7 +1420,7 @@ dependencies = [ [[package]] name = "mqdb-agent" -version = "0.8.24" +version = "0.8.25" dependencies = [ "arc-swap", "argon2", diff --git a/crates/mqdb-agent/Cargo.toml b/crates/mqdb-agent/Cargo.toml index 1817e265..eccf287e 100644 --- a/crates/mqdb-agent/Cargo.toml +++ b/crates/mqdb-agent/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "mqdb-agent" -version = "0.8.24" +version = "0.8.25" edition.workspace = true license = "Apache-2.0" authors.workspace = true diff --git a/crates/mqdb-agent/src/database/background.rs b/crates/mqdb-agent/src/database/background.rs index 59151326..38695c90 100644 --- a/crates/mqdb-agent/src/database/background.rs +++ b/crates/mqdb-agent/src/database/background.rs @@ -5,6 +5,7 @@ use super::Database; use crate::consumer_group::ConsumerGroup; use crate::dispatcher::EventDispatcher; use crate::outbox_processor::OutboxProcessor; +use mqdb_core::constraint::ConstraintManager; use mqdb_core::entity::Entity; use mqdb_core::events::ChangeEvent; use mqdb_core::index::IndexManager; @@ -16,46 +17,50 @@ use std::sync::Arc; use tokio::sync::{RwLock, watch}; use tokio::task::JoinHandle; +pub(super) struct BackgroundDeps<'a> { + pub config: &'a mqdb_core::config::DatabaseConfig, + pub outbox: &'a Arc, + pub dispatcher: &'a Arc, + pub storage: &'a Arc, + pub index_manager: &'a Arc>, + pub constraint_manager: &'a Arc>, + pub consumer_groups: &'a Arc>>, + pub shutdown_rx: &'a watch::Receiver, +} + impl Database { - pub(super) fn spawn_background_tasks( - config: &mqdb_core::config::DatabaseConfig, - outbox: &Arc, - dispatcher: &Arc, - storage: &Arc, - index_manager: &Arc>, - consumer_groups: &Arc>>, - shutdown_rx: &watch::Receiver, - ) -> Vec> { + pub(super) fn spawn_background_tasks(deps: &BackgroundDeps<'_>) -> Vec> { let mut handles = Vec::new(); - if !config.spawn_background_tasks { + if !deps.config.spawn_background_tasks { return handles; } - if config.outbox.enabled { + if deps.config.outbox.enabled { handles.push(Self::spawn_outbox_processor( - outbox, - dispatcher, - &config.outbox, - shutdown_rx, + deps.outbox, + deps.dispatcher, + &deps.config.outbox, + deps.shutdown_rx, )); } - if let Some(interval_secs) = config.ttl_cleanup_interval_secs { + if let Some(interval_secs) = deps.config.ttl_cleanup_interval_secs { handles.push(Self::spawn_ttl_cleanup( - storage, - dispatcher, - outbox, - index_manager, + deps.storage, + deps.dispatcher, + deps.outbox, + deps.index_manager, + deps.constraint_manager, interval_secs, )); } - if config.shared_subscription.consumer_timeout_ms > 0 { + if deps.config.shared_subscription.consumer_timeout_ms > 0 { handles.push(Self::spawn_consumer_timeout_cleanup( - consumer_groups, - config.shared_subscription.consumer_timeout_ms, - shutdown_rx, + deps.consumer_groups, + deps.config.shared_subscription.consumer_timeout_ms, + deps.shutdown_rx, )); } @@ -92,25 +97,60 @@ impl Database { dispatcher: &Arc, outbox: &Arc, index_manager: &Arc>, + constraint_manager: &Arc>, interval_secs: u64, ) -> JoinHandle<()> { - let storage_clone = Arc::clone(storage); - let dispatcher_clone = Arc::clone(dispatcher); - let outbox_clone = Arc::clone(outbox); - let index_manager_clone = Arc::clone(index_manager); + let ctx = TtlSweepCtx { + storage: Arc::clone(storage), + dispatcher: Arc::clone(dispatcher), + outbox: Arc::clone(outbox), + index_manager: Arc::clone(index_manager), + constraint_manager: Arc::clone(constraint_manager), + }; tokio::spawn(async move { - ttl_cleanup_task( - storage_clone, - dispatcher_clone, - outbox_clone, - index_manager_clone, - interval_secs, - ) - .await; + ttl_cleanup_task(ctx, interval_secs).await; }) } + #[cfg(test)] + pub(crate) async fn ttl_cleanup_pass_for_test(&self, now: u64) -> usize { + self.ttl_sweep_ctx().pass(now).await + } + + #[cfg(test)] + pub(crate) fn raw_row_for_test( + &self, + entity: &str, + id: &str, + ) -> Option<(Vec, Vec, Entity)> { + let key = mqdb_core::keys::encode_data_key(entity, id); + let value = self.storage.get(&key).ok()??; + let entity = Entity::deserialize(entity.to_string(), id.to_string(), &value).ok()?; + Some((key, value, entity)) + } + + #[cfg(test)] + pub(crate) async fn reap_one_for_test( + &self, + key: &[u8], + value: &[u8], + entity: &Entity, + ) -> bool { + self.ttl_sweep_ctx().reap(key, value, entity).await + } + + #[cfg(test)] + fn ttl_sweep_ctx(&self) -> TtlSweepCtx { + TtlSweepCtx { + storage: Arc::clone(&self.storage), + dispatcher: Arc::clone(&self.dispatcher), + outbox: Arc::clone(&self.outbox), + index_manager: Arc::clone(&self.index_manager), + constraint_manager: Arc::clone(&self.constraint_manager), + } + } + fn spawn_consumer_timeout_cleanup( consumer_groups: &Arc>>, timeout_ms: u64, @@ -147,13 +187,15 @@ impl Database { } } -async fn ttl_cleanup_task( +struct TtlSweepCtx { storage: Arc, dispatcher: Arc, outbox: Arc, index_manager: Arc>, - interval_secs: u64, -) { + constraint_manager: Arc>, +} + +async fn ttl_cleanup_task(ctx: TtlSweepCtx, interval_secs: u64) { let mut interval = tokio::time::interval(tokio::time::Duration::from_secs(interval_secs)); loop { @@ -167,78 +209,98 @@ async fn ttl_cleanup_task( } }; - let prefix = b"data/"; - let Ok(items) = storage.prefix_scan(prefix) else { - continue; - }; + ctx.pass(now).await; + } +} - let mut expired_entities = Vec::new(); +impl TtlSweepCtx { + async fn pass(&self, now: u64) -> usize { + let Ok(items) = self.storage.prefix_scan(b"data/") else { + return 0; + }; + let mut expired = Vec::new(); for (key, value) in items { let Ok(key_str) = std::str::from_utf8(&key) else { continue; }; - let parts: Vec<&str> = key_str.split('/').collect(); if parts.len() != 3 { continue; } - - let entity_name = parts[1]; - let id = parts[2]; - - let Ok(entity) = Entity::deserialize(entity_name.to_string(), id.to_string(), &value) + let Ok(entity) = + Entity::deserialize(parts[1].to_string(), parts[2].to_string(), &value) else { continue; }; - if let Some(expires_at) = entity.data.get("_expires_at").and_then(Value::as_u64) && expires_at <= now { - expired_entities.push((key, entity)); + expired.push((key, value, entity)); } } - if expired_entities.is_empty() { - continue; + let expired_count = expired.len(); + let mut reaped = 0usize; + for (key, value, entity) in expired { + if self.reap(&key, &value, &entity).await { + reaped += 1; + } } + if expired_count > 0 { + tracing::debug!(reaped, expired = expired_count, "TTL cleanup processed"); + } + reaped + } + + async fn reap(&self, key: &[u8], value: &[u8], entity: &Entity) -> bool { let operation_id = uuid::Uuid::new_v4().to_string(); - let mut batch = storage.batch(); - let mut events = Vec::new(); + let mut batch = self.storage.batch(); - for (key, entity) in &expired_entities { - batch.remove(key.clone()); + batch.expect_value(key.to_vec(), value.to_vec()); + batch.remove(key.to_vec()); - let index_mgr = index_manager.read().await; + { + let index_mgr = self.index_manager.read().await; index_mgr.remove_indexes(&mut batch, entity); - - events.push(ChangeEvent::delete( - entity.name.clone(), - entity.id.clone(), - entity.data.clone(), - )); } - outbox.enqueue_events(&mut batch, &operation_id, &events); - - if let Err(e) = batch.commit() { - tracing::error!(err = %e, "TTL cleanup batch commit failed"); - continue; + { + let constraint_mgr = self.constraint_manager.read().await; + if let Err(e) = constraint_mgr.release_unique_guards(entity, &mut batch) { + tracing::warn!( + entity = %entity.name, + id = %entity.id, + err = %e, + "TTL cleanup: release_unique_guards failed, skipping" + ); + return false; + } } - for event in events { - let _ = dispatcher.dispatch(event).await; - } + let event = + ChangeEvent::delete(entity.name.clone(), entity.id.clone(), entity.data.clone()); + self.outbox + .enqueue_events(&mut batch, &operation_id, std::slice::from_ref(&event)); - if let Err(e) = outbox.mark_delivered(&operation_id) { - tracing::warn!(op_id = %operation_id, err = %e, "TTL cleanup mark_delivered failed"); + match batch.commit() { + Ok(()) => { + let _ = self.dispatcher.dispatch(event).await; + if let Err(e) = self.outbox.mark_delivered(&operation_id) { + tracing::warn!(op_id = %operation_id, err = %e, "TTL cleanup mark_delivered failed"); + } + true + } + Err(e) => { + tracing::debug!( + entity = %entity.name, + id = %entity.id, + err = %e, + "TTL cleanup: skipped expired row (renewed or deleted concurrently)" + ); + false + } } - - tracing::debug!( - count = expired_entities.len(), - op_id = %operation_id, - "TTL cleanup processed expired entities" - ); } } diff --git a/crates/mqdb-agent/src/database/crud.rs b/crates/mqdb-agent/src/database/crud.rs index 53ce5ab6..3c52bd30 100644 --- a/crates/mqdb-agent/src/database/crud.rs +++ b/crates/mqdb-agent/src/database/crud.rs @@ -940,6 +940,115 @@ mod concurrency_tests { .expect("value must be reusable after the owner is deleted"); } + #[tokio::test(flavor = "multi_thread")] + async fn ttl_sweep_releases_unique_guard_and_reclaims_seat() { + let (_tmp, db) = test_db().await; + db.add_unique_constraint("hold".to_string(), vec!["seat_id".to_string()]) + .await + .unwrap(); + let scope = ScopeConfig::default(); + + db.create( + "hold".to_string(), + json!({ "id": "a", "seat_id": "S1", "ttl_secs": 1 }), + None, + None, + None, + &scope, + ) + .await + .unwrap(); + + assert!( + db.create( + "hold".to_string(), + json!({ "id": "b", "seat_id": "S1" }), + None, + None, + None, + &scope, + ) + .await + .is_err(), + "the live hold must block a duplicate claim" + ); + + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap() + .as_secs() + + 100; + assert_eq!(db.ttl_cleanup_pass_for_test(now).await, 1); + + db.create( + "hold".to_string(), + json!({ "id": "c", "seat_id": "S1" }), + None, + None, + None, + &scope, + ) + .await + .expect("seat must be reclaimable after the TTL sweep reaps the expired hold"); + } + + #[tokio::test(flavor = "multi_thread")] + async fn ttl_reap_skips_row_renewed_after_scan() { + let (_tmp, db) = test_db().await; + db.add_unique_constraint("hold".to_string(), vec!["seat_id".to_string()]) + .await + .unwrap(); + let scope = ScopeConfig::default(); + + db.create( + "hold".to_string(), + json!({ "id": "a", "seat_id": "S1", "ttl_secs": 1 }), + None, + None, + None, + &scope, + ) + .await + .unwrap(); + + let (key, stale_value, stale_entity) = db.raw_row_for_test("hold", "a").unwrap(); + + let caller = CallerContext { + sender: None, + client_id: None, + scope_config: &scope, + }; + db.update( + "hold".to_string(), + "a".to_string(), + json!({ "note": "renewed" }), + None, + &caller, + ) + .await + .unwrap(); + + assert!( + !db.reap_one_for_test(&key, &stale_value, &stale_entity) + .await, + "a row renewed after the scan must not be reaped" + ); + + assert!( + db.create( + "hold".to_string(), + json!({ "id": "b", "seat_id": "S1" }), + None, + None, + None, + &scope, + ) + .await + .is_err(), + "the renewed hold must still hold the seat" + ); + } + #[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/database/mod.rs b/crates/mqdb-agent/src/database/mod.rs index aea3c9ff..35fcd371 100644 --- a/crates/mqdb-agent/src/database/mod.rs +++ b/crates/mqdb-agent/src/database/mod.rs @@ -154,15 +154,16 @@ impl Database { registry.load().await?; let (shutdown_tx, shutdown_rx) = watch::channel(false); - let handles = Self::spawn_background_tasks( - &config, - &outbox, - &dispatcher, - &storage, - &index_manager, - &consumer_groups, - &shutdown_rx, - ); + let handles = Self::spawn_background_tasks(&background::BackgroundDeps { + config: &config, + outbox: &outbox, + dispatcher: &dispatcher, + storage: &storage, + index_manager: &index_manager, + constraint_manager: &constraint_manager, + consumer_groups: &consumer_groups, + shutdown_rx: &shutdown_rx, + }); Ok(Self { storage, diff --git a/docs/design/hold-reclaim.md b/docs/design/hold-reclaim.md new file mode 100644 index 00000000..6a5f58a3 --- /dev/null +++ b/docs/design/hold-reclaim.md @@ -0,0 +1,175 @@ +# Reclaiming silently-abandoned exclusive holds + +## Problem + +An application (e.g. ticketing seat holds) builds an *exclusive hold* on mqdb: a row +with a `unique(seatId)` constraint, one hold per seat. Two buyers racing to claim the +same seat is already solved — the unique guard rejects one with a hard 409 (see +`docs/design/cluster-unique-hardening.md`). The unsolved case is **silent +abandonment**: a holder crashes / closes the tab / loses the network and never returns. +No status change is ever written, so the hold lingers and the seat is stuck forever. + +This is a **liveness** problem, orthogonal to the **safety** (no-oversell) problem the +unique-hardening work solved. Detecting "the owner is gone" is an *absence*, and an +absence produces no write for the unique machinery to arbitrate. Only two things can +observe it: **time** (a lease/TTL) or **connection state** (presence). The feature adds +the missing detection axis and hands the actual seat hand-off back to the existing +unique fence. + +## Design (app-side janitor) + +mqdb provides three primitives; the application runs a trusted (admin-authenticated) +*janitor* service. The janitor subscribes to presence, holds the `client_id → hold` +binding, applies a short grace (cancelled if the holder returns), then releases the +abandoned hold with a fenced (CAS) delete. Reclaimers then race a fresh unique-guarded +create — exactly one wins. + +### The load-bearing invariant (TLA-proven) + +Modeled in `specs/AbandonedHoldReclaim.tla` (+ `.cfg`, `_noreassert.cfg`, `_nocas.cfg`). +The reap is correct **iff both** hold: + +1. **Terminal per-attempt CAS.** The janitor's release is conditioned on the hold not + having changed since it observed (a compare on `_version`), and this check is + *terminal* — a mismatch is a non-retryable error, not something the write loop + retries past. +2. **Reassert-on-reconnect.** A returning holder performs a **write** (re-asserting the + hold: bump `_version`, refresh `_bound_client_id` / a lease timestamp) within the + grace window. That write is what makes the janitor's CAS fail. + +The model shows: with both, `InvNoFalseRelease` holds over the full state space; +dropping either one false-releases a live holder. **Presence connect-event cancellation +is only a latency optimization — the correctness fence is the reassert-write + CAS.** + +Oversell stays inherited-safe: reclaim is release-then-guarded-create through the +existing unique fence, and guard-release is in the same storage batch as the row remove, +so there is never a two-live-holds window. + +## The three mqdb primitives + +### A. Correct TTL backstop `[P0 — real bug, ship first]` + +Presence is best-effort; the TTL sweep is the mandatory backstop for every missed +disconnect and every hard crash. Today it is broken in two (agent) / four (cluster) ways. + +- **Agent** (`crates/mqdb-agent/src/database/background.rs`, `ttl_cleanup_task`): the + sweep does a bare `batch.remove(key)` missing both `release_unique_guards` + (`crud.rs:381`) — so an expired unique hold leaks its guard and the seat is + **permanently unsellable** — and `expect_value` (`crud.rs:374`) — so a hold renewed + between the scan and the commit is deleted on the stale snapshot (**data loss**). Fix: + add both, using the *scanned bytes* for `expect_value` (vault byte-exactness), and + switch to per-entity batches (the precondition is all-or-nothing). +- **Cluster** (`cluster_agent/event_loop.rs::handle_ttl_cleanup` → + `cluster/db/data_store.rs::cleanup_expired_ttl`): worse — a raw in-memory + `entities.remove` with **no ChangeEvent, no replication, no guard release, no FK + cascade**. The missing ChangeEvent means the change-feed reclaim signal never fires in + cluster mode. Fix: rewrite to take a write lock, gate on `is_primary_for_partition`, + and route each expired row through the real delete path (`db_delete_prepare` + + `db_commit` + `release_unique_for_deleted_record` + `publish_and_deliver_change_event`) + with a version precondition. +- Optional (defer): lazy read-time expiry (`read()` returns not-found for expired) plus + create-path opportunistic reap (evict an expired holder on unique-collision) so a seat + unblocks before the sweep. A short sweep interval covers the common case without this. + +Known limitations of the agent sweep (PR 1a), deliberate / pre-existing: +- **No FK cascade.** The reap releases the *reaped row's* unique guards but does not run + FK cascade / set-null / restrict, so expiring a TTL'd row that is an FK *parent* orphans + its children and leaves the children's guards. This is pre-existing (the old sweep also + bare-removed) and does not affect holds (leaf entities); wiring the sweep through the + full `validate_delete` cascade path is a follow-up (and the cluster rewire in 1b needs it + too). +- **Per-row batches on purpose.** Each expired row is reaped in its own batch so a single + concurrently-renewed row (failing `expect_value`) doesn't abort the rest — this trades one + commit for N. A grouped-commit fast path for the no-contention case is a possible later + optimization. +- **Guard-release failure skips the row.** If `release_unique_guards` errors (a malformed + unique value), the reap is skipped rather than removing the row without releasing its + guard (skipping is the safer of the two); the row is retried on the next sweep. A + permanent failure would retry (and re-log) each interval. + +### B. Client CAS (`expected_version`, terminal, both paths) + +- Shared: add `expected_version: Option` to `Request::Update`/`Request::Delete` + (`mqdb-core/src/transport.rs`) and carve it out of the payload in `build_request` + (else it merges as a data field). Add a non-`Conflict` error variant + (`PreconditionFailed`) + code. +- **The retry-loop trap:** the existing update/delete loops retry on `Conflict`, + re-read, and write anyway — a naive pre-loop CAS is silently defeated. The check must + live *inside* the per-attempt body, compared against the fresh read, and return the + non-retryable error. +- Agent: `crates/mqdb-agent/src/database/crud.rs` (`try_update_once`/`try_delete_once`). +- Cluster: `crates/mqdb-cluster/src/cluster/node_controller/db_ops.rs` + (`handle_json_update_local` right after `db_get`, **before** the unique-reserve round + that drops the lock; delete before `db_delete_prepare`). Reuse the existing + `__mqdb_fk_expected` terminal-mismatch precedent. +- Parity: payload parsing is shared (mqdb-core); enforcement is duplicated across the two + paths and must be added to both. + +### C. Presence feed (both modes) + +- Cluster (`cluster/event_handler/broker_events.rs`): emit connect + disconnect events, + fanned out the **LWT way** (`forward_publish_to_remotes` + topic-index targets) — + `queue_local_publish` and change events are local-only and are the wrong template. The + publish target must not be swallowed by the DB-op handler (carve `_presence` out of the + `handle_db_publish` interception, or publish outside `$DB`), and a shared + topic-protection rule (`$DB/_presence/# ReadOnly`, mirroring `$DB/+/events/#` in + `mqdb-agent/src/topic_rules.rs`) lets a non-admin janitor subscribe. Keep the `mqdb-` + internal-client skip. +- Agent: wire a `BrokerEventHandler` (none exists today) + a `spawn_presence_task` + publishing locally. +- Payload: `{client_id, user_id, event, unexpected, ts}`. Treat as a **hint** — dropped + under load, flap-reordered; `ts` authoritative. + +## Phasing + +| 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. | +| 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). | + +The TTL fix is split into 1a (agent) and 1b (cluster) because the cluster sweep does a +raw in-memory `entities.remove` with no replication — a correct fix must route each +expired row through the replicated delete path (primary-gated), which is a substantially +larger change than the agent's batch tweak. + +## Open decisions (recommendations; confirm before PR 3) + +1. **Per-connection vs per-user binding.** Per-connection false-reclaims a user who still + has another live connection. **Rec: per-user** (`_bound_user_id` + user-level presence + aggregation). +2. **mqtt-lib `user_id`.** ✅ RESOLVED — shipped in mqtt5 `0.38.4` (pin `mqtt5 = "0.38.4"`, + replacing the git dep). Both `ClientConnectEvent` and `ClientDisconnectEvent` now carry + `pub user_id: Option>` (mirrors `ClientPublishEvent::user_id`; the authenticated + username, `None` for anonymous). Read as `event.user_id.as_deref()`. `ClientDisconnectEvent` + also has `client_id: Arc` + `unexpected: bool`, so the presence payload + `{client_id, user_id, event, unexpected, ts}` is fully available (generate `ts`). This + unblocks the per-user binding in decision 1. The pin is a workspace dependency bump done + as its own step (it swaps a git dep for the published crate and needs its own build/test + verification), landing with or before PR 3. +3. **Crash-detection latency.** Sub-second reclaim holds only for **graceful** disconnects; + a hard crash is detected at the MQTT keepalive timeout. **Rec: short keepalive on hold + connections + document** the crash-path bound. + +## Verification + +- **TLA (done):** `specs/AbandonedHoldReclaim.tla` — safe with CAS + reassert; + false-releases without either. +- **Tests per phase:** TTL guard-release + version-precondition (both modes); CAS terminal + + not-defeated-by-retry + agent/cluster parity; presence emit + cross-node delivery + + internal-client skip. +- **Live E2E:** `mqdb dev` cluster — presence cross-node, CAS reject, end-to-end TTL + reclaim. + +## Why this is different from the double-sale work + +The unique-hardening program proved a **safety** invariant (`NoOversell`: at most one +committed unique claim per value, across failover / reconfig / correlated loss). It is +triggered entirely by writes. A silently-vanished holder writes nothing, so that +machinery never fires — and the very durability that makes a claim survive failover is +what makes a dead claim stick. This feature adds the missing **liveness/detection** axis +(presence + TTL) and a **fenced release** (CAS), then routes the seat hand-off back +through the unique fence it already trusts. It reuses the double-sale guarantee; it does +not replace or weaken it. diff --git a/specs/AbandonedHoldReclaim.cfg b/specs/AbandonedHoldReclaim.cfg new file mode 100644 index 00000000..b9ab0a53 --- /dev/null +++ b/specs/AbandonedHoldReclaim.cfg @@ -0,0 +1,7 @@ +CONSTANTS + Guarded = TRUE + ReconnectReasserts = TRUE + Buyers = {b1, b2} +SPECIFICATION Spec +INVARIANT TypeOK +INVARIANT InvNoFalseRelease diff --git a/specs/AbandonedHoldReclaim.tla b/specs/AbandonedHoldReclaim.tla new file mode 100644 index 00000000..76ae27d5 --- /dev/null +++ b/specs/AbandonedHoldReclaim.tla @@ -0,0 +1,113 @@ +---- MODULE AbandonedHoldReclaim ---- +(***************************************************************************) +(* Sub-second on-disconnect reclaim of an abandoned exclusive hold *) +(* (ticketing seat holds), app-side-janitor architecture. *) +(* *) +(* One seat. A hold is held by a buyer via a keep-alive connection. When *) +(* that connection drops, a janitor (after a grace delay, abstracted as *) +(* "may happen any time after disconnect") reaps the hold so the seat is *) +(* reclaimable. The danger is a FALSE RELEASE: reaping a holder who has in *) +(* fact come back. Oversell (two live holds) is inherited-safe: reclaim is *) +(* release-then-create through the existing unique fence (Reclaim fires *) +(* only when the seat is free), so the novel property to check is *) +(* NoFalseRelease. *) +(* *) +(* The janitor's CAS is modeled by a token: Observe takes a fresh snapshot *) +(* (reassertedSinceObs := FALSE) and a reassert write flips it TRUE, so the *) +(* CAS "hold unchanged since I looked" holds iff ~reassertedSinceObs -- the *) +(* finite, artifact-free stand-in for an unbounded version compare. *) +(* *) +(* Two knobs model the design decisions: *) +(* Guarded = janitor reaps with a CAS (reap only if the hold *) +(* was not re-asserted since it observed). *) +(* ReconnectReasserts = a returning holder re-asserts its hold (a write *) +(* that would fail the janitor's CAS) on reconnect. *) +(* *) +(* Expected: Guarded /\ ReconnectReasserts = SAFE; dropping either one *) +(* reintroduces false release -- BOTH the CAS fence and the *) +(* reassert-on-reconnect contract are necessary. *) +(***************************************************************************) +EXTENDS Naturals + +CONSTANTS Guarded, ReconnectReasserts, Buyers + +VARIABLES holder, present, obsPending, reassertedSinceObs, falseRelease + +vars == <> + +None == "none" + +TypeOK == + /\ holder \in (Buyers \cup {None}) + /\ present \in BOOLEAN + /\ obsPending \in BOOLEAN + /\ reassertedSinceObs \in BOOLEAN + /\ falseRelease \in BOOLEAN + +Init == + /\ holder \in Buyers + /\ present = TRUE + /\ obsPending = FALSE + /\ reassertedSinceObs = FALSE + /\ falseRelease = FALSE + +\* The holder's keep-alive connection drops. +Disconnect == + /\ holder # None + /\ present = TRUE + /\ present' = FALSE + /\ UNCHANGED <> + +\* The janitor reacts to the disconnect event (after its grace delay) and +\* snapshots the hold to CAS on later. +Observe == + /\ holder # None + /\ present = FALSE + /\ ~obsPending + /\ obsPending' = TRUE + /\ reassertedSinceObs' = FALSE + /\ UNCHANGED <> + +\* The holder comes back. Under the janitor contract it re-asserts its hold (a +\* write that would fail the janitor's pending CAS); without that contract it +\* only restores presence. +Reconnect == + /\ holder # None + /\ present = FALSE + /\ present' = TRUE + /\ reassertedSinceObs' = IF ReconnectReasserts THEN TRUE ELSE reassertedSinceObs + /\ UNCHANGED <> + +\* The janitor's grace expired: reap. Guarded => only if the hold was not +\* re-asserted since it observed (CAS). Reaping a present holder is a false +\* release. +ReapCAS == + /\ obsPending + /\ LET reap == (holder # None) /\ (~Guarded \/ ~reassertedSinceObs) + IN /\ holder' = IF reap THEN None ELSE holder + /\ falseRelease' = IF reap /\ present THEN TRUE ELSE falseRelease + /\ obsPending' = FALSE + /\ reassertedSinceObs' = FALSE + /\ UNCHANGED present + +\* A new buyer reclaims the freed seat. The unique fence permits this only +\* when the seat is free, so at most one live hold ever exists. +Reclaim(b) == + /\ holder = None + /\ holder' = b + /\ present' = TRUE + /\ obsPending' = FALSE + /\ reassertedSinceObs' = FALSE + /\ UNCHANGED falseRelease + +Next == + \/ Disconnect + \/ Observe + \/ Reconnect + \/ ReapCAS + \/ \E b \in Buyers : Reclaim(b) + +Spec == Init /\ [][Next]_vars + +InvNoFalseRelease == ~falseRelease +==== diff --git a/specs/AbandonedHoldReclaim_nocas.cfg b/specs/AbandonedHoldReclaim_nocas.cfg new file mode 100644 index 00000000..435d2ea0 --- /dev/null +++ b/specs/AbandonedHoldReclaim_nocas.cfg @@ -0,0 +1,7 @@ +CONSTANTS + Guarded = FALSE + ReconnectReasserts = TRUE + Buyers = {b1, b2} +SPECIFICATION Spec +INVARIANT TypeOK +INVARIANT InvNoFalseRelease diff --git a/specs/AbandonedHoldReclaim_noreassert.cfg b/specs/AbandonedHoldReclaim_noreassert.cfg new file mode 100644 index 00000000..6a173a44 --- /dev/null +++ b/specs/AbandonedHoldReclaim_noreassert.cfg @@ -0,0 +1,7 @@ +CONSTANTS + Guarded = TRUE + ReconnectReasserts = FALSE + Buyers = {b1, b2} +SPECIFICATION Spec +INVARIANT TypeOK +INVARIANT InvNoFalseRelease