Files
MicrOBU/app/src/main/java/com/hawhamburg/micr0bu/viewmodel/MqttViewModel.kt
T
Ashin Walpola 034ef22336 Decode raw v2x/rx on the CiT One path, and stop tracking our own CAM pings
The Use Case app's v2x-uca/output/json topics are a rate-limited and
lossy view: traffic the OBU's radio actually heard, the ESP32's CAM
pinger among it, never reached the app. The raw v2x/rx topics carry
everything, as RecvV2XMessage protobuf with the ITS-G5 PDU in one bytes
field (CI-CiT MQTT API section 2.4).

RecvV2xMessage is a minimal protobuf wire-format reader for the three
fields needed: btpHeader type and destination port, the GeoNetworking
destination-area radius, and the payload. Hand-written for the same
reason the ASN.1 codecs are, rather than adding protoc and the protobuf
Gradle plugin and vendoring a third-party .proto into this repository.
Field numbers are pinned by a byte fixture written out by hand from the
encoding rules, not generated by our own encoder.

Raw payloads now travel as bytes rather than String. The previous UTF-8
round trip replaced every byte that is not valid UTF-8, leaving a
payload that still looked plausible in a log and decoded to nothing.

CAM, DENM and SPATEM from both transports now meet in shared handlers,
so everything downstream is transport-agnostic. SPATEM works on the CiT
One path for the first time, and DENM gains its relevance radius there.
Where both sources describe the same event the decoded one wins: remote
CAMs from the processed topic are suppressed while the raw topic is
live, and DENMs dedup on ETSI's actionID with the decoded list last.
The processed topics stay subscribed as a fallback for an OBU whose
configuration does not publish the raw ones.

Two defects found while testing this:

CamPinger transmits under a fixed bench station id, deliberately
distinct from the persisted one, but the self-heard filter only knew
the persisted id. Every ping therefore came back through the ESP32's
promiscuous receive as a remote road user sitting exactly on top of the
ego position, moving at the ego's own speed and heading, and was handed
to the detection engine as a collision partner for itself. The rule now
lives in OwnStationIds, covers both ids, and has tests, so a third
transmit path cannot reintroduce the same gap quietly.

Self-heard frames are now counted and reported on the pinger card
instead of being discarded. That round trip is the only direct evidence
the serial link, the ESP32's transmit path and its receive path all
work, which is what the bench pinger exists to demonstrate.

Also: the stationType warning banner no longer shows in ESP32-C5 mode.
It reads a value from the CiT One's obu_gnss topic, which that hardware
never publishes, so it stayed on screen reporting on an OBU that was no
longer in use.
2026-09-02 15:25:31 +02:00

400 lines
20 KiB
Kotlin

package com.hawhamburg.micr0bu.viewmodel
import androidx.lifecycle.ViewModel
import androidx.lifecycle.viewModelScope
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
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
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
import dagger.hilt.android.lifecycle.HiltViewModel
import kotlinx.coroutines.delay
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
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,
) : ViewModel() {
// ── MQTT connection & messages ────────────────────────────────────────────
val connectionState: StateFlow<MqttConnectionState> = repo.connectionState
val topicMessages: StateFlow<Map<String, List<MqttMessage>>> = repo.topicMessages
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) }
}
/** 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
/**
* Own transmissions heard back off the air, null until one is.
*
* This is the pinger's actual proof of life. [camPingerSentCount] only says frames were
* handed to the ESP32; this says they went out and came back, which is the round trip the
* bench test is there to demonstrate. See
* [com.hawhamburg.micr0bu.domain.cam.OwnTxLoopback].
*/
val ownTxLoopback: StateFlow<com.hawhamburg.micr0bu.domain.cam.OwnTxLoopback?> =
camUseCaseRepository.ownTxLoopback
fun startCamPinger() {
// Reset first, so the tally counts this run rather than accumulating across runs and
// making the comparison against sent count meaningless.
camUseCaseRepository.resetOwnTxLoopback()
camPinger.start()
}
fun stopCamPinger() = camPinger.stop()
// ── 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.
*/
val obuStationType: StateFlow<Int?> = _obuStationType.asStateFlow()
/**
* True when the CiT One 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.
*
* Suppressed in ESP32-C5 mode. The value behind it comes from the CiT One's
* `v2x/rx/obu_gnss` topic, which the ESP32-C5 does not publish, so a warning raised before a
* mode switch would otherwise stay on screen reporting on an OBU that is no longer in use.
* There is nothing for it to warn about on that path either: the phone builds its own CAM
* ([com.hawhamburg.micr0bu.domain.cam.PhoneCamBuilder]), which sets stationType to cyclist
* locally rather than reading it back from an OBU.
*
* The underlying [obuStationType] is deliberately not cleared on the switch. It remains the
* last thing that OBU actually said, and obu_gnss refreshes it at ~4 Hz on returning to the
* CiT One path, so the warning re-evaluates against fresh data within a fraction of a second.
*/
val obuStationTypeWarning: StateFlow<Boolean> = combine(
_obuStationType,
repo.obuHardware,
) { stationType, hardware ->
hardware == ObuHardware.CIT_ONE && stationType != null && stationType != 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 Use Case app's `v2x-uca/output/json/denm` MQTT topic
* (parsed by [DenmParser]), and UPER decoded by
* [com.hawhamburg.micr0bu.domain.asn1.DenmUperCodec] from whichever raw path is live, the
* ESP32-C5 serial link or the CiT One's `v2x/rx/denm` protobuf topic.
*
* Where both describe the same hazard, the decoded one wins. Both key on ETSI's actionID, so
* the `associateBy` below collapses them to one entry, and the decoded list is concatenated
* second so it is the one that survives. That is the intended preference: the Use Case app
* rate-limits and drops messages, and reduces what it does publish to the fields it cared
* about, so it can only ever be a lossier account of the same event.
*
* 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.decodedDenm
.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),
) { fromUseCaseApp, fromDecoder, _ ->
val now = System.currentTimeMillis()
(fromUseCaseApp + fromDecoder)
.filterNot { it.isTermination } // the hazard is over - stop drawing it
.associateBy { it.dedupKey } // last write wins, so the decoded one is kept
.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].
*
* Both hardware paths: SPATEM arrives over the air on BTP port 2004 via the ESP32-C5 serial
* link, or on the CiT One's `v2x/rx/spatem` protobuf topic. The CiT One's processed
* `v2x-uca/output/json/spat` topic is not used, since the raw topic carries every repetition.
*
* 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.decodedSpat
.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
}
// ── DENM transmission ─────────────────────────────────────────────────────
/** True while a DENM use case is actively broadcasting on the OBU. */
val denmActive: StateFlow<Boolean> = repo.denmActive
/** JSON string of the most recently transmitted DENM control message. */
val lastDenmPayload: StateFlow<String?> = repo.lastDenmPayload
/** 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)
}
// ── 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 ->
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()
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).
*/
fun sendDenm(useCase: String = DenmUseCase.STATIONARY.id) {
repo.activateDenm(useCase)
}
/**
* Publish a uca-denmctrl deactivate message (retained) for whichever use case is
* currently active.
*/
fun stopDenm() {
repo.deactivateDenm()
}
// ── 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.
}
}