From c13c3891c95476000d84841d1f37c7d9c709be22 Mon Sep 17 00:00:00 2001 From: Ashin Walpola Date: Fri, 2 Oct 2026 16:16:02 +0200 Subject: [PATCH] Add a V2X data logger to Settings > Developer Records every received V2X message (CAM, DENM, SPATEM, MAPEM, VAM) into one SQLite file per run, with a timestamp, the phone's GNSS fix, the OBU's GNSS fix where there is one (CiT One), and RSSI where the hardware reports it (ESP32-C5). Works on both hardware paths: V2X_RX frames from the ESP32 link, and the raw v2x/rx plus processed v2x-uca topics over MQTT. The UPER bytes are stored untouched; CAM, DENM and SPATEM also get a few decoded columns. A foreground service keeps it recording with the screen locked. Subscribes to v2x/rx/mapem so MAPEM is received on the CiT One. The current ESP32-C5 firmware forwards only BTP ports 2001, 2002 and 2004, so MAPEM and VAM appear in an ESP32 log only once it accepts 2003 and 2018. Checked on a Pixel 9 Pro against a simulated car and an RSU: 10 minutes with the screen locked and forced deep Doze, no gaps, files pass integrity_check. --- TODO.md | 19 + app/src/main/AndroidManifest.xml | 6 + .../com/hawhamburg/micr0bu/MainActivity.kt | 17 + .../micr0bu/data/cam/CamUseCaseRepository.kt | 3 + .../micr0bu/data/log/V2xLogWriter.kt | 235 ++++++++++ .../hawhamburg/micr0bu/data/log/V2xLogger.kt | 413 ++++++++++++++++++ .../micr0bu/data/mqtt/MqttRepository.kt | 6 +- .../micr0bu/domain/log/ItsMessageKind.kt | 79 ++++ .../micr0bu/service/V2xLogService.kt | 135 ++++++ .../micr0bu/ui/navigation/AppNavigation.kt | 1 + .../micr0bu/ui/screens/SettingsScreen.kt | 12 +- .../ui/screens/V2xLogSettingsScreen.kt | 247 +++++++++++ .../micr0bu/viewmodel/V2xLogViewModel.kt | 45 ++ app/src/main/res/values-de/strings.xml | 24 + app/src/main/res/values/strings.xml | 24 + app/src/main/res/xml/file_paths.xml | 2 + .../hawhamburg/micr0bu/ItsMessageKindTest.kt | 73 ++++ 17 files changed, 1336 insertions(+), 5 deletions(-) create mode 100644 app/src/main/java/com/hawhamburg/micr0bu/data/log/V2xLogWriter.kt create mode 100644 app/src/main/java/com/hawhamburg/micr0bu/data/log/V2xLogger.kt create mode 100644 app/src/main/java/com/hawhamburg/micr0bu/domain/log/ItsMessageKind.kt create mode 100644 app/src/main/java/com/hawhamburg/micr0bu/service/V2xLogService.kt create mode 100644 app/src/main/java/com/hawhamburg/micr0bu/ui/screens/V2xLogSettingsScreen.kt create mode 100644 app/src/main/java/com/hawhamburg/micr0bu/viewmodel/V2xLogViewModel.kt create mode 100644 app/src/test/java/com/hawhamburg/micr0bu/ItsMessageKindTest.kt diff --git a/TODO.md b/TODO.md index cb34d85..89fbbbf 100644 --- a/TODO.md +++ b/TODO.md @@ -5,6 +5,25 @@ Engineering to-do list. The reviewer-facing open items live in ## Waiting on hardware +### On-air check of the V2X data logger (added 2026-10-02) + +Needs the phone, a second station sending (the simulated car works) and, for the CiT One part, the OBU. +Settings > Developer > V2X data logger. 110 unit tests pass. ESP32-C5 path checked on the Pixel 9 Pro on +2026-10-02 (see the ticked item); the rest has not run on a device. + +- [x] **ESP32-C5 (done 2026-10-02, simulated car + RSU on air):** 153 CAMs and 60 SPATEMs in 30 s, RSSI -56..-36 dBm, + phone GNSS on every row, `obu_*` NULL, file integrity ok. MAPEM 0 as expected (firmware drops port 2003). Original steps: start a log, let the simulated car's CAMs arrive, stop. Pull the `.db` from + `files/v2x_logs/` and check `messages`: `msg_type='CAM'`, `source='esp32'`, `rssi_dbm` filled, + `phone_lat/lon` filled, `obu_*` NULL, `payload` decodes with `asn1tools`. +- [ ] **CiT One:** same, `source='cit_one'`, `rssi_dbm` NULL, `obu_*` filled from `v2x/rx/obu_gnss`. + Expect both `channel='raw'` and `channel='processed'` rows for the same CAM. +- [ ] **Screen off for a few minutes** with the log running: rows keep arriving (foreground service). +- [ ] **MAPEM and VAM on the ESP32-C5 need a firmware change.** `gn_unwrap.c` accepts only BTP ports + 2001, 2002, 2004; add 2003 (MAPEM) and 2018 (VAM) and reflash. Real MAPEMs are mostly over the + 512-byte link limit, so most will still be counted as oversize drops. The app side already logs + any port. On the CiT One, MAPEM arrives via `v2x/rx/mapem` (subscribed now) or the processed + `v2x-uca/output/json/map` topic; the API lists no VAM topic, so no VAM rows are expected there. + ### Signed-TX firmware (vanetza-idf port), VAM and BLE: first on-air checks (added 2026-09-23) obu-firmware is now a port of the colleague's `microbu-esp32c5` station (vanetza-idf, TS 103 097 diff --git a/app/src/main/AndroidManifest.xml b/app/src/main/AndroidManifest.xml index 3a52451..a791d76 100644 --- a/app/src/main/AndroidManifest.xml +++ b/app/src/main/AndroidManifest.xml @@ -72,6 +72,12 @@ android:foregroundServiceType="location" android:exported="false" /> + + + handleSpatUper(envelope.payload, rssiDbm = null, source = "mqtt") + // No MAPEM decoder yet; it is received for the V2X data logger only. + RAW_MAPEM_TOPIC -> Unit else -> Log.w(TAG, "rawV2x: unexpected topic ${raw.topic}") } } diff --git a/app/src/main/java/com/hawhamburg/micr0bu/data/log/V2xLogWriter.kt b/app/src/main/java/com/hawhamburg/micr0bu/data/log/V2xLogWriter.kt new file mode 100644 index 0000000..4df35d3 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/data/log/V2xLogWriter.kt @@ -0,0 +1,235 @@ +package com.hawhamburg.micr0bu.data.log + +import android.database.sqlite.SQLiteDatabase +import android.database.sqlite.SQLiteStatement +import com.hawhamburg.micr0bu.domain.log.ItsMessageKind +import java.io.File + +/** One GNSS fix, as stored next to a logged message. Phone and OBU fixes share this shape. */ +data class LogGnss( + val latitude: Double, + val longitude: Double, + val altitudeM: Double? = null, + val speedMps: Double? = null, + val headingDeg: Double? = null, + val accuracyM: Double? = null, + /** Epoch ms the fix was taken, so a stale fix shows as one instead of passing for current. */ + val fixMs: Long, +) + +/** One received V2X message, ready to be written. */ +class LogMessage( + val tsMs: Long, + /** `esp32` or `cit_one`: which link delivered it. */ + val source: String, + /** `raw` (UPER bytes in [payload]) or `processed` (the Use Case app's JSON in [payloadText]). */ + val channel: String, + val kind: ItsMessageKind, + val messageId: Int?, + val btpPort: Int?, + val stationId: Long?, + val isOwn: Boolean?, + val rssiDbm: Int?, + val geoRadiusM: Int?, + val senderLat: Double?, + val senderLon: Double?, + val senderSpeedMps: Double?, + val senderHeadingDeg: Double?, + val summary: String?, + val payload: ByteArray?, + val payloadText: String?, + val phone: LogGnss?, + val obu: LogGnss?, +) + +/** One GNSS fix for the `gnss` table: the receiver's own track, independent of any message. */ +class LogGnssSample(val tsMs: Long, val source: String, val fix: LogGnss) + +/** + * A V2X log: one self-contained SQLite file. Not Room, on purpose: a log is a file the developer + * pulls off the phone and opens in any SQLite tool or pandas, so it carries its own schema and + * needs no migration path or app-side entity classes. The app never reads one back. + * + * Not thread-safe; [V2xLogger] drives it from one writer coroutine. Each [writeBatch] is one + * transaction on the default rollback journal, so a killed process leaves a file that is + * complete up to the last batch, and a single file with no `-wal` sidecar to forget when sharing. + * + * Columns that do not apply to a row are NULL rather than invented: `rssi_dbm` is NULL for the + * CiT One, `obu_*` is NULL when no OBU position is available (always on the ESP32-C5, which has no + * GNSS), and the `sender_*` fields are NULL for messages the app has no decoder for. + */ +class V2xLogWriter(val file: File, meta: Map) { + + private val db: SQLiteDatabase = SQLiteDatabase.openOrCreateDatabase(file, null) + private val insertMessage: SQLiteStatement + private val insertGnss: SQLiteStatement + + init { + db.execSQL( + """ + CREATE TABLE meta (key TEXT PRIMARY KEY, value TEXT) + """.trimIndent() + ) + db.execSQL( + """ + CREATE TABLE messages ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ts_ms INTEGER NOT NULL, + source TEXT NOT NULL, + channel TEXT NOT NULL, + msg_type TEXT NOT NULL, + message_id INTEGER, + btp_port INTEGER, + station_id INTEGER, + is_own INTEGER, + rssi_dbm INTEGER, + geo_radius_m INTEGER, + sender_lat REAL, + sender_lon REAL, + sender_speed_mps REAL, + sender_heading_deg REAL, + summary TEXT, + payload BLOB, + payload_text TEXT, + phone_lat REAL, + phone_lon REAL, + phone_alt_m REAL, + phone_speed_mps REAL, + phone_heading_deg REAL, + phone_accuracy_m REAL, + phone_fix_ms INTEGER, + obu_lat REAL, + obu_lon REAL, + obu_alt_m REAL, + obu_speed_mps REAL, + obu_heading_deg REAL, + obu_fix_ms INTEGER + ) + """.trimIndent() + ) + db.execSQL("CREATE INDEX idx_messages_ts ON messages(ts_ms)") + db.execSQL("CREATE INDEX idx_messages_type ON messages(msg_type)") + db.execSQL( + """ + CREATE TABLE gnss ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + ts_ms INTEGER NOT NULL, + source TEXT NOT NULL, + lat REAL NOT NULL, + lon REAL NOT NULL, + alt_m REAL, + speed_mps REAL, + heading_deg REAL, + accuracy_m REAL, + fix_ms INTEGER NOT NULL + ) + """.trimIndent() + ) + db.execSQL("CREATE INDEX idx_gnss_ts ON gnss(ts_ms)") + // Same rows, with the timestamp also as readable UTC, for anyone browsing in a SQLite GUI. + db.execSQL( + "CREATE VIEW messages_v AS SELECT *, " + + "strftime('%Y-%m-%d %H:%M:%f', ts_ms / 1000.0, 'unixepoch') AS ts_utc FROM messages" + ) + db.execSQL("PRAGMA user_version = $SCHEMA_VERSION") + for ((k, v) in meta) setMeta(k, v) + setMeta("schema_version", SCHEMA_VERSION.toString()) + + insertMessage = db.compileStatement( + "INSERT INTO messages (ts_ms, source, channel, msg_type, message_id, btp_port, station_id, " + + "is_own, rssi_dbm, geo_radius_m, sender_lat, sender_lon, sender_speed_mps, " + + "sender_heading_deg, summary, payload, payload_text, " + + "phone_lat, phone_lon, phone_alt_m, phone_speed_mps, phone_heading_deg, " + + "phone_accuracy_m, phone_fix_ms, obu_lat, obu_lon, obu_alt_m, obu_speed_mps, " + + "obu_heading_deg, obu_fix_ms) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)" + ) + insertGnss = db.compileStatement( + "INSERT INTO gnss (ts_ms, source, lat, lon, alt_m, speed_mps, heading_deg, accuracy_m, fix_ms) " + + "VALUES (?,?,?,?,?,?,?,?,?)" + ) + } + + fun setMeta(key: String, value: String) { + db.execSQL("INSERT OR REPLACE INTO meta (key, value) VALUES (?, ?)", arrayOf(key, value)) + } + + /** Writes [messages] and [gnss] in one transaction. */ + fun writeBatch(messages: List, gnss: List) { + if (messages.isEmpty() && gnss.isEmpty()) return + db.beginTransaction() + try { + for (m in messages) bindAndRun(m) + for (g in gnss) bindAndRun(g) + db.setTransactionSuccessful() + } finally { + db.endTransaction() + } + } + + fun close() { + runCatching { insertMessage.close() } + runCatching { insertGnss.close() } + runCatching { db.close() } + } + + private fun bindAndRun(m: LogMessage) { + val s = insertMessage + s.clearBindings() + s.bind(1, m.tsMs) + s.bind(2, m.source) + s.bind(3, m.channel) + s.bind(4, m.kind.label) + s.bind(5, m.messageId?.toLong()) + s.bind(6, m.btpPort?.toLong()) + s.bind(7, m.stationId) + s.bind(8, m.isOwn?.let { if (it) 1L else 0L }) + s.bind(9, m.rssiDbm?.toLong()) + s.bind(10, m.geoRadiusM?.toLong()) + s.bind(11, m.senderLat) + s.bind(12, m.senderLon) + s.bind(13, m.senderSpeedMps) + s.bind(14, m.senderHeadingDeg) + s.bind(15, m.summary) + s.bind(16, m.payload) + s.bind(17, m.payloadText) + s.bind(18, m.phone?.latitude) + s.bind(19, m.phone?.longitude) + s.bind(20, m.phone?.altitudeM) + s.bind(21, m.phone?.speedMps) + s.bind(22, m.phone?.headingDeg) + s.bind(23, m.phone?.accuracyM) + s.bind(24, m.phone?.fixMs) + s.bind(25, m.obu?.latitude) + s.bind(26, m.obu?.longitude) + s.bind(27, m.obu?.altitudeM) + s.bind(28, m.obu?.speedMps) + s.bind(29, m.obu?.headingDeg) + s.bind(30, m.obu?.fixMs) + s.executeInsert() + } + + private fun bindAndRun(g: LogGnssSample) { + val s = insertGnss + s.clearBindings() + s.bind(1, g.tsMs) + s.bind(2, g.source) + s.bind(3, g.fix.latitude) + s.bind(4, g.fix.longitude) + s.bind(5, g.fix.altitudeM) + s.bind(6, g.fix.speedMps) + s.bind(7, g.fix.headingDeg) + s.bind(8, g.fix.accuracyM) + s.bind(9, g.fix.fixMs) + s.executeInsert() + } + + // clearBindings() leaves every parameter NULL, so a null value is simply not bound. + private fun SQLiteStatement.bind(i: Int, v: Long?) { if (v != null) bindLong(i, v) } + private fun SQLiteStatement.bind(i: Int, v: Double?) { if (v != null) bindDouble(i, v) } + private fun SQLiteStatement.bind(i: Int, v: String?) { if (v != null) bindString(i, v) } + private fun SQLiteStatement.bind(i: Int, v: ByteArray?) { if (v != null) bindBlob(i, v) } + + companion object { + const val SCHEMA_VERSION = 1 + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/data/log/V2xLogger.kt b/app/src/main/java/com/hawhamburg/micr0bu/data/log/V2xLogger.kt new file mode 100644 index 0000000..c59ebf3 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/data/log/V2xLogger.kt @@ -0,0 +1,413 @@ +package com.hawhamburg.micr0bu.data.log + +import android.content.Context +import android.os.Build +import android.util.Log +import com.hawhamburg.micr0bu.data.GnssReading +import com.hawhamburg.micr0bu.data.SensorRepository +import com.hawhamburg.micr0bu.data.cam.CamUseCaseRepository +import com.hawhamburg.micr0bu.data.mqtt.MessageDirection +import com.hawhamburg.micr0bu.data.mqtt.MqttRepository +import com.hawhamburg.micr0bu.data.mqtt.RAW_CAM_TOPIC +import com.hawhamburg.micr0bu.data.mqtt.RAW_DENM_TOPIC +import com.hawhamburg.micr0bu.data.mqtt.RAW_MAPEM_TOPIC +import com.hawhamburg.micr0bu.data.mqtt.RAW_SPATEM_TOPIC +import com.hawhamburg.micr0bu.data.mqtt.RecvV2xMessage +import com.hawhamburg.micr0bu.data.transport.Esp32Link +import com.hawhamburg.micr0bu.data.transport.Esp32LinkState +import com.hawhamburg.micr0bu.data.transport.ObuHardware +import com.hawhamburg.micr0bu.data.transport.SerialFrameType +import com.hawhamburg.micr0bu.data.transport.V2xRxFrame +import com.hawhamburg.micr0bu.domain.asn1.DenmUperCodec +import com.hawhamburg.micr0bu.domain.asn1.RealAsn1UperCodec +import com.hawhamburg.micr0bu.domain.asn1.SpatemUperCodec +import com.hawhamburg.micr0bu.domain.cam.ObuGnssParser +import com.hawhamburg.micr0bu.domain.log.ItsMessageKind +import com.hawhamburg.micr0bu.domain.log.classifyItsPdu +import dagger.hilt.android.qualifiers.ApplicationContext +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancelAndJoin +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.flow.asStateFlow +import kotlinx.coroutines.flow.update +import kotlinx.coroutines.launch +import kotlinx.coroutines.withTimeoutOrNull +import java.io.File +import java.text.SimpleDateFormat +import java.util.Date +import java.util.Locale +import javax.inject.Inject +import javax.inject.Singleton + +private const val TAG = "V2xLogger" +private const val LOG_DIR = "v2x_logs" +private const val PROCESSED_PREFIX = "v2x-uca/output/json/" +private const val OBU_GNSS_TOPIC = "v2x/rx/obu_gnss" +private const val SOURCE_ESP32 = "esp32" +private const val SOURCE_CIT_ONE = "cit_one" +private const val SOURCE_PHONE = "phone" +private const val SOURCE_OBU = "obu" + +/** How long the writer collects rows before it commits them. */ +private const val FLUSH_INTERVAL_MS = 250L + +/** What the logging screen shows about the log in progress. */ +data class V2xLogState( + val active: Boolean = false, + val fileName: String? = null, + val startedMs: Long = 0L, + /** Rows written so far, by kind. Kinds never heard are absent. */ + val counts: Map = emptyMap(), + val total: Int = 0, + val fileBytes: Long = 0L, + /** Age of the newest phone / OBU fix at the time of the last update; null if none yet. */ + val phoneFixAgeMs: Long? = null, + val obuFixAgeMs: Long? = null, + /** Set when the file could not be created or written. */ + val error: String? = null, +) + +/** A log file on the device. */ +data class V2xLogFile(val file: File, val sizeBytes: Long, val modifiedMs: Long, val active: Boolean) { + val name: String get() = file.name +} + +/** + * Developer-tool logger: records every V2X message the phone receives, with a timestamp, the + * phone's GNSS fix, the OBU's GNSS fix where there is one, and RSSI where the hardware reports it, + * into one SQLite file per run (see [V2xLogWriter] for the schema). + * + * Works on both hardware paths and listens to whichever is selected: + * - **ESP32-C5**: the V2X_RX frames from [Esp32Link], with the firmware's RSSI. No OBU position, + * since that board has no GNSS. + * - **CiT One**: the raw `v2x/rx` protobuf topics over MQTT, and the Use Case app's processed + * `v2x-uca/output/json` topics as a second channel (the only source of MAPEM when the OBU does + * not publish `v2x/rx/mapem`). The OBU's position comes from `v2x/rx/obu_gnss`. The API carries + * no RSSI, so that column stays NULL. + * + * Every message is stored, including kinds the app cannot decode (MAPEM, VAM): the UPER bytes go + * in the file untouched, and where a decoder exists a few fields are decoded alongside for + * convenience. A message is never skipped because decoding failed. + * + * What it can hear is bounded by the hardware, not by this class: the current ESP32-C5 firmware + * forwards only BTP ports 2001, 2002 and 2004 (see `gn_unwrap.c`), so MAPEM and VAM appear in an + * ESP32 log only once the firmware accepts ports 2003 and 2018. + * + * A singleton, like the repositories it listens to. [com.hawhamburg.micr0bu.service.V2xLogService] + * keeps the process alive and the CPU awake while a log runs; this class does the work. + */ +@Singleton +class V2xLogger @Inject constructor( + @ApplicationContext private val context: Context, + private val mqttRepository: MqttRepository, + private val esp32Link: Esp32Link, + private val camUseCaseRepository: CamUseCaseRepository, + private val camCodec: RealAsn1UperCodec, +) { + private val scope = CoroutineScope(SupervisorJob() + Dispatchers.IO) + + private val _state = MutableStateFlow(V2xLogState()) + val state: StateFlow = _state.asStateFlow() + + private val _files = MutableStateFlow>(emptyList()) + /** Logs on the device, newest first. Call [refreshFiles] after anything that changes the folder. */ + val files: StateFlow> = _files.asStateFlow() + + private val logDir: File get() = File(context.filesDir, LOG_DIR).apply { mkdirs() } + + private var session: Job? = null + private var rows = Channel(Channel.UNLIMITED) + + @Volatile private var phoneFix: LogGnss? = null + @Volatile private var obuFix: LogGnss? = null + + fun start() { + if (session?.isActive == true) return + val startedMs = System.currentTimeMillis() + val hardware = mqttRepository.obuHardware.value + val file = File(logDir, fileNameFor(startedMs, hardware)) + val writer = try { + V2xLogWriter(file, metaFor(startedMs, hardware)) + } catch (e: Exception) { + Log.e(TAG, "cannot create ${file.name}", e) + file.delete() + _state.value = V2xLogState(error = e.message ?: e.javaClass.simpleName) + return + } + phoneFix = null + obuFix = null + rows = Channel(Channel.UNLIMITED) + _state.value = V2xLogState(active = true, fileName = file.name, startedMs = startedMs) + refreshFiles() + + session = scope.launch { + val sources = listOf( + launch { collectEsp32() }, + launch { collectCitOneRaw() }, + launch { collectCitOneMessages() }, + launch { collectPhoneGnss() }, + ) + try { + runWriter(writer, file) + } finally { + sources.forEach { it.cancel() } + runCatching { writer.setMeta("ended_ms", System.currentTimeMillis().toString()) } + writer.close() + } + } + } + + /** Stops the log and waits for the last rows to reach the file. */ + suspend fun stop() { + val running = session ?: return + // Closing the channel lets the writer drain everything already queued before it ends. + rows.close() + withTimeoutOrNull(5_000L) { running.join() } ?: running.cancelAndJoin() + session = null + _state.update { it.copy(active = false) } + refreshFiles() + } + + fun delete(log: V2xLogFile): Boolean { + if (log.active) return false + val ok = log.file.delete() + refreshFiles() + return ok + } + + fun refreshFiles() { + val activeName = _state.value.fileName.takeIf { _state.value.active } + _files.value = logDir.listFiles { f -> f.isFile && f.extension == "db" } + .orEmpty() + .map { V2xLogFile(it, it.length(), it.lastModified(), it.name == activeName) } + .sortedByDescending { it.modifiedMs } + } + + // ── Sources ─────────────────────────────────────────────────────────────── + + private suspend fun collectEsp32() { + esp32Link.incomingFrames.collect { frame -> + if (frame.type != SerialFrameType.V2X_RX) return@collect + if (mqttRepository.obuHardware.value != ObuHardware.ESP32_C5) return@collect + if (esp32Link.state.value != Esp32LinkState.CONNECTED) return@collect + val v2x = V2xRxFrame.parse(frame.payload) ?: return@collect + enqueueRaw( + source = SOURCE_ESP32, uper = v2x.uper, btpPort = v2x.btpPort, + rssiDbm = v2x.rssiDbm, geoRadiusM = v2x.geoArea?.radiusMeters, + ) + } + } + + private suspend fun collectCitOneRaw() { + mqttRepository.rawV2x.collect { raw -> + if (mqttRepository.obuHardware.value != ObuHardware.CIT_ONE) return@collect + val envelope = RecvV2xMessage.parse(raw.bytes) + if (envelope == null) { + Log.w(TAG, "unparseable RecvV2XMessage on ${raw.topic}, ${raw.bytes.size} bytes") + return@collect + } + enqueueRaw( + source = SOURCE_CIT_ONE, uper = envelope.payload, + // The envelope's own port when it has one; the topic says the same thing otherwise. + btpPort = envelope.destinationPort ?: portOfRawTopic(raw.topic), + rssiDbm = null, geoRadiusM = envelope.destAreaRadiusM, + tsMs = raw.timestamp, + ) + } + } + + private suspend fun collectCitOneMessages() { + mqttRepository.messages.collect { msg -> + if (msg.direction != MessageDirection.RX) return@collect + if (mqttRepository.obuHardware.value != ObuHardware.CIT_ONE) return@collect + if (msg.topic == OBU_GNSS_TOPIC) { + val ego = ObuGnssParser.parseEgo(msg.payload, timestamp = msg.timestamp) ?: return@collect + val fix = LogGnss( + latitude = ego.latitude, longitude = ego.longitude, + speedMps = ego.speedMps, headingDeg = ego.headingDeg, fixMs = msg.timestamp, + ) + obuFix = fix + rows.trySend(LogGnssSample(msg.timestamp, SOURCE_OBU, fix)) + return@collect + } + if (!msg.topic.startsWith(PROCESSED_PREFIX)) return@collect + val kind = ItsMessageKind.fromProcessedTopicName(msg.topic.removePrefix(PROCESSED_PREFIX)) + ?: return@collect + rows.trySend( + LogMessage( + tsMs = msg.timestamp, source = SOURCE_CIT_ONE, channel = "processed", kind = kind, + messageId = null, btpPort = null, stationId = null, isOwn = null, rssiDbm = null, + geoRadiusM = null, senderLat = null, senderLon = null, senderSpeedMps = null, + senderHeadingDeg = null, summary = null, payload = null, payloadText = msg.payload, + phone = phoneFix, obu = obuFix, + ) + ) + } + } + + /** + * Phone GNSS, re-subscribed in a loop for the same reason [CamUseCaseRepository] does: a log + * can be started before the location permission is granted, and a single attempt would then + * leave every row without a phone position for the whole run. + */ + private suspend fun collectPhoneGnss() { + val sensors = SensorRepository(context) + while (true) { + runCatching { + sensors.gnssFlow().collect { reading -> + val fix = reading.toLogGnss() + phoneFix = fix + rows.trySend(LogGnssSample(System.currentTimeMillis(), SOURCE_PHONE, fix)) + } + } + delay(5_000L) + } + } + + // ── Rows ────────────────────────────────────────────────────────────────── + + private fun enqueueRaw( + source: String, + uper: ByteArray, + btpPort: Int?, + rssiDbm: Int?, + geoRadiusM: Int?, + tsMs: Long = System.currentTimeMillis(), + ) { + val cls = classifyItsPdu(uper, btpPort) + val decoded = decode(cls.kind, uper, rssiDbm, geoRadiusM, tsMs) + rows.trySend( + LogMessage( + tsMs = tsMs, source = source, channel = "raw", kind = cls.kind, + messageId = cls.messageId, btpPort = btpPort, stationId = cls.stationId, + isOwn = cls.stationId?.let { camUseCaseRepository.isOwnStationId(it) }, + rssiDbm = rssiDbm, geoRadiusM = geoRadiusM, + senderLat = decoded.lat, senderLon = decoded.lon, + senderSpeedMps = decoded.speedMps, senderHeadingDeg = decoded.headingDeg, + summary = decoded.summary, payload = uper, payloadText = null, + phone = phoneFix, obu = obuFix, + ) + ) + } + + private class Decoded( + val lat: Double? = null, val lon: Double? = null, + val speedMps: Double? = null, val headingDeg: Double? = null, + val summary: String? = null, + ) + + /** Decoded convenience fields for the kinds the app can decode. A failure just leaves them NULL. */ + private fun decode(kind: ItsMessageKind, uper: ByteArray, rssiDbm: Int?, geoRadiusM: Int?, tsMs: Long): Decoded = + runCatching { + when (kind) { + ItsMessageKind.CAM -> camCodec.decodeCam(uper, tsMs)?.let { + Decoded(it.latitude, it.longitude, it.speedMps, it.headingDeg, "stationType=${it.stationType}") + } + ItsMessageKind.DENM -> DenmUperCodec.decode(uper, tsMs, rssiDbm, geoRadiusM)?.let { + Decoded( + it.latitude, it.longitude, + summary = "seq=${it.sequenceNumber} cause=${it.causeCode}/${it.subCauseCode}" + + if (it.isTermination) " termination" else "", + ) + } + ItsMessageKind.SPATEM -> SpatemUperCodec.decode(uper, tsMs, rssiDbm)?.let { s -> + Decoded(summary = "intersections=" + s.intersections.joinToString(",") { it.key } + + " movements=" + s.intersections.sumOf { it.movements.size }) + } + else -> null + } + }.getOrNull() ?: Decoded() + + // ── Writer ──────────────────────────────────────────────────────────────── + + /** Drains [rows] into [writer] in short batches, until the channel is closed and empty. */ + private suspend fun runWriter(writer: V2xLogWriter, file: File) { + val counts = HashMap() + var total = 0 + var failure: String? = null + val messages = ArrayList() + val gnss = ArrayList() + + while (true) { + // Block for the first row of a batch, give more a moment to arrive, then take + // whatever is queued. No timeout around receive: a cancelled receive can lose the + // element it was handing over, which a logger must not do. + val first = rows.receiveCatching().getOrNull() ?: break + add(first, messages, gnss) + delay(FLUSH_INTERVAL_MS) + while (true) { + val next = rows.tryReceive().getOrNull() ?: break + add(next, messages, gnss) + } + try { + writer.writeBatch(messages, gnss) + } catch (e: Exception) { + // Disk full or the file went away. Say so on screen rather than logging into a void. + Log.e(TAG, "write failed", e) + failure = e.message ?: e.javaClass.simpleName + rows.close() + break + } + for (m in messages) counts.merge(m.kind, 1, Int::plus) + total += messages.size + messages.clear(); gnss.clear() + + val now = System.currentTimeMillis() + _state.update { + it.copy( + counts = counts.toMap(), total = total, fileBytes = file.length(), + phoneFixAgeMs = phoneFix?.let { f -> now - f.fixMs }, + obuFixAgeMs = obuFix?.let { f -> now - f.fixMs }, + ) + } + } + _state.update { it.copy(fileBytes = file.length(), error = failure ?: it.error, active = it.active && failure == null) } + } + + private fun add(row: Any, messages: MutableList, gnss: MutableList) { + when (row) { + is LogMessage -> messages += row + is LogGnssSample -> gnss += row + } + } + + // ── Helpers ─────────────────────────────────────────────────────────────── + + private fun portOfRawTopic(topic: String): Int? = when (topic) { + RAW_CAM_TOPIC -> 2001 + RAW_DENM_TOPIC -> 2002 + RAW_MAPEM_TOPIC -> 2003 + RAW_SPATEM_TOPIC -> 2004 + else -> null + } + + private fun GnssReading.toLogGnss() = LogGnss( + latitude = latitude, longitude = longitude, altitudeM = altitude, + speedMps = speedMs.toDouble(), headingDeg = bearingDeg.toDouble(), + accuracyM = accuracyM.toDouble(), fixMs = timestamp, + ) + + private fun fileNameFor(startedMs: Long, hardware: ObuHardware): String { + val stamp = SimpleDateFormat("yyyyMMdd_HHmmss", Locale.US).format(Date(startedMs)) + val mode = if (hardware == ObuHardware.ESP32_C5) "esp32" else "citone" + return "v2x_log_${stamp}_$mode.db" + } + + private fun metaFor(startedMs: Long, hardware: ObuHardware) = mapOf( + "started_ms" to startedMs.toString(), + "hardware" to if (hardware == ObuHardware.ESP32_C5) "esp32_c5" else "cit_one", + "phone" to "${Build.MANUFACTURER} ${Build.MODEL}", + "android_sdk" to Build.VERSION.SDK_INT.toString(), + "app_version" to runCatching { + context.packageManager.getPackageInfo(context.packageName, 0).versionName ?: "?" + }.getOrDefault("?"), + "note" to "ts_ms is phone wall-clock epoch ms. payload is the ITS PDU as UPER bytes. " + + "rssi_dbm is NULL on the CiT One; obu_* is NULL on the ESP32-C5.", + ) +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/MqttRepository.kt b/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/MqttRepository.kt index b5a5025..e58bf23 100644 --- a/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/MqttRepository.kt +++ b/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/MqttRepository.kt @@ -48,6 +48,9 @@ private val SUBSCRIBED_TOPICS = listOf( RAW_CAM_TOPIC, RAW_DENM_TOPIC, RAW_SPATEM_TOPIC, + // Nothing decodes MAPEM yet; subscribed so the V2X data logger records it. An OBU product + // configuration that does not publish it simply never delivers anything here. + RAW_MAPEM_TOPIC, // Kept subscribed as a fallback for an OBU whose product configuration does not publish the // raw topics, and because the Use Case app is still the only source of its own alert output. "v2x-uca/output/json/cam", @@ -61,8 +64,9 @@ private val SUBSCRIBED_TOPICS = listOf( const val RAW_CAM_TOPIC = "v2x/rx/cam" const val RAW_DENM_TOPIC = "v2x/rx/denm" const val RAW_SPATEM_TOPIC = "v2x/rx/spatem" +const val RAW_MAPEM_TOPIC = "v2x/rx/mapem" -private val RAW_V2X_TOPICS = setOf(RAW_CAM_TOPIC, RAW_DENM_TOPIC, RAW_SPATEM_TOPIC) +private val RAW_V2X_TOPICS = setOf(RAW_CAM_TOPIC, RAW_DENM_TOPIC, RAW_SPATEM_TOPIC, RAW_MAPEM_TOPIC) /** * A message straight off a `v2x/rx` topic, before the protobuf envelope is opened. diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/log/ItsMessageKind.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/log/ItsMessageKind.kt new file mode 100644 index 0000000..0afdafd --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/log/ItsMessageKind.kt @@ -0,0 +1,79 @@ +package com.hawhamburg.micr0bu.domain.log + +/** + * Which ITS message a received PDU is, for the V2X data logger. + * + * The logger records every message it hears, including ones the app has no decoder for (MAPEM, + * VAM), so classification cannot depend on a codec. It reads the two things every ETSI ITS PDU + * has in common instead: the `ItsPduHeader` at the start of the UPER bytes + * (`protocolVersion` 8 bits, `messageID` 8 bits, `stationID` 32 bits, none of it extensible, so + * it sits byte-aligned at the front), and the BTP destination port the transport reports. + * + * messageID is the primary source. The port is the fallback for a payload too short to carry a + * header. Mind the crossover: SPATEM is messageID 4 on port 2004 but MAPEM is messageID 5 on port + * 2003; the two numbering schemes are unrelated. + */ +enum class ItsMessageKind(val label: String, val messageId: Int, val btpPort: Int?) { + DENM("DENM", 1, 2002), + CAM("CAM", 2, 2001), + SPATEM("SPATEM", 4, 2004), + MAPEM("MAPEM", 5, 2003), + CPM("CPM", 14, null), + VAM("VAM", 16, 2018), + OTHER("OTHER", -1, null); + + companion object { + /** The five kinds the logger counts and shows on screen. */ + val TRACKED = listOf(CAM, DENM, SPATEM, MAPEM, VAM) + + fun fromMessageId(id: Int): ItsMessageKind? = entries.firstOrNull { it.messageId == id && it != OTHER } + + fun fromBtpPort(port: Int): ItsMessageKind? = entries.firstOrNull { it.btpPort == port } + + /** Matches the `v2x-uca/output/json/` suffix, for the CiT One's processed topics. */ + fun fromProcessedTopicName(name: String): ItsMessageKind? = when (name) { + "cam" -> CAM + "denm" -> DENM + "spat" -> SPATEM + "map" -> MAPEM + "cpm" -> CPM + else -> null + } + } +} + +/** The fixed 6-byte header at the start of every ETSI ITS PDU. */ +data class ItsPduHeader(val protocolVersion: Int, val messageId: Int, val stationId: Long) { + companion object { + const val SIZE = 6 + + /** Null when [uper] is too short to hold a header. */ + fun parse(uper: ByteArray): ItsPduHeader? { + if (uper.size < SIZE) return null + var station = 0L + for (i in 2 until SIZE) station = (station shl 8) or (uper[i].toLong() and 0xFF) + return ItsPduHeader( + protocolVersion = uper[0].toInt() and 0xFF, + messageId = uper[1].toInt() and 0xFF, + stationId = station, + ) + } + } +} + +/** What [classifyItsPdu] worked out about one received PDU. */ +data class ItsClassification(val kind: ItsMessageKind, val messageId: Int?, val stationId: Long?) + +/** + * Classifies [uper], preferring its header's messageID over [btpPort]. + * + * A message whose id and port disagree is classified by the id, since the id is what the sender + * actually put in the PDU; the port is only where it happened to arrive. + */ +fun classifyItsPdu(uper: ByteArray, btpPort: Int?): ItsClassification { + val header = ItsPduHeader.parse(uper) + val kind = header?.let { ItsMessageKind.fromMessageId(it.messageId) } + ?: btpPort?.let { ItsMessageKind.fromBtpPort(it) } + ?: ItsMessageKind.OTHER + return ItsClassification(kind, header?.messageId, header?.stationId) +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/service/V2xLogService.kt b/app/src/main/java/com/hawhamburg/micr0bu/service/V2xLogService.kt new file mode 100644 index 0000000..e911b1d --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/service/V2xLogService.kt @@ -0,0 +1,135 @@ +package com.hawhamburg.micr0bu.service + +import android.app.NotificationChannel +import android.app.NotificationManager +import android.app.PendingIntent +import android.app.Service +import android.content.Context +import android.content.Intent +import android.content.pm.ServiceInfo +import android.os.IBinder +import android.os.PowerManager +import androidx.core.app.NotificationCompat +import androidx.core.content.ContextCompat +import com.hawhamburg.micr0bu.MainActivity +import com.hawhamburg.micr0bu.R +import com.hawhamburg.micr0bu.data.log.V2xLogger +import dagger.hilt.android.AndroidEntryPoint +import javax.inject.Inject +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.launch + +/** + * Keeps the process alive, and the CPU awake, while [V2xLogger] records. + * + * The logger does the work; this exists only because a log is most useful on a drive, with the + * screen off, and a process with no foreground service is one the system may kill and whose GNSS + * callbacks it throttles. Same shape as [TripRecordingService], which does this for trips. + * + * Stops itself when the logger stops, including when the logger stops on a write failure. + */ +@AndroidEntryPoint +class V2xLogService : Service() { + + @Inject lateinit var logger: V2xLogger + + private val job = SupervisorJob() + private val scope = CoroutineScope(Dispatchers.Default + job) + private var watchJob: Job? = null + private var wakeLock: PowerManager.WakeLock? = null + + override fun onCreate() { + super.onCreate() + (getSystemService(NOTIFICATION_SERVICE) as NotificationManager).createNotificationChannel( + NotificationChannel(CHANNEL_ID, "V2X logging", NotificationManager.IMPORTANCE_LOW).apply { + description = "Shown while the developer V2X data logger is recording" + } + ) + } + + override fun onStartCommand(intent: Intent?, flags: Int, startId: Int): Int { + when (intent?.action) { + ACTION_START -> begin() + ACTION_STOP -> end() + } + return START_NOT_STICKY + } + + override fun onBind(intent: Intent?): IBinder? = null + + override fun onDestroy() { + wakeLock?.let { if (it.isHeld) it.release() } + job.cancel() + super.onDestroy() + } + + private fun begin() { + // Foreground first: Android allows five seconds from startForegroundService. + startForeground(NOTIFICATION_ID, notification(0), ServiceInfo.FOREGROUND_SERVICE_TYPE_LOCATION) + wakeLock = (getSystemService(POWER_SERVICE) as PowerManager) + .newWakeLock(PowerManager.PARTIAL_WAKE_LOCK, "micr0bu:V2xLog") + .also { it.acquire() } + logger.start() + if (!logger.state.value.active) { // the file could not be created; the screen shows why + end() + return + } + + watchJob?.cancel() + watchJob = scope.launch { + logger.state.collect { state -> + if (!state.active) { + stopForeground(STOP_FOREGROUND_REMOVE) + stopSelf() + } else { + (getSystemService(NOTIFICATION_SERVICE) as NotificationManager) + .notify(NOTIFICATION_ID, notification(state.total)) + } + } + } + } + + private fun end() { + scope.launch { + logger.stop() + stopForeground(STOP_FOREGROUND_REMOVE) + stopSelf() + } + } + + private fun notification(messages: Int) = + NotificationCompat.Builder(this, CHANNEL_ID) + .setContentTitle("Logging V2X messages") + .setContentText("$messages messages recorded") + .setSmallIcon(R.mipmap.ic_launcher_foreground) + .setOngoing(true) + .setOnlyAlertOnce(true) + .setContentIntent( + PendingIntent.getActivity( + this, 0, + Intent(this, MainActivity::class.java).apply { + flags = Intent.FLAG_ACTIVITY_SINGLE_TOP or Intent.FLAG_ACTIVITY_CLEAR_TOP + }, + PendingIntent.FLAG_IMMUTABLE or PendingIntent.FLAG_UPDATE_CURRENT, + ) + ) + .build() + + companion object { + const val ACTION_START = "com.hawhamburg.micr0bu.V2X_LOG_START" + const val ACTION_STOP = "com.hawhamburg.micr0bu.V2X_LOG_STOP" + private const val NOTIFICATION_ID = 9002 + private const val CHANNEL_ID = "v2x_logging" + + fun start(context: Context) = ContextCompat.startForegroundService( + context, Intent(context, V2xLogService::class.java).setAction(ACTION_START) + ) + + fun stop(context: Context) { + context.startService(Intent(context, V2xLogService::class.java).setAction(ACTION_STOP)) + } + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/ui/navigation/AppNavigation.kt b/app/src/main/java/com/hawhamburg/micr0bu/ui/navigation/AppNavigation.kt index 36bc42c..c6b7d30 100644 --- a/app/src/main/java/com/hawhamburg/micr0bu/ui/navigation/AppNavigation.kt +++ b/app/src/main/java/com/hawhamburg/micr0bu/ui/navigation/AppNavigation.kt @@ -51,6 +51,7 @@ sealed class Screen(val route: String, val labelRes: Int) { data object SettingsMqtt : Screen("settings/mqtt", R.string.settings_mqtt_section) data object SettingsUseCaseAlerts : Screen("settings/usecases", R.string.settings_usecase_alerts) data object SettingsDeveloper : Screen("settings/developer", R.string.settings_developer) + data object SettingsV2xLog : Screen("settings/developer/v2x_log", R.string.v2xlog_title) data object SettingsAbout : Screen("settings/about", R.string.settings_about) } diff --git a/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/SettingsScreen.kt b/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/SettingsScreen.kt index b79d8b3..e078dee 100644 --- a/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/SettingsScreen.kt +++ b/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/SettingsScreen.kt @@ -137,7 +137,7 @@ private fun MenuRow(label: String, onClick: () -> Unit) { // ── Sub-screen shell ──────────────────────────────────────────────────────── @Composable -private fun SubScreen(title: String, onBack: () -> Unit, content: @Composable () -> Unit) { +internal fun SubScreen(title: String, onBack: () -> Unit, content: @Composable () -> Unit) { Column( modifier = Modifier .fillMaxSize() @@ -448,6 +448,7 @@ fun DeveloperSettingsScreen( state: SensorUiState, onDeveloperMode: (Boolean) -> Unit, onOpenSensorMonitor: () -> Unit, + onOpenV2xLogger: () -> Unit, onBack: () -> Unit, ) { SubScreen(stringResource(R.string.settings_developer), onBack) { @@ -457,6 +458,9 @@ fun DeveloperSettingsScreen( // Sensor Monitor lives here rather than in the bottom nav: a live phone-sensor feed is // a bench-diagnosis tool, not something a rider needs mid-ride. MenuRow(stringResource(R.string.sensor_monitor_title), onOpenSensorMonitor) + Divider() + // Records every received V2X message to a database file, on either OBU hardware. + MenuRow(stringResource(R.string.v2xlog_title), onOpenV2xLogger) if (state.developerMode) { Divider() DisabledRow(stringResource(R.string.settings_wifi), stringResource(R.string.settings_wifi_val)) @@ -481,7 +485,7 @@ fun AboutSettingsScreen(onBack: () -> Unit) { // ── Internal components ───────────────────────────────────────────────────── @Composable -private fun SectionCard(content: @Composable () -> Unit) { +internal fun SectionCard(content: @Composable () -> Unit) { Card( modifier = Modifier.fillMaxWidth(), colors = CardDefaults.cardColors(containerColor = MaterialTheme.colorScheme.surfaceVariant), @@ -699,7 +703,7 @@ private fun TwoWayChoice( private fun RowDivider() = HorizontalDivider(color = MaterialTheme.colorScheme.outline.copy(alpha = 0.4f)) @Composable -private fun Divider() = HorizontalDivider(color = MaterialTheme.colorScheme.outline.copy(alpha = 0.4f)) +internal fun Divider() = HorizontalDivider(color = MaterialTheme.colorScheme.outline.copy(alpha = 0.4f)) @Composable private fun SettingToggleRow( @@ -723,7 +727,7 @@ private fun SettingToggleRow( } @Composable -private fun DisabledRow(label: String, value: String) { +internal fun DisabledRow(label: String, value: String) { Row( modifier = Modifier.fillMaxWidth().padding(vertical = 10.dp), horizontalArrangement = Arrangement.SpaceBetween, diff --git a/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/V2xLogSettingsScreen.kt b/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/V2xLogSettingsScreen.kt new file mode 100644 index 0000000..0fe2645 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/V2xLogSettingsScreen.kt @@ -0,0 +1,247 @@ +package com.hawhamburg.micr0bu.ui.screens + +import android.content.Context +import android.content.Intent +import androidx.compose.foundation.layout.Arrangement +import androidx.compose.foundation.layout.Column +import androidx.compose.foundation.layout.Row +import androidx.compose.foundation.layout.Spacer +import androidx.compose.foundation.layout.fillMaxWidth +import androidx.compose.foundation.layout.height +import androidx.compose.foundation.layout.padding +import androidx.compose.material.icons.Icons +import androidx.compose.material.icons.filled.Delete +import androidx.compose.material.icons.filled.FiberManualRecord +import androidx.compose.material.icons.filled.Share +import androidx.compose.material.icons.filled.Stop +import androidx.compose.material3.AlertDialog +import androidx.compose.material3.Button +import androidx.compose.material3.ButtonDefaults +import androidx.compose.material3.Icon +import androidx.compose.material3.IconButton +import androidx.compose.material3.MaterialTheme +import androidx.compose.material3.Text +import androidx.compose.material3.TextButton +import androidx.compose.runtime.Composable +import androidx.compose.runtime.getValue +import androidx.compose.runtime.mutableStateOf +import androidx.compose.runtime.remember +import androidx.compose.runtime.setValue +import androidx.compose.ui.Alignment +import androidx.compose.ui.Modifier +import androidx.compose.ui.graphics.Color +import androidx.compose.ui.platform.LocalContext +import androidx.compose.ui.res.stringResource +import androidx.compose.ui.text.font.FontFamily +import androidx.compose.ui.text.font.FontWeight +import androidx.compose.ui.unit.dp +import androidx.core.content.FileProvider +import com.hawhamburg.micr0bu.R +import com.hawhamburg.micr0bu.data.log.V2xLogFile +import com.hawhamburg.micr0bu.data.log.V2xLogState +import com.hawhamburg.micr0bu.data.transport.ObuHardware +import com.hawhamburg.micr0bu.domain.log.ItsMessageKind +import java.text.SimpleDateFormat +import java.util.Date +import java.util.Locale + +private val modifiedFormat = SimpleDateFormat("yyyy-MM-dd HH:mm:ss", Locale.getDefault()) + +/** + * Settings > Developer > V2X data logger. + * + * Start and stop one log at a time, watch what it is catching, and share or delete the finished + * files. What is recorded, and from where, is described on [com.hawhamburg.micr0bu.data.log.V2xLogger]. + */ +@Composable +fun V2xLogSettingsScreen( + state: V2xLogState, + files: List, + obuHardware: ObuHardware, + onStart: () -> Unit, + onStop: () -> Unit, + onDelete: (V2xLogFile) -> Unit, + onBack: () -> Unit, +) { + val context = LocalContext.current + var pendingDelete by remember { mutableStateOf(null) } + val isEsp32 = obuHardware == ObuHardware.ESP32_C5 + + pendingDelete?.let { log -> + AlertDialog( + onDismissRequest = { pendingDelete = null }, + title = { Text(stringResource(R.string.v2xlog_delete_title)) }, + text = { Text(stringResource(R.string.v2xlog_delete_text, log.name)) }, + confirmButton = { + TextButton(onClick = { onDelete(log); pendingDelete = null }) { + Text(stringResource(R.string.log_delete_confirm), color = MaterialTheme.colorScheme.error) + } + }, + dismissButton = { + TextButton(onClick = { pendingDelete = null }) { Text(stringResource(R.string.log_cancel)) } + }, + ) + } + + SubScreen(stringResource(R.string.v2xlog_title), onBack) { + Text( + stringResource(R.string.v2xlog_desc), + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + + SectionCard { + Spacer(Modifier.height(8.dp)) + Text( + stringResource(if (isEsp32) R.string.v2xlog_mode_esp32 else R.string.v2xlog_mode_cit), + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + // Said up front, not discovered later as two rows of zeros: the firmware decides this. + if (isEsp32) { + Spacer(Modifier.height(4.dp)) + Text( + stringResource(R.string.v2xlog_esp32_firmware_note), + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + } + Spacer(Modifier.height(12.dp)) + Button( + onClick = if (state.active) onStop else onStart, + modifier = Modifier.fillMaxWidth(), + colors = if (state.active) ButtonDefaults.buttonColors(containerColor = Color(0xFFB3261E)) + else ButtonDefaults.buttonColors(), + ) { + Icon( + if (state.active) Icons.Default.Stop else Icons.Default.FiberManualRecord, + contentDescription = null, + ) + Spacer(Modifier.padding(horizontal = 4.dp)) + Text(stringResource(if (state.active) R.string.v2xlog_stop else R.string.v2xlog_start)) + } + state.error?.let { + Spacer(Modifier.height(8.dp)) + Text( + stringResource(R.string.v2xlog_error, it), + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.error, + ) + } + Spacer(Modifier.height(12.dp)) + } + + if (state.active) { + SectionCard { + DisabledRow(stringResource(R.string.v2xlog_file), state.fileName.orEmpty()) + Divider() + DisabledRow( + stringResource(R.string.v2xlog_total), + stringResource(R.string.v2xlog_total_val, state.total, formatBytes(state.fileBytes)), + ) + ItsMessageKind.TRACKED.forEach { kind -> + Divider() + DisabledRow(kind.label, (state.counts[kind] ?: 0).toString()) + } + Divider() + DisabledRow(stringResource(R.string.v2xlog_phone_gnss), fixAge(state.phoneFixAgeMs)) + Divider() + DisabledRow( + stringResource(R.string.v2xlog_obu_gnss), + if (isEsp32) stringResource(R.string.v2xlog_not_available) else fixAge(state.obuFixAgeMs), + ) + } + if (state.total == 0) { + Text( + stringResource(R.string.v2xlog_waiting), + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + } + } + + Text( + stringResource(R.string.v2xlog_files_title, files.size), + style = MaterialTheme.typography.labelSmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + if (files.isEmpty()) { + Text( + stringResource(R.string.v2xlog_none), + style = MaterialTheme.typography.bodySmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + } else { + SectionCard { + files.forEachIndexed { i, log -> + if (i > 0) Divider() + LogFileRow( + log = log, + onShare = { shareLog(context, log) }, + onDelete = { pendingDelete = log }, + ) + } + } + } + } +} + +@Composable +private fun LogFileRow(log: V2xLogFile, onShare: () -> Unit, onDelete: () -> Unit) { + Row( + modifier = Modifier.fillMaxWidth().padding(vertical = 4.dp), + verticalAlignment = Alignment.CenterVertically, + horizontalArrangement = Arrangement.spacedBy(4.dp), + ) { + Column(modifier = Modifier.weight(1f)) { + Text( + log.name, + style = MaterialTheme.typography.bodySmall, + fontFamily = FontFamily.Monospace, + fontWeight = FontWeight.SemiBold, + ) + Text( + "${modifiedFormat.format(Date(log.modifiedMs))} · ${formatBytes(log.sizeBytes)}", + style = MaterialTheme.typography.labelSmall, + color = MaterialTheme.colorScheme.onSurfaceVariant, + ) + } + // The file being written is open and incomplete; sharing or deleting it would hand over a + // half-written database or pull it out from under the writer. + IconButton(onClick = onShare, enabled = !log.active) { + Icon(Icons.Default.Share, contentDescription = stringResource(R.string.v2xlog_share_cd)) + } + IconButton(onClick = onDelete, enabled = !log.active) { + Icon( + Icons.Default.Delete, + contentDescription = stringResource(R.string.v2xlog_delete_cd), + tint = if (log.active) LocalContentColorDisabled() else MaterialTheme.colorScheme.error, + ) + } + } +} + +@Composable +private fun LocalContentColorDisabled(): Color = MaterialTheme.colorScheme.onSurface.copy(alpha = 0.38f) + +@Composable +private fun fixAge(ageMs: Long?): String = + if (ageMs == null) stringResource(R.string.v2xlog_no_fix) + else stringResource(R.string.v2xlog_fix_age, (ageMs / 1000).toInt()) + +private fun formatBytes(bytes: Long): String = when { + bytes < 1024 -> "$bytes B" + bytes < 1024 * 1024 -> "%.1f kB".format(Locale.US, bytes / 1024.0) + else -> "%.1f MB".format(Locale.US, bytes / (1024.0 * 1024.0)) +} + +private fun shareLog(context: Context, log: V2xLogFile) { + val uri = FileProvider.getUriForFile(context, "${context.packageName}.fileprovider", log.file) + val intent = Intent(Intent.ACTION_SEND).apply { + type = "application/vnd.sqlite3" + putExtra(Intent.EXTRA_STREAM, uri) + putExtra(Intent.EXTRA_SUBJECT, "MicrOBU V2X log - ${log.name}") + addFlags(Intent.FLAG_GRANT_READ_URI_PERMISSION) + } + context.startActivity(Intent.createChooser(intent, context.getString(R.string.v2xlog_share_chooser))) +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/viewmodel/V2xLogViewModel.kt b/app/src/main/java/com/hawhamburg/micr0bu/viewmodel/V2xLogViewModel.kt new file mode 100644 index 0000000..18feaf2 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/viewmodel/V2xLogViewModel.kt @@ -0,0 +1,45 @@ +package com.hawhamburg.micr0bu.viewmodel + +import android.app.Application +import androidx.lifecycle.AndroidViewModel +import androidx.lifecycle.viewModelScope +import com.hawhamburg.micr0bu.data.log.V2xLogFile +import com.hawhamburg.micr0bu.data.log.V2xLogState +import com.hawhamburg.micr0bu.data.log.V2xLogger +import com.hawhamburg.micr0bu.service.V2xLogService +import dagger.hilt.android.lifecycle.HiltViewModel +import kotlinx.coroutines.delay +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.launch +import javax.inject.Inject + +/** Backs Settings > Developer > V2X data logger. */ +@HiltViewModel +class V2xLogViewModel @Inject constructor( + application: Application, + private val logger: V2xLogger, +) : AndroidViewModel(application) { + + val state: StateFlow = logger.state + val files: StateFlow> = logger.files + + init { + logger.refreshFiles() + // The running file grows between state updates, so keep the list's sizes current while a + // log is on screen. + viewModelScope.launch { + while (true) { + delay(2_000L) + if (logger.state.value.active) logger.refreshFiles() + } + } + } + + fun start() = V2xLogService.start(getApplication()) + + fun stop() = V2xLogService.stop(getApplication()) + + fun delete(log: V2xLogFile) { + logger.delete(log) + } +} diff --git a/app/src/main/res/values-de/strings.xml b/app/src/main/res/values-de/strings.xml index 68e719a..6edeb00 100644 --- a/app/src/main/res/values-de/strings.xml +++ b/app/src/main/res/values-de/strings.xml @@ -332,4 +332,28 @@ Signieren %1$s · Tickets %2$d · signiert %3$d · abgelehnt %4$d · gesendet %5$d an aus + V2X-Datenlogger + Zeichnet jede empfangene V2X-Nachricht (CAM, DENM, SPATEM, MAPEM, VAM) in einer Datenbankdatei auf, mit Zeitstempel, der GNSS-Position des Telefons und, wenn vorhanden, der GNSS-Position der OBU und der Signalstärke (RSSI). + ESP32-C5: protokolliert, was das Board über die Verbindung weiterleitet, mit RSSI. Das Board hat kein GNSS, daher gibt es keine OBU-Position. + CiT One: protokolliert die MQTT-Topics. Die API liefert keine RSSI; die OBU-Position stammt aus v2x/rx/obu_gnss. + Die aktuelle ESP32-C5-Firmware leitet nur CAM, DENM und SPATEM weiter; MAPEM und VAM erscheinen erst, wenn die Firmware sie annimmt. + Aufzeichnung starten + Aufzeichnung beenden + Aufzeichnung fehlgeschlagen: %1$s + Datei + Aufgezeichnet + %1$d Nachrichten · %2$s + Telefon-GNSS + OBU-GNSS + Fix vor %1$d s + kein Fix + nicht verfügbar + Warte auf Nachrichten. Ist die OBU noch nicht verbunden, verbinde sie auf dem Verbindungsbildschirm. + Logdateien (%1$d) + Noch keine Logs + Log teilen + Log löschen + V2X-Log teilen + Log löschen? + %1$s wird dauerhaft vom Gerät entfernt. diff --git a/app/src/main/res/values/strings.xml b/app/src/main/res/values/strings.xml index 24ae426..2b0930e 100644 --- a/app/src/main/res/values/strings.xml +++ b/app/src/main/res/values/strings.xml @@ -339,4 +339,28 @@ Signing %1$s · tickets %2$d · signed %3$d · refused %4$d · on air %5$d on off + V2X data logger + Records every V2X message received (CAM, DENM, SPATEM, MAPEM, VAM) into a database file, with a timestamp, the phone GNSS fix and, where available, the OBU GNSS fix and the signal strength (RSSI). + ESP32-C5: logging what the board forwards over the link, with RSSI. The board has no GNSS, so there is no OBU position. + CiT One: logging the MQTT topics. The API carries no RSSI; the OBU position comes from v2x/rx/obu_gnss. + The current ESP32-C5 firmware forwards only CAM, DENM and SPATEM, so MAPEM and VAM appear only once the firmware accepts them. + Start logging + Stop logging + Logging failed: %1$s + File + Recorded + %1$d messages · %2$s + Phone GNSS + OBU GNSS + fix %1$d s ago + no fix + not available + Waiting for messages. If the OBU is not connected yet, connect it on the Connection screen. + Log files (%1$d) + No logs yet + Share log + Delete log + Share V2X log + Delete log? + This permanently removes %1$s from the device. diff --git a/app/src/main/res/xml/file_paths.xml b/app/src/main/res/xml/file_paths.xml index 89c386c..56cd5db 100644 --- a/app/src/main/res/xml/file_paths.xml +++ b/app/src/main/res/xml/file_paths.xml @@ -2,4 +2,6 @@ + + diff --git a/app/src/test/java/com/hawhamburg/micr0bu/ItsMessageKindTest.kt b/app/src/test/java/com/hawhamburg/micr0bu/ItsMessageKindTest.kt new file mode 100644 index 0000000..23737fd --- /dev/null +++ b/app/src/test/java/com/hawhamburg/micr0bu/ItsMessageKindTest.kt @@ -0,0 +1,73 @@ +package com.hawhamburg.micr0bu + +import com.hawhamburg.micr0bu.domain.log.ItsMessageKind +import com.hawhamburg.micr0bu.domain.log.ItsPduHeader +import com.hawhamburg.micr0bu.domain.log.classifyItsPdu +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNull +import org.junit.Test + +/** + * Pins how the V2X data logger names a message it may have no decoder for. The CAM fixture is the + * golden UPER frame from [CamEncodeGoldenTest] (protocolVersion 2, messageID 2, stationID 999999); + * the others only need a valid ItsPduHeader, which is all the classifier reads. + */ +class ItsMessageKindTest { + + private fun String.hexToBytes() = chunked(2).map { it.toInt(16).toByte() }.toByteArray() + + private val goldenCam = + "0202000f423f3700402ab215af6e286477dffffffc23b7743e0027ffc0d0fe0118329337feebfff6000000".hexToBytes() + + /** An ItsPduHeader with [messageId] and stationID 0x01020304, then a stand-in body. */ + private fun pdu(messageId: Int) = + byteArrayOf(2, messageId.toByte(), 1, 2, 3, 4, 0x55, 0x66) + + @Test fun `header of the golden CAM is read`() { + val h = ItsPduHeader.parse(goldenCam)!! + assertEquals(2, h.protocolVersion) + assertEquals(2, h.messageId) + assertEquals(999_999L, h.stationId) + } + + @Test fun `station id above 2^31 is not sign-extended`() { + val h = ItsPduHeader.parse(byteArrayOf(2, 2, 0xFF.toByte(), 0xFF.toByte(), 0xFF.toByte(), 0xFE.toByte()))!! + assertEquals(0xFFFFFFFEL, h.stationId) + } + + @Test fun `every tracked kind is recognised by messageID alone`() { + assertEquals(ItsMessageKind.DENM, classifyItsPdu(pdu(1), null).kind) + assertEquals(ItsMessageKind.CAM, classifyItsPdu(pdu(2), null).kind) + assertEquals(ItsMessageKind.SPATEM, classifyItsPdu(pdu(4), null).kind) + assertEquals(ItsMessageKind.MAPEM, classifyItsPdu(pdu(5), null).kind) + assertEquals(ItsMessageKind.VAM, classifyItsPdu(pdu(16), null).kind) + } + + @Test fun `messageID wins over the BTP port, and the SPATEM-MAPEM numbers are not swapped`() { + // SPATEM is port 2004 but messageID 4; MAPEM is port 2003 but messageID 5. + assertEquals(ItsMessageKind.SPATEM, classifyItsPdu(pdu(4), 2004).kind) + assertEquals(ItsMessageKind.MAPEM, classifyItsPdu(pdu(5), 2003).kind) + // A PDU that disagrees with the port it arrived on is what its header says it is. + assertEquals(ItsMessageKind.MAPEM, classifyItsPdu(pdu(5), 2004).kind) + } + + @Test fun `a PDU too short for a header falls back to the BTP port`() { + val c = classifyItsPdu(byteArrayOf(2, 5), 2018) + assertEquals(ItsMessageKind.VAM, c.kind) + assertNull(c.messageId) + assertNull(c.stationId) + } + + @Test fun `an unknown messageID is kept as OTHER with its id`() { + val c = classifyItsPdu(pdu(99), null) + assertEquals(ItsMessageKind.OTHER, c.kind) + assertEquals(99, c.messageId) + assertEquals(0x01020304L, c.stationId) + } + + @Test fun `processed topic names map to kinds`() { + assertEquals(ItsMessageKind.MAPEM, ItsMessageKind.fromProcessedTopicName("map")) + assertEquals(ItsMessageKind.SPATEM, ItsMessageKind.fromProcessedTopicName("spat")) + assertNull(ItsMessageKind.fromProcessedTopicName("unknown")) + } +}