diff --git a/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleAudioEngine.kt b/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleAudioEngine.kt index ceceaa0a0..cb9445faf 100644 --- a/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleAudioEngine.kt +++ b/mobile/android/app/src/main/kotlin/xyz/block/buzz/mobile/HuddleAudioEngine.kt @@ -175,7 +175,12 @@ internal class HuddleAudioEngine( check(isSupported()) { "This device does not expose Android Opus encode/decode codecs." } val record = createAudioRecord() - val encoder = createEncoder() + val encoder = try { + createEncoder() + } catch (error: Throwable) { + runCatching { record.release() } + throw error + } try { encoder.start() record.startRecording() @@ -206,7 +211,27 @@ internal class HuddleAudioEngine( } fun setInterrupted(value: Boolean) { - interrupted.set(value) + val record = audioRecord + if (value) { + interrupted.set(true) + // 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() } + return + } + if (!running.get() || record == null) { + interrupted.set(false) + return + } + try { + record.startRecording() + check(record.recordingState == AudioRecord.RECORDSTATE_RECORDING) { + "Microphone recording did not resume." + } + interrupted.set(false) + } catch (error: Throwable) { + reportFailure("audio_resume_failed", error) + } } fun enqueueRemote(packet: HuddleRemoteOpusPacket) { @@ -498,17 +523,39 @@ internal class HuddleAudioEngine( onFailure(code, error.message ?: "Native Huddle audio failed.") } - private class PeerPlayback(val peerIndex: Int) { - private val decoder = createDecoder() - private val track = createAudioTrack() + private class PeerPlayback private constructor( + val peerIndex: Int, + private val decoder: MediaCodec, + private val track: AudioTrack, + ) { private val jitterQueue = HuddlePacketJitterQueue( capacity = PER_PEER_JITTER_CAPACITY, startPackets = JITTER_START_PACKETS, ) + companion object { + operator fun invoke(peerIndex: Int): PeerPlayback { + val decoder = createDecoder() + val track = try { + createAudioTrack() + } catch (error: Throwable) { + runCatching { decoder.release() } + throw error + } + return PeerPlayback(peerIndex, decoder, track) + } + } + init { - decoder.start() - track.play() + try { + decoder.start() + track.play() + } catch (error: Throwable) { + runCatching { track.release() } + runCatching { decoder.stop() } + runCatching { decoder.release() } + throw error + } } fun enqueue(packet: HuddleRemoteOpusPacket) { @@ -629,8 +676,9 @@ internal class HuddleAudioEngine( ) .setBufferSizeInBytes(max(minimum, FRAME_BYTES * 8)) .build() - check(record.state == AudioRecord.STATE_INITIALIZED) { - "Android microphone capture could not be initialized." + if (record.state != AudioRecord.STATE_INITIALIZED) { + record.release() + error("Android microphone capture could not be initialized.") } return record } @@ -645,8 +693,13 @@ internal class HuddleAudioEngine( setInteger(MediaFormat.KEY_MAX_INPUT_SIZE, FRAME_BYTES) setInteger(MediaFormat.KEY_PCM_ENCODING, AudioFormat.ENCODING_PCM_16BIT) } - return MediaCodec.createEncoderByType(MediaFormat.MIMETYPE_AUDIO_OPUS).apply { - configure(format, null, null, MediaCodec.CONFIGURE_FLAG_ENCODE) + val encoder = MediaCodec.createEncoderByType(MediaFormat.MIMETYPE_AUDIO_OPUS) + try { + encoder.configure(format, null, null, MediaCodec.CONFIGURE_FLAG_ENCODE) + return encoder + } catch (error: Throwable) { + runCatching { encoder.release() } + throw error } } @@ -662,8 +715,13 @@ internal class HuddleAudioEngine( setByteBuffer("csd-1", nativeLongBuffer(0)) setByteBuffer("csd-2", nativeLongBuffer(OPUS_SEEK_PREROLL_NS)) } - return MediaCodec.createDecoderByType(MediaFormat.MIMETYPE_AUDIO_OPUS).apply { - configure(format, null, null, 0) + val decoder = MediaCodec.createDecoderByType(MediaFormat.MIMETYPE_AUDIO_OPUS) + try { + decoder.configure(format, null, null, 0) + return decoder + } catch (error: Throwable) { + runCatching { decoder.release() } + throw error } } @@ -691,8 +749,9 @@ internal class HuddleAudioEngine( .setTransferMode(AudioTrack.MODE_STREAM) .setBufferSizeInBytes(max(minimum, FRAME_BYTES * 8)) .build() - check(track.state == AudioTrack.STATE_INITIALIZED) { - "Android Huddle audio output could not be initialized." + if (track.state != AudioTrack.STATE_INITIALIZED) { + track.release() + error("Android Huddle audio output could not be initialized.") } return track } diff --git a/mobile/ios/Runner/HuddleAudioEngine.swift b/mobile/ios/Runner/HuddleAudioEngine.swift index 738d721cc..c6f979388 100644 --- a/mobile/ios/Runner/HuddleAudioEngine.swift +++ b/mobile/ios/Runner/HuddleAudioEngine.swift @@ -859,13 +859,16 @@ final class HuddleAudioEngine { stateLock.unlock() return } + let message = + (error as? LocalizedError)?.errorDescription + ?? error.localizedDescription failureReported = true stateLock.unlock() - onFailure( - code, - (error as? LocalizedError)?.errorDescription - ?? error.localizedDescription - ) + // Fail closed on the native queue. Flutter will perform the durable session + // teardown too, but microphone and playout lifetime must not depend on a + // platform-channel callback making a successful round trip. + stopOnQueue() + onFailure(code, message) } private func matchesTargetPCM(_ format: AVAudioFormat) -> Bool { @@ -909,7 +912,7 @@ final class HuddleAudioEngine { } private final class HuddlePeerPlayback { - private let player = AVAudioPlayerNode() + private let player: AVAudioPlayerNode private let decoder: HuddleOpusDecoder private var jitterQueue = HuddlePacketJitterQueue( capacity: 10, @@ -921,9 +924,18 @@ private final class HuddlePeerPlayback { pcmFormat: AVAudioFormat ) throws { decoder = try HuddleOpusDecoder() + let player = AVAudioPlayerNode() + self.player = player engine.attach(player) - engine.connect(player, to: engine.mainMixerNode, format: pcmFormat) - player.play() + do { + engine.connect(player, to: engine.mainMixerNode, format: pcmFormat) + player.play() + } catch { + player.stop() + engine.disconnectNodeOutput(player) + engine.detach(player) + throw error + } } func enqueue(_ packet: HuddleRemoteOpusPacket) { diff --git a/mobile/ios/Runner/HuddleMediaPlugin.swift b/mobile/ios/Runner/HuddleMediaPlugin.swift index 518739110..ec3f9dd63 100644 --- a/mobile/ios/Runner/HuddleMediaPlugin.swift +++ b/mobile/ios/Runner/HuddleMediaPlugin.swift @@ -172,6 +172,11 @@ final class HuddleMediaPlugin { ]) } catch { audioSessionPrepared = false + try? audioSession.overrideOutputAudioPort(.none) + try? audioSession.setActive( + false, + options: [.notifyOthersOnDeactivation] + ) result( FlutterError( code: "audio_session_failed", diff --git a/mobile/lib/features/channels/channel_detail_page/huddle_sheet.dart b/mobile/lib/features/channels/channel_detail_page/huddle_sheet.dart index 44c630048..1ba28ed57 100644 --- a/mobile/lib/features/channels/channel_detail_page/huddle_sheet.dart +++ b/mobile/lib/features/channels/channel_detail_page/huddle_sheet.dart @@ -29,12 +29,14 @@ class _HuddleParticipantProfileUpdates extends Notifier { ephemeralChannelId: state.ephemeralChannelId, currentPubkey: state.currentPubkey, participantPubkeys: state.participantPubkeys, + wasAdmitted: state.wasAdmitted, ), ), ); - final members = - ref.watch(channelMembersProvider(channelId)).value ?? - const []; + final members = session.wasAdmitted + ? const [] + : ref.watch(channelMembersProvider(channelId)).value ?? + const []; final relayState = ref.watch(relaySessionProvider); final participantPubkeys = _huddleParticipantPubkeys( sessionParticipantPubkeys: session.ephemeralChannelId == channelId @@ -266,7 +268,9 @@ class _HuddleJoinSurface extends ConsumerWidget { DateTime.fromMillisecondsSinceEpoch(message.createdAt * 1000), ) > _huddleLifetime; - final canJoin = isMember && !isArchived && !ended && !stale; + final blockedByOtherHuddle = session.isInSession && !isThisSession; + final canJoin = + isMember && !isArchived && !ended && !stale && !blockedByOtherHuddle; final label = unavailable ? 'Start new' @@ -547,7 +551,7 @@ class _MobileHuddleCallPage extends ConsumerWidget { final participantPubkeys = _huddleParticipantPubkeys( sessionParticipantPubkeys: session.participantPubkeys, currentPubkey: localPubkey, - members: backingMembers, + members: session.wasAdmitted ? const [] : backingMembers, ); final remotePubkeys = participantPubkeys .where((pubkey) => pubkey != localPubkey) diff --git a/mobile/lib/shared/huddle/huddle_transport.dart b/mobile/lib/shared/huddle/huddle_transport.dart index 2d6a406b3..28df48b87 100644 --- a/mobile/lib/shared/huddle/huddle_transport.dart +++ b/mobile/lib/shared/huddle/huddle_transport.dart @@ -151,6 +151,8 @@ final class HuddleTransport implements HuddleTransportClient { int _generation = 0; bool _disposed = false; bool _intentionalClose = false; + bool _hasAdmitted = false; + Map? _admissionPreviousPeers; HuddleTransport({ required this.parameters, @@ -205,7 +207,15 @@ final class HuddleTransport implements HuddleTransportClient { final generation = ++_generation; _intentionalClose = false; - _emitState(HuddleTransportState(phase: HuddleTransportPhase.connecting)); + _admissionPreviousPeers = Map.from(_state.peers); + _emitState( + HuddleTransportState( + phase: HuddleTransportPhase.connecting, + localPeerIndex: _state.localPeerIndex, + rosterRevision: _state.rosterRevision, + peers: _state.peers, + ), + ); try { final channel = _channelFactory(parameters.audioWebSocketUri); @@ -227,7 +237,12 @@ final class HuddleTransport implements HuddleTransportClient { cancelOnError: false, ); _emitState( - HuddleTransportState(phase: HuddleTransportPhase.awaitingChallenge), + HuddleTransportState( + phase: HuddleTransportPhase.awaitingChallenge, + localPeerIndex: _state.localPeerIndex, + rosterRevision: _state.rosterRevision, + peers: _state.peers, + ), ); _handshakeTimer = Timer(handshakeTimeout, () { _fail( @@ -389,7 +404,12 @@ final class HuddleTransport implements HuddleTransportClient { ); _channel?.sink.add(jsonEncode(auth)); _emitState( - HuddleTransportState(phase: HuddleTransportPhase.authenticating), + HuddleTransportState( + phase: HuddleTransportPhase.authenticating, + localPeerIndex: _state.localPeerIndex, + rosterRevision: _state.rosterRevision, + peers: _state.peers, + ), ); } catch (error) { _fail( @@ -423,6 +443,9 @@ final class HuddleTransport implements HuddleTransportClient { final initialAdmission = _state.phase == HuddleTransportPhase.authenticating; + final previousAdmissionPeers = initialAdmission + ? (_admissionPreviousPeers ?? _state.peers) + : _state.peers; final revision = message['revision']; if (revision != null && (revision is! int || revision < 0)) { _handleProtocolProblem('Malformed Huddle joined revision.', generation); @@ -470,6 +493,9 @@ final class HuddleTransport implements HuddleTransportClient { if (!initialAdmission || revision == null) { peers[peerIndex] = peer; } + if (initialAdmission && _hasAdmitted) { + _emitAdmissionRosterChanges(previousAdmissionPeers, peers); + } if (!initialAdmission && replaced != null && replaced.pubkey != peer.pubkey) { @@ -503,11 +529,38 @@ final class HuddleTransport implements HuddleTransportClient { } _handshakeTimer?.cancel(); _handshakeTimer = null; + if (initialAdmission) { + _hasAdmitted = true; + _admissionPreviousPeers = null; + } if (_handshakeCompleter case final completer? when !completer.isCompleted) { completer.complete(); } } + void _emitAdmissionRosterChanges( + Map previous, + Map current, + ) { + final removedOrReplacedIndices = {}; + for (final entry in previous.entries) { + final replacement = current[entry.key]; + if (replacement == null || replacement.pubkey != entry.value.pubkey) { + removedOrReplacedIndices.add(entry.key); + _peerController.add( + HuddlePeerEvent( + type: replacement == null + ? HuddlePeerEventType.left + : HuddlePeerEventType.replaced, + peer: entry.value, + replacement: replacement, + ), + ); + } + } + _purgeAudioIngress(removedOrReplacedIndices); + } + void _handleLeft(Map message, int generation) { if (_state.phase != HuddleTransportPhase.connected) { _handleProtocolProblem('Unexpected Huddle left message.', generation); diff --git a/mobile/test/features/channels/channel_detail_page_test.dart b/mobile/test/features/channels/channel_detail_page_test.dart index 0e6a8b781..eaeb89204 100644 --- a/mobile/test/features/channels/channel_detail_page_test.dart +++ b/mobile/test/features/channels/channel_detail_page_test.dart @@ -46,7 +46,6 @@ import 'package:buzz/shared/widgets/frosted_scaffold.dart'; import 'package:buzz/shared/widgets/flapping_bee.dart'; import 'package:buzz/shared/widgets/keyboard_dismiss_on_drag.dart'; import 'package:buzz/shared/widgets/masked_avatar_badge.dart'; -import 'package:buzz/shared/widgets/avatar_image.dart'; import 'package:buzz/shared/widgets/skeleton.dart'; import 'package:shared_preferences/shared_preferences.dart'; @@ -114,6 +113,7 @@ NostrEvent _huddleMsg({ required int kind, String pubkey = 'alice', int createdAt = 1000, + String ephemeralChannelId = _huddleChannelId, }) => NostrEvent( id: id, pubkey: pubkey, @@ -122,7 +122,7 @@ NostrEvent _huddleMsg({ tags: [ ['h', _channelId], ], - content: jsonEncode({'ephemeral_channel_id': _huddleChannelId}), + content: jsonEncode({'ephemeral_channel_id': ephemeralChannelId}), sig: '', ); @@ -3258,6 +3258,53 @@ void main() { expect(join.onPressed, isNotNull); }); + testWidgets('disables a different Huddle card during an active call', ( + tester, + ) async { + const otherHuddleChannelId = 'other-huddle-channel'; + final now = DateTime.now().millisecondsSinceEpoch ~/ 1000; + await tester.pumpWidget( + _buildTestable( + messages: [ + _huddleMsg( + id: 'current-huddle', + kind: EventKind.huddleStarted, + pubkey: 'self', + createdAt: now, + ), + _huddleMsg( + id: 'other-huddle', + kind: EventKind.huddleStarted, + pubkey: 'alice', + createdAt: now, + ephemeralChannelId: otherHuddleChannelId, + ), + ], + users: const { + 'alice': UserProfile(pubkey: 'alice', displayName: 'Alice'), + 'self': UserProfile(pubkey: 'self', displayName: 'Self'), + }, + relayConfigNotifier: _HuddleRelayConfigNotifier(), + huddleCurrentPubkey: 'self', + huddleMediaFactory: _HuddleTestMedia.new, + huddleTransportFactory: (_) => _HuddleTestTransport(), + ), + ); + await tester.pumpAndSettle(); + + await tester.tap( + find.byKey(const ValueKey('huddle-Join-$_huddleChannelId')), + ); + await tester.pumpAndSettle(); + await tester.tap(find.byKey(const ValueKey('huddle-minimize'))); + await tester.pumpAndSettle(); + + final otherJoin = tester.widget( + find.byKey(const ValueKey('huddle-Join-$otherHuddleChannelId')), + ); + expect(otherJoin.onPressed, isNull); + }); + testWidgets('marks an expired Huddle card as ended', (tester) async { final now = DateTime.now().millisecondsSinceEpoch ~/ 1000; await tester.pumpWidget( @@ -3510,7 +3557,13 @@ void main() { (tester) async { final now = DateTime.now().millisecondsSinceEpoch ~/ 1000; final media = _HuddleTestMedia(); - final transport = _HuddleTestTransport(); + final transport = _HuddleTestTransport( + peers: const { + 1: HuddlePeer(pubkey: 'desktop', peerIndex: 1), + 2: HuddlePeer(pubkey: 'self', peerIndex: 2), + 3: HuddlePeer(pubkey: 'agent', peerIndex: 3), + }, + ); final navigator = _RecordingNavigatorObserver(); String? leftChannelId; @@ -4051,7 +4104,14 @@ void main() { // participant exercises the densest supported call. final remotePubkeys = List.generate(24, (index) => 'guest-$index'); final transport = _HuddleTestTransport( - peers: const {2: HuddlePeer(pubkey: 'self', peerIndex: 2)}, + peers: { + 0: const HuddlePeer(pubkey: 'self', peerIndex: 0), + for (var index = 0; index < remotePubkeys.length; index++) + index + 1: HuddlePeer( + pubkey: remotePubkeys[index], + peerIndex: index + 1, + ), + }, ); await tester.pumpWidget( @@ -4194,7 +4254,9 @@ void main() { (tester) async { const guestPubkey = 'guest'; final now = DateTime.now().millisecondsSinceEpoch ~/ 1000; - final membersNotifier = _MutableHuddleMembersNotifier(const []); + final transport = _HuddleTestTransport( + peers: const {2: HuddlePeer(pubkey: 'self', peerIndex: 2)}, + ); await tester.pumpWidget( _buildTestable( @@ -4213,14 +4275,11 @@ void main() { displayName: 'Guest', ), }, - huddleMembersNotifier: membersNotifier, relayConfigNotifier: _HuddleRelayConfigNotifier(), relaySessionNotifier: _ReconnectingRelaySession(), huddleCurrentPubkey: 'self', huddleMediaFactory: _HuddleTestMedia.new, - huddleTransportFactory: (_) => _HuddleTestTransport( - peers: const {2: HuddlePeer(pubkey: 'self', peerIndex: 2)}, - ), + huddleTransportFactory: (_) => transport, ), ); await tester.pumpAndSettle(); @@ -4234,14 +4293,9 @@ void main() { final soloCenter = tester.getCenter(localAvatar).dy; expect(soloCenter, closeTo(tester.getCenter(stage).dy, 1)); - membersNotifier.replace([ - ChannelMember( - pubkey: guestPubkey, - role: 'member', - joinedAt: DateTime(2025), - ), - ]); - await tester.pump(); + transport.emitPeerJoin( + const HuddlePeer(pubkey: guestPubkey, peerIndex: 1), + ); await tester.pump(); await tester.pump(const Duration(milliseconds: 130)); final movingCenter = tester.getCenter(localAvatar).dy; @@ -4261,92 +4315,65 @@ void main() { }, ); - testWidgets( - 'adds backing-channel members and applies profile updates without reopening', - (tester) async { - const latePubkey = - 'aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa'; - const avatarUrl = - 'data:image/png;base64,iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mNk+A8AAQUBAScY42YAAAAASUVORK5CYII='; - final now = DateTime.now().millisecondsSinceEpoch ~/ 1000; - final membersNotifier = _MutableHuddleMembersNotifier(const []); - final users = _FakeUserCacheNotifier(const { - 'desktop': UserProfile(pubkey: 'desktop', displayName: 'Miles'), - 'self': UserProfile(pubkey: 'self', displayName: 'Self'), - }); + testWidgets('does not add backing-channel members after relay admission', ( + tester, + ) async { + const staleMemberPubkey = + 'aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa'; + final now = DateTime.now().millisecondsSinceEpoch ~/ 1000; + final membersNotifier = _MutableHuddleMembersNotifier(const []); - await tester.pumpWidget( - _buildTestable( - messages: [ - _huddleMsg( - id: 'live-roster-call', - kind: EventKind.huddleStarted, - pubkey: 'self', - createdAt: now, - ), - ], - userCacheNotifier: users, - huddleMembersNotifier: membersNotifier, - relayConfigNotifier: _HuddleRelayConfigNotifier(), - relaySessionNotifier: _ReconnectingRelaySession(), - huddleCurrentPubkey: 'self', - huddleMediaFactory: _HuddleTestMedia.new, - huddleTransportFactory: (_) => _HuddleTestTransport(), - ), - ); - await tester.pumpAndSettle(); - await tester.tap(find.widgetWithText(FilledButton, 'Join')); - await tester.pumpAndSettle(); + await tester.pumpWidget( + _buildTestable( + messages: [ + _huddleMsg( + id: 'authoritative-live-roster', + kind: EventKind.huddleStarted, + pubkey: 'self', + createdAt: now, + ), + ], + users: const { + 'desktop': UserProfile(pubkey: 'desktop', displayName: 'Miles'), + 'self': UserProfile(pubkey: 'self', displayName: 'Self'), + staleMemberPubkey: UserProfile( + pubkey: staleMemberPubkey, + displayName: 'Stale member', + ), + }, + huddleMembersNotifier: membersNotifier, + relayConfigNotifier: _HuddleRelayConfigNotifier(), + relaySessionNotifier: _ReconnectingRelaySession(), + huddleCurrentPubkey: 'self', + huddleMediaFactory: _HuddleTestMedia.new, + huddleTransportFactory: (_) => _HuddleTestTransport(), + ), + ); + await tester.pumpAndSettle(); + await tester.tap(find.widgetWithText(FilledButton, 'Join')); + await tester.pumpAndSettle(); - final lateAvatar = find.byKey( - const ValueKey('huddle-participant-avatar-$latePubkey'), - ); - expect(lateAvatar, findsNothing); + membersNotifier.replace([ + ChannelMember( + pubkey: staleMemberPubkey, + role: 'member', + joinedAt: DateTime(2025), + ), + ]); + await tester.pump(); + await tester.pump(); - membersNotifier.replace([ - ChannelMember( - pubkey: latePubkey, - role: 'bot', - joinedAt: DateTime(2025), - ), - ]); - await tester.pump(); - await tester.pump(); - - expect(lateAvatar, findsOneWidget); - expect( - find.descendant(of: lateAvatar, matching: find.text('N')), - findsNothing, - ); - expect( - find.descendant( - of: lateAvatar, - matching: find.byIcon(LucideIcons.userRound), - ), - findsOneWidget, - ); - - users.replace( - const UserProfile( - pubkey: latePubkey, - displayName: 'Ada', - avatarUrl: avatarUrl, - ), - ); - await tester.pump(const Duration(milliseconds: 60)); - - final avatar = tester.widget( - find.descendant(of: lateAvatar, matching: find.byType(AvatarImage)), - ); - expect(avatar.imageUrl, avatarUrl); - expect(find.text('Ada'), findsNothing); - - await tester.tap(lateAvatar); - await tester.pump(); - expect(find.text('Ada'), findsOneWidget); - expect(find.byKey(const ValueKey('huddle-leave')), findsOneWidget); - }, - ); + expect( + find.byKey( + const ValueKey('huddle-participant-avatar-$staleMemberPubkey'), + ), + findsNothing, + ); + expect( + find.byKey(const ValueKey('huddle-participant-avatar-desktop')), + findsOneWidget, + ); + }); testWidgets('top-right call end leaves audio and the backing channel', ( tester, @@ -9039,15 +9066,15 @@ final class _HuddleTestTransport implements HuddleTransportClient { _HuddleTestTransport({ this.connectError, this.connectGate, - this.peers = const { + Map peers = const { 1: HuddlePeer(pubkey: 'desktop', peerIndex: 1), 2: HuddlePeer(pubkey: 'self', peerIndex: 2), }, - }); + }) : _peers = Map.from(peers); final HuddleTransportError? connectError; final Future? connectGate; - final Map peers; + final Map _peers; final _states = StreamController.broadcast(sync: true); final _remoteFrames = StreamController.broadcast( sync: true, @@ -9056,6 +9083,19 @@ final class _HuddleTestTransport implements HuddleTransportClient { final _issues = StreamController.broadcast(sync: true); HuddleTransportState _state = HuddleTransportState.idle(); + void emitPeerJoin(HuddlePeer peer) { + _peers[peer.peerIndex] = peer; + _state = HuddleTransportState( + phase: HuddleTransportPhase.connected, + localPeerIndex: _state.localPeerIndex, + peers: _peers, + ); + _states.add(_state); + _peerEvents.add( + HuddlePeerEvent(type: HuddlePeerEventType.joined, peer: peer), + ); + } + void emitRemoteAudio({int levelDbov = -30, int sequence = 1}) { _remoteFrames.add( HuddleRemoteAudioFrame( @@ -9093,7 +9133,7 @@ final class _HuddleTestTransport implements HuddleTransportClient { _state = HuddleTransportState( phase: HuddleTransportPhase.connected, localPeerIndex: 2, - peers: peers, + peers: _peers, ); _states.add(_state); } diff --git a/mobile/test/shared/huddle/huddle_transport_test.dart b/mobile/test/shared/huddle/huddle_transport_test.dart index fcf96b709..f85b34370 100644 --- a/mobile/test/shared/huddle/huddle_transport_test.dart +++ b/mobile/test/shared/huddle/huddle_transport_test.dart @@ -430,6 +430,69 @@ void main() { }, ); + test( + 're-admission reports departed and reused peer indexes before new media', + () async { + final firstChannel = _ControlledWebSocketChannel(); + final secondChannel = _ControlledWebSocketChannel(); + var connection = 0; + final transport = HuddleTransport( + parameters: HuddleConnectionParameters( + relayWebSocketUrl: 'wss://buzz.example', + nsec: _privateKey, + parentChannelId: _parentChannelId, + ephemeralChannelId: _ephemeralChannelId, + ), + channelFactory: (_) => connection++ == 0 ? firstChannel : secondChannel, + connectTimeout: const Duration(seconds: 1), + handshakeTimeout: const Duration(seconds: 1), + ); + addTearDown(transport.dispose); + await _connect( + firstChannel, + transport, + revision: 1, + peers: const [ + {'pubkey': 'old', 'peer_index': 4}, + {'pubkey': 'departed', 'peer_index': 5}, + ], + ); + final events = []; + final subscription = transport.peerEvents.listen(events.add); + addTearDown(subscription.cancel); + + firstChannel.emitError(StateError('socket lost')); + await _waitForPhase(transport, HuddleTransportPhase.failed); + final reconnect = transport.connect(); + await _waitForPhase(transport, HuddleTransportPhase.awaitingChallenge); + secondChannel.emitText( + jsonEncode({'type': 'challenge', 'challenge': 'reconnect'}), + ); + await _waitForPhase(transport, HuddleTransportPhase.authenticating); + secondChannel.emitText( + jsonEncode({ + 'type': 'joined', + 'revision': 4, + 'pubkey': 'self', + 'peer_index': 3, + 'peers': [ + {'pubkey': 'new', 'peer_index': 4}, + {'pubkey': 'self', 'peer_index': 3}, + ], + }), + ); + await reconnect; + + expect(transport.state.peers[4]?.pubkey, 'new'); + expect(events, hasLength(2)); + expect(events[0].type, HuddlePeerEventType.replaced); + expect(events[0].peer.pubkey, 'old'); + expect(events[0].replacement?.pubkey, 'new'); + expect(events[1].type, HuddlePeerEventType.left); + expect(events[1].peer.pubkey, 'departed'); + }, + ); + test( 'intentional disconnect reaches disconnected without an error', () async { @@ -522,6 +585,7 @@ final class _ControlledWebSocketChannel implements WebSocketChannel { void emitText(String value) => _streamController.add(value); void emitBinary(Uint8List value) => _streamController.add(value); + void emitError(Object error) => _streamController.addError(error); @override Future get ready => Future.value();