From cca0dbc6069d7a7ce52a19cdc290317b14fbdb6a Mon Sep 17 00:00:00 2001 From: Artem Grintsevich Date: Tue, 22 Sep 2026 11:32:15 +0200 Subject: [PATCH 1/2] fix: fixed pc close/dispose race --- .../WebRTCModule/PeerConnectionObserver.java | 50 ++++- .../com/oney/WebRTCModule/WebRTCModule.java | 174 +++++++++++++----- 2 files changed, 167 insertions(+), 57 deletions(-) diff --git a/android/src/main/java/com/oney/WebRTCModule/PeerConnectionObserver.java b/android/src/main/java/com/oney/WebRTCModule/PeerConnectionObserver.java index d060ff3b1..29156265b 100644 --- a/android/src/main/java/com/oney/WebRTCModule/PeerConnectionObserver.java +++ b/android/src/main/java/com/oney/WebRTCModule/PeerConnectionObserver.java @@ -66,12 +66,25 @@ void setPeerConnection(PeerConnection peerConnection) { } void close() { - Log.d(TAG, "PeerConnection.close() for " + id); + final PeerConnection pc = peerConnection; + if (pc == null) { + Log.d(TAG, "PeerConnection.close() for " + id + " ignored; already disposed"); + return; + } - peerConnection.close(); + Log.d(TAG, "PeerConnection.close() for " + id); + pc.close(); } void dispose() { + final PeerConnection pc = peerConnection; + peerConnection = null; + + if (pc == null) { + Log.d(TAG, "PeerConnection.dispose() for " + id + " ignored; already disposed"); + return; + } + Log.d(TAG, "PeerConnection.dispose() for " + id); // Remove track adapters for remote tracks @@ -85,7 +98,7 @@ void dispose() { } // Remove video track adapters for local tracks (from senders) - for (RtpSender sender : this.peerConnection.getSenders()) { + for (RtpSender sender : pc.getSenders()) { MediaStreamTrack track = sender.track(); if (track instanceof VideoTrack) { videoTrackAdapters.removeAdapter((VideoTrack) track); @@ -102,7 +115,7 @@ void dispose() { // At this point there should be no local MediaStreams in the associated // PeerConnection. Call dispose() to free all remaining resources held // by the PeerConnection instance (RtpReceivers, RtpSenders, etc.) - peerConnection.dispose(); + pc.dispose(); videoTrackAdapters.dispose(); @@ -305,7 +318,13 @@ public void onIceCandidate(final IceCandidate candidate) { params.putMap("candidate", candidateParams); - SessionDescription newSdp = peerConnection.getLocalDescription(); + final PeerConnection pc = peerConnection; + if (pc == null) { + Log.d(TAG, "onIceCandidate for " + id + " skipped; peer connection already disposed"); + return; + } + + SessionDescription newSdp = pc.getLocalDescription(); WritableMap newSdpMap = Arguments.createMap(); // Can happen when doing a rollback. @@ -356,7 +375,13 @@ public void onIceGatheringChange(PeerConnection.IceGatheringState iceGatheringSt params.putString("iceGatheringState", iceGatheringStateString(iceGatheringState)); if (iceGatheringState == PeerConnection.IceGatheringState.COMPLETE) { - SessionDescription newSdp = peerConnection.getLocalDescription(); + final PeerConnection pc = peerConnection; + if (pc == null) { + Log.d(TAG, + "onIceGatheringChange for " + id + " skipped; peer connection already disposed"); + return; + } + SessionDescription newSdp = pc.getLocalDescription(); WritableMap newSdpMap = Arguments.createMap(); // Can happen when doing a rollback. @@ -429,8 +454,14 @@ public void onAddTrack(final RtpReceiver receiver, final MediaStream[] mediaStre Log.d(TAG, "onAddTrack"); ThreadUtils.runOnExecutor(() -> { + final PeerConnection pc = peerConnection; + if (pc == null) { + Log.d(TAG, "onAddTrack for " + id + " skipped; peer connection already disposed"); + return; + } + RtpTransceiver transceiver = null; - for (RtpTransceiver t : this.peerConnection.getTransceivers()) { + for (RtpTransceiver t : pc.getTransceivers()) { if (Objects.equals(t.getReceiver().id(), receiver.id())) { transceiver = t; break; @@ -502,6 +533,11 @@ public void onTrack(final RtpTransceiver transceiver) {} @Override public void onRemoveTrack(RtpReceiver receiver) { ThreadUtils.runOnExecutor(() -> { + if (peerConnection == null) { + Log.d(TAG, "onRemoveTrack for " + id + " skipped; peer connection already disposed"); + return; + } + // Tear down track adapters so a subsequent onAddTrack with the // same trackId (SFU participant rejoin) creates a fresh adapter // on the new MediaStreamTrack object. Without this, the old sink diff --git a/android/src/main/java/com/oney/WebRTCModule/WebRTCModule.java b/android/src/main/java/com/oney/WebRTCModule/WebRTCModule.java index 4c903d8bc..b311fe9ba 100644 --- a/android/src/main/java/com/oney/WebRTCModule/WebRTCModule.java +++ b/android/src/main/java/com/oney/WebRTCModule/WebRTCModule.java @@ -6,6 +6,7 @@ import androidx.annotation.NonNull; import androidx.annotation.Nullable; +import androidx.core.util.Consumer; import com.facebook.react.bridge.Arguments; import com.facebook.react.bridge.Callback; @@ -134,15 +135,27 @@ public void createCallFactory(ReadableMap options, Promise promise) { && options.getBoolean("bypassVoiceProcessing"); boolean stereoInputEnabled = options != null && options.hasKey("stereoInputEnabled") && options.getBoolean("stereoInputEnabled"); - - // This makes default factory being disposed in a proper sequence. + + final Runnable create = () -> { + try { + factoryRegistry.create(bypassVoiceProcessing, stereoInputEnabled); + promise.resolve(null); + } catch (Exception e) { + Log.e(TAG, "createCallFactory() failed", e); + promise.reject("E_FACTORY_CREATE", e); + } + }; + + // Tear a stale bare-fork default down in order first. The teardown is two-phase, so + // the new factory must be built from the completion callback — building it inline + // would create it while the old factory's PeerConnections were still alive. if (factoryRegistry.isBareForkDefaultLive()) { Log.d(TAG, "createCallFactory(): tearing down stale bare-fork default (ordered) " + "before creating the call factory"); - disposeCurrentFactoryOrdered(); + disposeCurrentFactoryOrdered(disposed -> create.run()); + } else { + create.run(); } - factoryRegistry.create(bypassVoiceProcessing, stereoInputEnabled); - promise.resolve(null); } catch (Exception e) { Log.e(TAG, "createCallFactory() failed", e); promise.reject("E_FACTORY_CREATE", e); @@ -152,7 +165,7 @@ public void createCallFactory(ReadableMap options, Promise promise) { @ReactMethod public void disposeCallFactory(Promise promise) { - ThreadUtils.runOnExecutor(() -> promise.resolve(disposeCurrentFactoryOrdered())); + ThreadUtils.runOnExecutor(() -> disposeCurrentFactoryOrdered(promise::resolve)); } /** @@ -161,47 +174,70 @@ public void disposeCallFactory(Promise promise) { * PCs/tracks is a use-after-free); streams go first so {@code removeTrack()} runs while their * tracks are still alive. Also clears {@code localStreams}, otherwise only released in * {@link #invalidate()} (else it leaks across join/leave). No-op unless this is the last - * reference; returns whether it disposed the factory. + * reference; {@code onDisposed} receives whether the factory was actually disposed. + * + *

Split across two executor tasks: phase 1 only closes the PeerConnections, phase 2 frees + * them. libwebrtc callbacks arrive on its own threads and hand their work to this executor by + * appending a task to its queue, holding raw {@code PeerConnection} / {@code RtpReceiver} + * handles. Phase 2 is appended to that same queue at the end of phase 1, so every callback + * queued up to that point sits ahead of it and runs first — while the PeerConnections are still + * alive. Freeing inline at the end of phase 1 would instead run before that backlog drained, + * leaving those callbacks to dereference freed memory. */ - private boolean disposeCurrentFactoryOrdered() { + private void disposeCurrentFactoryOrdered(Consumer onDisposed) { if (!factoryRegistry.releaseReference()) { - return false; + onDisposed.accept(false); + return; } - for (Map.Entry entry : localStreams.entrySet()) { + for (int pcId : factoryRegistry.currentOwnedPcIds()) { try { - MediaStream stream = entry.getValue(); - for (AudioTrack t : new ArrayList<>(stream.audioTracks)) stream.removeTrack(t); - for (VideoTrack t : new ArrayList<>(stream.videoTracks)) stream.removeTrack(t); - stream.dispose(); + PeerConnectionObserver pco = mPeerConnectionObservers.get(pcId); + if (pco != null) { + pco.close(); + } } catch (Exception e) { - Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing stream " + entry.getKey(), e); + Log.w(TAG, "disposeCurrentFactoryOrdered(): error closing pc " + pcId, e); } } - localStreams.clear(); - for (int pcId : factoryRegistry.currentOwnedPcIds()) { - try { - PeerConnectionObserver pco = mPeerConnectionObservers.get(pcId); - if (pco != null && pco.getPeerConnection() != null) { - pco.dispose(); + ThreadUtils.runOnExecutor(() -> { + for (Map.Entry entry : localStreams.entrySet()) { + try { + MediaStream stream = entry.getValue(); + for (AudioTrack t : new ArrayList<>(stream.audioTracks)) stream.removeTrack(t); + for (VideoTrack t : new ArrayList<>(stream.videoTracks)) stream.removeTrack(t); + stream.dispose(); + } catch (Exception e) { + Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing stream " + entry.getKey(), e); + } + } + localStreams.clear(); + + for (int pcId : factoryRegistry.currentOwnedPcIds()) { + try { + PeerConnectionObserver pco = mPeerConnectionObservers.get(pcId); + if (pco != null) { + pco.dispose(); + } + } catch (Exception e) { + Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing pc " + pcId, e); + } finally { mPeerConnectionObservers.remove(pcId); + factoryRegistry.unbindPeerConnection(pcId); } - factoryRegistry.unbindPeerConnection(pcId); - } catch (Exception e) { - Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing pc " + pcId, e); } - } - for (String trackId : factoryRegistry.currentOwnedTrackIds()) { - try { - getUserMediaImpl.disposeTrack(trackId); - } catch (Exception e) { - Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing track " + trackId, e); + for (String trackId : factoryRegistry.currentOwnedTrackIds()) { + try { + getUserMediaImpl.disposeTrack(trackId); + } catch (Exception e) { + Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing track " + trackId, e); + } } - } - return factoryRegistry.disposeCurrent(); + onDisposed.accept(factoryRegistry.disposeCurrent()); + }); } @Override @@ -1259,13 +1295,20 @@ public void onCreateFailure(String s) { @Override public void onCreateSuccess(SessionDescription sdp) { ThreadUtils.runOnExecutor(() -> { + final PeerConnection livePc = pco.getPeerConnection(); + if (livePc == null) { + Log.d(TAG, "onCreateSuccess: pc " + id + " disposed before the callback ran"); + promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); + return; + } + WritableMap params = Arguments.createMap(); WritableMap sdpInfo = Arguments.createMap(); sdpInfo.putString("sdp", sdp.description); sdpInfo.putString("type", sdp.type.canonicalForm()); - params.putArray("transceiversInfo", getTransceiversInfo(peerConnection)); + params.putArray("transceiversInfo", getTransceiversInfo(livePc)); params.putMap("sdpInfo", sdpInfo); WritableArray newTransceivers = Arguments.createArray(); @@ -1299,7 +1342,8 @@ public void onSetSuccess() {} @ReactMethod public void peerConnectionCreateAnswer(int id, ReadableMap options, Promise promise) { ThreadUtils.runOnExecutor(() -> { - PeerConnection peerConnection = getPeerConnection(id); + PeerConnectionObserver pco = mPeerConnectionObservers.get(id); + PeerConnection peerConnection = pco == null ? null : pco.getPeerConnection(); if (peerConnection == null) { Log.d(TAG, "peerConnectionCreateAnswer() peerConnection is null"); @@ -1316,13 +1360,20 @@ public void onCreateFailure(String s) { @Override public void onCreateSuccess(SessionDescription sdp) { ThreadUtils.runOnExecutor(() -> { + final PeerConnection livePc = pco.getPeerConnection(); + if (livePc == null) { + Log.d(TAG, "onCreateSuccess: pc " + id + " disposed before the callback ran"); + promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); + return; + } + WritableMap params = Arguments.createMap(); WritableMap sdpInfo = Arguments.createMap(); sdpInfo.putString("sdp", sdp.description); sdpInfo.putString("type", sdp.type.canonicalForm()); - params.putArray("transceiversInfo", getTransceiversInfo(peerConnection)); + params.putArray("transceiversInfo", getTransceiversInfo(livePc)); params.putMap("sdpInfo", sdpInfo); promise.resolve(params); @@ -1343,7 +1394,8 @@ public void onSetSuccess() {} @ReactMethod public void peerConnectionSetLocalDescription(int pcId, ReadableMap desc, Promise promise) { ThreadUtils.runOnExecutor(() -> { - PeerConnection peerConnection = getPeerConnection(pcId); + PeerConnectionObserver pco = mPeerConnectionObservers.get(pcId); + PeerConnection peerConnection = pco == null ? null : pco.getPeerConnection(); if (peerConnection == null) { Log.d(TAG, "peerConnectionSetLocalDescription() peerConnection is null"); promise.reject(new Exception("PeerConnection not found")); @@ -1357,10 +1409,17 @@ public void onCreateSuccess(SessionDescription sdp) {} @Override public void onSetSuccess() { ThreadUtils.runOnExecutor(() -> { + final PeerConnection livePc = pco.getPeerConnection(); + if (livePc == null) { + Log.d(TAG, "peerConnectionSetLocalDescription() onSetSuccess: pc " + pcId + + " disposed before the callback ran"); + promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); + return; + } WritableMap newSdpMap = Arguments.createMap(); WritableMap params = Arguments.createMap(); - SessionDescription newSdp = peerConnection.getLocalDescription(); + SessionDescription newSdp = livePc.getLocalDescription(); // Can happen when doing a rollback. if (newSdp != null) { newSdpMap.putString("type", newSdp.type.canonicalForm()); @@ -1368,7 +1427,7 @@ public void onSetSuccess() { } params.putMap("sdpInfo", newSdpMap); - params.putArray("transceiversInfo", getTransceiversInfo(peerConnection)); + params.putArray("transceiversInfo", getTransceiversInfo(livePc)); promise.resolve(params); }); @@ -1422,17 +1481,24 @@ public void onCreateSuccess(final SessionDescription sdp) {} @Override public void onSetSuccess() { ThreadUtils.runOnExecutor(() -> { + final PeerConnection livePc = pco.getPeerConnection(); + if (livePc == null) { + Log.d(TAG, "peerConnectionSetRemoteDescription() onSetSuccess: pc " + id + + " disposed before the callback ran"); + promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); + return; + } WritableMap newSdpMap = Arguments.createMap(); WritableMap params = Arguments.createMap(); - SessionDescription newSdp = peerConnection.getRemoteDescription(); + SessionDescription newSdp = livePc.getRemoteDescription(); // Be defensive for the rollback cases. if (newSdp != null) { newSdpMap.putString("type", newSdp.type.canonicalForm()); newSdpMap.putString("sdp", newSdp.description); } - params.putArray("transceiversInfo", getTransceiversInfo(peerConnection)); + params.putArray("transceiversInfo", getTransceiversInfo(livePc)); params.putMap("sdpInfo", newSdpMap); WritableArray newTransceivers = Arguments.createArray(); @@ -1544,7 +1610,8 @@ public void senderGetStats(int pcId, String senderId, Promise promise) { @ReactMethod public void peerConnectionAddICECandidate(int pcId, ReadableMap candidateMap, Promise promise) { ThreadUtils.runOnExecutor(() -> { - PeerConnection peerConnection = getPeerConnection(pcId); + final PeerConnectionObserver pco = mPeerConnectionObservers.get(pcId); + PeerConnection peerConnection = pco == null ? null : pco.getPeerConnection(); if (peerConnection == null) { Log.d(TAG, "peerConnectionAddICECandidate() peerConnection is null"); promise.reject(new Exception("PeerConnection not found")); @@ -1568,8 +1635,15 @@ public void peerConnectionAddICECandidate(int pcId, ReadableMap candidateMap, Pr @Override public void onAddSuccess() { ThreadUtils.runOnExecutor(() -> { + final PeerConnection livePc = pco.getPeerConnection(); + if (livePc == null) { + Log.d(TAG, "peerConnectionAddICECandidate() onAddSuccess: pc " + pcId + + " disposed before the callback ran"); + promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); + return; + } WritableMap newSdpMap = Arguments.createMap(); - SessionDescription newSdp = peerConnection.getRemoteDescription(); + SessionDescription newSdp = livePc.getRemoteDescription(); newSdpMap.putString("type", newSdp.type.canonicalForm()); newSdpMap.putString("sdp", newSdp.description); promise.resolve(newSdpMap); @@ -1613,16 +1687,16 @@ public void peerConnectionClose(int id) { public void peerConnectionDispose(int id) { ThreadUtils.runOnExecutor(() -> { PeerConnectionObserver pco = mPeerConnectionObservers.get(id); - // Null-safe: the PC may already have been disposed (e.g. by - // disposeCallFactory, which tears down the factory's owned PCs first). Skip the - // dispose in that case instead of NPEing, but always clear the factory binding. - if (pco == null) { - Log.d(TAG, "peerConnectionDispose() peerConnection observer is null"); - } else { - pco.dispose(); + try { + if (pco == null) { + Log.d(TAG, "peerConnectionDispose() peerConnection observer is null"); + } else { + pco.dispose(); + } + } finally { mPeerConnectionObservers.remove(id); + factoryRegistry.unbindPeerConnection(id); } - factoryRegistry.unbindPeerConnection(id); }); } From cbc87df88e269d28cc02ae4abeab2faccc98ba97 Mon Sep 17 00:00:00 2001 From: Santhosh Vaiyapuri <3846977+santhoshvai@users.noreply.github.com> Date: Tue, 22 Sep 2026 15:27:30 +0200 Subject: [PATCH 2/2] fix(android): serialize factory teardown and share callback guards (#71) A `createCallFactory()` queued during two-phase teardown can reuse a factory that the pending teardown then destroys. Concurrent creates replacing a bare-fork default can likewise destroy the first successful create's factory. Queue factory lifecycle requests until teardown completes, preserving request order and allowing peer-connection callbacks to drain between phases. Centralize the five success-callback disposal checks in `runWithPeerConnection()`. The helper checks the connection inside the executor task and passes it to the callback. Also share peer-connection disposal, observer removal, and factory unbinding between the explicit and factory teardown paths. This draft targets `pc-dispose-race`, the branch of #70. Validation: - Android Java/Kotlin compilation passed with Gradle 9.4.1, AGP 9.2.1, and React Native 0.87.0 using a temporary host project. - Eight local JVM regression scenarios passed using the production lifecycle methods, registry, and executor with native dependencies stubbed: overlapping dispose/create, concurrent default replacement, reference counts, FIFO request completion, close callback draining, disposed/live callback handling, and cleanup after a PC disposal exception. - `git diff --check` passed. The regression harness was local; no Android device runtime test was performed. --- .../com/oney/WebRTCModule/WebRTCModule.java | 166 +++++++++--------- 1 file changed, 82 insertions(+), 84 deletions(-) diff --git a/android/src/main/java/com/oney/WebRTCModule/WebRTCModule.java b/android/src/main/java/com/oney/WebRTCModule/WebRTCModule.java index b311fe9ba..24092baae 100644 --- a/android/src/main/java/com/oney/WebRTCModule/WebRTCModule.java +++ b/android/src/main/java/com/oney/WebRTCModule/WebRTCModule.java @@ -27,6 +27,7 @@ import org.webrtc.*; import org.webrtc.audio.AudioDeviceModule; +import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Arrays; import java.util.HashMap; @@ -55,6 +56,10 @@ public class WebRTCModule extends ReactContextBaseJavaModule { private final SparseArray mPeerConnectionObservers; final Map localStreams; + // Accessed only on ThreadUtils' executor. Lifecycle requests wait for both teardown phases. + private final ArrayDeque pendingFactoryOperations = new ArrayDeque<>(); + private boolean factoryDisposalPending; + // Store generated certificates by ID to avoid exposing private keys to JS private static final Map mCertificates = new HashMap<>(); @@ -129,7 +134,7 @@ public WebRTCModule(ReactApplicationContext reactContext) { @ReactMethod public void createCallFactory(ReadableMap options, Promise promise) { - ThreadUtils.runOnExecutor(() -> { + runFactoryOperation(() -> { try { boolean bypassVoiceProcessing = options != null && options.hasKey("bypassVoiceProcessing") && options.getBoolean("bypassVoiceProcessing"); @@ -165,7 +170,20 @@ public void createCallFactory(ReadableMap options, Promise promise) { @ReactMethod public void disposeCallFactory(Promise promise) { - ThreadUtils.runOnExecutor(() -> disposeCurrentFactoryOrdered(promise::resolve)); + runFactoryOperation(() -> disposeCurrentFactoryOrdered(promise::resolve)); + } + + private void runFactoryOperation(Runnable operation) { + ThreadUtils.runOnExecutor(() -> { + pendingFactoryOperations.addLast(operation); + drainFactoryOperations(); + }); + } + + private void drainFactoryOperations() { + while (!factoryDisposalPending && !pendingFactoryOperations.isEmpty()) { + pendingFactoryOperations.removeFirst().run(); + } } /** @@ -190,6 +208,7 @@ private void disposeCurrentFactoryOrdered(Consumer onDisposed) { return; } + factoryDisposalPending = true; for (int pcId : factoryRegistry.currentOwnedPcIds()) { try { PeerConnectionObserver pco = mPeerConnectionObservers.get(pcId); @@ -202,41 +221,40 @@ private void disposeCurrentFactoryOrdered(Consumer onDisposed) { } ThreadUtils.runOnExecutor(() -> { - for (Map.Entry entry : localStreams.entrySet()) { - try { - MediaStream stream = entry.getValue(); - for (AudioTrack t : new ArrayList<>(stream.audioTracks)) stream.removeTrack(t); - for (VideoTrack t : new ArrayList<>(stream.videoTracks)) stream.removeTrack(t); - stream.dispose(); - } catch (Exception e) { - Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing stream " + entry.getKey(), e); + try { + for (Map.Entry entry : localStreams.entrySet()) { + try { + MediaStream stream = entry.getValue(); + for (AudioTrack t : new ArrayList<>(stream.audioTracks)) stream.removeTrack(t); + for (VideoTrack t : new ArrayList<>(stream.videoTracks)) stream.removeTrack(t); + stream.dispose(); + } catch (Exception e) { + Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing stream " + entry.getKey(), e); + } } - } - localStreams.clear(); + localStreams.clear(); - for (int pcId : factoryRegistry.currentOwnedPcIds()) { - try { - PeerConnectionObserver pco = mPeerConnectionObservers.get(pcId); - if (pco != null) { - pco.dispose(); + for (int pcId : factoryRegistry.currentOwnedPcIds()) { + try { + disposePeerConnection(pcId); + } catch (Exception e) { + Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing pc " + pcId, e); } - } catch (Exception e) { - Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing pc " + pcId, e); - } finally { - mPeerConnectionObservers.remove(pcId); - factoryRegistry.unbindPeerConnection(pcId); } - } - for (String trackId : factoryRegistry.currentOwnedTrackIds()) { - try { - getUserMediaImpl.disposeTrack(trackId); - } catch (Exception e) { - Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing track " + trackId, e); + for (String trackId : factoryRegistry.currentOwnedTrackIds()) { + try { + getUserMediaImpl.disposeTrack(trackId); + } catch (Exception e) { + Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing track " + trackId, e); + } } - } - onDisposed.accept(factoryRegistry.disposeCurrent()); + onDisposed.accept(factoryRegistry.disposeCurrent()); + } finally { + factoryDisposalPending = false; + drainFactoryOperations(); + } }); } @@ -350,6 +368,19 @@ private PeerConnection getPeerConnection(int id) { return (pco == null) ? null : pco.getPeerConnection(); } + private void runWithPeerConnection( + PeerConnectionObserver pco, Promise promise, Consumer action) { + ThreadUtils.runOnExecutor(() -> { + PeerConnection pc = pco.getPeerConnection(); + if (pc == null) { + Log.d(TAG, "PeerConnection disposed before callback ran"); + promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); + return; + } + action.accept(pc); + }); + } + void sendEvent(String eventName, @Nullable ReadableMap params) { if (getReactApplicationContext().hasActiveReactInstance()) { getReactApplicationContext() @@ -1294,14 +1325,7 @@ public void onCreateFailure(String s) { @Override public void onCreateSuccess(SessionDescription sdp) { - ThreadUtils.runOnExecutor(() -> { - final PeerConnection livePc = pco.getPeerConnection(); - if (livePc == null) { - Log.d(TAG, "onCreateSuccess: pc " + id + " disposed before the callback ran"); - promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); - return; - } - + runWithPeerConnection(pco, promise, livePc -> { WritableMap params = Arguments.createMap(); WritableMap sdpInfo = Arguments.createMap(); @@ -1312,7 +1336,7 @@ public void onCreateSuccess(SessionDescription sdp) { params.putMap("sdpInfo", sdpInfo); WritableArray newTransceivers = Arguments.createArray(); - for (RtpTransceiver transceiver : peerConnection.getTransceivers()) { + for (RtpTransceiver transceiver : livePc.getTransceivers()) { if (!receiversIds.contains(transceiver.getReceiver().id())) { WritableMap newTransceiver = Arguments.createMap(); newTransceiver.putInt("transceiverOrder", pco.getNextTransceiverId()); @@ -1359,14 +1383,7 @@ public void onCreateFailure(String s) { @Override public void onCreateSuccess(SessionDescription sdp) { - ThreadUtils.runOnExecutor(() -> { - final PeerConnection livePc = pco.getPeerConnection(); - if (livePc == null) { - Log.d(TAG, "onCreateSuccess: pc " + id + " disposed before the callback ran"); - promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); - return; - } - + runWithPeerConnection(pco, promise, livePc -> { WritableMap params = Arguments.createMap(); WritableMap sdpInfo = Arguments.createMap(); @@ -1408,14 +1425,7 @@ public void onCreateSuccess(SessionDescription sdp) {} @Override public void onSetSuccess() { - ThreadUtils.runOnExecutor(() -> { - final PeerConnection livePc = pco.getPeerConnection(); - if (livePc == null) { - Log.d(TAG, "peerConnectionSetLocalDescription() onSetSuccess: pc " + pcId - + " disposed before the callback ran"); - promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); - return; - } + runWithPeerConnection(pco, promise, livePc -> { WritableMap newSdpMap = Arguments.createMap(); WritableMap params = Arguments.createMap(); @@ -1480,14 +1490,7 @@ public void onCreateSuccess(final SessionDescription sdp) {} @Override public void onSetSuccess() { - ThreadUtils.runOnExecutor(() -> { - final PeerConnection livePc = pco.getPeerConnection(); - if (livePc == null) { - Log.d(TAG, "peerConnectionSetRemoteDescription() onSetSuccess: pc " + id - + " disposed before the callback ran"); - promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); - return; - } + runWithPeerConnection(pco, promise, livePc -> { WritableMap newSdpMap = Arguments.createMap(); WritableMap params = Arguments.createMap(); @@ -1502,7 +1505,7 @@ public void onSetSuccess() { params.putMap("sdpInfo", newSdpMap); WritableArray newTransceivers = Arguments.createArray(); - for (RtpTransceiver transceiver : peerConnection.getTransceivers()) { + for (RtpTransceiver transceiver : livePc.getTransceivers()) { if (!receiversIds.contains(transceiver.getReceiver().id())) { WritableMap newTransceiver = Arguments.createMap(); newTransceiver.putInt("transceiverOrder", pco.getNextTransceiverId()); @@ -1634,14 +1637,7 @@ public void peerConnectionAddICECandidate(int pcId, ReadableMap candidateMap, Pr peerConnection.addIceCandidate(candidate, new AddIceObserver() { @Override public void onAddSuccess() { - ThreadUtils.runOnExecutor(() -> { - final PeerConnection livePc = pco.getPeerConnection(); - if (livePc == null) { - Log.d(TAG, "peerConnectionAddICECandidate() onAddSuccess: pc " + pcId - + " disposed before the callback ran"); - promise.reject("E_PC_DISPOSED", "PeerConnection disposed"); - return; - } + runWithPeerConnection(pco, promise, livePc -> { WritableMap newSdpMap = Arguments.createMap(); SessionDescription newSdp = livePc.getRemoteDescription(); newSdpMap.putString("type", newSdp.type.canonicalForm()); @@ -1685,19 +1681,21 @@ public void peerConnectionClose(int id) { @ReactMethod public void peerConnectionDispose(int id) { - ThreadUtils.runOnExecutor(() -> { - PeerConnectionObserver pco = mPeerConnectionObservers.get(id); - try { - if (pco == null) { - Log.d(TAG, "peerConnectionDispose() peerConnection observer is null"); - } else { - pco.dispose(); - } - } finally { - mPeerConnectionObservers.remove(id); - factoryRegistry.unbindPeerConnection(id); + ThreadUtils.runOnExecutor(() -> disposePeerConnection(id)); + } + + private void disposePeerConnection(int id) { + PeerConnectionObserver pco = mPeerConnectionObservers.get(id); + try { + if (pco == null) { + Log.d(TAG, "peerConnectionDispose() peerConnection observer is null"); + } else { + pco.dispose(); } - }); + } finally { + mPeerConnectionObservers.remove(id); + factoryRegistry.unbindPeerConnection(id); + } } @ReactMethod