197 lines
8.7 KiB
Kotlin
197 lines
8.7 KiB
Kotlin
package com.hawhamburg.micr0bu.viewmodel
|
|||
|
|
|
||
|
|
import androidx.lifecycle.ViewModel
|
||
|
|
import androidx.lifecycle.viewModelScope
|
||
|
|
import com.hawhamburg.micr0bu.data.mqtt.MessageDirection
|
||
|
|
import com.hawhamburg.micr0bu.data.mqtt.MqttConnectionState
|
||
|
|
import com.hawhamburg.micr0bu.data.mqtt.MqttMessage
|
||
|
|
import com.hawhamburg.micr0bu.data.mqtt.MqttPreferences
|
||
|
|
import com.hawhamburg.micr0bu.data.mqtt.MqttPrefs
|
||
|
|
import com.hawhamburg.micr0bu.data.mqtt.MqttRepository
|
||
|
|
import com.hawhamburg.micr0bu.data.transport.TransportType
|
||
|
|
import com.hawhamburg.micr0bu.data.transport.UsbNetworkDetector
|
||
|
|
import dagger.hilt.android.lifecycle.HiltViewModel
|
||
|
|
import kotlinx.coroutines.flow.MutableStateFlow
|
||
|
|
import kotlinx.coroutines.flow.SharingStarted
|
||
|
|
import kotlinx.coroutines.flow.StateFlow
|
||
|
|
import kotlinx.coroutines.flow.asStateFlow
|
||
|
|
import kotlinx.coroutines.flow.map
|
||
|
|
import kotlinx.coroutines.flow.stateIn
|
||
|
|
import kotlinx.coroutines.flow.update
|
||
|
|
import kotlinx.coroutines.launch
|
||
|
|
import org.json.JSONObject
|
||
|
|
import javax.inject.Inject
|
||
|
|
|
||
|
|
private const val MAX_MESSAGES_PER_TOPIC = 50
|
||
|
|
|
||
|
|
private const val DENM_CTRL_TOPIC = "v2x-uca/input/denmtrg"
|
||
|
|
private const val DENM_USE_CASE = "hln-sv"
|
||
|
|
|
||
|
|
@HiltViewModel
|
||
|
|
class MqttViewModel @Inject constructor(
|
||
|
|
private val repo: MqttRepository,
|
||
|
|
private val prefs: MqttPreferences,
|
||
|
|
private val usbDetector: UsbNetworkDetector,
|
||
|
|
) : ViewModel() {
|
||
|
|
|
||
|
|
// ── MQTT connection & messages ────────────────────────────────────────────
|
||
|
|
|
||
|
|
val connectionState: StateFlow<MqttConnectionState> = repo.connectionState
|
||
|
|
|
||
|
|
private val _topicMessages = MutableStateFlow<Map<String, List<MqttMessage>>>(emptyMap())
|
||
|
|
val topicMessages: StateFlow<Map<String, List<MqttMessage>>> = _topicMessages.asStateFlow()
|
||
|
|
|
||
|
|
private val _selectedTopic = MutableStateFlow<String?>(null)
|
||
|
|
val selectedTopic: StateFlow<String?> = _selectedTopic.asStateFlow()
|
||
|
|
|
||
|
|
private val _autoScroll = MutableStateFlow(true)
|
||
|
|
val autoScroll: StateFlow<Boolean> = _autoScroll.asStateFlow()
|
||
|
|
|
||
|
|
// ── Transport & USB ───────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
val activeTransport: StateFlow<TransportType> = repo.activeTransport
|
||
|
|
|
||
|
|
/** True when a 192.168.42.x USB-C tethering network is detected. */
|
||
|
|
val usbConnected: StateFlow<Boolean> = usbDetector.usbNetwork
|
||
|
|
.map { it != null }
|
||
|
|
.stateIn(viewModelScope, SharingStarted.Eagerly, false)
|
||
|
|
|
||
|
|
/** Auto-detected OBU gateway IP on the USB interface. */
|
||
|
|
val detectedObuIp: StateFlow<String?> = usbDetector.detectedGatewayIp
|
||
|
|
|
||
|
|
// ── Prefs ─────────────────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
val mqttPrefs: StateFlow<MqttPrefs> = prefs.prefsFlow.stateIn(
|
||
|
|
viewModelScope,
|
||
|
|
SharingStarted.Eagerly,
|
||
|
|
MqttPrefs(),
|
||
|
|
)
|
||
|
|
|
||
|
|
// ── OBU identity (parsed from v2x/rx/obu_gnss own_info) ──────────────────
|
||
|
|
|
||
|
|
private val _obuStationType = MutableStateFlow<Int?>(null)
|
||
|
|
/**
|
||
|
|
* stationType from the OBU's own_info (v2x/rx/obu_gnss). Should be 2 (cyclist/VRU) per
|
||
|
|
* ETSI EN 302 637-2 Table 1. Null until the first obu_gnss message arrives.
|
||
|
|
*/
|
||
|
|
val obuStationType: StateFlow<Int?> = _obuStationType.asStateFlow()
|
||
|
|
|
||
|
|
/**
|
||
|
|
* True when the OBU has reported a stationType other than 2 (cyclist).
|
||
|
|
* Triggers a persistent warning banner — an incorrect stationType means this OBU will
|
||
|
|
* not be detected as a VRU at equipped intersections.
|
||
|
|
*/
|
||
|
|
val obuStationTypeWarning: StateFlow<Boolean> = _obuStationType
|
||
|
|
.map { it != null && it != 2 }
|
||
|
|
.stateIn(viewModelScope, SharingStarted.Eagerly, false)
|
||
|
|
|
||
|
|
// ── DENM transmission ─────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
private val _denmActive = MutableStateFlow(false)
|
||
|
|
/** True while the DENM use case is actively broadcasting on the OBU. */
|
||
|
|
val denmActive: StateFlow<Boolean> = _denmActive.asStateFlow()
|
||
|
|
|
||
|
|
private val _lastDenmPayload = MutableStateFlow<String?>(null)
|
||
|
|
/** JSON string of the most recently transmitted DENM control message. */
|
||
|
|
val lastDenmPayload: StateFlow<String?> = _lastDenmPayload.asStateFlow()
|
||
|
|
|
||
|
|
/** Tracks which use case is currently active so Stop uses the same usecase field. */
|
||
|
|
private var activeDenmUseCase: String? = null
|
||
|
|
|
||
|
|
// ── Init ──────────────────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
init {
|
||
|
|
viewModelScope.launch {
|
||
|
|
repo.messages.collect { msg ->
|
||
|
|
// Route into per-topic message lists
|
||
|
|
_topicMessages.update { current ->
|
||
|
|
val updated = ((current[msg.topic] ?: emptyList()) + msg)
|
||
|
|
.takeLast(MAX_MESSAGES_PER_TOPIC)
|
||
|
|
current + (msg.topic to updated)
|
||
|
|
}
|
||
|
|
|
||
|
|
// Check stationType from own_info — must be 2 (cyclist/VRU)
|
||
|
|
if (msg.topic == "v2x/rx/obu_gnss") {
|
||
|
|
runCatching {
|
||
|
|
val stType = JSONObject(msg.payload)
|
||
|
|
.optJSONObject("own_info")
|
||
|
|
?.optInt("stationType", -1) ?: -1
|
||
|
|
if (stType >= 0) _obuStationType.value = stType
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// ── Actions ───────────────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
fun connect() = repo.connect()
|
||
|
|
fun disconnect() = repo.disconnect()
|
||
|
|
|
||
|
|
fun selectTopic(topic: String?) { _selectedTopic.value = topic }
|
||
|
|
fun setAutoScroll(enabled: Boolean) { _autoScroll.value = enabled }
|
||
|
|
|
||
|
|
fun updatePrefs(newPrefs: MqttPrefs) {
|
||
|
|
viewModelScope.launch { prefs.update(newPrefs) }
|
||
|
|
}
|
||
|
|
|
||
|
|
// ── DENM actions ──────────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
/**
|
||
|
|
* Publish a uca-denmctrl activate message with the retain flag so the OBU's Use Case app
|
||
|
|
* receives it on any (re)connect. Only one use case may be active at a time.
|
||
|
|
*/
|
||
|
|
fun sendDenm(useCase: String = DENM_USE_CASE) {
|
||
|
|
if (_denmActive.value) return // prevent concurrent activations
|
||
|
|
val payload = buildDenmPayload(useCase = useCase, active = true)
|
||
|
|
_lastDenmPayload.value = payload
|
||
|
|
_denmActive.value = true
|
||
|
|
activeDenmUseCase = useCase
|
||
|
|
emitTxMessage(DENM_CTRL_TOPIC, payload)
|
||
|
|
repo.publishRetained(DENM_CTRL_TOPIC, payload)
|
||
|
|
}
|
||
|
|
|
||
|
|
/**
|
||
|
|
* Publish a uca-denmctrl deactivate message (retained).
|
||
|
|
* Uses the same usecase string as the activate call so the OBU phases out the correct DENM.
|
||
|
|
*/
|
||
|
|
fun stopDenm() {
|
||
|
|
if (!_denmActive.value) return
|
||
|
|
val useCase = activeDenmUseCase ?: DENM_USE_CASE
|
||
|
|
val payload = buildDenmPayload(useCase = useCase, active = false)
|
||
|
|
_lastDenmPayload.value = payload
|
||
|
|
_denmActive.value = false
|
||
|
|
activeDenmUseCase = null
|
||
|
|
emitTxMessage(DENM_CTRL_TOPIC, payload)
|
||
|
|
repo.publishRetained(DENM_CTRL_TOPIC, payload)
|
||
|
|
}
|
||
|
|
|
||
|
|
private fun buildDenmPayload(useCase: String, active: Boolean): String =
|
||
|
|
"""{"type":"uca-denmctrl","active":$active,"usecase":"$useCase","params":{}}"""
|
||
|
|
|
||
|
|
/**
|
||
|
|
* Inject a synthetic TX [MqttMessage] into the local topic map so outbound messages
|
||
|
|
* appear in the V2X Monitor alongside received traffic.
|
||
|
|
*/
|
||
|
|
private fun emitTxMessage(topic: String, payload: String) {
|
||
|
|
val msg = MqttMessage(
|
||
|
|
topic = topic,
|
||
|
|
payload = payload,
|
||
|
|
timestamp = System.currentTimeMillis(),
|
||
|
|
direction = MessageDirection.TX,
|
||
|
|
)
|
||
|
|
_topicMessages.update { current ->
|
||
|
|
val updated = ((current[topic] ?: emptyList()) + msg)
|
||
|
|
.takeLast(MAX_MESSAGES_PER_TOPIC)
|
||
|
|
current + (topic to updated)
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
// ── Lifecycle ─────────────────────────────────────────────────────────────
|
||
|
|
|
||
|
|
override fun onCleared() {
|
||
|
|
super.onCleared()
|
||
|
|
repo.disconnect()
|
||
|
|
}
|
||
|
|
}
|