Harden mobile Huddle reliability

Co-authored-by: Kenny Lopez <klopez4212@gmail.com>
Co-authored-by: Carl <3c4caeafb646d23867f1c4832e68211d77e2561946171625f75c3ce1a3f2670f@buzz.block.builderlab.xyz>
Signed-off-by: Kenny Lopez <klopez4212@gmail.com>
This commit is contained in:
Kenny Lopez
2026-08-17 17:32:39 +01:00
committed by Carl
co-authored by Carl
parent b9ebeb3492
commit ddaaef418b
18 changed files with 1566 additions and 161 deletions
+119 -29
View File
@@ -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<serde_json::Value> = 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<serde_json::Value>, 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<WsMessage>,
ctrl_tx: mpsc::Sender<WsMessage>,
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;
+18 -13
View File
@@ -1292,14 +1292,16 @@ impl<D: HuddleDirectory + ?Sized> HuddleControlAcceptor<D> {
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<D: HuddleDirectory + ?Sized> HuddleControlAcceptor<D> {
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<D: HuddleDirectory + ?Sized> HuddleControlAcceptor<D> {
registered: &mut std::collections::HashMap<String, Uuid>,
) -> 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<D: HuddleDirectory + ?Sized> HuddleControlAcceptor<D> {
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.
+52 -36
View File
@@ -229,7 +229,16 @@ impl Room {
&self,
pubkey: String,
requested_version: u8,
) -> Result<(Uuid, u8, mpsc::Receiver<Bytes>, mpsc::Receiver<PeerCtrl>), AdmissionError> {
) -> Result<
(
Uuid,
u8,
mpsc::Receiver<Bytes>,
mpsc::Receiver<PeerCtrl>,
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<Bytes>, mpsc::Receiver<PeerCtrl>), AdmissionError> {
) -> Result<(Uuid, mpsc::Receiver<Bytes>, mpsc::Receiver<PeerCtrl>, 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<RosterDelta> {
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!(
@@ -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<Int, Activity>()
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<Activity> { 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<Int> = 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<HuddleRemoteOpusPacket>(PLAYBACK_QUEUE_CAPACITY)
private val playbackQueue = ArrayBlockingQueue<HuddlePlaybackCommand>(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<Int, PeerPlayback>()
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()
@@ -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()
@@ -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<Int>(), selector.indices())
selector.activate(7, -20)
selector.remove(7)
assertEquals(emptySet<Int>(), selector.indices())
assertNull(selector.activate(7, -20))
}
}
+96 -8
View File
@@ -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)
+36
View File
@@ -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]
+25
View File
@@ -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()
@@ -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;
@@ -70,16 +70,38 @@ final huddleHumanCountProvider = Provider<HuddleHumanCountLoader>((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<bool> {
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<bool> {
Future<void> 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<bool> {
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<bool> {
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<bool> {
currentPubkey.toLowerCase() == startedBy.toLowerCase(),
startedEventId: startedEventId,
);
final session = ref.read(huddleSessionProvider);
if (session.wasAdmitted) {
_admissionToken = _HuddleAdmissionToken(admissionEpoch);
}
}
Future<void>? _backgroundLeave;
Future<void> _leaveForBackground() =>
_backgroundLeave ??= leave().whenComplete(() => _backgroundLeave = null);
Future<void> 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<bool> {
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<bool> {
}
}
Future<void> _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<void> _finishLeaveLifecycle({
required _HuddleAdmissionToken admissionToken,
required String? parentChannelId,
required String backingChannelId,
required Future<int> 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<bool> {
ephemeralChannelId: backingChannelId,
);
} catch (_) {}
if (admissionToken.epoch != _admissionEpoch) return;
try {
await actions.archiveChannel(backingChannelId);
} catch (_) {}
@@ -235,16 +294,29 @@ final class MobileHuddleController extends Notifier<bool> {
++_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<bool> {
} 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;
}
@@ -237,6 +237,7 @@ abstract interface class HuddleMedia {
Future<void> setMuted(bool muted);
Future<void> setSpeakerEnabled(bool enabled);
Future<void> playRemoteFrame(HuddleRemoteAudioFrame frame);
Future<void> removeRemotePeer(int peerIndex);
Future<void> stop();
Future<void> dispose();
}
@@ -534,6 +535,21 @@ final class MethodChannelHuddleMedia implements HuddleMedia {
}
}
@override
Future<void> removeRemotePeer(int peerIndex) async {
_ensureNotDisposed();
if (_state.phase != HuddleMediaPhase.active) return;
try {
await _channel.invokeMethod<void>('removeRemotePeer', {
'peerIndex': peerIndex,
});
} on PlatformException catch (error) {
final failure = _platformFailure(error);
_fail(failure);
throw failure;
}
}
@override
Future<void> stop() async {
if (_disposed ||
+90 -16
View File
@@ -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<HuddleSessionState> {
var _generation = 0;
var _receivedFrames = 0;
var _sentFrames = 0;
Future<void> _playbackTail = Future<void>.value();
static const _playbackQueueCapacityPerPeer = 10;
final Map<int, Queue<HuddleRemoteAudioFrame>> _playbackQueues = {};
Future<void>? _playbackDrain;
int? _lastPlaybackPeerIndex;
final Map<String, Timer> _speakerTimers = {};
final Map<String, double> _pendingSpeakerLevels = {};
Timer? _speakerLevelFlushTimer;
@@ -344,12 +349,12 @@ final class HuddleSessionNotifier extends Notifier<HuddleSessionState> {
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<HuddleSessionState> {
_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<HuddleSessionState> {
});
_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<HuddleSessionState> {
});
}
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<void> _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<String> _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<HuddleSessionState> {
Future<void> _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<void> _disposeResources() async {
@@ -627,7 +699,9 @@ final class HuddleSessionNotifier extends Notifier<HuddleSessionState> {
}),
);
}
_playbackTail = Future<void>.value();
_playbackQueues.clear();
_playbackDrain = null;
_lastPlaybackPeerIndex = null;
if (failure != null) {
Error.throwWithStackTrace(failure, failureStackTrace!);
}
+265 -28
View File
@@ -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<int, HuddlePeer> peers;
final HuddleTransportError? error;
HuddleTransportState({
required this.phase,
this.localPeerIndex,
this.rosterRevision,
Map<int, HuddlePeer> peers = const {},
this.error,
}) : peers = UnmodifiableMapView(Map<int, HuddlePeer>.from(peers));
@@ -128,6 +135,9 @@ final class HuddleTransport implements HuddleTransportClient {
StreamController.broadcast(sync: true);
final StreamController<HuddleRemoteAudioFrame> _audioController =
StreamController.broadcast(sync: true);
static const _audioIngressCapacity = 50;
final Queue<HuddleRemoteAudioFrame> _audioIngress = Queue();
Timer? _audioIngressTimer;
final StreamController<HuddlePeerEvent> _peerController =
StreamController.broadcast(sync: true);
final StreamController<HuddleTransportError> _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<int, HuddlePeer>.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<int, HuddlePeer>.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<String, dynamic> 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 = <int>{};
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<String, dynamic> 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<int, HuddlePeer>? _parsePeerSnapshot(List<dynamic> snapshot) {
final peers = <int, HuddlePeer>{};
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<String, dynamic> 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<int> 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,
),
+22 -1
View File
@@ -139,6 +139,7 @@ class RelaySessionNotifier extends Notifier<SessionState> {
bool _hasConnectedOnce = false;
int _connectionGeneration = 0;
final Map<Object, String> _visibleChannelsByOwner = {};
final Map<Object, Future<void> Function()> _beforePauseCallbacks = {};
bool _socketConnected = false;
bool _closedRetryReplayScheduled = false;
@@ -398,11 +399,30 @@ class RelaySessionNotifier extends Notifier<SessionState> {
await _connect(config);
}
/// Registers work that must settle before the background grace disconnect.
void Function() registerBeforePause(Future<void> 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<void> _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<SessionState> {
void _dispose() {
_disposed = true;
_beforePauseCallbacks.clear();
_connectionGeneration++;
_reconnectTimer?.cancel();
_flushTimer?.cancel();
@@ -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<void>();
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<void>? huddleEndPublishGate;
final huddleEndPublishStarted = Completer<void>();
final List<int> 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<void>? stopGate;
final stopStarted = Completer<void>();
final _states = StreamController<HuddleMediaState>.broadcast(sync: true);
final _localFrames = StreamController<HuddleLocalAudioFrame>.broadcast(
sync: true,
@@ -8885,8 +8985,13 @@ final class _HuddleTestMedia implements HuddleMedia {
@override
Future<void> playRemoteFrame(HuddleRemoteAudioFrame frame) async {}
@override
Future<void> removeRemotePeer(int peerIndex) async {}
@override
Future<void> stop() async {
if (!stopStarted.isCompleted) stopStarted.complete();
if (stopGate case final gate?) await gate;
_emit(
HuddleMediaState(
phase: HuddleMediaPhase.stopped,
@@ -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<void>.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<void>.delayed(Duration.zero);
expect(media.removedPeers, [1]);
media.releaseNextPlayback();
await Future<void>.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<void>();
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<void>.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<void>? disposeGate;
final _states = StreamController<HuddleMediaState>.broadcast(sync: true);
final _localFrames = StreamController<HuddleLocalAudioFrame>.broadcast(
sync: true,
);
final List<HuddleRemoteAudioFrame> playedFrames = [];
final List<int> removedPeers = [];
final Queue<Completer<void>> _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<void> playRemoteFrame(HuddleRemoteAudioFrame frame) async {
playedFrames.add(frame);
if (blockPlayback) {
final completer = Completer<void>();
_playbackCompleters.addLast(completer);
await completer.future;
}
}
void releaseNextPlayback() => _playbackCompleters.removeFirst().complete();
@override
Future<void> removeRemotePeer(int peerIndex) async {
removedPeers.add(peerIndex);
}
@override
@@ -324,7 +460,11 @@ final class _FakeMedia implements HuddleMedia {
}
@override
Future<void> dispose() => stop();
Future<void> 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<HuddleTransportError>.broadcast(sync: true);
final List<HuddleLocalAudioFrame> 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<void> dispose() => disconnect();
Future<void> dispose() async {
disposeCalls += 1;
await disconnect();
}
}
Future<void> _waitUntil(bool Function() predicate) async {
@@ -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<void>.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 = <HuddlePeerEvent>[];
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 = <HuddlePeerEvent>[];
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 = <int>[];
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<void>.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 = <int>[];
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<void>.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<void> _connect(
_ControlledWebSocketChannel channel,
HuddleTransport transport,
) async {
HuddleTransport transport, {
int? revision,
List<Map<String, Object>> 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 = <String, Object>{
'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<void> _waitForPhase(
HuddleTransport transport,
HuddleTransportPhase phase,
@@ -202,6 +505,17 @@ Future<void> _waitForPhase(
fail('Transport never reached ${phase.name}; was ${transport.state.phase}.');
}
Future<void> _waitForRevision(HuddleTransport transport, int revision) async {
for (var attempt = 0; attempt < 100; attempt++) {
if (transport.state.rosterRevision == revision) return;
await Future<void>.delayed(const Duration(milliseconds: 1));
}
fail(
'Transport never reached roster revision $revision; '
'was ${transport.state.rosterRevision}.',
);
}
final class _ControlledWebSocketChannel implements WebSocketChannel {
final StreamController<dynamic> _streamController = StreamController();
final _RecordingWebSocketSink _sink = _RecordingWebSocketSink();