From ddaaef418b86415ef861a0c99c28f015331edec0 Mon Sep 17 00:00:00 2001 From: Kenny Lopez Date: Mon, 17 Aug 2026 17:17:31 +0100 Subject: [PATCH] Harden mobile Huddle reliability Co-authored-by: Kenny Lopez Co-authored-by: Carl <3c4caeafb646d23867f1c4832e68211d77e2561946171625f75c3ce1a3f2670f@buzz.block.builderlab.xyz> Signed-off-by: Kenny Lopez --- crates/buzz-relay/src/audio/handler.rs | 148 ++++++-- crates/buzz-relay/src/audio/join.rs | 31 +- crates/buzz-relay/src/audio/room.rs | 88 +++-- .../block/buzz/mobile/HuddleAudioEngine.kt | 99 +++++- .../block/buzz/mobile/HuddleMediaPlugin.kt | 19 + .../mobile/HuddleActiveTalkerSelectorTest.kt | 50 +++ mobile/ios/Runner/HuddleAudioEngine.swift | 104 +++++- mobile/ios/Runner/HuddleMediaPlugin.swift | 36 ++ mobile/ios/RunnerTests/RunnerTests.swift | 25 ++ .../channel_detail_page/huddle_drawer.dart | 1 + .../channels/mobile_huddle_controller.dart | 85 ++++- mobile/lib/shared/huddle/huddle_media.dart | 16 + mobile/lib/shared/huddle/huddle_session.dart | 106 +++++- .../lib/shared/huddle/huddle_transport.dart | 293 +++++++++++++-- mobile/lib/shared/relay/relay_session.dart | 23 +- .../channels/channel_detail_page_test.dart | 107 +++++- .../shared/huddle/huddle_session_test.dart | 160 ++++++++- .../shared/huddle/huddle_transport_test.dart | 336 +++++++++++++++++- 18 files changed, 1566 insertions(+), 161 deletions(-) create mode 100644 mobile/android/app/src/test/kotlin/xyz/block/buzz/mobile/HuddleActiveTalkerSelectorTest.kt diff --git a/crates/buzz-relay/src/audio/handler.rs b/crates/buzz-relay/src/audio/handler.rs index 4c158eab0..dbbb1b2f9 100644 --- a/crates/buzz-relay/src/audio/handler.rs +++ b/crates/buzz-relay/src/audio/handler.rs @@ -511,11 +511,11 @@ async fn handle_active_audio_connection( let admission = if let Some(session) = remote_session.as_ref() { room.add_peer_at_index(pubkey_hex.clone(), requested_version, session.peer_index()) - .map(|(id, audio, ctrl)| (id, session.peer_index(), audio, ctrl)) + .map(|(id, audio, ctrl, revision)| (id, session.peer_index(), audio, ctrl, revision)) } else { room.add_peer(pubkey_hex.clone(), requested_version) }; - let (peer_id, peer_index, audio_rx, peer_ctrl_rx) = match admission { + let (peer_id, peer_index, audio_rx, peer_ctrl_rx, admission_revision) = match admission { Ok(v) => v, Err(crate::audio::room::AdmissionError::Full) => { warn!(channel_id = %channel_id, "audio room full (255 peers exhausted)"); @@ -608,22 +608,37 @@ async fn handle_active_audio_connection( // Remote registration and owner-assigned ingress admission completed above. - let peers_snapshot: Vec = if let Some(session) = remote_session.as_ref() { - session - .roster() - .peers - .iter() - .map(|peer| serde_json::json!({"pubkey": peer.pubkey, "peer_index": peer.peer_index})) - .collect() - } else { - room.peer_pubkeys() - .into_iter() - .map(|(pk, idx)| serde_json::json!({"pubkey": pk, "peer_index": idx})) - .collect() - }; + let (peers_snapshot, roster_revision): (Vec, u64) = + if let Some(session) = remote_session.as_ref() { + ( + session + .roster() + .peers + .iter() + .map(|peer| { + serde_json::json!({"pubkey": peer.pubkey, "peer_index": peer.peer_index}) + }) + .collect(), + session.roster().revision, + ) + } else { + let snapshot = room.roster_snapshot(); + ( + snapshot + .peers + .into_iter() + .map(|peer| { + serde_json::json!({"pubkey": peer.pubkey, "peer_index": peer.peer_index}) + }) + .collect(), + snapshot.revision, + ) + }; + debug_assert!(roster_revision >= admission_revision); let joined_msg = serde_json::json!({ "type": "joined", + "revision": roster_revision, "pubkey": pubkey_hex, "peer_index": peer_index, "peers": peers_snapshot, @@ -684,6 +699,7 @@ async fn handle_active_audio_connection( data_tx, ctrl_tx.clone(), fwd_cancel, + cancel.clone(), )); // Non-owner path: own the owner's `HuddleControl` stream in a reader task. @@ -811,23 +827,27 @@ async fn handle_active_audio_connection( // AdmissionGuard lock across index recycling AND the is_empty + ended=true // check. Ingress mirrors never archive authoritative huddle state; they // remove locally and let the owner decide room lifetime. - let should_auto_end = if remote_session.is_some() { - room.remove_peer(peer_id); - false + let removal = if remote_session.is_some() { + room.remove_peer(peer_id).map(|delta| (delta, false)) } else { room.remove_peer_and_check_ended(peer_id) - .map(|(_, ended)| ended) - .unwrap_or(false) }; + let should_auto_end = removal.as_ref().map(|(_, ended)| *ended).unwrap_or(false); - let left_msg = serde_json::json!({ - "type": "left", - "pubkey": pubkey_hex, - "peer_index": peer_index, - }) - .to_string(); if remote_session.is_none() { - room.broadcast_control(left_msg); + if let Some((delta, _)) = removal { + let left = delta + .left + .expect("peer removal deltas always carry the removed peer"); + let left_msg = serde_json::json!({ + "type": "left", + "revision": delta.revision, + "pubkey": left.pubkey, + "peer_index": left.peer_index, + }) + .to_string(); + room.broadcast_control(left_msg); + } } emit_participant_event( @@ -1115,6 +1135,7 @@ async fn audio_forward_loop( data_tx: mpsc::Sender, ctrl_tx: mpsc::Sender, cancel: CancellationToken, + connection_cancel: CancellationToken, ) { loop { tokio::select! { @@ -1124,9 +1145,18 @@ async fn audio_forward_loop( msg = peer_ctrl_rx.recv() => { match msg { Some(PeerCtrl::Json(json)) => { - let _ = ctrl_tx.try_send(WsMessage::Text(json.into())); + if ctrl_tx.try_send(WsMessage::Text(json.into())).is_err() { + // State-bearing roster control may not be dropped. + // Closing the connection forces admission to replay + // a fresh authoritative snapshot. + connection_cancel.cancel(); + break; + } + } + Some(PeerCtrl::Close) | None => { + connection_cancel.cancel(); + break; } - Some(PeerCtrl::Close) | None => break, } } frame = audio_rx.recv() => { @@ -1433,6 +1463,66 @@ mod tests { received } + #[tokio::test] + async fn saturated_websocket_control_queue_cancels_the_audio_connection() { + let (_audio_tx, audio_rx) = mpsc::channel(1); + let (peer_ctrl_tx, peer_ctrl_rx) = mpsc::channel(2); + let (data_tx, _data_rx) = mpsc::channel(1); + let (ctrl_tx, _ctrl_rx) = mpsc::channel(1); + ctrl_tx + .try_send(WsMessage::Ping(Bytes::new())) + .expect("fill websocket control queue"); + peer_ctrl_tx + .try_send(PeerCtrl::Json("{}".into())) + .expect("queue state-bearing control"); + let task_cancel = CancellationToken::new(); + let connection_cancel = CancellationToken::new(); + + audio_forward_loop( + audio_rx, + peer_ctrl_rx, + data_tx, + ctrl_tx, + task_cancel, + connection_cancel.clone(), + ) + .await; + + assert!( + connection_cancel.is_cancelled(), + "saturated websocket control must force a fresh roster admission" + ); + } + + #[tokio::test] + async fn closed_peer_control_queue_cancels_the_audio_connection() { + let (_audio_tx, audio_rx) = mpsc::channel(1); + let (peer_ctrl_tx, peer_ctrl_rx) = mpsc::channel(1); + let (data_tx, _data_rx) = mpsc::channel(1); + let (ctrl_tx, _ctrl_rx) = mpsc::channel(1); + let task_cancel = CancellationToken::new(); + let connection_cancel = CancellationToken::new(); + + let forward = tokio::spawn(audio_forward_loop( + audio_rx, + peer_ctrl_rx, + data_tx, + ctrl_tx, + task_cancel, + connection_cancel.clone(), + )); + drop(peer_ctrl_tx); + + tokio::time::timeout(Duration::from_secs(1), forward) + .await + .expect("forwarder exits when its state-bearing queue closes") + .expect("forwarder task completes cleanly"); + assert!( + connection_cancel.is_cancelled(), + "lost control state must tear down the WebSocket for a fresh roster" + ); + } + #[tokio::test] async fn audio_send_loop_sends_policy_close_when_community_is_deleted() { use futures_util::Sink; diff --git a/crates/buzz-relay/src/audio/join.rs b/crates/buzz-relay/src/audio/join.rs index ddadb13f7..0fd7261b9 100644 --- a/crates/buzz-relay/src/audio/join.rs +++ b/crates/buzz-relay/src/audio/join.rs @@ -1292,14 +1292,16 @@ impl HuddleControlAcceptor { self.rooms .get(CommunityId::from_uuid(community_id), session_id) }) { - let peer_index = room.peers.get(&peer_id).map(|peer| peer.peer_index); - room.remove_peer(peer_id); - if let Some(peer_index) = peer_index { + if let Some(delta) = room.remove_peer(peer_id) { + let left = delta + .left + .expect("peer removal deltas always carry the removed peer"); room.broadcast_control( serde_json::json!({ "type": "left", - "pubkey": pubkey, - "peer_index": peer_index, + "revision": delta.revision, + "pubkey": left.pubkey, + "peer_index": left.peer_index, }) .to_string(), ); @@ -1351,15 +1353,17 @@ impl HuddleControlAcceptor { self.rooms .get(CommunityId::from_uuid(community_id), session_id) }) { - for (pubkey, peer_id) in registered { - let peer_index = room.peers.get(&peer_id).map(|peer| peer.peer_index); - room.remove_peer(peer_id); - if let Some(peer_index) = peer_index { + for (_pubkey, peer_id) in registered { + if let Some(delta) = room.remove_peer(peer_id) { + let left = delta + .left + .expect("peer removal deltas always carry the removed peer"); room.broadcast_control( serde_json::json!({ "type": "left", - "pubkey": pubkey, - "peer_index": peer_index, + "revision": delta.revision, + "pubkey": left.pubkey, + "peer_index": left.peer_index, }) .to_string(), ); @@ -1381,7 +1385,7 @@ impl HuddleControlAcceptor { registered: &mut std::collections::HashMap, ) -> HuddleControlMsg { match room.add_peer(pubkey.to_string(), protocol_version) { - Ok((peer_id, peer_index, audio_rx, _peer_ctrl_rx)) => { + Ok((peer_id, peer_index, audio_rx, _peer_ctrl_rx, roster_revision)) => { registered.insert(pubkey.to_string(), peer_id); // The owner's Room fans out to this remote peer's `audio_tx`; // the sink drains `audio_rx` and ships each frame as a datagram @@ -1389,6 +1393,7 @@ impl HuddleControlAcceptor { spawn_remote_peer_sink(Arc::clone(&self.transport), from, fenced, audio_rx); let joined = serde_json::json!({ "type": "joined", + "revision": roster_revision, "pubkey": pubkey, "peer_index": peer_index, "peers": [{"pubkey": pubkey, "peer_index": peer_index}], @@ -2277,7 +2282,7 @@ mod tests { let fenced = fenced_owned_by(owner_rt, session_id); let rooms = Arc::new(AudioRoomManager::new()); let room = rooms.get_or_create(community(), session_id); - let (_local_id, _local_index, _audio_rx, mut local_ctrl_rx) = + let (_local_id, _local_index, _audio_rx, mut local_ctrl_rx, _revision) = room.add_peer("owner-local".into(), 2).unwrap(); // Discard the local peer's own roster delta; this assertion targets the // websocket-compatible control fanout below. diff --git a/crates/buzz-relay/src/audio/room.rs b/crates/buzz-relay/src/audio/room.rs index d5c428698..5a688094c 100644 --- a/crates/buzz-relay/src/audio/room.rs +++ b/crates/buzz-relay/src/audio/room.rs @@ -229,7 +229,16 @@ impl Room { &self, pubkey: String, requested_version: u8, - ) -> Result<(Uuid, u8, mpsc::Receiver, mpsc::Receiver), AdmissionError> { + ) -> Result< + ( + Uuid, + u8, + mpsc::Receiver, + mpsc::Receiver, + u64, + ), + AdmissionError, + > { let mut g = self.guard.lock().map_err( |_| AdmissionError::Ended, /* poisoned ≈ shutting down */ )?; @@ -265,14 +274,15 @@ impl Room { }, ); g.roster_revision = g.roster_revision.wrapping_add(1); + let revision = g.roster_revision; let delta = RosterDelta { - revision: g.roster_revision, + revision, joined: Some(RosterPeer { pubkey, peer_index }), left: None, }; let _ = self.roster_tx.send(delta); drop(g); // Release lock after ordered roster publication. - Ok((peer_id, peer_index, audio_rx, ctrl_rx)) + Ok((peer_id, peer_index, audio_rx, ctrl_rx, revision)) } /// Add a non-owner ingress peer at the index already allocated by the @@ -283,7 +293,7 @@ impl Room { pubkey: String, requested_version: u8, peer_index: u8, - ) -> Result<(Uuid, mpsc::Receiver, mpsc::Receiver), AdmissionError> { + ) -> Result<(Uuid, mpsc::Receiver, mpsc::Receiver, u64), AdmissionError> { let mut g = self.guard.lock().map_err(|_| AdmissionError::Ended)?; if g.ended { return Err(AdmissionError::Ended); @@ -323,43 +333,45 @@ impl Room { }, ); g.roster_revision = g.roster_revision.wrapping_add(1); + let revision = g.roster_revision; let delta = RosterDelta { - revision: g.roster_revision, + revision, joined: Some(RosterPeer { pubkey, peer_index }), left: None, }; let _ = self.roster_tx.send(delta); drop(g); - Ok((peer_id, audio_rx, ctrl_rx)) + Ok((peer_id, audio_rx, ctrl_rx, revision)) } - /// Remove a peer and recycle its index. - pub fn remove_peer(&self, peer_id: Uuid) { + /// Remove a peer and recycle its index. Returns the ordered roster delta + /// when the peer existed. + pub fn remove_peer(&self, peer_id: Uuid) -> Option { let Ok(mut g) = self.guard.lock() else { - return; + return None; }; - if let Some((_, peer)) = self.peers.remove(&peer_id) { - g.release(peer.peer_index); - g.roster_revision = g.roster_revision.wrapping_add(1); - let delta = RosterDelta { - revision: g.roster_revision, - joined: None, - left: Some(RosterPeer { - pubkey: peer.pubkey, - peer_index: peer.peer_index, - }), - }; - let _ = self.roster_tx.send(delta); - drop(g); - } + let (_, peer) = self.peers.remove(&peer_id)?; + g.release(peer.peer_index); + g.roster_revision = g.roster_revision.wrapping_add(1); + let delta = RosterDelta { + revision: g.roster_revision, + joined: None, + left: Some(RosterPeer { + pubkey: peer.pubkey, + peer_index: peer.peer_index, + }), + }; + let _ = self.roster_tx.send(delta.clone()); + drop(g); + Some(delta) } /// Remove a peer AND atomically check if the room should end. /// If the room is now empty, sets `ended = true` under the same lock /// acquisition that recycles the index — no window for a concurrent /// `add_peer` to sneak in between removal and the ended flag. - /// Returns `(peer_index, should_auto_end)`. - pub fn remove_peer_and_check_ended(&self, peer_id: Uuid) -> Option<(u8, bool)> { + /// Returns `(roster_delta, should_auto_end)`. + pub fn remove_peer_and_check_ended(&self, peer_id: Uuid) -> Option<(RosterDelta, bool)> { let mut g = self.guard.lock().ok()?; let (_, peer) = self.peers.remove(&peer_id)?; let peer_index = peer.peer_index; @@ -382,9 +394,9 @@ impl Room { } else { false }; - let _ = self.roster_tx.send(delta); + let _ = self.roster_tx.send(delta.clone()); drop(g); - Some((peer_index, should_end)) + Some((delta, should_end)) } /// Fan-out a binary frame to all peers except the sender. @@ -431,19 +443,23 @@ impl Room { /// Send a JSON control message to all peers via the control channel. /// Separate from audio so control is never starved by audio backpressure. /// Control messages (joined/left) are state-bearing — the client's - /// peer_index→pubkey map depends on receiving every one. The channel is - /// sized generously (32 slots) so drops should never happen in practice; - /// if they do, we log a warning so the issue is visible. + /// peer_index→pubkey map depends on receiving every one. Saturation is + /// therefore terminal for that receiver: dropping its sender closes the + /// queue, forcing a reconnect with a fresh authoritative admission snapshot. pub fn broadcast_control(&self, json: String) { - for entry in self.peers.iter() { + for mut entry in self.peers.iter_mut() { if entry .ctrl_tx .try_send(PeerCtrl::Json(json.clone())) .is_err() { + let (replacement_tx, replacement_rx) = mpsc::channel(1); + drop(replacement_rx); + let old_tx = std::mem::replace(&mut entry.ctrl_tx, replacement_tx); + drop(old_tx); tracing::warn!( peer_id = %entry.key(), - "control channel full — dropped state-bearing message (peer map may desync)" + "control channel full — closing receiver for authoritative roster resync" ); } } @@ -568,7 +584,7 @@ mod tests { let (_local_id, local_index, ..) = room.add_peer("owner-local".into(), 2).unwrap(); assert_eq!(local_index, 0); - let (remote_id, _audio, _ctrl) = room + let (remote_id, _audio, _ctrl, _revision) = room .add_peer_at_index("remote".into(), 2, 7) .expect("owner-assigned index admits"); assert_eq!(room.peers.get(&remote_id).unwrap().peer_index, 7); @@ -674,7 +690,7 @@ mod tests { let channel_id = Uuid::new_v4(); let room1 = manager.get_or_create(community_id, channel_id); - let (peer_id, _, _, _) = room1 + let (peer_id, _, _, _, _) = room1 .add_peer("alice".to_string(), 2) .expect("first peer admits"); // Last peer leaves and ends the room atomically. @@ -736,14 +752,14 @@ mod tests { #[test] fn version_pin_persists_across_peer_churn() { let room = fresh_room(); - let (alice_id, alice_idx, _, _) = + let (alice_id, alice_idx, _, _, _) = room.add_peer("alice".to_string(), 2).expect("alice admits"); room.remove_peer(alice_id); // Room is non-empty thanks to nothing yet — wait, alice left and // nobody else is here. Add bob with the same version: should work. // Then add carol with a different version: should fail with the // *original* pin, even though alice already left. - let (_, bob_idx, _, _) = room + let (_, bob_idx, _, _, _) = room .add_peer("bob".to_string(), 2) .expect("bob admits at v=2"); assert_eq!( diff --git a/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleAudioEngine.kt b/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleAudioEngine.kt index a39a17e40..b01076c6f 100644 --- a/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleAudioEngine.kt +++ b/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleAudioEngine.kt @@ -38,9 +38,55 @@ internal data class HuddleRemoteOpusPacket( val peerIndex: Int, val sequence: Int, val timestamp48k: Long, + val levelDbov: Int, val opus: ByteArray, ) +/// Fixed-capacity active-talker selector. Roster membership stays in Dart; +/// native decoder/track resources exist only for peers that actually send. +internal class HuddleActiveTalkerSelector( + private val capacity: Int, +) { + internal data class Activity( + val peerIndex: Int, + val lastPacketOrdinal: Long, + val levelDbov: Int, + ) + + private val active = mutableMapOf() + private var ordinal = 0L + + fun activate(peerIndex: Int, levelDbov: Int): Int? { + ordinal += 1 + if (active.containsKey(peerIndex)) { + active[peerIndex] = Activity(peerIndex, ordinal, levelDbov) + return null + } + val evicted = if (active.size >= capacity) { + active.values.minWithOrNull( + compareBy { it.lastPacketOrdinal } + .thenBy { it.levelDbov } + .thenBy { it.peerIndex }, + )?.peerIndex.also { index -> index?.let(active::remove) } + } else { + null + } + active[peerIndex] = Activity(peerIndex, ordinal, levelDbov) + return evicted + } + + fun remove(peerIndex: Int) { + active.remove(peerIndex) + } + + fun indices(): Set = active.keys.toSet() +} + +internal sealed interface HuddlePlaybackCommand { + data class Packet(val packet: HuddleRemoteOpusPacket) : HuddlePlaybackCommand + data class RemovePeer(val peerIndex: Int) : HuddlePlaybackCommand +} + /** Aggregate, debug-only microphone telemetry. No PCM or encoded audio is retained. */ internal data class HuddleCaptureDiagnostics( val frameCount: Long, @@ -76,7 +122,7 @@ internal class HuddleAudioEngine( ) { private val running = AtomicBoolean(false) private val failureReported = AtomicBoolean(false) - private val playbackQueue = ArrayBlockingQueue(PLAYBACK_QUEUE_CAPACITY) + private val playbackQueue = ArrayBlockingQueue(PLAYBACK_QUEUE_CAPACITY) private val interrupted = AtomicBoolean(false) @Volatile @@ -132,10 +178,28 @@ internal class HuddleAudioEngine( fun enqueueRemote(packet: HuddleRemoteOpusPacket) { check(running.get()) { "Huddle audio is not running." } - if (!playbackQueue.offer(packet)) { - playbackQueue.poll() - playbackQueue.offer(packet) + offerPlayback(HuddlePlaybackCommand.Packet(packet)) + } + + fun removeRemotePeer(peerIndex: Int) { + check(running.get()) { "Huddle audio is not running." } + playbackQueue.removeIf { command -> + command is HuddlePlaybackCommand.Packet && command.packet.peerIndex == peerIndex } + offerPlayback(HuddlePlaybackCommand.RemovePeer(peerIndex)) + } + + private fun offerPlayback(command: HuddlePlaybackCommand) { + if (playbackQueue.offer(command)) return + + // Packets are expendable realtime data; peer-removal commands are not. + // Make room by evicting the oldest packet only, so an audio flood can + // never resurrect stale decoder/track state after index reuse. + val oldestPacket = playbackQueue.firstOrNull { + it is HuddlePlaybackCommand.Packet + } + if (oldestPacket != null) playbackQueue.remove(oldestPacket) + playbackQueue.offer(command) } fun stop() { @@ -355,23 +419,34 @@ internal class HuddleAudioEngine( private fun runPlayback() { Process.setThreadPriority(Process.THREAD_PRIORITY_AUDIO) val playbacks = mutableMapOf() + val activeTalkers = HuddleActiveTalkerSelector(MAX_REMOTE_PEERS) try { while (running.get()) { - val packet = try { + val command = try { playbackQueue.poll(FRAME_DURATION_US, TimeUnit.MICROSECONDS) } catch (_: InterruptedException) { break } - if (packet != null) { - val playback = playbacks[packet.peerIndex] ?: run { - check(playbacks.size < MAX_REMOTE_PEERS) { - "This device reached the mobile Huddle participant limit." - } - PeerPlayback(packet.peerIndex).also { + when (command) { + is HuddlePlaybackCommand.Packet -> { + val packet = command.packet + val evicted = activeTalkers.activate( + packet.peerIndex, + packet.levelDbov, + ) + evicted?.let { playbacks.remove(it)?.release() } + val playback = playbacks[packet.peerIndex] ?: PeerPlayback( + packet.peerIndex, + ).also { playbacks[packet.peerIndex] = it } + playback.enqueue(packet) } - playback.enqueue(packet) + is HuddlePlaybackCommand.RemovePeer -> { + activeTalkers.remove(command.peerIndex) + playbacks.remove(command.peerIndex)?.release() + } + null -> Unit } if (!interrupted.get()) { for (playback in playbacks.values) playback.drainOne() diff --git a/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleMediaPlugin.kt b/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleMediaPlugin.kt index 3946b589d..dd67e65c9 100644 --- a/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleMediaPlugin.kt +++ b/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleMediaPlugin.kt @@ -64,6 +64,7 @@ internal class HuddleMediaPlugin( "setMuted" -> setMuted(call.arguments, result) "setSpeakerEnabled" -> setSpeakerEnabled(call.arguments, result) "playRemoteOpusFrame" -> playRemoteOpusFrame(call.arguments, result) + "removeRemotePeer" -> removeRemotePeer(call.arguments, result) "stop" -> stop(result) else -> result.notImplemented() } @@ -253,10 +254,12 @@ internal class HuddleMediaPlugin( val peerIndex = (values?.get("peerIndex") as? Number)?.toInt() val sequence = (values?.get("sequence") as? Number)?.toInt() val timestamp48k = (values?.get("timestamp48k") as? Number)?.toLong() + val levelDbov = (values?.get("levelDbov") as? Number)?.toInt() val opus = values?.get("opus") as? ByteArray if (peerIndex == null || peerIndex !in 0..255 || sequence == null || sequence !in 0..0xffff || timestamp48k == null || timestamp48k !in 0..0xffff_ffffL || + levelDbov == null || levelDbov !in -127..0 || opus == null || opus.isEmpty() || opus.size > MAX_OPUS_PACKET_BYTES ) { result.error("invalid_arguments", "Malformed remote Huddle Opus packet.", null) @@ -270,6 +273,7 @@ internal class HuddleMediaPlugin( peerIndex = peerIndex, sequence = sequence, timestamp48k = timestamp48k, + levelDbov = levelDbov, opus = opus.copyOf(), ), ) @@ -279,6 +283,21 @@ internal class HuddleMediaPlugin( } } + private fun removeRemotePeer(arguments: Any?, result: MethodChannel.Result) { + val peerIndex = ((arguments as? Map<*, *>)?.get("peerIndex") as? Number)?.toInt() + if (peerIndex == null || peerIndex !in 0..255) { + result.error("invalid_arguments", "Missing Huddle peer index.", null) + return + } + try { + val engine = audioEngine ?: error("Huddle audio is not running.") + engine.removeRemotePeer(peerIndex) + result.success(null) + } catch (error: Throwable) { + result.error("playback_failed", error.message, null) + } + } + private fun stop(result: MethodChannel.Result) { try { audioEngine?.stop() diff --git a/mobile/android/app/src/test/kotlin/xyz/block/buzz/mobile/HuddleActiveTalkerSelectorTest.kt b/mobile/android/app/src/test/kotlin/xyz/block/buzz/mobile/HuddleActiveTalkerSelectorTest.kt new file mode 100644 index 000000000..ae71b8d29 --- /dev/null +++ b/mobile/android/app/src/test/kotlin/xyz/block/buzz/mobile/HuddleActiveTalkerSelectorTest.kt @@ -0,0 +1,50 @@ +package xyz.block.buzz.mobile + +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNull +import org.junit.Test + +class HuddleActiveTalkerSelectorTest { + @Test + fun `sixteen senders retain exactly the active talker capacity`() { + val selector = HuddleActiveTalkerSelector(capacity = 15) + + repeat(15) { peer -> assertNull(selector.activate(peer, -20)) } + assertEquals(0, selector.activate(15, -20)) + assertEquals((1..15).toSet(), selector.indices()) + } + + @Test + fun `recent activity wins and an evicted peer can reactivate`() { + val selector = HuddleActiveTalkerSelector(capacity = 2) + + selector.activate(4, -40) + selector.activate(2, -20) + selector.activate(4, -10) + assertEquals(2, selector.activate(9, -30)) + assertEquals(setOf(4, 9), selector.indices()) + + assertEquals(4, selector.activate(2, -15)) + assertEquals(setOf(2, 9), selector.indices()) + } + + @Test + fun `stable peer index breaks equal activity ties`() { + val selector = HuddleActiveTalkerSelector(capacity = 1) + + selector.activate(7, -30) + assertEquals(7, selector.activate(3, -30)) + assertEquals(setOf(3), selector.indices()) + } + + @Test + fun `remove clears the slot without allocating roster-only peers`() { + val selector = HuddleActiveTalkerSelector(capacity = 2) + + assertEquals(emptySet(), selector.indices()) + selector.activate(7, -20) + selector.remove(7) + assertEquals(emptySet(), selector.indices()) + assertNull(selector.activate(7, -20)) + } +} diff --git a/mobile/ios/Runner/HuddleAudioEngine.swift b/mobile/ios/Runner/HuddleAudioEngine.swift index 618f92740..5c9cbd6d3 100644 --- a/mobile/ios/Runner/HuddleAudioEngine.swift +++ b/mobile/ios/Runner/HuddleAudioEngine.swift @@ -13,6 +13,7 @@ struct HuddleRemoteOpusPacket { let peerIndex: Int let sequence: Int let timestamp48k: Int64 + let levelDbov: Int let opus: Data } @@ -265,6 +266,63 @@ enum HuddleAudioLevels { } } +/// Fixed-capacity active-talker selector. The UI roster is independent; native +/// decoder/player resources exist only for peers that actually send packets. +struct HuddleActiveTalkerSelector { + struct Activity { + let peerIndex: Int + let lastPacketOrdinal: UInt64 + let levelDbov: Int + } + + let capacity: Int + private(set) var active: [Int: Activity] = [:] + private var ordinal: UInt64 = 0 + + init(capacity: Int) { + self.capacity = capacity + } + + mutating func activate(peerIndex: Int, levelDbov: Int) -> Int? { + ordinal &+= 1 + if active[peerIndex] != nil { + active[peerIndex] = Activity( + peerIndex: peerIndex, + lastPacketOrdinal: ordinal, + levelDbov: levelDbov + ) + return nil + } + let evicted = active.count >= capacity + ? active.values.min { left, right in + if left.lastPacketOrdinal != right.lastPacketOrdinal { + return left.lastPacketOrdinal < right.lastPacketOrdinal + } + if left.levelDbov != right.levelDbov { + return left.levelDbov < right.levelDbov + } + return left.peerIndex < right.peerIndex + }?.peerIndex + : nil + if let evicted { active.removeValue(forKey: evicted) } + active[peerIndex] = Activity( + peerIndex: peerIndex, + lastPacketOrdinal: ordinal, + levelDbov: levelDbov + ) + return evicted + } + + mutating func remove(_ peerIndex: Int) { + active.removeValue(forKey: peerIndex) + } + + mutating func removeAll() { + active.removeAll() + ordinal = 0 + } +} + /// Foreground-only iOS realtime media engine for the fixed Huddle Opus v2 path. /// /// AVAudioEngine owns voice-processed capture and per-peer mixed playout. @@ -288,6 +346,7 @@ final class HuddleAudioEngine { private var interrupted = false private var failureReported = false private var pendingCaptureBuffers = 0 + private var pendingRemotePackets = 0 private var captureConverter: AVAudioConverter? private var captureSourceFormat: AVAudioFormat? private var captureTapInstalled = false @@ -300,6 +359,7 @@ final class HuddleAudioEngine { private var rmsDbovHistogram: [Int: Int64] = [:] private var peakDbovHistogram: [Int: Int64] = [:] private var peerPlaybacks: [Int: HuddlePeerPlayback] = [:] + private var activeTalkers = HuddleActiveTalkerSelector(capacity: 15) private var playbackTimer: DispatchSourceTimer? private var configurationObserver: NSObjectProtocol? @@ -395,25 +455,35 @@ final class HuddleAudioEngine { func enqueueRemote(_ packet: HuddleRemoteOpusPacket) throws { stateLock.lock() - let isRunning = running - stateLock.unlock() - guard isRunning else { + guard running else { + stateLock.unlock() throw HuddleNativeMediaError.invalidState( "Huddle audio is not running." ) } + guard pendingRemotePackets < 50 else { + stateLock.unlock() + return + } + pendingRemotePackets += 1 + stateLock.unlock() processingQueue.async { [weak self] in guard let self else { return } + defer { completeRemotePacket() } do { + let evictedIndex = activeTalkers.activate( + peerIndex: packet.peerIndex, + levelDbov: packet.levelDbov + ) + if let evictedIndex, + let evicted = peerPlaybacks.removeValue(forKey: evictedIndex) + { + evicted.release(from: audioEngine) + } let playback: HuddlePeerPlayback if let existing = peerPlaybacks[packet.peerIndex] { playback = existing } else { - guard peerPlaybacks.count < 15 else { - throw HuddleNativeMediaError.invalidState( - "This device reached the mobile Huddle participant limit." - ) - } playback = try HuddlePeerPlayback( engine: audioEngine, pcmFormat: pcmFormat @@ -427,6 +497,22 @@ final class HuddleAudioEngine { } } + private func completeRemotePacket() { + stateLock.lock() + pendingRemotePackets = max(0, pendingRemotePackets - 1) + stateLock.unlock() + } + + func removeRemotePeer(_ peerIndex: Int) { + processingQueue.async { [weak self] in + guard let self else { return } + activeTalkers.remove(peerIndex) + guard let playback = peerPlaybacks.removeValue(forKey: peerIndex) + else { return } + playback.release(from: audioEngine) + } + } + func stop() { processingQueue.sync { stopOnQueue() @@ -484,6 +570,7 @@ final class HuddleAudioEngine { muted = false interrupted = false pendingCaptureBuffers = 0 + pendingRemotePackets = 0 stateLock.unlock() playbackTimer?.cancel() @@ -496,6 +583,7 @@ final class HuddleAudioEngine { playback.release(from: audioEngine) } peerPlaybacks.removeAll() + activeTalkers.removeAll() audioEngine.stop() audioEngine.reset() try? audioEngine.inputNode.setVoiceProcessingEnabled(false) diff --git a/mobile/ios/Runner/HuddleMediaPlugin.swift b/mobile/ios/Runner/HuddleMediaPlugin.swift index 407bbde55..518739110 100644 --- a/mobile/ios/Runner/HuddleMediaPlugin.swift +++ b/mobile/ios/Runner/HuddleMediaPlugin.swift @@ -76,6 +76,8 @@ final class HuddleMediaPlugin { setSpeakerEnabled(arguments: call.arguments, result: result) case "playRemoteOpusFrame": playRemoteOpusFrame(arguments: call.arguments, result: result) + case "removeRemotePeer": + removeRemotePeer(arguments: call.arguments, result: result) case "stop": stop(result: result) default: @@ -314,6 +316,8 @@ final class HuddleMediaPlugin { (0...0xffff).contains(sequence), let timestamp = (values["timestamp48k"] as? NSNumber)?.int64Value, (0...Int64(UInt32.max)).contains(timestamp), + let levelDbov = (values["levelDbov"] as? NSNumber)?.intValue, + (-127...0).contains(levelDbov), let typedData = values["opus"] as? FlutterStandardTypedData, !typedData.data.isEmpty, typedData.data.count <= HuddleAudioFormats.maximumOpusPacketBytes @@ -339,6 +343,7 @@ final class HuddleMediaPlugin { peerIndex: peerIndex, sequence: sequence, timestamp48k: timestamp, + levelDbov: levelDbov, opus: typedData.data ) ) @@ -354,6 +359,37 @@ final class HuddleMediaPlugin { } } + private func removeRemotePeer( + arguments: Any?, + result: @escaping FlutterResult + ) { + guard let values = arguments as? [String: Any], + let peerIndex = (values["peerIndex"] as? NSNumber)?.intValue, + (0...255).contains(peerIndex) + else { + result( + FlutterError( + code: "invalid_arguments", + message: "Missing Huddle peer index.", + details: nil + ) + ) + return + } + guard let audioEngine else { + result( + FlutterError( + code: "invalid_state", + message: "Huddle audio is not running.", + details: nil + ) + ) + return + } + audioEngine.removeRemotePeer(peerIndex) + result(nil) + } + private func handleInterruption(_ notification: Notification) { guard audioSessionPrepared, let rawType = notification.userInfo?[AVAudioSessionInterruptionTypeKey] diff --git a/mobile/ios/RunnerTests/RunnerTests.swift b/mobile/ios/RunnerTests/RunnerTests.swift index b275036e9..2ca38c4df 100644 --- a/mobile/ios/RunnerTests/RunnerTests.swift +++ b/mobile/ios/RunnerTests/RunnerTests.swift @@ -7,6 +7,31 @@ import XCTest class RunnerTests: XCTestCase { + func testHuddleActiveTalkerSelectorBoundsAndReactivates() { + var selector = HuddleActiveTalkerSelector(capacity: 15) + + for peer in 0..<15 { + XCTAssertNil(selector.activate(peerIndex: peer, levelDbov: -20)) + } + XCTAssertEqual(selector.activate(peerIndex: 15, levelDbov: -20), 0) + XCTAssertEqual(Set(selector.active.keys), Set(1...15)) + + selector.remove(7) + XCTAssertFalse(selector.active.keys.contains(7)) + XCTAssertNil(selector.activate(peerIndex: 7, levelDbov: -10)) + XCTAssertTrue(selector.active.keys.contains(7)) + } + + func testHuddleActiveTalkerSelectorUsesRecentActivity() { + var selector = HuddleActiveTalkerSelector(capacity: 2) + + XCTAssertNil(selector.activate(peerIndex: 4, levelDbov: -40)) + XCTAssertNil(selector.activate(peerIndex: 2, levelDbov: -20)) + XCTAssertNil(selector.activate(peerIndex: 4, levelDbov: -10)) + XCTAssertEqual(selector.activate(peerIndex: 9, levelDbov: -30), 2) + XCTAssertEqual(Set(selector.active.keys), Set([4, 9])) + } + func testHuddleOpusCodecRoundTripsFixedV2Frame() throws { let encoder = try HuddleOpusEncoder() let decoder = try HuddleOpusDecoder() diff --git a/mobile/lib/features/channels/channel_detail_page/huddle_drawer.dart b/mobile/lib/features/channels/channel_detail_page/huddle_drawer.dart index 8bffac15c..3c0bbf339 100644 --- a/mobile/lib/features/channels/channel_detail_page/huddle_drawer.dart +++ b/mobile/lib/features/channels/channel_detail_page/huddle_drawer.dart @@ -19,6 +19,7 @@ class MobileHuddleShell extends HookConsumerWidget { @override Widget build(BuildContext context, WidgetRef ref) { + ref.read(mobileHuddleControllerProvider.notifier); final presentation = ref.watch(mobileHuddlePresentationProvider); final session = ref.watch(huddleSessionProvider); final relayStatus = ref.watch(relaySessionProvider).status; diff --git a/mobile/lib/features/channels/mobile_huddle_controller.dart b/mobile/lib/features/channels/mobile_huddle_controller.dart index eacefa269..62460157e 100644 --- a/mobile/lib/features/channels/mobile_huddle_controller.dart +++ b/mobile/lib/features/channels/mobile_huddle_controller.dart @@ -70,16 +70,38 @@ final huddleHumanCountProvider = Provider((ref) { }; }); +/// One successful relay admission. Its epoch changes only when a newer join or +/// start attempt begins, so duplicate teardown calls cannot invalidate the +/// cleanup already in flight for the same admission. +final class _HuddleAdmissionToken { + final int epoch; + + const _HuddleAdmissionToken(this.epoch); +} + /// Coordinates the Nostr lifecycle around the foreground audio session. final class MobileHuddleController extends Notifier { var _generation = 0; + var _admissionEpoch = 0; + _HuddleAdmissionToken? _admissionToken; + final Set<_HuddleAdmissionToken> _finishedLifecycleAdmissions = {}; @override bool build() { + final unregisterBeforePause = ref + .read(relaySessionProvider.notifier) + .registerBeforePause(_leaveForBackground); + ref.onDispose(unregisterBeforePause); + ref.listen(huddleSessionProvider, (previous, next) { + if (previous?.wasAdmitted == true && + next.phase == HuddleSessionPhase.failed) { + unawaited(_cleanupAfterLocalTeardown(next)); + } + }); ref.listen(appLifecycleProvider, (_, next) { if (next == AppLifecycleState.paused || next == AppLifecycleState.detached) { - unawaited(leave()); + unawaited(_leaveForBackground()); } }); return false; @@ -88,6 +110,7 @@ final class MobileHuddleController extends Notifier { Future start({required String parentChannelId}) async { if (state) return; final generation = ++_generation; + final admissionEpoch = ++_admissionEpoch; state = true; final actions = ref.read(channelActionsProvider); String? backingChannelId; @@ -118,6 +141,7 @@ final class MobileHuddleController extends Notifier { if (session.phase == HuddleSessionPhase.failed) { throw StateError(session.error ?? 'Unable to join the new Huddle.'); } + _admissionToken = _HuddleAdmissionToken(admissionEpoch); } catch (_) { if (backingChannelId != null) { if (announced) { @@ -149,6 +173,7 @@ final class MobileHuddleController extends Notifier { required String startedEventId, }) async { ++_generation; + final admissionEpoch = ++_admissionEpoch; final currentPubkey = ref.read(currentPubkeyProvider); await ref .read(huddleSessionProvider.notifier) @@ -163,12 +188,23 @@ final class MobileHuddleController extends Notifier { currentPubkey.toLowerCase() == startedBy.toLowerCase(), startedEventId: startedEventId, ); + final session = ref.read(huddleSessionProvider); + if (session.wasAdmitted) { + _admissionToken = _HuddleAdmissionToken(admissionEpoch); + } } + Future? _backgroundLeave; + + Future _leaveForBackground() => + _backgroundLeave ??= leave().whenComplete(() => _backgroundLeave = null); + Future leave() async { ++_generation; state = false; final session = ref.read(huddleSessionProvider); + final admissionToken = session.wasAdmitted ? _admissionToken : null; + if (identical(_admissionToken, admissionToken)) _admissionToken = null; final parentChannelId = session.parentChannelId; final backingChannelId = session.ephemeralChannelId; final humanCount = !session.wasAdmitted || backingChannelId == null @@ -190,8 +226,11 @@ final class MobileHuddleController extends Notifier { localFailureStackTrace = stackTrace; } - if (backingChannelId != null && humanCount != null) { + if (backingChannelId != null && + humanCount != null && + admissionToken != null) { await _finishLeaveLifecycle( + admissionToken: admissionToken, parentChannelId: parentChannelId, backingChannelId: backingChannelId, humanCount: humanCount, @@ -202,13 +241,32 @@ final class MobileHuddleController extends Notifier { } } + Future _cleanupAfterLocalTeardown(HuddleSessionState session) async { + final backingChannelId = session.ephemeralChannelId; + final admissionToken = session.wasAdmitted ? _admissionToken : null; + if (admissionToken == null || backingChannelId == null) return; + if (identical(_admissionToken, admissionToken)) _admissionToken = null; + final humanCount = ref + .read(huddleHumanCountProvider)(backingChannelId) + .catchError((_) => 2); + await _finishLeaveLifecycle( + admissionToken: admissionToken, + parentChannelId: session.parentChannelId, + backingChannelId: backingChannelId, + humanCount: humanCount, + ); + } + Future _finishLeaveLifecycle({ + required _HuddleAdmissionToken admissionToken, required String? parentChannelId, required String backingChannelId, required Future humanCount, }) async { + if (!_finishedLifecycleAdmissions.add(admissionToken)) return; final actions = ref.read(channelActionsProvider); final humansRemaining = await humanCount; + if (admissionToken.epoch != _admissionEpoch) return; if (humansRemaining <= 1 && parentChannelId != null) { // Desktop auto-ends when the departing person is the last human. // Both lifecycle publication and archival are best effort there. @@ -218,6 +276,7 @@ final class MobileHuddleController extends Notifier { ephemeralChannelId: backingChannelId, ); } catch (_) {} + if (admissionToken.epoch != _admissionEpoch) return; try { await actions.archiveChannel(backingChannelId); } catch (_) {} @@ -235,16 +294,29 @@ final class MobileHuddleController extends Notifier { ++_generation; state = false; final session = ref.read(huddleSessionProvider); + final admissionToken = session.wasAdmitted ? _admissionToken : null; + if (identical(_admissionToken, admissionToken)) _admissionToken = null; final parentChannelId = session.parentChannelId; final backingChannelId = session.ephemeralChannelId; if (!session.isCreator || + admissionToken == null || parentChannelId == null || backingChannelId == null) { throw StateError('Only the Huddle creator can end it.'); } final actions = ref.read(channelActionsProvider); + if (!_finishedLifecycleAdmissions.add(admissionToken)) return; Object? failure; + try { + await ref.read(huddleSessionProvider.notifier).leave(); + } catch (error) { + failure = error; + } + if (admissionToken.epoch != _admissionEpoch) { + if (failure != null) throw failure; + return; + } try { await actions.announceHuddleEnded( parentChannelId: parentChannelId, @@ -253,12 +325,19 @@ final class MobileHuddleController extends Notifier { } catch (error) { failure = error; } + if (admissionToken.epoch != _admissionEpoch) { + if (failure != null) throw failure; + return; + } try { await actions.archiveChannel(backingChannelId); } catch (error) { failure ??= error; } - await ref.read(huddleSessionProvider.notifier).leave(); + if (admissionToken.epoch != _admissionEpoch) { + if (failure != null) throw failure; + return; + } if (failure != null) throw failure; } diff --git a/mobile/lib/shared/huddle/huddle_media.dart b/mobile/lib/shared/huddle/huddle_media.dart index 19891a591..4dcf73263 100644 --- a/mobile/lib/shared/huddle/huddle_media.dart +++ b/mobile/lib/shared/huddle/huddle_media.dart @@ -237,6 +237,7 @@ abstract interface class HuddleMedia { Future setMuted(bool muted); Future setSpeakerEnabled(bool enabled); Future playRemoteFrame(HuddleRemoteAudioFrame frame); + Future removeRemotePeer(int peerIndex); Future stop(); Future dispose(); } @@ -534,6 +535,21 @@ final class MethodChannelHuddleMedia implements HuddleMedia { } } + @override + Future removeRemotePeer(int peerIndex) async { + _ensureNotDisposed(); + if (_state.phase != HuddleMediaPhase.active) return; + try { + await _channel.invokeMethod('removeRemotePeer', { + 'peerIndex': peerIndex, + }); + } on PlatformException catch (error) { + final failure = _platformFailure(error); + _fail(failure); + throw failure; + } + } + @override Future stop() async { if (_disposed || diff --git a/mobile/lib/shared/huddle/huddle_session.dart b/mobile/lib/shared/huddle/huddle_session.dart index b49af484c..ada04c7ce 100644 --- a/mobile/lib/shared/huddle/huddle_session.dart +++ b/mobile/lib/shared/huddle/huddle_session.dart @@ -1,4 +1,5 @@ import 'dart:async'; +import 'dart:collection'; import 'package:flutter/foundation.dart'; import 'package:hooks_riverpod/hooks_riverpod.dart'; @@ -175,7 +176,11 @@ final class HuddleSessionNotifier extends Notifier { var _generation = 0; var _receivedFrames = 0; var _sentFrames = 0; - Future _playbackTail = Future.value(); + static const _playbackQueueCapacityPerPeer = 10; + + final Map> _playbackQueues = {}; + Future? _playbackDrain; + int? _lastPlaybackPeerIndex; final Map _speakerTimers = {}; final Map _pendingSpeakerLevels = {}; Timer? _speakerLevelFlushTimer; @@ -344,12 +349,12 @@ final class HuddleSessionNotifier extends Notifier { state.phase == HuddleSessionPhase.leaving) { return; } - _generation += 1; + final generation = ++_generation; state = state.copyWith(phase: HuddleSessionPhase.leaving, error: null); try { await _disposeResources(); } finally { - state = HuddleSessionState.idle; + if (_isCurrent(generation)) state = HuddleSessionState.idle; } } @@ -394,14 +399,7 @@ final class HuddleSessionNotifier extends Notifier { _recordSpeaker(frame, transport, generation); _receivedFrames += 1; _emitStatsIfNeeded(sent: false); - _playbackTail = _playbackTail.then((_) => media.playRemoteFrame(frame)); - unawaited( - _playbackTail.catchError((Object error) async { - if (_isCurrent(generation)) { - await _fail(_messageFor(error), generation); - } - }), - ); + _enqueuePlayback(media, frame, generation); }); _transportStateSubscription = transport.states.listen((transportState) { if (!_isCurrent(generation)) return; @@ -434,8 +432,19 @@ final class HuddleSessionNotifier extends Notifier { }); _peerEventSubscription = transport.peerEvents.listen((event) { if (!_isCurrent(generation)) return; - if (event.type == HuddlePeerEventType.left) { + if (event.type == HuddlePeerEventType.left || + event.type == HuddlePeerEventType.replaced) { _clearSpeaker(event.peer.pubkey); + _clearPlaybackPeer(event.peer.peerIndex); + unawaited( + media.removeRemotePeer(event.peer.peerIndex).catchError(( + Object error, + ) async { + if (_isCurrent(generation)) { + await _fail(_messageFor(error), generation); + } + }), + ); } }); _transportIssueSubscription = transport.issues.listen((issue) { @@ -443,6 +452,60 @@ final class HuddleSessionNotifier extends Notifier { }); } + void _enqueuePlayback( + HuddleMedia media, + HuddleRemoteAudioFrame frame, + int generation, + ) { + final queue = _playbackQueues.putIfAbsent(frame.peerIndex, Queue.new); + if (queue.length == _playbackQueueCapacityPerPeer) { + queue.removeFirst(); + } + queue.addLast(frame); + if (_playbackDrain != null) return; + + final drain = _drainPlayback(media, generation); + _playbackDrain = drain; + unawaited( + drain.whenComplete(() { + if (identical(_playbackDrain, drain)) _playbackDrain = null; + }), + ); + } + + Future _drainPlayback(HuddleMedia media, int generation) async { + try { + while (_isCurrent(generation)) { + final peerIndexes = _playbackQueues.keys.toList()..sort(); + if (peerIndexes.isEmpty) return; + final next = _lastPlaybackPeerIndex == null + ? 0 + : peerIndexes.indexWhere( + (index) => index > _lastPlaybackPeerIndex!, + ); + final start = next < 0 ? 0 : next; + HuddleRemoteAudioFrame? frame; + for (var offset = 0; offset < peerIndexes.length; offset++) { + final index = peerIndexes[(start + offset) % peerIndexes.length]; + final queue = _playbackQueues[index]; + if (queue != null && queue.isNotEmpty) { + frame = queue.removeFirst(); + _lastPlaybackPeerIndex = index; + break; + } + } + if (frame == null) return; + await media.playRemoteFrame(frame); + } + } catch (error) { + if (_isCurrent(generation)) await _fail(_messageFor(error), generation); + } + } + + void _clearPlaybackPeer(int peerIndex) { + _playbackQueues.remove(peerIndex)?.clear(); + } + List _participantPubkeys(HuddleTransportState transportState) { final peers = transportState.peers.values.toList() ..sort((left, right) => left.peerIndex.compareTo(right.peerIndex)); @@ -562,13 +625,22 @@ final class HuddleSessionNotifier extends Notifier { Future _fail(String message, int generation) async { if (!_isCurrent(generation)) return; - _generation += 1; - state = state.copyWith( + final failureGeneration = ++_generation; + final failedState = state.copyWith( phase: HuddleSessionPhase.failed, isMuted: true, error: message, ); - await _disposeResources(); + // Failure teardown follows the same local-first contract as an explicit + // hangup. Publish the failed state only after capture and transport are + // closed, so lifecycle observers cannot begin relay work while native + // audio is still active. + try { + await _disposeResources(); + } catch (error) { + debugPrint('Unable to fully release failed Huddle resources: $error'); + } + if (_isCurrent(failureGeneration)) state = failedState; } Future _disposeResources() async { @@ -627,7 +699,9 @@ final class HuddleSessionNotifier extends Notifier { }), ); } - _playbackTail = Future.value(); + _playbackQueues.clear(); + _playbackDrain = null; + _lastPlaybackPeerIndex = null; if (failure != null) { Error.throwWithStackTrace(failure, failureStackTrace!); } diff --git a/mobile/lib/shared/huddle/huddle_transport.dart b/mobile/lib/shared/huddle/huddle_transport.dart index b60a8558f..2d6a406b3 100644 --- a/mobile/lib/shared/huddle/huddle_transport.dart +++ b/mobile/lib/shared/huddle/huddle_transport.dart @@ -67,26 +67,33 @@ final class HuddlePeer { const HuddlePeer({required this.pubkey, required this.peerIndex}); } -enum HuddlePeerEventType { joined, left } +enum HuddlePeerEventType { joined, left, replaced } @immutable final class HuddlePeerEvent { final HuddlePeerEventType type; final HuddlePeer peer; + final HuddlePeer? replacement; - const HuddlePeerEvent({required this.type, required this.peer}); + const HuddlePeerEvent({ + required this.type, + required this.peer, + this.replacement, + }); } @immutable final class HuddleTransportState { final HuddleTransportPhase phase; final int? localPeerIndex; + final int? rosterRevision; final UnmodifiableMapView peers; final HuddleTransportError? error; HuddleTransportState({ required this.phase, this.localPeerIndex, + this.rosterRevision, Map peers = const {}, this.error, }) : peers = UnmodifiableMapView(Map.from(peers)); @@ -128,6 +135,9 @@ final class HuddleTransport implements HuddleTransportClient { StreamController.broadcast(sync: true); final StreamController _audioController = StreamController.broadcast(sync: true); + static const _audioIngressCapacity = 50; + final Queue _audioIngress = Queue(); + Timer? _audioIngressTimer; final StreamController _peerController = StreamController.broadcast(sync: true); final StreamController _issueController = @@ -306,6 +316,9 @@ final class HuddleTransport implements HuddleTransportClient { if (_disposed) return; await disconnect(); _disposed = true; + _audioIngressTimer?.cancel(); + _audioIngressTimer = null; + _audioIngress.clear(); await _stateController.close(); await _audioController.close(); await _peerController.close(); @@ -343,6 +356,8 @@ final class HuddleTransport implements HuddleTransportClient { _handleJoined(message, generation); case 'left': _handleLeft(message, generation); + case 'roster': + _handleRoster(message, generation); case 'error': _handleRelayError(message, generation); default: @@ -406,42 +421,85 @@ final class HuddleTransport implements HuddleTransportClient { return; } + final initialAdmission = + _state.phase == HuddleTransportPhase.authenticating; + final revision = message['revision']; + if (revision != null && (revision is! int || revision < 0)) { + _handleProtocolProblem('Malformed Huddle joined revision.', generation); + return; + } + final currentRevision = _state.rosterRevision; + final staleAdmission = + initialAdmission && + currentRevision != null && + revision is int && + revision < currentRevision; final peers = Map.from(_state.peers); + final replaced = peers[peerIndex]; final snapshot = message['peers']; if (snapshot is List) { - peers.clear(); - for (final candidate in snapshot) { - if (candidate is! Map) continue; - final candidatePubkey = candidate['pubkey']; - final candidateIndex = candidate['peer_index']; - if (candidatePubkey is String && - candidatePubkey.isNotEmpty && - candidateIndex is int && - candidateIndex >= 0 && - candidateIndex <= 255) { - peers[candidateIndex] = HuddlePeer( - pubkey: candidatePubkey, - peerIndex: candidateIndex, - ); - } + final parsed = _parsePeerSnapshot(snapshot); + if (parsed == null) { + _handleProtocolProblem('Malformed Huddle joined roster.', generation); + return; + } + if (initialAdmission && !staleAdmission) { + peers + ..clear() + ..addAll(parsed); + } else if (!initialAdmission && revision == null) { + // Legacy, unrevisioned relays used this field as a best-effort roster + // snapshot. Revisioned joined controls are deltas; merging their + // payload before inspecting the occupied slot would hide index reuse. + peers.addAll(parsed); + } + } + if (!initialAdmission && + !_acceptDeltaRevision(revision as int?, generation)) { + return; + } + if (initialAdmission && revision == null) { + final selfPresent = peers.values.any( + (candidate) => candidate.pubkey.toLowerCase() == pubkey.toLowerCase(), + ); + if (!selfPresent) { + peers[peerIndex] = HuddlePeer(pubkey: pubkey, peerIndex: peerIndex); } } final peer = HuddlePeer(pubkey: pubkey, peerIndex: peerIndex); - peers[peerIndex] = peer; + if (!initialAdmission || revision == null) { + peers[peerIndex] = peer; + } + if (!initialAdmission && + replaced != null && + replaced.pubkey != peer.pubkey) { + _purgeAudioIngress({peerIndex}); + } - final initialAdmission = - _state.phase == HuddleTransportPhase.authenticating; _emitState( HuddleTransportState( phase: HuddleTransportPhase.connected, localPeerIndex: initialAdmission ? peerIndex : _state.localPeerIndex, + rosterRevision: initialAdmission + ? (staleAdmission ? currentRevision : revision as int?) + : (revision as int? ?? _state.rosterRevision), peers: peers, ), ); if (!initialAdmission) { - _peerController.add( - HuddlePeerEvent(type: HuddlePeerEventType.joined, peer: peer), - ); + if (replaced != null && replaced.pubkey != peer.pubkey) { + _peerController.add( + HuddlePeerEvent( + type: HuddlePeerEventType.replaced, + peer: replaced, + replacement: peer, + ), + ); + } else { + _peerController.add( + HuddlePeerEvent(type: HuddlePeerEventType.joined, peer: peer), + ); + } } _handshakeTimer?.cancel(); _handshakeTimer = null; @@ -457,22 +515,31 @@ final class HuddleTransport implements HuddleTransportClient { } final pubkey = message['pubkey']; final peerIndex = message['peer_index']; + final revision = message['revision']; if (pubkey is! String || peerIndex is! int || peerIndex < 0 || - peerIndex > 255) { + peerIndex > 255 || + (revision != null && (revision is! int || revision < 0))) { _handleProtocolProblem('Malformed Huddle left message.', generation); return; } + if (!_acceptDeltaRevision(revision as int?, generation)) return; + final existing = _state.peers[peerIndex]; + if (existing != null && existing.pubkey != pubkey) { + _requestRosterResync(generation); + return; + } final peers = Map.from(_state.peers); - final removed = - peers.remove(peerIndex) ?? - HuddlePeer(pubkey: pubkey, peerIndex: peerIndex); + final removed = peers.remove(peerIndex); + if (removed == null) return; + _purgeAudioIngress({peerIndex}); _emitState( HuddleTransportState( phase: HuddleTransportPhase.connected, localPeerIndex: _state.localPeerIndex, + rosterRevision: revision ?? _state.rosterRevision, peers: peers, ), ); @@ -481,6 +548,141 @@ final class HuddleTransport implements HuddleTransportClient { ); } + void _handleRoster(Map message, int generation) { + if (_state.phase == HuddleTransportPhase.authenticating) { + _handleInitialRoster(message, generation); + return; + } + if (_state.phase != HuddleTransportPhase.connected) { + _handleProtocolProblem('Unexpected Huddle roster message.', generation); + return; + } + final revision = message['revision']; + final snapshot = message['peers']; + if (revision is! int || revision < 0 || snapshot is! List) { + _handleProtocolProblem('Malformed Huddle roster message.', generation); + return; + } + final previousRevision = _state.rosterRevision; + if (previousRevision != null && revision <= previousRevision) return; + final peers = _parsePeerSnapshot(snapshot); + if (peers == null) { + _handleProtocolProblem('Malformed Huddle roster message.', generation); + return; + } + + final previous = _state.peers; + final removedOrReplacedIndices = {}; + for (final entry in previous.entries) { + final replacement = peers[entry.key]; + if (replacement == null || replacement.pubkey != entry.value.pubkey) { + removedOrReplacedIndices.add(entry.key); + } + } + _purgeAudioIngress(removedOrReplacedIndices); + for (final entry in previous.entries) { + final replacement = peers[entry.key]; + if (replacement == null) { + _peerController.add( + HuddlePeerEvent(type: HuddlePeerEventType.left, peer: entry.value), + ); + } else if (replacement.pubkey != entry.value.pubkey) { + _peerController.add( + HuddlePeerEvent( + type: HuddlePeerEventType.replaced, + peer: entry.value, + replacement: replacement, + ), + ); + } + } + for (final entry in peers.entries) { + if (!previous.containsKey(entry.key)) { + _peerController.add( + HuddlePeerEvent(type: HuddlePeerEventType.joined, peer: entry.value), + ); + } + } + _emitState( + HuddleTransportState( + phase: HuddleTransportPhase.connected, + localPeerIndex: _state.localPeerIndex, + rosterRevision: revision, + peers: peers, + ), + ); + } + + void _handleInitialRoster(Map message, int generation) { + final revision = message['revision']; + final snapshot = message['peers']; + if (revision is! int || revision < 0 || snapshot is! List) { + _handleProtocolProblem('Malformed Huddle roster message.', generation); + return; + } + final peers = _parsePeerSnapshot(snapshot); + if (peers == null) { + _handleProtocolProblem('Malformed Huddle roster message.', generation); + return; + } + // Mesh owners may repair a snapshot-to-delta race before the admission + // message reaches this client. Retain that newer authoritative snapshot; + // the following joined admission is allowed to establish localPeerIndex + // without replacing it with older state. + _emitState( + HuddleTransportState( + phase: HuddleTransportPhase.authenticating, + rosterRevision: revision, + peers: peers, + ), + ); + } + + bool _acceptDeltaRevision(int? revision, int generation) { + final current = _state.rosterRevision; + // Single-node relays do not revision their legacy joined/left controls. + // Once a revisioned snapshot/delta is observed, every later delta must be + // exactly sequential; gaps are repaired by an authoritative snapshot. + if (revision == null) return current == null; + if (current == null || revision == current + 1) return true; + if (revision <= current) return false; + _requestRosterResync(generation); + return false; + } + + void _requestRosterResync(int generation) { + if (!_isCurrent(generation) || + _state.phase != HuddleTransportPhase.connected) { + return; + } + _fail( + const HuddleTransportError( + code: HuddleTransportErrorCode.protocolViolation, + message: 'Huddle roster revision gap; reconnecting for a fresh roster.', + ), + generation, + ); + } + + Map? _parsePeerSnapshot(List snapshot) { + final peers = {}; + for (final candidate in snapshot) { + if (candidate is! Map) return null; + final pubkey = candidate['pubkey']; + final peerIndex = candidate['peer_index']; + if (pubkey is! String || + pubkey.isEmpty || + peerIndex is! int || + peerIndex < 0 || + peerIndex > 255 || + peers.containsKey(peerIndex)) { + return null; + } + peers[peerIndex] = HuddlePeer(pubkey: pubkey, peerIndex: peerIndex); + } + return peers; + } + void _handleRelayError(Map message, int generation) { final relayCode = message['code'] is String ? message['code'] as String @@ -507,7 +709,15 @@ final class HuddleTransport implements HuddleTransportClient { return; } try { - _audioController.add(HuddleWireV2.decodeRelayFrame(bytes)); + final frame = HuddleWireV2.decodeRelayFrame(bytes); + // The authoritative control roster owns the routing table. Packets from + // absent/recycled slots are stale and must never allocate playback state. + if (!_state.peers.containsKey(frame.peerIndex)) return; + if (_audioIngress.length == _audioIngressCapacity) { + _audioIngress.removeFirst(); + } + _audioIngress.addLast(frame); + _audioIngressTimer ??= Timer(Duration.zero, _drainAudioIngress); } on HuddleWireException catch (error) { _issueController.add( HuddleTransportError( @@ -519,6 +729,32 @@ final class HuddleTransport implements HuddleTransportClient { } } + void _purgeAudioIngress(Set peerIndices) { + if (peerIndices.isEmpty || _audioIngress.isEmpty) return; + _audioIngress.removeWhere((frame) => peerIndices.contains(frame.peerIndex)); + } + + void _drainAudioIngress() { + _audioIngressTimer = null; + if (_disposed) { + _audioIngress.clear(); + return; + } + while (_audioIngress.isNotEmpty) { + final frame = _audioIngress.removeFirst(); + // Control and media are handled by the same socket callback, but queued + // media drains on later event-loop turns. Revalidate here so a removal + // or index replacement wins over every frame queued before that control. + if (_state.peers.containsKey(frame.peerIndex)) { + _audioController.add(frame); + break; + } + } + if (_audioIngress.isNotEmpty) { + _audioIngressTimer = Timer(Duration.zero, _drainAudioIngress); + } + } + void _handleProtocolProblem(String message, int generation) { final error = HuddleTransportError( code: HuddleTransportErrorCode.protocolViolation, @@ -565,6 +801,7 @@ final class HuddleTransport implements HuddleTransportClient { HuddleTransportState( phase: HuddleTransportPhase.failed, localPeerIndex: _state.localPeerIndex, + rosterRevision: _state.rosterRevision, peers: _state.peers, error: error, ), diff --git a/mobile/lib/shared/relay/relay_session.dart b/mobile/lib/shared/relay/relay_session.dart index d6c094d82..c861b706d 100644 --- a/mobile/lib/shared/relay/relay_session.dart +++ b/mobile/lib/shared/relay/relay_session.dart @@ -139,6 +139,7 @@ class RelaySessionNotifier extends Notifier { bool _hasConnectedOnce = false; int _connectionGeneration = 0; final Map _visibleChannelsByOwner = {}; + final Map Function()> _beforePauseCallbacks = {}; bool _socketConnected = false; bool _closedRetryReplayScheduled = false; @@ -398,11 +399,30 @@ class RelaySessionNotifier extends Notifier { await _connect(config); } + /// Registers work that must settle before the background grace disconnect. + void Function() registerBeforePause(Future Function() callback) { + final owner = Object(); + _beforePauseCallbacks[owner] = callback; + return () => _beforePauseCallbacks.remove(owner); + } + /// Called by the app lifecycle provider when the app goes to background. void onAppPaused() { _backgroundedAt = _now(); _backgroundGraceTimer?.cancel(); - _backgroundGraceTimer = Timer(_backgroundGraceDuration, _pauseNow); + _backgroundGraceTimer = Timer(_backgroundGraceDuration, () { + unawaited(_pauseAfterCallbacks()); + }); + } + + Future _pauseAfterCallbacks() async { + final callbacks = _beforePauseCallbacks.values.toList(); + try { + await Future.wait(callbacks.map((callback) => callback())); + } catch (error) { + debugPrint('Background cleanup failed: $error'); + } + if (_backgroundedAt != null) _pauseNow(); } void _pauseNow() { @@ -917,6 +937,7 @@ class RelaySessionNotifier extends Notifier { void _dispose() { _disposed = true; + _beforePauseCallbacks.clear(); _connectionGeneration++; _reconnectTimer?.cancel(); _flushTimer?.cancel(); diff --git a/mobile/test/features/channels/channel_detail_page_test.dart b/mobile/test/features/channels/channel_detail_page_test.dart index b39c5202a..a6ac6eebe 100644 --- a/mobile/test/features/channels/channel_detail_page_test.dart +++ b/mobile/test/features/channels/channel_detail_page_test.dart @@ -1,4 +1,5 @@ import 'dart:async'; +import 'dart:collection'; import 'dart:convert'; import 'package:flutter/foundation.dart'; @@ -295,6 +296,7 @@ Widget _buildTestable({ huddleHumanCountProvider.overrideWithValue(huddleHumanCountLoader), if (huddleCurrentPubkey != null) currentPubkeyProvider.overrideWith((ref) => huddleCurrentPubkey), + appLifecycleProvider.overrideWith(_TestAppLifecycleNotifier.new), // Compose bar drafts persist through SharedPreferences. savedPrefsProvider.overrideWithValue(_testPrefs), ], @@ -4459,7 +4461,7 @@ void main() { }); testWidgets( - 'top-right call end closes and releases media before member lookup', + 'duplicate hangup still completes the admitted lifecycle after local teardown', (tester) async { final now = DateTime.now().millisecondsSinceEpoch ~/ 1000; final media = _HuddleTestMedia(); @@ -4497,6 +4499,12 @@ void main() { await tester.tap(find.widgetWithText(FilledButton, 'Join')); await tester.pumpAndSettle(); await tester.tap(find.byKey(const ValueKey('huddle-leave'))); + final huddleContainer = ProviderScope.containerOf( + tester.element(find.byType(MobileHuddleShell)), + ); + unawaited( + huddleContainer.read(mobileHuddleControllerProvider.notifier).leave(), + ); for (var attempt = 0; attempt < 20; attempt++) { await tester.pump(); } @@ -4520,6 +4528,77 @@ void main() { }, ); + testWidgets( + 'creator end publish superseded by rejoin cannot archive new admission', + (tester) async { + final now = DateTime.now().millisecondsSinceEpoch ~/ 1000; + final endPublishGate = Completer(); + final media = Queue<_HuddleTestMedia>.of([ + _HuddleTestMedia(), + _HuddleTestMedia(), + ]); + final transports = Queue<_HuddleTestTransport>.of([ + _HuddleTestTransport(), + _HuddleTestTransport(), + ]); + final relaySession = _ReconnectingRelaySession( + huddleEndPublishGate: endPublishGate.future, + ); + String? archivedChannelId; + + await tester.pumpWidget( + _buildTestable( + messages: [ + _huddleMsg( + id: 'stale-creator-end', + kind: EventKind.huddleStarted, + pubkey: 'self', + createdAt: now, + ), + ], + users: const {'self': UserProfile(pubkey: 'self')}, + relayConfigNotifier: _HuddleRelayConfigNotifier(), + relaySessionNotifier: relaySession, + huddleCurrentPubkey: 'self', + huddleMediaFactory: media.removeFirst, + huddleTransportFactory: (_) => transports.removeFirst(), + createChannelActions: (ref) => _FakeChannelActions( + ref, + onArchiveChannel: (channelId) async => + archivedChannelId = channelId, + ), + ), + ); + await tester.pumpAndSettle(); + await tester.tap(find.widgetWithText(FilledButton, 'Join')); + await tester.pumpAndSettle(); + + final container = ProviderScope.containerOf( + tester.element(find.byType(MobileHuddleShell)), + ); + final controller = container.read( + mobileHuddleControllerProvider.notifier, + ); + final staleEnd = controller.end(); + await relaySession.huddleEndPublishStarted.future; + final rejoin = controller.join( + parentChannelId: _channelId, + ephemeralChannelId: _huddleChannelId, + startedBy: 'self', + startedEventId: 'stale-creator-end', + ); + endPublishGate.complete(); + await staleEnd; + await rejoin; + await tester.pump(); + + expect(relaySession.publishedKinds, contains(EventKind.huddleEnded)); + expect(archivedChannelId, isNull); + await tester.pumpWidget(const SizedBox.shrink()); + await tester.pump(const Duration(milliseconds: 100)); + }, + ); + testWidgets('disables Join after the matching huddle end event', ( tester, ) async { @@ -8572,7 +8651,16 @@ class _ErrorMessagesNotifier extends ChannelMessagesNotifier { AsyncError('Connection failed', StackTrace.current); } +class _TestAppLifecycleNotifier extends AppLifecycleNotifier { + @override + AppLifecycleState build() => AppLifecycleState.resumed; +} + class _ReconnectingRelaySession extends RelaySessionNotifier { + _ReconnectingRelaySession({this.huddleEndPublishGate}); + + final Future? huddleEndPublishGate; + final huddleEndPublishStarted = Completer(); final List publishedKinds = []; @override @@ -8591,6 +8679,14 @@ class _ReconnectingRelaySession extends RelaySessionNotifier { Duration timeout = const Duration(seconds: 8), }) async { publishedKinds.add(event.kind); + if (event.kind == EventKind.huddleEnded) { + if (huddleEndPublishGate case final gate?) { + if (!huddleEndPublishStarted.isCompleted) { + huddleEndPublishStarted.complete(); + } + await gate; + } + } return event; } @@ -8799,6 +8895,10 @@ class _HuddleRelayConfigNotifier extends RelayConfigNotifier { } final class _HuddleTestMedia implements HuddleMedia { + _HuddleTestMedia({this.stopGate}); + + final Future? stopGate; + final stopStarted = Completer(); final _states = StreamController.broadcast(sync: true); final _localFrames = StreamController.broadcast( sync: true, @@ -8885,8 +8985,13 @@ final class _HuddleTestMedia implements HuddleMedia { @override Future playRemoteFrame(HuddleRemoteAudioFrame frame) async {} + @override + Future removeRemotePeer(int peerIndex) async {} + @override Future stop() async { + if (!stopStarted.isCompleted) stopStarted.complete(); + if (stopGate case final gate?) await gate; _emit( HuddleMediaState( phase: HuddleMediaPhase.stopped, diff --git a/mobile/test/shared/huddle/huddle_session_test.dart b/mobile/test/shared/huddle/huddle_session_test.dart index a2efc7b9d..cd89b860f 100644 --- a/mobile/test/shared/huddle/huddle_session_test.dart +++ b/mobile/test/shared/huddle/huddle_session_test.dart @@ -1,4 +1,5 @@ import 'dart:async'; +import 'dart:collection'; import 'dart:typed_data'; import 'package:buzz/shared/huddle/huddle.dart'; @@ -106,6 +107,59 @@ void main() { await controller.leave(); }); + test('bounds playback per peer, drops oldest, and drains fairly', () async { + final media = _FakeMedia(blockPlayback: true); + final transport = _FakeTransport(); + final container = ProviderContainer( + overrides: [ + huddleMediaFactoryProvider.overrideWithValue(() => media), + huddleTransportFactoryProvider.overrideWithValue((_) => transport), + ], + ); + addTearDown(container.dispose); + await container.read(huddleSessionProvider.notifier).join(_parameters()); + + for (var sequence = 0; sequence < 20; sequence++) { + transport.emitRemote(_remoteFrame(peerIndex: 1, sequence: sequence)); + } + transport.emitRemote(_remoteFrame(peerIndex: 2, sequence: 100)); + await Future.delayed(Duration.zero); + expect(media.playedFrames.single.header.sequence, 0); + + media.releaseNextPlayback(); + await _waitUntil(() => media.playedFrames.length == 2); + expect(media.playedFrames[1].peerIndex, 2); + media.releaseNextPlayback(); + await _waitUntil(() => media.playedFrames.length == 3); + expect(media.playedFrames[2].header.sequence, 10); + }); + + test('peer replacement clears queued and native playback first', () async { + final media = _FakeMedia(blockPlayback: true); + final transport = _FakeTransport(); + final container = ProviderContainer( + overrides: [ + huddleMediaFactoryProvider.overrideWithValue(() => media), + huddleTransportFactoryProvider.overrideWithValue((_) => transport), + ], + ); + addTearDown(container.dispose); + await container.read(huddleSessionProvider.notifier).join(_parameters()); + + transport.emitRemote(_remoteFrame(peerIndex: 1, sequence: 1)); + transport.emitRemote(_remoteFrame(peerIndex: 1, sequence: 2)); + transport.emitReplacement( + const HuddlePeer(pubkey: 'desktop', peerIndex: 1), + const HuddlePeer(pubkey: 'new', peerIndex: 1), + ); + await Future.delayed(Duration.zero); + expect(media.removedPeers, [1]); + + media.releaseNextPlayback(); + await Future.delayed(Duration.zero); + expect(media.playedFrames.map((frame) => frame.header.sequence), [1]); + }); + test('keeps media alive through a bounded transport reconnect', () async { final media = _FakeMedia(); final transport = _FakeTransport(); @@ -173,6 +227,41 @@ void main() { expect(transport.connectCalls, 0); }); + test( + 'native failure releases media and transport before publishing failure', + () async { + final release = Completer(); + final media = _FakeMedia(disposeGate: release.future); + final transport = _FakeTransport(); + final container = ProviderContainer( + overrides: [ + huddleMediaFactoryProvider.overrideWithValue(() => media), + huddleTransportFactoryProvider.overrideWithValue((_) => transport), + ], + ); + addTearDown(container.dispose); + + await container.read(huddleSessionProvider.notifier).join(_parameters()); + media.emitFailure(); + await Future.delayed(Duration.zero); + + expect(media.disposeCalls, 1); + expect(transport.disposeCalls, 0); + expect( + container.read(huddleSessionProvider).phase, + HuddleSessionPhase.connected, + ); + + release.complete(); + await _waitUntil( + () => + container.read(huddleSessionProvider).phase == + HuddleSessionPhase.failed, + ); + expect(transport.disposeCalls, 1); + }, + ); + test( 'retains admission evidence after an established transport fails', () async { @@ -204,6 +293,20 @@ void main() { ); } +HuddleRemoteAudioFrame _remoteFrame({ + required int peerIndex, + required int sequence, +}) => HuddleRemoteAudioFrame( + peerIndex: peerIndex, + header: HuddleAudioHeader( + sequence: sequence, + timestamp48k: sequence * 960, + levelDbov: -20, + flags: 0, + ), + opusPayload: Uint8List.fromList([sequence & 0xff]), +); + HuddleConnectionParameters _parameters() => HuddleConnectionParameters( relayWebSocketUrl: 'wss://buzz.example', nsec: _privateKey, @@ -212,18 +315,27 @@ HuddleConnectionParameters _parameters() => HuddleConnectionParameters( ); final class _FakeMedia implements HuddleMedia { - _FakeMedia({this.permission = HuddleMicrophonePermission.granted}); + _FakeMedia({ + this.permission = HuddleMicrophonePermission.granted, + this.blockPlayback = false, + this.disposeGate, + }); final HuddleMicrophonePermission permission; + final bool blockPlayback; + final Future? disposeGate; final _states = StreamController.broadcast(sync: true); final _localFrames = StreamController.broadcast( sync: true, ); final List playedFrames = []; + final List removedPeers = []; + final Queue> _playbackCompleters = Queue(); HuddleMediaState _state = const HuddleMediaState( phase: HuddleMediaPhase.idle, ); var startCalls = 0; + var disposeCalls = 0; void emitLocal(HuddleLocalAudioFrame frame) => _localFrames.add(frame); @@ -238,6 +350,18 @@ final class _FakeMedia implements HuddleMedia { _states.add(_state); } + void emitFailure() { + _state = HuddleMediaState( + phase: HuddleMediaPhase.failed, + capabilities: _state.capabilities, + error: const HuddleMediaError( + code: HuddleMediaErrorCode.platformFailure, + message: 'native playback failed', + ), + ); + _states.add(_state); + } + @override HuddleMediaState get state => _state; @@ -312,6 +436,18 @@ final class _FakeMedia implements HuddleMedia { @override Future playRemoteFrame(HuddleRemoteAudioFrame frame) async { playedFrames.add(frame); + if (blockPlayback) { + final completer = Completer(); + _playbackCompleters.addLast(completer); + await completer.future; + } + } + + void releaseNextPlayback() => _playbackCompleters.removeFirst().complete(); + + @override + Future removeRemotePeer(int peerIndex) async { + removedPeers.add(peerIndex); } @override @@ -324,7 +460,11 @@ final class _FakeMedia implements HuddleMedia { } @override - Future dispose() => stop(); + Future dispose() async { + disposeCalls += 1; + if (disposeGate case final gate?) await gate; + await stop(); + } } final class _FakeTransport implements HuddleTransportClient { @@ -336,10 +476,21 @@ final class _FakeTransport implements HuddleTransportClient { final _issues = StreamController.broadcast(sync: true); final List sentFrames = []; var connectCalls = 0; + var disposeCalls = 0; HuddleTransportState _state = HuddleTransportState.idle(); void emitRemote(HuddleRemoteAudioFrame frame) => _remoteFrames.add(frame); + void emitReplacement(HuddlePeer peer, HuddlePeer replacement) { + _peerEvents.add( + HuddlePeerEvent( + type: HuddlePeerEventType.replaced, + peer: peer, + replacement: replacement, + ), + ); + } + void emitUnexpectedFailure() { _state = HuddleTransportState( phase: HuddleTransportPhase.failed, @@ -398,7 +549,10 @@ final class _FakeTransport implements HuddleTransportClient { } @override - Future dispose() => disconnect(); + Future dispose() async { + disposeCalls += 1; + await disconnect(); + } } Future _waitUntil(bool Function() predicate) async { diff --git a/mobile/test/shared/huddle/huddle_transport_test.dart b/mobile/test/shared/huddle/huddle_transport_test.dart index 95a2fd279..fcf96b709 100644 --- a/mobile/test/shared/huddle/huddle_transport_test.dart +++ b/mobile/test/shared/huddle/huddle_transport_test.dart @@ -53,7 +53,13 @@ void main() { final channel = _ControlledWebSocketChannel(); final transport = _transport(channel); addTearDown(transport.dispose); - await _connect(channel, transport); + await _connect( + channel, + transport, + peers: const [ + {'pubkey': 'remote', 'peer_index': 4}, + ], + ); final inbound = expectLater( transport.remoteAudioFrames, @@ -107,6 +113,211 @@ void main() { expect(transport.state.phase, HuddleTransportPhase.connected); }); + test( + 'revisioned roster is authoritative and deltas are sequential', + () async { + final channel = _ControlledWebSocketChannel(); + final transport = _transport(channel); + addTearDown(transport.dispose); + await _connect( + channel, + transport, + revision: 4, + peers: const [ + {'pubkey': 'desktop', 'peer_index': 1}, + ], + ); + expect(transport.state.peers.keys, [1]); + expect(transport.state.rosterRevision, 4); + + channel.emitText( + jsonEncode({ + 'type': 'joined', + 'revision': 5, + 'pubkey': 'agent', + 'peer_index': 2, + 'peers': [ + {'pubkey': 'agent', 'peer_index': 2}, + ], + }), + ); + await _waitForRevision(transport, 5); + expect(transport.state.peers.keys, containsAll([1, 2])); + + channel.emitText( + jsonEncode({ + 'type': 'roster', + 'revision': 6, + 'peers': [ + {'pubkey': 'replacement', 'peer_index': 2}, + ], + }), + ); + await _waitForRevision(transport, 6); + expect(transport.state.peers.keys, [2]); + expect(transport.state.peers[2]?.pubkey, 'replacement'); + + channel.emitText( + jsonEncode({'type': 'roster', 'revision': 5, 'peers': const []}), + ); + await Future.delayed(Duration.zero); + expect(transport.state.rosterRevision, 6); + expect(transport.state.peers[2]?.pubkey, 'replacement'); + }, + ); + + test( + 'newer roster arriving before admission remains authoritative', + () async { + final channel = _ControlledWebSocketChannel(); + final transport = _transport(channel); + addTearDown(transport.dispose); + + final connect = transport.connect(); + await _waitForPhase(transport, HuddleTransportPhase.awaitingChallenge); + channel.emitText(jsonEncode({'type': 'challenge', 'challenge': 'abc'})); + await _waitForPhase(transport, HuddleTransportPhase.authenticating); + channel.emitText( + jsonEncode({ + 'type': 'roster', + 'revision': 6, + 'peers': [ + {'pubkey': 'new', 'peer_index': 4}, + ], + }), + ); + channel.emitText( + jsonEncode({ + 'type': 'joined', + 'revision': 5, + 'pubkey': 'self', + 'peer_index': 3, + 'peers': [ + {'pubkey': 'old', 'peer_index': 1}, + ], + }), + ); + await connect; + + expect(transport.state.localPeerIndex, 3); + expect(transport.state.rosterRevision, 6); + expect(transport.state.peers.keys, [4]); + expect(transport.state.peers[4]?.pubkey, 'new'); + }, + ); + + test('revision gap reconnects without applying the gap', () async { + final channel = _ControlledWebSocketChannel(); + final transport = _transport(channel); + addTearDown(transport.dispose); + await _connect( + channel, + transport, + revision: 2, + peers: const [ + {'pubkey': 'desktop', 'peer_index': 1}, + ], + ); + + channel.emitText( + jsonEncode({ + 'type': 'left', + 'revision': 4, + 'pubkey': 'desktop', + 'peer_index': 1, + }), + ); + await _waitForPhase(transport, HuddleTransportPhase.failed); + + expect(transport.state.rosterRevision, 2); + expect(transport.state.peers[1]?.pubkey, 'desktop'); + expect(transport.state.error?.message, contains('revision gap')); + }); + + test('revisioned admission does not inject an absent local user', () async { + final channel = _ControlledWebSocketChannel(); + final transport = _transport(channel); + addTearDown(transport.dispose); + await _connect( + channel, + transport, + revision: 8, + peers: const [ + {'pubkey': 'desktop', 'peer_index': 1}, + ], + ); + + expect(transport.state.localPeerIndex, 3); + expect(transport.state.peers.keys, [1]); + }); + + test('revisioned joined delta reports index reuse', () async { + final channel = _ControlledWebSocketChannel(); + final transport = _transport(channel); + addTearDown(transport.dispose); + await _connect( + channel, + transport, + revision: 1, + peers: const [ + {'pubkey': 'old', 'peer_index': 4}, + ], + ); + final events = []; + final subscription = transport.peerEvents.listen(events.add); + addTearDown(subscription.cancel); + + channel.emitText( + jsonEncode({ + 'type': 'joined', + 'revision': 2, + 'pubkey': 'new', + 'peer_index': 4, + 'peers': [ + {'pubkey': 'new', 'peer_index': 4}, + ], + }), + ); + await _waitForRevision(transport, 2); + + expect(events.single.type, HuddlePeerEventType.replaced); + expect(events.single.peer.pubkey, 'old'); + expect(events.single.replacement?.pubkey, 'new'); + expect(transport.state.peers[4]?.pubkey, 'new'); + }); + + test('index replacement is emitted before new media is admitted', () async { + final channel = _ControlledWebSocketChannel(); + final transport = _transport(channel); + addTearDown(transport.dispose); + await _connect( + channel, + transport, + revision: 1, + peers: const [ + {'pubkey': 'old', 'peer_index': 4}, + ], + ); + final events = []; + final subscription = transport.peerEvents.listen(events.add); + addTearDown(subscription.cancel); + + channel.emitText( + jsonEncode({ + 'type': 'roster', + 'revision': 2, + 'peers': [ + {'pubkey': 'new', 'peer_index': 4}, + ], + }), + ); + await _waitForRevision(transport, 2); + + expect(events.single.type, HuddlePeerEventType.replaced); + expect(events.single.peer.pubkey, 'old'); + expect(events.single.replacement?.pubkey, 'new'); + }); + test('exposes relay rejection code and failed lifecycle state', () async { final channel = _ControlledWebSocketChannel(); final transport = _transport(channel); @@ -143,6 +354,82 @@ void main() { expect(transport.state.phase, HuddleTransportPhase.failed); }); + test('bounds transport audio ingress', () async { + final channel = _ControlledWebSocketChannel(); + final transport = _transport(channel); + addTearDown(transport.dispose); + await _connect( + channel, + transport, + revision: 1, + peers: const [ + {'pubkey': 'remote', 'peer_index': 4}, + ], + ); + final received = []; + final subscription = transport.remoteAudioFrames.listen( + (frame) => received.add(frame.header.sequence), + ); + addTearDown(subscription.cancel); + + for (var sequence = 0; sequence < 100; sequence++) { + channel.emitBinary(_relayFrame(peerIndex: 4, sequence: sequence)); + } + await Future.delayed(const Duration(milliseconds: 20)); + + expect(received, hasLength(50)); + expect(received.first, 50); + expect(received.last, 99); + }); + + test( + 'authoritative left purges queued media before same-index rejoin', + () async { + final channel = _ControlledWebSocketChannel(); + final transport = _transport(channel); + addTearDown(transport.dispose); + await _connect( + channel, + transport, + revision: 1, + peers: const [ + {'pubkey': 'old', 'peer_index': 4}, + ], + ); + final received = []; + final subscription = transport.remoteAudioFrames.listen( + (frame) => received.add(frame.header.sequence), + ); + addTearDown(subscription.cancel); + + for (var sequence = 0; sequence < 50; sequence++) { + channel.emitBinary(_relayFrame(peerIndex: 4, sequence: sequence)); + } + channel.emitText( + jsonEncode({ + 'type': 'left', + 'revision': 2, + 'pubkey': 'old', + 'peer_index': 4, + }), + ); + channel.emitText( + jsonEncode({ + 'type': 'joined', + 'revision': 3, + 'pubkey': 'new', + 'peer_index': 4, + }), + ); + channel.emitBinary(_relayFrame(peerIndex: 4, sequence: 50)); + await _waitForRevision(transport, 3); + await Future.delayed(const Duration(milliseconds: 20)); + + expect(transport.state.peers[4]?.pubkey, 'new'); + expect(received, [50]); + }, + ); + test( 'intentional disconnect reaches disconnected without an error', () async { @@ -174,23 +461,39 @@ HuddleTransport _transport(_ControlledWebSocketChannel channel) => Future _connect( _ControlledWebSocketChannel channel, - HuddleTransport transport, -) async { + HuddleTransport transport, { + int? revision, + List> peers = const [], +}) async { final connect = transport.connect(); await _waitForPhase(transport, HuddleTransportPhase.awaitingChallenge); channel.emitText(jsonEncode({'type': 'challenge', 'challenge': 'abc'})); await _waitForPhase(transport, HuddleTransportPhase.authenticating); - channel.emitText( - jsonEncode({ - 'type': 'joined', - 'pubkey': 'self', - 'peer_index': 3, - 'peers': const [], - }), - ); + final admission = { + 'type': 'joined', + 'pubkey': 'self', + 'peer_index': 3, + 'peers': peers, + }; + if (revision != null) admission['revision'] = revision; + channel.emitText(jsonEncode(admission)); await connect; } +Uint8List _relayFrame({required int peerIndex, required int sequence}) { + final header = HuddleAudioHeader( + sequence: sequence, + timestamp48k: sequence * 960, + levelDbov: -30, + flags: 0, + ); + final clientFrame = HuddleWireV2.encodeClientFrame( + header, + Uint8List.fromList([sequence & 0xff]), + ); + return Uint8List.fromList([peerIndex, ...clientFrame]); +} + Future _waitForPhase( HuddleTransport transport, HuddleTransportPhase phase, @@ -202,6 +505,17 @@ Future _waitForPhase( fail('Transport never reached ${phase.name}; was ${transport.state.phase}.'); } +Future _waitForRevision(HuddleTransport transport, int revision) async { + for (var attempt = 0; attempt < 100; attempt++) { + if (transport.state.rosterRevision == revision) return; + await Future.delayed(const Duration(milliseconds: 1)); + } + fail( + 'Transport never reached roster revision $revision; ' + 'was ${transport.state.rosterRevision}.', + ); +} + final class _ControlledWebSocketChannel implements WebSocketChannel { final StreamController _streamController = StreamController(); final _RecordingWebSocketSink _sink = _RecordingWebSocketSink();