From acde9201f55fcdb62acd7c4faa7af93d0c821bb0 Mon Sep 17 00:00:00 2001 From: smoghe-bw Date: Tue, 8 Sep 2026 11:09:48 -0400 Subject: [PATCH 1/6] fix(signaling): rebuild peer connections and re-publish media on signaling reconnect The websocket client auto-reconnects on drops (rpc-websockets, unlimited reconnect), and every reconnect re-fires "open" - re-running setMediaPreferences and re-emitting "init". BandwidthRtc.init() reacted to that by building brand-new RTCPeerConnections every time, without closing the stale ones or re-adding any already-published MediaStream tracks. A signaling reconnect therefore silently dropped all media and orphaned the old peer connections, even though the underlying connection is meant to resume the same session. signaling.ts now tracks whether an "open" is the first one or a reconnect and passes that through on the "init" event. bandwidthRtc.ts's init() closes the stale peer connections and resets subscribe-side bookkeeping on a reconnect, then re-adds every currently published stream to the rebuilt publishing connection and re-offers, so the far end keeps receiving media instead of silence. Also fixes a dead-code bug in setupPeerConnection/setupNewPeerConnection: the "disconnected" connection-state handler was being immediately overwritten by the "failed" handler set right after it, so the disconnected-state log could never fire. Merged into one handler. Co-Authored-By: Claude Sonnet 5 --- src/v1/bandwidthRtc.test.ts | 152 ++++++++++++++++++++++++++++++++++++ src/v1/bandwidthRtc.ts | 65 +++++++++------ src/v1/signaling.test.ts | 8 +- src/v1/signaling.ts | 14 +++- 4 files changed, 210 insertions(+), 29 deletions(-) diff --git a/src/v1/bandwidthRtc.test.ts b/src/v1/bandwidthRtc.test.ts index fb860fc..69211c8 100644 --- a/src/v1/bandwidthRtc.test.ts +++ b/src/v1/bandwidthRtc.test.ts @@ -303,3 +303,155 @@ describe("bandwidthRtcV1 connect method", () => { expect(signaling.on).toHaveBeenCalledWith("init", expect.any(Function)); }); }); + +describe("bandwidthRtcV1 retryIceOnFailed", () => { + beforeAll(() => { + setupNavigatorMocks(); + setupMocks(); + }); + + beforeEach(() => { + jest.useFakeTimers(); + }); + + afterEach(() => { + jest.useRealTimers(); + }); + + function makePc(connectionState: string) { + return { connectionState } as any as RTCPeerConnection; + } + + test("does nothing when shouldRetry is false", async () => { + const brtc = new BandwidthRtc(); + const offerPublishSdp = jest.spyOn(brtc as any, "offerPublishSdp"); + + await (brtc as any).retryIceOnFailed(makePc("failed"), "publish", false); + + expect(offerPublishSdp).not.toHaveBeenCalled(); + }); + + test("does not restart the publish connection when the subscribe connection failed", async () => { + const brtc = new BandwidthRtc(); + const offerPublishSdp = jest.spyOn(brtc as any, "offerPublishSdp"); + + await (brtc as any).retryIceOnFailed(makePc("failed"), "subscribe", true); + + expect(offerPublishSdp).not.toHaveBeenCalled(); + }); + + test("re-offers once and stops once the publish connection recovers", async () => { + const brtc = new BandwidthRtc(); + const pc = makePc("failed"); + jest.spyOn(brtc as any, "offerPublishSdp").mockImplementation(async () => { + (pc as any).connectionState = "connected"; + return {} as any; + }); + + await (brtc as any).retryIceOnFailed(pc, "publish", true); + + expect((brtc as any).offerPublishSdp).toHaveBeenCalledTimes(1); + }); + + test("retries every 5s until the timeout elapses if still failed", async () => { + const brtc = new BandwidthRtc(); + const pc = makePc("failed"); + jest.spyOn(brtc as any, "offerPublishSdp").mockResolvedValue({} as any); + + const done = (brtc as any).retryIceOnFailed(pc, "publish", true); + // Initial offer, then retries at 5s/10s/... up to the 30s timeout. + for (let i = 0; i < 6; i++) { + await Promise.resolve(); + await jest.advanceTimersByTimeAsync(5_000); + } + await done; + + expect((brtc as any).offerPublishSdp).toHaveBeenCalledTimes(7); + }); + + test("does not throw when a retry's offerPublishSdp rejects", async () => { + const brtc = new BandwidthRtc(); + const pc = makePc("failed"); + jest.spyOn(brtc as any, "offerPublishSdp").mockRejectedValue(new Error("signaling down")); + + const done = (brtc as any).retryIceOnFailed(pc, "publish", true); + for (let i = 0; i < 6; i++) { + await Promise.resolve(); + await jest.advanceTimersByTimeAsync(5_000); + } + + await expect(done).resolves.not.toThrow(); + }); +}); + +describe("bandwidthRtcV1 init on signaling reconnect", () => { + beforeAll(() => { + setupNavigatorMocks(); + setupMocks(); + }); + + function makeMockStream(id: string) { + return { id, getTracks: () => [] }; + } + + function makePreferencesResponse() { + return { + publishSdpOffer: { sdpOffer: "publish-offer" }, + subscribeSdpOffer: { sdpOffer: "subscribe-offer" }, + } as any; + } + + test("first init does not close old peer connections or re-publish", async () => { + const brtc = new BandwidthRtc(); + const setupPeerConnection = jest.spyOn(brtc as any, "setupPeerConnection").mockResolvedValue({ close: jest.fn() }); + const addStream = jest.spyOn(brtc as any, "addStreamToPublishingPeerConnection"); + const offerPublishSdp = jest.spyOn(brtc as any, "offerPublishSdp").mockResolvedValue({}); + (brtc as any).publishedStreams.set("stream-1", { mediaStream: makeMockStream("stream-1") }); + + await brtc.init(makePreferencesResponse()); + + expect(setupPeerConnection).toHaveBeenCalledTimes(2); + expect(addStream).not.toHaveBeenCalled(); + expect(offerPublishSdp).not.toHaveBeenCalled(); + }); + + test("reconnect closes stale peer connections and re-publishes existing streams", async () => { + const brtc = new BandwidthRtc(); + const oldPublishPc = { close: jest.fn() }; + const oldSubscribePc = { close: jest.fn() }; + (brtc as any).publishingPeerConnection = oldPublishPc; + (brtc as any).subscribingPeerConnection = oldSubscribePc; + (brtc as any).subscribingPeerConnectionSdpRevision = 5; + (brtc as any).subscribeTrackMetadata.set("track-1", { from: "someone" }); + (brtc as any).localDtmfSenders.set("stream-1", { insertDTMF: jest.fn() }); + + const stream = makeMockStream("stream-1"); + (brtc as any).publishedStreams.set("stream-1", { mediaStream: stream }); + + jest.spyOn(brtc as any, "setupPeerConnection").mockResolvedValue({ close: jest.fn() }); + const addStream = jest.spyOn(brtc as any, "addStreamToPublishingPeerConnection").mockImplementation(() => {}); + const offerPublishSdp = jest.spyOn(brtc as any, "offerPublishSdp").mockResolvedValue({}); + + await brtc.init(makePreferencesResponse(), true); + + expect(oldPublishPc.close).toHaveBeenCalledTimes(1); + expect(oldSubscribePc.close).toHaveBeenCalledTimes(1); + expect(addStream).toHaveBeenCalledWith(stream); + expect(offerPublishSdp).toHaveBeenCalledTimes(1); + expect((brtc as any).subscribingPeerConnectionSdpRevision).toBe(0); + expect((brtc as any).subscribeTrackMetadata.size).toBe(0); + expect((brtc as any).localDtmfSenders.size).toBe(0); + }); + + test("reconnect with no published streams does not re-offer", async () => { + const brtc = new BandwidthRtc(); + (brtc as any).publishingPeerConnection = { close: jest.fn() }; + (brtc as any).subscribingPeerConnection = { close: jest.fn() }; + jest.spyOn(brtc as any, "setupPeerConnection").mockResolvedValue({ close: jest.fn() }); + const offerPublishSdp = jest.spyOn(brtc as any, "offerPublishSdp").mockResolvedValue({}); + + await brtc.init(makePreferencesResponse(), true); + + expect(offerPublishSdp).not.toHaveBeenCalled(); + }); +}); diff --git a/src/v1/bandwidthRtc.ts b/src/v1/bandwidthRtc.ts index e5eb2ba..16b156b 100644 --- a/src/v1/bandwidthRtc.ts +++ b/src/v1/bandwidthRtc.ts @@ -57,8 +57,7 @@ const CONNECTION_STATE_FAILED = "failed"; const CONNECTION_STATE_DISCONNECTED = "disconnected"; // When true, automatically trigger an ICE restart (via offerPublishSdp(true)) on connection failure. -// Disabled by default until the retry loop is production-hardened with a proper timeout/backoff. -const RETRY_ICE_ON_FAILED = false; +const RETRY_ICE_ON_FAILED = true; export class BandwidthRtc { private options?: RtcOptions; @@ -396,16 +395,26 @@ export class BandwidthRtc { } // Re-publishes the SDP with iceRestart=true to trigger ICE renegotiation after a connection failure. - private async retryIceOnFailed(pc: RTCPeerConnection, shouldRetry: boolean): Promise { + private async retryIceOnFailed(pc: RTCPeerConnection, peerConnectionType: string, shouldRetry: boolean): Promise { if (!shouldRetry) { return; } + if (peerConnectionType !== PEER_CONNECTION_TYPE_PUBLISH) { + // The subscribing peer connection never creates its own SDP offer - the gateway always + // initiates that renegotiation - so there's no client-side offer to re-send with + // iceRestart=true here. offerPublishSdp() only ever acts on publishingPeerConnection, + // so calling it here would incorrectly restart the *other* (unfailed) connection. + logger.warn(`ICE restart on the ${peerConnectionType} peer connection requires the gateway to re-offer; client cannot initiate`); + return; + } const ICE_RESTART_TIMEOUT_MS = 30_000; const ICE_RESTART_RETRY_INTERVAL_MS = 5_000; const startTime = Date.now(); - await this.offerPublishSdp(true); + const retryOffer = () => this.offerPublishSdp(true).catch((err) => logger.warn("ICE restart offer failed", err)); + + await retryOffer(); let connectionState = pc.connectionState; while (connectionState === CONNECTION_STATE_FAILED) { if (Date.now() - startTime >= ICE_RESTART_TIMEOUT_MS) { @@ -414,7 +423,7 @@ export class BandwidthRtc { } await new Promise((resolve) => setTimeout(resolve, ICE_RESTART_RETRY_INTERVAL_MS)); // Don't block on this, we should try multiple times - this.offerPublishSdp(true); + retryOffer(); connectionState = pc.connectionState; } } @@ -519,7 +528,20 @@ export class BandwidthRtc { } } - public async init(setMediaPreferencesResponse: SetMediaPreferencesWebRtcResponse) { + public async init(setMediaPreferencesResponse: SetMediaPreferencesWebRtcResponse, isReconnect: boolean = false) { + if (isReconnect) { + // The signaling websocket reconnected (e.g. a gateway-initiated 1001 that expects the + // same endpoint to keep going, not a fresh connect()). The old peer connections are + // stale, so rebuild them and re-publish whatever was already being sent instead of + // silently losing media. + logger.info("Signaling reconnected; rebuilding peer connections and re-publishing existing streams"); + this.publishingPeerConnection?.close(); + this.subscribingPeerConnection?.close(); + this.subscribingPeerConnectionSdpRevision = 0; + this.subscribeTrackMetadata.clear(); + this.localDtmfSenders.clear(); + } + const publishOnTrackHandler = (event: RTCTrackEvent) => { logger.debug("publish ontrack event", event); }; @@ -529,6 +551,16 @@ export class BandwidthRtc { setMediaPreferencesResponse.publishSdpOffer.sdpOffer, ); + if (isReconnect && this.publishedStreams.size > 0) { + // The new publishing peer connection starts with no tracks; re-add whatever was + // already published so the gateway (and far end) keep receiving media instead of + // silence. Per-stream codecPreferences aren't retained across a reconnect. + for (const publishedStream of this.publishedStreams.values()) { + this.addStreamToPublishingPeerConnection(publishedStream.mediaStream); + } + await this.offerPublishSdp(); + } + let streamTracks: Map> = new Map(); const subscriptionOnTrackHandler = (event: RTCTrackEvent) => { @@ -629,9 +661,11 @@ export class BandwidthRtc { const pc = event.target as RTCPeerConnection; const connectionState = pc.connectionState; logger.debug("onconnectionstatechange", connectionState, pc); - if (connectionState === CONNECTION_STATE_FAILED) { + if (connectionState === CONNECTION_STATE_DISCONNECTED) { + logger.warn("Peer disconnected, connection may be reestablished"); + } else if (connectionState === CONNECTION_STATE_FAILED) { logger.warn("Connection failed, ICE restart required"); - await this.retryIceOnFailed(pc, RETRY_ICE_ON_FAILED); + await this.retryIceOnFailed(pc, peerConnectionType, RETRY_ICE_ON_FAILED); } } catch (err) { if (globalThis.window) { @@ -677,21 +711,6 @@ export class BandwidthRtc { } }; - peerConnection.onconnectionstatechange = (event) => { - try { - const pc = event.target as RTCPeerConnection; - logger.debug("onconnectionstatechange", pc.connectionState, pc); - const connectionState = pc.connectionState; - if (connectionState === CONNECTION_STATE_DISCONNECTED) { - logger.warn("Peer disconnected, connection may be reestablished"); - } - } catch (err) { - if (globalThis.window) { - logger.warn("onconnectionstatechange error", err); - } - } - }; - peerConnection.oniceconnectionstatechange = (event) => { try { const pc = event.target as RTCPeerConnection; diff --git a/src/v1/signaling.test.ts b/src/v1/signaling.test.ts index a993a41..4c78616 100644 --- a/src/v1/signaling.test.ts +++ b/src/v1/signaling.test.ts @@ -136,15 +136,19 @@ describe("Signaling websocket event handlers", () => { return ws.on.mock.calls.find((call: any) => call[0] === event)?.[1]; } - test("should emit init and set up ping interval on open", async () => { + test("should emit init with isReconnect false on the first open, then true on subsequent opens", async () => { const emitSpy = jest.spyOn(signaling, "emit"); const openCallback = getWsCallback("open"); expect(openCallback).toBeDefined(); await openCallback(); - expect(emitSpy).toHaveBeenCalledWith("init", expect.anything()); + expect(emitSpy).toHaveBeenCalledWith("init", expect.anything(), false); expect((signaling as any).pingInterval).toBeDefined(); + + await openCallback(); + + expect(emitSpy).toHaveBeenCalledWith("init", expect.anything(), true); }); test("should reject with error and disconnect on 403 error", async () => { diff --git a/src/v1/signaling.ts b/src/v1/signaling.ts index cd369bf..132792d 100644 --- a/src/v1/signaling.ts +++ b/src/v1/signaling.ts @@ -52,6 +52,10 @@ class Signaling extends EventEmitter { connect(authParams: RtcAuthParams, options?: RtcOptions) { let rpc_id = 1; + // rpc-websockets auto-reconnects with a brand new underlying WebSocket (same + // JsonRpcClient instance), so "open" fires again on every reconnect. Scoped to + // this connect() call so a fresh top-level connect() always starts as false. + let hasConnectedOnce = false; return new Promise((resolve, reject) => { if (this.ws) { @@ -95,16 +99,18 @@ class Signaling extends EventEmitter { ws.on("open", async () => { logger.debug("Websocket open"); - if (globalThis.addEventListener) { + const isReconnect = hasConnectedOnce; + hasConnectedOnce = true; + if (!isReconnect && globalThis.addEventListener) { globalThis.addEventListener("beforeunload", (event) => { this.disconnect(); }); } - // TODO: handle reconnections let preferencesResponse = await this.setMediaPreferences(); // logger.debug(`Media preferences set`, preferencesResponse); - // Setup Peers - this.emit("init", preferencesResponse); + // Setup Peers. isReconnect tells the caller whether existing peer connections/media + // need to be rebuilt and re-published, rather than created for the first time. + this.emit("init", preferencesResponse, isReconnect); this.pingInterval = setInterval(() => { ws.call("ping", {}); From 7a1e4fd126625b275adbfee64d3c7297c895d7d5 Mon Sep 17 00:00:00 2001 From: smoghe-bw Date: Tue, 8 Sep 2026 16:24:14 -0400 Subject: [PATCH 2/6] refactor: convert ICE restart retry to async/await with try-catch --- src/v1/bandwidthRtc.ts | 8 +++++++- 1 file changed, 7 insertions(+), 1 deletion(-) diff --git a/src/v1/bandwidthRtc.ts b/src/v1/bandwidthRtc.ts index 16b156b..56ab069 100644 --- a/src/v1/bandwidthRtc.ts +++ b/src/v1/bandwidthRtc.ts @@ -412,7 +412,13 @@ export class BandwidthRtc { const ICE_RESTART_RETRY_INTERVAL_MS = 5_000; const startTime = Date.now(); - const retryOffer = () => this.offerPublishSdp(true).catch((err) => logger.warn("ICE restart offer failed", err)); + const retryOffer = async () => { + try { + await this.offerPublishSdp(true); + } catch (err) { + logger.warn("ICE restart offer failed", err); + } + }; await retryOffer(); let connectionState = pc.connectionState; From 2b12424250165cecebf9f0935ad93218d7d39e60 Mon Sep 17 00:00:00 2001 From: smoghe-bw Date: Tue, 8 Sep 2026 16:30:36 -0400 Subject: [PATCH 3/6] fix(bandwidthRtc): always log caught peer-connection handler errors Every handler in setupNewPeerConnection (and the merged onconnectionstatechange handler in setupPeerConnection) caught errors but only logged them when globalThis.window was set, silently discarding them in any non-browser environment. logger.warn has no browser dependency (just console + EventEmitter), so there was nothing this guard was protecting against - it just meant every error here vanished with zero trace outside a browser. Log unconditionally. Also switched retryOffer in retryIceOnFailed to an async/await function with braces instead of an implicit-return arrow expression, matching the project's style. Co-Authored-By: Claude Sonnet 5 --- src/v1/bandwidthRtc.ts | 20 +++++--------------- 1 file changed, 5 insertions(+), 15 deletions(-) diff --git a/src/v1/bandwidthRtc.ts b/src/v1/bandwidthRtc.ts index 56ab069..b683389 100644 --- a/src/v1/bandwidthRtc.ts +++ b/src/v1/bandwidthRtc.ts @@ -674,9 +674,7 @@ export class BandwidthRtc { await this.retryIceOnFailed(pc, peerConnectionType, RETRY_ICE_ON_FAILED); } } catch (err) { - if (globalThis.window) { - logger.warn("onconnectionstatechange error", err); - } + logger.warn("onconnectionstatechange error", err); } }; logger.debug("Initial SDP offer", initialSdpOffer); @@ -722,9 +720,7 @@ export class BandwidthRtc { const pc = event.target as RTCPeerConnection; logger.debug("oniceconnectionstatechange", pc.iceConnectionState, pc); } catch (err) { - if (globalThis.window) { - logger.warn("oniceconnectionstatechange error", err); - } + logger.warn("oniceconnectionstatechange error", err); } }; @@ -733,9 +729,7 @@ export class BandwidthRtc { const pc = event.target as RTCPeerConnection; logger.debug("onicegatheringstatechange", pc.iceGatheringState, pc); } catch (err) { - if (globalThis.window) { - logger.warn("onicegatheringstatechange error", err); - } + logger.warn("onicegatheringstatechange error", err); } }; @@ -743,9 +737,7 @@ export class BandwidthRtc { try { logger.debug("onnegotiationneeded", event.target); } catch (err) { - if (globalThis.window) { - logger.warn("onnegotiationneeded error", err); - } + logger.warn("onnegotiationneeded error", err); } }; @@ -754,9 +746,7 @@ export class BandwidthRtc { const pc = event.target as RTCPeerConnection; logger.debug("onsignalingstatechange", pc.signalingState, pc); } catch (err) { - if (globalThis.window) { - logger.warn("onsignalingstatechange error", err); - } + logger.warn("onsignalingstatechange error", err); } }; From 9285f352ff90ec82f1d5b1451c68a09af685b6fb Mon Sep 17 00:00:00 2001 From: smoghe-bw Date: Tue, 8 Sep 2026 16:37:27 -0400 Subject: [PATCH 4/6] fix(signaling): log rejected RPC calls instead of failing silently requestOutboundConnection, hangupConnection, acceptStream, declineStream, offerSdp, and answerSdp only logged the outgoing call. If the gateway rejected the RPC, the rejection propagated with no trace in this SDK's own logs - silent unless the calling application happened to catch and log it itself. Log a warning on rejection and rethrow, so callers still see the same rejected promise. Co-Authored-By: Claude Sonnet 5 --- src/v1/signaling.ts | 72 ++++++++++++++++++++++++++++++++------------- 1 file changed, 51 insertions(+), 21 deletions(-) diff --git a/src/v1/signaling.ts b/src/v1/signaling.ts index 132792d..dd76696 100644 --- a/src/v1/signaling.ts +++ b/src/v1/signaling.ts @@ -223,43 +223,73 @@ class Signaling extends EventEmitter { this._disconnect(true); } - requestOutboundConnection(id: string, type: EndpointType): Promise { + async requestOutboundConnection(id: string, type: EndpointType): Promise { logger.debug(`Calling "requestOutboundConnection"`, { id: id, type: type }); - return this.ws?.call("requestOutboundConnection", { - id: id, - type: type, - }) as Promise; + try { + return await (this.ws?.call("requestOutboundConnection", { + id: id, + type: type, + }) as Promise); + } catch (err) { + logger.warn(`"requestOutboundConnection" rejected`, err); + throw err; + } } - hangupConnection(endpoint: string, type: EndpointType): Promise { + async hangupConnection(endpoint: string, type: EndpointType): Promise { logger.debug(`Calling "hangupConnection"`, { endpoint: endpoint, type: type }); - return this.ws?.call("hangupConnection", { - endpoint: endpoint, - type: type, - }) as Promise; + try { + return await (this.ws?.call("hangupConnection", { + endpoint: endpoint, + type: type, + }) as Promise); + } catch (err) { + logger.warn(`"hangupConnection" rejected`, err); + throw err; + } } - acceptStream(): Promise { + async acceptStream(): Promise { logger.debug(`Calling "acceptStream"`); - return this.ws?.call("acceptStream", {}) as Promise; + try { + return await (this.ws?.call("acceptStream", {}) as Promise); + } catch (err) { + logger.warn(`"acceptStream" rejected`, err); + throw err; + } } - declineStream(): Promise { + async declineStream(): Promise { logger.debug(`Calling "declineStream"`); - return this.ws?.call("declineStream", {}) as Promise; + try { + return await (this.ws?.call("declineStream", {}) as Promise); + } catch (err) { + logger.warn(`"declineStream" rejected`, err); + throw err; + } } - offerSdp(peerType: string, sdpOffer: string): Promise { + async offerSdp(peerType: string, sdpOffer: string): Promise { logger.debug(`Calling "offerSdp"`, { sdpOffer: sdpOffer, peerType: peerType }); - return this.ws?.call("offerSdp", { sdpOffer: sdpOffer, peerType: peerType }) as Promise; + try { + return await (this.ws?.call("offerSdp", { sdpOffer: sdpOffer, peerType: peerType }) as Promise); + } catch (err) { + logger.warn(`"offerSdp" rejected`, err); + throw err; + } } - answerSdp(sdpAnswer: string, peerType: string): Promise { + async answerSdp(sdpAnswer: string, peerType: string): Promise { logger.debug(`Calling "answerSdp"`, { sdpAnswer: sdpAnswer }); - return this.ws?.call("answerSdp", { - peerType: peerType, - sdpAnswer: sdpAnswer, - }) as Promise; + try { + return await (this.ws?.call("answerSdp", { + peerType: peerType, + sdpAnswer: sdpAnswer, + }) as Promise); + } catch (err) { + logger.warn(`"answerSdp" rejected`, err); + throw err; + } } private sendDiagnostics(diagnostics: Diagnostics): Promise { From ded48aed052d07ed923c4f0b6f8b29629925cdbe Mon Sep 17 00:00:00 2001 From: smoghe-bw Date: Wed, 9 Sep 2026 14:34:39 -0400 Subject: [PATCH 5/6] revert(signaling): drop async/try-catch wrapping from RPC passthrough methods requestOutboundConnection, hangupConnection, acceptStream, declineStream, offerSdp, and answerSdp were wrapped in async/try-catch+rethrow, out of scope for this PR (signaling reconnect + ICE restart retry) and a regression: wrapping a plain `this.ws?.call(...)` passthrough in `async` means `await undefined` (when `this.ws` is null) resolves silently instead of leaving the caller with a non-Promise value that blows up immediately on `.then()`. Reverts to the direct passthrough. Co-Authored-By: Claude Sonnet 5 --- src/v1/signaling.ts | 72 +++++++++++++-------------------------------- 1 file changed, 21 insertions(+), 51 deletions(-) diff --git a/src/v1/signaling.ts b/src/v1/signaling.ts index dd76696..132792d 100644 --- a/src/v1/signaling.ts +++ b/src/v1/signaling.ts @@ -223,73 +223,43 @@ class Signaling extends EventEmitter { this._disconnect(true); } - async requestOutboundConnection(id: string, type: EndpointType): Promise { + requestOutboundConnection(id: string, type: EndpointType): Promise { logger.debug(`Calling "requestOutboundConnection"`, { id: id, type: type }); - try { - return await (this.ws?.call("requestOutboundConnection", { - id: id, - type: type, - }) as Promise); - } catch (err) { - logger.warn(`"requestOutboundConnection" rejected`, err); - throw err; - } + return this.ws?.call("requestOutboundConnection", { + id: id, + type: type, + }) as Promise; } - async hangupConnection(endpoint: string, type: EndpointType): Promise { + hangupConnection(endpoint: string, type: EndpointType): Promise { logger.debug(`Calling "hangupConnection"`, { endpoint: endpoint, type: type }); - try { - return await (this.ws?.call("hangupConnection", { - endpoint: endpoint, - type: type, - }) as Promise); - } catch (err) { - logger.warn(`"hangupConnection" rejected`, err); - throw err; - } + return this.ws?.call("hangupConnection", { + endpoint: endpoint, + type: type, + }) as Promise; } - async acceptStream(): Promise { + acceptStream(): Promise { logger.debug(`Calling "acceptStream"`); - try { - return await (this.ws?.call("acceptStream", {}) as Promise); - } catch (err) { - logger.warn(`"acceptStream" rejected`, err); - throw err; - } + return this.ws?.call("acceptStream", {}) as Promise; } - async declineStream(): Promise { + declineStream(): Promise { logger.debug(`Calling "declineStream"`); - try { - return await (this.ws?.call("declineStream", {}) as Promise); - } catch (err) { - logger.warn(`"declineStream" rejected`, err); - throw err; - } + return this.ws?.call("declineStream", {}) as Promise; } - async offerSdp(peerType: string, sdpOffer: string): Promise { + offerSdp(peerType: string, sdpOffer: string): Promise { logger.debug(`Calling "offerSdp"`, { sdpOffer: sdpOffer, peerType: peerType }); - try { - return await (this.ws?.call("offerSdp", { sdpOffer: sdpOffer, peerType: peerType }) as Promise); - } catch (err) { - logger.warn(`"offerSdp" rejected`, err); - throw err; - } + return this.ws?.call("offerSdp", { sdpOffer: sdpOffer, peerType: peerType }) as Promise; } - async answerSdp(sdpAnswer: string, peerType: string): Promise { + answerSdp(sdpAnswer: string, peerType: string): Promise { logger.debug(`Calling "answerSdp"`, { sdpAnswer: sdpAnswer }); - try { - return await (this.ws?.call("answerSdp", { - peerType: peerType, - sdpAnswer: sdpAnswer, - }) as Promise); - } catch (err) { - logger.warn(`"answerSdp" rejected`, err); - throw err; - } + return this.ws?.call("answerSdp", { + peerType: peerType, + sdpAnswer: sdpAnswer, + }) as Promise; } private sendDiagnostics(diagnostics: Diagnostics): Promise { From 6a15d99c39e1e46d5826b26fa3060da434c54957 Mon Sep 17 00:00:00 2001 From: smoghe-bw Date: Wed, 23 Sep 2026 15:48:53 -0400 Subject: [PATCH 6/6] fix(v1): route gateway sdpOffer by peerType; leave ICE restart to the gateway - Disable the client ICE-restart retry again: the gateway rejects client offers until its peer is connected, so a retry on a failed peer can never succeed, and the gateway runs its own ICE restart. - Apply "publish" sdpOffer notifications (the gateway's ICE-restart offer) to the publishing connection under publishMutex and answer them as "publish". Before this every sdpOffer went to the subscribing connection. A missing peerType still means subscribe. - Track a separate publish SDP revision and reset it on every init(). - Stop a stale retry loop once init() has replaced the publishing connection. Co-Authored-By: Claude Opus 5.5 (1M context) --- src/v1/bandwidthRtc.test.ts | 156 ++++++++++++++++++++++++++++++++++++ src/v1/bandwidthRtc.ts | 72 ++++++++++++++++- src/v1/types.ts | 2 + 3 files changed, 227 insertions(+), 3 deletions(-) diff --git a/src/v1/bandwidthRtc.test.ts b/src/v1/bandwidthRtc.test.ts index 9f0e022..f03df0f 100644 --- a/src/v1/bandwidthRtc.test.ts +++ b/src/v1/bandwidthRtc.test.ts @@ -697,6 +697,8 @@ describe("bandwidthRtcV1 retryIceOnFailed", () => { test("re-offers once and stops once the publish connection recovers", async () => { const brtc = new BandwidthRtc(); const pc = makePc("failed"); + // retryIceOnFailed only acts on the current publishing connection. + (brtc as any).publishingPeerConnection = pc; jest.spyOn(brtc as any, "offerPublishSdp").mockImplementation(async () => { (pc as any).connectionState = "connected"; return {} as any; @@ -710,6 +712,7 @@ describe("bandwidthRtcV1 retryIceOnFailed", () => { test("retries every 5s until the timeout elapses if still failed", async () => { const brtc = new BandwidthRtc(); const pc = makePc("failed"); + (brtc as any).publishingPeerConnection = pc; jest.spyOn(brtc as any, "offerPublishSdp").mockResolvedValue({} as any); const done = (brtc as any).retryIceOnFailed(pc, "publish", true); @@ -726,6 +729,7 @@ describe("bandwidthRtcV1 retryIceOnFailed", () => { test("does not throw when a retry's offerPublishSdp rejects", async () => { const brtc = new BandwidthRtc(); const pc = makePc("failed"); + (brtc as any).publishingPeerConnection = pc; jest.spyOn(brtc as any, "offerPublishSdp").mockRejectedValue(new Error("signaling down")); const done = (brtc as any).retryIceOnFailed(pc, "publish", true); @@ -736,6 +740,18 @@ describe("bandwidthRtcV1 retryIceOnFailed", () => { await expect(done).resolves.not.toThrow(); }); + + test("sends no offer when pc is no longer the publishing peer connection (e.g. init() replaced it)", async () => { + const brtc = new BandwidthRtc(); + const pc = makePc("failed"); + // A different object is now the live publishing connection. + (brtc as any).publishingPeerConnection = {}; + const offerPublishSdp = jest.spyOn(brtc as any, "offerPublishSdp"); + + await (brtc as any).retryIceOnFailed(pc, "publish", true); + + expect(offerPublishSdp).not.toHaveBeenCalled(); + }); }); describe("bandwidthRtcV1 init on signaling reconnect", () => { @@ -808,4 +824,144 @@ describe("bandwidthRtcV1 init on signaling reconnect", () => { expect(offerPublishSdp).not.toHaveBeenCalled(); }); + + test("every init resets the publishing peer's ICE-restart revision, not only a reconnect", async () => { + const brtc = new BandwidthRtc(); + jest.spyOn(brtc as any, "setupPeerConnection").mockResolvedValue({ close: jest.fn() }); + (brtc as any).publishingPeerConnectionSdpRevision = 3; + + await brtc.init(makePreferencesResponse()); + + expect((brtc as any).publishingPeerConnectionSdpRevision).toBe(0); + }); +}); + +describe("bandwidthRtcV1 sdpOffer peerType routing", () => { + beforeAll(() => { + setupNavigatorMocks(); + setupMocks(); + }); + + function makePc() { + return { + setRemoteDescription: jest.fn().mockResolvedValue(undefined), + createAnswer: jest.fn().mockResolvedValue({ sdp: "answer-sdp" }), + setLocalDescription: jest.fn().mockResolvedValue(undefined), + } as any; + } + + function withMockedAnswerSdp(brtc: BandwidthRtc) { + const answerSdp = jest.fn().mockResolvedValue(undefined); + (brtc as any).signaling.answerSdp = answerSdp; + return answerSdp; + } + + test("publish peerType offer is applied to the publishing connection and answered with answerSdp(sdp, publish); subscribing connection untouched", async () => { + const brtc = new BandwidthRtc(); + const publishPc = makePc(); + const subscribePc = makePc(); + (brtc as any).publishingPeerConnection = publishPc; + (brtc as any).subscribingPeerConnection = subscribePc; + const answerSdp = withMockedAnswerSdp(brtc); + + await (brtc as any).handleSdpOffer({ peerType: "publish", sdpOffer: "offer-sdp", sdpRevision: 1 }); + + expect(publishPc.setRemoteDescription).toHaveBeenCalledWith({ type: "offer", sdp: "offer-sdp" }); + expect(publishPc.setLocalDescription).toHaveBeenCalledWith({ sdp: "answer-sdp" }); + expect(answerSdp).toHaveBeenCalledWith("answer-sdp", "publish"); + expect(subscribePc.setRemoteDescription).not.toHaveBeenCalled(); + }); + + test.each([["subscribe"], [undefined]])("peerType %s routes to the subscribing connection (existing behaviour)", async (peerType) => { + const brtc = new BandwidthRtc(); + const publishPc = makePc(); + const subscribePc = makePc(); + (brtc as any).publishingPeerConnection = publishPc; + (brtc as any).subscribingPeerConnection = subscribePc; + const answerSdp = withMockedAnswerSdp(brtc); + + await (brtc as any).handleSdpOffer({ peerType, sdpOffer: "offer-sdp", sdpRevision: 1 }); + + expect(subscribePc.setRemoteDescription).toHaveBeenCalledWith({ type: "offer", sdp: "offer-sdp" }); + expect(answerSdp).toHaveBeenCalledWith("answer-sdp", "subscribe"); + expect(publishPc.setRemoteDescription).not.toHaveBeenCalled(); + }); + + test("a publish offer with a revision <= the last applied one is ignored", async () => { + const brtc = new BandwidthRtc(); + const publishPc = makePc(); + (brtc as any).publishingPeerConnection = publishPc; + (brtc as any).publishingPeerConnectionSdpRevision = 2; + withMockedAnswerSdp(brtc); + + await (brtc as any).handleSdpOffer({ peerType: "publish", sdpOffer: "offer-sdp", sdpRevision: 2 }); + + expect(publishPc.setRemoteDescription).not.toHaveBeenCalled(); + }); + + test("init resets the publish revision, so a revision-1 offer is applied again after a second init", async () => { + const brtc = new BandwidthRtc(); + jest.spyOn(brtc as any, "setupPeerConnection").mockResolvedValue({ close: jest.fn() }); + await brtc.init({ publishSdpOffer: {}, subscribeSdpOffer: {} } as any); + // Simulate a restart already applied on the first publishing connection. + (brtc as any).publishingPeerConnectionSdpRevision = 1; + + await brtc.init({ publishSdpOffer: {}, subscribeSdpOffer: {} } as any); + expect((brtc as any).publishingPeerConnectionSdpRevision).toBe(0); + + const publishPc = makePc(); + (brtc as any).publishingPeerConnection = publishPc; + withMockedAnswerSdp(brtc); + + await (brtc as any).handleSdpOffer({ peerType: "publish", sdpOffer: "offer-sdp", sdpRevision: 1 }); + + expect(publishPc.setRemoteDescription).toHaveBeenCalled(); + }); + + test("a publish offer waits for an in-flight publish negotiation holding publishMutex", async () => { + const brtc = new BandwidthRtc(); + const publishPc = makePc(); + (brtc as any).publishingPeerConnection = publishPc; + withMockedAnswerSdp(brtc); + + let releaseMutex: () => void = () => {}; + const heldMutexTask = new Promise((resolve) => { + releaseMutex = resolve; + }); + const mutexPromise = (brtc as any).publishMutex.runExclusive(() => heldMutexTask); + + const offerPromise = (brtc as any).handleSdpOffer({ peerType: "publish", sdpOffer: "offer-sdp", sdpRevision: 1 }); + + // Give handleSdpOffer a chance to run; it should still be blocked on the mutex. + await new Promise((resolve) => setTimeout(resolve, 10)); + expect(publishPc.setRemoteDescription).not.toHaveBeenCalled(); + + releaseMutex(); + await mutexPromise; + await offerPromise; + + expect(publishPc.setRemoteDescription).toHaveBeenCalled(); + }); +}); + +describe("bandwidthRtcV1 onconnectionstatechange wiring", () => { + beforeAll(() => { + setupNavigatorMocks(); + setupMocks(); + }); + + test("connection 'failed' does not send an offer (retry disabled)", async () => { + const brtc = new BandwidthRtc(); + const fakePc: any = { connectionState: "connected" }; + jest.spyOn(brtc as any, "createPeerConnection").mockReturnValue(fakePc); + const offerPublishSdp = jest.spyOn(brtc as any, "offerPublishSdp").mockResolvedValue({}); + + await (brtc as any).setupPeerConnection("publish", () => {}); + (brtc as any).publishingPeerConnection = fakePc; + + fakePc.connectionState = "failed"; + await fakePc.onconnectionstatechange({ target: fakePc }); + + expect(offerPublishSdp).not.toHaveBeenCalled(); + }); }); diff --git a/src/v1/bandwidthRtc.ts b/src/v1/bandwidthRtc.ts index 0263c99..e5554d5 100644 --- a/src/v1/bandwidthRtc.ts +++ b/src/v1/bandwidthRtc.ts @@ -63,8 +63,8 @@ const CONNECTION_STATE_CLOSED = "closed"; const PUBLISH_ICE_CONNECT_TIMEOUT_MS = 10_000; const PUBLISH_ICE_CONNECT_POLL_INTERVAL_MS = 100; -// When true, automatically trigger an ICE restart (via offerPublishSdp(true)) on connection failure. -const RETRY_ICE_ON_FAILED = true; +// The gateway owns ICE restart and rejects client offers until the peer is connected again. +const RETRY_ICE_ON_FAILED = false; export class BandwidthRtc { private options?: RtcOptions; @@ -98,6 +98,8 @@ export class BandwidthRtc { // Current SDP revision for the subscribing peer; used to reject outdated SDP offers private subscribingPeerConnectionSdpRevision = 0; + // Current SDP revision for the publishing peer's gateway-initiated ICE restarts; same purpose. + private publishingPeerConnectionSdpRevision = 0; private localDtmfSenders: Map = new Map(); @@ -142,7 +144,7 @@ export class BandwidthRtc { logger.info("Connecting to Bandwidth WebRTC"); this.signaling.on("ready", this.handleReady.bind(this)); - this.signaling.on("sdpOffer", this.handleSubscribeSdpOffer.bind(this)); + this.signaling.on("sdpOffer", this.handleSdpOffer.bind(this)); this.signaling.on("init", this.init.bind(this)); this.signaling.on("fatalError", this.handleError.bind(this)); @@ -438,6 +440,12 @@ export class BandwidthRtc { logger.warn(`ICE restart on the ${peerConnectionType} peer connection requires the gateway to re-offer; client cannot initiate`); return; } + if (pc !== this.publishingPeerConnection) { + // init() already replaced the publishing connection (e.g. a reconnect ran while this + // failed pc's loop was still active); it would otherwise send ICE-restart offers on the + // new, healthy connection instead of the one that actually failed. + return; + } const ICE_RESTART_TIMEOUT_MS = 30_000; const ICE_RESTART_RETRY_INTERVAL_MS = 5_000; @@ -458,6 +466,9 @@ export class BandwidthRtc { logger.warn("ICE restart timed out"); break; } + if (pc !== this.publishingPeerConnection) { + return; + } await new Promise((resolve) => setTimeout(resolve, ICE_RESTART_RETRY_INTERVAL_MS)); // Don't block on this, we should try multiple times retryOffer(); @@ -519,6 +530,58 @@ export class BandwidthRtc { } } + // Routes a gateway sdpOffer notification by peerType: "publish" is a gateway-initiated ICE + // restart on the publishing connection; anything else (including missing peerType, for older + // gateways) is the existing subscribe renegotiation. + private async handleSdpOffer(offer: SubscribeSdpOffer): Promise { + if (offer.peerType === PEER_CONNECTION_TYPE_PUBLISH) { + await this.handlePublishSdpOffer(offer); + } else { + await this.handleSubscribeSdpOffer(offer); + } + } + + // A gateway-initiated ICE restart on the publishing connection (see handleSdpOffer). Runs + // under publishMutex so it can't interleave with the SDK's own offerPublishSdp negotiations. + private async handlePublishSdpOffer(offer: SubscribeSdpOffer): Promise { + try { + await this.publishMutex.runExclusive(async () => { + logger.info("Received publish SDP offer", offer); + logger.debug("Current publish SDP revision", this.publishingPeerConnectionSdpRevision); + if (offer.sdpRevision <= this.publishingPeerConnectionSdpRevision) { + logger.debug( + `Revision on publish SDP offer (${offer.sdpRevision}) is less than current revision (${this.publishingPeerConnectionSdpRevision}), ignoring`, + ); + return; + } + + if (!this.publishingPeerConnection) { + throw new BandwidthRtcError("No publishing RTCPeerConnection, cannot handle SDP offer"); + } + + await this.publishingPeerConnection!.setRemoteDescription({ + type: "offer", + sdp: offer.sdpOffer, + }); + + let localSdpAnswer = await this.publishingPeerConnection!.createAnswer(); + if (!localSdpAnswer.sdp) { + throw new BandwidthRtcError(`RTCPeerConnection.createAnswer returned ${JSON.stringify(localSdpAnswer)}`); + } + + await this.publishingPeerConnection!.setLocalDescription(localSdpAnswer); + await this.signaling.answerSdp(localSdpAnswer.sdp, PEER_CONNECTION_TYPE_PUBLISH); + + this.publishingPeerConnectionSdpRevision = offer.sdpRevision; + logger.debug(`set current publish SDP revision to ${this.publishingPeerConnectionSdpRevision}`); + }); + } catch (err) { + // A failed ICE-restart answer leaves publish media down, unlike a dropped subscribe + // renegotiation, so this warrants more than debug-level logging. + logger.warn("error in handlePublishSdpOffer", err); + } + } + private async handleSubscribeSdpOffer(subscribeSdpOffer: SubscribeSdpOffer): Promise { try { await this.subscribeMutex.runExclusive(async () => { @@ -585,6 +648,9 @@ export class BandwidthRtc { const publishOnTrackHandler = (event: RTCTrackEvent) => { logger.debug("publish ontrack event", event); }; + // Every init() builds a new gateway-side publishing connection whose ICE-restart revision + // counter starts at 0, not only on reconnect, so this resets unconditionally. + this.publishingPeerConnectionSdpRevision = 0; // A reconnect re-runs init() against a fresh offer; the previous peer connection (if any) // is already dead on the far end, but nothing local closes it, so it leaks otherwise. this.publishingPeerConnection?.close(); diff --git a/src/v1/types.ts b/src/v1/types.ts index 0c2f898..8325f99 100644 --- a/src/v1/types.ts +++ b/src/v1/types.ts @@ -19,6 +19,8 @@ export interface SdpAnswer { export interface SubscribeSdpOffer { sdpOffer: string; sdpRevision: number; + /** "publish" routes to the publishing connection; absent means subscribe, for older gateways. */ + peerType?: string; /** * Per-track stream metadata keyed by track id, sent alongside the offer that * adds a call's subscribe track. The SDK derives stream-available from the