PTT bidireccional funcional con WebRTC y ActionCable

This commit is contained in:
Miguel Angel 2026-09-17 15:38:38 -05:00
parent 0d87c0e4d9
commit 1a6dba6a32
9 changed files with 2209 additions and 36 deletions

View file

@ -1,4 +1,4 @@
plugins { plugins {
id("com.android.application") id("com.android.application")
id("org.jetbrains.kotlin.android") id("org.jetbrains.kotlin.android")
id("org.jetbrains.kotlin.plugin.serialization") id("org.jetbrains.kotlin.plugin.serialization")
@ -89,6 +89,9 @@ dependencies {
// Coroutines // Coroutines
implementation("org.jetbrains.kotlinx:kotlinx-coroutines-android:1.7.3") implementation("org.jetbrains.kotlinx:kotlinx-coroutines-android:1.7.3")
// WebSocket / ActionCable para PTT
implementation("com.squareup.okhttp3:okhttp:4.12.0")
// ViewModel // ViewModel
implementation("androidx.lifecycle:lifecycle-viewmodel-compose:2.6.2") implementation("androidx.lifecycle:lifecycle-viewmodel-compose:2.6.2")
implementation("androidx.lifecycle:lifecycle-runtime-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-tooling")
debugImplementation("androidx.compose.ui:ui-test-manifest") debugImplementation("androidx.compose.ui:ui-test-manifest")
} }

View file

@ -8,19 +8,20 @@
<uses-permission android:name="android.permission.CHANGE_WIFI_STATE" /> <uses-permission android:name="android.permission.CHANGE_WIFI_STATE" />
<uses-permission android:name="android.permission.CHANGE_NETWORK_STATE" /> <uses-permission android:name="android.permission.CHANGE_NETWORK_STATE" />
<!-- Permisos para configuración WiFi (requieren ubicación) --> <!-- Permisos para configuración WiFi (requieren ubicación) -->
<uses-permission android:name="android.permission.ACCESS_FINE_LOCATION" /> <uses-permission android:name="android.permission.ACCESS_FINE_LOCATION" />
<uses-permission android:name="android.permission.ACCESS_COARSE_LOCATION" /> <uses-permission android:name="android.permission.ACCESS_COARSE_LOCATION" />
<!-- Permisos para cámara y streaming --> <!-- Permisos para cámara y streaming -->
<uses-permission android:name="android.permission.CAMERA" /> <uses-permission android:name="android.permission.CAMERA" />
<uses-permission android:name="android.permission.RECORD_AUDIO" /> <uses-permission android:name="android.permission.RECORD_AUDIO" />
<uses-permission android:name="android.permission.MODIFY_AUDIO_SETTINGS" />
<uses-permission android:name="android.permission.VIBRATE" /> <uses-permission android:name="android.permission.VIBRATE" />
<uses-permission android:name="android.permission.FOREGROUND_SERVICE" /> <uses-permission android:name="android.permission.FOREGROUND_SERVICE" />
<!-- Características de hardware --> <!-- Características de hardware -->
<uses-feature android:name="android.hardware.camera" android:required="false" /> <uses-feature android:name="android.hardware.camera" android:required="false" />
<uses-feature android:name="android.hardware.camera.autofocus" android:required="false" /> <uses-feature android:name="android.hardware.camera.autofocus" android:required="false" />
<uses-feature android:name="android.hardware.camera.any" android:required="false" /> <uses-feature android:name="android.hardware.camera.any" android:required="false" />
@ -74,16 +75,7 @@
android:name="android.support.FILE_PROVIDER_PATHS" android:name="android.support.FILE_PROVIDER_PATHS"
android:resource="@xml/file_paths" /> android:resource="@xml/file_paths" />
</provider> </provider>
<receiver
android:name=".receiver.M530PttReceiver"
android:exported="true">
<intent-filter>
<action android:name="android.intent.action.ACTION_PTTKEY_DOWN" />
<action android:name="android.intent.action.ACTION_PTTKEY_LONG_PRESS" />
<action android:name="android.intent.action.ACTION_PTTKEY_UP" />
</intent-filter>
</receiver>
</application> </application>
</manifest> </manifest>

View file

@ -1,4 +1,4 @@
package com.bodycamera.twentyfoulabs package com.bodycamera.twentyfoulabs
import android.Manifest import android.Manifest
import android.content.BroadcastReceiver import android.content.BroadcastReceiver
@ -27,6 +27,7 @@ import androidx.navigation.compose.composable
import androidx.navigation.compose.rememberNavController import androidx.navigation.compose.rememberNavController
import com.bodycamera.twentyfoulabs.data.device.M530DeviceController import com.bodycamera.twentyfoulabs.data.device.M530DeviceController
import com.bodycamera.twentyfoulabs.data.auth.AuthSessionStore 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.LocalNavigationHandler
import com.bodycamera.twentyfoulabs.navigation.rememberNavigationHandler import com.bodycamera.twentyfoulabs.navigation.rememberNavigationHandler
import com.bodycamera.twentyfoulabs.ui.screens.ConnectionSettingsScreen import com.bodycamera.twentyfoulabs.ui.screens.ConnectionSettingsScreen
@ -173,6 +174,7 @@ class MainActivity : ComponentActivity() {
@Inject lateinit var deviceController: M530DeviceController @Inject lateinit var deviceController: M530DeviceController
@Inject lateinit var authSessionStore: AuthSessionStore @Inject lateinit var authSessionStore: AuthSessionStore
@Inject lateinit var pttManager: PttManager
@Volatile private var isAuthenticated = false @Volatile private var isAuthenticated = false
private var frontPreviewOpened = false private var frontPreviewOpened = false
@ -210,15 +212,8 @@ class MainActivity : ComponentActivity() {
return return
} }
val result = deviceController.startPtt() val accepted = pttManager.pressToTalk()
Log.i(TAG, if (accepted) "PTT LIVE DOWN - solicitud enviada" else "PTT LIVE DOWN - PTT no preparado")
Log.i(
TAG,
if (result.isSuccess)
"PTT BROADCAST DOWN - PTT iniciado"
else
"PTT BROADCAST DOWN - ERROR: ${result.exceptionOrNull()?.message}"
)
} }
ACTION_PTT_LONG_PRESS -> { ACTION_PTT_LONG_PRESS -> {
@ -232,15 +227,8 @@ class MainActivity : ComponentActivity() {
return return
} }
val result = deviceController.stopPtt() pttManager.releaseTalk()
Log.i(TAG, "PTT LIVE UP - transmision liberada")
Log.i(
TAG,
if (result.isSuccess)
"PTT BROADCAST UP - PTT detenido"
else
"PTT BROADCAST UP - ERROR: ${result.exceptionOrNull()?.message}"
)
} }
} }
} }
@ -263,6 +251,16 @@ class MainActivity : ComponentActivity() {
window.setFlags(WindowManager.LayoutParams.FLAG_KEEP_SCREEN_ON, WindowManager.LayoutParams.FLAG_KEEP_SCREEN_ON) window.setFlags(WindowManager.LayoutParams.FLAG_KEEP_SCREEN_ON, WindowManager.LayoutParams.FLAG_KEEP_SCREEN_ON)
requestPermissions() 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) { onBackPressedDispatcher.addCallback(this, object : OnBackPressedCallback(true) {
override fun handleOnBackPressed() = Unit override fun handleOnBackPressed() = Unit
}) })
@ -271,6 +269,12 @@ class MainActivity : ComponentActivity() {
authSessionStore.token.collect { token -> authSessionStore.token.collect { token ->
isAuthenticated = !token.isNullOrBlank() isAuthenticated = !token.isNullOrBlank()
if (!token.isNullOrBlank()) {
pttManager.start(token)
} else {
pttManager.stop()
}
if (isAuthenticated && !frontPreviewOpened) { if (isAuthenticated && !frontPreviewOpened) {
frontPreviewOpened = true frontPreviewOpened = true
runOnUiThread { runOnUiThread {
@ -647,3 +651,9 @@ fun AppNavigation(viewModel: MainViewModel, onOpenPhotoCamera: () -> Unit, onOpe
} }
} }
} }

View file

@ -1,4 +1,4 @@
package com.bodycamera.twentyfoulabs.data.backend package com.bodycamera.twentyfoulabs.data.backend
/** Centralized endpoints for the M530 PoC. */ /** Centralized endpoints for the M530 PoC. */
object BackendConfig { object BackendConfig {
@ -10,10 +10,11 @@ object BackendConfig {
const val STREAM_ID = "webrtc_camera_stream" const val STREAM_ID = "webrtc_camera_stream"
// Rails mobile API exposed through ngrok HTTPS. // 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_SERIAL = "BC-M530-001"
const val DEVICE_MODEL = "Recoda M530" const val DEVICE_MODEL = "Recoda M530"
val httpBaseUrl: String get() = "http://$HOST:$HTTP_PORT" val httpBaseUrl: String get() = "http://$HOST:$HTTP_PORT"
val whipUrl: String get() = "http://$HOST:8889/$STREAM_ID/whip" val whipUrl: String get() = "http://$HOST:8889/$STREAM_ID/whip"
} }

View file

@ -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
)

View file

@ -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<Long, ReceiverSession>()
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<out IceCandidate>?
) {
}
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<out MediaStream>?
) {
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<AudioTrack>()
.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
)
}
}

