fix(pubsub): decide unsubscribes from desire re-read under the establishment lock

Review (Wren, ratified by Eva) found a renewal-vs-withdrawal race in the
reconciler. `reconcile()` snapshotted `desired()` and `established` once, then
ran both loops against those snapshots with awaited Redis round-trips in
between. Its own doc comment claimed the opposite — "read immediately before
each command" — describing a contract the code did not implement. I wrote that
comment, and it is exactly how the next reader would have re-introduced this.

The lost interleaving: desire drops community C, the pass records `have`
containing C, and before the unsubscribe loop reaches C a new authz resolution
calls `admit(C)`. `admit` sees the still-present establishment bit, permits a
cache insert, and the reconciler then withdraws establishment and unsubscribes
from its stale view. The entry stays readable for its whole 10s TTL with no
invalidation channel behind it — the silent stale-authz hole this lane exists
to close, reached through our own unsubscribe path instead of AWS's
restriction. A per-command re-read does not fix it: renewal can still land
between the check and the removal. `admit()`'s wake shortens the gap to
milliseconds but cannot close it, because pub/sub has no replay — the tick
restores state, nothing restores a dropped message.

Fix: `CommunityTopics::withdraw_undesired` re-reads `desired()` *inside* the
establishment write lock and returns only the communities to unsubscribe.
`admit()` already recorded residency before reading establishment, so the two
critical sections are mutually exclusive and only two orders exist: either the
renewal's residency write precedes the withdrawal's re-read (interest is seen,
the subscription survives) or it follows the withdrawal (establishment reads
false, nothing is admitted). An admitted entry can no longer outlive its
subscription. `remove_established` is gone, so no caller can withdraw outside
the lock; the ordering requirement is documented at both call sites, since each
is correct alone and only the pairing carries the property.

Subscribe stays on the pass's initial read: holding a subscription nobody wants
costs one idle channel until the next pass and can never make an entry readable
without coverage. Both families inherit the invariant because it lives in the
shared reconciler; conn-control's milder exposure (a lost disconnect leaves a
socket the DB ban row still refuses at next authz) is now stated in code rather
than implied.

Tests, all mutation-verified to bite:
- `a_renewal_racing_withdrawal_never_sees_establishment_it_loses` parks the
  `desired` closure inside the withdrawal critical section and releases a
  renewal thread from there, forcing the interleaving rather than sleeping for
  it. Moving the `desired()` read back outside the lock fails it.
- `interest_renewed_mid_reconcile_keeps_its_exact_subscription` proves it
  end-to-end through the real path against live Redis, asserting on
  `PUBSUB NUMSUB`. Its first draft did NOT bite the reinstated bug twice over:
  renewing on the pass's first read meant no withdrawal was ever attempted, and
  a coarse end-state assertion passed vacuously because the next tick
  re-subscribes. It now discriminates on read count within a window well under
  `RECONCILE_INTERVAL` and asserts establishment never transiently drops — the
  transient loss is the bug, the repaired end state is not proof.
- `withdrawal_reads_live_desire_not_the_passs_stale_view`,
  `withdrawal_still_removes_communities_nobody_wants`, and relay-side
  `admitted_residency_is_visible_to_a_concurrent_withdrawal` cover both
  directions and the cross-crate pairing.

Also corrects two stale comments Eva flagged: `state.rs` named the renamed
`may_cache` in an intra-doc link (a real rustdoc `broken_intra_doc_links`
warning I introduced), and `event.rs` test prose said "PSUBSCRIBEs" for a path
that has always used exact `SUBSCRIBE`.

Verified at this tree: buzz-pubsub 52/52 `--include-ignored` against standalone
Redis 8.6.2 and against a cluster-mode rig (`cluster_state:ok`, 16384 slots);
dead-port control fails all 6 Redis-backed reconciliation tests by name while
the 7 pure-logic tests still pass; buzz-relay --lib 807 passed / 0 failed;
workspace `clippy --all-targets -D warnings` rc=0; `cargo fmt --all --check`
clean.

