mirror of
https://github.com/vitorpamplona/amethyst.git
synced 2026-10-06 03:38:23 +00:00
Merge pull request #2211 from vitorpamplona/claude/review-webrtc-calls-Vc6v0
Claude/review webrtc calls vc6v0
This commit is contained in:
@@ -103,6 +103,8 @@ class CallController(
|
||||
private val cleanedUp = AtomicBoolean(false)
|
||||
private val videoSenders = ConcurrentHashMap<HexKey, org.webrtc.RtpSender>()
|
||||
|
||||
private val pendingRenegotiation = ConcurrentHashMap<HexKey, Boolean>()
|
||||
|
||||
private val connectivityManager = context.getSystemService(Context.CONNECTIVITY_SERVICE) as ConnectivityManager
|
||||
private var networkCallbackRegistered = false
|
||||
private val networkCallback =
|
||||
@@ -220,6 +222,7 @@ class CallController(
|
||||
|
||||
callManager.beginOffering(callId, peerPubKeys, callType)
|
||||
|
||||
var successCount = 0
|
||||
for (peerPubKey in peerPubKeys) {
|
||||
try {
|
||||
val webRtcSession = withContext(Dispatchers.IO) { createWebRtcSession(peerPubKey) }
|
||||
@@ -230,10 +233,16 @@ class CallController(
|
||||
callManager.publishOfferToPeer(peerPubKey, peerPubKeys, callType, callId, sdp.description)
|
||||
}
|
||||
}
|
||||
successCount++
|
||||
} catch (e: Exception) {
|
||||
Log.e(TAG, "Failed to create PeerConnection for ${peerPubKey.take(8)}", e)
|
||||
}
|
||||
}
|
||||
if (successCount == 0) {
|
||||
Log.e(TAG, "All PeerConnection creations failed, hanging up")
|
||||
_errorMessage.value = "Failed to start call: could not create any connections"
|
||||
callManager.hangup()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -402,10 +411,19 @@ class CallController(
|
||||
}
|
||||
|
||||
private fun performRenegotiation(peerPubKey: HexKey) {
|
||||
val webRtcSession = webRtcSession(peerPubKey) ?: return
|
||||
if (pendingRenegotiation.putIfAbsent(peerPubKey, true) != null) return
|
||||
val webRtcSession =
|
||||
webRtcSession(peerPubKey) ?: run {
|
||||
pendingRenegotiation.remove(peerPubKey)
|
||||
return
|
||||
}
|
||||
val state = callManager.state.value
|
||||
if (state !is CallState.Connected && state !is CallState.Connecting) return
|
||||
if (state !is CallState.Connected && state !is CallState.Connecting) {
|
||||
pendingRenegotiation.remove(peerPubKey)
|
||||
return
|
||||
}
|
||||
webRtcSession.createOffer { sdp ->
|
||||
pendingRenegotiation.remove(peerPubKey)
|
||||
scope.launch { callManager.sendRenegotiation(sdp.description, peerPubKey) }
|
||||
}
|
||||
}
|
||||
@@ -533,7 +551,9 @@ class CallController(
|
||||
Log.d(TAG) { "Peer ${peerPubKey.take(8)} connected!" }
|
||||
scope.launch {
|
||||
callManager.onPeerConnected()
|
||||
ensureForegroundService()
|
||||
if (callManager.state.value is CallState.Connected) {
|
||||
ensureForegroundService()
|
||||
}
|
||||
}
|
||||
},
|
||||
onRemoteVideoTrack = { track -> videoMonitor.onRemoteVideoTrack(peerPubKey, track) },
|
||||
@@ -576,6 +596,7 @@ class CallController(
|
||||
// ---- Per-peer cleanup ----
|
||||
|
||||
fun disposePeerSession(peerPubKey: HexKey) {
|
||||
videoSenders.remove(peerPubKey)
|
||||
val entry = peerSessionMgr.removeSession(peerPubKey)
|
||||
if (entry != null) {
|
||||
try {
|
||||
@@ -619,6 +640,7 @@ class CallController(
|
||||
_isAudioMuted.value = false
|
||||
videoPausedByProximity = false
|
||||
videoSenders.clear()
|
||||
pendingRenegotiation.clear()
|
||||
cleanedUp.set(false)
|
||||
}
|
||||
|
||||
|
||||
@@ -117,6 +117,7 @@ class CallMediaManager(
|
||||
}
|
||||
}
|
||||
|
||||
@Synchronized
|
||||
fun createVideoResources() {
|
||||
if (localVideoSource != null) return
|
||||
val factory = peerConnectionFactory ?: return
|
||||
|
||||
+32
-22
@@ -85,46 +85,56 @@ class RemoteVideoMonitor(
|
||||
private val perPeerLastFrameTimeMs = ConcurrentHashMap<HexKey, AtomicLong>()
|
||||
private var groupVideoMonitorJob: Job? = null
|
||||
|
||||
/** Protects compound read-modify-write on [_remoteVideoTracks] and
|
||||
* [_remoteVideoTrack] which can be called from WebRTC callback threads. */
|
||||
private val trackLock = Any()
|
||||
|
||||
fun onRemoteVideoTrack(
|
||||
peerPubKey: HexKey,
|
||||
track: VideoTrack,
|
||||
) {
|
||||
Log.d(TAG) { "Remote video track from ${peerPubKey.take(8)}" }
|
||||
_remoteVideoTracks.value = _remoteVideoTracks.value + (peerPubKey to track)
|
||||
if (_remoteVideoTrack.value == null) {
|
||||
_remoteVideoTrack.value = track
|
||||
startPrimaryMonitor(track)
|
||||
synchronized(trackLock) {
|
||||
_remoteVideoTracks.value = _remoteVideoTracks.value + (peerPubKey to track)
|
||||
if (_remoteVideoTrack.value == null) {
|
||||
_remoteVideoTrack.value = track
|
||||
startPrimaryMonitor(track)
|
||||
}
|
||||
}
|
||||
startPeerMonitor(peerPubKey, track)
|
||||
}
|
||||
|
||||
fun onPeerRemoved(peerPubKey: HexKey) {
|
||||
stopPeerMonitor(peerPubKey)
|
||||
val currentTracks = _remoteVideoTracks.value
|
||||
if (peerPubKey in currentTracks) {
|
||||
_remoteVideoTracks.value = currentTracks - peerPubKey
|
||||
if (_remoteVideoTrack.value == currentTracks[peerPubKey]) {
|
||||
stopPrimaryMonitor()
|
||||
val nextTrack = _remoteVideoTracks.value.values.firstOrNull()
|
||||
_remoteVideoTrack.value = nextTrack
|
||||
if (nextTrack != null) {
|
||||
startPrimaryMonitor(nextTrack)
|
||||
synchronized(trackLock) {
|
||||
val currentTracks = _remoteVideoTracks.value
|
||||
if (peerPubKey in currentTracks) {
|
||||
_remoteVideoTracks.value = currentTracks - peerPubKey
|
||||
if (_remoteVideoTrack.value == currentTracks[peerPubKey]) {
|
||||
stopPrimaryMonitor()
|
||||
val nextTrack = _remoteVideoTracks.value.values.firstOrNull()
|
||||
_remoteVideoTrack.value = nextTrack
|
||||
if (nextTrack != null) {
|
||||
startPrimaryMonitor(nextTrack)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fun dispose() {
|
||||
stopPrimaryMonitor()
|
||||
stopGroupMonitor()
|
||||
for (peerPubKey in perPeerFrameSinks.keys.toList()) {
|
||||
stopPeerMonitor(peerPubKey)
|
||||
synchronized(trackLock) {
|
||||
stopPrimaryMonitor()
|
||||
stopGroupMonitor()
|
||||
for (peerPubKey in perPeerFrameSinks.keys.toList()) {
|
||||
stopPeerMonitor(peerPubKey)
|
||||
}
|
||||
_remoteVideoTrack.value = null
|
||||
_remoteVideoTracks.value = emptyMap()
|
||||
_isRemoteVideoActive.value = false
|
||||
_remoteVideoAspectRatio.value = null
|
||||
_activePeerVideos.value = emptySet()
|
||||
}
|
||||
_remoteVideoTrack.value = null
|
||||
_remoteVideoTracks.value = emptyMap()
|
||||
_isRemoteVideoActive.value = false
|
||||
_remoteVideoAspectRatio.value = null
|
||||
_activePeerVideos.value = emptySet()
|
||||
}
|
||||
|
||||
private fun startPrimaryMonitor(track: VideoTrack) {
|
||||
|
||||
@@ -46,9 +46,10 @@ import com.vitorpamplona.amethyst.ui.StringResSetup
|
||||
import com.vitorpamplona.amethyst.ui.screen.ManageRelayServices
|
||||
import com.vitorpamplona.amethyst.ui.screen.ManageWebOkHttp
|
||||
import com.vitorpamplona.amethyst.ui.theme.AmethystTheme
|
||||
import kotlinx.coroutines.NonCancellable
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.withContext
|
||||
|
||||
class CallActivity : AppCompatActivity() {
|
||||
val isInPipMode = mutableStateOf(false)
|
||||
@@ -66,6 +67,7 @@ class CallActivity : AppCompatActivity() {
|
||||
}
|
||||
|
||||
private var pendingAcceptIsVideo = false
|
||||
private var hangupInitiated = false
|
||||
|
||||
private val pipActionReceiver =
|
||||
object : BroadcastReceiver() {
|
||||
@@ -198,13 +200,17 @@ class CallActivity : AppCompatActivity() {
|
||||
// We must NOT hang up when the user simply presses Home from the full-screen
|
||||
// call UI (that enters PiP via onUserLeaveHint instead).
|
||||
if (wasInPipMode && !isInPictureInPictureMode) {
|
||||
val state = CallSessionBridge.callManager?.state?.value
|
||||
hangupInitiated = true
|
||||
val manager = CallSessionBridge.callManager
|
||||
val state = manager?.state?.value
|
||||
if (state is CallState.Connected || state is CallState.Connecting || state is CallState.Offering) {
|
||||
lifecycleScope.launch {
|
||||
withContext(NonCancellable) { CallSessionBridge.callManager?.hangup() }
|
||||
CoroutineScope(SupervisorJob() + Dispatchers.Main.immediate).launch {
|
||||
manager.hangup()
|
||||
finishAndRemoveTask()
|
||||
}
|
||||
} else {
|
||||
finishAndRemoveTask()
|
||||
}
|
||||
finishAndRemoveTask()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -213,26 +219,27 @@ class CallActivity : AppCompatActivity() {
|
||||
|
||||
// Safety net: if the Activity is destroyed while a call is still
|
||||
// ringing/offering, ensure the call is hung up so audio stops.
|
||||
// Use NonCancellable so the signaling event is published even
|
||||
// though the lifecycle scope is being cancelled.
|
||||
val manager = CallSessionBridge.callManager
|
||||
when (manager?.state?.value) {
|
||||
is CallState.IncomingCall -> {
|
||||
lifecycleScope.launch {
|
||||
withContext(NonCancellable) { manager.rejectCall() }
|
||||
// Skip if onStop already initiated the hangup to avoid double signaling.
|
||||
if (!hangupInitiated) {
|
||||
val manager = CallSessionBridge.callManager
|
||||
when (manager?.state?.value) {
|
||||
is CallState.IncomingCall -> {
|
||||
CoroutineScope(SupervisorJob() + Dispatchers.Main.immediate).launch {
|
||||
manager.rejectCall()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
is CallState.Offering,
|
||||
is CallState.Connecting,
|
||||
is CallState.Connected,
|
||||
-> {
|
||||
lifecycleScope.launch {
|
||||
withContext(NonCancellable) { manager.hangup() }
|
||||
is CallState.Offering,
|
||||
is CallState.Connecting,
|
||||
is CallState.Connected,
|
||||
-> {
|
||||
CoroutineScope(SupervisorJob() + Dispatchers.Main.immediate).launch {
|
||||
manager.hangup()
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
else -> {}
|
||||
else -> {}
|
||||
}
|
||||
}
|
||||
|
||||
super.onDestroy()
|
||||
|
||||
Reference in New Issue
Block a user