Address Huddle review feedback

Co-authored-by: Princess Donut <b238ea756dee4d98afa5883fc7f1de61eeabe65bf700e3a5a5a80db5e42e2c2b@buzz.block.builderlab.xyz>
Signed-off-by: Princess Donut <b238ea756dee4d98afa5883fc7f1de61eeabe65bf700e3a5a5a80db5e42e2c2b@buzz.block.builderlab.xyz>
This commit is contained in:
Princess Donut
2026-08-17 20:03:21 +01:00
parent 2368ed605a
commit 1c9180bda8
14 changed files with 387 additions and 36 deletions
@@ -70,6 +70,22 @@ test("ignores other rooms and preserves local and agent participants", () => {
);
});
test("preserves delivery order for same-second leave then rejoin", () => {
const events = [
event({ id: "left", kind: 48102, participant: "mobile", createdAt: 2 }),
event({ id: "joined", kind: 48101, participant: "mobile", createdAt: 2 }),
];
assert.deepEqual(
reconstructHuddleParticipantRoster({
ephemeralChannelId: room,
events,
fallbackParticipants: ["mobile"],
}),
["mobile"],
);
});
test("an ended lifecycle does not retain stale membership", () => {
const events = [
event({ id: "started", kind: 48100, pubkey: "desktop", createdAt: 1 }),
@@ -60,12 +60,9 @@ export function reconstructHuddleParticipantRoster({
);
const sorted = [...events]
.filter((event) => lifecycleChannelId(event) === ephemeralChannelId)
.sort(
(left, right) =>
left.created_at - right.created_at ||
left.kind - right.kind ||
left.id.localeCompare(right.id),
);
// Equal-second lifecycle events must retain relay delivery order: kind/id
// are not causal and can invert a leave followed by an immediate rejoin.
.sort((left, right) => left.created_at - right.created_at);
let ended = false;
for (const event of sorted) {
@@ -87,6 +87,39 @@ internal sealed interface HuddlePlaybackCommand {
data class RemovePeer(val peerIndex: Int) : HuddlePlaybackCommand
}
internal class HuddlePacketJitterQueue(
private val capacity: Int,
private val startPackets: Int,
) {
private val packets = mutableListOf<HuddleRemoteOpusPacket>()
private var lastDrainedSequence: Int? = null
private var started = false
fun enqueue(packet: HuddleRemoteOpusPacket) {
val lastDrained = lastDrainedSequence
if (lastDrained != null && !sequenceAfter(packet.sequence, lastDrained)) return
if (packets.any { it.sequence == packet.sequence }) return
val insertionIndex = packets.indexOfFirst { sequenceBefore(packet.sequence, it.sequence) }
if (insertionIndex < 0) packets.add(packet) else packets.add(insertionIndex, packet)
if (packets.size > capacity) packets.removeAt(0)
if (packets.size >= startPackets) started = true
}
fun drainOne(): HuddleRemoteOpusPacket? {
if (!started || packets.isEmpty()) return null
return packets.removeAt(0).also { lastDrainedSequence = it.sequence }
}
fun sequences(): List<Int> = packets.map { it.sequence }
private fun sequenceAfter(sequence: Int, reference: Int): Boolean {
val distance = (sequence - reference) and 0xffff
return distance in 1..0x7fff
}
private fun sequenceBefore(left: Int, right: Int): Boolean = sequenceAfter(right, left)
}
/** Aggregate, debug-only microphone telemetry. No PCM or encoded audio is retained. */
internal data class HuddleCaptureDiagnostics(
val frameCount: Long,
@@ -468,9 +501,10 @@ internal class HuddleAudioEngine(
private class PeerPlayback(val peerIndex: Int) {
private val decoder = createDecoder()
private val track = createAudioTrack()
private val jitterQueue = ArrayDeque<HuddleRemoteOpusPacket>()
private var playoutStarted = false
private var lastSequence: Int? = null
private val jitterQueue = HuddlePacketJitterQueue(
capacity = PER_PEER_JITTER_CAPACITY,
startPackets = JITTER_START_PACKETS,
)
init {
decoder.start()
@@ -478,18 +512,11 @@ internal class HuddleAudioEngine(
}
fun enqueue(packet: HuddleRemoteOpusPacket) {
if (lastSequence == packet.sequence) return
lastSequence = packet.sequence
if (jitterQueue.size == PER_PEER_JITTER_CAPACITY) {
jitterQueue.removeFirst()
}
jitterQueue.addLast(packet)
if (jitterQueue.size >= JITTER_START_PACKETS) playoutStarted = true
jitterQueue.enqueue(packet)
}
fun drainOne() {
if (!playoutStarted || jitterQueue.isEmpty()) return
decode(jitterQueue.removeFirst())
jitterQueue.drainOne()?.let(::decode)
}
private fun decode(packet: HuddleRemoteOpusPacket) {
@@ -37,6 +37,26 @@ class HuddleActiveTalkerSelectorTest {
assertEquals(setOf(3), selector.indices())
}
@Test
fun `jitter queue reorders packets and rejects stale duplicates`() {
val queue = HuddlePacketJitterQueue(capacity = 3, startPackets = 2)
queue.enqueue(packet(11))
queue.enqueue(packet(10))
assertEquals(listOf(10, 11), queue.sequences())
assertEquals(10, queue.drainOne()?.sequence)
queue.enqueue(packet(10))
queue.enqueue(packet(12))
assertEquals(listOf(11, 12), queue.sequences())
}
private fun packet(sequence: Int) = HuddleRemoteOpusPacket(
peerIndex = 1,
sequence = sequence,
timestamp48k = sequence.toLong() * 960,
levelDbov = -20,
opus = byteArrayOf(1),
)
@Test
fun `remove clears the slot without allocating roster-only peers`() {
val selector = HuddleActiveTalkerSelector(capacity = 2)
+51 -15
View File
@@ -17,6 +17,50 @@ struct HuddleRemoteOpusPacket {
let opus: Data
}
struct HuddlePacketJitterQueue {
private(set) var packets: [HuddleRemoteOpusPacket] = []
private var lastDrainedSequence: Int?
private var started = false
let capacity: Int
let startPackets: Int
init(capacity: Int, startPackets: Int) {
self.capacity = capacity
self.startPackets = startPackets
}
mutating func enqueue(_ packet: HuddleRemoteOpusPacket) {
if let last = lastDrainedSequence, !sequenceAfter(packet.sequence, last) {
return
}
guard !packets.contains(where: { $0.sequence == packet.sequence }) else {
return
}
let index = packets.firstIndex {
sequenceBefore(packet.sequence, $0.sequence)
} ?? packets.endIndex
packets.insert(packet, at: index)
if packets.count > capacity { packets.removeFirst() }
if packets.count >= startPackets { started = true }
}
mutating func drainOne() -> HuddleRemoteOpusPacket? {
guard started, !packets.isEmpty else { return nil }
let packet = packets.removeFirst()
lastDrainedSequence = packet.sequence
return packet
}
private func sequenceAfter(_ sequence: Int, _ reference: Int) -> Bool {
let distance = (sequence - reference) & 0xffff
return distance > 0 && distance <= 0x7fff
}
private func sequenceBefore(_ left: Int, _ right: Int) -> Bool {
sequenceAfter(right, left)
}
}
struct HuddleCaptureDiagnostics {
let frameCount: Int64
let rmsDbovHistogram: [String: Int64]
@@ -867,9 +911,10 @@ final class HuddleAudioEngine {
private final class HuddlePeerPlayback {
private let player = AVAudioPlayerNode()
private let decoder: HuddleOpusDecoder
private var jitterQueue: [HuddleRemoteOpusPacket] = []
private var playoutStarted = false
private var lastSequence: Int?
private var jitterQueue = HuddlePacketJitterQueue(
capacity: 10,
startPackets: 3
)
init(
engine: AVAudioEngine,
@@ -882,20 +927,11 @@ private final class HuddlePeerPlayback {
}
func enqueue(_ packet: HuddleRemoteOpusPacket) {
guard lastSequence != packet.sequence else { return }
lastSequence = packet.sequence
if jitterQueue.count == 10 {
jitterQueue.removeFirst()
}
jitterQueue.append(packet)
if jitterQueue.count >= 3 {
playoutStarted = true
}
jitterQueue.enqueue(packet)
}
func drainOne() throws {
guard playoutStarted, !jitterQueue.isEmpty else { return }
let packet = jitterQueue.removeFirst()
guard let packet = jitterQueue.drainOne() else { return }
let decoded = try decoder.decode(packet.opus)
player.scheduleBuffer(decoded)
if !player.isPlaying {
@@ -919,6 +955,6 @@ private final class HuddlePeerPlayback {
player.stop()
engine.disconnectNodeOutput(player)
engine.detach(player)
jitterQueue.removeAll()
jitterQueue = HuddlePacketJitterQueue(capacity: 10, startPackets: 3)
}
}
+21
View File
@@ -32,6 +32,27 @@ class RunnerTests: XCTestCase {
XCTAssertEqual(Set(selector.active.keys), Set([4, 9]))
}
func testHuddlePacketJitterQueueReordersAndRejectsStaleDuplicates() {
var queue = HuddlePacketJitterQueue(capacity: 3, startPackets: 2)
queue.enqueue(remotePacket(sequence: 11))
queue.enqueue(remotePacket(sequence: 10))
XCTAssertEqual(queue.packets.map(\.sequence), [10, 11])
XCTAssertEqual(queue.drainOne()?.sequence, 10)
queue.enqueue(remotePacket(sequence: 10))
queue.enqueue(remotePacket(sequence: 12))
XCTAssertEqual(queue.packets.map(\.sequence), [11, 12])
}
private func remotePacket(sequence: Int) -> HuddleRemoteOpusPacket {
HuddleRemoteOpusPacket(
peerIndex: 1,
sequence: sequence,
timestamp48k: Int64(sequence * 960),
levelDbov: -20,
opus: Data([1])
)
}
func testHuddleOpusCodecRoundTripsFixedV2Frame() throws {
let encoder = try HuddleOpusEncoder()
let decoder = try HuddleOpusDecoder()
@@ -155,6 +155,8 @@ class ChannelDetailPage extends HookConsumerWidget {
final detailsAsync = ref.watch(channelDetailsProvider(channel.id));
final channelsAsync = ref.watch(channelsProvider);
final messagesState = ref.watch(channelMessagesProvider(channel.id));
final huddleLifecycle =
ref.watch(huddleLifecycleProvider(channel.id)).value ?? const [];
final sessionStatus = ref.watch(relaySessionProvider).status;
final readState = ref.watch(readStateProvider);
final channelsNotifier = ref.read(channelsProvider.notifier);
@@ -361,7 +363,7 @@ class ChannelDetailPage extends HookConsumerWidget {
if (showsComposer)
_HuddleButton(
channel: resolvedChannel,
events: messagesState.value ?? const [],
events: [...messagesState.value ?? const [], ...huddleLifecycle],
),
if (_showsMembersAction(resolvedChannel))
_MembersButton(
@@ -3,6 +3,7 @@ import 'dart:async';
import 'package:flutter/widgets.dart';
import 'package:hooks_riverpod/hooks_riverpod.dart';
import '../../shared/community/community_provider.dart';
import '../../shared/huddle/huddle.dart';
import '../../shared/relay/relay.dart';
import 'channel_management_provider.dart';
@@ -70,6 +71,18 @@ final huddleHumanCountProvider = Provider<HuddleHumanCountLoader>((ref) {
};
});
/// Complete parent-channel Huddle lifecycle, independent of timeline paging.
final huddleLifecycleProvider = FutureProvider.family<List<NostrEvent>, String>(
(ref, channelId) async {
if (ref.watch(relaySessionProvider).status != SessionStatus.connected) {
return const [];
}
return ref
.read(relaySessionProvider.notifier)
.fetchHistory(NostrFilters.huddleLifecycle(channelId));
},
);
/// 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.
@@ -91,7 +104,11 @@ final class MobileHuddleController extends Notifier<bool> {
final unregisterBeforePause = ref
.read(relaySessionProvider.notifier)
.registerBeforePause(_leaveForBackground);
final unregisterBeforeCommunityTransition = ref
.read(communityTransitionProvider)
.register(_leaveForTransition);
ref.onDispose(unregisterBeforePause);
ref.onDispose(unregisterBeforeCommunityTransition);
ref.listen(huddleSessionProvider, (previous, next) {
if (previous?.wasAdmitted == true &&
next.phase == HuddleSessionPhase.failed) {
@@ -199,6 +216,11 @@ final class MobileHuddleController extends Notifier<bool> {
Future<void> _leaveForBackground() =>
_backgroundLeave ??= leave().whenComplete(() => _backgroundLeave = null);
Future<void> _leaveForTransition() async {
if (!ref.read(huddleSessionProvider).isInSession) return;
await _leaveForBackground();
}
Future<void> leave() async {
++_generation;
state = false;
@@ -54,6 +54,7 @@ class AuthNotifier extends AsyncNotifier<AuthState> {
/// Writes to storage directly to avoid circular dependency with community
/// providers.
Future<void> authenticateWithCommunity(Community community) async {
await ref.read(communityTransitionProvider).run();
final storage = ref.read(communityStorageProvider);
await storage.save(community);
await storage.saveActiveId(community.id);
@@ -68,6 +69,7 @@ class AuthNotifier extends AsyncNotifier<AuthState> {
}
Future<void> signOut() async {
await ref.read(communityTransitionProvider).run();
final storage = ref.read(communityStorageProvider);
final activeId = await storage.loadActiveId();
if (activeId != null) {
@@ -4,6 +4,24 @@ import '../auth/auth_provider.dart';
import 'community.dart';
import 'community_storage.dart';
final class CommunityTransitionCoordinator {
final Map<Object, Future<void> Function()> _callbacks = {};
void Function() register(Future<void> Function() callback) {
final owner = Object();
_callbacks[owner] = callback;
return () => _callbacks.remove(owner);
}
Future<void> run() async {
await Future.wait(_callbacks.values.map((callback) => callback()));
}
}
final communityTransitionProvider = Provider<CommunityTransitionCoordinator>(
(ref) => CommunityTransitionCoordinator(),
);
final communityStorageProvider = Provider<CommunityStorage>((ref) {
return CommunityStorage();
});
@@ -46,17 +64,23 @@ class CommunityListNotifier extends AsyncNotifier<List<Community>> {
Future<void> removeCommunity(String id) async {
final storage = ref.read(communityStorageProvider);
final activeId = await storage.loadActiveId();
if (activeId == id) {
await ref.read(communityTransitionProvider).run();
}
await storage.remove(id);
final current = state.value ?? [];
state = AsyncData(current.where((w) => w.id != id).toList());
// If we removed the active community, switch to another or sign out.
final activeId = await storage.loadActiveId();
if (activeId == id) {
final remaining = state.value ?? [];
if (remaining.isNotEmpty) {
await switchCommunity(remaining.first.id);
await storage.saveActiveId(remaining.first.id);
// Reassign list state so activeCommunityProvider picks up the new ID.
state = AsyncData([...remaining]);
ref.invalidate(authProvider);
} else {
await storage.clearActiveId();
// Invalidate auth so it re-evaluates against the now-empty storage
@@ -68,6 +92,9 @@ class CommunityListNotifier extends AsyncNotifier<List<Community>> {
Future<void> switchCommunity(String id) async {
final storage = ref.read(communityStorageProvider);
final activeId = await storage.loadActiveId();
if (activeId == id) return;
await ref.read(communityTransitionProvider).run();
await storage.saveActiveId(id);
// Reassign list state to trigger activeCommunityProvider (which watches
// communityListProvider.future) to rebuild and pick up the new active ID.
@@ -49,6 +49,15 @@ abstract final class NostrFilters {
until: until,
);
/// Huddle lifecycle state for one parent channel, independent of timeline paging.
static NostrFilter huddleLifecycle(String channelId) => NostrFilter(
kinds: [EventKind.huddleStarted, EventKind.huddleEnded],
tags: {
'#h': [channelId],
},
limit: 200,
);
/// Reactions (kind:7) on a specific event.
static NostrFilter reactions(String eventId) => NostrFilter(
kinds: [7],
@@ -209,6 +209,7 @@ Widget _buildTestable({
HuddleMediaFactory? huddleMediaFactory,
HuddleTransportFactory? huddleTransportFactory,
HuddleHumanCountLoader? huddleHumanCountLoader,
List<NostrEvent> huddleLifecycle = const [],
String? huddleCurrentPubkey,
http.Client? mediaClient,
}) {
@@ -294,6 +295,9 @@ Widget _buildTestable({
),
if (huddleHumanCountLoader != null)
huddleHumanCountProvider.overrideWithValue(huddleHumanCountLoader),
huddleLifecycleProvider(
_channelId,
).overrideWith((ref) async => huddleLifecycle),
if (huddleCurrentPubkey != null)
currentPubkeyProvider.overrideWith((ref) => huddleCurrentPubkey),
appLifecycleProvider.overrideWith(_TestAppLifecycleNotifier.new),
@@ -3205,6 +3209,28 @@ void main() {
);
});
testWidgets('top action discovers a Huddle outside the timeline window', (
tester,
) async {
final now = DateTime.now().millisecondsSinceEpoch ~/ 1000;
await tester.pumpWidget(
_buildTestable(
messages: const [],
huddleLifecycle: [
_huddleMsg(
id: 'off-window-huddle',
kind: EventKind.huddleStarted,
createdAt: now,
),
],
),
);
await tester.pumpAndSettle();
expect(find.byTooltip('Open Huddle'), findsOneWidget);
expect(find.text('Huddle in progress'), findsNothing);
});
testWidgets('offers Join for a recent desktop-started huddle', (
tester,
) async {
@@ -1,3 +1,5 @@
import 'dart:async';
import 'package:flutter_test/flutter_test.dart';
import 'package:hooks_riverpod/hooks_riverpod.dart';
import 'package:nostr/nostr.dart' as nostr;
@@ -9,6 +11,77 @@ import 'package:buzz/shared/community/community_storage.dart';
import '../community/community_storage_test.dart';
void main() {
test('waits for transition teardown before replacing credentials', () async {
final storage = CommunityStorage(secure: FakeSecureStorage());
final active = Community.create(
name: 'Active',
relayUrl: 'https://active.example',
nsec: nostr.Keys.generate().nsec,
);
final replacement = Community.create(
name: 'Replacement',
relayUrl: 'https://replacement.example',
nsec: nostr.Keys.generate().nsec,
);
await storage.save(active);
await storage.saveActiveId(active.id);
final container = ProviderContainer(
overrides: [communityStorageProvider.overrideWithValue(storage)],
);
addTearDown(container.dispose);
await container.read(authProvider.future);
final teardown = Completer<void>();
container.read(communityTransitionProvider).register(() => teardown.future);
final authenticating = container
.read(authProvider.notifier)
.authenticateWithCommunity(replacement);
await Future<void>.delayed(Duration.zero);
expect(await storage.loadActiveId(), active.id);
expect(
(await storage.loadAll()).map((community) => community.id),
isNot(contains(replacement.id)),
);
teardown.complete();
await authenticating;
expect(await storage.loadActiveId(), replacement.id);
expect(
(await storage.loadAll()).map((community) => community.id),
contains(replacement.id),
);
});
test('waits for transition teardown before removing credentials', () async {
final storage = CommunityStorage(secure: FakeSecureStorage());
final community = Community.create(
name: 'Active',
relayUrl: 'https://relay.example',
nsec: nostr.Keys.generate().nsec,
);
await storage.save(community);
await storage.saveActiveId(community.id);
final container = ProviderContainer(
overrides: [communityStorageProvider.overrideWithValue(storage)],
);
addTearDown(container.dispose);
await container.read(authProvider.future);
final teardown = Completer<void>();
container.read(communityTransitionProvider).register(() => teardown.future);
final signingOut = container.read(authProvider.notifier).signOut();
await Future<void>.delayed(Duration.zero);
expect(await storage.loadActiveId(), community.id);
final saved = await storage.loadAll();
expect(saved, hasLength(1));
expect(saved.single.id, community.id);
teardown.complete();
await signingOut;
expect(await storage.loadActiveId(), isNull);
expect(await storage.loadAll(), isEmpty);
});
test(
'removes an invalid saved community instead of authenticating',
() async {
@@ -1,3 +1,5 @@
import 'dart:async';
import 'package:flutter_test/flutter_test.dart';
import 'package:hooks_riverpod/hooks_riverpod.dart';
import 'package:buzz/shared/community/community.dart';
@@ -63,6 +65,47 @@ void main() {
expect(communities, isEmpty);
});
test(
'waits for transition teardown before removing active community',
() async {
container = createContainer();
await container.read(communityListProvider.future);
final ws1 = Community.create(
name: 'One',
relayUrl: 'https://one.example.com',
);
final ws2 = Community.create(
name: 'Two',
relayUrl: 'https://two.example.com',
);
final notifier = container.read(communityListProvider.notifier);
await notifier.addCommunity(ws1);
await notifier.addCommunity(ws2);
await notifier.switchCommunity(ws1.id);
final teardown = Completer<void>();
container
.read(communityTransitionProvider)
.register(() => teardown.future);
final removing = notifier.removeCommunity(ws1.id);
await Future<void>.delayed(Duration.zero);
expect(await communityStorage.loadActiveId(), ws1.id);
expect(
(await communityStorage.loadAll()).map((community) => community.id),
contains(ws1.id),
);
teardown.complete();
await removing;
expect(await communityStorage.loadActiveId(), ws2.id);
expect(
(await communityStorage.loadAll()).map((community) => community.id),
isNot(contains(ws1.id)),
);
},
);
test('renameCommunity updates name', () async {
container = createContainer();
await container.read(communityListProvider.future);
@@ -80,6 +123,36 @@ void main() {
expect(communities.first.name, 'Renamed');
});
test('waits for transition teardown before updating active ID', () async {
container = createContainer();
await container.read(communityListProvider.future);
final ws1 = Community.create(
name: 'One',
relayUrl: 'https://one.example.com',
);
final ws2 = Community.create(
name: 'Two',
relayUrl: 'https://two.example.com',
);
final notifier = container.read(communityListProvider.notifier);
await notifier.addCommunity(ws1);
await notifier.addCommunity(ws2);
await notifier.switchCommunity(ws1.id);
final teardown = Completer<void>();
container
.read(communityTransitionProvider)
.register(() => teardown.future);
final switching = notifier.switchCommunity(ws2.id);
await Future<void>.delayed(Duration.zero);
expect(await communityStorage.loadActiveId(), ws1.id);
teardown.complete();
await switching;
expect(await communityStorage.loadActiveId(), ws2.id);
});
test('switchCommunity updates active ID', () async {
container = createContainer();
await container.read(communityListProvider.future);