Merge pull request #2075 from vitorpamplona/claude/fix-video-freeze-on-share-stop-V8BX3
Add remote video activity monitoring to detect peer disconnections
This commit is contained in:
@@ -28,10 +28,12 @@ import com.vitorpamplona.amethyst.service.notifications.NotificationUtils
|
|||||||
import com.vitorpamplona.quartz.nip59Giftwrap.wraps.GiftWrapEvent
|
import com.vitorpamplona.quartz.nip59Giftwrap.wraps.GiftWrapEvent
|
||||||
import com.vitorpamplona.quartz.nipACWebRtcCalls.WebRtcCallFactory
|
import com.vitorpamplona.quartz.nipACWebRtcCalls.WebRtcCallFactory
|
||||||
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallIceCandidateEvent
|
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallIceCandidateEvent
|
||||||
|
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallRenegotiateEvent
|
||||||
import com.vitorpamplona.quartz.nipACWebRtcCalls.tags.CallType
|
import com.vitorpamplona.quartz.nipACWebRtcCalls.tags.CallType
|
||||||
import com.vitorpamplona.quartz.utils.Log
|
import com.vitorpamplona.quartz.utils.Log
|
||||||
import kotlinx.coroutines.CoroutineScope
|
import kotlinx.coroutines.CoroutineScope
|
||||||
import kotlinx.coroutines.Dispatchers
|
import kotlinx.coroutines.Dispatchers
|
||||||
|
import kotlinx.coroutines.delay
|
||||||
import kotlinx.coroutines.flow.MutableStateFlow
|
import kotlinx.coroutines.flow.MutableStateFlow
|
||||||
import kotlinx.coroutines.flow.StateFlow
|
import kotlinx.coroutines.flow.StateFlow
|
||||||
import kotlinx.coroutines.flow.asStateFlow
|
import kotlinx.coroutines.flow.asStateFlow
|
||||||
@@ -39,9 +41,12 @@ import kotlinx.coroutines.launch
|
|||||||
import kotlinx.coroutines.withContext
|
import kotlinx.coroutines.withContext
|
||||||
import org.webrtc.IceCandidate
|
import org.webrtc.IceCandidate
|
||||||
import org.webrtc.SessionDescription
|
import org.webrtc.SessionDescription
|
||||||
|
import org.webrtc.VideoFrame
|
||||||
|
import org.webrtc.VideoSink
|
||||||
import org.webrtc.VideoTrack
|
import org.webrtc.VideoTrack
|
||||||
import java.util.UUID
|
import java.util.UUID
|
||||||
import java.util.concurrent.CopyOnWriteArrayList
|
import java.util.concurrent.CopyOnWriteArrayList
|
||||||
|
import java.util.concurrent.atomic.AtomicLong
|
||||||
|
|
||||||
private const val TAG = "CallController"
|
private const val TAG = "CallController"
|
||||||
|
|
||||||
@@ -70,6 +75,16 @@ class CallController(
|
|||||||
private val _errorMessage = MutableStateFlow<String?>(null)
|
private val _errorMessage = MutableStateFlow<String?>(null)
|
||||||
val errorMessage: StateFlow<String?> = _errorMessage.asStateFlow()
|
val errorMessage: StateFlow<String?> = _errorMessage.asStateFlow()
|
||||||
|
|
||||||
|
// Remote video frame monitoring — detects when peer stops sending video
|
||||||
|
private val _isRemoteVideoActive = MutableStateFlow(false)
|
||||||
|
val isRemoteVideoActive: StateFlow<Boolean> = _isRemoteVideoActive.asStateFlow()
|
||||||
|
private val lastRemoteFrameTimeMs = AtomicLong(0L)
|
||||||
|
private var remoteVideoMonitorJob: kotlinx.coroutines.Job? = null
|
||||||
|
private val remoteFrameSink =
|
||||||
|
VideoSink { _: VideoFrame ->
|
||||||
|
lastRemoteFrameTimeMs.set(System.currentTimeMillis())
|
||||||
|
}
|
||||||
|
|
||||||
// Audio/video toggle state (UI concerns, not domain state)
|
// Audio/video toggle state (UI concerns, not domain state)
|
||||||
private val _isAudioMuted = MutableStateFlow(false)
|
private val _isAudioMuted = MutableStateFlow(false)
|
||||||
val isAudioMuted: StateFlow<Boolean> = _isAudioMuted.asStateFlow()
|
val isAudioMuted: StateFlow<Boolean> = _isAudioMuted.asStateFlow()
|
||||||
@@ -79,6 +94,8 @@ class CallController(
|
|||||||
val isBluetoothAvailable: StateFlow<Boolean> = audioManager.isBluetoothAvailable
|
val isBluetoothAvailable: StateFlow<Boolean> = audioManager.isBluetoothAvailable
|
||||||
|
|
||||||
init {
|
init {
|
||||||
|
callManager.onRenegotiationOfferReceived = { event -> onRenegotiationOfferReceived(event) }
|
||||||
|
|
||||||
scope.launch {
|
scope.launch {
|
||||||
callManager.state.collect { state ->
|
callManager.state.collect { state ->
|
||||||
when (state) {
|
when (state) {
|
||||||
@@ -268,20 +285,73 @@ class CallController(
|
|||||||
_errorMessage.value = null
|
_errorMessage.value = null
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private fun onRenegotiationOfferReceived(event: CallRenegotiateEvent) {
|
||||||
|
val session = webRtcSession ?: return
|
||||||
|
val sdpOffer = event.sdpOffer()
|
||||||
|
Log.d(TAG) { "Renegotiation offer received, SDP length=${sdpOffer.length}" }
|
||||||
|
|
||||||
|
scope.launch {
|
||||||
|
session.setRemoteDescription(SessionDescription(SessionDescription.Type.OFFER, sdpOffer))
|
||||||
|
session.createAnswer { sdp ->
|
||||||
|
scope.launch {
|
||||||
|
callManager.sendRenegotiationAnswer(sdp.description)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun performRenegotiation() {
|
||||||
|
val session = webRtcSession ?: return
|
||||||
|
val state = callManager.state.value
|
||||||
|
if (state !is CallState.Connected && state !is CallState.Connecting) return
|
||||||
|
|
||||||
|
Log.d(TAG) { "Starting renegotiation" }
|
||||||
|
session.createOffer { sdp ->
|
||||||
|
scope.launch {
|
||||||
|
callManager.sendRenegotiation(sdp.description)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fun cleanup() {
|
fun cleanup() {
|
||||||
audioManager.release()
|
audioManager.release()
|
||||||
stopForegroundService()
|
stopForegroundService()
|
||||||
NotificationUtils.cancelCallNotification(context)
|
NotificationUtils.cancelCallNotification(context)
|
||||||
|
stopRemoteVideoMonitor()
|
||||||
_remoteVideoTrack.value = null
|
_remoteVideoTrack.value = null
|
||||||
_localVideoTrack.value = null
|
_localVideoTrack.value = null
|
||||||
_isAudioMuted.value = false
|
_isAudioMuted.value = false
|
||||||
_isVideoEnabled.value = false
|
_isVideoEnabled.value = false
|
||||||
|
_isRemoteVideoActive.value = false
|
||||||
webRtcSession?.dispose()
|
webRtcSession?.dispose()
|
||||||
webRtcSession = null
|
webRtcSession = null
|
||||||
remoteDescriptionSet.set(false)
|
remoteDescriptionSet.set(false)
|
||||||
pendingIceCandidates.clear()
|
pendingIceCandidates.clear()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private fun startRemoteVideoMonitor(track: VideoTrack) {
|
||||||
|
stopRemoteVideoMonitor()
|
||||||
|
lastRemoteFrameTimeMs.set(System.currentTimeMillis())
|
||||||
|
track.addSink(remoteFrameSink)
|
||||||
|
remoteVideoMonitorJob =
|
||||||
|
scope.launch {
|
||||||
|
while (true) {
|
||||||
|
delay(1500)
|
||||||
|
val elapsed = System.currentTimeMillis() - lastRemoteFrameTimeMs.get()
|
||||||
|
_isRemoteVideoActive.value = elapsed < 2000
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private fun stopRemoteVideoMonitor() {
|
||||||
|
remoteVideoMonitorJob?.cancel()
|
||||||
|
remoteVideoMonitorJob = null
|
||||||
|
try {
|
||||||
|
_remoteVideoTrack.value?.removeSink(remoteFrameSink)
|
||||||
|
} catch (_: Exception) {
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
private fun createWebRtcSession() {
|
private fun createWebRtcSession() {
|
||||||
val session =
|
val session =
|
||||||
WebRtcCallSession(
|
WebRtcCallSession(
|
||||||
@@ -292,9 +362,13 @@ class CallController(
|
|||||||
callManager.onPeerConnected()
|
callManager.onPeerConnected()
|
||||||
startForegroundService()
|
startForegroundService()
|
||||||
},
|
},
|
||||||
onRemoteVideoTrack = { track -> _remoteVideoTrack.value = track },
|
onRemoteVideoTrack = { track ->
|
||||||
|
_remoteVideoTrack.value = track
|
||||||
|
startRemoteVideoMonitor(track)
|
||||||
|
},
|
||||||
onDisconnected = { scope.launch { callManager.hangup() } },
|
onDisconnected = { scope.launch { callManager.hangup() } },
|
||||||
onError = { error -> _errorMessage.value = error },
|
onError = { error -> _errorMessage.value = error },
|
||||||
|
onRenegotiationNeeded = { performRenegotiation() },
|
||||||
)
|
)
|
||||||
try {
|
try {
|
||||||
session.initialize()
|
session.initialize()
|
||||||
|
|||||||
@@ -52,6 +52,7 @@ class WebRtcCallSession(
|
|||||||
private val onRemoteVideoTrack: (VideoTrack) -> Unit,
|
private val onRemoteVideoTrack: (VideoTrack) -> Unit,
|
||||||
private val onDisconnected: () -> Unit,
|
private val onDisconnected: () -> Unit,
|
||||||
private val onError: (String) -> Unit = {},
|
private val onError: (String) -> Unit = {},
|
||||||
|
private val onRenegotiationNeeded: () -> Unit = {},
|
||||||
) {
|
) {
|
||||||
private var peerConnectionFactory: PeerConnectionFactory? = null
|
private var peerConnectionFactory: PeerConnectionFactory? = null
|
||||||
private var peerConnection: PeerConnection? = null
|
private var peerConnection: PeerConnection? = null
|
||||||
@@ -131,7 +132,10 @@ class WebRtcCallSession(
|
|||||||
|
|
||||||
override fun onDataChannel(channel: DataChannel?) {}
|
override fun onDataChannel(channel: DataChannel?) {}
|
||||||
|
|
||||||
override fun onRenegotiationNeeded() {}
|
override fun onRenegotiationNeeded() {
|
||||||
|
Log.d(TAG) { "Renegotiation needed" }
|
||||||
|
this@WebRtcCallSession.onRenegotiationNeeded()
|
||||||
|
}
|
||||||
|
|
||||||
override fun onAddTrack(
|
override fun onAddTrack(
|
||||||
receiver: RtpReceiver?,
|
receiver: RtpReceiver?,
|
||||||
|
|||||||
@@ -353,6 +353,7 @@ private fun ConnectedCallUI(
|
|||||||
val localVideoTrack by (callController?.localVideoTrack ?: emptyVideoFlow).collectAsState()
|
val localVideoTrack by (callController?.localVideoTrack ?: emptyVideoFlow).collectAsState()
|
||||||
val defaultFalse = remember { kotlinx.coroutines.flow.MutableStateFlow(false) }
|
val defaultFalse = remember { kotlinx.coroutines.flow.MutableStateFlow(false) }
|
||||||
val defaultTrue = remember { kotlinx.coroutines.flow.MutableStateFlow(true) }
|
val defaultTrue = remember { kotlinx.coroutines.flow.MutableStateFlow(true) }
|
||||||
|
val isRemoteVideoActive by (callController?.isRemoteVideoActive ?: defaultFalse).collectAsState()
|
||||||
val defaultRoute = remember { kotlinx.coroutines.flow.MutableStateFlow(AudioRoute.EARPIECE) }
|
val defaultRoute = remember { kotlinx.coroutines.flow.MutableStateFlow(AudioRoute.EARPIECE) }
|
||||||
val isAudioMuted by (callController?.isAudioMuted ?: defaultFalse).collectAsState()
|
val isAudioMuted by (callController?.isAudioMuted ?: defaultFalse).collectAsState()
|
||||||
val isVideoEnabled by (callController?.isVideoEnabled ?: defaultTrue).collectAsState()
|
val isVideoEnabled by (callController?.isVideoEnabled ?: defaultTrue).collectAsState()
|
||||||
@@ -364,14 +365,16 @@ private fun ConnectedCallUI(
|
|||||||
.fillMaxSize()
|
.fillMaxSize()
|
||||||
.background(Color.Black),
|
.background(Color.Black),
|
||||||
) {
|
) {
|
||||||
// Remote video (full screen background)
|
// Remote video (full screen background) — only when peer is actively sending
|
||||||
remoteVideoTrack?.let { track ->
|
if (isRemoteVideoActive) {
|
||||||
VideoRenderer(
|
remoteVideoTrack?.let { track ->
|
||||||
videoTrack = track,
|
VideoRenderer(
|
||||||
eglBase = callController?.getEglBase(),
|
videoTrack = track,
|
||||||
modifier = Modifier.fillMaxSize(),
|
eglBase = callController?.getEglBase(),
|
||||||
mirror = false,
|
modifier = Modifier.fillMaxSize(),
|
||||||
)
|
mirror = false,
|
||||||
|
)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Local video (small pip in corner) — only when camera is active
|
// Local video (small pip in corner) — only when camera is active
|
||||||
@@ -390,8 +393,8 @@ private fun ConnectedCallUI(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// If no video, show avatar
|
// If no video or peer stopped sharing, show avatar
|
||||||
if (remoteVideoTrack == null) {
|
if (!isRemoteVideoActive) {
|
||||||
Column(
|
Column(
|
||||||
modifier = Modifier.fillMaxSize(),
|
modifier = Modifier.fillMaxSize(),
|
||||||
horizontalAlignment = Alignment.CenterHorizontally,
|
horizontalAlignment = Alignment.CenterHorizontally,
|
||||||
|
|||||||
+49
-5
@@ -30,6 +30,7 @@ import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallHangupEvent
|
|||||||
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallIceCandidateEvent
|
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallIceCandidateEvent
|
||||||
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallOfferEvent
|
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallOfferEvent
|
||||||
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallRejectEvent
|
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallRejectEvent
|
||||||
|
import com.vitorpamplona.quartz.nipACWebRtcCalls.events.CallRenegotiateEvent
|
||||||
import com.vitorpamplona.quartz.nipACWebRtcCalls.tags.CallType
|
import com.vitorpamplona.quartz.nipACWebRtcCalls.tags.CallType
|
||||||
import com.vitorpamplona.quartz.utils.Log
|
import com.vitorpamplona.quartz.utils.Log
|
||||||
import com.vitorpamplona.quartz.utils.TimeUtils
|
import com.vitorpamplona.quartz.utils.TimeUtils
|
||||||
@@ -54,6 +55,7 @@ class CallManager(
|
|||||||
|
|
||||||
var onAnswerReceived: ((CallAnswerEvent) -> Unit)? = null
|
var onAnswerReceived: ((CallAnswerEvent) -> Unit)? = null
|
||||||
var onIceCandidateReceived: ((CallIceCandidateEvent) -> Unit)? = null
|
var onIceCandidateReceived: ((CallIceCandidateEvent) -> Unit)? = null
|
||||||
|
var onRenegotiationOfferReceived: ((CallRenegotiateEvent) -> Unit)? = null
|
||||||
|
|
||||||
private var timeoutJob: Job? = null
|
private var timeoutJob: Job? = null
|
||||||
private var resetJob: Job? = null
|
private var resetJob: Job? = null
|
||||||
@@ -119,12 +121,26 @@ class CallManager(
|
|||||||
|
|
||||||
fun onCallAnswered(event: CallAnswerEvent) {
|
fun onCallAnswered(event: CallAnswerEvent) {
|
||||||
val current = _state.value
|
val current = _state.value
|
||||||
if (current !is CallState.Offering) return
|
val callId = event.callId()
|
||||||
if (event.callId() != current.callId) return
|
|
||||||
|
|
||||||
_state.value = CallState.Connecting(current.callId, current.peerPubKey, current.callType)
|
when (current) {
|
||||||
cancelTimeout()
|
is CallState.Offering -> {
|
||||||
onAnswerReceived?.invoke(event)
|
if (callId != current.callId) return
|
||||||
|
_state.value = CallState.Connecting(current.callId, current.peerPubKey, current.callType)
|
||||||
|
cancelTimeout()
|
||||||
|
onAnswerReceived?.invoke(event)
|
||||||
|
}
|
||||||
|
|
||||||
|
is CallState.Connected -> {
|
||||||
|
// Renegotiation answer (e.g., peer accepted our video upgrade offer)
|
||||||
|
if (callId != current.callId) return
|
||||||
|
onAnswerReceived?.invoke(event)
|
||||||
|
}
|
||||||
|
|
||||||
|
else -> {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fun onCallRejected(event: CallRejectEvent) {
|
fun onCallRejected(event: CallRejectEvent) {
|
||||||
@@ -139,6 +155,33 @@ class CallManager(
|
|||||||
onIceCandidateReceived?.invoke(event)
|
onIceCandidateReceived?.invoke(event)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fun onRenegotiate(event: CallRenegotiateEvent) {
|
||||||
|
val current = _state.value
|
||||||
|
val callId = event.callId()
|
||||||
|
val currentCallId =
|
||||||
|
when (current) {
|
||||||
|
is CallState.Connecting -> current.callId
|
||||||
|
is CallState.Connected -> current.callId
|
||||||
|
else -> return
|
||||||
|
}
|
||||||
|
if (callId != currentCallId) return
|
||||||
|
onRenegotiationOfferReceived?.invoke(event)
|
||||||
|
}
|
||||||
|
|
||||||
|
suspend fun sendRenegotiation(sdpOffer: String) {
|
||||||
|
val callId = currentCallId() ?: return
|
||||||
|
val peerPubKey = currentPeerPubKey() ?: return
|
||||||
|
val result = factory.createRenegotiate(sdpOffer, peerPubKey, callId, signer)
|
||||||
|
publishEvent(result.wrap)
|
||||||
|
}
|
||||||
|
|
||||||
|
suspend fun sendRenegotiationAnswer(sdpAnswer: String) {
|
||||||
|
val callId = currentCallId() ?: return
|
||||||
|
val peerPubKey = currentPeerPubKey() ?: return
|
||||||
|
val result = factory.createCallAnswer(sdpAnswer, peerPubKey, callId, signer)
|
||||||
|
publishEvent(result.wrap)
|
||||||
|
}
|
||||||
|
|
||||||
fun onPeerConnected() {
|
fun onPeerConnected() {
|
||||||
val current = _state.value
|
val current = _state.value
|
||||||
if (current !is CallState.Connecting) return
|
if (current !is CallState.Connecting) return
|
||||||
@@ -213,6 +256,7 @@ class CallManager(
|
|||||||
is CallRejectEvent -> onCallRejected(event)
|
is CallRejectEvent -> onCallRejected(event)
|
||||||
is CallHangupEvent -> onPeerHangup(event)
|
is CallHangupEvent -> onPeerHangup(event)
|
||||||
is CallIceCandidateEvent -> onIceCandidate(event)
|
is CallIceCandidateEvent -> onIceCandidate(event)
|
||||||
|
is CallRenegotiateEvent -> onRenegotiate(event)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user