View file

@ -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)
}
}

View file

@ -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<Long, Participant>()
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
}
}

View file

@ -1,6 +1,6 @@
<?xml version="1.0" encoding="utf-8"?> <?xml version="1.0" encoding="utf-8"?>
<network-security-config> <network-security-config>
<!-- Permitir tráfico HTTP para redes locales (LAN) --> <!-- Permitir tráfico HTTP para redes locales (LAN) -->
<domain-config cleartextTrafficPermitted="true"> <domain-config cleartextTrafficPermitted="true">
<!-- IPs locales (192.168.x.x, 10.x.x.x, etc.) --> <!-- IPs locales (192.168.x.x, 10.x.x.x, etc.) -->
<domain includeSubdomains="true">192.168.0.25</domain> <domain includeSubdomains="true">192.168.0.25</domain>
@ -9,6 +9,7 @@
<domain includeSubdomains="true">172.16</domain> <domain includeSubdomains="true">172.16</domain>
<domain includeSubdomains="true">localhost</domain> <domain includeSubdomains="true">localhost</domain>
<domain includeSubdomains="true">127.0.0.1</domain> <domain includeSubdomains="true">127.0.0.1</domain>
<domain includeSubdomains="true">35.224.197.127</domain>
</domain-config> </domain-config>
<!-- Permitir todo en debug para pruebas --> <!-- Permitir todo en debug para pruebas -->
@ -19,3 +20,4 @@
</trust-anchors> </trust-anchors>
</base-config> </base-config>
</network-security-config> </network-security-config>