commit 3d2cd67529e8c074febac78e2620dd792cdc0a3a
parent b86468be9147bef986e2f18e08793a97d931933a
Author: triesap <tyson@radroots.org>
Date: Fri, 2 Oct 2026 05:19:15 +0000
runtime: fix silent final-update loss in the snapshot stream
- Retain one latest change per subscriber before nonblocking enqueue
- Preserve queued initial order and monotonic predecessor metadata
- Cover saturation, independent consumers and cancellation with seven regressions
- Verify application tests, workspace checks, Clippy and foundation guards
Diffstat:
1 file changed, 277 insertions(+), 15 deletions(-)
diff --git a/core/crates/harvestcircle_application/src/change_stream.rs b/core/crates/harvestcircle_application/src/change_stream.rs
@@ -1,7 +1,7 @@
use std::collections::BTreeMap;
use std::num::{NonZeroU64, NonZeroUsize};
-use tokio::sync::mpsc;
+use tokio::sync::{mpsc, watch};
use crate::{AppSnapshot, SnapshotRevision};
@@ -51,18 +51,60 @@ impl SnapshotChange {
pub struct SnapshotChangeReceiver {
receiver: mpsc::Receiver<SnapshotChange>,
+ latest: watch::Receiver<SnapshotChange>,
+ last_delivered_revision: Option<SnapshotRevision>,
}
impl SnapshotChangeReceiver {
pub async fn receive(&mut self) -> Option<SnapshotChange> {
- self.receiver.recv().await
+ loop {
+ let queue_open = match self.receiver.try_recv() {
+ Ok(change) => {
+ if let Some(change) = self.deliver_newer(change) {
+ return Some(change);
+ }
+ continue;
+ }
+ Err(mpsc::error::TryRecvError::Empty) => true,
+ Err(mpsc::error::TryRecvError::Disconnected) => false,
+ };
+ let retained = self.latest.borrow().clone();
+ if let Some(change) = self.deliver_newer(retained) {
+ return Some(change);
+ }
+ if !queue_open {
+ return None;
+ }
+ if let Some(change) = self.receiver.recv().await
+ && let Some(change) = self.deliver_newer(change)
+ {
+ return Some(change);
+ }
+ }
}
+
+ fn deliver_newer(&mut self, change: SnapshotChange) -> Option<SnapshotChange> {
+ let revision = change.revision();
+ if self
+ .last_delivered_revision
+ .is_some_and(|delivered| revision <= delivered)
+ {
+ return None;
+ }
+ self.last_delivered_revision = Some(revision);
+ Some(change)
+ }
+}
+
+struct SnapshotChangeSubscriber {
+ sender: mpsc::Sender<SnapshotChange>,
+ latest: watch::Sender<SnapshotChange>,
}
pub struct OrderedSnapshotChanges {
latest: AppSnapshot,
next_subscription: u64,
- subscribers: BTreeMap<ChangeSubscriptionId, mpsc::Sender<SnapshotChange>>,
+ subscribers: BTreeMap<ChangeSubscriptionId, SnapshotChangeSubscriber>,
closed: bool,
}
@@ -96,15 +138,23 @@ impl OrderedSnapshotChanges {
}
let id = ChangeSubscriptionId(NonZeroU64::new(self.next_subscription)?);
self.next_subscription = self.next_subscription.checked_add(1)?;
+ let initial = SnapshotChange {
+ snapshot: self.latest.clone(),
+ previous_revision: None,
+ };
let (sender, receiver) = mpsc::channel(capacity.get());
- sender
- .try_send(SnapshotChange {
- snapshot: self.latest.clone(),
- previous_revision: None,
- })
- .ok()?;
- self.subscribers.insert(id, sender);
- Some((id, SnapshotChangeReceiver { receiver }))
+ let (latest, retained) = watch::channel(initial.clone());
+ sender.try_send(initial).ok()?;
+ self.subscribers
+ .insert(id, SnapshotChangeSubscriber { sender, latest });
+ Some((
+ id,
+ SnapshotChangeReceiver {
+ receiver,
+ latest: retained,
+ last_delivered_revision: None,
+ },
+ ))
}
#[must_use]
@@ -121,11 +171,14 @@ impl OrderedSnapshotChanges {
snapshot,
};
self.latest = change.snapshot.clone();
- self.subscribers
- .retain(|_, sender| match sender.try_send(change.clone()) {
+ self.subscribers.retain(|_, subscriber| {
+ // Retain before enqueue so saturation and sender closure preserve the tail.
+ drop(subscriber.latest.send_replace(change.clone()));
+ match subscriber.sender.try_send(change.clone()) {
Ok(()) | Err(mpsc::error::TrySendError::Full(_)) => true,
Err(mpsc::error::TrySendError::Closed(_)) => false,
- });
+ }
+ });
}
pub fn close(&mut self) {
@@ -136,12 +189,18 @@ impl OrderedSnapshotChanges {
#[cfg(test)]
mod tests {
+ use std::future::Future;
use std::num::NonZeroUsize;
+ use std::task::Poll;
+ use std::time::Duration;
use crate::{
- AppSnapshot, OrderedSnapshotChanges, RelayConfiguration, SessionState, SnapshotRevision,
+ AppSnapshot, OrderedSnapshotChanges, RelayConfiguration, SessionState, SnapshotChange,
+ SnapshotChangeReceiver, SnapshotRevision,
};
+ const RECEIVE_TIMEOUT: Duration = Duration::from_secs(1);
+
#[tokio::test]
async fn change_stream_publishes_monotonic_revisions_to_multiple_consumers() {
let mut changes = OrderedSnapshotChanges::new(snapshot(0));
@@ -216,6 +275,209 @@ mod tests {
);
}
+ #[tokio::test]
+ async fn capacity_one_recovers_final_tail_without_later_publication() {
+ let mut changes = OrderedSnapshotChanges::new(snapshot(0));
+ let (_, mut receiver) = changes
+ .subscribe(NonZeroUsize::new(1).expect("capacity"))
+ .expect("subscription");
+ let initial = next_change(&mut receiver, "initial").await;
+ assert_eq!(initial.revision(), revision(0));
+
+ changes.publish(snapshot(1));
+ changes.publish(snapshot(2));
+ let queued = next_change(&mut receiver, "queued revision 1").await;
+ assert_eq!(queued.revision(), revision(1));
+ assert_eq!(queued.previous_revision(), Some(revision(0)));
+ let tail = next_change(&mut receiver, "final revision 2 without another publish").await;
+ assert_eq!(tail.revision(), revision(2));
+ assert_eq!(tail.snapshot(), &snapshot(2));
+ assert_eq!(tail.previous_revision(), Some(revision(1)));
+ assert!(!tail.recovers_gap_after(queued.revision()));
+ assert_eq!(changes.last_revision(), revision(2));
+ }
+
+ #[tokio::test]
+ async fn sustained_overflow_coalesces_latest_tail_with_original_gap_metadata() {
+ let mut changes = OrderedSnapshotChanges::new(snapshot(0));
+ let (_, mut receiver) = changes
+ .subscribe(NonZeroUsize::new(1).expect("capacity"))
+ .expect("subscription");
+ next_change(&mut receiver, "initial").await;
+
+ for value in 1..=128 {
+ changes.publish(snapshot(value));
+ }
+ assert_eq!(changes.last_revision(), revision(128));
+ let queued = next_change(&mut receiver, "queued revision 1").await;
+ assert_eq!(queued.revision(), revision(1));
+ let latest = next_change(&mut receiver, "coalesced revision 128").await;
+ assert_eq!(latest.revision(), revision(128));
+ assert_eq!(latest.snapshot(), &snapshot(128));
+ assert_eq!(latest.previous_revision(), Some(revision(127)));
+ assert!(latest.recovers_gap_after(queued.revision()));
+
+ changes.publish(snapshot(128));
+ changes.publish(snapshot(127));
+ changes.close();
+ assert!(next_delivery(&mut receiver).await.is_none());
+ }
+
+ #[tokio::test]
+ async fn queued_initial_snapshot_precedes_coalesced_latest_tail() {
+ let mut changes = OrderedSnapshotChanges::new(snapshot(0));
+ let (_, mut receiver) = changes
+ .subscribe(NonZeroUsize::new(1).expect("capacity"))
+ .expect("subscription");
+ changes.publish(snapshot(1));
+ changes.publish(snapshot(2));
+
+ let initial = next_change(&mut receiver, "queued initial snapshot").await;
+ assert_eq!(initial.revision(), revision(0));
+ assert_eq!(initial.previous_revision(), None);
+ let latest = next_change(&mut receiver, "latest after queued initial").await;
+ assert_eq!(latest.revision(), revision(2));
+ assert_eq!(latest.previous_revision(), Some(revision(1)));
+ assert!(latest.recovers_gap_after(initial.revision()));
+ changes.close();
+ assert!(next_delivery(&mut receiver).await.is_none());
+ }
+
+ #[tokio::test]
+ async fn independently_paced_consumers_recover_their_own_latest_tail() {
+ let mut changes = OrderedSnapshotChanges::new(snapshot(0));
+ let (_, mut slow) = changes
+ .subscribe(NonZeroUsize::new(1).expect("capacity"))
+ .expect("slow subscription");
+ let (_, mut fast) = changes
+ .subscribe(NonZeroUsize::new(1).expect("capacity"))
+ .expect("fast subscription");
+ next_change(&mut slow, "slow initial").await;
+ next_change(&mut fast, "fast initial").await;
+
+ for value in 1..=3 {
+ changes.publish(snapshot(value));
+ let delivered = next_change(&mut fast, "fast consumer revision").await;
+ assert_eq!(delivered.revision(), revision(value));
+ assert_eq!(delivered.previous_revision(), Some(revision(value - 1)));
+ assert!(!delivered.recovers_gap_after(revision(value - 1)));
+ }
+ let queued = next_change(&mut slow, "slow queued revision 1").await;
+ assert_eq!(queued.revision(), revision(1));
+ let tail = next_change(&mut slow, "slow retained revision 3").await;
+ assert_eq!(tail.revision(), revision(3));
+ assert_eq!(tail.previous_revision(), Some(revision(2)));
+ assert!(tail.recovers_gap_after(queued.revision()));
+ changes.close();
+ assert!(next_delivery(&mut slow).await.is_none());
+ assert!(next_delivery(&mut fast).await.is_none());
+ }
+
+ #[tokio::test]
+ async fn full_close_drains_queued_change_and_final_tail_before_none() {
+ let mut changes = OrderedSnapshotChanges::new(snapshot(0));
+ let (_, mut receiver) = changes
+ .subscribe(NonZeroUsize::new(1).expect("capacity"))
+ .expect("subscription");
+ next_change(&mut receiver, "initial").await;
+ changes.publish(snapshot(1));
+ changes.publish(snapshot(2));
+ changes.close();
+ changes.publish(snapshot(3));
+
+ assert_eq!(
+ next_change(&mut receiver, "queued revision before close")
+ .await
+ .revision(),
+ revision(1)
+ );
+ let tail = next_change(&mut receiver, "retained final revision before close").await;
+ assert_eq!(tail.revision(), revision(2));
+ assert_eq!(tail.previous_revision(), Some(revision(1)));
+ assert!(next_delivery(&mut receiver).await.is_none());
+ assert!(next_delivery(&mut receiver).await.is_none());
+ assert_eq!(changes.last_revision(), revision(2));
+ }
+
+ #[tokio::test]
+ async fn unsubscribe_drains_retained_tail_without_future_updates() {
+ let mut changes = OrderedSnapshotChanges::new(snapshot(0));
+ let (id, mut receiver) = changes
+ .subscribe(NonZeroUsize::new(1).expect("capacity"))
+ .expect("subscription");
+ next_change(&mut receiver, "initial").await;
+ changes.publish(snapshot(1));
+ changes.publish(snapshot(2));
+ assert!(changes.unsubscribe(id));
+ assert!(!changes.unsubscribe(id));
+ changes.publish(snapshot(3));
+
+ assert_eq!(
+ next_change(&mut receiver, "queued revision before unsubscribe")
+ .await
+ .revision(),
+ revision(1)
+ );
+ let tail = next_change(&mut receiver, "retained revision before unsubscribe").await;
+ assert_eq!(tail.revision(), revision(2));
+ assert_eq!(tail.previous_revision(), Some(revision(1)));
+ assert!(next_delivery(&mut receiver).await.is_none());
+ assert_eq!(changes.last_revision(), revision(3));
+ }
+
+ #[tokio::test]
+ async fn cancelled_pending_receive_preserves_future_queued_and_retained_tail() {
+ let mut changes = OrderedSnapshotChanges::new(snapshot(0));
+ let (_, mut receiver) = changes
+ .subscribe(NonZeroUsize::new(1).expect("capacity"))
+ .expect("subscription");
+ next_change(&mut receiver, "initial").await;
+ {
+ let pending = receiver.receive();
+ tokio::pin!(pending);
+ tokio::time::timeout(
+ RECEIVE_TIMEOUT,
+ std::future::poll_fn(|context| {
+ assert!(pending.as_mut().poll(context).is_pending());
+ Poll::Ready(())
+ }),
+ )
+ .await
+ .expect("empty receive was polled before cancellation");
+ }
+
+ changes.publish(snapshot(1));
+ changes.publish(snapshot(2));
+ assert_eq!(
+ next_change(&mut receiver, "queued revision after cancellation")
+ .await
+ .revision(),
+ revision(1)
+ );
+ let tail = next_change(&mut receiver, "retained revision after cancellation").await;
+ assert_eq!(tail.revision(), revision(2));
+ assert_eq!(tail.previous_revision(), Some(revision(1)));
+ changes.publish(snapshot(3));
+ assert_eq!(
+ next_change(&mut receiver, "future revision after recovered tail")
+ .await
+ .revision(),
+ revision(3)
+ );
+ changes.close();
+ assert!(next_delivery(&mut receiver).await.is_none());
+ }
+
+ async fn next_delivery(receiver: &mut SnapshotChangeReceiver) -> Option<SnapshotChange> {
+ tokio::time::timeout(RECEIVE_TIMEOUT, receiver.receive())
+ .await
+ .expect("snapshot receive completed within its bound")
+ }
+
+ async fn next_change(receiver: &mut SnapshotChangeReceiver, label: &str) -> SnapshotChange {
+ next_delivery(receiver).await.expect(label)
+ }
+
fn revision(value: u64) -> SnapshotRevision {
SnapshotRevision::from_value(value)
}