Files
MicrOBU/app/src/main/java/com/hawhamburg/micr0bu/viewmodel/MqttViewModel.kt
T

363 lines
18 KiB
Kotlin
Raw Normal View History

2026-06-03 14:25:20 +02:00
package com.hawhamburg.micr0bu.viewmodel
import androidx.lifecycle.ViewModel
import androidx.lifecycle.viewModelScope
import com.hawhamburg.micr0bu.data.cam.CamUseCaseRepository
2026-06-03 14:25:20 +02:00
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.mqtt.ObuHardwarePreferences
import com.hawhamburg.micr0bu.data.transport.EspLinkStatus
import com.hawhamburg.micr0bu.data.transport.ObuHardware
2026-06-03 14:25:20 +02:00
import com.hawhamburg.micr0bu.data.transport.TransportType
import com.hawhamburg.micr0bu.data.transport.UsbNetworkDetector
import com.hawhamburg.micr0bu.data.transport.UsbSerialState
import com.hawhamburg.micr0bu.data.transport.UsbSerialTransport
import com.hawhamburg.micr0bu.domain.denm.DenmEvent
import com.hawhamburg.micr0bu.domain.denm.DenmParser
import com.hawhamburg.micr0bu.domain.spat.SpatIntersection
import com.hawhamburg.micr0bu.domain.denm.DenmUseCase
import com.hawhamburg.micr0bu.domain.usecase.UseCaseAlert
import com.hawhamburg.micr0bu.domain.usecase.UseCaseType
import com.hawhamburg.micr0bu.service.CamPinger
2026-06-03 14:25:20 +02:00
import dagger.hilt.android.lifecycle.HiltViewModel
import kotlinx.coroutines.delay
2026-06-03 14:25:20 +02:00
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.combine
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.runningFold
2026-06-03 14:25:20 +02:00
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.launch
import org.json.JSONObject
import javax.inject.Inject
@HiltViewModel
class MqttViewModel @Inject constructor(
private val repo: MqttRepository,
private val prefs: MqttPreferences,
private val usbDetector: UsbNetworkDetector,
private val camUseCaseRepository: CamUseCaseRepository,
private val obuHardwarePrefs: ObuHardwarePreferences,
private val usbSerialTransport: UsbSerialTransport,
private val camPinger: CamPinger,
2026-06-03 14:25:20 +02:00
) : ViewModel() {
// ── MQTT connection & messages ────────────────────────────────────────────
val connectionState: StateFlow<MqttConnectionState> = repo.connectionState
val topicMessages: StateFlow<Map<String, List<MqttMessage>>> = repo.topicMessages
2026-06-03 14:25:20 +02:00
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
/** Which physical OBU (Section 13) is currently selected — CiT One or ESP32-C5. */
val obuHardware: StateFlow<ObuHardware> = repo.obuHardware
fun setObuHardware(hardware: ObuHardware) {
viewModelScope.launch { obuHardwarePrefs.setObuHardware(hardware) }
}
2026-06-03 14:25:20 +02:00
/** 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
/** ESP32-C5 USB-serial link state (Phase 03) — see [UsbSerialTransport]. */
val usbSerialState: StateFlow<UsbSerialState> = usbSerialTransport.state
/** Latest firmware heartbeat + drop counters, null until the first STATUS frame arrives. */
val espLinkStatus: StateFlow<EspLinkStatus?> = usbSerialTransport.linkStatus
/** Non-zero means CAMs are being built and dropped — see [UsbSerialTransport.sendCamTx]. */
val camSendFailures: StateFlow<Int> = usbSerialTransport.consecutiveWriteFailures
// ── ESP32-C5 CAM pinger (manual bench test, Phase 03) ─────────────────────
// The ESP32-C5-path equivalent of the CiT One's manual DENM trigger below — a fixed-
// location 1 Hz CAM ping the user starts/stops from the V2X Monitor screen to verify the
// serial link + ESP32 TX/RX radio path independent of GNSS movement or trip recording.
// See CamPinger's KDoc.
val camPingerActive: StateFlow<Boolean> = camPinger.isActive
val camPingerSentCount: StateFlow<Int> = camPinger.sentCount
/** False while the pinger runs without a GNSS fix — it has no position to build a CAM from. */
val camPingerHasFix: StateFlow<Boolean> = camPinger.hasFix
fun startCamPinger() = camPinger.start()
fun stopCamPinger() = camPinger.stop()
2026-06-03 14:25:20 +02:00
// ── Prefs ─────────────────────────────────────────────────────────────────
val mqttPrefs: StateFlow<MqttPrefs> = prefs.prefsFlow.stateIn(
viewModelScope,
SharingStarted.Eagerly,
MqttPrefs(),
)
// ── OBU identity (parsed from v2x/rx/obu_gnss own_info) ──────────────────
/**
* 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.
*/
2026-06-03 14:25:20 +02:00
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 reception (live map hazard pins) ─────────────────────────────────
/**
* Hazards received from other stations, newest first, deduped by [DenmEvent.dedupKey] so a
* repeating DENM about the same hazard stays one pin instead of stacking up.
*
* Two sources, merged: the CiT One path's `v2x-uca/output/json/denm` MQTT topic (parsed by
* [DenmParser]), and the ESP32-C5 path's over-the-air DENMs (GeoBroadcast, BTP port 2002,
* decoded by [com.hawhamburg.micr0bu.domain.asn1.DenmUperCodec]). Only one is ever active at a
* time since the hardware selection decides the transport, so merging costs nothing and keeps
* the UI transport-agnostic.
*
* Events carrying `termination` are filtered out rather than shown — the hazard is over.
*/
val denmEvents: StateFlow<List<DenmEvent>> = combine(
repo.topicMessages.map { byTopic ->
(byTopic[DENM_RX_TOPIC] ?: emptyList())
.mapNotNull { DenmParser.parse(it.payload, it.timestamp) }
},
// Air DENMs accumulate here rather than being a snapshot: the serial path delivers one
// event at a time, so runningFold keeps the set of hazards heard so far.
camUseCaseRepository.airDenm
.runningFold(emptyMap<String, DenmEvent>()) { acc, denm -> acc + (denm.dedupKey to denm) }
.map { it.values.toList() },
// Expiry has to be driven by a clock, not by arrivals. Both upstream flows only re-emit
// when a DENM arrives, so a sender that simply stops transmitting - drives away, loses
// power, leaves range - would otherwise leave its hazard on the map forever: there is no
// further emission to recompute the list. This tick is what makes a hazard fade.
tickerFlow(DENM_EXPIRY_TICK_MS),
) { fromMqtt, fromAir, _ ->
val now = System.currentTimeMillis()
(fromMqtt + fromAir)
.filterNot { it.isTermination } // the hazard is over - stop drawing it
.associateBy { it.dedupKey } // last write wins = most recent per hazard
.values
// Not heard from in DENM_TTL_MS: treat as gone. DENMs repeat at roughly 1 Hz, so a
// full minute of silence is ~60 missed repetitions - well past "we briefly lost one".
.filter { now - it.timestamp <= DENM_TTL_MS }
.sortedByDescending { it.timestamp }
}.stateIn(viewModelScope, SharingStarted.Eagerly, emptyList())
/**
* Live signal state per intersection, newest first, keyed by [IntersectionSignalState.key].
*
* ESP32-C5 path only: SPATEM arrives over the air on BTP port 2004. The CiT One path publishes
* SPATEM on its own MQTT topic in a different (protobuf-wrapped) shape, which is not wired up.
*
* One entry per intersection, not per message: SPATEM repeats at ~2 Hz per RSU, so a log would
* grow without telling anyone anything. Entries expire like DENMs do - an intersection left
* behind stops transmitting, and the same clock-driven argument applies.
*/
val spatIntersections: StateFlow<List<SpatIntersection>> = combine(
camUseCaseRepository.airSpat
.runningFold(emptyMap<String, SpatIntersection>()) { acc, spat ->
acc + spat.intersections.associate { i ->
i.key to SpatIntersection(i, spat.stationId, spat.rssiDbm, spat.timestamp)
}
},
tickerFlow(SPAT_EXPIRY_TICK_MS),
) { byKey, _ ->
val now = System.currentTimeMillis()
byKey.values
.filter { now - it.timestamp <= SPAT_TTL_MS }
.sortedByDescending { it.timestamp }
}.stateIn(viewModelScope, SharingStarted.Eagerly, emptyList())
/** Emits immediately, then every [periodMs], purely to re-trigger a time-dependent combine. */
private fun tickerFlow(periodMs: Long): Flow<Long> = flow {
while (true) {
emit(System.currentTimeMillis())
delay(periodMs)
}
}
private companion object {
/** Use Case API topic carrying received DENMs (CiT One path only). */
const val DENM_RX_TOPIC = "v2x-uca/output/json/denm"
/**
* How long a hazard stays listed after its last repetition. A DENM has no "still here"
* guarantee beyond the sender repeating it, and its own validityDuration is not decoded
* yet, so silence is the only expiry signal available.
*/
const val DENM_TTL_MS = 60_000L
/** How often the list is re-evaluated for expiry. Sets the worst-case lateness of a fade. */
const val DENM_EXPIRY_TICK_MS = 5_000L
/**
* SPATEM repeats at ~2 Hz, so 15 s of silence is ~30 missed repetitions: the RSU is out of
* range. Much shorter than the DENM window because a stale traffic light is more
* misleading than a stale hazard - a light that stopped updating is not "still green".
*/
const val SPAT_TTL_MS = 15_000L
const val SPAT_EXPIRY_TICK_MS = 2_000L
/** RSU CAMs arrive at ~2 Hz, same as any other station, so the same window applies. */
const val RSU_TTL_MS = 15_000L
const val RSU_EXPIRY_TICK_MS = 2_000L
}
2026-06-03 14:25:20 +02:00
// ── DENM transmission ─────────────────────────────────────────────────────
/** True while a DENM use case is actively broadcasting on the OBU. */
val denmActive: StateFlow<Boolean> = repo.denmActive
2026-06-03 14:25:20 +02:00
/** JSON string of the most recently transmitted DENM control message. */
val lastDenmPayload: StateFlow<String?> = repo.lastDenmPayload
2026-06-03 14:25:20 +02:00
/** 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
/** Ego bike's latest known position, for the V2X Monitor live map view (Section 13). */
val ownCamPosition: StateFlow<com.hawhamburg.micr0bu.domain.cam.Cam?> = camUseCaseRepository.ownPosition
/** Latest known CAM per tracked remote road user, for the live map view (Section 13). */
val remoteCamPositions: StateFlow<Map<Long, com.hawhamburg.micr0bu.domain.cam.Cam>> = camUseCaseRepository.remotePositions
/**
* Every station to draw: road users from the detection engine, plus roadside units, which are
* tracked outside it (see [com.hawhamburg.micr0bu.data.cam.CamUseCaseRepository.rsuStations]).
*
* The engine prunes its own stale entries; nothing prunes the RSU map, so the staleness window
* is applied here. As with hazards and signals, expiry has to be clock-driven - an RSU that
* goes out of range simply stops transmitting, and no further emission would arrive to
* recompute the list.
*/
val stationsInRange: StateFlow<Map<Long, com.hawhamburg.micr0bu.domain.cam.Cam>> = combine(
camUseCaseRepository.remotePositions,
camUseCaseRepository.rsuStations,
tickerFlow(RSU_EXPIRY_TICK_MS),
) { roadUsers, rsus, _ ->
val now = System.currentTimeMillis()
roadUsers + rsus.filterValues { now - it.timestamp <= RSU_TTL_MS }
}.stateIn(viewModelScope, SharingStarted.Eagerly, emptyMap())
/** 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)
}
2026-06-03 14:25:20 +02:00
// ── 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.
2026-06-03 14:25:20 +02:00
repo.messages.collect { msg ->
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()
/** Connect/disconnect the ESP32-C5 USB-serial link — separate from [connect]/[disconnect],
* which drive the CiT One's MQTT-over-USB-C/Wi-Fi path. See [ConnectionSetupScreen]. */
fun connectUsbSerial() = usbSerialTransport.connect()
fun disconnectUsbSerial() = usbSerialTransport.disconnect()
2026-06-03 14:25:20 +02:00
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
* (no-op if another use case, manual or automatic, is already active).
2026-06-03 14:25:20 +02:00
*/
fun sendDenm(useCase: String = DenmUseCase.STATIONARY.id) {
repo.activateDenm(useCase)
2026-06-03 14:25:20 +02:00
}
/**
* Publish a uca-denmctrl deactivate message (retained) for whichever use case is
* currently active.
2026-06-03 14:25:20 +02:00
*/
fun stopDenm() {
repo.deactivateDenm()
2026-06-03 14:25:20 +02:00
}
// ── Lifecycle ─────────────────────────────────────────────────────────────
override fun onCleared() {
super.onCleared()
repo.disconnect()
camPinger.stop()
// Deliberately NOT usbSerialTransport.disconnect(): the transport is an app-scoped
// @Singleton also held by the foreground TripRecordingService (via CamTransmitLoop).
// Closing it here would tear the port down when the Activity goes away — e.g. swiping
// the app from Recents mid-recording — leaving the still-running service beaconing into
// a dead port. The port closes on explicit user Disconnect, on USB detach (handled
// inside the transport), or with the process. See UsbSerialTransport's "Ownership" KDoc.
2026-06-03 14:25:20 +02:00
}
}