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.EspRxMode 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.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.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.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 = 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) } } /** * ESP32-C5-only: whether received CAM traffic is processed or discarded — see [EspRxMode]'s * KDoc for the important caveat that this doesn't actually disable the ESP32's receiver * (it can't, without also breaking TX). */ val espRxMode: StateFlow = obuHardwarePrefs.espRxModeFlow.stateIn( viewModelScope, SharingStarted.Eagerly, EspRxMode.SEND_AND_RECEIVE, ) fun setEspRxMode(mode: EspRxMode) { viewModelScope.launch { obuHardwarePrefs.setEspRxMode(mode) } } /** 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 USB-serial link state (Phase 03) — see [UsbSerialTransport]. */ val usbSerialState: StateFlow = usbSerialTransport.state /** Latest firmware heartbeat + drop counters, null until the first STATUS frame arrives. */ val espLinkStatus: StateFlow = usbSerialTransport.linkStatus /** Non-zero means CAMs are being built and dropped — see [UsbSerialTransport.sendCamTx]. */ val camSendFailures: StateFlow = 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 = 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 fun startCamPinger() = 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 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 = _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. * * Derived from the raw `v2x-uca/output/json/denm` messages the repository already buffers, * rather than a second subscription — the repository caps each topic's history, so this is * bounded by construction. * * Always empty on the ESP32-C5 path: that firmware forwards BTP-B port 2001 (CAM) only and * drops DENM before it reaches the phone. See [DenmEvent]'s KDoc. */ val denmEvents: StateFlow> = repo.topicMessages .map { byTopic -> (byTopic[DENM_RX_TOPIC] ?: emptyList()) .mapNotNull { DenmParser.parse(it.payload, it.timestamp) } .associateBy { it.dedupKey } // last write wins = most recent per hazard .values .sortedByDescending { it.timestamp } } .stateIn(viewModelScope, SharingStarted.Eagerly, emptyList()) private companion object { /** Use Case API topic carrying received DENMs (CiT One path only). */ const val DENM_RX_TOPIC = "v2x-uca/output/json/denm" } // ── 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 /** 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. } }