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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,17 @@ All notable changes to this project will be documented in this file.
The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/),
and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html).

## [mqtt5 0.43.1] - 2026-09-27

### Fixed

- **A clean-start reconnect no longer receives a message routed to the session it replaced.** A publish routed to a client just before that client reconnected with `clean_start=1` could land in the new session's queue after the queue had been cleared, delivering a message for a subscription the new session never made (#150). Each client queue now counts its full clears, and a message routed before the latest one is dropped. A resumed session still receives it.
- **The file backend no longer writes a queued message to disk after it was delivered or cleared.** A message delivered or cleared at the same moment it was queued could have its delete reach the storage writer before its write, leaving the file on disk; after a restart it was loaded again and redelivered. Writes are now handed to the storage writer before the entry becomes visible to delivery or clearing.

### Added

- `ClientQueue::epoch` and `ClientQueue::push_in_epoch`.

## [mqtt5 0.43.0] - 2026-09-27

### Breaking
Expand Down
2 changes: 1 addition & 1 deletion crates/mqtt5/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[package]
name = "mqtt5"
version = "0.43.0"
version = "0.43.1"
edition.workspace = true
rust-version.workspace = true
authors.workspace = true
Expand Down
154 changes: 131 additions & 23 deletions crates/mqtt5/src/broker/router.rs
Original file line number Diff line number Diff line change
Expand Up @@ -326,12 +326,14 @@ enum DeliveryPlan {
target_flow: Option<u64>,
lanes: DeliveryLanes,
queue: QueueHandle,
epoch: u64,
},
Behind {
client_id: String,
message: PublishPacket,
target_flow: Option<u64>,
queue: QueueHandle,
epoch: u64,
},
}

