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.
This commit is contained in:
Ashin Walpola
2026-09-02 15:25:31 +02:00
parent ebe1c9edfd
commit 034ef22336
12 changed files with 868 additions and 56 deletions
@@ -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<Map<UseCaseType, Boolean>> = prefs.enabledMapFlow.stateIn(
scope, SharingStarted.Eagerly, UseCaseType.entries.associateWith { true },
@@ -140,21 +165,42 @@ class CamUseCaseRepository @Inject constructor(
*/
val rsuStations: StateFlow<Map<Long, Cam>> = _rsuStations.asStateFlow()
private val _airSpat = MutableSharedFlow<SpatEvent>(replay = 16, extraBufferCapacity = 32)
private val _ownTxLoopback = MutableStateFlow<OwnTxLoopback?>(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<SpatEvent> = _airSpat.asSharedFlow()
val ownTxLoopback: StateFlow<OwnTxLoopback?> = _ownTxLoopback.asStateFlow()
private val _airDenm = MutableSharedFlow<DenmEvent>(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<SpatEvent>(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<DenmEvent> = _airDenm.asSharedFlow()
val decodedSpat: SharedFlow<SpatEvent> = _decodedSpat.asSharedFlow()
private val _decodedDenm = MutableSharedFlow<DenmEvent>(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<DenmEvent> = _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 ""
}
@@ -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<MqttMessage> = _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<RawV2xMqttMessage>(
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<RawV2xMqttMessage> = _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,
)
)
}
@@ -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
}
}
@@ -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
}
}
@@ -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,
)
@@ -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.
}
}
@@ -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(
@@ -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<Boolean> = 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<com.hawhamburg.micr0bu.domain.cam.OwnTxLoopback?> =
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<Int?> = _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<Boolean> = _obuStationType
.map { it != null && it != 2 }
.stateIn(viewModelScope, SharingStarted.Eagerly, false)
val obuStationTypeWarning: StateFlow<Boolean> = 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<String, DenmEvent>()) { 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<List<SpatIntersection>> = combine(
camUseCaseRepository.airSpat
camUseCaseRepository.decodedSpat
.runningFold(emptyMap<String, SpatIntersection>()) { acc, spat ->
acc + spat.intersections.associate { i ->
i.key to SpatIntersection(i, spat.stationId, spat.rssiDbm, spat.timestamp)