feat: CAM-based use case detection engine, decouple DENM from sensor triggers
Added: - UseCaseDetectionEngine: IMA-B, IMA-S, RTW-B, LTW-B, SMVA/BCW-B via per-station CAM history, yaw-rate/heading-trend turn detection, Info/Awareness/Warning alerts - Ego state from obu_gnss (~4Hz, yaw rate) with phone GPS fallback when stale - Human-readable alert narratives + technical detail in V2X Monitor - Settings > Use Case Alerts per-use-case toggles Fixed: - DENM no longer auto-fires on braking/stopping; manual-only (antenna/RSU test) - Removed unwired CamBuilder.kt (dead code, contradicted CAM-only architecture) - TripRecordingService no longer depends on MQTT/OBU connectivity
This commit is contained in:
@@ -2,7 +2,7 @@ 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.cam.CamUseCaseRepository
|
||||
import com.hawhamburg.micr0bu.data.mqtt.MqttConnectionState
|
||||
import com.hawhamburg.micr0bu.data.mqtt.MqttMessage
|
||||
import com.hawhamburg.micr0bu.data.mqtt.MqttPreferences
|
||||
@@ -10,6 +10,9 @@ 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 com.hawhamburg.micr0bu.domain.denm.DenmUseCase
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseAlert
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseType
|
||||
import dagger.hilt.android.lifecycle.HiltViewModel
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.SharingStarted
|
||||
@@ -17,29 +20,23 @@ 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,
|
||||
private val camUseCaseRepository: CamUseCaseRepository,
|
||||
) : 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()
|
||||
val topicMessages: StateFlow<Map<String, List<MqttMessage>>> = repo.topicMessages
|
||||
|
||||
private val _selectedTopic = MutableStateFlow<String?>(null)
|
||||
val selectedTopic: StateFlow<String?> = _selectedTopic.asStateFlow()
|
||||
@@ -69,11 +66,16 @@ class MqttViewModel @Inject constructor(
|
||||
|
||||
// ── 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.
|
||||
*/
|
||||
private val _obuStationType = MutableStateFlow<Int?>(null)
|
||||
|
||||
/**
|
||||
* The raw stationType value last reported by the OBU. Null until the first
|
||||
* obu_gnss message arrives. Exposed so the UI can show the actual value.
|
||||
*/
|
||||
val obuStationType: StateFlow<Int?> = _obuStationType.asStateFlow()
|
||||
|
||||
/**
|
||||
@@ -87,30 +89,42 @@ class MqttViewModel @Inject constructor(
|
||||
|
||||
// ── 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()
|
||||
/** True while a DENM use case is actively broadcasting on the OBU. */
|
||||
val denmActive: StateFlow<Boolean> = repo.denmActive
|
||||
|
||||
private val _lastDenmPayload = MutableStateFlow<String?>(null)
|
||||
/** JSON string of the most recently transmitted DENM control message. */
|
||||
val lastDenmPayload: StateFlow<String?> = _lastDenmPayload.asStateFlow()
|
||||
val lastDenmPayload: StateFlow<String?> = repo.lastDenmPayload
|
||||
|
||||
/** Tracks which use case is currently active so Stop uses the same usecase field. */
|
||||
private var activeDenmUseCase: String? = null
|
||||
/** The use case id currently active on the OBU, if any (e.g. for showing in the UI). */
|
||||
val activeDenmUseCase: StateFlow<String?> = repo.activeDenmUseCase
|
||||
|
||||
// ── CAM-based Use Case Detection (Section 10.2 / 10.4) ────────────────────
|
||||
// Entirely separate from the DENM transmission above: this consumes CAM only, raises
|
||||
// local HMI alerts, and never triggers an outbound V2X message.
|
||||
|
||||
/** The ego OBU's own station ID, learned from v2x/rx/obu_gnss's own_info. */
|
||||
val ownStationId: StateFlow<Long?> = camUseCaseRepository.ownStationId
|
||||
|
||||
/** Active CAM-based use case alerts, filtered to the use cases enabled in Settings. */
|
||||
val useCaseAlerts: StateFlow<List<UseCaseAlert>> = camUseCaseRepository.enabledAlerts
|
||||
|
||||
/** Per-use-case enable/disable state (Settings > Use Case Alerts). */
|
||||
val useCaseEnabledMap: StateFlow<Map<UseCaseType, Boolean>> = camUseCaseRepository.enabledMap
|
||||
|
||||
/** True if [stationId] is the ego OBU's own — used for OWN/REMOTE badges in the raw message list. */
|
||||
fun isOwnStationId(stationId: Long): Boolean = camUseCaseRepository.isOwnStationId(stationId)
|
||||
|
||||
fun setUseCaseEnabled(type: UseCaseType, enabled: Boolean) {
|
||||
camUseCaseRepository.setUseCaseEnabled(type, enabled)
|
||||
}
|
||||
|
||||
// ── Init ──────────────────────────────────────────────────────────────────
|
||||
|
||||
init {
|
||||
viewModelScope.launch {
|
||||
// Per-topic message lists now live in MqttRepository (survives screen close);
|
||||
// here we just watch for own_info to track the OBU's reported stationType.
|
||||
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)
|
||||
@@ -139,52 +153,19 @@ class MqttViewModel @Inject constructor(
|
||||
|
||||
/**
|
||||
* 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.
|
||||
* receives it on any (re)connect. Only one use case may be active at a time
|
||||
* (no-op if another use case, manual or automatic, is already active).
|
||||
*/
|
||||
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)
|
||||
fun sendDenm(useCase: String = DenmUseCase.STATIONARY.id) {
|
||||
repo.activateDenm(useCase)
|
||||
}
|
||||
|
||||
/**
|
||||
* Publish a uca-denmctrl deactivate message (retained).
|
||||
* Uses the same usecase string as the activate call so the OBU phases out the correct DENM.
|
||||
* Publish a uca-denmctrl deactivate message (retained) for whichever use case is
|
||||
* currently active.
|
||||
*/
|
||||
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)
|
||||
}
|
||||
repo.deactivateDenm()
|
||||
}
|
||||
|
||||
// ── Lifecycle ─────────────────────────────────────────────────────────────
|
||||
|
||||
Reference in New Issue
Block a user