Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
eaaafc9
fix(protocol): a new confirmation probe supersedes the last
mizanxali Oct 5, 2026
9eace23
Merge branch 'main' into fix/confirmation-probe-supersede
mizanxali Oct 6, 2026
44ba62e
chore(docs): update CHANGELOG.md
mizanxali Oct 6, 2026
ae4d4f5
fix(protocol): every confirmation probe supersedes, and none is restored
mizanxali Oct 6, 2026
734f91d
fix(protocol): withdraw the outstanding probe with the confirmation t…
mizanxali Oct 6, 2026
de77122
fix(protocol): keep confirmation probes out of persistent storage
mizanxali Oct 6, 2026
2691dcb
docs(protocol): rewrap the probe section
mizanxali Oct 6, 2026
e96f984
docs(protocol): a relay_pushed verdict backs probes off too
mizanxali Oct 6, 2026
e23b24f
fix(protocol): clear a probe whose send failed after it was queued
mizanxali Oct 6, 2026
e8a3b45
fix(protocol): a late relay verdict on a superseded probe still backs…
mizanxali Oct 6, 2026
e011573
fix(protocol): a late verdict on a superseded probe records the fact too
mizanxali Oct 6, 2026
1abf92d
docs(protocol): bound the late-verdict window, and log a drop only wh…
mizanxali Oct 6, 2026
b3e87e8
docs(protocol): name every penalty-free use of the outbound teardown
mizanxali Oct 6, 2026
18784f9
Merge remote-tracking branch 'origin/main' into fix/confirmation-prob…
bahdotsh Oct 6, 2026
712d518
docs(changelog): move the probe fix out of the shipped 0.28.0 notes
bahdotsh Oct 6, 2026
c1e25db
test(protocol): split the probe supersede test, and pin what it missed
bahdotsh Oct 6, 2026
d3caaa6
fix(protocol): pin that one restore can clear a full outbox of probes
bahdotsh Oct 6, 2026
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
14 changes: 14 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,19 @@ archived by series under [docs/changelog/](docs/changelog/); see the

### Fixed

- **Confirmation probes no longer crowd real messages out of the outbox.** A
session the peer never confirms is probed every five seconds, and each probe
used to enter the outbox with its own retry ladder on top of the unanswered
ones before it. Only a relay verdict backed that off, so a mesh-only link
never did. An account with ten such contacts stacked about two probes a
second, and once the outbox hit its 500-entry cap, eviction failed the user's
own messages with `Outbox capacity exceeded`. Seen on a device test between
an Android and an iPhone, offline over BLE. Each new probe now supersedes
the last, quietly and without counting against the carrier, whether the
periodic scan or the Welcome fast path sent it. Probes are no longer written
to storage, and those an older build wrote are dropped at restore. A peer
holds one probe in the outbox and is still probed on the same cadence.

