From 1a6dba6a3249276f9376d2597abb70da49b46abf Mon Sep 17 00:00:00 2001 From: Miguel Grandez Date: Thu, 17 Sep 2026 15:38:38 -0500 Subject: [PATCH] PTT bidireccional funcional con WebRTC y ActionCable --- app/build.gradle.kts | 7 +- app/src/main/AndroidManifest.xml | 18 +- .../bodycamera/twentyfoulabs/MainActivity.kt | 48 +- .../data/backend/BackendConfig.kt | 5 +- .../twentyfoulabs/data/ptt/PttApiClient.kt | 155 +++ .../twentyfoulabs/data/ptt/PttAudioClient.kt | 916 ++++++++++++++++++ .../twentyfoulabs/data/ptt/PttCableClient.kt | 349 +++++++ .../twentyfoulabs/data/ptt/PttManager.kt | 743 ++++++++++++++ .../main/res/xml/network_security_config.xml | 4 +- 9 files changed, 2209 insertions(+), 36 deletions(-) create mode 100644 app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttApiClient.kt create mode 100644 app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttAudioClient.kt create mode 100644 app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttCableClient.kt create mode 100644 app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttManager.kt diff --git a/app/build.gradle.kts b/app/build.gradle.kts index a95ce0f..889e440 100644 --- a/app/build.gradle.kts +++ b/app/build.gradle.kts @@ -1,4 +1,4 @@ -plugins { +plugins { id("com.android.application") id("org.jetbrains.kotlin.android") id("org.jetbrains.kotlin.plugin.serialization") @@ -89,6 +89,9 @@ dependencies { // Coroutines implementation("org.jetbrains.kotlinx:kotlinx-coroutines-android:1.7.3") + // WebSocket / ActionCable para PTT + implementation("com.squareup.okhttp3:okhttp:4.12.0") + // ViewModel implementation("androidx.lifecycle:lifecycle-viewmodel-compose:2.6.2") implementation("androidx.lifecycle:lifecycle-runtime-compose:2.6.2") @@ -113,3 +116,5 @@ implementation("com.google.mlkit:barcode-scanning:17.3.0") debugImplementation("androidx.compose.ui:ui-tooling") debugImplementation("androidx.compose.ui:ui-test-manifest") } + + diff --git a/app/src/main/AndroidManifest.xml b/app/src/main/AndroidManifest.xml index 80776b0..2ff2d28 100644 --- a/app/src/main/AndroidManifest.xml +++ b/app/src/main/AndroidManifest.xml @@ -8,19 +8,20 @@ - + - + + - + @@ -74,16 +75,7 @@ android:name="android.support.FILE_PROVIDER_PATHS" android:resource="@xml/file_paths" /> - - - - - - - - + diff --git a/app/src/main/kotlin/com/bodycamera/twentyfoulabs/MainActivity.kt b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/MainActivity.kt index 04e3919..7b57b9e 100644 --- a/app/src/main/kotlin/com/bodycamera/twentyfoulabs/MainActivity.kt +++ b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/MainActivity.kt @@ -1,4 +1,4 @@ -package com.bodycamera.twentyfoulabs +package com.bodycamera.twentyfoulabs import android.Manifest import android.content.BroadcastReceiver @@ -27,6 +27,7 @@ import androidx.navigation.compose.composable import androidx.navigation.compose.rememberNavController import com.bodycamera.twentyfoulabs.data.device.M530DeviceController import com.bodycamera.twentyfoulabs.data.auth.AuthSessionStore +import com.bodycamera.twentyfoulabs.data.ptt.PttManager import com.bodycamera.twentyfoulabs.navigation.LocalNavigationHandler import com.bodycamera.twentyfoulabs.navigation.rememberNavigationHandler import com.bodycamera.twentyfoulabs.ui.screens.ConnectionSettingsScreen @@ -173,6 +174,7 @@ class MainActivity : ComponentActivity() { @Inject lateinit var deviceController: M530DeviceController @Inject lateinit var authSessionStore: AuthSessionStore + @Inject lateinit var pttManager: PttManager @Volatile private var isAuthenticated = false private var frontPreviewOpened = false @@ -210,15 +212,8 @@ class MainActivity : ComponentActivity() { return } - val result = deviceController.startPtt() - - Log.i( - TAG, - if (result.isSuccess) - "PTT BROADCAST DOWN - PTT iniciado" - else - "PTT BROADCAST DOWN - ERROR: ${result.exceptionOrNull()?.message}" - ) + val accepted = pttManager.pressToTalk() + Log.i(TAG, if (accepted) "PTT LIVE DOWN - solicitud enviada" else "PTT LIVE DOWN - PTT no preparado") } ACTION_PTT_LONG_PRESS -> { @@ -232,15 +227,8 @@ class MainActivity : ComponentActivity() { return } - val result = deviceController.stopPtt() - - Log.i( - TAG, - if (result.isSuccess) - "PTT BROADCAST UP - PTT detenido" - else - "PTT BROADCAST UP - ERROR: ${result.exceptionOrNull()?.message}" - ) + pttManager.releaseTalk() + Log.i(TAG, "PTT LIVE UP - transmision liberada") } } } @@ -263,6 +251,16 @@ class MainActivity : ComponentActivity() { window.setFlags(WindowManager.LayoutParams.FLAG_KEEP_SCREEN_ON, WindowManager.LayoutParams.FLAG_KEEP_SCREEN_ON) requestPermissions() + pttManager.initialize( + object : PttManager.Listener { + override fun onPttReady() = Unit + override fun onPttDisconnected() = Unit + override fun onTransmitChanged(transmitting: Boolean) = Unit + override fun onChannelBusy(holderName: String?) = Unit + override fun onError(message: String) = Unit + } + ) + onBackPressedDispatcher.addCallback(this, object : OnBackPressedCallback(true) { override fun handleOnBackPressed() = Unit }) @@ -271,6 +269,12 @@ class MainActivity : ComponentActivity() { authSessionStore.token.collect { token -> isAuthenticated = !token.isNullOrBlank() + if (!token.isNullOrBlank()) { + pttManager.start(token) + } else { + pttManager.stop() + } + if (isAuthenticated && !frontPreviewOpened) { frontPreviewOpened = true runOnUiThread { @@ -647,3 +651,9 @@ fun AppNavigation(viewModel: MainViewModel, onOpenPhotoCamera: () -> Unit, onOpe } } } + + + + + + diff --git a/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/backend/BackendConfig.kt b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/backend/BackendConfig.kt index 4b9d5a1..393967d 100644 --- a/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/backend/BackendConfig.kt +++ b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/backend/BackendConfig.kt @@ -1,4 +1,4 @@ -package com.bodycamera.twentyfoulabs.data.backend +package com.bodycamera.twentyfoulabs.data.backend /** Centralized endpoints for the M530 PoC. */ object BackendConfig { @@ -10,10 +10,11 @@ object BackendConfig { const val STREAM_ID = "webrtc_camera_stream" // Rails mobile API exposed through ngrok HTTPS. - const val API_BASE_URL = "https://5ea0-2001-1388-7829-eaff-f431-27c4-8f8e-5b5a.ngrok-free.app" + const val API_BASE_URL = "https://bodycam-prd-web-nefjf74quq-uc.a.run.app" const val DEVICE_SERIAL = "BC-M530-001" const val DEVICE_MODEL = "Recoda M530" val httpBaseUrl: String get() = "http://$HOST:$HTTP_PORT" val whipUrl: String get() = "http://$HOST:8889/$STREAM_ID/whip" } + diff --git a/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttApiClient.kt b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttApiClient.kt new file mode 100644 index 0000000..92ab41c --- /dev/null +++ b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttApiClient.kt @@ -0,0 +1,155 @@ +package com.bodycamera.twentyfoulabs.data.ptt + +import com.bodycamera.twentyfoulabs.data.backend.BackendConfig +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.withContext +import org.json.JSONObject +import java.net.HttpURLConnection +import java.net.URL +import javax.inject.Inject +import javax.inject.Singleton + +@Singleton +class PttApiClient @Inject constructor() { + + suspend fun joinRoom( + token: String, + room: String = "general" + ): PttJoinResult = withContext(Dispatchers.IO) { + + val url = URL( + BackendConfig.API_BASE_URL.trimEnd('/') + + "/mobile/v1/ptt/rooms/$room/join" + ) + + val connection = + (url.openConnection() as HttpURLConnection).apply { + requestMethod = "POST" + connectTimeout = 15_000 + readTimeout = 20_000 + doOutput = true + + setRequestProperty( + "Accept", + "application/json" + ) + + setRequestProperty( + "Content-Type", + "application/json" + ) + + setRequestProperty( + "Authorization", + "Bearer $token" + ) + + setRequestProperty( + "X-Device-Serial", + BackendConfig.DEVICE_SERIAL + ) + } + + try { + connection.outputStream + .bufferedWriter(Charsets.UTF_8) + .use { + it.write("{}") + } + + val code = + connection.responseCode + + val stream = + if (code in 200..299) { + connection.inputStream + } else { + connection.errorStream + } + + val responseBody = + stream + ?.bufferedReader(Charsets.UTF_8) + ?.use { + it.readText() + } + .orEmpty() + + if (code !in 200..299) { + throw IllegalStateException( + "PTT join HTTP $code: $responseBody" + ) + } + + parseJoin( + responseBody + ) + + } finally { + connection.disconnect() + } + } + + private fun parseJoin( + body: String + ): PttJoinResult { + + val root = + JSONObject(body) + + if ( + !root.optBoolean( + "success", + false + ) + ) { + throw IllegalStateException( + root.optString( + "error", + "PTT join sin success" + ) + ) + } + + val room = + root.getJSONObject( + "room" + ) + + val cable = + root.getJSONObject( + "cable" + ) + + return PttJoinResult( + roomSlug = + room.optString( + "slug", + "general" + ), + + cableChannel = + cable.optString( + "channel", + "PttChannel" + ), + + cableRoom = + cable + .optJSONObject( + "params" + ) + ?.optString( + "room", + "general" + ) + ?: "general" + ) + } +} + +data class PttJoinResult( + val roomSlug: String, + val cableChannel: String, + val cableRoom: String +) \ No newline at end of file diff --git a/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttAudioClient.kt b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttAudioClient.kt new file mode 100644 index 0000000..dc16712 --- /dev/null +++ b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttAudioClient.kt @@ -0,0 +1,916 @@ +package com.bodycamera.twentyfoulabs.data.ptt + +import android.content.Context +import android.media.AudioManager +import dagger.hilt.android.qualifiers.ApplicationContext +import android.util.Log +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.delay +import kotlinx.coroutines.launch +import org.webrtc.AudioSource +import org.webrtc.AudioTrack +import org.webrtc.IceCandidate +import org.webrtc.MediaConstraints +import org.webrtc.MediaStream +import org.webrtc.PeerConnection +import org.webrtc.PeerConnectionFactory +import org.webrtc.RtpReceiver +import org.webrtc.SdpObserver +import org.webrtc.SessionDescription +import java.io.BufferedReader +import java.io.InputStreamReader +import java.io.OutputStreamWriter +import java.net.HttpURLConnection +import java.net.URI +import java.net.URL +import java.util.concurrent.ConcurrentHashMap +import javax.inject.Inject +import javax.inject.Singleton + +@Singleton +class PttAudioClient @Inject constructor( + @ApplicationContext private val context: Context +) { + + companion object { + private const val TAG = "PttAudioClient" + + private const val MEDIA_BASE_URL = + "http://35.224.197.127:8889" + + private const val STUN_URL = + "stun:stun.l.google.com:19302" + + private const val MAX_WHEP_RETRIES = 15 + private const val WHEP_RETRY_DELAY_MS = 120L + } + + interface Listener { + fun onPublisherReady() + fun onReceiverReady(userId: Long) + fun onReceiverClosed(userId: Long) + fun onError(message: String) + } + + private val scope = + CoroutineScope( + Dispatchers.IO + SupervisorJob() + ) + + private var factory: PeerConnectionFactory? = null + + private var publisherPeer: PeerConnection? = null + private var audioSource: AudioSource? = null + private var localAudioTrack: AudioTrack? = null + private var publisherSessionUrl: String? = null + + private val receivers = + ConcurrentHashMap() + + private var listener: Listener? = null + + @Volatile + private var initialized = false + + @Volatile + private var transmitting = false + + private data class ReceiverSession( + val peerConnection: PeerConnection, + var sessionUrl: String? = null + ) + + fun initialize(listener: Listener) { + this.listener = listener + + if (initialized) { + return + } + + val audioManager = context.getSystemService(Context.AUDIO_SERVICE) as AudioManager + audioManager.mode = AudioManager.MODE_IN_COMMUNICATION + audioManager.isSpeakerphoneOn = true + Log.i(TAG, "PTT audio routing: mode=${audioManager.mode} speaker=${audioManager.isSpeakerphoneOn}") + + val options = + PeerConnectionFactory.InitializationOptions + .builder(context) + .setEnableInternalTracer(false) + .createInitializationOptions() + + PeerConnectionFactory.initialize(options) + + factory = + PeerConnectionFactory.builder() + .setOptions( + PeerConnectionFactory.Options() + ) + .createPeerConnectionFactory() + + initialized = true + + Log.i(TAG, "PTT WebRTC inicializado") + } + + fun startPublisher( + userId: Long, + whipUrl: String? = null + ) { + scope.launch { + try { + closePublisher() + + val pc = + createPeerConnection( + tag = "TX-$userId", + onRemoteAudioTrack = null + ) + + publisherPeer = pc + + val constraints = + MediaConstraints() + + audioSource = + factory?.createAudioSource( + constraints + ) + + val source = + audioSource + ?: throw IllegalStateException( + "No se pudo crear AudioSource" + ) + + localAudioTrack = + factory?.createAudioTrack( + "ptt_audio_$userId", + source + ) + + val track = + localAudioTrack + ?: throw IllegalStateException( + "No se pudo crear AudioTrack" + ) + + // El microfono comienza muteado. + track.setEnabled(false) + transmitting = false + + pc.addTrack( + track, + listOf("ptt") + ) + + val targetUrl = + normalizeMediaUrl( + whipUrl, + "ptt/u/$userId/whip" + ) + + createOffer( + peerConnection = pc, + receiveAudio = false + ) { offer -> + + scope.launch { + try { + val response = + postSdp( + targetUrl, + offer.description + ) + + publisherSessionUrl = + response.sessionUrl + + setRemoteAnswer( + pc, + response.answer + ) { + Log.i( + TAG, + "WHIP publisher conectado user=$userId" + ) + + listener?.onPublisherReady() + } + + } catch (t: Throwable) { + reportError( + "WHIP publisher: ${t.message}", + t + ) + } + } + } + + } catch (t: Throwable) { + reportError( + "Error iniciando publisher PTT: ${t.message}", + t + ) + } + } + } + + fun connectReceiver( + userId: Long, + whepUrl: String? = null + ) { + if (receivers.containsKey(userId)) { + Log.d( + TAG, + "WHEP user=$userId ya existe" + ) + return + } + + scope.launch { + try { + val pc = + createPeerConnection( + tag = "RX-$userId", + onRemoteAudioTrack = { track -> + Log.i(TAG, "RX-$userId AudioTrack recibido id=${track.id()} state=${track.state()} transmitting=$transmitting") + track.setEnabled(!transmitting) + Log.i(TAG, "RX-$userId AudioTrack enabled=${track.enabled()}") + } + ) + + val session = + ReceiverSession(pc) + + receivers[userId] = session + + val targetUrl = + normalizeMediaUrl( + whepUrl, + "ptt/u/$userId/whep" + ) + + createOffer( + peerConnection = pc, + receiveAudio = true + ) { offer -> + + scope.launch { + connectWhepWithRetry( + userId = userId, + targetUrl = targetUrl, + offer = offer, + session = session + ) + } + } + + } catch (t: Throwable) { + receivers.remove(userId) + + reportError( + "Error creando WHEP user=$userId: ${t.message}", + t + ) + } + } + } + + private suspend fun connectWhepWithRetry( + userId: Long, + targetUrl: String, + offer: SessionDescription, + session: ReceiverSession + ) { + var lastError: Throwable? = null + + for (attempt in 1..MAX_WHEP_RETRIES) { + try { + val response = + postSdp( + targetUrl, + offer.description + ) + + session.sessionUrl = + response.sessionUrl + + setRemoteAnswer( + session.peerConnection, + response.answer + ) { + Log.i( + TAG, + "WHEP conectado user=$userId" + ) + + listener?.onReceiverReady( + userId + ) + } + + return + + } catch (t: Throwable) { + lastError = t + + if (attempt < + MAX_WHEP_RETRIES + ) { + delay( + WHEP_RETRY_DELAY_MS + ) + } + } + } + + closeReceiver(userId) + + reportError( + "WHEP user=$userId no disponible tras " + + "$MAX_WHEP_RETRIES intentos: " + + "${lastError?.message}", + lastError + ) + } + + fun setTransmitEnabled( + enabled: Boolean + ) { + transmitting = enabled + + localAudioTrack?.setEnabled( + enabled + ) + + receivers.values.forEach { + setRemoteAudioEnabled( + it.peerConnection, + !enabled + ) + } + + Log.i( + TAG, + if (enabled) { + "PTT TX ON / RX OFF" + } else { + "PTT TX OFF / RX ON" + } + ) + } + + fun forceReceiveMode() { + setTransmitEnabled(false) + } + + fun closeReceiver( + userId: Long + ) { + val session = + receivers.remove(userId) + ?: return + + scope.launch { + deleteSession( + session.sessionUrl + ) + + try { + session.peerConnection.close() + session.peerConnection.dispose() + } catch (_: Throwable) { + } + + listener?.onReceiverClosed( + userId + ) + + Log.i( + TAG, + "WHEP cerrado user=$userId" + ) + } + } + + private suspend fun closePublisher() { + setTransmitEnabled(false) + + deleteSession( + publisherSessionUrl + ) + + publisherSessionUrl = null + + try { + publisherPeer?.close() + publisherPeer?.dispose() + } catch (_: Throwable) { + } + + publisherPeer = null + + try { + localAudioTrack?.dispose() + } catch (_: Throwable) { + } + + localAudioTrack = null + + try { + audioSource?.dispose() + } catch (_: Throwable) { + } + + audioSource = null + } + + fun stopAll() { + scope.launch { + closePublisher() + + val ids = + receivers.keys.toList() + + ids.forEach { + closeReceiver(it) + } + + transmitting = false + + Log.i( + TAG, + "Sesiones PTT cerradas" + ) + } + } + + fun release() { + stopAll() + + scope.launch { + delay(200) + + try { + factory?.dispose() + } catch (_: Throwable) { + } + + factory = null + initialized = false + listener = null + + scope.cancel() + } + } + + private fun createPeerConnection( + tag: String, + onRemoteAudioTrack: ((AudioTrack) -> Unit)? + ): PeerConnection { + + val iceServer = + PeerConnection.IceServer + .builder(STUN_URL) + .createIceServer() + + val rtcConfig = + PeerConnection.RTCConfiguration( + listOf(iceServer) + ) + + rtcConfig.iceTransportsType = + PeerConnection.IceTransportsType.ALL + + val observer = + object : PeerConnection.Observer { + + override fun onSignalingChange( + state: PeerConnection.SignalingState? + ) { + Log.d( + TAG, + "$tag signaling=$state" + ) + } + + override fun onIceConnectionChange( + state: PeerConnection.IceConnectionState? + ) { + Log.i( + TAG, + "$tag ICE=$state" + ) + } + + override fun onIceConnectionReceivingChange( + receiving: Boolean + ) { + } + + override fun onIceGatheringChange( + state: PeerConnection.IceGatheringState? + ) { + Log.d( + TAG, + "$tag gathering=$state" + ) + } + + override fun onIceCandidate( + candidate: IceCandidate? + ) { + // MediaMTX puede usar el SDP completo + // una vez finalizado ICE gathering. + } + + override fun onIceCandidatesRemoved( + candidates: Array? + ) { + } + + override fun onAddStream( + stream: MediaStream? + ) { + stream + ?.audioTracks + ?.forEach { track -> + onRemoteAudioTrack?.invoke( + track + ) + } + } + + override fun onRemoveStream( + stream: MediaStream? + ) { + } + + override fun onDataChannel( + channel: org.webrtc.DataChannel? + ) { + } + + override fun onRenegotiationNeeded() { + } + + override fun onAddTrack( + receiver: RtpReceiver?, + mediaStreams: Array? + ) { + val track = + receiver?.track() + + if (track is AudioTrack) { + onRemoteAudioTrack?.invoke( + track + ) + } + } + } + + return factory?.createPeerConnection( + rtcConfig, + observer + ) ?: throw IllegalStateException( + "No se pudo crear PeerConnection" + ) + } + + private fun createOffer( + peerConnection: PeerConnection, + receiveAudio: Boolean, + onOffer: (SessionDescription) -> Unit + ) { + val constraints = + MediaConstraints().apply { + mandatory.add( + MediaConstraints.KeyValuePair( + "OfferToReceiveAudio", + if (receiveAudio) { + "true" + } else { + "false" + } + ) + ) + + mandatory.add( + MediaConstraints.KeyValuePair( + "OfferToReceiveVideo", + "false" + ) + ) + } + + peerConnection.createOffer( + object : SdpObserver { + + override fun onCreateSuccess( + sdp: SessionDescription? + ) { + if (sdp == null) { + reportError( + "SDP offer PTT vacio", + null + ) + return + } + + peerConnection.setLocalDescription( + object : SdpObserver { + + override fun onSetSuccess() { + onOffer(sdp) + } + + override fun onSetFailure( + error: String? + ) { + reportError( + "No se pudo establecer SDP local: $error", + null + ) + } + + override fun onCreateSuccess( + p0: SessionDescription? + ) { + } + + override fun onCreateFailure( + p0: String? + ) { + } + }, + sdp + ) + } + + override fun onCreateFailure( + error: String? + ) { + reportError( + "No se pudo crear SDP offer: $error", + null + ) + } + + override fun onSetSuccess() { + } + + override fun onSetFailure( + p0: String? + ) { + } + }, + constraints + ) + } + + private fun setRemoteAnswer( + peerConnection: PeerConnection, + answer: String, + onReady: () -> Unit + ) { + val description = + SessionDescription( + SessionDescription.Type.ANSWER, + answer + ) + + peerConnection.setRemoteDescription( + object : SdpObserver { + + override fun onSetSuccess() { + onReady() + } + + override fun onSetFailure( + error: String? + ) { + reportError( + "No se pudo establecer SDP remoto: $error", + null + ) + } + + override fun onCreateSuccess( + p0: SessionDescription? + ) { + } + + override fun onCreateFailure( + p0: String? + ) { + } + }, + description + ) + } + + private data class SdpResponse( + val answer: String, + val sessionUrl: String? + ) + + private fun postSdp( + endpoint: String, + sdp: String + ): SdpResponse { + + val connection = + URL(endpoint) + .openConnection() + as HttpURLConnection + + connection.requestMethod = "POST" + connection.connectTimeout = 15_000 + connection.readTimeout = 20_000 + connection.doOutput = true + + connection.setRequestProperty( + "Content-Type", + "application/sdp" + ) + + connection.setRequestProperty( + "Accept", + "application/sdp" + ) + + try { + OutputStreamWriter( + connection.outputStream, + Charsets.UTF_8 + ).use { + it.write(sdp) + it.flush() + } + + val code = + connection.responseCode + + if (code !in 200..299) { + val error = + connection.errorStream + ?.bufferedReader() + ?.use { it.readText() } + .orEmpty() + + throw IllegalStateException( + "HTTP $code: $error" + ) + } + + val answer = + BufferedReader( + InputStreamReader( + connection.inputStream + ) + ).use { + it.readText() + } + + val location = + connection.getHeaderField( + "Location" + ) + + return SdpResponse( + answer = answer, + sessionUrl = + resolveSessionUrl( + endpoint, + location + ) + ) + + } finally { + connection.disconnect() + } + } + + private fun deleteSession( + sessionUrl: String? + ) { + if (sessionUrl.isNullOrBlank()) { + return + } + + try { + val connection = + URL(sessionUrl) + .openConnection() + as HttpURLConnection + + connection.requestMethod = + "DELETE" + + connection.connectTimeout = + 10_000 + + connection.readTimeout = + 10_000 + + try { + val code = + connection.responseCode + + Log.d( + TAG, + "DELETE session HTTP $code" + ) + } finally { + connection.disconnect() + } + + } catch (t: Throwable) { + Log.w( + TAG, + "No se pudo cerrar session: ${t.message}" + ) + } + } + + private fun resolveSessionUrl( + endpoint: String, + location: String? + ): String? { + + if (location.isNullOrBlank()) { + return null + } + + return try { + URI(endpoint) + .resolve(location) + .toString() + } catch (_: Throwable) { + location + } + } + + private fun normalizeMediaUrl( + suppliedUrl: String?, + fallbackPath: String + ): String { + + if (suppliedUrl.isNullOrBlank()) { + return "$MEDIA_BASE_URL/$fallbackPath" + } + + return when { + suppliedUrl.startsWith( + "http://" + ) || + suppliedUrl.startsWith( + "https://" + ) -> suppliedUrl + + else -> + "$MEDIA_BASE_URL/" + + suppliedUrl.trimStart('/') + } + } + + private fun setRemoteAudioEnabled( + peerConnection: PeerConnection, + enabled: Boolean + ) { + peerConnection.receivers + .mapNotNull { + it.track() + } + .filterIsInstance() + .forEach { + it.setEnabled(enabled) + } + } + + private fun reportError( + message: String, + throwable: Throwable? + ) { + if (throwable != null) { + Log.e( + TAG, + message, + throwable + ) + } else { + Log.e( + TAG, + message + ) + } + + listener?.onError( + message + ) + } +} + diff --git a/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttCableClient.kt b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttCableClient.kt new file mode 100644 index 0000000..57db518 --- /dev/null +++ b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttCableClient.kt @@ -0,0 +1,349 @@ +package com.bodycamera.twentyfoulabs.data.ptt + +import android.os.Handler +import android.os.Looper +import android.util.Log +import okhttp3.OkHttpClient +import okhttp3.Request +import okhttp3.Response +import okhttp3.WebSocket +import okhttp3.WebSocketListener +import org.json.JSONObject +import java.util.concurrent.TimeUnit +import javax.inject.Inject +import javax.inject.Singleton + +@Singleton +class PttCableClient @Inject constructor() { + + companion object { + private const val TAG = "PttCableClient" + private const val CABLE_URL = + "wss://bodycam-prd-web-nefjf74quq-uc.a.run.app/cable" + private const val RECONNECT_DELAY_MS = 3000L + } + + interface Listener { + fun onConnected() + fun onDisconnected() + fun onEvent(event: String, payload: JSONObject) + fun onError(message: String) + } + + private val client = + OkHttpClient.Builder() + .pingInterval(20, TimeUnit.SECONDS) + .connectTimeout(15, TimeUnit.SECONDS) + .build() + + private val reconnectHandler = Handler(Looper.getMainLooper()) + + private var socket: WebSocket? = null + private var listener: Listener? = null + private var identifier: String? = null + + private var lastToken: String? = null + private var lastRoom: String = "general" + + @Volatile + private var subscribed = false + + @Volatile + private var manualDisconnect = false + + @Volatile + private var reconnectScheduled = false + + fun connect( + token: String, + room: String = "general", + listener: Listener + ) { + disconnectInternal(clearSession = false) + + this.listener = listener + this.lastToken = token + this.lastRoom = room + + manualDisconnect = false + reconnectScheduled = false + + connectInternal(token, room) + } + + private fun connectInternal( + token: String, + room: String + ) { + if (manualDisconnect) return + + subscribed = false + + identifier = + JSONObject() + .put("channel", "PttChannel") + .put("room", room) + .toString() + + val request = + Request.Builder() + .url("$CABLE_URL?token=$token") + .build() + + Log.i(TAG, "Conectando ActionCable room=$room") + + val newSocket = + client.newWebSocket( + request, + object : WebSocketListener() { + + override fun onOpen( + webSocket: WebSocket, + response: Response + ) { + Log.i(TAG, "WebSocket abierto") + } + + override fun onMessage( + webSocket: WebSocket, + text: String + ) { + handleMessage(webSocket, text) + } + + override fun onClosing( + webSocket: WebSocket, + code: Int, + reason: String + ) { + Log.i( + TAG, + "WebSocket cerrando code=$code reason=$reason" + ) + webSocket.close(code, reason) + } + + override fun onClosed( + webSocket: WebSocket, + code: Int, + reason: String + ) { + if (socket !== webSocket) return + + subscribed = false + socket = null + + Log.i( + TAG, + "WebSocket cerrado code=$code reason=$reason" + ) + + listener?.onDisconnected() + + if (!manualDisconnect) { + scheduleReconnect() + } + } + + override fun onFailure( + webSocket: WebSocket, + t: Throwable, + response: Response? + ) { + if (socket !== webSocket) return + + subscribed = false + socket = null + + val message = + "ActionCable error: ${t.message}" + + Log.e(TAG, message, t) + + listener?.onError(message) + listener?.onDisconnected() + + if (!manualDisconnect) { + scheduleReconnect() + } + } + } + ) + + socket = newSocket + } + + private fun scheduleReconnect() { + if (manualDisconnect || reconnectScheduled) return + if (lastToken == null || listener == null) return + + reconnectScheduled = true + + Log.i( + TAG, + "Reconexión ActionCable programada en ${RECONNECT_DELAY_MS}ms" + ) + + reconnectHandler.postDelayed( + { + reconnectScheduled = false + + if (manualDisconnect) { + return@postDelayed + } + + val token = lastToken ?: return@postDelayed + val currentListener = listener ?: return@postDelayed + + Log.i(TAG, "Reintentando conexión ActionCable...") + + this.listener = currentListener + connectInternal(token, lastRoom) + }, + RECONNECT_DELAY_MS + ) + } + + private fun handleMessage( + webSocket: WebSocket, + text: String + ) { + try { + val root = JSONObject(text) + + when (root.optString("type")) { + + "welcome" -> { + Log.i(TAG, "ActionCable welcome") + subscribe(webSocket) + return + } + + "ping" -> return + + "confirm_subscription" -> { + subscribed = true + reconnectScheduled = false + Log.i(TAG, "PttChannel suscrito") + listener?.onConnected() + return + } + + "reject_subscription" -> { + subscribed = false + listener?.onError( + "PttChannel rechazo la suscripcion" + ) + + if (!manualDisconnect) { + scheduleReconnect() + } + return + } + } + + val message = + root.optJSONObject("message") + ?: return + + val event = + message.optString("event") + + if (event.isNotBlank()) { + Log.i(TAG, "Evento PTT: $event payload=${message}") + listener?.onEvent(event, message) + } + + } catch (t: Throwable) { + Log.e(TAG, "Error procesando ActionCable", t) + + listener?.onError( + "ActionCable parse error: ${t.message}" + ) + } + } + + private fun subscribe( + webSocket: WebSocket + ) { + val id = identifier ?: return + + val command = + JSONObject() + .put("command", "subscribe") + .put("identifier", id) + + webSocket.send(command.toString()) + } + + fun perform( + action: String, + data: JSONObject = JSONObject() + ): Boolean { + + if (!subscribed) { + Log.w( + TAG, + "perform $action ignorado: no suscrito" + ) + return false + } + + val id = identifier ?: return false + + val actionData = + JSONObject(data.toString()) + .put("action", action) + + val command = + JSONObject() + .put("command", "message") + .put("identifier", id) + .put("data", actionData.toString()) + + Log.i(TAG, "perform PTT: $action") + + return socket?.send( + command.toString() + ) == true + } + + fun publisherReady(): Boolean = + perform("publisher_ready") + + fun requestFloor(): Boolean = + perform("ptt_request") + + fun releaseFloor(): Boolean = + perform("ptt_release") + + private fun disconnectInternal( + clearSession: Boolean + ) { + subscribed = false + + reconnectHandler.removeCallbacksAndMessages(null) + reconnectScheduled = false + + val oldSocket = socket + socket = null + + oldSocket?.close( + 1000, + "PTT disconnect" + ) + + identifier = null + + if (clearSession) { + lastToken = null + lastRoom = "general" + listener = null + } + } + + fun disconnect() { + manualDisconnect = true + disconnectInternal(clearSession = true) + } +} \ No newline at end of file diff --git a/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttManager.kt b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttManager.kt new file mode 100644 index 0000000..aa07b60 --- /dev/null +++ b/app/src/main/kotlin/com/bodycamera/twentyfoulabs/data/ptt/PttManager.kt @@ -0,0 +1,743 @@ +package com.bodycamera.twentyfoulabs.data.ptt + +import android.util.Log +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.launch +import org.json.JSONArray +import org.json.JSONObject +import javax.inject.Inject +import javax.inject.Singleton + +@Singleton +class PttManager @Inject constructor( + private val apiClient: PttApiClient, + private val cableClient: PttCableClient, + private val audioClient: PttAudioClient +) { + + companion object { + private const val TAG = "PttManager" + private const val DEFAULT_ROOM = "general" + } + + interface Listener { + fun onPttReady() + fun onPttDisconnected() + fun onTransmitChanged(transmitting: Boolean) + fun onChannelBusy(holderName: String?) + fun onError(message: String) + } + + private val scope = + CoroutineScope( + Dispatchers.IO + SupervisorJob() + ) + + private var listener: Listener? = null + + @Volatile + private var started = false + + @Volatile + private var cableConnected = false + + @Volatile + private var publisherReady = false + + @Volatile + private var buttonPressed = false + + private var myUserId: Long? = null + + private val participants = + mutableMapOf() + + private data class Participant( + val userId: Long, + val name: String?, + val publishing: Boolean, + val whipUrl: String?, + val whepUrl: String? + ) + + fun initialize( + listener: Listener + ) { + this.listener = listener + + audioClient.initialize( + object : PttAudioClient.Listener { + + override fun onPublisherReady() { + publisherReady = true + + Log.i( + TAG, + "Publisher WHIP listo" + ) + + cableClient.publisherReady() + + listener.onPttReady() + } + + override fun onReceiverReady( + userId: Long + ) { + Log.i( + TAG, + "WHEP listo user=$userId" + ) + } + + override fun onReceiverClosed( + userId: Long + ) { + Log.i( + TAG, + "WHEP cerrado user=$userId" + ) + } + + override fun onError( + message: String + ) { + Log.e( + TAG, + "Audio PTT: $message" + ) + + listener.onError( + message + ) + } + } + ) + } + + fun start( + token: String, + room: String = DEFAULT_ROOM + ) { + if (token.isBlank()) { + listener?.onError( + "Token PTT vacio" + ) + return + } + + if (started) { + Log.d( + TAG, + "PTT ya iniciado" + ) + return + } + + started = true + cableConnected = false + publisherReady = false + buttonPressed = false + myUserId = null + participants.clear() + + scope.launch { + try { + Log.i( + TAG, + "Uniendo PTT room=$room" + ) + + val join = + apiClient.joinRoom( + token = token, + room = room + ) + + if (!started) { + return@launch + } + + Log.i( + TAG, + "Join PTT OK room=${join.roomSlug}" + ) + + connectCable( + token = token, + room = join.cableRoom + ) + + } catch (t: Throwable) { + started = false + + val message = + "No se pudo unir al PTT: ${t.message}" + + Log.e( + TAG, + message, + t + ) + + listener?.onError( + message + ) + } + } + } + + private fun connectCable( + token: String, + room: String + ) { + cableClient.connect( + token = token, + room = room, + listener = + object : PttCableClient.Listener { + + override fun onConnected() { + cableConnected = true + + Log.i( + TAG, + "ActionCable conectado" + ) + } + + override fun onDisconnected() { + cableConnected = false + + Log.i( + TAG, + "ActionCable desconectado" + ) + + listener?.onPttDisconnected() + } + + override fun onEvent( + event: String, + payload: JSONObject + ) { + handleCableEvent( + event, + payload + ) + } + + override fun onError( + message: String + ) { + Log.e( + TAG, + message + ) + + listener?.onError( + message + ) + } + } + ) + } + + private fun handleCableEvent( + event: String, + payload: JSONObject + ) { + Log.i( + TAG, + "Evento servidor: $event" + ) + + when (event) { + + "joined" -> + handleJoined( + payload + ) + + "participant_joined" -> + handleParticipantJoined( + payload + ) + + "publisher_ready" -> + handlePublisherReady( + payload + ) + + "participant_left" -> + handleParticipantLeft( + payload + ) + + "floor_granted" -> + handleFloorGranted( + payload + ) + + "floor_released" -> { + buttonPressed = false + + audioClient.forceReceiveMode() + + listener?.onTransmitChanged( + false + ) + } + + "floor_denied" -> + handleFloorDenied( + payload + ) + } + } + + private fun handleJoined( + payload: JSONObject + ) { + val ownJson = + payload.optJSONObject( + "participant" + ) + + if (ownJson == null) { + listener?.onError( + "Evento joined sin participant" + ) + return + } + + val own = + parseParticipant( + ownJson + ) + + if (own == null) { + listener?.onError( + "Participant propio invalido" + ) + return + } + + myUserId = + own.userId + + participants[own.userId] = + own + + Log.i( + TAG, + "PTT joined myUserId=${own.userId}" + ) + + /* + * Abrimos nuestro publisher WHIP una sola vez. + * El AudioTrack comienza muteado. + */ + audioClient.startPublisher( + userId = own.userId, + whipUrl = own.whipUrl + ) + + /* + * Si al entrar ya hay participantes publicando, + * abrimos sus WHEP. + */ + val others = + payload.optJSONArray( + "participants" + ) ?: JSONArray() + + for (index in 0 until others.length()) { + + val json = + others.optJSONObject(index) + ?: continue + + val participant = + parseParticipant(json) + ?: continue + + participants[participant.userId] = + participant + + if ( + participant.userId != + myUserId && + participant.publishing + ) { + connectReceiver( + participant + ) + } + } + } + + private fun handleParticipantJoined( + payload: JSONObject + ) { + val json = + payload.optJSONObject( + "participant" + ) + + if (json != null) { + parseParticipant(json) + ?.let { + participants[it.userId] = + it + } + } + + /* + * El backend indica que si ya publicamos, + * debemos volver a anunciar publisher_ready + * para que el nuevo participante nos escuche. + */ + if (publisherReady) { + cableClient.publisherReady() + } + } + + private fun handlePublisherReady( + payload: JSONObject + ) { + val json = + payload.optJSONObject( + "participant" + ) + + if (json == null) { + Log.w( + TAG, + "publisher_ready sin participant" + ) + return + } + + val participant = + parseParticipant(json) + ?: return + + participants[participant.userId] = + participant + + if ( + participant.userId == + myUserId + ) { + return + } + + connectReceiver( + participant + ) + } + + private fun connectReceiver( + participant: Participant + ) { + Log.i( + TAG, + "Abriendo WHEP user=${participant.userId}" + ) + + audioClient.connectReceiver( + userId = participant.userId, + whepUrl = participant.whepUrl + ) + } + + private fun handleParticipantLeft( + payload: JSONObject + ) { + val userId = + extractUserId( + payload + ) + + if (userId <= 0L) { + return + } + + participants.remove( + userId + ) + + if ( + userId != + myUserId + ) { + audioClient.closeReceiver( + userId + ) + } + } + + private fun handleFloorGranted( + payload: JSONObject + ) { + val floor = + payload + .optJSONObject("room") + ?.optJSONObject("floor") + ?: payload.optJSONObject( + "floor" + ) + + val holderId = + floor?.optLong( + "holder_user_id", + -1L + ) ?: -1L + + val mine = + holderId > 0L && + holderId == myUserId + + if (mine) { + /* + * Normalmente ya estamos transmitiendo + * porque F4 DOWN hace unmute local inmediato. + */ + if (buttonPressed) { + audioClient.setTransmitEnabled( + true + ) + + listener?.onTransmitChanged( + true + ) + } + } else { + /* + * Otro participante tiene el floor. + */ + buttonPressed = false + + audioClient.forceReceiveMode() + + listener?.onTransmitChanged( + false + ) + } + } + + private fun handleFloorDenied( + payload: JSONObject + ) { + buttonPressed = false + + audioClient.forceReceiveMode() + + listener?.onTransmitChanged( + false + ) + + val holderName = + payload + .optJSONObject("room") + ?.optJSONObject("floor") + ?.optString( + "holder_name" + ) + ?.takeIf { + it.isNotBlank() + } + + listener?.onChannelBusy( + holderName + ) + } + + /** + * F4 DOWN. + * + * El cambio de audio es inmediato. + * No esperamos respuesta del servidor para + * que el boton se sienta como radio. + */ + fun pressToTalk(): Boolean { + + if ( + !started || + !cableConnected || + !publisherReady + ) { + Log.w( + TAG, + "F4 DOWN ignorado: PTT no listo" + ) + return false + } + + if (buttonPressed) { + return true + } + + buttonPressed = true + + audioClient.setTransmitEnabled( + true + ) + + listener?.onTransmitChanged( + true + ) + + val sent = + cableClient.requestFloor() + + if (!sent) { + buttonPressed = false + + audioClient.forceReceiveMode() + + listener?.onTransmitChanged( + false + ) + } + + return sent + } + + /** + * F4 UP. + * + * Primero muteamos localmente y volvemos + * a escuchar; despues notificamos al backend. + */ + fun releaseTalk(): Boolean { + + buttonPressed = false + + audioClient.forceReceiveMode() + + listener?.onTransmitChanged( + false + ) + + if ( + !started || + !cableConnected + ) { + return false + } + + return cableClient.releaseFloor() + } + + fun stop() { + buttonPressed = false + publisherReady = false + cableConnected = false + started = false + + audioClient.forceReceiveMode() + audioClient.stopAll() + + cableClient.disconnect() + + participants.clear() + myUserId = null + + listener?.onPttDisconnected() + + Log.i( + TAG, + "PTT detenido" + ) + } + + fun isReady(): Boolean = + started && + cableConnected && + publisherReady + + private fun parseParticipant( + json: JSONObject + ): Participant? { + + val userId = + json.optLong( + "user_id", + -1L + ) + + if (userId <= 0L) { + return null + } + + val media = + json.optJSONObject( + "media" + ) + + return Participant( + userId = userId, + + name = + json.optString( + "name" + ).takeIf { + it.isNotBlank() + }, + + publishing = + json.optBoolean( + "publishing", + false + ), + + whipUrl = + media + ?.optString( + "whip_url" + ) + ?.takeIf { + it.isNotBlank() + }, + + whepUrl = + media + ?.optString( + "whep_url" + ) + ?.takeIf { + it.isNotBlank() + } + ) + } + + private fun extractUserId( + payload: JSONObject + ): Long { + + val direct = + payload.optLong( + "user_id", + -1L + ) + + if (direct > 0L) { + return direct + } + + return payload + .optJSONObject( + "participant" + ) + ?.optLong( + "user_id", + -1L + ) + ?: -1L + } +} + diff --git a/app/src/main/res/xml/network_security_config.xml b/app/src/main/res/xml/network_security_config.xml index 7e29385..950f75b 100644 --- a/app/src/main/res/xml/network_security_config.xml +++ b/app/src/main/res/xml/network_security_config.xml @@ -1,6 +1,6 @@ - + 192.168.0.25 @@ -9,6 +9,7 @@ 172.16 localhost 127.0.0.1 + 35.224.197.127 @@ -19,3 +20,4 @@ +