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/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 0b128535..56b03bdc 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 { @@ -200,74 +199,50 @@ impl SubscriptionCache { &self, topic_index: &TopicIndex, wildcard_store: &WildcardStore, + holds_partition: impl Fn(PartitionId) -> bool, ) -> 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() { + let client_id = snapshot.client_id_str(); + if !holds_partition(session_partition(client_id)) { + continue; + } result.clients_checked += 1; - 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,79 +689,192 @@ 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(); - 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.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 result = cache.reconcile(&topic_index, &wildcard_store, |_| true); - 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) + .add_subscription("client1", "app/resp/client1", 0) .unwrap(); - wildcard_store - .subscribe_mqtt("sensors/+/temp", "client1", 1) + 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); + let result = cache.reconcile(&topic_index, &wildcard_store, |partition| partition == held); assert_eq!(result.clients_checked, 1); - assert_eq!(result.subscriptions_removed, 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); + } - 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, |_| true) + .index_entries_restored, + 2 + ); + assert_eq!( + cache + .reconcile(&topic_index, &wildcard_store, |_| true) + .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, |_| true); + + 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..244b2edf 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; } @@ -233,14 +229,16 @@ impl ClusteredAgent { } if became_primary { let stores = ctrl.stores(); - let result = stores - .subscriptions - .reconcile(&stores.topics, &stores.wildcards); - if result.subscriptions_added > 0 || result.subscriptions_removed > 0 { + 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, - added = result.subscriptions_added, - removed = result.subscriptions_removed, + restored = result.index_entries_restored, "reconciled subscriptions after partition takeover" ); } @@ -595,50 +593,21 @@ 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; let stores = ctrl.stores(); if stores.subscriptions.needs_reconciliation(now) { - let result = stores - .subscriptions - .reconcile(&stores.topics, &stores.wildcards); - if result.subscriptions_added > 0 || result.subscriptions_removed > 0 { + 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, - 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..0eb476a9 100644 --- a/docs/distributed-design.md +++ b/docs/distributed-design.md @@ -156,26 +156,26 @@ 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 - `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 - `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 @@ -951,8 +951,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 +2241,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 +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 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: 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 new file mode 100644 index 00000000..94e10166 --- /dev/null +++ b/specs/ClusterSubReconcile.cfg @@ -0,0 +1,10 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + MaxOps = 2 + Mode = "new" + Lossy = TRUE + Restrict = TRUE +INVARIANTS TypeOK InvNoReconcileLoss InvNoResurrection InvNoStaleCopyGhost +PROPERTY RoutingRepaired +CHECK_DEADLOCK FALSE diff --git a/specs/ClusterSubReconcile.tla b/specs/ClusterSubReconcile.tla new file mode 100644 index 00000000..06bb147e --- /dev/null +++ b/specs/ClusterSubReconcile.tla @@ -0,0 +1,194 @@ +------------------------- 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). 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" = 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, 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 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, Restrict + +ASSUME Mode \in {"old", "new", "none"} +ASSUME MaxOps \in Nat /\ Lossy \in BOOLEAN /\ Restrict \in BOOLEAN + +Nodes == {"P", "R", "O", "S"} +Holders == {"P", "R"} +Cached == {"P", "R", "S"} +Remote == {"R", "O", "S"} + +Msg == [k : {"repl", "idx"}, op : {"add", "del"}, t : Topics] + +VARIABLES truth, rec, idx, ch, ops, sHolds, lost, resurrected, staleGhost, holderGhost + +vars == <> + +TypeOK == + /\ truth \subseteq 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 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) == + [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 + /\ 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) + +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 + /\ 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 <> + +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] + +Reconcile(n) == + /\ Mode # "none" + /\ rec[n] # {} + /\ (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 Cached : Reconcile(n) + +Spec == + /\ Init /\ [][Next]_vars + /\ \A n \in Remote : WF_vars(Deliver(n)) + /\ \A n \in Holders : WF_vars(Reconcile(n)) + +---------------------------------------------------------------------------- +\* 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] +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_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 new file mode 100644 index 00000000..12a2d5a9 --- /dev/null +++ b/specs/ClusterSubReconcile_nonholder.cfg @@ -0,0 +1,10 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + 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 new file mode 100644 index 00000000..0c7e3d9e --- /dev/null +++ b/specs/ClusterSubReconcile_norepair.cfg @@ -0,0 +1,10 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + 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 new file mode 100644 index 00000000..f52b171a --- /dev/null +++ b/specs/ClusterSubReconcile_oldloss.cfg @@ -0,0 +1,9 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + 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 new file mode 100644 index 00000000..25ff44cd --- /dev/null +++ b/specs/ClusterSubReconcile_oldresurrect.cfg @@ -0,0 +1,9 @@ +SPECIFICATION Spec +CONSTANTS + Topics = {"a", "b"} + 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