- **A message carried through a dense cluster no longer dies inside it**
([#510](https://github.com/Offline-Protocol/offline-protocol-sdk/issues/510)).
A forwarder that received a second copy of a frame while its own forward
Expand Down Expand Up @@ -488,6 +501,7 @@ archived by series under [docs/changelog/](docs/changelog/); see the
unchanged. TypeScript's
`SecurityWarningCode` union gains the code, so an exhaustive `switch` over it
needs a new arm.

- **An old message no longer comes back as a push notification, again and
again.** A direct message the relay had pushed stayed in the outbox well past
a week, because each probe re-send refreshed its lifetime and so did each
Expand Down
21 changes: 21 additions & 0 deletions crates/offline-protocol/src/protocol/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -709,6 +709,25 @@ pub struct OfflineProtocol {
/// flap the connection. Reset on a reachability edge (`on_peer_presence`).
confirmation_probe_unreachable_parks: HashMap<String, u32>,

/// The confirmation probe each pending peer last sent. At most one probe
/// per peer lives in the outbox: a new one supersedes the last. Stacking a
/// fresh probe (new id, full retry ladder) every
/// `CONFIRMATION_PROBE_INTERVAL_SECS` on top of the unanswered ones filled
/// the outbox whenever no relay verdict (`unreachable` or
/// `relay_pushed`) arrived to back the schedule off, as on a mesh-only
/// carrier. Capacity eviction then failed real messages.
confirmation_probe_outstanding: HashMap<String, MessageId>,

/// The probe each pending peer's current one superseded. A relay verdict
/// finds its peer through the outbox entry, which the supersede removed,
/// so a verdict slower than `CONFIRMATION_PROBE_INTERVAL_SECS` would
/// otherwise be dropped and never back the schedule off, leaving the
/// peer probed every 5s against a relay that has said it is not there.
/// One id per peer, so this covers a verdict up to two intervals late;
/// one later than that is dropped, and the next probe's verdict backs
/// the schedule off instead.
confirmation_probe_superseded: HashMap<String, MessageId>,

/// Token bucket that caps how fast the RETRY/PROBE path
/// (`process_retry_queue`) re-sends, so a large backlog of resends to
/// unreachable peers cannot burst past the relay's per-connection rate limit
Expand Down Expand Up @@ -1176,6 +1195,8 @@ impl OfflineProtocol {
confirmation_retry_due_at: HashMap::new(),
confirmation_probe_due_at: HashMap::new(),
confirmation_probe_unreachable_parks: HashMap::new(),
confirmation_probe_outstanding: HashMap::new(),
confirmation_probe_superseded: HashMap::new(),
retry_drain_tokens: 0.0,
retry_drain_last: None,
rekey_due_at: HashMap::new(),
Expand Down
44 changes: 36 additions & 8 deletions crates/offline-protocol/src/protocol/send.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3277,12 +3277,28 @@ impl OfflineProtocol {
// nothing, and keeping it would only grow the map.
self.custody_receipts.remove(message_id);
if let Some(entry) = self.outbox.remove(message_id) {
self.clear_outbox_entry_from_storage(message_id);
if !Self::is_confirmation_probe(&entry.message) {
self.clear_outbox_entry_from_storage(message_id);
}
return Some(entry);
}
self.media_outbox.remove(message_id)
}

/// Takes `message_id` out of the three places that can send it again:
/// the retry queue, the pending-ACK tracker and the outbox. One without
/// the others resurrects it (see [`Self::retire_undeliverable_message`]).
/// No failure is recorded against the carrier: that is the give-up
/// path's, and a message withdrawn for being superseded did not fail.
pub(super) fn forget_outbound_message(
&mut self,
message_id: &MessageId,
) -> Option<OutboxEntry> {
self.retry_queue.remove(&message_id.as_str());
self.ack_manager.remove_ack(message_id);
self.remove_outbox_entry(message_id)
}

/// Takes a message this device has given up on out of everything that
/// could send it again, returning its outbox entry if it still had one.
///
Expand All @@ -3307,10 +3323,8 @@ impl OfflineProtocol {
message_id: &MessageId,
media_reason: &str,
) -> Option<OutboxEntry> {
self.retry_queue.remove(&message_id.as_str());
self.ack_manager.remove_ack(message_id);
self.handle_outbound_media_chunk_failed(message_id, media_reason);
let entry = self.remove_outbox_entry(message_id);
let entry = self.forget_outbound_message(message_id);
if let Some(transport) = entry.as_ref().and_then(|entry| entry.last_transport) {
self.transport_manager.record_delivery_failure(transport);
}
Expand Down Expand Up @@ -3534,6 +3548,17 @@ impl OfflineProtocol {
message.content_type == ContentType::FileChunk
}

/// A session confirmation probe lives in memory only. It carries no user
/// data, and a restarted process does not know which probe was the last,
/// so a restored one could never be superseded; the next scan asks the
/// same question. Persisting it cost a secure-storage write per probe and
/// a delete per supersede, under the protocol lock.
pub(super) fn is_confirmation_probe(message: &Message) -> bool {
message
.content
.starts_with(internal_prefixes::SESSION_CONFIRM_PROBE)
}

// ========================================================================
// DELIVERY TRACKING
// ========================================================================
Expand Down Expand Up @@ -4111,10 +4136,12 @@ impl OfflineProtocol {
None => match self.media_outbox.get(&parsed_id) {
Some(entry) => (entry, true),
None => {
debug!(
message_id = %message_id,
"Unreachable verdict for message without outbox entry, dropping"
);
if !self.note_superseded_probe_verdict(&parsed_id, carrier) {
debug!(
message_id = %message_id,
"Unreachable verdict for message without outbox entry, dropping"
);
}
return;
}
},
Expand Down Expand Up @@ -4270,6 +4297,7 @@ impl OfflineProtocol {
return;
};
if !self.is_parkable_plain_dm(&parsed_id) {
self.note_superseded_probe_verdict(&parsed_id, carrier);
return;
}
let Some(entry) = self.outbox.get_mut(&parsed_id) else {
Expand Down
114 changes: 106 additions & 8 deletions crates/offline-protocol/src/protocol/session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ use crate::mls_observability::MlsOperationContext;
use crate::protocol::reachability::{Claim, FactSource};
use crate::{Error, EstablishmentState, Event, Result, SessionStateError};
use chrono::{Duration as ChronoDuration, Utc};
use offline_protocol_core::{Message, MessagePriority};
use offline_protocol_core::{Message, MessageId, MessagePriority};
use offline_protocol_mls::{MlsManager, WelcomeMessage};
use offline_protocol_transport::TransportType;
use std::collections::HashSet;
Expand Down Expand Up @@ -449,6 +449,45 @@ impl OfflineProtocol {
self.confirmation_retry_due_at.remove(peer_id);
self.confirmation_probe_due_at.remove(peer_id);
self.confirmation_probe_unreachable_parks.remove(peer_id);
// Withdrawn with the tracking: once the map forgets it, no later
// probe can supersede it, and it would keep its full retry ladder.
if let Some(probe) = self.confirmation_probe_outstanding.remove(peer_id) {
self.forget_outbound_message(&probe);
}
self.confirmation_probe_superseded.remove(peer_id);
}

/// Applies a relay verdict on a probe that has since been superseded
/// (`confirmation_probe_superseded`) the way the verdict handlers apply
/// one on a queued probe: the reachability fact, then the probe back-off.
/// Consumed, so one probe's verdict escalates once, as it did while the
/// probe was still queued. A no-op for any other id; returns whether
/// the verdict was applied.
pub(super) fn note_superseded_probe_verdict(
&mut self,
message_id: &MessageId,
carrier: Option<TransportType>,
) -> bool {
let Some(peer) = self
.confirmation_probe_superseded
.iter()
.find(|(_, probe)| *probe == message_id)
.map(|(peer, _)| peer.clone())
else {
return false;
};
self.confirmation_probe_superseded.remove(&peer);
if let Some(carrier) = carrier {
self.reachability.record(
&peer,
carrier,
Claim::Unreachable,
FactSource::GatewayVerdict,
Instant::now(),
);
}
self.note_confirmation_probe_unreachable(&peer);
true
}

/// Fold a relay "recipient unreachable" verdict into the session-confirmation
Expand All @@ -457,7 +496,7 @@ impl OfflineProtocol {
/// [`Self::apply_recipient_unreachable_failure`]). Without this, a peer that
/// established an MLS session then vanished before confirming is re-probed
/// every `CONFIRMATION_PROBE_INTERVAL_SECS` (5s) indefinitely, because
/// `kick_pending_session_reconciliation` re-arms it every scan and
/// `send_session_confirmation_probe` re-arms it on every probe and
/// `on_transport_send_failed` otherwise never consults this scheduler. A
/// handful of such peers pushes aggregate relay traffic past the
/// per-connection rate limit, which disconnects the socket on a loop.
Expand Down Expand Up @@ -512,12 +551,44 @@ impl OfflineProtocol {
}

pub(super) fn send_session_confirmation_probe(&mut self, peer_id: &str, source_event: &str) {
// Superseded, not failed: the new probe asks the same question, so
// the old one's retry ladder would only stack onto it. Here rather
// than in the scan because the Welcome fast path sends one too, and
// a probe it replaced would keep its full ladder. Silent, and no
// delivery failure is recorded against the carrier. Before the send,
// not after: if the send errors the peer has no probe in flight
// until the next scan, but a new probe can never evict a message at
// outbox capacity to make room beside the one it replaces. The
// superseded id is kept so a relay verdict that arrives for it late
// still backs the schedule off (`note_superseded_probe_verdict`).
if let Some(previous) = self.confirmation_probe_outstanding.remove(peer_id) {
self.forget_outbound_message(&previous);
self.confirmation_probe_superseded
.insert(peer_id.to_string(), previous);
}
// The sender owns the cadence, so a fast-path probe is not superseded
// by the next scan before its ACK can arrive, and a relay verdict for
// it can back the schedule off (`note_confirmation_probe_unreachable`
// only escalates a peer that has a due time). Stamped before the send
// so a verdict delivered during it is not overwritten.
//
// A later due time is kept: the scan only sends once the due time has
// passed, but the fast path can fire while the schedule is backed off,
// and resetting it would undo the ladder a relay verdict built.
let next = Utc::now() + ChronoDuration::seconds(CONFIRMATION_PROBE_INTERVAL_SECS);
let due_at = self
.confirmation_probe_due_at
.entry(peer_id.to_string())
.or_insert(next);
*due_at = (*due_at).max(next);
match self.send_internal_message(
peer_id,
internal_prefixes::SESSION_CONFIRM_PROBE.to_string(),
MessagePriority::High,
) {
Ok(_) => {
Ok(message_id) => {
self.confirmation_probe_outstanding
.insert(peer_id.to_string(), message_id);
info!(
event = "session_confirmation_probe_sent",
session_or_group_id = %peer_id,
Expand All @@ -526,6 +597,22 @@ impl OfflineProtocol {
);
}
Err(err) => {
// The send can fail after the probe entered the outbox (a
// full ACK tracker refuses it), and then no id comes back to
// track. Any probe still queued for this peer is such an
// orphan, since the tracked one was withdrawn above.
let orphans: Vec<_> = self
.outbox
.values()
.filter(|entry| {
entry.message.recipient.as_str() == peer_id
&& Self::is_confirmation_probe(&entry.message)
})
.map(|entry| entry.message.id.clone())
.collect();
for orphan in orphans {
self.forget_outbound_message(&orphan);
}
warn!(
event = "session_confirmation_probe_failed",
session_or_group_id = %peer_id,
Expand Down Expand Up @@ -596,6 +683,22 @@ impl OfflineProtocol {
.retain(|peer, _| pending_set.contains(peer));
self.confirmation_probe_unreachable_parks
.retain(|peer, _| pending_set.contains(peer));
self.confirmation_probe_superseded
.retain(|peer, _| pending_set.contains(peer));
// Not a plain `retain`: a probe whose peer left the pending set
// without passing through `clear_confirmation_recovery_tracking`
// would be forgotten while still queued.
let settled: Vec<String> = self
.confirmation_probe_outstanding
.keys()
.filter(|peer| !pending_set.contains(*peer))
.cloned()
.collect();
for peer in settled {
if let Some(probe) = self.confirmation_probe_outstanding.remove(&peer) {
self.forget_outbound_message(&probe);
}
}

for peer_id in pending_peers {
let due_at = self
Expand All @@ -606,12 +709,7 @@ impl OfflineProtocol {
if due_at > now {
continue;
}

self.send_session_confirmation_probe(&peer_id, source_event);
self.confirmation_probe_due_at.insert(
peer_id,
now + ChronoDuration::seconds(CONFIRMATION_PROBE_INTERVAL_SECS),
);
}
}

Expand Down
36 changes: 35 additions & 1 deletion crates/offline-protocol/src/protocol/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -882,6 +882,16 @@ const _: () = assert!(
underflows to zero and it can never prune at all"
);

// The outbox walk charges a legacy confirmation probe's delete to its pool
// (`restore_outbox`). A build that stacked probes could leave a full outbox
// of them, and if the pool is smaller than that the walk breaks on probes and
// leaves real messages unrestored until a later launch.
const _: () = assert!(
MAX_RESTORE_PRUNE_DELETES >= MAX_OUTBOX_ENTRIES,
"an outbox full of legacy confirmation probes must clear in one launch, \
or the walk breaks on them and defers restoring real messages"
);

/// A pool of durable restore-path deletes, drawn on by one or more walks.
///
/// The unit the [`MAX_RESTORE_PRUNE_DELETES`] bound is actually about. A pool
Expand Down Expand Up @@ -3800,7 +3810,9 @@ impl OfflineProtocol {
let Some(storage) = &self.protocol_state_storage else {
return;
};
if Self::is_media_outbox_message(&entry.message) {
if Self::is_media_outbox_message(&entry.message)
|| Self::is_confirmation_probe(&entry.message)
{
return;
}
match serde_json::to_vec(entry) {
Expand Down Expand Up @@ -3941,6 +3953,7 @@ impl OfflineProtocol {
// `PruneAllowance::pool`.
let mut budget = allowance.counting();
let mut prune_bound_reached = false;
let mut legacy_probes_dropped = 0usize;
for message_id in message_ids.into_iter().take(OUTBOX_RESTORE_KEY_CAP) {
if budget.is_spent() {
prune_bound_reached = true;
Expand Down Expand Up @@ -3997,8 +4010,29 @@ impl OfflineProtocol {
continue;
}

// Probes are no longer persisted (`is_confirmation_probe`); drop
// any an older build wrote. Charged to the pool like any delete:
// the first launch after upgrading from the build that stacked
// them can spend most of it here, deferring the expiry and
// capacity prunes below to the next launch. That is the bound
// working, not a reason to make this delete unbudgeted.
if Self::is_confirmation_probe(&entry.message) {
budget.claim();
self.delete_outbox_key(&message_id);
legacy_probes_dropped += 1;
continue;
}

restored.push(entry);
}
if legacy_probes_dropped > 0 {
info!(
event = "outbox_entry_dropped",
repair_action = "legacy_confirmation_probe",
count = legacy_probes_dropped,
"outbox_entry_dropped"
);
}

// Drop entries past the absolute lifetime cap before anything else:
// the carrier-relative refresh below would otherwise re-grant them a
Expand Down
Loading
Loading