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.Esp32LinkState import com.hawhamburg.micr0bu.data.transport.Esp32Link import com.hawhamburg.micr0bu.data.transport.Esp32Transport import com.hawhamburg.micr0bu.data.transport.OutgoingMessage import com.hawhamburg.micr0bu.data.transport.StationStatus 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 esp32Link: Esp32Link, private val camPinger: CamPinger, ) : ViewModel() { // ── MQTT connection & messages ──────────────────────────────────────────── val connectionState: StateFlow = repo.connectionState val topicMessages: StateFlow>> = repo.topicMessages private val _selectedTopic = MutableStateFlow(null) val selectedTopic: StateFlow = _selectedTopic.asStateFlow() private val _autoScroll = MutableStateFlow(true) val autoScroll: StateFlow = _autoScroll.asStateFlow() // ── Transport & USB ─────────────────────────────────────────────────────── val activeTransport: StateFlow = repo.activeTransport /** Which physical OBU (Section 13) is currently selected — CiT One or ESP32-C5. */ val obuHardware: StateFlow = 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 = usbDetector.usbNetwork .map { it != null } .stateIn(viewModelScope, SharingStarted.Eagerly, false) /** Auto-detected OBU gateway IP on the USB interface. */ val detectedObuIp: StateFlow = usbDetector.detectedGatewayIp /** ESP32-C5 link state, over USB or BLE per [esp32Transport] — see [Esp32Link]. */ val esp32LinkState: StateFlow = esp32Link.state /** Latest firmware heartbeat + drop counters, null until the first STATUS frame arrives. */ val espLinkStatus: StateFlow = esp32Link.linkStatus /** Non-zero means CAMs are being built and dropped — see [Esp32Link.send]. */ val camSendFailures: StateFlow = esp32Link.consecutiveWriteFailures /** Signing and radio counters of the current obu-firmware; null with the previous firmware. */ val stationStatus: StateFlow = esp32Link.stationStatus /** One line about the link session (pairing passkey, provisioning, refusals); null when quiet. */ val esp32Detail: StateFlow = esp32Link.detail // ── ESP32-C5 settings ───────────────────────────────────────────────────── val esp32Transport: StateFlow = esp32Link.transport fun setEsp32Transport(transport: Esp32Transport) { viewModelScope.launch { obuHardwarePrefs.setEsp32Transport(transport) } } val outgoingMessage: StateFlow = obuHardwarePrefs.outgoingMessageFlow .stateIn(viewModelScope, SharingStarted.Eagerly, OutgoingMessage.CAM) fun setOutgoingMessage(message: OutgoingMessage) { viewModelScope.launch { obuHardwarePrefs.setOutgoingMessage(message) } } val signOutgoing: StateFlow = obuHardwarePrefs.signOutgoingFlow .stateIn(viewModelScope, SharingStarted.Eagerly, true) fun setSignOutgoing(sign: Boolean) { viewModelScope.launch { obuHardwarePrefs.setSignOutgoing(sign) } } // ── 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 = camPinger.isActive val camPingerSentCount: StateFlow = camPinger.sentCount /** False while the pinger runs without a GNSS fix — it has no position to build a CAM from. */ val camPingerHasFix: StateFlow = 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 = 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 = 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(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 = _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 = 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> = 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()) { 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> = combine( camUseCaseRepository.decodedSpat .runningFold(emptyMap()) { 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 = 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 = repo.denmActive /** JSON string of the most recently transmitted DENM control message. */ val lastDenmPayload: StateFlow = repo.lastDenmPayload /** The use case id currently active on the OBU, if any (e.g. for showing in the UI). */ val activeDenmUseCase: StateFlow = 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 = camUseCaseRepository.ownStationId /** Active CAM-based use case alerts, filtered to the use cases enabled in Settings. */ val useCaseAlerts: StateFlow> = camUseCaseRepository.enabledAlerts /** Per-use-case enable/disable state (Settings > Use Case Alerts). */ val useCaseEnabledMap: StateFlow> = camUseCaseRepository.enabledMap /** Ego bike's latest known position, for the V2X Monitor live map view (Section 13). */ val ownCamPosition: StateFlow = camUseCaseRepository.ownPosition /** Latest known CAM per tracked remote road user, for the live map view (Section 13). */ val remoteCamPositions: StateFlow> = 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> = 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 link (USB or BLE per [esp32Transport]) — separate from * [connect]/[disconnect], which drive the CiT One's MQTT-over-USB-C/Wi-Fi path. */ fun connectEsp32() = esp32Link.connect() fun disconnectEsp32() = esp32Link.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 esp32Link.disconnect(): the link 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. } }