From 034ef223366214c356f2b56796c4ba792ffb8f4e Mon Sep 17 00:00:00 2001 From: Ashin Walpola Date: Wed, 2 Sep 2026 15:25:31 +0200 Subject: [PATCH] Decode raw v2x/rx on the CiT One path, and stop tracking our own CAM pings The Use Case app's v2x-uca/output/json topics are a rate-limited and lossy view: traffic the OBU's radio actually heard, the ESP32's CAM pinger among it, never reached the app. The raw v2x/rx topics carry everything, as RecvV2XMessage protobuf with the ITS-G5 PDU in one bytes field (CI-CiT MQTT API section 2.4). RecvV2xMessage is a minimal protobuf wire-format reader for the three fields needed: btpHeader type and destination port, the GeoNetworking destination-area radius, and the payload. Hand-written for the same reason the ASN.1 codecs are, rather than adding protoc and the protobuf Gradle plugin and vendoring a third-party .proto into this repository. Field numbers are pinned by a byte fixture written out by hand from the encoding rules, not generated by our own encoder. Raw payloads now travel as bytes rather than String. The previous UTF-8 round trip replaced every byte that is not valid UTF-8, leaving a payload that still looked plausible in a log and decoded to nothing. CAM, DENM and SPATEM from both transports now meet in shared handlers, so everything downstream is transport-agnostic. SPATEM works on the CiT One path for the first time, and DENM gains its relevance radius there. Where both sources describe the same event the decoded one wins: remote CAMs from the processed topic are suppressed while the raw topic is live, and DENMs dedup on ETSI's actionID with the decoded list last. The processed topics stay subscribed as a fallback for an OBU whose configuration does not publish the raw ones. Two defects found while testing this: CamPinger transmits under a fixed bench station id, deliberately distinct from the persisted one, but the self-heard filter only knew the persisted id. Every ping therefore came back through the ESP32's promiscuous receive as a remote road user sitting exactly on top of the ego position, moving at the ego's own speed and heading, and was handed to the detection engine as a collision partner for itself. The rule now lives in OwnStationIds, covers both ids, and has tests, so a third transmit path cannot reintroduce the same gap quietly. Self-heard frames are now counted and reported on the pinger card instead of being discarded. That round trip is the only direct evidence the serial link, the ESP32's transmit path and its receive path all work, which is what the bench pinger exists to demonstrate. Also: the stationType warning banner no longer shows in ESP32-C5 mode. It reads a value from the CiT One's obu_gnss topic, which that hardware never publishes, so it stayed on screen reporting on an OBU that was no longer in use. --- .../micr0bu/data/cam/CamUseCaseRepository.kt | 201 ++++++++++++--- .../micr0bu/data/mqtt/MqttRepository.kt | 70 ++++- .../micr0bu/data/mqtt/RecvV2xMessage.kt | 240 ++++++++++++++++++ .../micr0bu/domain/cam/OwnStationIds.kt | 54 ++++ .../micr0bu/domain/cam/OwnTxLoopback.kt | 30 +++ .../hawhamburg/micr0bu/service/CamPinger.kt | 14 +- .../ui/screens/MqttTopicViewerScreen.kt | 23 +- .../micr0bu/viewmodel/MqttViewModel.kt | 71 ++++-- app/src/main/res/values-de/strings.xml | 2 + app/src/main/res/values/strings.xml | 2 + .../hawhamburg/micr0bu/OwnStationIdsTest.kt | 66 +++++ .../hawhamburg/micr0bu/RecvV2xMessageTest.kt | 151 +++++++++++ 12 files changed, 868 insertions(+), 56 deletions(-) create mode 100644 app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/RecvV2xMessage.kt create mode 100644 app/src/main/java/com/hawhamburg/micr0bu/domain/cam/OwnStationIds.kt create mode 100644 app/src/main/java/com/hawhamburg/micr0bu/domain/cam/OwnTxLoopback.kt create mode 100644 app/src/test/java/com/hawhamburg/micr0bu/OwnStationIdsTest.kt create mode 100644 app/src/test/java/com/hawhamburg/micr0bu/RecvV2xMessageTest.kt diff --git a/app/src/main/java/com/hawhamburg/micr0bu/data/cam/CamUseCaseRepository.kt b/app/src/main/java/com/hawhamburg/micr0bu/data/cam/CamUseCaseRepository.kt index bc70ec6..cb17660 100644 --- a/app/src/main/java/com/hawhamburg/micr0bu/data/cam/CamUseCaseRepository.kt +++ b/app/src/main/java/com/hawhamburg/micr0bu/data/cam/CamUseCaseRepository.kt @@ -7,6 +7,10 @@ import com.hawhamburg.micr0bu.data.SensorRepository import com.hawhamburg.micr0bu.data.mqtt.MqttConnectionState import com.hawhamburg.micr0bu.data.mqtt.MqttRepository import com.hawhamburg.micr0bu.data.mqtt.ObuHardwarePreferences +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_SPATEM_TOPIC +import com.hawhamburg.micr0bu.data.mqtt.RecvV2xMessage import com.hawhamburg.micr0bu.data.mqtt.UseCaseAlertPreferences import com.hawhamburg.micr0bu.data.transport.ObuHardware import com.hawhamburg.micr0bu.data.transport.BtpPort @@ -20,6 +24,8 @@ import com.hawhamburg.micr0bu.domain.asn1.SpatemUperCodec import com.hawhamburg.micr0bu.domain.cam.Cam import com.hawhamburg.micr0bu.domain.cam.CamParser import com.hawhamburg.micr0bu.domain.cam.ObuGnssParser +import com.hawhamburg.micr0bu.domain.cam.OwnStationIds +import com.hawhamburg.micr0bu.domain.cam.OwnTxLoopback import com.hawhamburg.micr0bu.domain.cam.StationType import com.hawhamburg.micr0bu.domain.denm.DenmEvent import com.hawhamburg.micr0bu.domain.spat.SpatEvent @@ -40,6 +46,7 @@ import kotlinx.coroutines.flow.asSharedFlow import kotlinx.coroutines.flow.asStateFlow import kotlinx.coroutines.flow.combine import kotlinx.coroutines.flow.stateIn +import kotlinx.coroutines.flow.update import kotlinx.coroutines.launch import javax.inject.Inject import javax.inject.Singleton @@ -55,6 +62,16 @@ private const val PRUNE_INTERVAL_MS = 1_000L // missed updates, not just normal jitter between samples. private const val OBU_GNSS_STALE_MS = 2_500L +/** + * How long a raw `v2x/rx/cam` message keeps the Use Case app's CAM topic suppressed. + * + * The two topics carry the same traffic, but `v2x-uca/output/json/cam` is rate-limited and drops + * messages, so while the raw topic is arriving there is nothing the processed one can add. A few + * seconds is many missed repetitions at CAM rates, so this only lapses if the raw topic really + * has stopped, which is what makes the fallback automatic on an OBU that does not publish it. + */ +private const val RAW_PREFERRED_WINDOW_MS = 5_000L + /** * Bridges the raw MQTT CAM stream (plus the ego's own obu_gnss/phone GNSS state) to * [UseCaseDetectionEngine] and exposes the resulting CAM-based Use Case Alerts to the UI @@ -72,7 +89,12 @@ private const val OBU_GNSS_STALE_MS = 2_500L * A singleton so detection keeps running (and alert state survives) even while no screen is * collecting it — same rationale as [MqttRepository]'s per-topic message log. * - * DENM is decoded from the ESP32-C5 serial path (see [airDenm]) but deliberately kept out of + * **Two decode sources, one funnel.** UPER arrives either from the ESP32-C5 serial link or, on + * the CiT One path, from the raw `v2x/rx` protobuf topics ([RecvV2xMessage]). Both end up in + * the same handlers, so everything downstream is transport-agnostic. The CiT One's processed + * `v2x-uca/output/json` topics remain a fallback for an OBU that does not publish the raw ones. + * + * DENM is decoded from both (see [decodedDenm]) but deliberately kept out of * [UseCaseDetectionEngine] — that engine reasons about moving road users from CAM kinematics. */ @Singleton @@ -98,6 +120,9 @@ class CamUseCaseRepository @Inject constructor( @Volatile private var lastOwnStationType: Int = StationType.CYCLIST @Volatile private var lastObuGnssTimestamp: Long = 0L + /** When a raw `v2x/rx/cam` message last arrived, for [rawCamPreferred]. */ + @Volatile private var lastRawCamMs: Long = 0L + /** Per-use-case enable/disable toggles (Settings > Use Case Alerts). */ val enabledMap: StateFlow> = prefs.enabledMapFlow.stateIn( scope, SharingStarted.Eagerly, UseCaseType.entries.associateWith { true }, @@ -140,21 +165,42 @@ class CamUseCaseRepository @Inject constructor( */ val rsuStations: StateFlow> = _rsuStations.asStateFlow() - private val _airSpat = MutableSharedFlow(replay = 16, extraBufferCapacity = 32) + private val _ownTxLoopback = MutableStateFlow(null) /** - * SPATEMs decoded from over-the-air traffic on the ESP32-C5 path. Replayed so a screen opened - * mid-stream sees the current signal state immediately rather than waiting up to half a second - * for the next repetition. + * Our own transmissions heard back off the air, or null until one is. + * + * These frames are dropped from the detection engine, correctly, since the phone is not a + * road user to itself. But dropping them silently threw away the one thing that proves the + * whole radio loop works: the frame went out over serial, the ESP32 transmitted it, and the + * ESP32 received it again. That is precisely what the bench pinger exists to demonstrate, so + * it is counted here and reported rather than discarded. + * + * ESP32-C5 path in practice. The CiT One does not normally hear its own transmissions. */ - val airSpat: SharedFlow = _airSpat.asSharedFlow() + val ownTxLoopback: StateFlow = _ownTxLoopback.asStateFlow() - private val _airDenm = MutableSharedFlow(replay = 32, extraBufferCapacity = 32) + /** Clears the loopback tally. Called when a fresh pinger run starts, so the count is per run. */ + fun resetOwnTxLoopback() { _ownTxLoopback.value = null } + + private val _decodedSpat = MutableSharedFlow(replay = 16, extraBufferCapacity = 32) /** - * DENMs decoded from over-the-air traffic on the ESP32-C5 path. `replay` so a screen opened - * after a hazard was first heard still sees it - DENMs repeat at ~1 Hz but a subscriber that - * missed the last repetition shouldn't have to wait for the next. + * SPATEMs decoded from UPER, from either hardware path: the ESP32-C5 serial link or the CiT + * One's `v2x/rx/spatem` topic. Replayed so a screen opened mid-stream sees the current signal + * state immediately rather than waiting up to half a second for the next repetition. + * + * The CiT One's own `v2x-uca/output/json/spat` topic is not a source here. It was never + * parsed, so before the raw topic was wired up this path produced no signal state at all. */ - val airDenm: SharedFlow = _airDenm.asSharedFlow() + val decodedSpat: SharedFlow = _decodedSpat.asSharedFlow() + + private val _decodedDenm = MutableSharedFlow(replay = 32, extraBufferCapacity = 32) + /** + * DENMs decoded from UPER, from either hardware path: the ESP32-C5 serial link or the CiT + * One's `v2x/rx/denm` topic. `replay` so a screen opened after a hazard was first heard still + * sees it - DENMs repeat at ~1 Hz but a subscriber that missed the last repetition shouldn't + * have to wait for the next. + */ + val decodedDenm: SharedFlow = _decodedDenm.asSharedFlow() init { scope.launch { @@ -166,6 +212,37 @@ class CamUseCaseRepository @Inject constructor( } } + // CiT One raw path: every message the OBU's radio heard, as protobuf, decoded here with + // the same codecs the serial path uses. This is what makes the CiT One see traffic the + // Use Case app filtered out, the ESP32-C5's CAM pinger among it, and it is the only + // source of SPATEM on this hardware. + scope.launch { + mqttRepository.rawV2x.collect { raw -> + val envelope = RecvV2xMessage.parse(raw.bytes) + if (envelope == null) { + Log.w(TAG, "rawV2x: unparseable RecvV2XMessage on ${raw.topic}, " + + "${raw.bytes.size} bytes - first bytes: ${raw.bytes.toHexPreview()}") + return@collect + } + when (raw.topic) { + RAW_CAM_TOPIC -> { + lastRawCamMs = raw.timestamp + handleCamUper(envelope.payload, rssiDbm = null, source = "mqtt") + } + // The GeoBroadcast radius comes off the GeoNetworking header the same way it + // does on the serial path, so a hazard's relevance area survives here too. + RAW_DENM_TOPIC -> handleDenmUper( + uper = envelope.payload, + rssiDbm = null, + relevanceRadiusM = envelope.destAreaRadiusM, + source = "mqtt", + ) + RAW_SPATEM_TOPIC -> handleSpatUper(envelope.payload, rssiDbm = null, source = "mqtt") + else -> Log.w(TAG, "rawV2x: unexpected topic ${raw.topic}") + } + } + } + // Phone GNSS fallback — only applied when obu_gnss has gone stale (see class KDoc). // Retries in a loop: this singleton can be created before the user grants location // permission (requested at app startup), so a single subscription attempt isn't @@ -257,8 +334,16 @@ class CamUseCaseRepository @Inject constructor( scope.launch { prefs.setEnabled(type, enabled) } } - /** True if [stationId] matches the ego OBU's own station ID (for OWN/REMOTE UI badges). */ - fun isOwnStationId(stationId: Long): Boolean = stationId != 0L && stationId == _ownStationId.value + /** + * True if [stationId] is one this phone transmits under, so a frame heard back off the air is + * recognised as our own rather than tracked as another road user. Also drives the OWN/REMOTE + * badges in the raw message list. + * + * The rule lives in [OwnStationIds], which explains why there are two such ids and what goes + * wrong when only one of them is checked. + */ + fun isOwnStationId(stationId: Long): Boolean = + OwnStationIds.isOwn(stationId, _ownStationId.value) /** * Primary ego state source: `v2x/rx/obu_gnss`, ~4 Hz, carries position/speed/heading/yaw @@ -301,16 +386,30 @@ class CamUseCaseRepository @Inject constructor( _processedCam.tryEmit(ego) } + /** True while `v2x/rx/cam` is arriving, in which case the processed CAM topic adds nothing. */ + private fun rawCamPreferred(now: Long): Boolean = + lastRawCamMs != 0L && now - lastRawCamMs <= RAW_PREFERRED_WINDOW_MS + private fun handleCam(payload: String, timestamp: Long) { val cam = CamParser.parse(payload, _ownStationId.value, timestamp) ?: return + // A bench ping the CiT One's radio picked up and relayed here. Not a road user, and not + // ego state either: the ping is built from the same phone GNSS the engine already has. + if (cam.stationId == OwnStationIds.BENCH_PING) return if (cam.isOwn) { - // Third fallback — the CAM topic's own low-rate entry. onOwnCam() keeps whichever + // Third fallback - the CAM topic's own low-rate entry. onOwnCam() keeps whichever // update is freshest, so this only actually wins when both obu_gnss and phone GNSS - // are unavailable/stale. + // are unavailable/stale. Deliberately still processed while the raw topic is live: + // v2x/rx/cam is a receive topic and never carries the ego station's own CAM, so + // suppressing this would remove the fallback without anything replacing it. engine.onOwnCam(cam) - } else { - engine.onRemoteCam(cam) + _processedCam.tryEmit(cam) + return } + // A remote CAM the raw topic has already delivered, in fuller form and without the Use + // Case app's rate limiting. Dropping it here rather than letting both reach the engine + // keeps one station from being fed by two sources at two different rates. + if (rawCamPreferred(timestamp)) return + engine.onRemoteCam(cam) _processedCam.tryEmit(cam) } @@ -325,7 +424,12 @@ class CamUseCaseRepository @Inject constructor( * its own just-transmitted frame (promiscuous capture of a local TX). Guarded the same way * the MQTT path guards against reprocessing "own" CAM: compare against [_ownStationId]. */ - private fun handleCamFromSerial(v2x: V2xRxFrame) { + private fun handleCamFromSerial(v2x: V2xRxFrame) = + handleCamUper(v2x.uper, v2x.rssiDbm, source = "serial") + + /** Shared by both transports: [rssiDbm] is null on the MQTT path, which does not report it. */ + private fun handleCamUper(uper: ByteArray, rssiDbm: Int?, source: String) { + val v2x = UperSource(uper, rssiDbm, source) val cam = camCodec.decodeCam(v2x.uper, System.currentTimeMillis())?.copy(rssiDbm = v2x.rssiDbm) if (cam == null) { // Logged, not silently dropped: "the app shows nothing" has two completely different @@ -333,14 +437,27 @@ class CamUseCaseRepository @Inject constructor( // without this line they're indistinguishable from the outside. Log.w( TAG, - "handleCamFromSerial: decode FAILED for ${v2x.uper.size}-byte CAM " + + "handleCamUper[${v2x.source}]: decode FAILED for ${v2x.uper.size}-byte CAM " + "(rssi=${v2x.rssiDbm} dBm) - first bytes: ${v2x.uper.toHexPreview()}", ) return } - Log.d(TAG, "handleCamFromSerial: decoded station=${cam.stationId} " + + Log.d(TAG, "handleCamUper[${v2x.source}]: decoded station=${cam.stationId} " + "lat=${cam.latitude} lon=${cam.longitude} speed=${cam.speedMps} rssi=${v2x.rssiDbm} dBm") - if (_ownStationId.value != null && cam.stationId == _ownStationId.value) return // self-heard TX + if (isOwnStationId(cam.stationId)) { + // Ours, on either station id. Kept out of the engine, but counted: this is the + // round trip completing, and it is the only direct evidence the radio path works. + _ownTxLoopback.update { prev -> + OwnTxLoopback( + frames = (prev?.frames ?: 0) + 1, + // Hold the last known reading rather than overwriting it with null on a + // transport that does not report RSSI, so the figure does not blink away. + lastRssiDbm = v2x.rssiDbm ?: prev?.lastRssiDbm, + lastHeardMs = System.currentTimeMillis(), + ) + } + return + } // Roadside units are infrastructure, not road users. Their CAM carries no kinematics (see // CamUperCodec's rsuContainerHighFrequency branch), so it reaches here as a permanently @@ -367,25 +484,38 @@ class CamUseCaseRepository @Inject constructor( * reasons about moving road users from CAM kinematics, and a static hazard is a different kind * of thing. DENMs go to the map and the message list only. */ - private fun handleDenmFromSerial(v2x: V2xRxFrame) { + private fun handleDenmFromSerial(v2x: V2xRxFrame) = handleDenmUper( + uper = v2x.uper, + rssiDbm = v2x.rssiDbm, + relevanceRadiusM = v2x.geoArea?.radiusMeters, + source = "serial", + ) + + private fun handleDenmUper( + uper: ByteArray, + rssiDbm: Int?, + relevanceRadiusM: Int?, + source: String, + ) { + val v2x = UperSource(uper, rssiDbm, source) val denm = DenmUperCodec.decode( bytes = v2x.uper, receivedAtEpochMs = System.currentTimeMillis(), rssiDbm = v2x.rssiDbm, - relevanceRadiusM = v2x.geoArea?.radiusMeters, + relevanceRadiusM = relevanceRadiusM, ) if (denm == null) { Log.w( TAG, - "handleDenmFromSerial: decode FAILED for ${v2x.uper.size}-byte DENM " + + "handleDenmUper[${v2x.source}]: decode FAILED for ${v2x.uper.size}-byte DENM " + "(rssi=${v2x.rssiDbm} dBm) - first bytes: ${v2x.uper.toHexPreview()}", ) return } - Log.d(TAG, "handleDenmFromSerial: decoded station=${denm.stationId}/${denm.sequenceNumber} " + + Log.d(TAG, "handleDenmUper[${v2x.source}]: decoded station=${denm.stationId}/${denm.sequenceNumber} " + "cause=${denm.causeCode}/${denm.subCauseCode} lat=${denm.latitude} lon=${denm.longitude} " + "radius=${denm.relevanceRadiusM}m termination=${denm.isTermination} rssi=${v2x.rssiDbm} dBm") - _airDenm.tryEmit(denm) + _decodedDenm.tryEmit(denm) } /** @@ -398,7 +528,11 @@ class CamUseCaseRepository @Inject constructor( * counts it as an oversize drop, so on real road RSUs (median 555 bytes) most will not arrive * until that cap is raised. The bench trigger's ~58-byte messages are unaffected. */ - private fun handleSpatFromSerial(v2x: V2xRxFrame) { + private fun handleSpatFromSerial(v2x: V2xRxFrame) = + handleSpatUper(v2x.uper, v2x.rssiDbm, source = "serial") + + private fun handleSpatUper(uper: ByteArray, rssiDbm: Int?, source: String) { + val v2x = UperSource(uper, rssiDbm, source) val spat = SpatemUperCodec.decode( bytes = v2x.uper, receivedAtEpochMs = System.currentTimeMillis(), @@ -407,17 +541,24 @@ class CamUseCaseRepository @Inject constructor( if (spat == null) { Log.w( TAG, - "handleSpatFromSerial: decode FAILED for ${v2x.uper.size}-byte SPATEM " + + "handleSpatUper[${v2x.source}]: decode FAILED for ${v2x.uper.size}-byte SPATEM " + "(rssi=${v2x.rssiDbm} dBm) - first bytes: ${v2x.uper.toHexPreview()}", ) return } - Log.d(TAG, "handleSpatFromSerial: decoded station=${spat.stationId} " + + Log.d(TAG, "handleSpatUper[${v2x.source}]: decoded station=${spat.stationId} " + "intersections=${spat.intersections.joinToString { it.key }} " + "movements=${spat.intersections.sumOf { it.movements.size }} rssi=${v2x.rssiDbm} dBm") - _airSpat.tryEmit(spat) + _decodedSpat.tryEmit(spat) } + /** + * The bits of a received frame the decoders and their log lines need, independent of whether + * it came off the serial link or an MQTT topic. [rssiDbm] is null on the MQTT path: the + * RecvV2XMessage envelope does not carry signal strength. + */ + private data class UperSource(val uper: ByteArray, val rssiDbm: Int?, val source: String) + private fun ByteArray.toHexPreview(limit: Int = 16): String = take(limit).joinToString(" ") { "%02x".format(it) } + if (size > limit) " ..." else "" } 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 aa5b111..b5a5025 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 @@ -43,6 +43,13 @@ private val SUBSCRIBED_TOPICS = listOf( "sys/state/heartbeat", "sys/state/cellular", "v2x/rx/obu_gnss", + // Everything the radio heard, as RecvV2XMessage protobuf (API section 2.4). Preferred over + // the v2x-uca topics below, which are a rate-limited and lossy view of the same traffic. + RAW_CAM_TOPIC, + RAW_DENM_TOPIC, + RAW_SPATEM_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", "v2x-uca/output/json/denm", "v2x-uca/output/json/spat", @@ -50,6 +57,31 @@ private val SUBSCRIBED_TOPICS = listOf( "v2x-uca/output/json/cpm", ) +/** Raw received-V2X topics, carrying protobuf rather than JSON. See [RecvV2xMessage]. */ +const val RAW_CAM_TOPIC = "v2x/rx/cam" +const val RAW_DENM_TOPIC = "v2x/rx/denm" +const val RAW_SPATEM_TOPIC = "v2x/rx/spatem" + +private val RAW_V2X_TOPICS = setOf(RAW_CAM_TOPIC, RAW_DENM_TOPIC, RAW_SPATEM_TOPIC) + +/** + * A message straight off a `v2x/rx` topic, before the protobuf envelope is opened. + * + * Carried as bytes, not [MqttMessage]: that type holds a String, and putting protobuf through + * a UTF-8 round trip replaces every byte that is not valid UTF-8 with U+FFFD. The payload + * survives looking plausible in a log and decodes to nothing. + */ +data class RawV2xMqttMessage(val topic: String, val bytes: ByteArray, val timestamp: Long) { + override fun equals(other: Any?): Boolean { + if (this === other) return true + if (other !is RawV2xMqttMessage) return false + return topic == other.topic && timestamp == other.timestamp && bytes.contentEquals(other.bytes) + } + + override fun hashCode(): Int = + 31 * (31 * topic.hashCode() + timestamp.hashCode()) + bytes.contentHashCode() +} + @Singleton class MqttRepository @Inject constructor( private val prefs: MqttPreferences, @@ -69,6 +101,21 @@ class MqttRepository @Inject constructor( ) val messages: SharedFlow = _messages.asSharedFlow() + // Same buffering rationale as [_messages], with more headroom: this stream carries every CAM + // the radio hears rather than the Use Case app's thinned-out selection, which at a busy + // intersection is a considerably higher rate. + private val _rawV2x = MutableSharedFlow( + replay = 0, + extraBufferCapacity = 512, + ) + + /** + * Undecoded `v2x/rx` protobuf messages. Consumed by + * [com.hawhamburg.micr0bu.data.cam.CamUseCaseRepository], which opens the envelope and runs + * the UPER decoders over the payload, exactly as it does for the ESP32-C5 serial path. + */ + val rawV2x: SharedFlow = _rawV2x.asSharedFlow() + // Per-topic message log, kept here (singleton) so it survives even when no screen is // collecting — e.g. DENM TX messages emitted by TripRecordingService while the V2X // Monitor screen isn't open. @@ -194,6 +241,12 @@ class MqttRepository @Inject constructor( ) } + /** A one-line, printable stand-in for a binary payload, for the raw topic log. */ + private fun describeBinary(bytes: ByteArray, limit: Int = 24): String { + val hex = bytes.take(limit).joinToString(" ") { "%02x".format(it) } + return "${bytes.size} bytes protobuf: $hex" + if (bytes.size > limit) " ..." else "" + } + /** * Record a message into both the live [messages] stream (for screens currently open) * and the persistent [topicMessages] log (survives even when no screen is collecting). @@ -301,11 +354,26 @@ class MqttRepository @Inject constructor( override fun connectionLost(cause: Throwable?) { lostSignal.complete(cause) } override fun messageArrived(topic: String, message: PahoMqttMessage) { + val now = System.currentTimeMillis() + if (topic in RAW_V2X_TOPICS) { + // Binary. The bytes go to the decoders untouched; the topic log gets a hex + // preview instead, because decoding these to a String would show the operator + // a screenful of replacement characters and imply the data was corrupt. + _rawV2x.tryEmit(RawV2xMqttMessage(topic, message.payload, now)) + recordMessage( + MqttMessage( + topic = topic, + payload = describeBinary(message.payload), + timestamp = now, + ) + ) + return + } recordMessage( MqttMessage( topic = topic, payload = message.payload.toString(Charsets.UTF_8), - timestamp = System.currentTimeMillis(), + timestamp = now, ) ) } diff --git a/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/RecvV2xMessage.kt b/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/RecvV2xMessage.kt new file mode 100644 index 0000000..a1b6c25 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/RecvV2xMessage.kt @@ -0,0 +1,240 @@ +package com.hawhamburg.micr0bu.data.mqtt + +/** + * The CiT One's raw received-V2X envelope, as published on the `v2x/rx` MQTT topics. + * + * These topics carry a `RecvV2XMessage` protobuf (CI-CiT MQTT API section 2.4), not JSON: the + * ITS-G5 PDU sits in one bytes field, and the GeoNetworking and BTP headers the stack stripped + * off travel alongside it. That is the CiT One's counterpart to the ESP32-C5 path's + * [com.hawhamburg.micr0bu.data.transport.V2xRxFrame], and it exists for the same reason: the + * app decodes the UPER itself instead of accepting somebody else's summary. + * + * **Why this rather than the Use Case app's JSON.** `v2x-uca/output/json` is a processed, + * rate-limited view. It drops messages, and what it does publish has already been reduced to + * the fields the Use Case app cared about. `v2x/rx` is everything the radio actually heard. + * + * **Why a hand-written reader.** Only three of this envelope's fields are used, protobuf's wire + * format is trivial to walk, and the alternative is adding protoc and the protobuf Gradle plugin + * to an Android build plus vendoring a third-party `.proto` into this repository. The same + * argument the ASN.1 codecs in `domain/asn1/` are built on applies here. + * + * Field numbers below come from consider it's `v2x_interface.proto`, V2X RX protocol v2.4.2. + * They are wire-format constants: changing them silently mis-parses every message, so they are + * pinned by `RecvV2xMessageTest` against a byte fixture rather than left to inspection. + */ +data class RecvV2xMessage( + /** + * `btpHeader.type`, the stack's own idea of which PDU this is: DENM 1, CAM 2, SPATEM 4, + * MAPEM 5. Null when the sender omitted the header. Advisory only, since every decoder + * re-checks the messageID in the ItsPduHeader itself. + */ + val pduType: Int?, + + /** `btpHeader.destinationPort`: 2001 CAM, 2002 DENM, 2003 MAPEM, 2004 SPATEM. */ + val destinationPort: Int?, + + /** + * `gnHeader.dest.area.distA`, metres: the radius of the GeoBroadcast destination area, so + * how far the sender meant its message to apply. Only DENM normally carries one. This is the + * MQTT path's equivalent of the serial prefix's + * [com.hawhamburg.micr0bu.data.transport.V2xRxFrame.GeoArea.radiusMeters]. + */ + val destAreaRadiusM: Int?, + + /** The ITS-G5 PDU as UPER, ItsPduHeader included. Empty when the field was absent. */ + val payload: ByteArray, +) { + // Generated equals/hashCode would compare the payload array by identity, which makes two + // decodes of the same bytes unequal and quietly breaks any test or set that holds these. + override fun equals(other: Any?): Boolean { + if (this === other) return true + if (other !is RecvV2xMessage) return false + return pduType == other.pduType && + destinationPort == other.destinationPort && + destAreaRadiusM == other.destAreaRadiusM && + payload.contentEquals(other.payload) + } + + override fun hashCode(): Int { + var result = pduType ?: 0 + result = 31 * result + (destinationPort ?: 0) + result = 31 * result + (destAreaRadiusM ?: 0) + result = 31 * result + payload.contentHashCode() + return result + } + + companion object { + // RecvV2XMessage + private const val F_BTP_HEADER = 1 + private const val F_GN_HEADER = 2 + private const val F_PAYLOAD = 3 + + // BasicTransportProtocolHeader + private const val F_BTP_TYPE = 1 + private const val F_BTP_DEST_PORT = 2 + + // GeoNetworkingHeader + private const val F_GN_DEST = 8 + + // GNDestination + private const val F_DEST_AREA = 1 + + // GeoNetworkingArea + private const val F_AREA_DIST_A = 3 + + /** + * Parses an MQTT payload from a `v2x/rx` topic, or null if it is not a readable + * `RecvV2XMessage` or carries no PDU. + * + * Unknown fields are skipped rather than treated as errors, which is what protobuf + * requires and what keeps this working if consider it adds fields in a later revision. + */ + fun parse(bytes: ByteArray): RecvV2xMessage? { + var pduType: Int? = null + var destPort: Int? = null + var radius: Int? = null + var payload: ByteArray? = null + + val reader = ProtoReader(bytes) + while (reader.hasNext()) { + val tag = reader.readTag() ?: return null + when { + tag.field == F_PAYLOAD && tag.wireType == WIRE_LENGTH_DELIMITED -> + payload = reader.readBytes() ?: return null + + tag.field == F_BTP_HEADER && tag.wireType == WIRE_LENGTH_DELIMITED -> { + val sub = reader.readBytes() ?: return null + val btp = ProtoReader(sub) + while (btp.hasNext()) { + val t = btp.readTag() ?: return null + when { + t.field == F_BTP_TYPE && t.wireType == WIRE_VARINT -> + pduType = btp.readVarint()?.toInt() ?: return null + t.field == F_BTP_DEST_PORT && t.wireType == WIRE_VARINT -> + destPort = btp.readVarint()?.toInt() ?: return null + else -> if (!btp.skip(t.wireType)) return null + } + } + } + + tag.field == F_GN_HEADER && tag.wireType == WIRE_LENGTH_DELIMITED -> { + val sub = reader.readBytes() ?: return null + radius = readDestAreaRadius(sub) + } + + else -> if (!reader.skip(tag.wireType)) return null + } + } + + // A message with no payload has nothing to decode. Returning it anyway would push an + // empty byte array into the ASN.1 decoders for them to reject one layer later. + val pdu = payload ?: return null + if (pdu.isEmpty()) return null + + return RecvV2xMessage( + pduType = pduType, + destinationPort = destPort, + destAreaRadiusM = radius, + payload = pdu, + ) + } + + /** GeoNetworkingHeader.dest.area.distA, walking two levels down. Null at any break. */ + private fun readDestAreaRadius(gnHeader: ByteArray): Int? { + val dest = nestedField(gnHeader, F_GN_DEST) ?: return null + val area = nestedField(dest, F_DEST_AREA) ?: return null + + val reader = ProtoReader(area) + while (reader.hasNext()) { + val tag = reader.readTag() ?: return null + if (tag.field == F_AREA_DIST_A && tag.wireType == WIRE_VARINT) { + return reader.readVarint()?.toInt() + } + if (!reader.skip(tag.wireType)) return null + } + return null + } + + /** The bytes of the first length-delimited field numbered [field], or null. */ + private fun nestedField(bytes: ByteArray, field: Int): ByteArray? { + val reader = ProtoReader(bytes) + while (reader.hasNext()) { + val tag = reader.readTag() ?: return null + if (tag.field == field && tag.wireType == WIRE_LENGTH_DELIMITED) { + return reader.readBytes() + } + if (!reader.skip(tag.wireType)) return null + } + return null + } + } +} + +private const val WIRE_VARINT = 0 +private const val WIRE_FIXED64 = 1 +private const val WIRE_LENGTH_DELIMITED = 2 +private const val WIRE_FIXED32 = 5 + +private data class ProtoTag(val field: Int, val wireType: Int) + +/** + * A minimal protobuf wire-format reader: enough to walk a message, read varints and + * length-delimited fields, and skip everything else. + * + * Every read returns null instead of throwing on a malformed or truncated buffer. These bytes + * arrive off a network topic and a decoder that throws on bad input is a decoder that takes the + * MQTT callback thread down with it. + */ +private class ProtoReader(private val buf: ByteArray) { + private var pos = 0 + + fun hasNext(): Boolean = pos < buf.size + + fun readTag(): ProtoTag? { + val raw = readVarint() ?: return null + val field = (raw ushr 3).toInt() + val wireType = (raw and 0x7L).toInt() + if (field <= 0) return null + return ProtoTag(field, wireType) + } + + /** + * Reads a base-128 varint. Capped at ten bytes: that is the longest a 64-bit value can be, + * and without the cap a run of 0x80 bytes would walk the reader off the end of the buffer. + */ + fun readVarint(): Long? { + var result = 0L + var shift = 0 + while (shift < 64) { + if (pos >= buf.size) return null + val b = buf[pos++].toInt() + result = result or ((b and 0x7F).toLong() shl shift) + if (b and 0x80 == 0) return result + shift += 7 + } + return null + } + + fun readBytes(): ByteArray? { + val len = readVarint()?.toInt() ?: return null + if (len < 0 || pos + len > buf.size) return null + val out = buf.copyOfRange(pos, pos + len) + pos += len + return out + } + + /** Advances past a field of [wireType]. False if the type is unknown or the buffer is short. */ + fun skip(wireType: Int): Boolean = when (wireType) { + WIRE_VARINT -> readVarint() != null + WIRE_FIXED64 -> advance(8) + WIRE_LENGTH_DELIMITED -> readBytes() != null + WIRE_FIXED32 -> advance(4) + else -> false // groups (3, 4) are not used by this schema + } + + private fun advance(n: Int): Boolean { + if (pos + n > buf.size) return false + pos += n + return true + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/OwnStationIds.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/OwnStationIds.kt new file mode 100644 index 0000000..f63efbc --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/OwnStationIds.kt @@ -0,0 +1,54 @@ +package com.hawhamburg.micr0bu.domain.cam + +/** + * Which station IDs belong to this phone, and therefore must never be treated as another road + * user when a frame comes back off the air. + * + * ## Why this exists + * On the ESP32-C5 path the radio receives promiscuously, so it hears the phone's own + * transmissions. Anything that decodes received CAMs has to recognise them, or the phone tracks + * itself: a station sitting exactly on top of the ego position, moving at the ego's own speed and + * heading, handed to [com.hawhamburg.micr0bu.domain.usecase.UseCaseDetectionEngine] as a + * collision partner for itself. + * + * ## Why two IDs + * The phone transmits under two different station IDs by design: + * + * - [com.hawhamburg.micr0bu.service.CamTransmitLoop] uses the persisted per-install ID from + * `ObuHardwarePreferences.getOrCreateOwnStationId()`, which is the real identity this station + * presents to the world. + * - [com.hawhamburg.micr0bu.service.CamPinger] uses [BENCH_PING], a fixed and recognisable value, + * so manual bench pings stay identifiable in captures and cannot be confused with the + * recording-driven stream when both run at once. + * + * That second ID is the whole reason this object exists. A filter that knew only the persisted ID + * let every bench ping return as a ghost road user, which is the bug this centralises the fix + * for. Keeping the rule in one place, in a layer with no Android dependencies, is what makes it + * testable and what stops the next transmit path from reintroducing the same gap. + */ +object OwnStationIds { + + /** + * The bench pinger's station ID. Fixed rather than derived so a ping is recognisable at a + * glance in a capture or a log line. + */ + const val BENCH_PING = 999_999L + + /** + * True when [stationId] is one this phone transmits under. + * + * [persistedOwnId] is the per-install station ID, or null before it has been loaded. Station + * ID 0 is never ours: it is the "not known yet" placeholder used while the ego identity is + * still being resolved, and matching on it would swallow real traffic. + * + * [BENCH_PING] counts as ours unconditionally, not merely while the pinger is running. A + * time-windowed check would still let a frame transmitted moments before Stop arrive + * afterwards and be tracked as a stranger. The cost is that a genuine remote station using + * this ID would be ignored, which is not a real risk at a lab site and is the bargain that + * reserving a fixed ID already implies. + */ + fun isOwn(stationId: Long, persistedOwnId: Long?): Boolean { + if (stationId == 0L) return false + return stationId == persistedOwnId || stationId == BENCH_PING + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/OwnTxLoopback.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/OwnTxLoopback.kt new file mode 100644 index 0000000..a531160 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/OwnTxLoopback.kt @@ -0,0 +1,30 @@ +package com.hawhamburg.micr0bu.domain.cam + +/** + * A tally of this phone's own transmissions heard back off the air. + * + * On the ESP32-C5 path the radio receives promiscuously, so a frame the phone sent out over the + * serial link comes back through the receive path a moment later. Those frames are deliberately + * kept out of the detection engine, since the phone is not a road user to itself, but they are + * worth counting: a frame completing that round trip is direct evidence that the serial link, the + * ESP32's transmit path and its receive path all work. That is exactly what + * [com.hawhamburg.micr0bu.service.CamPinger] exists to demonstrate. + * + * Compare [frames] against the pinger's own sent count to see the loop rate. Equal numbers mean + * every ping made it out and back; a shortfall means frames are being lost on air or dropped in + * the receive chain, which is a different fault from "nothing is being sent at all". + */ +data class OwnTxLoopback( + /** How many own frames have been heard back since the tally was last reset. */ + val frames: Int, + + /** + * Signal strength of the most recent one, dBm, or null if no transport reported it. Retained + * across frames that carry no reading rather than being cleared, so the figure does not blink + * in and out on screen. + */ + val lastRssiDbm: Int?, + + /** Wall-clock ms the most recent own frame was heard back. */ + val lastHeardMs: Long, +) diff --git a/app/src/main/java/com/hawhamburg/micr0bu/service/CamPinger.kt b/app/src/main/java/com/hawhamburg/micr0bu/service/CamPinger.kt index 3d9f5ec..94b934e 100644 --- a/app/src/main/java/com/hawhamburg/micr0bu/service/CamPinger.kt +++ b/app/src/main/java/com/hawhamburg/micr0bu/service/CamPinger.kt @@ -5,6 +5,7 @@ import com.hawhamburg.micr0bu.data.GnssReading import com.hawhamburg.micr0bu.data.SensorRepository import com.hawhamburg.micr0bu.data.transport.UsbSerialTransport import com.hawhamburg.micr0bu.domain.asn1.RealAsn1UperCodec +import com.hawhamburg.micr0bu.domain.cam.OwnStationIds import com.hawhamburg.micr0bu.domain.cam.PhoneCamBuilder import dagger.hilt.android.qualifiers.ApplicationContext import kotlinx.coroutines.CoroutineScope @@ -88,7 +89,7 @@ class CamPinger @Inject constructor( val cam = PhoneCamBuilder.build( gnss = gnss, gyroZRadPerSec = latestGyroZ, - stationId = PING_STATION_ID, + stationId = OwnStationIds.BENCH_PING, longitudinalAccelMps2 = longitudinalAccel(gnss), ) val bytes = codec.encodeCam(cam) @@ -131,11 +132,10 @@ class CamPinger @Inject constructor( private const val MIN_ACCEL_DT_SEC = 0.2 private const val MAX_ACCEL_DT_SEC = 3.0 - /** - * Recognizable station id, deliberately distinct from the persisted real one - * [CamTransmitLoop] uses, so manual bench pings stay identifiable in captures and can't be - * confused with the recording-driven stream if both happen to run at once. - */ - private const val PING_STATION_ID = 999_999L + // The station id these pings go out under lives in + // [com.hawhamburg.micr0bu.domain.cam.OwnStationIds.BENCH_PING], not here. It is not a + // private detail of this class: the ESP32 hears these frames back off the air, so the + // receive path has to recognise the same value, and a second copy of it is exactly how + // the two sides would drift apart. } } diff --git a/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/MqttTopicViewerScreen.kt b/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/MqttTopicViewerScreen.kt index d414ebd..cf7e036 100644 --- a/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/MqttTopicViewerScreen.kt +++ b/app/src/main/java/com/hawhamburg/micr0bu/ui/screens/MqttTopicViewerScreen.kt @@ -134,6 +134,7 @@ fun MqttTopicViewerScreen( val camPingerActive by viewModel.camPingerActive.collectAsState() val camPingerSentCount by viewModel.camPingerSentCount.collectAsState() val camPingerHasFix by viewModel.camPingerHasFix.collectAsState() + val ownTxLoopback by viewModel.ownTxLoopback.collectAsState() val camSendFailures by viewModel.camSendFailures.collectAsState() val espLinkStatus by viewModel.espLinkStatus.collectAsState() val denmEvents by viewModel.denmEvents.collectAsState() @@ -239,6 +240,7 @@ fun MqttTopicViewerScreen( camPingerActive = camPingerActive, camPingerSentCount = camPingerSentCount, camPingerHasFix = camPingerHasFix, + ownTxLoopback = ownTxLoopback, camSendFailures = camSendFailures, espLinkStatus = espLinkStatus, ownCamPosition = ownCamPosition, @@ -279,6 +281,7 @@ private fun TopicListPane( camPingerActive: Boolean = false, camPingerSentCount: Int = 0, camPingerHasFix: Boolean = false, + ownTxLoopback: com.hawhamburg.micr0bu.domain.cam.OwnTxLoopback? = null, camSendFailures: Int = 0, espLinkStatus: EspLinkStatus? = null, ownCamPosition: com.hawhamburg.micr0bu.domain.cam.Cam? = null, @@ -322,6 +325,7 @@ private fun TopicListPane( pingerActive = camPingerActive, sentCount = camPingerSentCount, hasFix = camPingerHasFix, + loopback = ownTxLoopback, sendFailures = camSendFailures, linkStatus = espLinkStatus, onStart = onStartCamPinger, @@ -1108,6 +1112,7 @@ private fun CamPingerCard( pingerActive: Boolean, sentCount: Int, hasFix: Boolean, + loopback: com.hawhamburg.micr0bu.domain.cam.OwnTxLoopback?, sendFailures: Int, linkStatus: EspLinkStatus?, onStart: () -> Unit, @@ -1176,8 +1181,24 @@ private fun CamPingerCard( // ── Link diagnostics ────────────────────────────────────────────── // "Sent: 240" is meaningless on its own if all 240 writes failed, or if the ESP32 - // accepted them and the radio rejected every one. These two lines are the difference + // accepted them and the radio rejected every one. These lines are the difference // between a bench session that tells you something and one that doesn't. + + // The round trip closing: sent over serial, transmitted, and heard again by the same + // radio. Compared against Sent above, a shortfall separates "nothing is going out" + // from "it goes out but is not coming back". + loopback?.takeIf { it.frames > 0 }?.let { lb -> + Spacer(Modifier.height(6.dp)) + Text( + text = lb.lastRssiDbm?.let { + stringResource(R.string.mqtt_cam_pinger_loopback, lb.frames, it) + } ?: stringResource(R.string.mqtt_cam_pinger_loopback_no_rssi, lb.frames), + style = MaterialTheme.typography.labelSmall, + color = ConnectedGreen, + fontFamily = FontFamily.Monospace, + ) + } + if (sendFailures > 0) { Spacer(Modifier.height(6.dp)) Text( diff --git a/app/src/main/java/com/hawhamburg/micr0bu/viewmodel/MqttViewModel.kt b/app/src/main/java/com/hawhamburg/micr0bu/viewmodel/MqttViewModel.kt index 2e7b9ce..673a548 100644 --- a/app/src/main/java/com/hawhamburg/micr0bu/viewmodel/MqttViewModel.kt +++ b/app/src/main/java/com/hawhamburg/micr0bu/viewmodel/MqttViewModel.kt @@ -101,7 +101,24 @@ class MqttViewModel @Inject constructor( /** 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() + /** + * 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 ───────────────────────────────────────────────────────────────── @@ -127,13 +144,27 @@ class MqttViewModel @Inject constructor( val obuStationType: StateFlow = _obuStationType.asStateFlow() /** - * True when the OBU has reported a stationType other than 2 (cyclist). + * 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 = _obuStationType - .map { it != null && it != 2 } - .stateIn(viewModelScope, SharingStarted.Eagerly, false) + 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) ───────────────────────────────── @@ -141,11 +172,16 @@ class MqttViewModel @Inject constructor( * 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 path's `v2x-uca/output/json/denm` MQTT topic (parsed by - * [DenmParser]), and the ESP32-C5 path's over-the-air DENMs (GeoBroadcast, BTP port 2002, - * decoded by [com.hawhamburg.micr0bu.domain.asn1.DenmUperCodec]). Only one is ever active at a - * time since the hardware selection decides the transport, so merging costs nothing and keeps - * the UI transport-agnostic. + * 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. */ @@ -156,7 +192,7 @@ class MqttViewModel @Inject constructor( }, // 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.airDenm + 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 @@ -164,11 +200,11 @@ class MqttViewModel @Inject constructor( // 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), - ) { fromMqtt, fromAir, _ -> + ) { fromUseCaseApp, fromDecoder, _ -> val now = System.currentTimeMillis() - (fromMqtt + fromAir) + (fromUseCaseApp + fromDecoder) .filterNot { it.isTermination } // the hazard is over - stop drawing it - .associateBy { it.dedupKey } // last write wins = most recent per hazard + .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". @@ -179,15 +215,16 @@ class MqttViewModel @Inject constructor( /** * Live signal state per intersection, newest first, keyed by [IntersectionSignalState.key]. * - * ESP32-C5 path only: SPATEM arrives over the air on BTP port 2004. The CiT One path publishes - * SPATEM on its own MQTT topic in a different (protobuf-wrapped) shape, which is not wired up. + * 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.airSpat + camUseCaseRepository.decodedSpat .runningFold(emptyMap()) { acc, spat -> acc + spat.intersections.associate { i -> i.key to SpatIntersection(i, spat.stationId, spat.rssiDbm, spat.timestamp) diff --git a/app/src/main/res/values-de/strings.xml b/app/src/main/res/values-de/strings.xml index 23a2525..c857167 100644 --- a/app/src/main/res/values-de/strings.xml +++ b/app/src/main/res/values-de/strings.xml @@ -247,6 +247,8 @@ Gesendet: %1$d Schreibfehler: %1$d in Folge - CAMs erreichen den ESP32 nicht ESP32: TX-Fehler %1$d · zu groß %2$d · CRC-Fehler %3$d + Eigene Sendung empfangen: %1$d Frames · %2$d dBm + Eigene Sendung empfangen: %1$d Frames Pinger starten Pinger stoppen diff --git a/app/src/main/res/values/strings.xml b/app/src/main/res/values/strings.xml index 77ac61c..d742eed 100644 --- a/app/src/main/res/values/strings.xml +++ b/app/src/main/res/values/strings.xml @@ -260,6 +260,8 @@ Sent: %1$d Write failures: %1$d consecutive - CAMs are not reaching the ESP32 ESP32: tx fail %1$d · oversize %2$d · crc err %3$d + Own TX heard back: %1$d frames · %2$d dBm + Own TX heard back: %1$d frames Start Pinger Stop Pinger diff --git a/app/src/test/java/com/hawhamburg/micr0bu/OwnStationIdsTest.kt b/app/src/test/java/com/hawhamburg/micr0bu/OwnStationIdsTest.kt new file mode 100644 index 0000000..b0ec764 --- /dev/null +++ b/app/src/test/java/com/hawhamburg/micr0bu/OwnStationIdsTest.kt @@ -0,0 +1,66 @@ +package com.hawhamburg.micr0bu + +import com.hawhamburg.micr0bu.domain.cam.OwnStationIds +import org.junit.Assert.assertFalse +import org.junit.Assert.assertNotEquals +import org.junit.Assert.assertTrue +import org.junit.Test + +/** + * Pins the rule that decides whether a received CAM is one this phone sent. + * + * ## The bug this exists to prevent + * The phone transmits under two station IDs: the persisted per-install one used by + * `CamTransmitLoop`, and a fixed bench ID used by `CamPinger` so pings stay identifiable in + * captures. The ESP32-C5 receives promiscuously, so both come straight back off the air. + * + * The filter originally checked only the persisted ID. Every bench ping therefore returned as a + * remote road user sitting exactly on top of the ego position, moving at the ego's own speed and + * heading, and was fed to the detection engine as a collision partner for itself. Nothing failed + * loudly: the app simply raised use case alerts against itself for as long as the pinger ran. + * + * These tests are what should fail if a third transmit path is ever added without teaching this + * rule about it. + */ +class OwnStationIdsTest { + + private val persisted = 1_691_338_363L + + @Test + fun `recognises the persisted transmit id`() { + assertTrue(OwnStationIds.isOwn(persisted, persisted)) + } + + @Test + fun `recognises the bench ping id even though it is not the persisted one`() { + // The regression. The pinger's id is deliberately different, which is exactly why a + // filter written around the persisted id alone let every ping through. + assertNotEquals( + "the bench id is meant to be distinct, or this test proves nothing", + persisted, + OwnStationIds.BENCH_PING, + ) + assertTrue(OwnStationIds.isOwn(OwnStationIds.BENCH_PING, persisted)) + } + + @Test + fun `recognises the bench ping id before the persisted id has loaded`() { + // The persisted id is read asynchronously, so it can still be null while the pinger is + // already transmitting. The ping must be recognised as ours regardless. + assertTrue(OwnStationIds.isOwn(OwnStationIds.BENCH_PING, null)) + } + + @Test + fun `treats a genuine remote station as remote`() { + assertFalse(OwnStationIds.isOwn(2_741_041_966L, persisted)) + assertFalse(OwnStationIds.isOwn(2_741_041_966L, null)) + } + + @Test + fun `station id zero is never ours`() { + // 0 is the "not resolved yet" placeholder for the ego identity. Matching on it would + // swallow real traffic from any station that reported 0. + assertFalse(OwnStationIds.isOwn(0L, null)) + assertFalse(OwnStationIds.isOwn(0L, 0L)) + } +} diff --git a/app/src/test/java/com/hawhamburg/micr0bu/RecvV2xMessageTest.kt b/app/src/test/java/com/hawhamburg/micr0bu/RecvV2xMessageTest.kt new file mode 100644 index 0000000..780aa12 --- /dev/null +++ b/app/src/test/java/com/hawhamburg/micr0bu/RecvV2xMessageTest.kt @@ -0,0 +1,151 @@ +package com.hawhamburg.micr0bu + +import com.hawhamburg.micr0bu.data.mqtt.RecvV2xMessage +import com.hawhamburg.micr0bu.domain.asn1.CamUperCodec +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNotNull +import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue +import org.junit.Test + +/** + * Pins [RecvV2xMessage] to the protobuf wire format of consider it's `RecvV2XMessage` + * (`v2x_interface.proto`, V2X RX protocol v2.4.2), the envelope the CiT One publishes on its raw + * `v2x/rx` topics. + * + * ## Where the fixtures come from + * The envelope bytes are written out here by hand from the protobuf encoding rules and the field + * numbers in that `.proto`, with the derivation in the comments, so a reviewer can check them + * without running anything. They are deliberately **not** produced by an encoder in this + * repository: a fixture generated by our own code would agree with our own reader no matter how + * wrong both were, which is exactly the failure mode the ASN.1 work in this project ran into + * three times. + * + * The CAM payload inside is the golden UPER frame from [CamEncodeGoldenTest], itself verified + * against `asn1tools` and the real ETSI modules in `asn1/`. + * + * ## Why this matters + * Field numbers are wire-format constants with no self-describing names on the wire. Reading + * field 2 where the schema says field 3 does not fail loudly, it silently yields a plausible + * looking byte string that decodes to nothing. These tests are what should fail if the constants + * in [RecvV2xMessage] are ever "tidied". + */ +class RecvV2xMessageTest { + + /** + * The golden CAM UPER, 43 bytes, from [CamEncodeGoldenTest]. Its ItsPduHeader reads + * protocolVersion 2, messageID 2 (CAM), stationID 0x000f423f = 999999. + */ + private val goldenCam = + "0202000f423f3700402ab215af6e286477dffffffc23b7743e0027ffc0d0fe0118329337feebfff6000000" + + /** + * A complete `RecvV2XMessage` carrying [goldenCam], byte by byte: + * + * ``` + * 0a 05 field 1 (btpHeader), length-delimited, 5 bytes + * 08 02 field 1 (type) varint = 2, CAM + * 10 d1 0f field 2 (destinationPort) varint = 2001 + * 12 07 field 2 (gnHeader), length-delimited, 7 bytes + * 42 05 field 8 (dest), length-delimited, 5 bytes + * 0a 03 field 1 (area), length-delimited, 3 bytes + * 18 f4 03 field 3 (distA) varint = 500 metres + * 1a 2b field 3 (payload), length-delimited, 0x2b = 43 bytes + * ``` + */ + private val camEnvelope = "0a05080210d10f120742050a0318f4031a2b" + goldenCam + + private fun String.hexToBytes(): ByteArray = + chunked(2).map { it.toInt(16).toByte() }.toByteArray() + + // ---- the happy path -------------------------------------------------------------------- + + @Test + fun `parses btp header, geo radius and payload from a full envelope`() { + val msg = RecvV2xMessage.parse(camEnvelope.hexToBytes()) + assertNotNull("envelope should parse", msg) + msg!! + + assertEquals("btpHeader.type: CAM", 2, msg.pduType) + assertEquals("btpHeader.destinationPort", 2001, msg.destinationPort) + assertEquals("gnHeader.dest.area.distA, metres", 500, msg.destAreaRadiusM) + assertTrue( + "payload must be the CAM UPER byte for byte", + msg.payload.contentEquals(goldenCam.hexToBytes()), + ) + } + + @Test + fun `extracted payload is decodable UPER, not a mangled copy`() { + val msg = RecvV2xMessage.parse(camEnvelope.hexToBytes())!! + // The whole point of carrying bytes rather than a String through the MQTT layer: a UTF-8 + // round trip would replace most of these bytes and this decode would fail. + val cam = CamUperCodec.decode(msg.payload, receivedAtEpochMs = 1_787_100_000_000L) + assertNotNull("payload should decode as a CAM", cam) + assertEquals("stationID from the ItsPduHeader", 999_999L, cam!!.stationId) + } + + @Test + fun `reads a DENM envelope's relevance radius`() { + // Same shape, DENM values: type 1, port 2002, distA 1000 m, a 2-byte stand-in payload. + // 0a 05 08 01 10 d2 0f | 12 07 42 05 0a 03 18 e8 07 | 1a 02 02 01 + val msg = RecvV2xMessage.parse("0a05080110d20f120742050a0318e8071a020201".hexToBytes()) + assertNotNull(msg) + assertEquals(1, msg!!.pduType) + assertEquals(2002, msg.destinationPort) + assertEquals(1000, msg.destAreaRadiusM) + } + + // ---- forward compatibility ------------------------------------------------------------- + + @Test + fun `skips unknown fields and does not depend on field order`() { + // payload first, then an unknown varint (field 7) and an unknown fixed32 (field 6) that + // this schema revision does not define, then the btpHeader. Protobuf permits all three, + // and a reader that assumed order or choked on unknowns would break the first time + // consider it added a field. + val bytes = ("1a2b" + goldenCam + "38b96035deadbeef0a05080210d10f").hexToBytes() + val msg = RecvV2xMessage.parse(bytes) + assertNotNull(msg) + assertEquals(2, msg!!.pduType) + assertEquals(2001, msg.destinationPort) + assertTrue(msg.payload.contentEquals(goldenCam.hexToBytes())) + } + + @Test + fun `accepts an envelope carrying nothing but a payload`() { + val msg = RecvV2xMessage.parse(("1a2b" + goldenCam).hexToBytes()) + assertNotNull(msg) + assertNull("no btpHeader was sent", msg!!.pduType) + assertNull("no gnHeader was sent", msg.destAreaRadiusM) + assertTrue(msg.payload.contentEquals(goldenCam.hexToBytes())) + } + + // ---- malformed input ------------------------------------------------------------------- + // These arrive off a network topic. A reader that throws takes the MQTT callback thread with + // it, so every one of these must return null instead. + + @Test + fun `returns null for a truncated envelope`() { + val full = camEnvelope.hexToBytes() + assertNull(RecvV2xMessage.parse(full.copyOfRange(0, full.size / 2))) + } + + @Test + fun `returns null when the payload field is present but empty`() { + assertNull(RecvV2xMessage.parse("1a00".hexToBytes())) + } + + @Test + fun `returns null when there is no payload field at all`() { + assertNull(RecvV2xMessage.parse("0a05080210d10f".hexToBytes())) + } + + @Test + fun `returns null for empty input and for bytes that are not protobuf`() { + assertNull(RecvV2xMessage.parse(ByteArray(0))) + // A run of continuation bytes: a varint that never terminates, which is what would walk + // an unguarded reader off the end of the buffer. + assertNull(RecvV2xMessage.parse(ByteArray(24) { 0xFF.toByte() })) + } +}