Harden Huddle reconnect and native cleanup

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:31:48 +01:00
parent 1c9180bda8
commit 8b950daef3
7 changed files with 372 additions and 135 deletions
@@ -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
}
+20 -8
View File
@@ -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) {
@@ -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",
@@ -29,12 +29,14 @@ class _HuddleParticipantProfileUpdates extends Notifier<int> {
ephemeralChannelId: state.ephemeralChannelId,
currentPubkey: state.currentPubkey,
participantPubkeys: state.participantPubkeys,
wasAdmitted: state.wasAdmitted,
),
),
);
final members =
ref.watch(channelMembersProvider(channelId)).value ??
const <ChannelMember>[];
final members = session.wasAdmitted
? const <ChannelMember>[]
: ref.watch(channelMembersProvider(channelId)).value ??
const <ChannelMember>[];
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)
+56 -3
View File
@@ -151,6 +151,8 @@ final class HuddleTransport implements HuddleTransportClient {
int _generation = 0;
bool _disposed = false;
bool _intentionalClose = false;
bool _hasAdmitted = false;
Map<int, HuddlePeer>? _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<int, HuddlePeer>.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<int, HuddlePeer> previous,
Map<int, HuddlePeer> current,
) {
final removedOrReplacedIndices = <int>{};
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<String, dynamic> message, int generation) {
if (_state.phase != HuddleTransportPhase.connected) {
_handleProtocolProblem('Unexpected Huddle left message.', generation);
@@ -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<FilledButton>(
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<AvatarImage>(
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<int, HuddlePeer> peers = const {
1: HuddlePeer(pubkey: 'desktop', peerIndex: 1),
2: HuddlePeer(pubkey: 'self', peerIndex: 2),
},
});
}) : _peers = Map<int, HuddlePeer>.from(peers);
final HuddleTransportError? connectError;
final Future<void>? connectGate;
final Map<int, HuddlePeer> peers;
final Map<int, HuddlePeer> _peers;
final _states = StreamController<HuddleTransportState>.broadcast(sync: true);
final _remoteFrames = StreamController<HuddleRemoteAudioFrame>.broadcast(
sync: true,
@@ -9056,6 +9083,19 @@ final class _HuddleTestTransport implements HuddleTransportClient {
final _issues = StreamController<HuddleTransportError>.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);
}
@@ -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 = <HuddlePeerEvent>[];
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<void> get ready => Future.value();