From d14621e72a1197a6c3b6920680ea4ab2217c1dad Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Tue, 22 Sep 2026 21:52:25 -0300 Subject: [PATCH 1/2] make replicated subscriptions authoritative in reconcile --- .../cluster/event_handler/broker_events.rs | 4 +- .../src/cluster/event_handler/mod.rs | 2 +- .../src/cluster/event_handler/routing.rs | 2 +- crates/mqdb-cluster/src/cluster/mod.rs | 8 +- .../src/cluster/store_manager/mod.rs | 4 - .../src/cluster/subscription_cache.rs | 227 +++++++++------ .../mqdb-cluster/src/cluster/topic_index.rs | 5 + .../src/cluster/wildcard_pending.rs | 261 ------------------ .../src/cluster_agent/event_loop.rs | 45 +-- crates/mqdb-cluster/src/lib.rs | 34 +-- docs/distributed-design.md | 7 +- specs/ClusterSubReconcile.cfg | 9 + specs/ClusterSubReconcile.tla | 139 ++++++++++ specs/ClusterSubReconcile_nonholder.cfg | 9 + specs/ClusterSubReconcile_norepair.cfg | 9 + specs/ClusterSubReconcile_oldloss.cfg | 9 + specs/ClusterSubReconcile_oldresurrect.cfg | 9 + 17 files changed, 356 insertions(+), 427 deletions(-) delete mode 100644 crates/mqdb-cluster/src/cluster/wildcard_pending.rs create mode 100644 specs/ClusterSubReconcile.cfg create mode 100644 specs/ClusterSubReconcile.tla create mode 100644 specs/ClusterSubReconcile_nonholder.cfg create mode 100644 specs/ClusterSubReconcile_norepair.cfg create mode 100644 specs/ClusterSubReconcile_oldloss.cfg create mode 100644 specs/ClusterSubReconcile_oldresurrect.cfg diff --git a/crates/mqdb-cluster/src/cluster/event_handler/broker_events.rs b/crates/mqdb-cluster/src/cluster/event_handler/broker_events.rs index cb9a0729..99fae37d 100644 --- a/crates/mqdb-cluster/src/cluster/event_handler/broker_events.rs +++ b/crates/mqdb-cluster/src/cluster/event_handler/broker_events.rs @@ -264,7 +264,7 @@ impl BrokerEventHandler for ClusterEventHandler BrokerEventHandler for ClusterEventHandler( } } else { let _ = ctrl.stores_mut().topics.unsubscribe(topic, client_id); - let is_response_topic = topic.starts_with("resp/") || topic.contains("/resp/"); + let is_response_topic = crate::cluster::is_response_topic(topic); if !is_response_topic { let broadcast = TopicSubscriptionBroadcast::unsubscribe(topic, client_id); ClusterEventHandler::::broadcast_topic_subscription(ctrl, broadcast).await; diff --git a/crates/mqdb-cluster/src/cluster/event_handler/routing.rs b/crates/mqdb-cluster/src/cluster/event_handler/routing.rs index 4c3dbfd7..3b3c115b 100644 --- a/crates/mqdb-cluster/src/cluster/event_handler/routing.rs +++ b/crates/mqdb-cluster/src/cluster/event_handler/routing.rs @@ -168,7 +168,7 @@ impl ClusterEventHandler { .stores_mut() .topics .subscribe(topic, client_id, client_partition, qos); - let is_response_topic = topic.starts_with("resp/") || topic.contains("/resp/"); + let is_response_topic = crate::cluster::is_response_topic(topic); if is_response_topic { trace!(topic, client_id, "skipping broadcast for response topic"); } else { diff --git a/crates/mqdb-cluster/src/cluster/mod.rs b/crates/mqdb-cluster/src/cluster/mod.rs index 3e33c249..a6696149 100644 --- a/crates/mqdb-cluster/src/cluster/mod.rs +++ b/crates/mqdb-cluster/src/cluster/mod.rs @@ -41,7 +41,6 @@ mod subscription_cache; mod topic_index; mod topic_trie; mod transport; -mod wildcard_pending; mod wildcard_store; mod write_log; @@ -123,8 +122,8 @@ pub use subscription_cache::{ mqtt_subscription_key, }; pub use topic_index::{ - SubscriberLocation, TopicIndex, TopicIndexEntry, TopicIndexError, topic_index_key, - topic_partition, + SubscriberLocation, TopicIndex, TopicIndexEntry, TopicIndexError, is_response_topic, + topic_index_key, topic_partition, }; pub use topic_trie::{ SubscriptionType, TopicTrie, WildcardSubscriber, is_wildcard_pattern, validate_pattern, @@ -132,9 +131,6 @@ pub use topic_trie::{ pub use transport::{ ClusterMessage, ClusterTransport, InboundMessage, TransportConfig, TransportError, }; -pub use wildcard_pending::{ - PendingWildcard, WILDCARD_RECONCILIATION_INTERVAL_MS, WildcardPendingStore, -}; pub use wildcard_store::{WildcardEntry, WildcardStore, WildcardStoreError, wildcard_key}; pub use write_log::PartitionWriteLog; diff --git a/crates/mqdb-cluster/src/cluster/store_manager/mod.rs b/crates/mqdb-cluster/src/cluster/store_manager/mod.rs index d5b22312..0fa97b46 100644 --- a/crates/mqdb-cluster/src/cluster/store_manager/mod.rs +++ b/crates/mqdb-cluster/src/cluster/store_manager/mod.rs @@ -33,7 +33,6 @@ use super::retained_store::RetainedStore; use super::session::SessionStore; use super::subscription_cache::SubscriptionCache; use super::topic_index::TopicIndex; -use super::wildcard_pending::WildcardPendingStore; use super::wildcard_store::WildcardStore; use mqdb_core::storage::StorageBackend; use std::sync::Arc; @@ -165,7 +164,6 @@ pub struct StoreManager { pub retained: RetainedStore, pub topics: TopicIndex, pub wildcards: WildcardStore, - pub wildcard_pending: WildcardPendingStore, pub inflight: InflightStore, pub offsets: OffsetStore, pub idempotency: IdempotencyStore, @@ -198,7 +196,6 @@ impl StoreManager { retained: RetainedStore::new(node_id), topics: TopicIndex::new(node_id), wildcards: WildcardStore::new(node_id), - wildcard_pending: WildcardPendingStore::new(node_id), inflight: InflightStore::new(node_id), offsets: OffsetStore::new(node_id), idempotency: IdempotencyStore::new(node_id), @@ -234,7 +231,6 @@ impl std::fmt::Debug for StoreManager { .field("retained", &self.retained.message_count()) .field("topics", &self.topics.topic_count()) .field("wildcards", &self.wildcards.pattern_count()) - .field("wildcard_pending", &self.wildcard_pending.count()) .field("inflight", &self.inflight.count()) .field("offsets", &self.offsets.count()) .field("idempotency", &self.idempotency.count()) diff --git a/crates/mqdb-cluster/src/cluster/subscription_cache.rs b/crates/mqdb-cluster/src/cluster/subscription_cache.rs index 0b128535..862d6030 100644 --- a/crates/mqdb-cluster/src/cluster/subscription_cache.rs +++ b/crates/mqdb-cluster/src/cluster/subscription_cache.rs @@ -2,7 +2,7 @@ // SPDX-License-Identifier: AGPL-3.0-only use super::protocol::Operation; -use super::{NodeId, PartitionId, TopicIndex, WildcardStore, session_partition}; +use super::{NodeId, PartitionId, TopicIndex, WildcardStore, is_response_topic, session_partition}; use bebytes::BeBytes; use std::collections::{HashMap, HashSet}; use std::sync::RwLock; @@ -153,8 +153,7 @@ pub struct SubscriptionCache { #[derive(Debug, Default)] pub struct ReconciliationResult { pub clients_checked: usize, - pub subscriptions_added: usize, - pub subscriptions_removed: usize, + pub index_entries_restored: usize, } impl SubscriptionCache { @@ -202,72 +201,44 @@ impl SubscriptionCache { wildcard_store: &WildcardStore, ) -> ReconciliationResult { let mut result = ReconciliationResult::default(); - let mut snapshots = self.snapshots.write().unwrap(); - - let client_ids: Vec = snapshots.keys().cloned().collect(); - for client_id in &client_ids { + for snapshot in self.all_snapshots() { result.clients_checked += 1; + let client_id = snapshot.client_id_str(); - let authoritative_exact: HashSet = topic_index + let indexed_exact: HashSet = topic_index .get_client_topics(client_id) .into_iter() .map(|(topic, _)| topic) .collect(); + for entry in snapshot.exact_subscriptions() { + if !indexed_exact.contains(entry.topic_str()) + && !is_response_topic(entry.topic_str()) + && topic_index + .subscribe( + entry.topic_str(), + client_id, + session_partition(client_id), + entry.qos, + ) + .is_ok() + { + result.index_entries_restored += 1; + } + } - let authoritative_wildcards: HashSet = wildcard_store + let indexed_wildcards: HashSet = wildcard_store .get_client_patterns(client_id) .into_iter() .map(|(pattern, _)| pattern) .collect(); - - if let Some(snapshot) = snapshots.get_mut(client_id) { - let cached_topics: HashSet = snapshot - .topics - .iter() - .map(|t| t.topic_str().to_string()) - .collect(); - - for cached_topic in &cached_topics { - let is_wildcard = cached_topic.contains('+') || cached_topic.contains('#'); - let exists = if is_wildcard { - authoritative_wildcards.contains(cached_topic) - } else { - authoritative_exact.contains(cached_topic) - }; - - if !exists { - let _ = snapshot.remove_subscription(cached_topic); - result.subscriptions_removed += 1; - } - } - - for topic in &authoritative_exact { - if !snapshot.has_subscription(topic) - && let Some((_, qos)) = topic_index - .get_client_topics(client_id) - .into_iter() - .find(|(t, _)| t == topic) - { - snapshot.add_subscription(topic, qos); - result.subscriptions_added += 1; - } - } - - for pattern in &authoritative_wildcards { - if !snapshot.has_subscription(pattern) - && let Some((_, qos)) = wildcard_store - .get_client_patterns(client_id) - .into_iter() - .find(|(p, _)| p == pattern) - { - snapshot.add_subscription(pattern, qos); - result.subscriptions_added += 1; - } - } - - if snapshot.topics.is_empty() { - snapshots.remove(client_id); + for entry in snapshot.wildcard_subscriptions() { + if !indexed_wildcards.contains(entry.topic_str()) + && wildcard_store + .subscribe_mqtt(entry.topic_str(), client_id, entry.qos) + .is_ok() + { + result.index_entries_restored += 1; } } } @@ -714,14 +685,16 @@ mod tests { } #[test] - fn reconcile_removes_stale_entries() { + fn reconcile_restores_missing_index_entries_from_cache() { let cache = SubscriptionCache::new(node(1)); let topic_index = TopicIndex::new(node(1)); let wildcard_store = WildcardStore::new(node(1)); cache.add_subscription("client1", "topic/a", 1).unwrap(); - cache.add_subscription("client1", "topic/b", 1).unwrap(); - + cache.add_subscription("client1", "topic/b", 2).unwrap(); + cache + .add_subscription("client1", "sensors/+/temp", 1) + .unwrap(); topic_index .subscribe("topic/a", "client1", partition(), 1) .unwrap(); @@ -729,64 +702,140 @@ mod tests { let result = cache.reconcile(&topic_index, &wildcard_store); assert_eq!(result.clients_checked, 1); - assert_eq!(result.subscriptions_removed, 1); - assert_eq!(result.subscriptions_added, 0); - - let subs = cache.get_subscriptions("client1"); - assert_eq!(subs.len(), 1); - assert_eq!(subs[0].topic_str(), "topic/a"); + assert_eq!(result.index_entries_restored, 2); + let mut indexed = topic_index.get_client_topics("client1"); + indexed.sort(); + assert_eq!( + indexed, + vec![("topic/a".to_string(), 1), ("topic/b".to_string(), 2)] + ); + assert_eq!( + topic_index.get_subscribers("topic/b")[0].partition(), + Some(session_partition("client1")) + ); + assert_eq!( + wildcard_store.get_client_patterns("client1"), + vec![("sensors/+/temp".to_string(), 1)] + ); + assert_eq!(cache.get_subscriptions("client1").len(), 3); } #[test] - fn reconcile_adds_missing_entries() { + fn reconcile_does_not_resurrect_unsubscribed_topics_from_index() { let cache = SubscriptionCache::new(node(1)); let topic_index = TopicIndex::new(node(1)); let wildcard_store = WildcardStore::new(node(1)); - cache.add_subscription("client1", "topic/a", 1).unwrap(); - + cache.add_subscription("client1", "topic/kept", 1).unwrap(); topic_index - .subscribe("topic/a", "client1", partition(), 1) + .subscribe("topic/kept", "client1", partition(), 1) .unwrap(); topic_index - .subscribe("topic/b", "client1", partition(), 2) + .subscribe("topic/dropped", "client1", partition(), 1) + .unwrap(); + wildcard_store + .subscribe_mqtt("dropped/+/pattern", "client1", 1) .unwrap(); let result = cache.reconcile(&topic_index, &wildcard_store); - assert_eq!(result.clients_checked, 1); - assert_eq!(result.subscriptions_removed, 0); - assert_eq!(result.subscriptions_added, 1); - - let subs = cache.get_subscriptions("client1"); - assert_eq!(subs.len(), 2); + assert_eq!(result.index_entries_restored, 0); + let topics: Vec = cache + .get_subscriptions("client1") + .iter() + .map(|t| t.topic_str().to_string()) + .collect(); + assert_eq!(topics, vec!["topic/kept".to_string()]); } #[test] - fn reconcile_handles_wildcards() { + fn reconcile_skips_response_topics() { let cache = SubscriptionCache::new(node(1)); let topic_index = TopicIndex::new(node(1)); let wildcard_store = WildcardStore::new(node(1)); cache - .add_subscription("client1", "sensors/+/temp", 1) + .add_subscription("client1", "resp/client1", 0) .unwrap(); cache - .add_subscription("client1", "stale/+/pattern", 1) - .unwrap(); - - wildcard_store - .subscribe_mqtt("sensors/+/temp", "client1", 1) + .add_subscription("client1", "app/resp/client1", 0) .unwrap(); let result = cache.reconcile(&topic_index, &wildcard_store); - assert_eq!(result.clients_checked, 1); - assert_eq!(result.subscriptions_removed, 1); + assert_eq!(result.index_entries_restored, 0); + assert!(topic_index.get_client_topics("client1").is_empty()); + } - let subs = cache.get_subscriptions("client1"); - assert_eq!(subs.len(), 1); - assert_eq!(subs[0].topic_str(), "sensors/+/temp"); + #[test] + fn reconcile_is_idempotent() { + let cache = SubscriptionCache::new(node(1)); + let topic_index = TopicIndex::new(node(1)); + let wildcard_store = WildcardStore::new(node(1)); + + cache.add_subscription("client1", "topic/a", 1).unwrap(); + cache.add_subscription("client1", "alerts/#", 1).unwrap(); + + assert_eq!( + cache + .reconcile(&topic_index, &wildcard_store) + .index_entries_restored, + 2 + ); + assert_eq!( + cache + .reconcile(&topic_index, &wildcard_store) + .index_entries_restored, + 0 + ); + assert_eq!(topic_index.get_subscribers("topic/a").len(), 1); + } + + #[test] + fn replicated_subscription_survives_reconcile_when_broadcast_was_missed() { + let replica = SubscriptionCache::new(node(2)); + let topic_index = TopicIndex::new(node(2)); + let wildcard_store = WildcardStore::new(node(2)); + + let mut snapshot = MqttSubscriptionSnapshot::create("client1"); + snapshot.add_subscription("sensors/a", 1); + snapshot.add_subscription("alerts/+", 1); + replica + .apply_replicated( + Operation::Insert, + "client1", + &SubscriptionCache::serialize(&snapshot), + ) + .unwrap(); + + let result = replica.reconcile(&topic_index, &wildcard_store); + + assert_eq!(result.index_entries_restored, 2); + assert_eq!( + topic_index.get_client_topics("client1"), + vec![("sensors/a".to_string(), 1)] + ); + assert_eq!( + wildcard_store.get_client_patterns("client1"), + vec![("alerts/+".to_string(), 1)] + ); + + let promoted = SubscriptionCache::new(node(3)); + promoted + .import_subscriptions(&replica.export_for_partition(session_partition("client1"))) + .unwrap(); + + let mut topics: Vec = promoted + .get_subscriptions("client1") + .iter() + .map(|t| t.topic_str().to_string()) + .collect(); + topics.sort(); + assert_eq!( + topics, + vec!["alerts/+".to_string(), "sensors/a".to_string()], + "replicated subscriptions were deleted by reconcile because the index broadcast was missed" + ); } #[test] diff --git a/crates/mqdb-cluster/src/cluster/topic_index.rs b/crates/mqdb-cluster/src/cluster/topic_index.rs index 0af720c6..618c49c8 100644 --- a/crates/mqdb-cluster/src/cluster/topic_index.rs +++ b/crates/mqdb-cluster/src/cluster/topic_index.rs @@ -111,6 +111,11 @@ pub fn topic_partition(topic: &str) -> PartitionId { PartitionId::new(partition_num).unwrap() } +#[must_use] +pub fn is_response_topic(topic: &str) -> bool { + topic.starts_with("resp/") || topic.contains("/resp/") +} + #[must_use] pub fn topic_index_key(topic: &str) -> String { let partition = topic_partition(topic); diff --git a/crates/mqdb-cluster/src/cluster/wildcard_pending.rs b/crates/mqdb-cluster/src/cluster/wildcard_pending.rs deleted file mode 100644 index c8b281b1..00000000 --- a/crates/mqdb-cluster/src/cluster/wildcard_pending.rs +++ /dev/null @@ -1,261 +0,0 @@ -// Copyright 2025-2026 LabOverWire. All rights reserved. -// SPDX-License-Identifier: AGPL-3.0-only - -use super::protocol::{Operation, WildcardBroadcast}; -use super::{NodeId, PartitionId, SubscriptionType}; -use std::collections::HashMap; -use std::sync::RwLock; - -pub const WILDCARD_RECONCILIATION_INTERVAL_MS: u64 = 60_000; - -#[derive(Debug, Clone)] -pub struct PendingWildcard { - pub pattern: String, - pub client_id: String, - pub client_partition: PartitionId, - pub qos: u8, - pub subscription_type: SubscriptionType, - pub operation: Operation, - pub failed_partitions: Vec, - pub created_at: u64, - pub retry_count: u32, -} - -impl PendingWildcard { - #[must_use] - pub fn new_subscribe( - pattern: &str, - client_id: &str, - client_partition: PartitionId, - qos: u8, - subscription_type: SubscriptionType, - now: u64, - ) -> Self { - Self { - pattern: pattern.to_string(), - client_id: client_id.to_string(), - client_partition, - qos, - subscription_type, - operation: Operation::Insert, - failed_partitions: Vec::new(), - created_at: now, - retry_count: 0, - } - } - - /// # Panics - /// Panics if partition 0 is invalid (should never happen). - #[must_use] - pub fn new_unsubscribe(pattern: &str, client_id: &str, now: u64) -> Self { - Self { - pattern: pattern.to_string(), - client_id: client_id.to_string(), - client_partition: PartitionId::ZERO, - qos: 0, - subscription_type: SubscriptionType::Mqtt, - operation: Operation::Delete, - failed_partitions: Vec::new(), - created_at: now, - retry_count: 0, - } - } - - #[must_use] - pub fn key(&self) -> String { - format!("{}:{}", self.pattern, self.client_id) - } - - #[must_use] - pub fn to_broadcast(&self) -> WildcardBroadcast { - match self.operation { - Operation::Insert => WildcardBroadcast::subscribe( - &self.pattern, - &self.client_id, - self.client_partition, - self.qos, - self.subscription_type as u8, - ), - Operation::Delete | Operation::Update => { - WildcardBroadcast::unsubscribe(&self.pattern, &self.client_id) - } - } - } - - pub fn add_failed_partition(&mut self, partition: PartitionId) { - if !self.failed_partitions.contains(&partition) { - self.failed_partitions.push(partition); - } - } - - pub fn mark_retry(&mut self) { - self.retry_count += 1; - self.failed_partitions.clear(); - } -} - -pub struct WildcardPendingStore { - node_id: NodeId, - pending: RwLock>, - last_reconciliation: RwLock, -} - -impl WildcardPendingStore { - #[must_use] - pub fn new(node_id: NodeId) -> Self { - Self { - node_id, - pending: RwLock::new(HashMap::new()), - last_reconciliation: RwLock::new(0), - } - } - - /// # Panics - /// Panics if the internal lock is poisoned. - pub fn add_pending(&self, pending: PendingWildcard) { - let key = pending.key(); - self.pending.write().unwrap().insert(key, pending); - } - - /// # Panics - /// Panics if the internal lock is poisoned. - pub fn remove_pending(&self, pattern: &str, client_id: &str) { - let key = format!("{pattern}:{client_id}"); - self.pending.write().unwrap().remove(&key); - } - - /// # Panics - /// Panics if the internal lock is poisoned. - pub fn mark_partition_failed(&self, pattern: &str, client_id: &str, partition: PartitionId) { - let key = format!("{pattern}:{client_id}"); - if let Some(pending) = self.pending.write().unwrap().get_mut(&key) { - pending.add_failed_partition(partition); - } - } - - /// # Panics - /// Panics if the internal lock is poisoned. - #[must_use] - pub fn needs_reconciliation(&self, now: u64) -> bool { - let last = *self.last_reconciliation.read().unwrap(); - now.saturating_sub(last) >= WILDCARD_RECONCILIATION_INTERVAL_MS - } - - /// # Panics - /// Panics if the internal lock is poisoned. - pub fn mark_reconciliation(&self, now: u64) { - *self.last_reconciliation.write().unwrap() = now; - } - - /// # Panics - /// Panics if the internal lock is poisoned. - #[must_use] - pub fn get_pending_for_retry(&self) -> Vec { - self.pending.read().unwrap().values().cloned().collect() - } - - /// # Panics - /// Panics if the internal lock is poisoned. - pub fn mark_retried(&self, pattern: &str, client_id: &str) { - let key = format!("{pattern}:{client_id}"); - if let Some(pending) = self.pending.write().unwrap().get_mut(&key) { - pending.mark_retry(); - } - } - - /// # Panics - /// Panics if the internal lock is poisoned. - #[must_use] - pub fn count(&self) -> usize { - self.pending.read().unwrap().len() - } - - /// # Panics - /// Panics if the internal lock is poisoned. - pub fn clear_old_entries(&self, now: u64, max_age_ms: u64) -> usize { - let mut pending = self.pending.write().unwrap(); - let before = pending.len(); - pending.retain(|_, p| now.saturating_sub(p.created_at) <= max_age_ms); - before - pending.len() - } -} - -impl std::fmt::Debug for WildcardPendingStore { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - f.debug_struct("WildcardPendingStore") - .field("node_id", &self.node_id) - .field("pending_count", &self.count()) - .finish_non_exhaustive() - } -} - -#[cfg(test)] -mod tests { - use super::*; - - fn node() -> NodeId { - NodeId::validated(1).unwrap() - } - - fn partition() -> PartitionId { - PartitionId::new(5).unwrap() - } - - #[test] - fn add_and_remove_pending() { - let store = WildcardPendingStore::new(node()); - - let pending = PendingWildcard::new_subscribe( - "sensors/+/temp", - "client1", - partition(), - 1, - SubscriptionType::Mqtt, - 1000, - ); - - store.add_pending(pending); - assert_eq!(store.count(), 1); - - store.remove_pending("sensors/+/temp", "client1"); - assert_eq!(store.count(), 0); - } - - #[test] - fn needs_reconciliation_after_interval() { - let store = WildcardPendingStore::new(node()); - - assert!(store.needs_reconciliation(60_000)); - - store.mark_reconciliation(1000); - assert!(!store.needs_reconciliation(30_000)); - assert!(!store.needs_reconciliation(60_000)); - assert!(store.needs_reconciliation(61_001)); - } - - #[test] - fn clear_old_entries_removes_expired() { - let store = WildcardPendingStore::new(node()); - - store.add_pending(PendingWildcard::new_subscribe( - "old/+", - "client1", - partition(), - 1, - SubscriptionType::Mqtt, - 1000, - )); - store.add_pending(PendingWildcard::new_subscribe( - "new/+", - "client2", - partition(), - 1, - SubscriptionType::Mqtt, - 50_000, - )); - - let removed = store.clear_old_entries(100_000, 60_000); - assert_eq!(removed, 1); - assert_eq!(store.count(), 1); - } -} diff --git a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs index 17734cef..16b2e0d5 100644 --- a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs +++ b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs @@ -76,7 +76,6 @@ impl ClusteredAgent { let mut tick_interval = interval(Duration::from_millis(10)); let mut cleanup_interval = interval(Duration::from_secs(CLEANUP_INTERVAL_SECS)); let mut ttl_cleanup_interval = interval(Duration::from_secs(TTL_CLEANUP_INTERVAL_SECS)); - let mut wildcard_reconciliation_interval = interval(Duration::from_mins(1)); let mut subscription_reconciliation_interval = interval(Duration::from_mins(5)); let mut mesh_check_interval = interval(Duration::from_secs(MESH_CHECK_INTERVAL_SECS)); let mut retained_sync_cleanup_interval = @@ -140,9 +139,6 @@ impl ClusteredAgent { _ = cleanup_interval.tick() => { self.handle_session_cleanup().await; } - _ = wildcard_reconciliation_interval.tick() => { - self.handle_wildcard_reconciliation().await; - } _ = subscription_reconciliation_interval.tick() => { self.handle_subscription_reconciliation().await; } @@ -236,11 +232,10 @@ impl ClusteredAgent { let result = stores .subscriptions .reconcile(&stores.topics, &stores.wildcards); - if result.subscriptions_added > 0 || result.subscriptions_removed > 0 { + if result.index_entries_restored > 0 { info!( clients = result.clients_checked, - added = result.subscriptions_added, - removed = result.subscriptions_removed, + restored = result.index_entries_restored, "reconciled subscriptions after partition takeover" ); } @@ -595,37 +590,6 @@ impl ClusteredAgent { ); } - async fn handle_wildcard_reconciliation(&self) { - let now = current_time_ms(); - let ctrl = self.controller.read().await; - let pending_store = &ctrl.stores().wildcard_pending; - if pending_store.needs_reconciliation(now) { - let pending = pending_store.get_pending_for_retry(); - if !pending.is_empty() { - debug!( - count = pending.len(), - "retrying pending wildcard broadcasts" - ); - for p in &pending { - let broadcast = p.to_broadcast(); - let msg = ClusterMessage::WildcardBroadcast(broadcast); - let _ = ctrl.transport().broadcast(msg).await; - pending_store.mark_retried(&p.pattern, &p.client_id); - } - info!( - count = pending.len(), - "rebroadcast pending wildcard subscriptions" - ); - } - pending_store.mark_reconciliation(now); - let max_age_ms = 5 * 60 * 1000; - let removed = pending_store.clear_old_entries(now, max_age_ms); - if removed > 0 { - debug!(removed, "cleared old wildcard pending entries"); - } - } - } - async fn handle_subscription_reconciliation(&self) { let now = current_time_ms(); let ctrl = self.controller.read().await; @@ -634,11 +598,10 @@ impl ClusteredAgent { let result = stores .subscriptions .reconcile(&stores.topics, &stores.wildcards); - if result.subscriptions_added > 0 || result.subscriptions_removed > 0 { + if result.index_entries_restored > 0 { info!( clients = result.clients_checked, - added = result.subscriptions_added, - removed = result.subscriptions_removed, + restored = result.index_entries_restored, "reconciled subscription cache" ); } diff --git a/crates/mqdb-cluster/src/lib.rs b/crates/mqdb-cluster/src/lib.rs index 5b5bd9d1..e57de4a9 100644 --- a/crates/mqdb-cluster/src/lib.rs +++ b/crates/mqdb-cluster/src/lib.rs @@ -27,30 +27,30 @@ pub use cluster::{ MigrationState, MqttSubscriptionSnapshot, MqttTopicEntry, NUM_PARTITIONS, NodeController, NodeId, NodeStatus, OffsetStore, OffsetStoreError, Operation, ParsedDbTopic, PartitionAssignment, PartitionCursor, PartitionId, PartitionMap, PartitionReassignment, - PartitionRole, PartitionStorage, PartitionWriteLog, PendingWildcard, PendingWrites, - ProcessingBatch, PublishRouteResult, PublishRouter, Qos2Direction, Qos2Phase, Qos2State, - Qos2Store, Qos2StoreError, QueryCoordinator, QueryRequest, QueryResponse, QueryResult, - QueryStatus, QuicDirectTransport, QuorumResult, QuorumTracker, RaftAdminCommand, RaftEvent, - RaftMessage, RaftStatus, RaftTask, RebalanceAck, RebalanceCommit, RebalanceConfig, - RebalanceCoordinator, RebalanceError, RebalanceProposal, RebalanceState, ReconciliationResult, - RecoveryStats, ReplicaRole, ReplicaState, ReplicationAck, ReplicationError, ReplicationWrite, - RetainedMessage, RetainedStore, RetainedStoreError, RoutingTarget, ScatterCursor, SessionData, - SessionError, SessionStore, SnapshotBuilder, SnapshotChunk, SnapshotComplete, SnapshotRequest, + PartitionRole, PartitionStorage, PartitionWriteLog, PendingWrites, ProcessingBatch, + PublishRouteResult, PublishRouter, Qos2Direction, Qos2Phase, Qos2State, Qos2Store, + Qos2StoreError, QueryCoordinator, QueryRequest, QueryResponse, QueryResult, QueryStatus, + QuicDirectTransport, QuorumResult, QuorumTracker, RaftAdminCommand, RaftEvent, RaftMessage, + RaftStatus, RaftTask, RebalanceAck, RebalanceCommit, RebalanceConfig, RebalanceCoordinator, + RebalanceError, RebalanceProposal, RebalanceState, ReconciliationResult, RecoveryStats, + ReplicaRole, ReplicaState, ReplicationAck, ReplicationError, ReplicationWrite, RetainedMessage, + RetainedStore, RetainedStoreError, RoutingTarget, ScatterCursor, SessionData, SessionError, + SessionStore, SnapshotBuilder, SnapshotChunk, SnapshotComplete, SnapshotRequest, SnapshotSender, SnapshotStatus, StoreApplyError, StoreManager, SubscriberLocation, SubscriptionCache, SubscriptionCacheError, SubscriptionType, TickOutput, TopicIndex, TopicIndexEntry, TopicIndexError, TopicSubscriptionBroadcast, TopicTrie, TransportConfig, TransportError, UniqueCommitRequest, UniqueCommitResponse, UniqueReleaseRequest, UniqueReleaseResponse, UniqueReserveRequest, UniqueReserveResponse, UniqueReserveStatus, - WildcardBroadcast, WildcardEntry, WildcardOp, WildcardPendingStore, WildcardStore, - WildcardStoreError, WildcardSubscriber, + WildcardBroadcast, WildcardEntry, WildcardOp, WildcardStore, WildcardStoreError, + WildcardSubscriber, }; pub use cluster::{ - SUBSCRIPTION_RECONCILIATION_INTERVAL_MS, WILDCARD_RECONCILIATION_INTERVAL_MS, - client_location_key, compute_balanced_assignments, compute_incremental_assignments, - compute_removal_assignments, data_partition, db_data_key, determine_lwt_action, effective_qos, - generate_lwt_token, idempotency_storage_key, inflight_key, is_wildcard_pattern, - mqtt_subscription_key, offset_key, qos2_state_key, retained_message_key, session_key, - session_partition, topic_index_key, topic_partition, validate_pattern, wildcard_key, + SUBSCRIPTION_RECONCILIATION_INTERVAL_MS, client_location_key, compute_balanced_assignments, + compute_incremental_assignments, compute_removal_assignments, data_partition, db_data_key, + determine_lwt_action, effective_qos, generate_lwt_token, idempotency_storage_key, inflight_key, + is_wildcard_pattern, mqtt_subscription_key, offset_key, qos2_state_key, retained_message_key, + session_key, session_partition, topic_index_key, topic_partition, validate_pattern, + wildcard_key, }; pub use cluster_agent::ClusterTransportKind; pub use cluster_agent::{ClusterConfig, ClusterInitError, ClusteredAgent, PeerConfig, QuicConfig}; diff --git a/docs/distributed-design.md b/docs/distributed-design.md index cf7a613d..16dc8e40 100644 --- a/docs/distributed-design.md +++ b/docs/distributed-design.md @@ -162,7 +162,6 @@ The `StoreManager` coordinates 17 distinct stores behind a common interface: - `retained` - Retained messages by topic - `topics` - TopicIndex (topic → subscribers) - `wildcards` - Wildcard subscription trie -- `wildcard_pending` - Pending wildcard operations - `qos2` - QoS 2 exactly-once state machines - `inflight` - QoS 1 messages awaiting acknowledgment - `offsets` - Consumer group offsets @@ -951,8 +950,7 @@ All cluster-managed data types (`crates/mqdb-cluster/src/cluster/entity.rs`): | Constant | Value | File | |----------|-------|------| -| `SUBSCRIPTION_RECONCILIATION_INTERVAL_MS` | 300,000 (5 min) | `subscription_cache.rs:7` | -| `WILDCARD_RECONCILIATION_INTERVAL_MS` | 60,000 (1 min) | `wildcard_pending.rs:6` | +| `SUBSCRIPTION_RECONCILIATION_INTERVAL_MS` | 300,000 (5 min) | `subscription_cache.rs:10` | | `CATCHUP_REQUEST_INTERVAL_MS` | 5,000 | `replication.rs:24` | | `OFFSET_STALE_TTL_MS` | 604,800,000 (7 days) | `offset_store.rs:86` | | `DEFAULT_QUERY_TIMEOUT_MS` | 10,000 | `query_coordinator.rs:7` | @@ -2242,7 +2240,6 @@ Peer format: `{node_id}@{address}:{port}` (comma-separated for multiple) 6. Register peers with Controller and Raft 7. Start main event loop: - 100ms: Controller tick, Raft tick, heartbeats - - 60s: Wildcard reconciliation - 300s: Subscription reconciliation - 3600s: Cleanup expired sessions/idempotency ``` @@ -2374,7 +2371,7 @@ delete itself. The following are implemented but documented in other files: - **LWT (Last Will Testament)**: Implemented in M9, see `docs/implementation-plan.md` -- **Subscription cache reconciliation**: `SubscriptionCache::reconcile()` runs every 5 minutes +- **Subscription reconciliation**: `SubscriptionCache::reconcile()` runs every 5 minutes and after a partition takeover. The replicated subscription record is authoritative: reconcile re-adds any recorded subscription missing from the local TopicIndex/WildcardStore (response topics excepted, since they are never broadcast) and never modifies the record. Model: `specs/ClusterSubReconcile.tla` - **Database operations ($DB/#)**: See README.md "MQTT Topic Structure" section - **Backup and restore**: See README.md "Admin Operations" section diff --git a/specs/ClusterSubReconcile.cfg b/specs/ClusterSubReconcile.cfg new file mode 100644 index 00000000..eed00389 --- /dev/null +++ b/specs/ClusterSubReconcile.cfg @@ -0,0 +1,9 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + MaxOps = 3 + Mode = "new" + Lossy = TRUE +INVARIANTS TypeOK InvNoReconcileLoss InvNoResurrection +PROPERTY RoutingRepaired +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile.tla b/specs/ClusterSubReconcile.tla new file mode 100644 index 00000000..5be5c56c --- /dev/null +++ b/specs/ClusterSubReconcile.tla @@ -0,0 +1,139 @@ +------------------------- MODULE ClusterSubReconcile ------------------------- +(***************************************************************************) +(* Design model for gh #141: SubscriptionCache::reconcile vs a missed *) +(* topic-index broadcast. One client, subscribed at its partition primary *) +(* P. The durable subscription record reaches the replica R by partition *) +(* replication; the topic-index entry reaches every other node only by a *) +(* one-hop broadcast (routing.rs handle_topic_subscribe). apply_subscription*) +(* never touches the index. Replication and the index broadcast share the *) +(* P->R bulk lane (FIFO, lossy under backpressure). O holds no record. *) +(* *) +(* Mode "old" = current reconcile: for a client present in the cache, the *) +(* record is overwritten with the local index. *) +(* Mode "new" = proposed: the index gains every recorded topic; the record *) +(* is never modified by reconcile. *) +(* Mode "none" = no reconcile (shows the liveness property is non-vacuous). *) +(* *) +(* Configs: *) +(* ClusterSubReconcile.cfg new -> safety holds, holders' index repaired *) +(* *_oldloss.cfg old -> InvNoReconcileLoss violated *) +(* *_oldresurrect.cfg old -> InvNoResurrection violated *) +(* *_norepair.cfg none -> RoutingRepaired violated *) +(* *_nonholder.cfg new -> NonHolderIndexed violated: a node *) +(* that holds no record and misses the broadcast is NOT repaired by *) +(* this fix (residual gap, out of scope, tracked separately). *) +(* *) +(* NOT MODELLED: the on-disk copy (reconcile edits memory, and memory feeds *) +(* export_partition); explicit promotion (reconcile after takeover is just *) +(* a Reconcile step); loss of the replicated write itself (async RF=2 *) +(* durability, independent of reconcile). Bounded by MaxOps. *) +(***************************************************************************) +EXTENDS Naturals, Sequences, FiniteSets + +CONSTANTS Topics, MaxOps, Mode, Lossy + +ASSUME Mode \in {"old", "new", "none"} +ASSUME MaxOps \in Nat /\ Lossy \in BOOLEAN + +Nodes == {"P", "R", "O"} +Holders == {"P", "R"} +Remote == {"R", "O"} + +Msg == [k : {"repl", "idx"}, op : {"add", "del"}, t : Topics] + +VARIABLES truth, rec, idx, ch, ops, lost, resurrected + +vars == <> + +TypeOK == + /\ truth \subseteq Topics + /\ rec \in [Holders -> SUBSET Topics] + /\ idx \in [Nodes -> SUBSET Topics] + /\ \A n \in Remote : ch[n] \in Seq(Msg) + /\ ops \in 0..MaxOps + /\ lost \in BOOLEAN /\ resurrected \in BOOLEAN + +Init == + /\ truth = {} + /\ rec = [n \in Holders |-> {}] + /\ idx = [n \in Nodes |-> {}] + /\ ch = [n \in Remote |-> <<>>] + /\ ops = 0 + /\ lost = FALSE + /\ resurrected = FALSE + +Apply(s, op, t) == IF op = "add" THEN s \cup {t} ELSE s \ {t} + +Emit(op, t) == + [ch EXCEPT !["R"] = Append(Append(@, [k |-> "repl", op |-> op, t |-> t]), + [k |-> "idx", op |-> op, t |-> t]), + !["O"] = Append(@, [k |-> "idx", op |-> op, t |-> t])] + +Change(op, t) == + /\ ops < MaxOps + /\ truth' = Apply(truth, op, t) + /\ rec' = [rec EXCEPT !["P"] = Apply(@, op, t)] + /\ idx' = [idx EXCEPT !["P"] = Apply(@, op, t)] + /\ ch' = Emit(op, t) + /\ ops' = ops + 1 + /\ UNCHANGED <> + +Sub(t) == t \notin truth /\ Change("add", t) +Unsub(t) == t \in truth /\ Change("del", t) + +Deliver(n) == + /\ ch[n] # <<>> + /\ LET m == Head(ch[n]) IN + /\ ch' = [ch EXCEPT ![n] = Tail(@)] + /\ IF m.k = "repl" + THEN /\ rec' = [rec EXCEPT ![n] = Apply(@, m.op, m.t)] + /\ UNCHANGED idx + ELSE /\ idx' = [idx EXCEPT ![n] = Apply(@, m.op, m.t)] + /\ UNCHANGED rec + /\ UNCHANGED <> + +Drop(n, i) == + /\ Lossy + /\ i \in 1..Len(ch[n]) + /\ ch' = [ch EXCEPT ![n] = SubSeq(@, 1, i - 1) \o SubSeq(@, i + 1, Len(@))] + /\ UNCHANGED <> + +\* The real reconcile iterates only clients present in the cache, so it +\* does nothing for a node whose record for the client is empty. +Reconcile(n) == + /\ Mode # "none" + /\ rec[n] # {} + /\ IF Mode = "old" + THEN /\ rec' = [rec EXCEPT ![n] = idx[n]] + /\ lost' = (lost \/ ((rec[n] \ idx[n]) \cap truth # {})) + /\ resurrected' = (resurrected \/ ((idx[n] \ rec[n]) \ truth # {})) + /\ UNCHANGED idx + ELSE /\ idx' = [idx EXCEPT ![n] = @ \cup rec[n]] + /\ UNCHANGED <> + /\ UNCHANGED <> + +Next == + \/ \E t \in Topics : Sub(t) \/ Unsub(t) + \/ \E n \in Remote : Deliver(n) + \/ \E n \in Remote : \E i \in 1..Len(ch[n]) : Drop(n, i) + \/ \E n \in Holders : Reconcile(n) + +Spec == + /\ Init /\ [][Next]_vars + /\ \A n \in Remote : WF_vars(Deliver(n)) + /\ \A n \in Holders : WF_vars(Reconcile(n)) + +---------------------------------------------------------------------------- +\* SAFETY: reconcile never deletes a subscription the client still holds, +\* and never brings back one the client dropped. +InvNoReconcileLoss == ~lost +InvNoResurrection == ~resurrected + +\* LIVENESS: every node that holds the record ends up with every recorded +\* topic in its index (routing on that node is repaired). +HoldersIndexed == \A n \in Holders : rec[n] \subseteq idx[n] +RoutingRepaired == <>[]HoldersIndexed + +\* SCOPE: a node holding no record is never repaired if it missed a broadcast. +NonHolderIndexed == <>[](truth \subseteq idx["O"]) +============================================================================= diff --git a/specs/ClusterSubReconcile_nonholder.cfg b/specs/ClusterSubReconcile_nonholder.cfg new file mode 100644 index 00000000..d3ba5dd3 --- /dev/null +++ b/specs/ClusterSubReconcile_nonholder.cfg @@ -0,0 +1,9 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + MaxOps = 3 + Mode = "new" + Lossy = TRUE +INVARIANTS TypeOK +PROPERTY NonHolderIndexed +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile_norepair.cfg b/specs/ClusterSubReconcile_norepair.cfg new file mode 100644 index 00000000..f1158dc3 --- /dev/null +++ b/specs/ClusterSubReconcile_norepair.cfg @@ -0,0 +1,9 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + MaxOps = 3 + Mode = "none" + Lossy = TRUE +INVARIANTS TypeOK +PROPERTY RoutingRepaired +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile_oldloss.cfg b/specs/ClusterSubReconcile_oldloss.cfg new file mode 100644 index 00000000..05f9c7f7 --- /dev/null +++ b/specs/ClusterSubReconcile_oldloss.cfg @@ -0,0 +1,9 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + MaxOps = 3 + Mode = "old" + Lossy = TRUE +INVARIANTS TypeOK InvNoReconcileLoss + +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile_oldresurrect.cfg b/specs/ClusterSubReconcile_oldresurrect.cfg new file mode 100644 index 00000000..bcac1a7c --- /dev/null +++ b/specs/ClusterSubReconcile_oldresurrect.cfg @@ -0,0 +1,9 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + MaxOps = 3 + Mode = "old" + Lossy = TRUE +INVARIANTS TypeOK InvNoResurrection + +CHECK_DEADLOCK FALSE From 6a51916bcfe852c2f0819294081bbedbeffbb759 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Fabr=C3=ADcio=20Bracht?= Date: Wed, 23 Sep 2026 14:35:12 -0300 Subject: [PATCH 2/2] restrict subscription reconcile to held partitions and skip response topics in recovery --- .../src/cluster/store_manager/recovery.rs | 11 +- .../src/cluster/store_manager/tests.rs | 33 ++++ .../src/cluster/subscription_cache.rs | 53 ++++++- .../src/cluster_agent/event_loop.rs | 18 ++- docs/distributed-design.md | 7 +- specs/ClusterSubReconcile.cfg | 5 +- specs/ClusterSubReconcile.tla | 141 ++++++++++++------ specs/ClusterSubReconcile_holderghost.cfg | 9 ++ specs/ClusterSubReconcile_nonholder.cfg | 3 +- specs/ClusterSubReconcile_norepair.cfg | 3 +- specs/ClusterSubReconcile_oldloss.cfg | 2 +- specs/ClusterSubReconcile_oldresurrect.cfg | 2 +- specs/ClusterSubReconcile_unrestricted.cfg | 9 ++ 13 files changed, 226 insertions(+), 70 deletions(-) create mode 100644 specs/ClusterSubReconcile_holderghost.cfg create mode 100644 specs/ClusterSubReconcile_unrestricted.cfg diff --git a/crates/mqdb-cluster/src/cluster/store_manager/recovery.rs b/crates/mqdb-cluster/src/cluster/store_manager/recovery.rs index 83aa465c..7c9fdbde 100644 --- a/crates/mqdb-cluster/src/cluster/store_manager/recovery.rs +++ b/crates/mqdb-cluster/src/cluster/store_manager/recovery.rs @@ -4,7 +4,7 @@ use super::{RecoveryField, RecoveryStats, StoreApplyError, StoreManager}; use crate::cluster::protocol::{Operation, ReplicationWrite}; use crate::cluster::session::session_partition; -use crate::cluster::{Epoch, SubscriptionType, entity}; +use crate::cluster::{Epoch, SubscriptionType, entity, is_response_topic}; impl StoreManager { /// # Errors @@ -68,10 +68,11 @@ impl StoreManager { let client_id = snapshot.client_id_str(); let partition = session_partition(client_id); for entry in snapshot.exact_subscriptions() { - if self - .topics - .subscribe(entry.topic_str(), client_id, partition, entry.qos) - .is_ok() + if !is_response_topic(entry.topic_str()) + && self + .topics + .subscribe(entry.topic_str(), client_id, partition, entry.qos) + .is_ok() { count += 1; } diff --git a/crates/mqdb-cluster/src/cluster/store_manager/tests.rs b/crates/mqdb-cluster/src/cluster/store_manager/tests.rs index 939eec75..b6735938 100644 --- a/crates/mqdb-cluster/src/cluster/store_manager/tests.rs +++ b/crates/mqdb-cluster/src/cluster/store_manager/tests.rs @@ -73,3 +73,36 @@ fn apply_unknown_entity_fails() { let result = manager.apply_write(&write); assert!(matches!(result, Err(StoreApplyError::UnknownEntity))); } + +#[test] +fn recovery_rebuilds_index_without_response_topics() { + let backend: std::sync::Arc = + std::sync::Arc::new(mqdb_core::MemoryBackend::new()); + let writer = StoreManager::new_with_storage(node(1), Some(backend.clone())); + + let mut snapshot = crate::cluster::MqttSubscriptionSnapshot::create("client1"); + snapshot.add_subscription("sensors/a", 1); + snapshot.add_subscription("resp/client1", 0); + let write = ReplicationWrite::new( + crate::cluster::session_partition("client1"), + Operation::Insert, + Epoch::new(1), + 1, + entity::SUBSCRIPTIONS.to_string(), + "client1".to_string(), + crate::cluster::SubscriptionCache::serialize(&snapshot), + ); + writer.apply_write(&write).unwrap(); + + let recovered = StoreManager::new_with_storage(node(1), Some(backend)); + recovered.recover().unwrap(); + + assert_eq!( + recovered.subscriptions.get_subscriptions("client1").len(), + 2 + ); + assert_eq!( + recovered.topics.get_client_topics("client1"), + vec![("sensors/a".to_string(), 1)] + ); +} diff --git a/crates/mqdb-cluster/src/cluster/subscription_cache.rs b/crates/mqdb-cluster/src/cluster/subscription_cache.rs index 862d6030..56b03bdc 100644 --- a/crates/mqdb-cluster/src/cluster/subscription_cache.rs +++ b/crates/mqdb-cluster/src/cluster/subscription_cache.rs @@ -199,12 +199,16 @@ impl SubscriptionCache { &self, topic_index: &TopicIndex, wildcard_store: &WildcardStore, + holds_partition: impl Fn(PartitionId) -> bool, ) -> ReconciliationResult { let mut result = ReconciliationResult::default(); for snapshot in self.all_snapshots() { - result.clients_checked += 1; let client_id = snapshot.client_id_str(); + if !holds_partition(session_partition(client_id)) { + continue; + } + result.clients_checked += 1; let indexed_exact: HashSet = topic_index .get_client_topics(client_id) @@ -699,7 +703,7 @@ mod tests { .subscribe("topic/a", "client1", partition(), 1) .unwrap(); - let result = cache.reconcile(&topic_index, &wildcard_store); + let result = cache.reconcile(&topic_index, &wildcard_store, |_| true); assert_eq!(result.clients_checked, 1); assert_eq!(result.index_entries_restored, 2); @@ -737,7 +741,7 @@ mod tests { .subscribe_mqtt("dropped/+/pattern", "client1", 1) .unwrap(); - let result = cache.reconcile(&topic_index, &wildcard_store); + let result = cache.reconcile(&topic_index, &wildcard_store, |_| true); assert_eq!(result.index_entries_restored, 0); let topics: Vec = cache @@ -761,10 +765,45 @@ mod tests { .add_subscription("client1", "app/resp/client1", 0) .unwrap(); - let result = cache.reconcile(&topic_index, &wildcard_store); + let result = cache.reconcile(&topic_index, &wildcard_store, |_| true); assert_eq!(result.index_entries_restored, 0); assert!(topic_index.get_client_topics("client1").is_empty()); + assert_eq!(cache.get_subscriptions("client1").len(), 2); + } + + #[test] + fn reconcile_ignores_records_for_partitions_not_held() { + let cache = SubscriptionCache::new(node(1)); + let topic_index = TopicIndex::new(node(1)); + let wildcard_store = WildcardStore::new(node(1)); + + cache.add_subscription("held-client", "topic/a", 1).unwrap(); + cache + .add_subscription("stale-client", "topic/a", 1) + .unwrap(); + cache + .add_subscription("stale-client", "stale/+/pattern", 1) + .unwrap(); + let held = session_partition("held-client"); + assert_ne!(held, session_partition("stale-client")); + + let result = cache.reconcile(&topic_index, &wildcard_store, |partition| partition == held); + + assert_eq!(result.clients_checked, 1); + assert_eq!(result.index_entries_restored, 1); + let subscribers: Vec = topic_index + .get_subscribers("topic/a") + .iter() + .map(|s| s.client_id_str().to_string()) + .collect(); + assert_eq!(subscribers, vec!["held-client".to_string()]); + assert!( + wildcard_store + .get_client_patterns("stale-client") + .is_empty() + ); + assert_eq!(cache.get_subscriptions("stale-client").len(), 2); } #[test] @@ -778,13 +817,13 @@ mod tests { assert_eq!( cache - .reconcile(&topic_index, &wildcard_store) + .reconcile(&topic_index, &wildcard_store, |_| true) .index_entries_restored, 2 ); assert_eq!( cache - .reconcile(&topic_index, &wildcard_store) + .reconcile(&topic_index, &wildcard_store, |_| true) .index_entries_restored, 0 ); @@ -808,7 +847,7 @@ mod tests { ) .unwrap(); - let result = replica.reconcile(&topic_index, &wildcard_store); + let result = replica.reconcile(&topic_index, &wildcard_store, |_| true); assert_eq!(result.index_entries_restored, 2); assert_eq!( diff --git a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs index 16b2e0d5..244b2edf 100644 --- a/crates/mqdb-cluster/src/cluster_agent/event_loop.rs +++ b/crates/mqdb-cluster/src/cluster_agent/event_loop.rs @@ -229,9 +229,12 @@ impl ClusteredAgent { } if became_primary { let stores = ctrl.stores(); - let result = stores - .subscriptions - .reconcile(&stores.topics, &stores.wildcards); + let result = + stores + .subscriptions + .reconcile(&stores.topics, &stores.wildcards, |partition| { + ctrl.can_serve_reads(partition) + }); if result.index_entries_restored > 0 { info!( clients = result.clients_checked, @@ -595,9 +598,12 @@ impl ClusteredAgent { let ctrl = self.controller.read().await; let stores = ctrl.stores(); if stores.subscriptions.needs_reconciliation(now) { - let result = stores - .subscriptions - .reconcile(&stores.topics, &stores.wildcards); + let result = + stores + .subscriptions + .reconcile(&stores.topics, &stores.wildcards, |partition| { + ctrl.can_serve_reads(partition) + }); if result.index_entries_restored > 0 { info!( clients = result.clients_checked, diff --git a/docs/distributed-design.md b/docs/distributed-design.md index 16dc8e40..0eb476a9 100644 --- a/docs/distributed-design.md +++ b/docs/distributed-design.md @@ -156,7 +156,7 @@ The key distinction: exact topic subscriptions are lightweight (1 replicated wri The `StoreManager` coordinates 17 distinct stores behind a common interface: -**MQTT Stores (11)**: +**MQTT Stores (10)**: - `sessions` - Client session lifecycle - `subscriptions` - Topic subscriptions per client - `retained` - Retained messages by topic @@ -168,13 +168,14 @@ The `StoreManager` coordinates 17 distinct stores behind a common interface: - `idempotency` - Request deduplication tokens - `client_locations` - Client → connected node mapping -**Database Stores (6)**: +**Database Stores (7)**: - `db_data` - Entity records - `db_schema` - Entity schemas with constraints - `db_index` - Secondary indexes - `db_unique` - Unique constraint reservations - `db_fk` - Foreign key validation - `db_constraints` - Constraint definitions +- `fk_reverse_index` - Reverse index from referenced records to referencing records ### A3.2 The apply_write() Dispatcher @@ -2371,7 +2372,7 @@ delete itself. The following are implemented but documented in other files: - **LWT (Last Will Testament)**: Implemented in M9, see `docs/implementation-plan.md` -- **Subscription reconciliation**: `SubscriptionCache::reconcile()` runs every 5 minutes and after a partition takeover. The replicated subscription record is authoritative: reconcile re-adds any recorded subscription missing from the local TopicIndex/WildcardStore (response topics excepted, since they are never broadcast) and never modifies the record. Model: `specs/ClusterSubReconcile.tla` +- **Subscription reconciliation**: `SubscriptionCache::reconcile()` runs every 5 minutes and after a partition takeover. The replicated subscription record is authoritative: for clients whose session partition this node holds as primary or replica, reconcile re-adds any recorded subscription missing from the local TopicIndex/WildcardStore and never modifies the record. Exact response topics are skipped, as they are never broadcast. Copies held by ex-owners or by the node a client used to connect through are ignored, so they cannot recreate index entries. The index is never pruned, so an entry whose unsubscribe broadcast was missed stays until the client's subscriptions are cleared; it only costs a wasted forward, since the receiving node delivers to its local subscribers only. Model: `specs/ClusterSubReconcile.tla` - **Database operations ($DB/#)**: See README.md "MQTT Topic Structure" section - **Backup and restore**: See README.md "Admin Operations" section diff --git a/specs/ClusterSubReconcile.cfg b/specs/ClusterSubReconcile.cfg index eed00389..94e10166 100644 --- a/specs/ClusterSubReconcile.cfg +++ b/specs/ClusterSubReconcile.cfg @@ -1,9 +1,10 @@ SPECIFICATION Spec CONSTANTS Topics = {"a", "b"} - MaxOps = 3 + MaxOps = 2 Mode = "new" Lossy = TRUE -INVARIANTS TypeOK InvNoReconcileLoss InvNoResurrection + Restrict = TRUE +INVARIANTS TypeOK InvNoReconcileLoss InvNoResurrection InvNoStaleCopyGhost PROPERTY RoutingRepaired CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile.tla b/specs/ClusterSubReconcile.tla index 5be5c56c..06bb147e 100644 --- a/specs/ClusterSubReconcile.tla +++ b/specs/ClusterSubReconcile.tla @@ -6,68 +6,110 @@ (* replication; the topic-index entry reaches every other node only by a *) (* one-hop broadcast (routing.rs handle_topic_subscribe). apply_subscription*) (* never touches the index. Replication and the index broadcast share the *) -(* P->R bulk lane (FIFO, lossy under backpressure). O holds no record. *) +(* P->R bulk lane (FIFO). Any in-flight message, replication writes *) +(* included, can be dropped (Drop). O never holds a record. S is a former *) +(* holder (ex-owner, or the node the client used to connect through): it *) +(* receives replication until LoseOwnership, then keeps a stale copy that *) +(* nothing deletes (SubscriptionCache::clear_partition has no production *) +(* caller), while still receiving index broadcasts. *) (* *) -(* Mode "old" = current reconcile: for a client present in the cache, the *) -(* record is overwritten with the local index. *) -(* Mode "new" = proposed: the index gains every recorded topic; the record *) -(* is never modified by reconcile. *) +(* Mode "old" = main before #151: for every client present in the cache, *) +(* the record is overwritten with the local index. *) +(* Mode "new" = #151: the index gains every recorded topic; the record is *) +(* never modified. Restrict = only for partitions this node *) +(* currently holds (primary or replica). *) (* Mode "none" = no reconcile (shows the liveness property is non-vacuous). *) (* *) +(* The loss/resurrection flags are computed from the record change of every *) +(* Reconcile step in every mode. In mode "new" the record never changes, so *) +(* those two invariants hold by construction; the checker results that *) +(* matter for "new" are RoutingRepaired and the ghost invariants. *) +(* *) (* Configs: *) -(* ClusterSubReconcile.cfg new -> safety holds, holders' index repaired *) -(* *_oldloss.cfg old -> InvNoReconcileLoss violated *) -(* *_oldresurrect.cfg old -> InvNoResurrection violated *) -(* *_norepair.cfg none -> RoutingRepaired violated *) -(* *_nonholder.cfg new -> NonHolderIndexed violated: a node *) -(* that holds no record and misses the broadcast is NOT repaired by *) -(* this fix (residual gap, out of scope, tracked separately). *) +(* ClusterSubReconcile.cfg new, Restrict -> loss/resurrection/stale- *) +(* copy-ghost invariants and RoutingRepaired hold (MaxOps = 2, the *) +(* largest bound where the liveness check is exhaustive) *) +(* *_oldloss.cfg old -> InvNoReconcileLoss violated *) +(* *_oldresurrect.cfg old -> InvNoResurrection violated *) +(* *_norepair.cfg none -> RoutingRepaired violated *) +(* *_nonholder.cfg new -> NonHolderIndexed violated: a node *) +(* that holds no record and misses the broadcast is never repaired *) +(* (tracked in gh #140) *) +(* *_unrestricted.cfg new, no Restrict -> InvNoStaleCopyGhost *) +(* violated: a former holder's stale copy recreates an index entry *) +(* for a topic the client dropped *) +(* *_holderghost.cfg new, Restrict -> InvNoHolderGhost violated:*) +(* a holder whose record lags the index (replication delete delayed *) +(* or dropped behind a delivered unsubscribe broadcast) re-adds an *) +(* entry; the index is never pruned, so it stays. Accepted: the *) +(* receiving node only delivers to its local subscribers, so the *) +(* cost is a wasted forward, not a misdelivery *) +(* *) +(* The old-mode guard rec[n] # {} under-approximates the real code, where *) +(* the replica keeps an empty snapshot after the last unsubscribe; it can *) +(* only hide old-mode violations, never create them. *) (* *) -(* NOT MODELLED: the on-disk copy (reconcile edits memory, and memory feeds *) -(* export_partition); explicit promotion (reconcile after takeover is just *) -(* a Reconcile step); loss of the replicated write itself (async RF=2 *) -(* durability, independent of reconcile). Bounded by MaxOps. *) +(* NOT MODELLED: the on-disk copy and restart recovery; explicit promotion *) +(* (reconcile after takeover is a Reconcile step); snapshot import of *) +(* TOPIC_INDEX (an extra repair source for topic-partition holders); *) +(* replication catch-up of a dropped write; a client connected to a node *) +(* other than P, where the record and the broadcast take different paths *) +(* and are not FIFO-ordered. Bounded by MaxOps. *) (***************************************************************************) EXTENDS Naturals, Sequences, FiniteSets -CONSTANTS Topics, MaxOps, Mode, Lossy +CONSTANTS Topics, MaxOps, Mode, Lossy, Restrict ASSUME Mode \in {"old", "new", "none"} -ASSUME MaxOps \in Nat /\ Lossy \in BOOLEAN +ASSUME MaxOps \in Nat /\ Lossy \in BOOLEAN /\ Restrict \in BOOLEAN -Nodes == {"P", "R", "O"} +Nodes == {"P", "R", "O", "S"} Holders == {"P", "R"} -Remote == {"R", "O"} +Cached == {"P", "R", "S"} +Remote == {"R", "O", "S"} Msg == [k : {"repl", "idx"}, op : {"add", "del"}, t : Topics] -VARIABLES truth, rec, idx, ch, ops, lost, resurrected +VARIABLES truth, rec, idx, ch, ops, sHolds, lost, resurrected, staleGhost, holderGhost -vars == <> +vars == <> TypeOK == /\ truth \subseteq Topics - /\ rec \in [Holders -> SUBSET Topics] + /\ rec \in [Cached -> SUBSET Topics] /\ idx \in [Nodes -> SUBSET Topics] /\ \A n \in Remote : ch[n] \in Seq(Msg) /\ ops \in 0..MaxOps + /\ sHolds \in BOOLEAN /\ lost \in BOOLEAN /\ resurrected \in BOOLEAN + /\ staleGhost \in BOOLEAN /\ holderGhost \in BOOLEAN Init == /\ truth = {} - /\ rec = [n \in Holders |-> {}] + /\ rec = [n \in Cached |-> {}] /\ idx = [n \in Nodes |-> {}] /\ ch = [n \in Remote |-> <<>>] /\ ops = 0 + /\ sHolds = TRUE /\ lost = FALSE /\ resurrected = FALSE + /\ staleGhost = FALSE + /\ holderGhost = FALSE + +Holds(n) == n \in Holders \/ (n = "S" /\ sHolds) Apply(s, op, t) == IF op = "add" THEN s \cup {t} ELSE s \ {t} +ReplMsg(op, t) == [k |-> "repl", op |-> op, t |-> t] +IdxMsg(op, t) == [k |-> "idx", op |-> op, t |-> t] + Emit(op, t) == - [ch EXCEPT !["R"] = Append(Append(@, [k |-> "repl", op |-> op, t |-> t]), - [k |-> "idx", op |-> op, t |-> t]), - !["O"] = Append(@, [k |-> "idx", op |-> op, t |-> t])] + [n \in Remote |-> + CASE n = "R" -> ch[n] \o <> + [] n = "S" -> IF sHolds + THEN ch[n] \o <> + ELSE Append(ch[n], IdxMsg(op, t)) + [] OTHER -> Append(ch[n], IdxMsg(op, t))] Change(op, t) == /\ ops < MaxOps @@ -76,11 +118,17 @@ Change(op, t) == /\ idx' = [idx EXCEPT !["P"] = Apply(@, op, t)] /\ ch' = Emit(op, t) /\ ops' = ops + 1 - /\ UNCHANGED <> + /\ UNCHANGED <> Sub(t) == t \notin truth /\ Change("add", t) Unsub(t) == t \in truth /\ Change("del", t) +LoseOwnership == + /\ sHolds + /\ sHolds' = FALSE + /\ ch' = [ch EXCEPT !["S"] = SelectSeq(@, LAMBDA m : m.k = "idx")] + /\ UNCHANGED <> + Deliver(n) == /\ ch[n] # <<>> /\ LET m == Head(ch[n]) IN @@ -90,33 +138,35 @@ Deliver(n) == /\ UNCHANGED idx ELSE /\ idx' = [idx EXCEPT ![n] = Apply(@, m.op, m.t)] /\ UNCHANGED rec - /\ UNCHANGED <> + /\ UNCHANGED <> Drop(n, i) == /\ Lossy /\ i \in 1..Len(ch[n]) /\ ch' = [ch EXCEPT ![n] = SubSeq(@, 1, i - 1) \o SubSeq(@, i + 1, Len(@))] - /\ UNCHANGED <> + /\ UNCHANGED <> + +NewRec(n) == IF Mode = "old" THEN idx[n] ELSE rec[n] +NewIdx(n) == IF Mode = "old" THEN idx[n] ELSE idx[n] \cup rec[n] -\* The real reconcile iterates only clients present in the cache, so it -\* does nothing for a node whose record for the client is empty. Reconcile(n) == /\ Mode # "none" /\ rec[n] # {} - /\ IF Mode = "old" - THEN /\ rec' = [rec EXCEPT ![n] = idx[n]] - /\ lost' = (lost \/ ((rec[n] \ idx[n]) \cap truth # {})) - /\ resurrected' = (resurrected \/ ((idx[n] \ rec[n]) \ truth # {})) - /\ UNCHANGED idx - ELSE /\ idx' = [idx EXCEPT ![n] = @ \cup rec[n]] - /\ UNCHANGED <> - /\ UNCHANGED <> + /\ (Mode = "new" /\ Restrict) => Holds(n) + /\ rec' = [rec EXCEPT ![n] = NewRec(n)] + /\ idx' = [idx EXCEPT ![n] = NewIdx(n)] + /\ lost' = (lost \/ (Holds(n) /\ (rec[n] \ NewRec(n)) \cap truth # {})) + /\ resurrected' = (resurrected \/ (Holds(n) /\ (NewRec(n) \ rec[n]) \ truth # {})) + /\ staleGhost' = (staleGhost \/ (~Holds(n) /\ (NewIdx(n) \ idx[n]) \ truth # {})) + /\ holderGhost' = (holderGhost \/ (Holds(n) /\ (NewIdx(n) \ idx[n]) \ truth # {})) + /\ UNCHANGED <> Next == \/ \E t \in Topics : Sub(t) \/ Unsub(t) + \/ LoseOwnership \/ \E n \in Remote : Deliver(n) \/ \E n \in Remote : \E i \in 1..Len(ch[n]) : Drop(n, i) - \/ \E n \in Holders : Reconcile(n) + \/ \E n \in Cached : Reconcile(n) Spec == /\ Init /\ [][Next]_vars @@ -124,11 +174,16 @@ Spec == /\ \A n \in Holders : WF_vars(Reconcile(n)) ---------------------------------------------------------------------------- -\* SAFETY: reconcile never deletes a subscription the client still holds, -\* and never brings back one the client dropped. +\* SAFETY on the record of a current holder: reconcile never deletes a +\* subscription the client still holds, and never brings back a dropped one. InvNoReconcileLoss == ~lost InvNoResurrection == ~resurrected +\* SAFETY on the index: a node that no longer holds the partition never +\* recreates an entry for a topic the client dropped. +InvNoStaleCopyGhost == ~staleGhost +InvNoHolderGhost == ~holderGhost + \* LIVENESS: every node that holds the record ends up with every recorded \* topic in its index (routing on that node is repaired). HoldersIndexed == \A n \in Holders : rec[n] \subseteq idx[n] diff --git a/specs/ClusterSubReconcile_holderghost.cfg b/specs/ClusterSubReconcile_holderghost.cfg new file mode 100644 index 00000000..35428195 --- /dev/null +++ b/specs/ClusterSubReconcile_holderghost.cfg @@ -0,0 +1,9 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + MaxOps = 3 + Mode = "new" + Lossy = TRUE + Restrict = TRUE +INVARIANTS TypeOK InvNoHolderGhost +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile_nonholder.cfg b/specs/ClusterSubReconcile_nonholder.cfg index d3ba5dd3..12a2d5a9 100644 --- a/specs/ClusterSubReconcile_nonholder.cfg +++ b/specs/ClusterSubReconcile_nonholder.cfg @@ -1,9 +1,10 @@ SPECIFICATION Spec CONSTANTS Topics = {"a", "b"} - MaxOps = 3 + MaxOps = 2 Mode = "new" Lossy = TRUE + Restrict = TRUE INVARIANTS TypeOK PROPERTY NonHolderIndexed CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile_norepair.cfg b/specs/ClusterSubReconcile_norepair.cfg index f1158dc3..0c7e3d9e 100644 --- a/specs/ClusterSubReconcile_norepair.cfg +++ b/specs/ClusterSubReconcile_norepair.cfg @@ -1,9 +1,10 @@ SPECIFICATION Spec CONSTANTS Topics = {"a", "b"} - MaxOps = 3 + MaxOps = 2 Mode = "none" Lossy = TRUE + Restrict = TRUE INVARIANTS TypeOK PROPERTY RoutingRepaired CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile_oldloss.cfg b/specs/ClusterSubReconcile_oldloss.cfg index 05f9c7f7..f52b171a 100644 --- a/specs/ClusterSubReconcile_oldloss.cfg +++ b/specs/ClusterSubReconcile_oldloss.cfg @@ -4,6 +4,6 @@ CONSTANTS MaxOps = 3 Mode = "old" Lossy = TRUE + Restrict = TRUE INVARIANTS TypeOK InvNoReconcileLoss - CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile_oldresurrect.cfg b/specs/ClusterSubReconcile_oldresurrect.cfg index bcac1a7c..25ff44cd 100644 --- a/specs/ClusterSubReconcile_oldresurrect.cfg +++ b/specs/ClusterSubReconcile_oldresurrect.cfg @@ -4,6 +4,6 @@ CONSTANTS MaxOps = 3 Mode = "old" Lossy = TRUE + Restrict = TRUE INVARIANTS TypeOK InvNoResurrection - CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile_unrestricted.cfg b/specs/ClusterSubReconcile_unrestricted.cfg new file mode 100644 index 00000000..00e5c0e9 --- /dev/null +++ b/specs/ClusterSubReconcile_unrestricted.cfg @@ -0,0 +1,9 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + MaxOps = 3 + Mode = "new" + Lossy = TRUE + Restrict = FALSE +INVARIANTS TypeOK InvNoStaleCopyGhost +CHECK_DEADLOCK FALSE