Expand Down Expand Up @@ -1418,12 +1420,7 @@ impl MessageRouter {
if self.storage.is_some() && sub.qos != QoS::AtMostOnce {
let mut message = publish.clone();
message.qos = sub.qos;
plans.push(DeliveryPlan::Behind {
client_id: sub.client_id.clone(),
message,
target_flow: sub.flow_id,
queue: self.queue_handle(&sub.client_id),
});
plans.push(self.plan_behind(sub, message));
}
}
}
Expand Down Expand Up @@ -1478,15 +1475,21 @@ impl MessageRouter {

fn queue_behind(
queue: &QueueHandle,
epoch: u64,
message: PublishPacket,
client_id: &str,
target_flow: Option<u64>,
) {
let qos = message.qos;
let outcome = queue.push(
QueuedMessage::new(message, client_id.to_string(), qos, None)
.with_target_flow(target_flow),
);
let queued = QueuedMessage::new(message, client_id.to_string(), qos, None)
.with_target_flow(target_flow);
let Some(outcome) = queue.push_in_epoch(queued, epoch) else {
debug!(
client_id,
"Dropped message routed to a session that has since been discarded"
);
return;
};
queue.notify();
trace!(
client_id,
Expand All @@ -1507,13 +1510,15 @@ impl MessageRouter {
message,
target_flow,
queue,
} => Self::queue_behind(&queue, message, &client_id, target_flow),
epoch,
} => Self::queue_behind(&queue, epoch, message, &client_id, target_flow),
DeliveryPlan::Online {
client_id,
message,
target_flow,
lanes,
queue,
epoch,
} => {
let routable = RoutableMessage {
publish: message,
Expand All @@ -1526,14 +1531,21 @@ impl MessageRouter {
return;
}
if queue.behind() {
Self::queue_behind(&queue, routable.publish, &client_id, routable.target_flow);
Self::queue_behind(
&queue,
epoch,
routable.publish,
&client_id,
routable.target_flow,
);
return;
}
let routable = match lanes.qos1_tx.try_send(routable) {
Ok(()) => return,
Err(mpsc::error::TrySendError::Closed(rejected)) => {
Self::queue_behind(
&queue,
epoch,
rejected.publish,
&client_id,
rejected.target_flow,
Expand All @@ -1543,7 +1555,13 @@ impl MessageRouter {
Err(mpsc::error::TrySendError::Full(rejected)) => rejected,
};
if publishing_client_id == Some(client_id.as_str()) {
Self::queue_behind(&queue, routable.publish, &client_id, routable.target_flow);
Self::queue_behind(
&queue,
epoch,
routable.publish,
&client_id,
routable.target_flow,
);
return;
}
let permit = match deadline {
Expand All @@ -1557,6 +1575,7 @@ impl MessageRouter {
Some(permit) => permit.send(routable),
None => Self::queue_behind(
&queue,
epoch,
routable.publish,
&client_id,
routable.target_flow,
Expand Down Expand Up @@ -1621,12 +1640,7 @@ impl MessageRouter {
if qos == QoS::AtMostOnce || self.storage.is_none() {
return None;
}
return Some(DeliveryPlan::Behind {
client_id: sub.client_id.clone(),
message: Self::prepare_message(publish, sub, qos),
target_flow: sub.flow_id,
queue: self.queue_handle(&sub.client_id),
});
return Some(self.plan_behind(sub, Self::prepare_message(publish, sub, qos)));
}
}

Expand All @@ -1649,6 +1663,7 @@ impl MessageRouter {
qos0_tx: client_info.qos0_tx.clone(),
},
queue: Arc::clone(&client_info.queue),
epoch: client_info.queue.epoch(),
});
}

Expand All @@ -1662,12 +1677,18 @@ impl MessageRouter {
);
return None;
}
Some(DeliveryPlan::Behind {
Some(self.plan_behind(sub, Self::prepare_message(publish, sub, sub.qos)))
}

fn plan_behind(&self, sub: &Subscription, message: PublishPacket) -> DeliveryPlan {
let queue = self.queue_handle(&sub.client_id);
DeliveryPlan::Behind {
client_id: sub.client_id.clone(),
message: Self::prepare_message(publish, sub, sub.qos),
message,
target_flow: sub.flow_id,
queue: self.queue_handle(&sub.client_id),
})
epoch: queue.epoch(),
queue,
}
}

pub async fn get_retained_messages(&self, topic_filter: &str) -> Vec<PublishPacket> {
Expand Down Expand Up @@ -2640,6 +2661,93 @@ mod tests {
.generation
}

async fn stall_publish_behind_full_lane(
router: &Arc<MessageRouter>,
client_id: &str,
lanes: &TestLanes,
) -> tokio::task::JoinHandle<()> {
register(router, client_id, lanes).await;
router
.subscribe(SubscriptionRequest::new(client_id, "t", QoS::AtLeastOnce))
.await
.unwrap();
let first = PublishPacket::new("t", &b"first"[..], QoS::AtLeastOnce);
router.route_message(&first, None).await;

let stalled = Arc::clone(router);
let handle = tokio::spawn(async move {
let second = PublishPacket::new("t", &b"second"[..], QoS::AtLeastOnce);
stalled
.route_message_with_deadline(
&second,
None,
Instant::now() + Duration::from_secs(30),
)
.await;
});
for _ in 0..10 {
tokio::task::yield_now().await;
}
assert!(
!handle.is_finished(),
"the publish must wait on the full lane"
);
handle
}

#[tokio::test]
async fn clean_start_drops_a_publish_planned_for_the_previous_session() {
let router = Arc::new(MessageRouter::new());
let previous = TestLanes::new(1);
let stalled = stall_publish_behind_full_lane(&router, "c", &previous).await;

let current = TestLanes::new(1);
let queue = router.queue_handle("c");
let (dtx, _drx) = tokio::sync::oneshot::channel();
router
.register_session(
"c".to_string(),
current.lanes(),
Arc::clone(&queue),
dtx,
true,
)
.await;
queue.clear(None);
drop(previous);
stalled.await.unwrap();

assert_eq!(
queue.count(),
0,
"a publish routed under the previous session must not reach the clean session"
);
}

#[tokio::test]
async fn resumed_session_keeps_a_publish_planned_for_the_previous_connection() {
let router = Arc::new(MessageRouter::new());
let previous = TestLanes::new(1);
let stalled = stall_publish_behind_full_lane(&router, "c", &previous).await;

let current = TestLanes::new(1);
let queue = router.queue_handle("c");
let (dtx, _drx) = tokio::sync::oneshot::channel();
router
.register_session(
"c".to_string(),
current.lanes(),
Arc::clone(&queue),
dtx,
false,
)
.await;
drop(previous);
stalled.await.unwrap();

assert_eq!(queue.count(), 1, "a resumed session must still receive it");
}

#[tokio::test]
async fn armed_will_is_claimed_exactly_once() {
let router = MessageRouter::new();
Expand Down
Loading
Loading