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..24092baae 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; @@ -26,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; @@ -54,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<>(); @@ -128,21 +134,33 @@ 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"); 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 +170,20 @@ public void createCallFactory(ReadableMap options, Promise promise) { @ReactMethod public void disposeCallFactory(Promise promise) { - ThreadUtils.runOnExecutor(() -> promise.resolve(disposeCurrentFactoryOrdered())); + 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(); + } } /** @@ -161,47 +192,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()) { - 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(); - + factoryDisposalPending = true; for (int pcId : factoryRegistry.currentOwnedPcIds()) { try { PeerConnectionObserver pco = mPeerConnectionObservers.get(pcId); - if (pco != null && pco.getPeerConnection() != null) { - pco.dispose(); - mPeerConnectionObservers.remove(pcId); + if (pco != null) { + pco.close(); } - factoryRegistry.unbindPeerConnection(pcId); } catch (Exception e) { - Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing pc " + pcId, e); + Log.w(TAG, "disposeCurrentFactoryOrdered(): error closing pc " + pcId, e); } } - for (String trackId : factoryRegistry.currentOwnedTrackIds()) { + ThreadUtils.runOnExecutor(() -> { try { - getUserMediaImpl.disposeTrack(trackId); - } catch (Exception e) { - Log.w(TAG, "disposeCurrentFactoryOrdered(): error disposing track " + trackId, e); - } - } + 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(); - return factoryRegistry.disposeCurrent(); + for (int pcId : factoryRegistry.currentOwnedPcIds()) { + try { + disposePeerConnection(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); + } + } + + onDisposed.accept(factoryRegistry.disposeCurrent()); + } finally { + factoryDisposalPending = false; + drainFactoryOperations(); + } + }); } @Override @@ -314,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() @@ -1258,18 +1325,18 @@ public void onCreateFailure(String s) { @Override public void onCreateSuccess(SessionDescription sdp) { - ThreadUtils.runOnExecutor(() -> { + runWithPeerConnection(pco, promise, livePc -> { 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(); - 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()); @@ -1299,7 +1366,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"); @@ -1315,14 +1383,14 @@ public void onCreateFailure(String s) { @Override public void onCreateSuccess(SessionDescription sdp) { - ThreadUtils.runOnExecutor(() -> { + runWithPeerConnection(pco, promise, livePc -> { 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 +1411,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")); @@ -1356,11 +1425,11 @@ public void onCreateSuccess(SessionDescription sdp) {} @Override public void onSetSuccess() { - ThreadUtils.runOnExecutor(() -> { + runWithPeerConnection(pco, promise, livePc -> { 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 +1437,7 @@ public void onSetSuccess() { } params.putMap("sdpInfo", newSdpMap); - params.putArray("transceiversInfo", getTransceiversInfo(peerConnection)); + params.putArray("transceiversInfo", getTransceiversInfo(livePc)); promise.resolve(params); }); @@ -1421,22 +1490,22 @@ public void onCreateSuccess(final SessionDescription sdp) {} @Override public void onSetSuccess() { - ThreadUtils.runOnExecutor(() -> { + runWithPeerConnection(pco, promise, livePc -> { 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(); - 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()); @@ -1544,7 +1613,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")); @@ -1567,9 +1637,9 @@ public void peerConnectionAddICECandidate(int pcId, ReadableMap candidateMap, Pr peerConnection.addIceCandidate(candidate, new AddIceObserver() { @Override public void onAddSuccess() { - ThreadUtils.runOnExecutor(() -> { + runWithPeerConnection(pco, promise, livePc -> { 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); @@ -1611,19 +1681,21 @@ public void peerConnectionClose(int id) { @ReactMethod 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. + 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(); - mPeerConnectionObservers.remove(id); } + } finally { + mPeerConnectionObservers.remove(id); factoryRegistry.unbindPeerConnection(id); - }); + } } @ReactMethod