Skip to content

Commit 67df841

Browse files
Keep group call tracks synchronized during reconnection (#106)
1 parent 9753fac commit 67df841

2 files changed

Lines changed: 273 additions & 4 deletions

File tree

‎apps/web/src/hooks/useWebRTCHost.test.ts‎

Lines changed: 259 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -847,6 +847,265 @@ describe('useWebRTCHost', () => {
847847
});
848848

849849
describe('audio relay lifecycle', () => {
850+
it.each(['connected', 'connecting', 'disconnected'])(
851+
'relays new participant audio after a destination transitions through %s',
852+
async (destinationState) => {
853+
const { result, joinHandler, signalHandler, unmount } = await initializeSubscribedHost();
854+
855+
try {
856+
await act(async () => {
857+
joinHandler({
858+
newPresences: [
859+
{ user_id: 'viewer-1', role: 'viewer' },
860+
{ user_id: 'viewer-2', role: 'viewer' },
861+
],
862+
});
863+
await flushAsyncWork();
864+
for (const viewerId of ['viewer-1', 'viewer-2']) {
865+
signalHandler({
866+
payload: {
867+
type: 'answer',
868+
sdp: `${viewerId}-initial-answer`,
869+
senderId: viewerId,
870+
targetId: 'host-1',
871+
timestamp: Date.now(),
872+
},
873+
});
874+
}
875+
await flushAsyncWork();
876+
});
877+
878+
const sourcePc = MockRTCPeerConnection.instances[0]!;
879+
const targetPc = MockRTCPeerConnection.instances[1]!;
880+
const sourceTrack = {
881+
id: 'viewer-1-new-audio',
882+
kind: 'audio',
883+
enabled: true,
884+
} as MediaStreamTrack;
885+
886+
await act(async () => {
887+
sourcePc.connectionState = 'connected';
888+
sourcePc.onconnectionstatechange?.();
889+
targetPc.connectionState = destinationState;
890+
targetPc.onconnectionstatechange?.();
891+
await flushAsyncWork();
892+
});
893+
expect(result.current.viewers.get('viewer-2')?.connectionState).toBe(
894+
destinationState === 'disconnected' ? 'reconnecting' : destinationState
895+
);
896+
897+
await act(async () => {
898+
sourcePc.ontrack?.({ track: sourceTrack, streams: [] });
899+
await flushAsyncWork();
900+
targetPc.connectionState = 'connected';
901+
targetPc.onconnectionstatechange?.();
902+
await flushAsyncWork();
903+
});
904+
905+
expect(result.current.viewers.get('viewer-2')?.connectionState).toBe('connected');
906+
expect(result.current.viewers.get('viewer-1')?.audioTrack).toBe(sourceTrack);
907+
expect(
908+
targetPc.addTrack.mock.calls.filter(([track]) => track === sourceTrack)
909+
).toHaveLength(1);
910+
} finally {
911+
unmount();
912+
}
913+
}
914+
);
915+
916+
describe('reconnecting destinations', () => {
917+
const connectViewers = async (answerTarget = true) => {
918+
const host = await initializeSubscribedHost();
919+
const answer = async (viewerId = 'viewer-2') => {
920+
await act(async () => {
921+
host.signalHandler({
922+
payload: {
923+
type: 'answer',
924+
sdp: `${viewerId}-answer`,
925+
senderId: viewerId,
926+
targetId: 'host-1',
927+
timestamp: Date.now(),
928+
},
929+
});
930+
await flushAsyncWork();
931+
});
932+
};
933+
await act(async () => {
934+
host.joinHandler({
935+
newPresences: ['viewer-1', 'viewer-2'].map((user_id) => ({ user_id, role: 'viewer' })),
936+
});
937+
await flushAsyncWork();
938+
});
939+
await answer('viewer-1');
940+
if (answerTarget) await answer();
941+
const source = MockRTCPeerConnection.instances[0]!;
942+
const target = MockRTCPeerConnection.instances[1]!;
943+
const setTargetState = (state: string) => {
944+
act(() => {
945+
target.connectionState = state;
946+
target.onconnectionstatechange?.();
947+
});
948+
};
949+
setTargetState('disconnected');
950+
const receive = async (id: string) => {
951+
const track = { id, kind: 'audio', enabled: true } as MediaStreamTrack;
952+
await act(async () => {
953+
source.ontrack?.({ track, streams: [] });
954+
await flushAsyncWork();
955+
});
956+
return track;
957+
};
958+
const screen = (id: string) => {
959+
const video = { id: `${id}-video`, kind: 'video', contentHint: '' } as MediaStreamTrack;
960+
const audio = { id: `${id}-audio`, kind: 'audio' } as MediaStreamTrack;
961+
return new MediaStream([video, audio]);
962+
};
963+
return { ...host, source, target, answer, setTargetState, receive, screen };
964+
};
965+
966+
it('queues relay changes behind the outstanding offer and targets the next offer', async () => {
967+
const host = await connectViewers(false);
968+
try {
969+
const track = await host.receive('queued-mic');
970+
expect(host.target.addTrack).toHaveBeenCalledWith(track, expect.anything());
971+
expect(host.target.createOffer).toHaveBeenCalledTimes(1);
972+
expect(host.target.signalingState).toBe('have-local-offer');
973+
await host.answer();
974+
expect(host.target.createOffer).toHaveBeenCalledTimes(2);
975+
expect(mockChannel.send).toHaveBeenLastCalledWith(
976+
expect.objectContaining({
977+
payload: expect.objectContaining({ type: 'offer', targetId: 'viewer-2' }),
978+
})
979+
);
980+
host.setTargetState('connected');
981+
expect(host.target.createOffer).toHaveBeenCalledTimes(2);
982+
} finally {
983+
host.unmount();
984+
}
985+
});
986+
987+
it('replaces a muted relay once while disconnected, without echoing it to its source', async () => {
988+
const host = await connectViewers();
989+
try {
990+
host.setTargetState('connected');
991+
const oldTrack = await host.receive('old-mic');
992+
await host.answer();
993+
const oldSender = host.target.getSenders().find((sender) => sender.track === oldTrack)!;
994+
act(() => host.result.current.muteViewer('viewer-1', true));
995+
host.setTargetState('disconnected');
996+
const replacement = await host.receive('new-mic');
997+
await act(async () => {
998+
host.source.ontrack?.({ track: replacement, streams: [] });
999+
await flushAsyncWork();
1000+
});
1001+
host.setTargetState('connected');
1002+
expect(replacement.enabled).toBe(false);
1003+
expect(host.target.removeTrack).toHaveBeenCalledWith(oldSender);
1004+
expect(
1005+
host.target.getSenders().filter((sender) => sender.track === replacement)
1006+
).toHaveLength(1);
1007+
expect(host.target.getSenders().some((sender) => sender.track === oldTrack)).toBe(false);
1008+
expect(
1009+
host.target.addTrack.mock.calls.filter(([track]) => track === replacement)
1010+
).toHaveLength(1);
1011+
expect(host.source.addTrack).not.toHaveBeenCalledWith(replacement, expect.anything());
1012+
} finally {
1013+
host.unmount();
1014+
}
1015+
});
1016+
1017+
it('removes a departed source while disconnected and does not resurrect it on recovery', async () => {
1018+
const host = await connectViewers();
1019+
try {
1020+
host.setTargetState('connected');
1021+
const track = await host.receive('departing-mic');
1022+
expect(host.target.getSenders().some((sender) => sender.track === track)).toBe(true);
1023+
host.setTargetState('disconnected');
1024+
await act(async () => {
1025+
host.leaveHandler({ leftPresences: [{ user_id: 'viewer-1', role: 'viewer' }] });
1026+
await flushAsyncWork();
1027+
});
1028+
await host.answer();
1029+
host.setTargetState('connected');
1030+
expect(host.target.getSenders().some((sender) => sender.track === track)).toBe(false);
1031+
expect(host.result.current.viewers.has('viewer-1')).toBe(false);
1032+
} finally {
1033+
host.unmount();
1034+
}
1035+
});
1036+
1037+
it.each(['publish', 'replace', 'unpublish'])(
1038+
'keeps %s synchronized during reconnection without removing host mic or relays',
1039+
async (operation) => {
1040+
const host = await connectViewers();
1041+
try {
1042+
host.setTargetState('connected');
1043+
const relay = await host.receive('retained-mic');
1044+
await host.answer();
1045+
const voiceSenders = host.target
1046+
.getSenders()
1047+
.filter((sender) => sender.track?.kind === 'audio');
1048+
const voiceTracks = voiceSenders.map((sender) => sender.track);
1049+
1050+
// Set up the prior state while connected so this test isolates only
1051+
// the next publish/unpublish operation during disconnection.
1052+
host.setTargetState('connected');
1053+
await act(async () => {
1054+
if (operation === 'publish') await host.result.current.unpublishStream();
1055+
else await host.result.current.publishStream(host.screen('previous'));
1056+
await flushAsyncWork();
1057+
});
1058+
await host.answer('viewer-1');
1059+
await host.answer();
1060+
host.setTargetState('disconnected');
1061+
const before = host.target
1062+
.getSenders()
1063+
.filter((sender) => sender.track && !voiceSenders.includes(sender));
1064+
const next = host.screen('next');
1065+
1066+
await act(async () => {
1067+
if (operation === 'unpublish') await host.result.current.unpublishStream();
1068+
else await host.result.current.publishStream(next);
1069+
await flushAsyncWork();
1070+
});
1071+
host.setTargetState('connected');
1072+
1073+
expect(voiceSenders.map((sender) => sender.track)).toEqual(voiceTracks);
1074+
expect(host.target.getSenders().some((sender) => sender.track === relay)).toBe(true);
1075+
for (const sender of before) expect(sender.track).toBeNull();
1076+
const remainingScreen = host.target
1077+
.getSenders()
1078+
.filter((sender) => sender.track && !voiceSenders.includes(sender));
1079+
expect(remainingScreen.map((sender) => sender.track)).toEqual(
1080+
operation === 'unpublish' ? [] : next.getTracks()
1081+
);
1082+
} finally {
1083+
host.unmount();
1084+
}
1085+
}
1086+
);
1087+
1088+
it.each(['failed', 'closed'])('does not add tracks to a %s peer', async (state) => {
1089+
const host = await connectViewers();
1090+
try {
1091+
host.setTargetState(state);
1092+
const additions = host.target.addTrack.mock.calls.length;
1093+
const offers = host.target.createOffer.mock.calls.length;
1094+
await host.receive('late-mic');
1095+
await act(async () => {
1096+
await host.result.current.publishStream(host.screen('late'));
1097+
await host.result.current.unpublishStream();
1098+
await flushAsyncWork();
1099+
});
1100+
expect(host.result.current.viewers.has('viewer-2')).toBe(false);
1101+
expect(host.target.addTrack).toHaveBeenCalledTimes(additions);
1102+
expect(host.target.createOffer).toHaveBeenCalledTimes(offers);
1103+
} finally {
1104+
host.unmount();
1105+
}
1106+
});
1107+
});
1108+
8501109
it('still relays audio and retries playback when local amplification fails', async () => {
8511110
const { result, joinHandler, signalHandler } = await initializeSubscribedHost();
8521111

‎apps/web/src/hooks/useWebRTCHost.ts‎

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -390,7 +390,8 @@ export function useWebRTCHost({
390390
[hostId]
391391
);
392392

393-
// Relay a viewer's audio track to all other connected viewers via renegotiation
393+
// Temporary media disconnection must not discard track updates. Signaling
394+
// still queues offers independently, so reconnecting peers stay in sync.
394395
const relayAudioToOtherViewers = useCallback(
395396
async (sourceViewerId: string, audioTrack: MediaStreamTrack) => {
396397
const audioStream = new MediaStream([audioTrack]);
@@ -401,7 +402,8 @@ export function useWebRTCHost({
401402
if (otherId === sourceViewerId) continue;
402403
if (
403404
otherViewer.connectionState !== 'connected' &&
404-
otherViewer.connectionState !== 'connecting'
405+
otherViewer.connectionState !== 'connecting' &&
406+
otherViewer.connectionState !== 'reconnecting'
405407
)
406408
continue;
407409

@@ -1032,7 +1034,11 @@ export function useWebRTCHost({
10321034
localStreamRef.current !== stream
10331035
)
10341036
break;
1035-
if (viewer.connectionState !== 'connected' && viewer.connectionState !== 'connecting')
1037+
if (
1038+
viewer.connectionState !== 'connected' &&
1039+
viewer.connectionState !== 'connecting' &&
1040+
viewer.connectionState !== 'reconnecting'
1041+
)
10361042
continue;
10371043

10381044
try {
@@ -1082,7 +1088,11 @@ export function useWebRTCHost({
10821088

10831089
for (const viewer of viewersRef.current.values()) {
10841090
if (publishedStreamVersionRef.current !== unpublishVersion) break;
1085-
if (viewer.connectionState !== 'connected' && viewer.connectionState !== 'connecting')
1091+
if (
1092+
viewer.connectionState !== 'connected' &&
1093+
viewer.connectionState !== 'connecting' &&
1094+
viewer.connectionState !== 'reconnecting'
1095+
)
10861096
continue;
10871097

10881098
try {

0 commit comments

Comments
 (0)