diff --git a/src/main/java/org/thoughtcrime/securesms/calls/CallCoordinator.java b/src/main/java/org/thoughtcrime/securesms/calls/CallCoordinator.java index 7bd7057f6..60b53bf62 100644 --- a/src/main/java/org/thoughtcrime/securesms/calls/CallCoordinator.java +++ b/src/main/java/org/thoughtcrime/securesms/calls/CallCoordinator.java @@ -42,6 +42,8 @@ import com.b44t.messenger.DcChat; import com.b44t.messenger.DcContext; import com.b44t.messenger.DcEvent; import java.util.List; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; import kotlin.Unit; import kotlin.coroutines.Continuation; import kotlin.coroutines.CoroutineContext; @@ -77,6 +79,7 @@ public class CallCoordinator implements DcEventCenter.DcEventDelegate { private final Rpc rpc; private CoroutineScope audioFlowScope; private final Handler mainHandler = new Handler(Looper.getMainLooper()); + private final ExecutorService eventExecutor = Executors.newSingleThreadExecutor(); // LiveData for Observable State private final MutableLiveData connectionState = @@ -121,7 +124,7 @@ public class CallCoordinator implements DcEventCenter.DcEventDelegate { private final Runnable outgoingRingtoneRunnable = () -> { synchronized (CallCoordinator.this) { - if (callService != null && activeCallId != null) { + if (callService != null && hasActiveCall()) { callService.startOutgoingRingtone(); } } @@ -823,43 +826,42 @@ public class CallCoordinator implements DcEventCenter.DcEventDelegate { // It seems this observer can fire on either main or background thread // Always move to background - new Thread( - () -> { - boolean hasVideo; + eventExecutor.execute( + () -> { + boolean hasVideo; - switch (eventId) { - case DcContext.DC_EVENT_INCOMING_CALL: - try { - hasVideo = this.rpc.callInfo(accId, callId).hasVideo; - } catch (RpcException e) { - Log.e(TAG, "Rpc.callInfo() failed", e); - hasVideo = false; - } - onIncomingCall(accId, callId, event.getData2Str(), hasVideo); - break; - case DcContext.DC_EVENT_INCOMING_CALL_ACCEPTED: - onIncomingCallAccepted(callId); - break; - case DcContext.DC_EVENT_OUTGOING_CALL_ACCEPTED: - String answerSDP = event.getData2Str(); - onOutgoingCallAccepted(callId, answerSDP); - break; - case DcContext.DC_EVENT_CALL_ENDED: - // This event is problematic because it can trigger in both directions, - // in addition to multiple other scenarios which cannot easily be distinguished - // May cause problems in edge cases - onCallEnded(accId, callId); - break; + switch (eventId) { + case DcContext.DC_EVENT_INCOMING_CALL: + try { + hasVideo = this.rpc.callInfo(accId, callId).hasVideo; + } catch (RpcException e) { + Log.e(TAG, "Rpc.callInfo() failed", e); + hasVideo = false; } - }) - .start(); + onIncomingCall(accId, callId, event.getData2Str(), hasVideo); + break; + case DcContext.DC_EVENT_INCOMING_CALL_ACCEPTED: + onIncomingCallAccepted(callId); + break; + case DcContext.DC_EVENT_OUTGOING_CALL_ACCEPTED: + String answerSDP = event.getData2Str(); + onOutgoingCallAccepted(callId, answerSDP); + break; + case DcContext.DC_EVENT_CALL_ENDED: + // This event is problematic because it can trigger in both directions, + // in addition to multiple other scenarios which cannot easily be distinguished + // May cause problems in edge cases + onCallEnded(accId, callId); + break; + } + }); } private synchronized void onIncomingCall( int accId, int callId, String offerSdp, boolean startsWithVideo) { Log.d(TAG, "onIncomingCall: accId=" + accId + ", callId=" + callId); - if (activeCallId != null) { + if (hasActiveCall()) { Log.w(TAG, "Already have an active call, ignoring incoming call"); return; } @@ -955,6 +957,27 @@ public class CallCoordinator implements DcEventCenter.DcEventDelegate { private synchronized void onCallEnded(int accId, int callId) { Log.d(TAG, "onCallEnded: accId=" + accId + ", callId=" + callId); + if (!hasActiveCall()) { + Log.w(TAG, "No active call, ignoring"); + return; + } + + if (!activeAccId.equals(accId) || !activeCallId.equals(callId)) { + Log.w( + TAG, + "Event IDs don't match active call " + + "(active: accId=" + + activeAccId + + " callId=" + + activeCallId + + ", event: accId=" + + accId + + " callId=" + + callId + + "), ignoring"); + return; + } + if (callService != null) { callService.stopRingtone(); } @@ -1002,7 +1025,7 @@ public class CallCoordinator implements DcEventCenter.DcEventDelegate { notifyBackendCallEnded(); // Cleanup - if (activeAccId != null && activeCallId != null) { + if (hasActiveCall()) { cleanupCall(activeAccId, activeCallId); } } @@ -1011,10 +1034,14 @@ public class CallCoordinator implements DcEventCenter.DcEventDelegate { public synchronized void cleanupCall(int accId, int callId) { Log.d(TAG, "cleanupCall: accId=" + accId + ", callId=" + callId); - if (activeCallId != null && !activeCallId.equals(callId) - || activeAccId != null && !activeAccId.equals(accId)) { - Log.w(TAG, "Cleanup accountId or callId doesn't match active call"); - // Clean up anyway. Otherwise, no new calls can happen. + if (!hasActiveCall()) { + Log.d(TAG, "No active call to clean up"); + return; + } + + if (!activeAccId.equals(accId) || !activeCallId.equals(callId)) { + Log.w(TAG, "Cleanup IDs don't match active call, aborting"); + return; } // Clear state @@ -1112,7 +1139,7 @@ public class CallCoordinator implements DcEventCenter.DcEventDelegate { public synchronized void initiateOutgoingCall(int accId, int chatId, boolean startsWithVideo) { Log.d(TAG, "Initiating outgoing call:accId=" + accId + ", chatId=" + chatId); - if (activeCallId != null) { + if (hasActiveCall()) { Log.w(TAG, "Already have an active call, cannot start new one"); return; } @@ -1230,10 +1257,7 @@ public class CallCoordinator implements DcEventCenter.DcEventDelegate { Log.d(TAG, "CallControlScope initialized"); activeCallControlScope = scope; - mainHandler.post( - () -> { - setupAudioEndpointCollection(scope); - }); + mainHandler.post(() -> setupAudioEndpointCollection(scope)); return Unit.INSTANCE; }, @@ -1472,7 +1496,7 @@ public class CallCoordinator implements DcEventCenter.DcEventDelegate { } public synchronized boolean hasActiveCall() { - return activeCallId != null; + return activeAccId != null && activeCallId != null; } public synchronized boolean hasOngoingCall() { diff --git a/src/main/java/org/thoughtcrime/securesms/calls/CallService.java b/src/main/java/org/thoughtcrime/securesms/calls/CallService.java index 1250e2412..fa18aabad 100644 --- a/src/main/java/org/thoughtcrime/securesms/calls/CallService.java +++ b/src/main/java/org/thoughtcrime/securesms/calls/CallService.java @@ -25,8 +25,10 @@ import org.webrtc.PeerConnection; import org.webrtc.VideoTrack; /** - * Foreground service for VoIP calls Required to post CallStyle notifications on Android 12+ Owns - * WebRTC resources and keeps call alive + * Foreground service for VoIP calls + * + *

Required to post CallStyle notifications on Android 12+. Owns WebRTC resources and keeps call + * alive. */ @RequiresApi(api = Build.VERSION_CODES.O) public class CallService extends Service implements WebRTCClient.Callbacks { @@ -100,6 +102,10 @@ public class CallService extends Service implements WebRTCClient.Callbacks { private void fetchIceServersAndSetup() { new Thread( () -> { + if (!callCoordinator.hasActiveCall()) { + Log.d(TAG, "Call ended before ICE fetch, aborting"); + return; + } try { String iceServersJson = callCoordinator.fetchIceServers(); Log.d(TAG, "ICE servers fetched: " + iceServersJson); @@ -117,8 +123,10 @@ public class CallService extends Service implements WebRTCClient.Callbacks { } /** - * Start camera/microphone capture Must be called when app is in foreground Called by coordinator - * when ViewModel/Activity is ready + * Start camera/microphone capture + * + *

Must be called when app is in foreground. Called by coordinator when ViewModel/Activity is + * ready. */ public void startMediaCapture() { Log.d(TAG, "startMediaCapture (Camera/Microphone)"); diff --git a/src/main/java/org/thoughtcrime/securesms/calls/MediaStreamManager.java b/src/main/java/org/thoughtcrime/securesms/calls/MediaStreamManager.java index 4ec94252d..44a436782 100644 --- a/src/main/java/org/thoughtcrime/securesms/calls/MediaStreamManager.java +++ b/src/main/java/org/thoughtcrime/securesms/calls/MediaStreamManager.java @@ -19,7 +19,7 @@ import org.webrtc.VideoTrack; public class MediaStreamManager { - private static final String TAG = CallUtil.class.getSimpleName(); + private static final String TAG = MediaStreamManager.class.getSimpleName(); private static final String STREAM_ID = "local_stream"; private static final String AUDIO_TRACK_ID = "audio_track"; private static final String VIDEO_TRACK_ID = "video_track"; diff --git a/src/main/java/org/thoughtcrime/securesms/webrtc/WebRTCClient.java b/src/main/java/org/thoughtcrime/securesms/webrtc/WebRTCClient.java index 5c9f21a84..f64fd2f59 100644 --- a/src/main/java/org/thoughtcrime/securesms/webrtc/WebRTCClient.java +++ b/src/main/java/org/thoughtcrime/securesms/webrtc/WebRTCClient.java @@ -156,8 +156,14 @@ public class WebRTCClient { } /** - * Expected JSON format: [ {"urls": "stun:stun.example.com:3478"}, {"urls": - * "turn:turn.example.com", "username": "user", "credential": "pass"} ] + * Expected JSON format: + * + *

{@code
+   * [
+   *   {"urls": "stun:stun.example.com:3478"},
+   *   {"urls": "turn:turn.example.com", "username": "user", "credential": "pass"}
+   * ]
+   * }
*/ public void configure(String iceServersJson) { this.iceServers = parseIceServers(iceServersJson); @@ -514,7 +520,7 @@ public class WebRTCClient { } } - /** Wait for enough ICE before sending offer/answer Mirrors TypeScript gatheredEnoughIce() */ + /** Wait for enough ICE before sending offer/answer. Mirrors TypeScript gatheredEnoughIce() */ private boolean waitForEnoughIce() { boolean hasTurnServer = false; boolean hasStunServer = false; @@ -637,24 +643,32 @@ public class WebRTCClient { () -> { boolean gotIce = waitForEnoughIce(); - if (!gotIce) { - Log.w(TAG, "Proceeding without optimal ICE candidates"); - } + synchronized (WebRTCClient.this) { + if (isEnded || peerConnection == null) { + Log.d(TAG, "Call ended during ICE gathering, aborting"); + return; + } - // Get final SDP - SessionDescription localDesc = peerConnection.getLocalDescription(); - if (localDesc != null) { - String finalSdp = localDesc.description; + if (!gotIce) { + Log.w(TAG, "Proceeding without optimal ICE candidates"); + } - // Enable trickling for additional candidates - enableIceTrickling = true; + // Get final SDP + SessionDescription localDesc = peerConnection.getLocalDescription(); + if (localDesc != null) { + String finalSdp = localDesc.description; - // Notify ViewModel - mainHandler.post(() -> callbacks.onOfferReady(finalSdp)); - } else { - mainHandler.post( - () -> - callbacks.onError("Local description is null after ICE gathering")); + // Enable trickling for additional candidates + enableIceTrickling = true; + + // Notify ViewModel + mainHandler.post(() -> callbacks.onOfferReady(finalSdp)); + } else { + mainHandler.post( + () -> + callbacks.onError( + "Local description is null after ICE gathering")); + } } }) .start(); @@ -788,24 +802,32 @@ public class WebRTCClient { () -> { boolean gotIce = waitForEnoughIce(); - if (!gotIce) { - Log.w(TAG, "Proceeding without optimal ICE candidates"); - } + synchronized (WebRTCClient.this) { + if (isEnded || peerConnection == null) { + Log.d(TAG, "Call ended during ICE gathering, aborting"); + return; + } - // Get final SDP - SessionDescription localDesc = peerConnection.getLocalDescription(); - if (localDesc != null) { - String finalSdp = localDesc.description; + if (!gotIce) { + Log.w(TAG, "Proceeding without optimal ICE candidates"); + } - // Enable trickling for additional candidates - enableIceTrickling = true; + // Get final SDP + SessionDescription localDesc = peerConnection.getLocalDescription(); + if (localDesc != null) { + String finalSdp = localDesc.description; - // Notify ViewModel - mainHandler.post(() -> callbacks.onAnswerReady(finalSdp)); - } else { - mainHandler.post( - () -> - callbacks.onError("Local description is null after ICE gathering")); + // Enable trickling for additional candidates + enableIceTrickling = true; + + // Notify ViewModel + mainHandler.post(() -> callbacks.onAnswerReady(finalSdp)); + } else { + mainHandler.post( + () -> + callbacks.onError( + "Local description is null after ICE gathering")); + } } }) .start(); @@ -916,7 +938,7 @@ public class WebRTCClient { // Cleanup - public void endCall() { + public synchronized void endCall() { if (isEnded) { Log.d(TAG, "endCall() already called, skipping"); return; @@ -925,6 +947,11 @@ public class WebRTCClient { isEnded = true; Log.d(TAG, "Ending call"); + // Unblock any thread waiting in waitForEnoughIce() + if (iceGatheringLatch != null) iceGatheringLatch.countDown(); + if (relayCandidateLatch != null) relayCandidateLatch.countDown(); + if (srflxCandidateLatch != null) srflxCandidateLatch.countDown(); + if (peerConnection != null) { peerConnection.close(); peerConnection.dispose();