fix(huddles): close lifecycle race conditions

Co-authored-by: Kenny Lopez <klopez4212@gmail.com>
Signed-off-by: Kenny Lopez <klopez4212@gmail.com>
This commit is contained in:
Princess Donut
2026-08-17 22:01:53 +01:00
co-authored by Kenny Lopez
parent 88ed4212c7
commit 5efcaada9c
12 changed files with 467 additions and 77 deletions
+45 -25
View File
@@ -47,7 +47,7 @@ use dashmap::DashMap;
use serde::{Deserialize, Serialize};
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use tracing::debug;
use tracing::{debug, warn};
use uuid::Uuid;
use super::mesh::spawn_remote_peer_sink;
@@ -1293,18 +1293,7 @@ impl<D: HuddleDirectory + ?Sized> HuddleControlAcceptor<D> {
.get(CommunityId::from_uuid(community_id), session_id)
}) {
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",
"revision": delta.revision,
"pubkey": left.pubkey,
"peer_index": left.peer_index,
})
.to_string(),
);
broadcast_peer_left(&room, delta, session_id);
}
}
}
@@ -1355,18 +1344,7 @@ impl<D: HuddleDirectory + ?Sized> HuddleControlAcceptor<D> {
}) {
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",
"revision": delta.revision,
"pubkey": left.pubkey,
"peer_index": left.peer_index,
})
.to_string(),
);
broadcast_peer_left(&room, delta, session_id);
}
}
}
@@ -1414,6 +1392,33 @@ impl<D: HuddleDirectory + ?Sized> HuddleControlAcceptor<D> {
}
}
fn broadcast_peer_left(room: &Room, delta: RoomRosterDelta, session_id: Uuid) {
let Some(left) = peer_left_control(delta, session_id) else {
return;
};
room.broadcast_control(left);
}
fn peer_left_control(delta: RoomRosterDelta, session_id: Uuid) -> Option<String> {
let Some(left) = delta.left else {
warn!(
%session_id,
revision = delta.revision,
"mesh audio peer removal delta did not include the removed peer"
);
return None;
};
Some(
serde_json::json!({
"type": "left",
"revision": delta.revision,
"pubkey": left.pubkey,
"peer_index": left.peer_index,
})
.to_string(),
)
}
fn roster_snapshot(room: &Room) -> RosterSnapshot {
let snapshot = room.roster_snapshot();
RosterSnapshot {
@@ -1822,6 +1827,21 @@ mod tests {
CommunityId::from_uuid(Uuid::from_u128(0xC0FFEE))
}
#[test]
fn missing_peer_in_removal_delta_is_non_fatal() {
assert_eq!(
peer_left_control(
RoomRosterDelta {
revision: 7,
joined: None,
left: None,
},
Uuid::from_u128(42),
),
None,
);
}
/// Scripted directory: `owner_of` returns a queued lookup, `acquire`
/// returns a queued outcome, `validate` returns a queued result. Records
/// call counts so ordering can be asserted.
@@ -382,6 +382,7 @@ export function HuddleBar({
ephemeralChannelId: barState?.ephemeral_channel_id ?? null,
fallbackParticipants: barState?.participants ?? [],
preservedParticipants: barState?.agent_pubkeys ?? [],
huddleThreadEventId: barState?.huddle_thread_event_id ?? null,
});
const participantSpeakerLevels = React.useMemo(() => {
const levels = { ...speakerLevels };
@@ -25,6 +25,7 @@ const KIND_HUDDLE_ENDED = 48103;
type ActiveHuddle = {
ephemeralChannelId: string;
huddleThreadEventId: string | null;
participants: Set<string>;
};
@@ -105,6 +106,7 @@ export function HuddleIndicator({
endedChannels.delete(ephId);
huddle = {
ephemeralChannelId: ephId,
huddleThreadEventId: ev.id,
participants: new Set([ev.pubkey]),
};
break;
@@ -120,6 +122,7 @@ export function HuddleIndicator({
if (!huddle || ephId !== huddle.ephemeralChannelId) {
huddle = {
ephemeralChannelId: ephId,
huddleThreadEventId: null,
participants: new Set(),
};
}
@@ -135,6 +138,7 @@ export function HuddleIndicator({
if (!huddle || ephId !== huddle.ephemeralChannelId) {
huddle = {
ephemeralChannelId: ephId,
huddleThreadEventId: null,
participants: new Set(),
};
}
@@ -227,7 +231,11 @@ export function HuddleIndicator({
if (activeHuddle) {
setIsJoining(true);
void joinHuddle(channelId, activeHuddle.ephemeralChannelId)
void joinHuddle(
channelId,
activeHuddle.ephemeralChannelId,
activeHuddle.huddleThreadEventId ?? undefined,
)
.then(() => {
void queryClient.invalidateQueries({ queryKey: ["channels"] });
})
@@ -309,7 +317,11 @@ export function HuddleIndicator({
if (!activeHuddle || isJoining) return;
setIsJoining(true);
try {
await joinHuddle(channelId, activeHuddle.ephemeralChannelId);
await joinHuddle(
channelId,
activeHuddle.ephemeralChannelId,
activeHuddle.huddleThreadEventId ?? undefined,
);
// Refetch channels so the ephemeral channel appears in the sidebar.
void queryClient.invalidateQueries({ queryKey: ["channels"] });
} catch (e) {
@@ -22,6 +22,7 @@ type HuddleRosterState = {
agent_voice_settings: Record<string, HuddleAgentVoiceSettings>;
parent_channel_id: string | null;
ephemeral_channel_id: string | null;
huddle_thread_event_id: string | null;
};
function isVisible(state: HuddleRosterState | null) {
@@ -42,6 +43,7 @@ export function HuddleRoomHeader() {
ephemeralChannelId: state?.ephemeral_channel_id ?? null,
fallbackParticipants: state?.participants ?? [],
preservedParticipants: state?.agent_pubkeys ?? [],
huddleThreadEventId: state?.huddle_thread_event_id ?? null,
});
const participantSpeakerLevels = React.useMemo(() => {
const levels = { ...speakerLevels };
@@ -1,7 +1,20 @@
import assert from "node:assert/strict";
import test from "node:test";
import { reconstructHuddleParticipantRoster } from "./useHuddleParticipantRoster.ts";
import { reconstructHuddleParticipantRoster as reconstructRoster } from "./useHuddleParticipantRoster.ts";
function reconstructHuddleParticipantRoster(options) {
const events = [...options.events];
return reconstructRoster({
...options,
events,
relaySelfPubkey: "relay",
huddleThreadEventId:
options.huddleThreadEventId ??
events.find((candidate) => candidate.kind === 48100)?.id ??
null,
});
}
const room = "huddle-room";
@@ -131,6 +144,74 @@ test("orders same-second leave then rejoin by roster revision", () => {
);
});
test("orders same-second unrevisioned remote leave before revised rejoin", () => {
const events = [
event({
id: "joined",
kind: 48101,
participant: "mobile",
createdAt: 2,
rosterRevision: 3,
}),
event({
id: "left",
kind: 48102,
participant: "mobile",
createdAt: 2,
}),
];
assert.deepEqual(
reconstructHuddleParticipantRoster({
ephemeralChannelId: room,
events,
fallbackParticipants: ["mobile"],
}),
["mobile"],
);
});
test("rejects lifecycle events without relay or creator provenance", () => {
const started = event({
id: "started",
kind: 48100,
pubkey: "creator",
createdAt: 1,
});
const events = [
started,
event({
id: "forged-left",
kind: 48102,
participant: "mobile",
pubkey: "attacker",
createdAt: 2,
}),
event({
id: "forged-start",
kind: 48100,
pubkey: "attacker",
createdAt: 3,
}),
event({
id: "forged-end",
kind: 48103,
pubkey: "attacker",
createdAt: 4,
}),
];
assert.deepEqual(
reconstructHuddleParticipantRoster({
ephemeralChannelId: room,
events,
fallbackParticipants: ["mobile"],
huddleThreadEventId: started.id,
}),
["creator"],
);
});
test("an ended lifecycle does not retain stale membership", () => {
const events = [
event({ id: "started", kind: 48100, pubkey: "desktop", createdAt: 1 }),
@@ -1,5 +1,6 @@
import * as React from "react";
import { useRelaySelfQuery } from "@/features/moderation/hooks";
import { relayClient } from "@/shared/api/relayClient";
import type { RelayEvent } from "@/shared/api/types";
import {
@@ -14,6 +15,8 @@ type ParticipantRosterOptions = {
events: Iterable<RelayEvent>;
fallbackParticipants?: readonly string[];
preservedParticipants?: readonly (string | null | undefined)[];
relaySelfPubkey?: string | null;
huddleThreadEventId?: string | null;
};
function normalizedPubkey(value: string | null | undefined): string | null {
@@ -64,8 +67,21 @@ function compareLifecycleEvents(left: RelayEvent, right: RelayEvent): number {
const phaseOrder = lifecyclePhase(left.kind) - lifecyclePhase(right.kind);
if (phaseOrder !== 0) return phaseOrder;
const leftParticipant = lifecycleParticipant(left);
const rightParticipant = lifecycleParticipant(right);
const sameParticipant =
leftParticipant !== null && leftParticipant === rightParticipant;
const leftRevision = lifecycleContent(left).rosterRevision;
const rightRevision = lifecycleContent(right).rosterRevision;
if (sameParticipant && leftRevision !== rightRevision) {
if (leftRevision === null) {
return left.kind === KIND_HUDDLE_PARTICIPANT_LEFT ? -1 : 1;
}
if (rightRevision === null) {
return right.kind === KIND_HUDDLE_PARTICIPANT_LEFT ? 1 : -1;
}
return leftRevision - rightRevision;
}
if (
leftRevision !== null &&
rightRevision !== null &&
@@ -83,6 +99,43 @@ function lifecycleParticipant(event: RelayEvent): string | null {
);
}
function authenticatedLifecycleEvents({
events,
relaySelfPubkey,
huddleThreadEventId,
}: Pick<
ParticipantRosterOptions,
"events" | "relaySelfPubkey" | "huddleThreadEventId"
>): RelayEvent[] {
const relaySelf = normalizedPubkey(relaySelfPubkey);
const lifecycle = [...events];
const started = huddleThreadEventId
? (lifecycle.find(
(event) =>
event.kind === KIND_HUDDLE_STARTED &&
event.id === huddleThreadEventId,
) ?? null)
: null;
const creator = started ? normalizedPubkey(started.pubkey) : null;
return lifecycle.filter((event) => {
if (
event.kind === KIND_HUDDLE_PARTICIPANT_JOINED ||
event.kind === KIND_HUDDLE_PARTICIPANT_LEFT
) {
return relaySelf !== null && normalizedPubkey(event.pubkey) === relaySelf;
}
if (event.kind === KIND_HUDDLE_STARTED) {
return started !== null && event.id === started.id;
}
if (event.kind === KIND_HUDDLE_ENDED) {
const signer = normalizedPubkey(event.pubkey);
return signer !== null && (signer === creator || signer === relaySelf);
}
return false;
});
}
/**
* Reconstruct the live media roster from relay-signed Huddle lifecycle events.
*
@@ -95,13 +148,19 @@ export function reconstructHuddleParticipantRoster({
events,
fallbackParticipants = [],
preservedParticipants = [],
relaySelfPubkey = null,
huddleThreadEventId = null,
}: ParticipantRosterOptions): string[] {
const participants = new Set(
fallbackParticipants
.map(normalizedPubkey)
.filter((pubkey): pubkey is string => pubkey !== null),
);
const sorted = [...events]
const sorted = authenticatedLifecycleEvents({
events,
relaySelfPubkey,
huddleThreadEventId,
})
.filter((event) => lifecycleChannelId(event) === ephemeralChannelId)
// Nostr timestamps have second precision and relay history arrives newest
// first. Session phases break start/end ties; owner roster revisions order
@@ -152,6 +211,7 @@ type UseHuddleParticipantRosterOptions = {
ephemeralChannelId: string | null;
fallbackParticipants: readonly string[];
preservedParticipants?: readonly (string | null | undefined)[];
huddleThreadEventId?: string | null;
};
/** Subscribe the active desktop roster to the canonical signed lifecycle. */
@@ -160,7 +220,9 @@ export function useHuddleParticipantRoster({
ephemeralChannelId,
fallbackParticipants,
preservedParticipants = [],
huddleThreadEventId = null,
}: UseHuddleParticipantRosterOptions): string[] {
const relaySelfPubkey = useRelaySelfQuery(Boolean(parentChannelId)).data;
const sessionKey =
parentChannelId && ephemeralChannelId
? `${parentChannelId}:${ephemeralChannelId}`
@@ -219,6 +281,8 @@ export function useHuddleParticipantRoster({
events: [],
fallbackParticipants,
preservedParticipants,
relaySelfPubkey,
huddleThreadEventId,
});
}
return reconstructHuddleParticipantRoster({
@@ -226,5 +290,7 @@ export function useHuddleParticipantRoster({
events,
fallbackParticipants,
preservedParticipants,
relaySelfPubkey,
huddleThreadEventId,
});
}
@@ -19,6 +19,7 @@ import java.util.ArrayDeque
import java.util.concurrent.ArrayBlockingQueue
import java.util.concurrent.TimeUnit
import java.util.concurrent.atomic.AtomicBoolean
import java.util.concurrent.atomic.AtomicInteger
import kotlin.concurrent.thread
import kotlin.math.log10
import kotlin.math.max
@@ -42,6 +43,14 @@ internal data class HuddleRemoteOpusPacket(
val opus: ByteArray,
)
internal data class HuddleTalkerSelection(
val accepted: Boolean,
val evictedPeerIndex: Int?,
) {
inline fun <T> allocateIfAccepted(allocate: () -> T): T? =
if (accepted) allocate() else null
}
/// Fixed-capacity active-talker selector. Roster membership stays in Dart;
/// native decoder/track resources exist only for peers that actually send.
internal class HuddleActiveTalkerSelector(
@@ -56,11 +65,11 @@ internal class HuddleActiveTalkerSelector(
private val active = mutableMapOf<Int, Activity>()
private var ordinal = 0L
fun activate(peerIndex: Int, levelDbov: Int): Int? {
fun activate(peerIndex: Int, levelDbov: Int): HuddleTalkerSelection {
ordinal += 1
if (active.containsKey(peerIndex)) {
active[peerIndex] = Activity(peerIndex, ordinal, levelDbov)
return null
return HuddleTalkerSelection(accepted = true, evictedPeerIndex = null)
}
val weakest = if (active.size >= capacity) {
active.values.minWithOrNull(
@@ -74,19 +83,19 @@ internal class HuddleActiveTalkerSelector(
// Equal-level packets (especially continuous -127 dBov silence) must
// not churn decoder/jitter state. A new peer earns a slot only when it
// is actually louder than the quietest selected peer.
if (weakest != null && levelDbov <= weakest.levelDbov) return null
if (weakest != null && levelDbov <= weakest.levelDbov) {
return HuddleTalkerSelection(accepted = false, evictedPeerIndex = null)
}
val evicted = weakest?.peerIndex
evicted?.let(active::remove)
active[peerIndex] = Activity(peerIndex, ordinal, levelDbov)
return evicted
return HuddleTalkerSelection(accepted = true, evictedPeerIndex = evicted)
}
fun remove(peerIndex: Int) {
active.remove(peerIndex)
}
fun contains(peerIndex: Int): Boolean = active.containsKey(peerIndex)
fun indices(): Set<Int> = active.keys.toSet()
}
@@ -165,6 +174,8 @@ internal class HuddleAudioEngine(
private val failureReported = AtomicBoolean(false)
private val playbackQueue = ArrayBlockingQueue<HuddlePlaybackCommand>(PLAYBACK_QUEUE_CAPACITY)
private val interrupted = AtomicBoolean(false)
private val interruptionEpoch = AtomicInteger(0)
private val interruptionLock = Object()
@Volatile
private var muted = false
@@ -222,6 +233,7 @@ internal class HuddleAudioEngine(
val record = audioRecord
if (value) {
interrupted.set(true)
interruptionEpoch.incrementAndGet()
// A focus loss must stop native capture, not merely discard frames
// after AudioRecord has already read them from the microphone.
if (record != null) runCatching { record.stop() }
@@ -229,6 +241,7 @@ internal class HuddleAudioEngine(
}
if (!running.get() || record == null) {
interrupted.set(false)
synchronized(interruptionLock) { interruptionLock.notifyAll() }
return
}
try {
@@ -237,6 +250,7 @@ internal class HuddleAudioEngine(
"Microphone recording did not resume."
}
interrupted.set(false)
synchronized(interruptionLock) { interruptionLock.notifyAll() }
} catch (error: Throwable) {
reportFailure("audio_resume_failed", error)
}
@@ -270,6 +284,7 @@ internal class HuddleAudioEngine(
fun stop() {
running.set(false)
synchronized(interruptionLock) { interruptionLock.notifyAll() }
runCatching { audioRecord?.stop() }
playbackThread?.interrupt()
runCatching { captureThread?.join(STOP_JOIN_TIMEOUT_MS) }
@@ -329,6 +344,18 @@ internal class HuddleAudioEngine(
try {
while (running.get()) {
// Snapshot before checking [interrupted]. If focus loss starts
// after that check, its epoch still marks the stopped read as
// expected; if it starts first, the loop waits below.
val readEpoch = interruptionEpoch.get()
if (interrupted.get()) {
synchronized(interruptionLock) {
while (running.get() && interrupted.get()) {
interruptionLock.wait(INTERRUPTION_WAIT_MS)
}
}
continue
}
val read = record.read(
readBuffer,
0,
@@ -336,6 +363,10 @@ internal class HuddleAudioEngine(
AudioRecord.READ_BLOCKING,
)
if (read < 0) {
// AudioRecord.stop() unblocks a pending read with an error.
// A focus interruption that began during this read owns that
// result even if focus was regained before we inspected it.
if (interruptionEpoch.get() != readEpoch) continue
throw IllegalStateException("AudioRecord read failed with code $read.")
}
var sourceOffset = 0
@@ -496,17 +527,20 @@ internal class HuddleAudioEngine(
when (command) {
is HuddlePlaybackCommand.Packet -> {
val packet = command.packet
val evicted = activeTalkers.activate(
val selection = activeTalkers.activate(
packet.peerIndex,
packet.levelDbov,
)
evicted?.let { playbacks.remove(it)?.release() }
if (!activeTalkers.contains(packet.peerIndex)) continue
val playback = playbacks[packet.peerIndex] ?: PeerPlayback(
packet.peerIndex,
).also {
playbacks[packet.peerIndex] = it
selection.evictedPeerIndex?.let {
playbacks.remove(it)?.release()
}
val playback = playbacks[packet.peerIndex]
?: selection.allocateIfAccepted {
PeerPlayback(packet.peerIndex).also {
playbacks[packet.peerIndex] = it
}
}
?: continue
playback.enqueue(packet)
}
is HuddlePlaybackCommand.RemovePeer -> {
@@ -650,6 +684,7 @@ internal class HuddleAudioEngine(
private const val JITTER_START_PACKETS = 3
private const val MAX_REMOTE_PEERS = 15
private const val STOP_JOIN_TIMEOUT_MS = 750L
private const val INTERRUPTION_WAIT_MS = 250L
private const val SILENT_LEVEL_DBOV = -127
private const val OPUS_SEEK_PREROLL_NS = 80_000_000L
private const val DIAGNOSTIC_EMIT_FRAME_INTERVAL = 50L
@@ -10,8 +10,18 @@ class HuddleActiveTalkerSelectorTest {
fun `sixteen senders retain exactly the active talker capacity`() {
val selector = HuddleActiveTalkerSelector(capacity = 15)
repeat(15) { peer -> assertNull(selector.activate(peer, -20)) }
assertNull(selector.activate(15, -20))
repeat(15) { peer ->
assertEquals(
HuddleTalkerSelection(accepted = true, evictedPeerIndex = null),
selector.activate(peer, -20),
)
}
val rejected = selector.activate(15, -20)
assertEquals(
HuddleTalkerSelection(accepted = false, evictedPeerIndex = null),
rejected,
)
assertNull(rejected.allocateIfAccepted { Any() })
assertEquals((0..14).toSet(), selector.indices())
}
@@ -22,10 +32,16 @@ class HuddleActiveTalkerSelectorTest {
selector.activate(4, -40)
selector.activate(2, -20)
selector.activate(4, -10)
assertEquals(2, selector.activate(9, -5))
assertEquals(
HuddleTalkerSelection(accepted = true, evictedPeerIndex = 2),
selector.activate(9, -5),
)
assertEquals(setOf(4, 9), selector.indices())
assertEquals(4, selector.activate(2, -1))
assertEquals(
HuddleTalkerSelection(accepted = true, evictedPeerIndex = 4),
selector.activate(2, -1),
)
assertEquals(setOf(2, 9), selector.indices())
}
@@ -34,7 +50,10 @@ class HuddleActiveTalkerSelectorTest {
val selector = HuddleActiveTalkerSelector(capacity = 1)
selector.activate(7, -30)
assertNull(selector.activate(3, -30))
assertEquals(
HuddleTalkerSelection(accepted = false, evictedPeerIndex = null),
selector.activate(3, -30),
)
assertEquals(setOf(7), selector.indices())
}
@@ -44,8 +63,11 @@ class HuddleActiveTalkerSelectorTest {
selector.activate(4, -127)
selector.activate(2, -127)
assertNull(selector.activate(9, -127))
assertFalse(selector.contains(9))
assertEquals(
HuddleTalkerSelection(accepted = false, evictedPeerIndex = null),
selector.activate(9, -127),
)
assertFalse(9 in selector.indices())
assertEquals(setOf(2, 4), selector.indices())
}
@@ -55,7 +77,10 @@ class HuddleActiveTalkerSelectorTest {
selector.activate(4, -40)
selector.activate(2, -20)
assertEquals(4, selector.activate(9, -30))
assertEquals(
HuddleTalkerSelection(accepted = true, evictedPeerIndex = 4),
selector.activate(9, -30),
)
assertEquals(setOf(2, 9), selector.indices())
}
@@ -87,6 +112,9 @@ class HuddleActiveTalkerSelectorTest {
selector.activate(7, -20)
selector.remove(7)
assertEquals(emptySet<Int>(), selector.indices())
assertNull(selector.activate(7, -20))
assertEquals(
HuddleTalkerSelection(accepted = true, evictedPeerIndex = null),
selector.activate(7, -20),
)
}
}
+27 -12
View File
@@ -310,6 +310,16 @@ enum HuddleAudioLevels {
}
}
struct HuddleTalkerSelection: Equatable {
let accepted: Bool
let evictedPeerIndex: Int?
func allocateIfAccepted<T>(_ allocate: () throws -> T) rethrows -> T? {
guard accepted else { return nil }
return try allocate()
}
}
/// Fixed-capacity active-talker selector. The UI roster is independent; native
/// decoder/player resources exist only for peers that actually send packets.
struct HuddleActiveTalkerSelector {
@@ -327,7 +337,7 @@ struct HuddleActiveTalkerSelector {
self.capacity = capacity
}
mutating func activate(peerIndex: Int, levelDbov: Int) -> Int? {
mutating func activate(peerIndex: Int, levelDbov: Int) -> HuddleTalkerSelection {
ordinal &+= 1
if active[peerIndex] != nil {
active[peerIndex] = Activity(
@@ -335,7 +345,7 @@ struct HuddleActiveTalkerSelector {
lastPacketOrdinal: ordinal,
levelDbov: levelDbov
)
return nil
return HuddleTalkerSelection(accepted: true, evictedPeerIndex: nil)
}
let weakest = active.count >= capacity
? active.values.min { left, right in
@@ -351,7 +361,9 @@ struct HuddleActiveTalkerSelector {
// Equal-level packets (especially continuous -127 dBov silence) must not
// churn decoder/jitter state. Admit a new peer only when it is louder than
// the quietest selected peer.
if let weakest, levelDbov <= weakest.levelDbov { return nil }
if let weakest, levelDbov <= weakest.levelDbov {
return HuddleTalkerSelection(accepted: false, evictedPeerIndex: nil)
}
let evicted = weakest?.peerIndex
if let evicted { active.removeValue(forKey: evicted) }
active[peerIndex] = Activity(
@@ -359,7 +371,7 @@ struct HuddleActiveTalkerSelector {
lastPacketOrdinal: ordinal,
levelDbov: levelDbov
)
return evicted
return HuddleTalkerSelection(accepted: true, evictedPeerIndex: evicted)
}
mutating func remove(_ peerIndex: Int) {
@@ -520,25 +532,28 @@ final class HuddleAudioEngine {
guard let self else { return }
defer { completeRemotePacket() }
do {
let evictedIndex = activeTalkers.activate(
let selection = activeTalkers.activate(
peerIndex: packet.peerIndex,
levelDbov: packet.levelDbov
)
if let evictedIndex,
if let evictedIndex = selection.evictedPeerIndex,
let evicted = peerPlaybacks.removeValue(forKey: evictedIndex)
{
evicted.release(from: audioEngine)
}
if activeTalkers.active[packet.peerIndex] == nil { return }
let playback: HuddlePeerPlayback
if let existing = peerPlaybacks[packet.peerIndex] {
playback = existing
} else {
playback = try HuddlePeerPlayback(
engine: audioEngine,
pcmFormat: pcmFormat
} else if let allocated = try selection.allocateIfAccepted({
try HuddlePeerPlayback(
engine: self.audioEngine,
pcmFormat: self.pcmFormat
)
peerPlaybacks[packet.peerIndex] = playback
}) {
playback = allocated
peerPlaybacks[packet.peerIndex] = allocated
} else {
return
}
playback.enqueue(packet)
} catch {
+49 -10
View File
@@ -11,33 +11,72 @@ class RunnerTests: XCTestCase {
var selector = HuddleActiveTalkerSelector(capacity: 15)
for peer in 0..<15 {
XCTAssertNil(selector.activate(peerIndex: peer, levelDbov: -20))
XCTAssertEqual(
selector.activate(peerIndex: peer, levelDbov: -20),
HuddleTalkerSelection(accepted: true, evictedPeerIndex: nil)
)
}
XCTAssertNil(selector.activate(peerIndex: 15, levelDbov: -20))
let rejected = selector.activate(peerIndex: 15, levelDbov: -20)
XCTAssertEqual(
rejected,
HuddleTalkerSelection(accepted: false, evictedPeerIndex: nil)
)
var allocationCount = 0
XCTAssertNil(rejected.allocateIfAccepted {
allocationCount += 1
return NSObject()
})
XCTAssertEqual(allocationCount, 0)
XCTAssertEqual(Set(selector.active.keys), Set(0..<15))
selector.remove(7)
XCTAssertFalse(selector.active.keys.contains(7))
XCTAssertNil(selector.activate(peerIndex: 7, levelDbov: -10))
XCTAssertEqual(
selector.activate(peerIndex: 7, levelDbov: -10),
HuddleTalkerSelection(accepted: true, evictedPeerIndex: nil)
)
XCTAssertTrue(selector.active.keys.contains(7))
}
func testHuddleActiveTalkerSelectorUsesRecentActivity() {
var selector = HuddleActiveTalkerSelector(capacity: 2)
XCTAssertNil(selector.activate(peerIndex: 4, levelDbov: -10))
XCTAssertNil(selector.activate(peerIndex: 2, levelDbov: -20))
XCTAssertNil(selector.activate(peerIndex: 4, levelDbov: -10))
XCTAssertEqual(selector.activate(peerIndex: 9, levelDbov: -5), 2)
XCTAssertEqual(
selector.activate(peerIndex: 4, levelDbov: -10),
HuddleTalkerSelection(accepted: true, evictedPeerIndex: nil)
)
XCTAssertEqual(
selector.activate(peerIndex: 2, levelDbov: -20),
HuddleTalkerSelection(accepted: true, evictedPeerIndex: nil)
)
XCTAssertEqual(
selector.activate(peerIndex: 4, levelDbov: -10),
HuddleTalkerSelection(accepted: true, evictedPeerIndex: nil)
)
XCTAssertEqual(
selector.activate(peerIndex: 9, levelDbov: -5),
HuddleTalkerSelection(accepted: true, evictedPeerIndex: 2)
)
XCTAssertEqual(Set(selector.active.keys), Set([4, 9]))
}
func testHuddleActiveTalkerSelectorDoesNotChurnAtEqualLevels() {
var selector = HuddleActiveTalkerSelector(capacity: 2)
XCTAssertNil(selector.activate(peerIndex: 4, levelDbov: -127))
XCTAssertNil(selector.activate(peerIndex: 2, levelDbov: -127))
XCTAssertNil(selector.activate(peerIndex: 9, levelDbov: -127))
XCTAssertEqual(
selector.activate(peerIndex: 4, levelDbov: -127),
HuddleTalkerSelection(accepted: true, evictedPeerIndex: nil)
)
XCTAssertEqual(
selector.activate(peerIndex: 2, levelDbov: -127),
HuddleTalkerSelection(accepted: true, evictedPeerIndex: nil)
)
let rejected = selector.activate(peerIndex: 9, levelDbov: -127)
XCTAssertEqual(
rejected,
HuddleTalkerSelection(accepted: false, evictedPeerIndex: nil)
)
XCTAssertNil(rejected.allocateIfAccepted { NSObject() })
XCTAssertNil(selector.active[9])
XCTAssertEqual(Set(selector.active.keys), Set([2, 4]))
}
@@ -124,7 +124,21 @@ final class MobileHuddleController extends Notifier<bool> {
return false;
}
Future<void> start({required String parentChannelId}) async {
Future<void>? _startFuture;
Future<void> start({required String parentChannelId}) {
final inFlight = _startFuture;
if (inFlight != null) return inFlight;
late final Future<void> startFuture;
startFuture = _start(parentChannelId: parentChannelId).whenComplete(() {
if (identical(_startFuture, startFuture)) _startFuture = null;
});
_startFuture = startFuture;
return startFuture;
}
Future<void> _start({required String parentChannelId}) async {
if (state) return;
final generation = ++_generation;
final admissionEpoch = ++_admissionEpoch;
@@ -217,8 +231,20 @@ final class MobileHuddleController extends Notifier<bool> {
_backgroundLeave ??= leave().whenComplete(() => _backgroundLeave = null);
Future<void> _leaveForTransition() async {
if (!ref.read(huddleSessionProvider).isInSession) return;
await _leaveForBackground();
final startFuture = _startFuture;
if (startFuture != null) {
++_generation;
state = false;
try {
await startFuture;
} catch (_) {
// Cancellation is observed by the in-flight start, whose rollback must
// settle before the relay is rebound to the next community.
}
}
if (ref.read(huddleSessionProvider).isInSession) {
await _leaveForBackground();
}
}
Future<void> leave() async {
@@ -36,6 +36,7 @@ import 'package:buzz/features/profile/profile_provider.dart';
import 'package:buzz/features/profile/user_cache_provider.dart';
import 'package:buzz/features/profile/user_profile.dart';
import 'package:buzz/features/profile/user_profile_sheet.dart';
import 'package:buzz/shared/community/community_provider.dart';
import 'package:buzz/shared/emoji/emoji_burst.dart';
import 'package:buzz/shared/mentions/agent_identity_provider.dart';
import 'package:buzz/shared/huddle/huddle.dart';
@@ -4581,6 +4582,57 @@ void main() {
},
);
testWidgets(
'community transition cancels and awaits an in-flight Huddle start',
(tester) async {
final createGate = Completer<void>();
final relaySession = _ReconnectingRelaySession(
huddleCreatePublishGate: createGate.future,
);
String? archivedChannelId;
await tester.pumpWidget(
_buildTestable(
messages: const [],
users: const {'self': UserProfile(pubkey: 'self')},
relayConfigNotifier: _HuddleRelayConfigNotifier(),
relaySessionNotifier: relaySession,
huddleCurrentPubkey: 'self',
huddleMediaFactory: _HuddleTestMedia.new,
huddleTransportFactory: (_) => _HuddleTestTransport(),
createChannelActions: (ref) => _FakeChannelActions(
ref,
onArchiveChannel: (channelId) async =>
archivedChannelId = channelId,
),
),
);
await tester.pumpAndSettle();
await tester.tap(find.byKey(const ValueKey('channel-huddle-button')));
await relaySession.huddleCreatePublishStarted.future;
var transitionCompleted = false;
final transition =
ProviderScope.containerOf(
tester.element(find.byType(MobileHuddleShell)),
).read(communityTransitionProvider).run().then((_) {
transitionCompleted = true;
});
await tester.pump();
expect(transitionCompleted, isFalse);
createGate.complete();
await transition;
expect(archivedChannelId, isNotNull);
expect(
relaySession.publishedKinds,
isNot(contains(EventKind.huddleStarted)),
);
await tester.pump(const Duration(seconds: 1));
},
);
testWidgets(
'creator end publish superseded by rejoin cannot archive new admission',
(tester) async {
@@ -8710,9 +8762,14 @@ class _TestAppLifecycleNotifier extends AppLifecycleNotifier {
}
class _ReconnectingRelaySession extends RelaySessionNotifier {
_ReconnectingRelaySession({this.huddleEndPublishGate});
_ReconnectingRelaySession({
this.huddleCreatePublishGate,
this.huddleEndPublishGate,
});
final Future<void>? huddleCreatePublishGate;
final Future<void>? huddleEndPublishGate;
final huddleCreatePublishStarted = Completer<void>();
final huddleEndPublishStarted = Completer<void>();
final List<int> publishedKinds = [];
@@ -8732,6 +8789,14 @@ class _ReconnectingRelaySession extends RelaySessionNotifier {
Duration timeout = const Duration(seconds: 8),
}) async {
publishedKinds.add(event.kind);
if (event.kind == 9007) {
if (huddleCreatePublishGate case final gate?) {
if (!huddleCreatePublishStarted.isCompleted) {
huddleCreatePublishStarted.complete();
}
await gate;
}
}
if (event.kind == EventKind.huddleEnded) {
if (huddleEndPublishGate case final gate?) {
if (!huddleEndPublishStarted.isCompleted) {