Initial Commit

This commit is contained in:
Ashin Walpola
2026-06-03 14:51:31 +02:00
parent ee122765e1
commit ae17d4fcba
74 changed files with 7126 additions and 0 deletions
@@ -0,0 +1,196 @@
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.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.transport.TransportType
import com.hawhamburg.micr0bu.data.transport.UsbNetworkDetector
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.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,
) : 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()
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
/** 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
// ── Prefs ─────────────────────────────────────────────────────────────────
val mqttPrefs: StateFlow<MqttPrefs> = prefs.prefsFlow.stateIn(
viewModelScope,
SharingStarted.Eagerly,
MqttPrefs(),
)
// ── 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.
*/
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 transmission ─────────────────────────────────────────────────────
private val _denmActive = MutableStateFlow(false)
/** True while the DENM use case is actively broadcasting on the OBU. */
val denmActive: StateFlow<Boolean> = _denmActive.asStateFlow()
private val _lastDenmPayload = MutableStateFlow<String?>(null)
/** JSON string of the most recently transmitted DENM control message. */
val lastDenmPayload: StateFlow<String?> = _lastDenmPayload.asStateFlow()
/** Tracks which use case is currently active so Stop uses the same usecase field. */
private var activeDenmUseCase: String? = null
// ── Init ──────────────────────────────────────────────────────────────────
init {
viewModelScope.launch {
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)
.optJSONObject("own_info")
?.optInt("stationType", -1) ?: -1
if (stType >= 0) _obuStationType.value = stType
}
}
}
}
}
// ── Actions ───────────────────────────────────────────────────────────────
fun connect() = repo.connect()
fun disconnect() = repo.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.
*/
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)
}
/**
* Publish a uca-denmctrl deactivate message (retained).
* Uses the same usecase string as the activate call so the OBU phases out the correct DENM.
*/
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)
}
}
// ── Lifecycle ─────────────────────────────────────────────────────────────
override fun onCleared() {
super.onCleared()
repo.disconnect()
}
}
@@ -0,0 +1,293 @@
package com.hawhamburg.micr0bu.viewmodel
import android.app.Application
import android.hardware.Sensor
import androidx.lifecycle.AndroidViewModel
import androidx.lifecycle.viewModelScope
import com.hawhamburg.micr0bu.data.GnssReading
import com.hawhamburg.micr0bu.data.ImuReading
import com.hawhamburg.micr0bu.data.SensorRepository
import com.hawhamburg.micr0bu.data.db.AppDatabase
import com.hawhamburg.micr0bu.data.db.SessionEntity
import kotlinx.coroutines.Job
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.launch
import java.io.BufferedWriter
import java.io.File
import java.io.FileWriter
import java.text.SimpleDateFormat
import java.util.Date
import java.util.Locale
import java.util.UUID
// UI model (no Room annotations)
data class RecordingSession(
val id: String,
val startTime: Long,
val endTime: Long,
val durationSeconds: Long,
val sampleCount: Int,
)
data class SensorUiState(
// Sensor readings
val gnss: GnssReading? = null,
val accel: ImuReading? = null,
val gyro: ImuReading? = null,
val magnet: ImuReading? = null,
val pressureHpa: Float? = null,
// Permissions & availability
val locationPermissionGranted: Boolean = false,
val hasGyroscope: Boolean = true,
val hasMagnetometer: Boolean = true,
val hasBarometer: Boolean = true,
// OBU connection (stub)
val obuConnected: Boolean = false,
val obuDeviceName: String? = null,
// Recording
val isRecording: Boolean = false,
val recordingStartTime: Long? = null,
val recordingDurationSeconds: Long = 0,
val recordingSampleCount: Int = 0,
val sessions: List<RecordingSession> = emptyList(),
// Settings
val accelEnabled: Boolean = true,
val gyroEnabled: Boolean = true,
val magnetEnabled: Boolean = true,
val barometerEnabled: Boolean = true,
val gnssEnabled: Boolean = true,
val developerMode: Boolean = true,
val useMetric: Boolean = true,
val darkTheme: Boolean = true,
)
private fun SessionEntity.toUiModel() = RecordingSession(
id = id,
startTime = startTime,
endTime = endTime,
durationSeconds = durationSeconds,
sampleCount = sampleCount,
)
class SensorViewModel(application: Application) : AndroidViewModel(application) {
private val repo = SensorRepository(application)
private val dao = AppDatabase.getInstance(application).sessionDao()
private val _state = MutableStateFlow(
SensorUiState(
hasGyroscope = repo.hasSensor(Sensor.TYPE_GYROSCOPE),
hasMagnetometer = repo.hasSensor(Sensor.TYPE_MAGNETIC_FIELD),
hasBarometer = repo.hasSensor(Sensor.TYPE_PRESSURE),
)
)
val state: StateFlow<SensorUiState> = _state.asStateFlow()
private var recordingTimerJob: Job? = null
private var recordingSessionId: String = ""
private var csvWriter: BufferedWriter? = null
private val isoFmt = SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss.SSS'Z'", Locale.US)
init {
viewModelScope.launch {
dao.getAllSessions().catch { }.collect { entities ->
_state.value = _state.value.copy(sessions = entities.map { it.toUiModel() })
}
}
}
// ── Permissions ───────────────────────────────────────────────────────────
fun onLocationPermissionGranted() {
_state.value = _state.value.copy(locationPermissionGranted = true)
startGnss()
}
// ── Sensor streams ────────────────────────────────────────────────────────
fun startImuStreams() {
viewModelScope.launch {
repo.accelerometerFlow().catch { }.collect { r ->
if (_state.value.accelEnabled) {
_state.value = _state.value.copy(accel = r)
if (_state.value.isRecording) {
_state.value = _state.value.copy(
recordingSampleCount = _state.value.recordingSampleCount + 1
)
writeImuRow("accel", r)
}
}
}
}
viewModelScope.launch {
repo.gyroscopeFlow().catch { }.collect { r ->
if (_state.value.gyroEnabled) {
_state.value = _state.value.copy(gyro = r)
if (_state.value.isRecording) writeImuRow("gyro", r)
}
}
}
viewModelScope.launch {
repo.magnetometerFlow().catch { }.collect { r ->
if (_state.value.magnetEnabled) {
_state.value = _state.value.copy(magnet = r)
if (_state.value.isRecording) writeImuRow("magnet", r)
}
}
}
viewModelScope.launch {
repo.barometerFlow().catch { }.collect { hpa ->
if (_state.value.barometerEnabled) {
_state.value = _state.value.copy(pressureHpa = hpa)
if (_state.value.isRecording) writeBaroRow(hpa)
}
}
}
}
private fun startGnss() {
viewModelScope.launch {
repo.gnssFlow().catch { }.collect { r ->
if (_state.value.gnssEnabled) {
_state.value = _state.value.copy(gnss = r)
if (_state.value.isRecording) {
_state.value = _state.value.copy(
recordingSampleCount = _state.value.recordingSampleCount + 1
)
writeGnssRow(r)
}
}
}
}
}
// ── Recording ─────────────────────────────────────────────────────────────
fun toggleRecording() {
if (_state.value.isRecording) stopRecording() else startRecording()
}
private fun startRecording() {
recordingSessionId = UUID.randomUUID().toString()
val startTime = System.currentTimeMillis()
// Open streaming CSV file
val sessionsDir = File(getApplication<Application>().filesDir, "sessions")
.also { it.mkdirs() }
runCatching {
csvWriter = BufferedWriter(FileWriter(File(sessionsDir, "$recordingSessionId.csv")))
csvWriter?.apply {
appendLine("# MicrOBU Session Export")
appendLine("# Generated by MicrOBU v0.2.0 — HAW Hamburg / Project MicrOBU")
appendLine("# Session ID,$recordingSessionId")
appendLine("# Start,${isoFmt.format(Date(startTime))}")
appendLine()
// Columns: type | timestamp_ms | timestamp_iso |
// lat | lon | alt_m | speed_ms | bearing_deg | accuracy_m |
// x | y | z | pressure_hpa
appendLine("type,timestamp_ms,timestamp_iso,lat,lon,alt_m,speed_ms,bearing_deg,accuracy_m,x,y,z,pressure_hpa")
}
}
_state.value = _state.value.copy(
isRecording = true,
recordingStartTime = startTime,
recordingDurationSeconds = 0,
recordingSampleCount = 0,
)
recordingTimerJob = viewModelScope.launch {
while (true) {
delay(1_000)
_state.value = _state.value.copy(
recordingDurationSeconds = _state.value.recordingDurationSeconds + 1
)
}
}
}
private fun stopRecording() {
recordingTimerJob?.cancel()
val s = _state.value
val endTime = System.currentTimeMillis()
// Append footer and close the CSV file
runCatching {
csvWriter?.apply {
appendLine()
appendLine("# End,${isoFmt.format(Date(endTime))}")
appendLine("# Duration (s),${s.recordingDurationSeconds}")
appendLine("# Samples captured,${s.recordingSampleCount}")
flush()
close()
}
}
csvWriter = null
val entity = SessionEntity(
id = recordingSessionId,
startTime = s.recordingStartTime ?: endTime,
endTime = endTime,
durationSeconds = s.recordingDurationSeconds,
sampleCount = s.recordingSampleCount,
)
viewModelScope.launch { dao.insertSession(entity) }
_state.value = s.copy(isRecording = false)
}
// ── CSV row writers ───────────────────────────────────────────────────────
// Columns: type, timestamp_ms, timestamp_iso,
// lat, lon, alt_m, speed_ms, bearing_deg, accuracy_m,
// x, y, z, pressure_hpa (13 total)
private fun writeGnssRow(r: GnssReading) {
val ts = r.timestamp
csvWriter?.appendLine(
"gnss,$ts,${isoFmt.format(Date(ts))},${r.latitude},${r.longitude},${r.altitude},${r.speedMs},${r.bearingDeg},${r.accuracyM},,,,"
)
}
private fun writeImuRow(type: String, r: ImuReading) {
val ts = r.timestamp
csvWriter?.appendLine(
"$type,$ts,${isoFmt.format(Date(ts))},,,,,,,${r.x},${r.y},${r.z},"
)
}
private fun writeBaroRow(hpa: Float) {
val ts = System.currentTimeMillis()
csvWriter?.appendLine(
"baro,$ts,${isoFmt.format(Date(ts))},,,,,,,,,,$hpa"
)
}
// ── Settings ──────────────────────────────────────────────────────────────
fun deleteSession(id: String) {
viewModelScope.launch { dao.deleteById(id) }
// Remove the associated CSV data file
runCatching {
File(File(getApplication<Application>().filesDir, "sessions"), "$id.csv").delete()
}
}
fun setAccelEnabled(v: Boolean) { _state.value = _state.value.copy(accelEnabled = v) }
fun setGyroEnabled(v: Boolean) { _state.value = _state.value.copy(gyroEnabled = v) }
fun setMagnetEnabled(v: Boolean) { _state.value = _state.value.copy(magnetEnabled = v) }
fun setBarometerEnabled(v: Boolean) { _state.value = _state.value.copy(barometerEnabled = v) }
fun setGnssEnabled(v: Boolean) {
_state.value = _state.value.copy(gnssEnabled = v)
if (!v) _state.value = _state.value.copy(gnss = null)
}
fun setDeveloperMode(v: Boolean) { _state.value = _state.value.copy(developerMode = v) }
fun setUseMetric(v: Boolean) { _state.value = _state.value.copy(useMetric = v) }
fun setDarkTheme(v: Boolean) { _state.value = _state.value.copy(darkTheme = v) }
}
@@ -0,0 +1,81 @@
package com.hawhamburg.micr0bu.viewmodel
import android.app.Application
import android.content.Intent
import androidx.lifecycle.AndroidViewModel
import androidx.lifecycle.viewModelScope
import com.hawhamburg.micr0bu.data.TripRepository
import com.hawhamburg.micr0bu.data.db.AppDatabase
import com.hawhamburg.micr0bu.data.db.DetectedEventEntity
import com.hawhamburg.micr0bu.data.db.RecordedTripEntity
import com.hawhamburg.micr0bu.service.TripRecordingService
import com.hawhamburg.micr0bu.service.TripServiceBus
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.launch
/**
* ViewModel for the Trip Recording feature (Phase A).
*
* Sends [ACTION_START] / [ACTION_STOP] intents to [TripRecordingService] and
* exposes service state via [TripServiceBus]. Also provides trip history and
* per-trip event lists for the History and Review screens.
*/
class TripRecordingViewModel(application: Application) : AndroidViewModel(application) {
private val repository = TripRepository(AppDatabase.getInstance(application))
// ── Service state (live during a recording) ───────────────────────────────
/** Live recording state forwarded from [TripServiceBus]. */
val serviceState: StateFlow<TripServiceBus.State> = TripServiceBus.state
// ── Trip history ─────────────────────────────────────────────────────────
/** All recorded trips, newest first. */
val trips: Flow<List<RecordedTripEntity>> = repository.getAllTrips()
// ── Trip review ───────────────────────────────────────────────────────────
private val _selectedTripEvents = MutableStateFlow<List<DetectedEventEntity>>(emptyList())
val selectedTripEvents: StateFlow<List<DetectedEventEntity>> = _selectedTripEvents.asStateFlow()
/** Load events for [tripId] into [selectedTripEvents]. */
fun loadTripEvents(tripId: Long) {
viewModelScope.launch {
repository.getEventsForTrip(tripId).collect { events ->
_selectedTripEvents.value = events
}
}
}
// ── Recording control ─────────────────────────────────────────────────────
/**
* Starts the foreground recording service if not already running,
* or stops it if a trip is already active.
*/
fun toggleRecording() {
if (serviceState.value.isRecording) stopRecording() else startRecording()
}
fun startRecording() {
val intent = Intent(getApplication(), TripRecordingService::class.java)
.apply { action = TripRecordingService.ACTION_START }
getApplication<Application>().startForegroundService(intent)
}
fun stopRecording() {
val intent = Intent(getApplication(), TripRecordingService::class.java)
.apply { action = TripRecordingService.ACTION_STOP }
getApplication<Application>().startService(intent)
}
// ── Trip management ───────────────────────────────────────────────────────
fun deleteTrip(tripId: Long) {
viewModelScope.launch { repository.deleteTrip(tripId) }
}
}