Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -264,7 +264,7 @@ impl<T: ClusterTransport + 'static> BrokerEventHandler for ClusterEventHandler<T

let qos = qos_to_u8(sub.qos);
let is_wildcard = topic.contains('+') || topic.contains('#');
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 {
if is_wildcard {
Expand Down Expand Up @@ -370,7 +370,7 @@ impl<T: ClusterTransport + 'static> BrokerEventHandler for ClusterEventHandler<T
}
} 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 {
trace!(
topic,
Expand Down
2 changes: 1 addition & 1 deletion crates/mqdb-cluster/src/cluster/event_handler/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -147,7 +147,7 @@ async fn clear_client_subscriptions<T: ClusterTransport + 'static>(
}
} 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::<T>::broadcast_topic_subscription(ctrl, broadcast).await;
Expand Down
2 changes: 1 addition & 1 deletion crates/mqdb-cluster/src/cluster/event_handler/routing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -168,7 +168,7 @@ impl<T: ClusterTransport + 'static> ClusterEventHandler<T> {
.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 {
Expand Down
8 changes: 2 additions & 6 deletions crates/mqdb-cluster/src/cluster/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,6 @@ mod subscription_cache;
mod topic_index;
mod topic_trie;
mod transport;
mod wildcard_pending;
mod wildcard_store;
mod write_log;

Expand Down Expand Up @@ -123,18 +122,15 @@ 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,
};
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;

Expand Down
4 changes: 0 additions & 4 deletions crates/mqdb-cluster/src/cluster/store_manager/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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),
Expand Down Expand Up @@ -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())
Expand Down
11 changes: 6 additions & 5 deletions crates/mqdb-cluster/src/cluster/store_manager/recovery.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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;
}
Expand Down
33 changes: 33 additions & 0 deletions crates/mqdb-cluster/src/cluster/store_manager/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<dyn mqdb_core::StorageBackend> =
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)]
);
}
Loading
Loading