Co-authored-by: npub17jjz49l9jjmhhk7cac63j8yt9z555n9cw8vk7v5jz4vzw4ppld5qgj57cc <f4a42a97e594b77bdbd8ee35191c8b28a94a4cb871d96f32921558275421fb68@buzz.block.builderlab.xyz>
Signed-off-by: npub17jjz49l9jjmhhk7cac63j8yt9z555n9cw8vk7v5jz4vzw4ppld5qgj57cc <f4a42a97e594b77bdbd8ee35191c8b28a94a4cb871d96f32921558275421fb68@buzz.block.builderlab.xyz>
This commit is contained in:
npub17jjz49l9jjmhhk7cac63j8yt9z555n9cw8vk7v5jz4vzw4ppld5qgj57cc
2026-07-31 16:04:43 -04:00
parent ab490cff5a
commit c05cdb15d1
3 changed files with 322 additions and 22 deletions
+276 -19
View File
@@ -8,8 +8,10 @@
//! * **desired** — owned by whoever holds the interest (live sockets for
//! connection control; local cache residency for cache invalidation),
//! supplied as a closure and recomputed from scratch on every reconcile. This
//! module never caches it, so a queued command can never be applied against a
//! stale view.
//! module never caches it. Crucially, the *withdrawal* direction re-reads it
//! under the same lock that removes establishment (see
//! [`CommunityTopics::withdraw_undesired`]), so interest renewed after the
//! diff was computed cannot have its subscription pulled out from under it.
//! * **established** — the communities Redis has acknowledged a `SUBSCRIBE` for
//! on the *current* connection. A community enters only after the ack returns,
//! and the whole set is cleared when the connection ends.
@@ -110,6 +112,50 @@ impl CommunityTopics {
self.insert_established(community);
}
/// Re-reads `desired` and withdraws establishment for everything it no
/// longer names, **atomically**, returning the communities whose channel the
/// caller must now `UNSUBSCRIBE`.
///
/// This is the ordering invariant the establishment gate rests on, and the
/// reason the decision cannot be made from a snapshot taken earlier in the
/// reconcile. Interest renewal ([`CacheResidency::admit`] in the relay)
/// records residency and *then* reads establishment under the read lock;
/// this method re-reads residency under the matching write lock. Mutual
/// exclusion makes the two orderings the only possibilities:
///
/// * renewal's establishment read precedes this critical section — then this
/// `desired()` call happens after the renewal recorded interest, sees the
/// community, and does not withdraw it. The admitted entry stays covered.
/// * renewal's establishment read follows this critical section — then the
/// withdrawal already happened, the read returns false, and nothing is
/// admitted in the first place.
///
/// So an entry admitted against establishment is never left readable behind
/// a withdrawn subscription. A `desired()` re-read *outside* the lock would
/// not give this: renewal could still land between the check and the removal.
///
/// `desired` must not acquire this lock, or this would deadlock. Both
/// production closures read only their own registry.
pub(crate) fn withdraw_undesired(&self, desired: &DesiredCommunities) -> Vec<CommunityId> {
let mut established = self
.established
.write()
.expect("community topics lock poisoned");
let want = desired();
let withdrawn: Vec<CommunityId> = established.difference(&want).copied().collect();
for community in &withdrawn {
established.remove(community);
}
withdrawn
}
/// Test seam for [`Self::withdraw_undesired`], whose invariant is shared with
/// the cache-residency gate in another crate.
#[doc(hidden)]
pub fn withdraw_undesired_for_test(&self, desired: &DesiredCommunities) -> Vec<CommunityId> {
self.withdraw_undesired(desired)
}
fn snapshot(&self) -> HashSet<CommunityId> {
self.established
.read()
@@ -124,13 +170,6 @@ impl CommunityTopics {
.insert(community);
}
fn remove_established(&self, community: CommunityId) {
self.established
.write()
.expect("community topics lock poisoned")
.remove(&community);
}
fn clear_established(&self) {
self.established
.write()
@@ -155,6 +194,17 @@ pub(crate) trait CommunityChannelFamily: Send + 'static {
/// Never returns — spawn it in a background task. Reconnects with exponential
/// backoff (1s → 2s → 4s → … → 30s max), rebuilding the exact subscription set
/// from `desired` on every connect.
///
/// The withdrawal ordering below is enforced for **both** families, not just
/// cache invalidation, because it lives in the shared reconciler. Cache
/// invalidation is what makes it a correctness requirement — a lost invalidation
/// leaves a readable authorization entry with nothing behind it. Connection
/// control's exposure is milder: a `disconnect` command lost in an unsubscribe
/// gap leaves a socket open that a ban has already closed elsewhere, and the DB
/// ban row still refuses that pubkey's next authorization, so the miss costs one
/// stale socket rather than a wrong access decision. The invariant is applied
/// uniformly anyway: one reconciler is less to get wrong than two, and nothing
/// here is per-family.
pub(crate) async fn run_community_subscriber<F: CommunityChannelFamily>(
redis_url: String,
family: F,
@@ -249,9 +299,20 @@ async fn connect_and_serve<F: CommunityChannelFamily>(
/// Bring the live connection's subscriptions in line with current desire.
///
/// `desired` is read here, inside the subscriber task, immediately before each
/// command — so an unsubscribe can never be applied against a stale view of
/// interest that has since been renewed.
/// The two directions are deliberately asymmetric, because only one of them can
/// break an invariant:
///
/// * **Subscribing** from the pass's initial `desired()` read is safe. Holding a
/// subscription nobody wants any more costs one idle channel until the next
/// pass withdraws it; it can never make an entry readable without coverage.
/// * **Unsubscribing** must not be decided from that same read. Between the read
/// and the command, interest can be renewed and an insert admitted against the
/// still-present establishment bit — and a pub/sub message lost in the
/// resulting gap is lost permanently, because pub/sub has no replay. The tick
/// restores *state*; nothing restores a dropped invalidation. So withdrawal
/// re-reads desire under the establishment write lock, via
/// [`CommunityTopics::withdraw_undesired`], and only the communities that
/// round-trip out of that critical section are unsubscribed.
async fn reconcile<F: CommunityChannelFamily>(
sink: &mut redis::aio::PubSubSink,
family: &F,
@@ -268,12 +329,12 @@ async fn reconcile<F: CommunityChannelFamily>(
topics.insert_established(*community);
}
for community in have.difference(&want) {
// Withdraw establishment before the unsubscribe, never after: a reader
// that consults `is_established` between the two would otherwise cache
// an entry this connection is about to stop protecting.
topics.remove_established(*community);
sink.unsubscribe(&family.channel(*community)).await?;
// Establishment is withdrawn inside the lock, so a reader that admits an
// entry either observed the bit before the withdrawal (and its renewed
// interest was seen here, keeping the subscription) or observes it after
// (and admits nothing).
for community in topics.withdraw_undesired(desired) {
sink.unsubscribe(&family.channel(community)).await?;
}
Ok(())
@@ -314,7 +375,9 @@ mod tests {
topics.insert_established(community(1));
topics.insert_established(community(2));
topics.remove_established(community(1));
let keep = community(2);
let desired: DesiredCommunities = Arc::new(move || HashSet::from([keep]));
assert_eq!(topics.withdraw_undesired(&desired), vec![community(1)]);
assert!(!topics.is_established(community(1)));
assert!(topics.is_established(community(2)));
@@ -332,6 +395,116 @@ mod tests {
.await
.expect("a wake raised before the wait must still complete");
}
/// Withdrawal must decide from desire read *inside* its own critical
/// section, not from a snapshot taken earlier in the reconcile pass.
///
/// This is the race Wren found in review. The subscribe loop awaits Redis
/// round-trips, so an arbitrary amount of time passes between the pass's
/// `desired()` read and the withdrawal — long enough for interest to be
/// renewed and a cache entry admitted against the still-present
/// establishment bit. Here the stale view says "withdraw c", live desire says
/// "keep c", and the assertion is that live desire wins.
#[test]
fn withdrawal_reads_live_desire_not_the_passs_stale_view() {
let topics = CommunityTopics::new("test");
let c = community(1);
topics.insert_established(c);
// The pass's initial read saw an empty desired set (c not wanted). By
// the time withdrawal runs, interest has been renewed.
let desired: DesiredCommunities = Arc::new(move || HashSet::from([c]));
assert!(
topics.withdraw_undesired(&desired).is_empty(),
"a community desired at withdrawal time must not be unsubscribed"
);
assert!(
topics.is_established(c),
"renewed interest must keep its establishment bit"
);
}
/// The other direction of the same call: still-undesired communities are
/// withdrawn, and establishment is dropped before the caller unsubscribes.
#[test]
fn withdrawal_still_removes_communities_nobody_wants() {
let topics = CommunityTopics::new("test");
let (kept, dropped) = (community(1), community(2));
topics.insert_established(kept);
topics.insert_established(dropped);
let desired: DesiredCommunities = Arc::new(move || HashSet::from([kept]));
assert_eq!(topics.withdraw_undesired(&desired), vec![dropped]);
assert!(
!topics.is_established(dropped),
"establishment must be gone before the UNSUBSCRIBE is issued"
);
assert!(topics.is_established(kept));
}
/// The invariant itself, forced deterministically: a renewal that observes
/// establishment can never have its subscription withdrawn by a
/// concurrently-running reconcile.
///
/// The `desired` closure parks *inside* the withdrawal critical section and
/// releases the renewal thread from there, so the renewal's
/// `is_established` read is guaranteed to happen while withdrawal holds the
/// write lock. That is the interleaving a sleep-based test only reaches by
/// luck. The reader must then observe `false` (withdrawal won the lock and
/// this admission is correctly refused) — never `true` while the community
/// also ends up unsubscribed, which is the stale-authz hole.
#[test]
fn a_renewal_racing_withdrawal_never_sees_establishment_it_loses() {
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc;
let topics = Arc::new(CommunityTopics::new("test"));
let c = community(1);
topics.insert_established(c);
// Signals the renewal thread that withdrawal now holds the write lock.
let (in_section_tx, in_section_rx) = mpsc::channel::<()>();
// Set by the renewal thread once its establishment read has returned.
let renewal_read = Arc::new(AtomicBool::new(false));
let read_result = Arc::new(AtomicBool::new(false));
let renewal_topics = Arc::clone(&topics);
let renewal_read_flag = Arc::clone(&renewal_read);
let renewal_result = Arc::clone(&read_result);
let renewal = std::thread::spawn(move || {
in_section_rx
.recv()
.expect("withdrawal must signal from inside its critical section");
// Blocks until withdrawal releases the write lock.
let admitted = renewal_topics.is_established(c);
renewal_result.store(admitted, Ordering::SeqCst);
renewal_read_flag.store(true, Ordering::SeqCst);
});
// Desire says "nobody wants c" — but hands the renewal thread its cue
// first, from inside the lock.
let desired: DesiredCommunities = Arc::new(move || {
in_section_tx.send(()).expect("renewal thread alive");
// Give the renewal thread every chance to reach its read while the
// write lock is still held; RwLock write exclusion is what makes
// the outcome deterministic regardless of how far it gets.
std::thread::sleep(Duration::from_millis(50));
HashSet::new()
});
assert_eq!(topics.withdraw_undesired(&desired), vec![c]);
renewal.join().expect("renewal thread panicked");
assert!(renewal_read.load(Ordering::SeqCst));
assert!(
!read_result.load(Ordering::SeqCst),
"a renewal whose interest was not seen by withdrawal must be refused, \
never admitted against an establishment bit that is being removed"
);
assert!(!topics.is_established(c));
}
}
/// Redis-backed reconciliation tests.
@@ -690,4 +863,88 @@ mod redis_tests {
topics.wake();
await_unestablished(&topics, c[0], "gate closes on withdrawal").await;
}
/// End-to-end against live Redis: interest renewed *during* a reconcile pass
/// must not lose its exact subscription.
///
/// The unit test above proves the lock discipline in isolation; this proves
/// the property survives the real `reconcile` path, where awaited Redis
/// round-trips sit between the pass's desire read and the withdrawal.
///
/// The closure is a phase machine keyed on read count, which is what makes
/// the interleaving deterministic: once the race is armed, the pass's *first*
/// read reports the community undesired (so a stale-snapshot reconciler would
/// queue an `UNSUBSCRIBE`), and every read after it reports the community
/// desired again (the renewal). Correct code re-reads inside the withdrawal
/// critical section, sees the renewal, and keeps the subscription. Verified
/// to bite: reinstating the `have.difference(&want)` withdrawal fails this
/// test on the `PUBSUB NUMSUB` assertion.
#[tokio::test]
#[ignore = "requires Redis"]
async fn interest_renewed_mid_reconcile_keeps_its_exact_subscription() {
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
let c = fresh(1)[0];
let armed = Arc::new(AtomicBool::new(false));
let reads = Arc::new(AtomicUsize::new(0));
let closure_armed = Arc::clone(&armed);
let closure_reads = Arc::clone(&reads);
let racing: DesiredCommunities = Arc::new(move || {
if !closure_armed.load(Ordering::SeqCst) {
return HashSet::from([c]);
}
// First read after arming: interest looks gone. Every later read
// (crucially, the one inside the withdrawal lock) sees it renewed.
if closure_reads.fetch_add(1, Ordering::SeqCst) == 0 {
HashSet::new()
} else {
HashSet::from([c])
}
});
let (topics, _rx) = spawn_subscriber(racing);
await_established(&topics, &[c], "pre-race establishment").await;
assert_eq!(exact_subscribers(c).await, 1);
armed.store(true, Ordering::SeqCst);
topics.wake();
// The discriminator, and why this bites where a coarser wait did not:
// correct code reads desire TWICE in one pass (the pass read, then the
// withdrawal read inside the lock), microseconds apart. A stale-snapshot
// reconciler reads once, unsubscribes, and only reads again at the next
// tick — a whole `RECONCILE_INTERVAL` later, by which point it has
// re-subscribed and a naive end-state assertion passes vacuously.
//
// So: poll far faster than the tick, require the second read inside a
// window well below it, and assert establishment never drops along the
// way. The transient loss is the bug; the repaired end state is not proof.
let window = RECONCILE_INTERVAL / 4;
let deadline = std::time::Instant::now() + window;
while reads.load(Ordering::SeqCst) < 2 {
assert!(
topics.is_established(c),
"establishment was withdrawn for renewed interest — an \
invalidation published now would be lost, and the later tick \
that re-subscribes cannot bring it back"
);
assert!(
std::time::Instant::now() < deadline,
"withdrawal did not re-read desire within its own pass ({window:?}); \
a second read this late means it decided from the stale snapshot"
);
tokio::time::sleep(Duration::from_millis(1)).await;
}
assert!(
topics.is_established(c),
"renewed interest must not lose establishment"
);
assert_eq!(
exact_subscribers(c).await,
1,
"renewed interest must still hold an exact SUBSCRIBE in Redis"
);
}
}
+2 -2
View File
@@ -1692,7 +1692,7 @@ mod tests {
register_presence_sub(&receiver, "receiver-presence");
// Under the community-scoped bus, Redis delivery is demand-driven:
// a relay only PSUBSCRIBEs `buzz:{community}:global` after it retains
// a relay only SUBSCRIBEs `buzz:{community}:global` after it retains
// interest in that topic. Both relays share one explicit tenant and
// retain Global before publishing — origin too, so the echo-
// suppression assertion still exercises `mark_local_event` against a
@@ -1710,7 +1710,7 @@ mod tests {
.retain_topic(&tenant, EventTopic::Global)
.await;
// Match buzz-pubsub's own Redis round-trip test: give PSUBSCRIBE a
// Match buzz-pubsub's own Redis round-trip test: give SUBSCRIBE a
// bounded moment to attach before publishing the single test event.
tokio::time::sleep(std::time::Duration::from_millis(200)).await;
+44 -1
View File
@@ -193,6 +193,16 @@ impl CacheResidency {
/// Gating **inserts only** is deliberate: gating reads would discard a live,
/// still-invalidatable cache on every reconnect and turn a pub/sub blip into
/// a DB stampede.
///
/// **The order of the two statements below is load-bearing.** Residency is
/// recorded *before* establishment is read, and the reconciler withdraws
/// establishment while re-reading residency under the matching write lock
/// (`CommunityTopics::withdraw_undesired`). Those two facts together are
/// what stop an admitted entry from outliving its subscription: whichever
/// side wins the lock, either this renewal is visible to the withdrawal
/// decision (so the subscription survives) or the withdrawal already
/// happened (so `is_established` is false and nothing is admitted). Reading
/// establishment first, or withdrawing outside the lock, reopens the gap.
#[must_use]
pub fn admit(&self, community: CommunityId) -> bool {
let deadline = Instant::now() + self.ttl;
@@ -201,6 +211,7 @@ impl CacheResidency {
// does not depend on this — the reconcile tick is the guarantee.
self.topics.wake();
}
// Must follow the residency write above; see the ordering note.
let admitted = self.topics.is_established(community);
if !admitted {
metrics::counter!("buzz_authz_cache_insert_skipped_total").increment(1);
@@ -928,7 +939,7 @@ impl AppState {
/// Check channel membership with a 10-second cache. Falls back to DB on miss.
///
/// The result is cached only while this pod's cache-invalidation channel
/// for `community_id` is established — see [`CacheResidency::may_cache`].
/// for `community_id` is established — see [`CacheResidency::admit`].
/// Ungated, a pod could cache an authorization decision it would never hear
/// the invalidation for, and serve it for the whole TTL.
pub async fn is_member_cached(
@@ -2117,6 +2128,38 @@ mod tests {
);
}
/// The cross-crate half of the ordering invariant: a community `admit()`
/// just made resident must survive a withdrawal decision taken concurrently
/// by the reconciler.
///
/// `admit` records residency *before* reading establishment, and the
/// reconciler re-reads residency under the write lock, so a renewal that got
/// its residency in cannot be unsubscribed. This test pairs the two real
/// call sites — the relay's residency closure and the pubsub withdrawal —
/// because each is correct alone and only the pairing carries the property.
#[test]
fn admitted_residency_is_visible_to_a_concurrent_withdrawal() {
let topics = Arc::new(CommunityTopics::new("test"));
let residency = Arc::new(CacheResidency::new(Arc::clone(&topics), AUTHZ_CACHE_TTL));
let community = CommunityId::from_uuid(Uuid::from_u128(0xa));
// Exactly the closure main.rs hands the cache-invalidation subscriber.
let for_closure = Arc::clone(&residency);
let desired: buzz_pubsub::community_topics::DesiredCommunities =
Arc::new(move || for_closure.resident_communities());
// Established and resident: the steady state.
topics.insert_established_for_test(community);
assert!(residency.admit(community));
// Reconciler asks whether to withdraw. Residency is live, so it must not.
assert!(
topics.withdraw_undesired_for_test(&desired).is_empty(),
"a resident community must never be unsubscribed"
);
assert!(residency.admit(community), "and it stays admissible");
}
/// A newly resident community wakes the cache-invalidation subscriber, so
/// the subscribe does not wait for the reconcile tick.
#[tokio::test]