From 1be3f2da79cc61b00f3a2c297003810e20ce3b94 Mon Sep 17 00:00:00 2001 From: Mathias Koch Date: Fri, 19 Jun 2026 09:37:16 +0200 Subject: [PATCH] feat(mqtt/greengrass): release pooled subscriptions via linger eviction MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Pooled IPC subscription streams were never terminated, so every distinct topic ever subscribed stayed subscribed at the cloud broker for the life of the process. Each is a cloud-side MQTT subscription counting toward AWS IoT's per-connection limit; past 50 the nucleus opens a `#2` connection that a thing-policy-variable IoT policy rejects, leaving the device online-but-unmanageable (factbird-edge-applications#191, section C). Give each pooled slot a lifecycle instead of immortality: - On the last subscriber's unsubscribe/drop, the slot goes idle and arms a `linger` timer rather than terminating immediately. Reuse within the window bumps a generation counter so the pending timer no-ops and the live stream is reused — no cloud churn, so the resubscribe race never arises for the common request/response pattern. - A slot idle for the full `linger` is evicted: removed from the pool and its forwarder aborted, dropping the StreamOperation which sends TERMINATE_STREAM. - Because terminate is un-acked, a `settle` guard delays the next subscribe to a just-terminated topic so the stale teardown is processed by the core first. Defaults: linger 5s, settle 500ms, both overridable via with_linger/with_settle. Adds pooled_subscription_count()/pooled_topics() for observability. The pool->slot cycle is broken with a Weak handle; eviction and claim serialize on the pool write lock and eviction is ptr_eq-guarded against replacement slots. --- Cargo.toml | 2 +- src/mqtt/greengrass.rs | 324 +++++++++++++++++++++++++++++++++++++---- 2 files changed, 299 insertions(+), 27 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index 8154450..08e4263 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -181,7 +181,7 @@ shadows_multi = [ # MQTT client implementations (all optional, feature-gated) mqtt_mqttrust = ["dep:mqttrust"] mqtt_rumqttc = ["std", "dep:rumqttc", "tokio/time"] -mqtt_greengrass = ["std", "dep:greengrass-ipc-rust", "dep:bytes", "dep:futures"] +mqtt_greengrass = ["std", "dep:greengrass-ipc-rust", "dep:bytes", "dep:futures", "tokio/time"] [patch.crates-io] serde-json-core = { git = "https://github.com/rust-embedded-community/serde-json-core", rev = "a435f78" } diff --git a/src/mqtt/greengrass.rs b/src/mqtt/greengrass.rs index fc43c0a..d878db4 100644 --- a/src/mqtt/greengrass.rs +++ b/src/mqtt/greengrass.rs @@ -12,17 +12,36 @@ //! window is lost. Unlike normal MQTT, `TERMINATE_STREAM` has no protocol //! acknowledgment, so the client cannot wait for the unsubscribe to complete. //! -//! To avoid this race, we pool IPC subscription streams per topic inside -//! [`GreengrassClient`]. A single long-lived stream per topic is kept open, -//! and a lightweight channel is swapped on each logical subscribe / unsubscribe -//! from rustot's core. The IPC stream is never terminated during normal -//! operation, so the MQTT subscription at the Greengrass core stays alive and -//! no messages are lost between rustot subscribe cycles. +//! We pool IPC subscription streams per topic inside [`GreengrassClient`]. A +//! single stream per topic is kept open and a lightweight channel is swapped +//! on each logical subscribe / unsubscribe from rustot's core, so back-to-back +//! subscribe cycles on the same topic never churn the cloud subscription. +//! +//! ## Lifecycle: linger eviction + settle guard +//! +//! Pooled streams are not kept forever — every cloud-side MQTT subscription +//! counts toward the broker's per-connection subscription limit (AWS IoT opens +//! a second client connection past 50, which a thing-policy-variable IoT policy +//! then rejects). So when a topic's last logical subscriber unsubscribes or is +//! dropped, the slot goes idle and arms a `linger` timer instead of being torn +//! down immediately. If it is reused before the timer fires, the live stream is +//! reused and nothing is terminated — eliminating the resubscribe race for the +//! common request/response pattern. Only a slot that stays idle for the full +//! `linger` window is evicted: it is removed from the pool and its forwarder is +//! aborted, dropping the [`StreamOperation`] and sending `TERMINATE_STREAM`. +//! +//! Because terminate is un-acked, a fresh subscribe to a *just-terminated* +//! topic could still overlap the in-flight teardown. A `settle` guard closes +//! that residual window: the next subscribe to a topic terminated less than +//! `settle` ago waits out the remainder before re-establishing, so the stale +//! `TERMINATE_STREAM` is processed by the core first. use std::collections::HashMap; use std::pin::Pin; -use std::sync::{Arc, Mutex, RwLock}; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::sync::{Arc, Mutex, RwLock, Weak}; use std::task::Poll; +use std::time::{Duration, Instant}; use bytes::Bytes; use tokio::sync::mpsc; @@ -30,14 +49,92 @@ use tokio::sync::mpsc; use crate::mqtt::{MqttMessage, MqttSubscription, PublishOptions, QoS, ToPayload}; type TopicChannel = mpsc::UnboundedSender; +type PoolMap = Arc>>>; +type TerminatedMap = Arc>>; + +/// Default time an idle pooled subscription is kept alive before its IPC +/// stream is terminated (releasing the cloud-side MQTT subscription). +const DEFAULT_LINGER: Duration = Duration::from_secs(5); + +/// Default time to wait before re-subscribing to a topic whose stream was just +/// terminated, so the un-acked `TERMINATE_STREAM` is processed by the core +/// before the new subscription is established. +const DEFAULT_SETTLE: Duration = Duration::from_millis(500); /// Per-topic pooled state. Shared between the forwarding task and every /// logical subscriber. struct PooledSlot { + /// Topic this slot subscribes to; also its key in the pool map. + topic: String, /// Sender for the currently active logical subscriber. `None` when the /// topic was previously subscribed and has since been unsubscribed. /// Accessed synchronously (short critical sections only). sender: Mutex>, + /// Bumped on every claim (subscribe). A pending eviction timer captures + /// the generation when armed and no-ops if it has since changed — i.e. + /// the slot was reused before the linger elapsed. + generation: AtomicU64, + /// Abort handle for the forwarder task. Aborting drops the IPC + /// `StreamOperation`, whose `Drop` sends `TERMINATE_STREAM`. + forwarder: Mutex>, + /// Weak handle to the shared pool so the slot can evict itself when its + /// linger expires. Weak to avoid a pool → slot → pool reference cycle. + pool: Weak>>>, + /// Shared map of recently-terminated topics, written on eviction and read + /// by the settle guard on the next subscribe. + terminated_at: TerminatedMap, + /// Linger duration captured from the client at construction. + linger: Duration, +} + +impl PooledSlot { + /// Arm a linger timer after the slot goes idle. When it fires, the slot is + /// evicted only if it has not been reused (generation unchanged) and is + /// still idle. + fn arm_eviction(self: &Arc) { + let armed_gen = self.generation.load(Ordering::SeqCst); + let slot = self.clone(); + tokio::spawn(async move { + tokio::time::sleep(slot.linger).await; + slot.try_evict(armed_gen); + }); + } + + /// Evict the slot if it was not reused since `armed_gen` and is still idle. + /// Removal from the pool and the generation check happen under the pool + /// write lock, which also serializes against [`GreengrassClient::ensure_slot`] + /// claiming the slot — so a slot is never handed out and evicted at once. + fn try_evict(self: &Arc, armed_gen: u64) { + let Some(pool) = self.pool.upgrade() else { + return; + }; + let mut pool = pool.write().unwrap(); + + // Reused since the timer was armed, or re-claimed without bumping the + // generation yet — either way, leave it alone. + if self.generation.load(Ordering::SeqCst) != armed_gen { + return; + } + if self.sender.lock().unwrap().is_some() { + return; + } + // Only evict if we are still the pooled slot for this topic. + match pool.get(&self.topic) { + Some(existing) if Arc::ptr_eq(existing, self) => {} + _ => return, + } + pool.remove(&self.topic); + + // Record the terminate time for the settle guard, then abort the + // forwarder — dropping the StreamOperation sends TERMINATE_STREAM. + self.terminated_at + .lock() + .unwrap() + .insert(self.topic.clone(), Instant::now()); + if let Some(handle) = self.forwarder.lock().unwrap().take() { + handle.abort(); + } + } } /// Wrapper around greengrass-ipc-rust client. @@ -50,7 +147,15 @@ pub struct GreengrassClient { client_id: String, /// Pool of active IPC subscription streams keyed by topic. /// Uses std RwLock — no lock is held across await points. - pool: Arc>>>, + pool: PoolMap, + /// Topics whose stream was terminated recently, with the terminate time. + /// Read by the settle guard before re-subscribing; pruned opportunistically. + terminated_at: TerminatedMap, + /// How long an idle pooled subscription is kept before its stream is + /// terminated. See module docs. + linger: Duration, + /// Delay applied before re-subscribing to a just-terminated topic. + settle: Duration, } impl GreengrassClient { @@ -86,9 +191,39 @@ impl GreengrassClient { client, client_id, pool: Arc::new(RwLock::new(HashMap::new())), + terminated_at: Arc::new(Mutex::new(HashMap::new())), + linger: DEFAULT_LINGER, + settle: DEFAULT_SETTLE, } } + /// Override how long an idle pooled subscription is kept before its stream + /// is terminated (default [`DEFAULT_LINGER`]). A longer linger reduces + /// churn; a shorter one reclaims subscriptions faster to stay under the + /// broker's per-connection limit. + pub fn with_linger(mut self, linger: Duration) -> Self { + self.linger = linger; + self + } + + /// Override the settle delay applied before re-subscribing to a topic whose + /// stream was just terminated (default [`DEFAULT_SETTLE`]). + pub fn with_settle(mut self, settle: Duration) -> Self { + self.settle = settle; + self + } + + /// Number of IPC subscription streams currently pooled — one per live + /// cloud-side MQTT subscription. Useful for staying under broker limits. + pub fn pooled_subscription_count(&self) -> usize { + self.pool.read().unwrap().len() + } + + /// Topics currently held in the subscription pool. + pub fn pooled_topics(&self) -> Vec { + self.pool.read().unwrap().keys().cloned().collect() + } + /// Ensure a long-lived IPC subscription exists for `topic` and return its /// pooled slot. Spawns a forwarding task on first subscribe. /// @@ -98,9 +233,30 @@ impl GreengrassClient { topic: &str, qos: QoS, ) -> Result, greengrass_ipc_rust::Error> { - // Fast path: already pooled. - if let Some(slot) = self.pool.read().unwrap().get(topic).cloned() { - return Ok(slot); + // Fast path: already pooled. Claim it (bump generation) under the write + // lock so any pending eviction timer for this slot no-ops. + { + let pool = self.pool.write().unwrap(); + if let Some(slot) = pool.get(topic) { + slot.generation.fetch_add(1, Ordering::SeqCst); + return Ok(slot.clone()); + } + } + + // Settle guard: if this topic's stream was terminated very recently, + // wait out the remainder so the stale (un-acked) TERMINATE_STREAM is + // processed by the core before we establish a new subscription. + let settle_wait = { + let mut term = self.terminated_at.lock().unwrap(); + let now = Instant::now(); + term.retain(|_, t| now.duration_since(*t) < self.settle); + term.get(topic) + .map(|t| self.settle.saturating_sub(now.duration_since(*t))) + }; + if let Some(wait) = settle_wait + && !wait.is_zero() + { + tokio::time::sleep(wait).await; } // Slow path: subscribe without holding the pool lock. @@ -113,29 +269,37 @@ impl GreengrassClient { .await?; let slot = Arc::new(PooledSlot { + topic: topic.to_string(), sender: Mutex::new(None), + generation: AtomicU64::new(1), + forwarder: Mutex::new(None), + pool: Arc::downgrade(&self.pool), + terminated_at: self.terminated_at.clone(), + linger: self.linger, }); // Insert into pool. Recheck in case of race — another task might // have subscribed to the same topic while we were awaiting. { let mut pool = self.pool.write().unwrap(); - if let Some(existing) = pool.get(topic).cloned() { + if let Some(existing) = pool.get(topic) { + existing.generation.fetch_add(1, Ordering::SeqCst); + let existing = existing.clone(); drop(stream); return Ok(existing); } pool.insert(topic.to_string(), slot.clone()); } - // Spawn the forwarder. It reads from the IPC stream forever and - // forwards each message to whichever subscriber's sender is - // currently installed in the slot. When the stream closes (e.g. - // IPC connection lost) the task removes the entry so a subsequent - // subscribe reopens the stream. + // Spawn the forwarder. It reads from the IPC stream and forwards each + // message to whichever subscriber's sender is currently installed in + // the slot. When the stream closes (e.g. IPC connection lost, or the + // slot is evicted via `AbortHandle`) the task removes the entry — but + // only if it is still the pooled slot — so a later subscribe reopens it. let slot_for_task = slot.clone(); let pool_for_task = self.pool.clone(); let topic_for_task = topic.to_string(); - tokio::spawn(async move { + let handle = tokio::spawn(async move { use futures::StreamExt; let mut stream = stream; while let Some(msg) = stream.next().await { @@ -144,8 +308,14 @@ impl GreengrassClient { let _ = tx.send(msg.message); } } - pool_for_task.write().unwrap().remove(&topic_for_task); + let mut pool = pool_for_task.write().unwrap(); + if let Some(existing) = pool.get(&topic_for_task) + && Arc::ptr_eq(existing, &slot_for_task) + { + pool.remove(&topic_for_task); + } }); + *slot.forwarder.lock().unwrap() = Some(handle.abort_handle()); Ok(slot) } @@ -234,16 +404,18 @@ struct SubscriptionReceiver { /// /// Holds one receiver per subscribed topic. On drop or explicit /// [`close()`](Self::close) / [`unsubscribe()`](MqttSubscription::unsubscribe), -/// the receivers are dropped and the pooled slots' senders are cleared. -/// The underlying IPC streams stay open and reusable for future subscribes. +/// the receivers are dropped and the pooled slots' senders are cleared. Each +/// now-idle slot arms a `linger` timer; if nothing resubscribes within the +/// window the IPC stream is terminated and the cloud subscription released, +/// otherwise the live stream is reused with no cloud churn. pub struct GreengrassSubscription { receivers: Vec, } impl GreengrassSubscription { /// Release this subscription. Clears the pooled slots' senders so that - /// future messages are discarded until someone subscribes again. Does - /// NOT terminate the underlying IPC streams — they stay in the pool. + /// future messages are discarded until someone subscribes again, and arms + /// each slot's linger timer (see [`GreengrassSubscription`]). pub async fn close(self) -> Result<(), greengrass_ipc_rust::Error> { // Delegated to Drop. Ok(()) @@ -293,10 +465,13 @@ impl Drop for GreengrassSubscription { // whether two senders reference the same underlying channel. for receiver in &self.receivers { let mut guard = receiver.slot.sender.lock().unwrap(); - if let Some(ref current) = *guard - && current.same_channel(&receiver.tx) - { + let is_ours = matches!(&*guard, Some(current) if current.same_channel(&receiver.tx)); + if is_ours { *guard = None; + drop(guard); + // Slot is now idle. Arm the linger timer; if no one + // resubscribes within `linger`, the stream is terminated. + receiver.slot.arm_eviction(); } } } @@ -373,3 +548,100 @@ pub use greengrass_ipc_rust::{ Error as GreengrassError, GreengrassCoreIPCClient, IoTCoreMessage, PublishToIoTCoreRequest, StreamOperation, SubscribeToIoTCoreRequest, }; + +#[cfg(test)] +mod tests { + use super::*; + + /// A forwarder stand-in: a parked task whose abort handle we can store in + /// the slot, mirroring the real forwarder that holds the IPC stream. + fn dummy_forwarder() -> tokio::task::AbortHandle { + tokio::spawn(futures::future::pending::<()>()).abort_handle() + } + + fn make_slot( + topic: &str, + pool: &PoolMap, + terminated_at: &TerminatedMap, + linger: Duration, + ) -> Arc { + let slot = Arc::new(PooledSlot { + topic: topic.to_string(), + sender: Mutex::new(None), + generation: AtomicU64::new(1), + forwarder: Mutex::new(Some(dummy_forwarder())), + pool: Arc::downgrade(pool), + terminated_at: terminated_at.clone(), + linger, + }); + pool.write() + .unwrap() + .insert(topic.to_string(), slot.clone()); + slot + } + + #[tokio::test] + async fn idle_slot_is_evicted_after_linger() { + let pool: PoolMap = Arc::new(RwLock::new(HashMap::new())); + let terminated: TerminatedMap = Arc::new(Mutex::new(HashMap::new())); + let slot = make_slot("t/a", &pool, &terminated, Duration::from_millis(20)); + + slot.arm_eviction(); + assert!(pool.read().unwrap().contains_key("t/a")); + + tokio::time::sleep(Duration::from_millis(60)).await; + + // Evicted: removed from the pool and recorded for the settle guard. + assert!(!pool.read().unwrap().contains_key("t/a")); + assert!(terminated.lock().unwrap().contains_key("t/a")); + } + + #[tokio::test] + async fn reuse_before_linger_cancels_eviction() { + let pool: PoolMap = Arc::new(RwLock::new(HashMap::new())); + let terminated: TerminatedMap = Arc::new(Mutex::new(HashMap::new())); + let slot = make_slot("t/b", &pool, &terminated, Duration::from_millis(40)); + + slot.arm_eviction(); + // Simulate a resubscribe within the linger window: a claim bumps the + // generation, so the pending timer must no-op. + slot.generation.fetch_add(1, Ordering::SeqCst); + + tokio::time::sleep(Duration::from_millis(80)).await; + + assert!(pool.read().unwrap().contains_key("t/b")); + assert!(terminated.lock().unwrap().is_empty()); + } + + #[tokio::test] + async fn active_slot_is_not_evicted() { + let pool: PoolMap = Arc::new(RwLock::new(HashMap::new())); + let terminated: TerminatedMap = Arc::new(Mutex::new(HashMap::new())); + let slot = make_slot("t/c", &pool, &terminated, Duration::from_millis(20)); + + // A subscriber is active (sender installed) — eviction must not fire. + let (tx, _rx) = mpsc::unbounded_channel(); + *slot.sender.lock().unwrap() = Some(tx); + + slot.arm_eviction(); + tokio::time::sleep(Duration::from_millis(60)).await; + + assert!(pool.read().unwrap().contains_key("t/c")); + } + + #[tokio::test] + async fn evicted_slot_does_not_remove_a_replacement() { + let pool: PoolMap = Arc::new(RwLock::new(HashMap::new())); + let terminated: TerminatedMap = Arc::new(Mutex::new(HashMap::new())); + let old = make_slot("t/d", &pool, &terminated, Duration::from_millis(20)); + + // A fresh slot replaces the old one in the pool (same topic). + let new = make_slot("t/d", &pool, &terminated, Duration::from_millis(20)); + assert!(Arc::ptr_eq(pool.read().unwrap().get("t/d").unwrap(), &new)); + + // The old slot's timer firing must not evict the replacement. + old.try_evict(old.generation.load(Ordering::SeqCst)); + + assert!(Arc::ptr_eq(pool.read().unwrap().get("t/d").unwrap(), &new)); + } +}