From a28957d79b84d862b6a411f250f33379aa4772bf Mon Sep 17 00:00:00 2001 From: Cea Stapleton Cordasco <261786559+cea-block@users.noreply.github.com> Date: Fri, 31 Jul 2026 12:17:47 -0500 Subject: [PATCH] fix(audio): release failed admission lease Signed-off-by: Cea Stapleton Cordasco <261786559+cea-block@users.noreply.github.com> --- crates/buzz-relay/src/audio/handler.rs | 68 ++++++++++-- crates/buzz-relay/src/audio/join.rs | 142 +++++++++++++++++++++++++ crates/buzz-relay/src/audio/room.rs | 4 + 3 files changed, 206 insertions(+), 8 deletions(-) diff --git a/crates/buzz-relay/src/audio/handler.rs b/crates/buzz-relay/src/audio/handler.rs index 692b2ac9b..7bfd6d4b2 100644 --- a/crates/buzz-relay/src/audio/handler.rs +++ b/crates/buzz-relay/src/audio/handler.rs @@ -151,6 +151,48 @@ fn default_protocol_version() -> u8 { 1 } +/// Remove a denied private admission and release only the exact owner lease +/// that this connection acquired. The room is sealed while it is still the +/// manager-visible instance, so a remote registration that already holds its +/// `Arc` cannot enter between the peer removal and the Redis release. +async fn cleanup_failed_private_audio_admission( + state: &Arc, + tenant: &TenantContext, + channel_id: Uuid, + room: &Arc, + peer_id: Uuid, + acquired_lease: &mut Option, +) { + let directory = state + .mesh() + .map(|mesh| &mesh.directory as &dyn crate::audio::join::HuddleDirectory); + match crate::audio::join::cleanup_failed_admission_lease( + directory, + acquired_lease, + &state.audio_rooms, + tenant.community(), + channel_id, + room, + peer_id, + ) + .await + { + Ok(Some(crate::audio::join::HuddleReleaseOutcome::Released)) | Ok(None) => {} + Ok(Some(crate::audio::join::HuddleReleaseOutcome::NotOwner)) => { + debug!( + channel_id = %channel_id, + "failed audio admission lease already moved; stale cleanup left current owner intact" + ); + } + Err(e) => { + warn!( + channel_id = %channel_id, + "failed audio admission could not release huddle owner lease: {e}" + ); + } + } +} + async fn handle_audio_connection( socket: WebSocket, state: Arc, @@ -672,10 +714,15 @@ async fn handle_active_audio_connection( ) .await; } - room.remove_peer(peer_id); - state - .audio_rooms - .cleanup_if_empty(tenant.community(), channel_id); + cleanup_failed_private_audio_admission( + &state, + &tenant, + channel_id, + &room, + peer_id, + &mut acquired_lease, + ) + .await; return; } }; @@ -709,10 +756,15 @@ async fn handle_active_audio_connection( ) .await; } - room.remove_peer(peer_id); - state - .audio_rooms - .cleanup_if_empty(tenant.community(), channel_id); + cleanup_failed_private_audio_admission( + &state, + &tenant, + channel_id, + &room, + peer_id, + &mut acquired_lease, + ) + .await; return; } }; diff --git a/crates/buzz-relay/src/audio/join.rs b/crates/buzz-relay/src/audio/join.rs index ddadb13f7..6cf60c311 100644 --- a/crates/buzz-relay/src/audio/join.rs +++ b/crates/buzz-relay/src/audio/join.rs @@ -233,6 +233,48 @@ pub enum HuddleReleaseOutcome { NotOwner, } +/// Remove a denied admission, atomically seal its room when it was the last +/// peer, and release only the exact freshly acquired owner token before the +/// empty room is evicted. Keeping the sealed room manager-visible during the +/// awaited release prevents a registration holding the old room `Arc` from +/// entering under that generation. The Redis release itself is owner- and +/// generation-matched, so stale cleanup cannot delete a replacement token. +pub async fn cleanup_failed_admission_lease( + directory: Option<&dyn HuddleDirectory>, + acquired_lease: &mut Option, + rooms: &AudioRoomManager, + community_id: CommunityId, + session_id: Uuid, + room: &Arc, + peer_id: Uuid, +) -> Result, MeshError> { + let sealed_empty = room + .remove_peer_and_check_ended(peer_id) + .map(|(_, ended)| ended) + .unwrap_or(false); + if !sealed_empty { + return Ok(None); + } + + let released = if let Some(lease) = acquired_lease.take() { + let result = match directory { + Some(directory) => directory.release(&lease).await, + None => Err(MeshError::Transport( + "acquired huddle lease has no directory".to_string(), + )), + }; + result.map(Some) + } else { + Ok(None) + }; + + // Preserve the pre-existing bounded Redis-error behavior: an unrenewed + // token expires at its TTL, while the empty local room is immediately + // reusable instead of becoming a permanent `ended` tombstone. + rooms.cleanup_if_empty(community_id, session_id); + released +} + /// Result of an ownership acquire attempt. #[derive(Clone, Debug, PartialEq, Eq)] pub enum AcquireOutcome { @@ -1831,6 +1873,7 @@ mod tests { // yields `Renewed` (lease holds). `release` returns the scripted value. renew_outcomes: Mutex>, release_outcome: Mutex>, + release_fails: Mutex, renew_calls: Mutex, release_calls: Mutex, } @@ -1906,6 +1949,9 @@ mod tests { } async fn release(&self, _lease: &HuddleLease) -> Result { *self.release_calls.lock().unwrap() += 1; + if *self.release_fails.lock().unwrap() { + return Err(MeshError::Transport("injected release failure".into())); + } Ok(self .release_outcome .lock() @@ -2571,6 +2617,102 @@ mod tests { ); } + /// Both private-admission failure exits share the same cleanup primitive: + /// seal and evict the failed room, release the exact freshly acquired lease + /// token once, and let an immediate retry acquire the next generation. + #[tokio::test] + async fn failed_identity_admissions_release_lease_and_allow_immediate_retry() { + for failure_case in ["identity_conflict", "identity_storage_failure"] { + let session = Uuid::new_v4(); + let rooms = AudioRoomManager::new(); + let room = rooms.get_or_create(community(), session); + let (peer_id, _, _, _) = room.add_peer(failure_case.into(), 1).unwrap(); + let dir = Arc::new(FakeDir::owned_by(Ownership { + owner_runtime_id: rt(1), + generation: 5, + })); + let mut acquired = Some(lease_for(session, 5)); + + assert_eq!( + cleanup_failed_admission_lease( + Some(&*dir), + &mut acquired, + &rooms, + community(), + session, + &room, + peer_id, + ) + .await + .unwrap(), + Some(HuddleReleaseOutcome::Released), + "{failure_case} must release its exact owner token" + ); + assert!(acquired.is_none()); + assert_eq!(*dir.release_calls.lock().unwrap(), 1); + assert!(rooms.get(community(), session).is_none()); + assert!(matches!( + room.add_peer("stale-room".into(), 1), + Err(AdmissionError::Ended) + )); + + // Model Redis's successful exact-token delete and monotonic next + // generation, then retry immediately in this same task. + *dir.owner.lock().unwrap() = None; + *dir.acquire.lock().unwrap() = Some(AcquireOutcome::Acquired(lease_for(session, 6))); + let registry = HuddleOwnerRegistry::new(); + let retried = tokio::time::timeout( + Duration::from_secs(2), + resolve_join_owner_ready(&*dir, community(), session, rt(1), ®istry), + ) + .await + .expect("retry must not wait for the 30-second lease TTL") + .expect("retry acquires the released huddle"); + assert_eq!(retried.outcome, JoinOutcome::LocalOwner { generation: 6 }); + assert_eq!( + retried.acquired.as_ref().map(HuddleLease::generation), + Some(6) + ); + } + } + + /// A Redis error preserves 037's bounded-TTL fallback without leaving the + /// local manager permanently pinned to the sealed failed room. + #[tokio::test] + async fn failed_admission_release_error_does_not_tombstone_room() { + let session = Uuid::new_v4(); + let rooms = AudioRoomManager::new(); + let room = rooms.get_or_create(community(), session); + let (peer_id, _, _, _) = room.add_peer("failed".into(), 1).unwrap(); + let dir = FakeDir::owned_by(Ownership { + owner_runtime_id: rt(1), + generation: 5, + }); + *dir.release_fails.lock().unwrap() = true; + let mut acquired = Some(lease_for(session, 5)); + + assert!(cleanup_failed_admission_lease( + Some(&dir), + &mut acquired, + &rooms, + community(), + session, + &room, + peer_id, + ) + .await + .is_err()); + assert!(acquired.is_none()); + assert_eq!(*dir.release_calls.lock().unwrap(), 1); + assert!(rooms.get(community(), session).is_none()); + + let replacement = rooms.get_or_create(community(), session); + assert!(!Arc::ptr_eq(&room, &replacement)); + replacement + .add_peer("retry".into(), 1) + .expect("release errors must not permanently tombstone the room"); + } + /// `drain` is generation-fenced like `release`, but unlike room-empty it /// also cancels the drain signal so local owner peers and remote control /// streams can rejoin with an explicit draining cause before the renewer diff --git a/crates/buzz-relay/src/audio/room.rs b/crates/buzz-relay/src/audio/room.rs index d5c428698..c7d95d43c 100644 --- a/crates/buzz-relay/src/audio/room.rs +++ b/crates/buzz-relay/src/audio/room.rs @@ -682,6 +682,10 @@ mod tests { .remove_peer_and_check_ended(peer_id) .expect("peer existed"); assert!(ended, "single-peer room should end on its last departure"); + let err = room1 + .add_peer("late-peer".to_string(), 2) + .expect_err("a stale room handle must not admit after empty cleanup seals it"); + assert!(matches!(err, AdmissionError::Ended)); assert!(manager.cleanup_if_empty(community_id, channel_id)); // Next joiner with a different version on the same channel id gets a