SPATEM receive, RSU CAM decode, and two ASN.1 encoding fixes
SPATEM over the air - gn_unwrap.c accepts BTP-B port 2004 alongside 2001/2002. The serial protocol already carries the port in its V2X_RX prefix, so nothing else changed there. Note the crossover that makes this easy to get wrong: SPATEM is port 2004 but messageID 4, while MAPEM is port 2003 and messageID 5. - SpatemUperCodec decodes SPAT down to per-signal-group phase and timing. The bit layout was validated by replaying 79,042 real SPATEMs - the whole 2026-03-18 drive across 7+ RSUs plus the bench trigger - against asn1tools using the ETSI modules. All 79,042 matched on every field, none hit an unsupported branch. Two traps are pinned by tests: TimeChangeDetails is the one SEQUENCE here that is NOT extensible (5 optional bits, no extension bit), and maneuverAssistList cannot be skipped when present - it is variable-length, so it has to be walked to find where the next movement starts. - The V2X list shows one row per intersection with each signal group coloured by phase and a countdown where the RSU supplies timing. TimeMark wraps hourly, so the countdown corrects for it; without that it reads hugely negative once an hour, precisely when someone is watching it. - Entries expire after 15 s, much shorter than DENM's window: a traffic light that stopped updating is not "still green". Size caveat, deliberately deferred: SERIAL_LINK_MAX_PAYLOAD is still 512, so a SPATEM over ~498 bytes is counted as an oversize drop. The bench RSU sends 58 bytes and is unaffected, but real road RSUs measured 555 median / 1243 max, so roughly 70% would not arrive. Raising the cap also requires enlarging RX_FRAME_MAX_LEN and moving rx_item_t off the WiFi driver's callback stack, where it would otherwise overflow. RSU CAM decode - HighFrequencyContainer is a CHOICE, and a roadside unit picks rsuContainerHighFrequency, which carries no kinematics at all. The decoder bailed on that branch, so every RSU CAM was dropped - including the bench RSU, which sends CAM and SPATEM from the same station id. It now decodes for position and stationType. - RSU CAMs are kept out of UseCaseDetectionEngine. They arrive as a permanently stationary station at a fixed point, which is exactly the shape the stopped-vehicle and intersection-movement use cases match, and would raise a standing false alert for as long as the RSU was in range. CAM transmit: yawRateConfidence - YawRateConfidence has nine enumerands (0..8), so UPER needs 4 bits and "unavailable" is 8. The encoder wrote 3 bits with value 7 - one bit short and the wrong symbol - shifting every field after yawRate for any standards-strict receiver. The decoder read 3 bits too, so phone and ESP32 agreed with each other and with nothing else. - This is the third instance of that exact failure mode in this project, after CurvatureCalculationMode and the GeoNetworking reserved bytes. A round-trip test through our own decoder structurally cannot catch it, so CamEncodeGolden Test asserts the bytes asn1tools produces instead: it decoded this encoder's output and re-encoded it byte-identically. Confirmed on air afterwards - 26 of our own CAMs captured back off the OBU's receiver, all 26 accepted, where the same decoder rejected them before. DENM - Hazards now expire 60 s after their last repetition. This needs a clock, not just a filter: both source flows only emit when a DENM arrives, so a sender that drives away or loses power would never trigger a recompute and its hazard would stay on screen indefinitely. - The MQTT path was dropping every DENM for two independent reasons, both found by checking the payload against CI-CiT-MQTT_API_Documentation-v6 listing 2.6 rather than guessing: the station id key is originatingStationId, and eventPosition IS a GeoJSON Point rather than an object containing one. Also parses termination (presence is the signal), sequenceNumber, stationType and the RFC3339 detectionTime. Note roadSideUnit is 15, not 12 - the enumeration has a gap after tram(11). V2X screen - The decoded CAM/DENM list now renders on the CiT One path too; it was gated to the ESP32-C5 path and CiT One fell through to the raw MQTT topic list. Those topics move to their own tab, hidden on the ESP32-C5 path where there is no broker. Testing - Adds org.json as a test-only dependency: the android.jar stub throws "not mocked" on every JSONObject call, which made the MQTT payload parsers untestable off-device. - 23 V2X tests pass. EventDetectorTest's 4 failures are pre-existing and untouched by this change.
This commit is contained in:
@@ -16,11 +16,13 @@ import com.hawhamburg.micr0bu.data.transport.UsbSerialState
|
||||
import com.hawhamburg.micr0bu.data.transport.UsbSerialTransport
|
||||
import com.hawhamburg.micr0bu.domain.asn1.DenmUperCodec
|
||||
import com.hawhamburg.micr0bu.domain.asn1.RealAsn1UperCodec
|
||||
import com.hawhamburg.micr0bu.domain.asn1.SpatemUperCodec
|
||||
import com.hawhamburg.micr0bu.domain.cam.Cam
|
||||
import com.hawhamburg.micr0bu.domain.cam.CamParser
|
||||
import com.hawhamburg.micr0bu.domain.cam.ObuGnssParser
|
||||
import com.hawhamburg.micr0bu.domain.cam.StationType
|
||||
import com.hawhamburg.micr0bu.domain.denm.DenmEvent
|
||||
import com.hawhamburg.micr0bu.domain.spat.SpatEvent
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseAlert
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseDetectionEngine
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseType
|
||||
@@ -126,6 +128,14 @@ class CamUseCaseRepository @Inject constructor(
|
||||
*/
|
||||
val processedCam: SharedFlow<Cam> = _processedCam.asSharedFlow()
|
||||
|
||||
private val _airSpat = MutableSharedFlow<SpatEvent>(replay = 16, extraBufferCapacity = 32)
|
||||
/**
|
||||
* 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.
|
||||
*/
|
||||
val airSpat: SharedFlow<SpatEvent> = _airSpat.asSharedFlow()
|
||||
|
||||
private val _airDenm = MutableSharedFlow<DenmEvent>(replay = 32, extraBufferCapacity = 32)
|
||||
/**
|
||||
* DENMs decoded from over-the-air traffic on the ESP32-C5 path. `replay` so a screen opened
|
||||
@@ -190,6 +200,7 @@ class CamUseCaseRepository @Inject constructor(
|
||||
when (v2x.btpPort) {
|
||||
BtpPort.CAM -> handleCamFromSerial(v2x)
|
||||
BtpPort.DENM -> handleDenmFromSerial(v2x)
|
||||
BtpPort.SPATEM -> handleSpatFromSerial(v2x)
|
||||
// The firmware only forwards ports it was told to accept, so anything else
|
||||
// means the two sides have drifted out of sync.
|
||||
else -> Log.w(TAG, "unexpected BTP port ${v2x.btpPort} from firmware")
|
||||
@@ -317,6 +328,17 @@ class CamUseCaseRepository @Inject constructor(
|
||||
Log.d(TAG, "handleCamFromSerial: 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
|
||||
|
||||
// Roadside units are infrastructure, not road users. Their CAM carries no kinematics (see
|
||||
// CamUperCodec's rsuContainerHighFrequency branch), so it reaches here as a permanently
|
||||
// stationary station at a fixed point - which is precisely the shape the stopped-vehicle
|
||||
// and intersection-movement use cases look for. Feeding it to the engine would raise a
|
||||
// standing false alert for as long as the RSU is in range.
|
||||
if (cam.stationType == StationType.ROAD_SIDE_UNIT) {
|
||||
_processedCam.tryEmit(cam)
|
||||
return
|
||||
}
|
||||
|
||||
engine.onRemoteCam(cam)
|
||||
_processedCam.tryEmit(cam)
|
||||
}
|
||||
@@ -347,6 +369,36 @@ class CamUseCaseRepository @Inject constructor(
|
||||
_airDenm.tryEmit(denm)
|
||||
}
|
||||
|
||||
/**
|
||||
* A SPATEM heard over the air: the live signal phase for one or more intersections.
|
||||
*
|
||||
* Like DENM, this is deliberately kept out of [UseCaseDetectionEngine] - a traffic light is
|
||||
* not a moving road user, and the CAM-based use cases reason about kinematics.
|
||||
*
|
||||
* Note the firmware drops any SPATEM whose UPER exceeds the 512-byte serial payload cap and
|
||||
* 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) {
|
||||
val spat = SpatemUperCodec.decode(
|
||||
bytes = v2x.uper,
|
||||
receivedAtEpochMs = System.currentTimeMillis(),
|
||||
rssiDbm = v2x.rssiDbm,
|
||||
)
|
||||
if (spat == null) {
|
||||
Log.w(
|
||||
TAG,
|
||||
"handleSpatFromSerial: 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} " +
|
||||
"intersections=${spat.intersections.joinToString { it.key }} " +
|
||||
"movements=${spat.intersections.sumOf { it.movements.size }} rssi=${v2x.rssiDbm} dBm")
|
||||
_airSpat.tryEmit(spat)
|
||||
}
|
||||
|
||||
private fun ByteArray.toHexPreview(limit: Int = 16): String =
|
||||
take(limit).joinToString(" ") { "%02x".format(it) } + if (size > limit) " ..." else ""
|
||||
}
|
||||
|
||||
@@ -99,6 +99,13 @@ object Crc16CcittFalse {
|
||||
object BtpPort {
|
||||
const val CAM = 2001
|
||||
const val DENM = 2002
|
||||
|
||||
/**
|
||||
* Watch the crossover: SPATEM is BTP port **2004** but ItsPduHeader messageID **4**, while
|
||||
* MAPEM is port 2003 and messageID 5. The two numbering schemes are unrelated, and swapping
|
||||
* them routes messages to the wrong decoder.
|
||||
*/
|
||||
const val SPATEM = 2004
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -28,6 +28,9 @@ import kotlin.math.roundToLong
|
||||
*/
|
||||
object CamUperCodec {
|
||||
|
||||
/** HighFrequencyContainer CHOICE index for rsuContainerHighFrequency. */
|
||||
private const val HF_CONTAINER_RSU = 1
|
||||
|
||||
/** Encode buffer size — matches `cam.c`'s `cam_payload[96]`, the known-sufficient size. */
|
||||
private const val ENCODE_BUFFER_BYTES = 96
|
||||
|
||||
@@ -144,7 +147,17 @@ object CamUperCodec {
|
||||
?.let { (it * 100.0).roundToInt().coerceIn(-32766, 32766) }
|
||||
?: YAW_RATE_UNAVAILABLE
|
||||
bw.putBits(yawRateCentiDegS - (-32766), 16)
|
||||
bw.putBits(7, 3) // yawRateConfidence: unavailable
|
||||
// YawRateConfidence has NINE enumerands (degSec-000-01(0) .. unavailable(8), see
|
||||
// cdd_1_3_1_1.asn), so UPER needs 4 bits and "unavailable" is 8. This wrote 3 bits with
|
||||
// value 7 - which is both one bit short and the wrong symbol (7 is outOfRange), shifting
|
||||
// every field after yawRate for any standards-compliant receiver.
|
||||
//
|
||||
// Exactly the same failure mode as the CurvatureCalculationMode note above, and caught
|
||||
// the same way it should have been the first time: asn1tools, decoding this project's own
|
||||
// transmitted CAM with the real ETSI modules, rejected it with
|
||||
// "yawRateConfidence: Expected enumeration index ...". The decoder below read 3 bits too,
|
||||
// so phone <-> ESP32 agreed with each other and with nothing else.
|
||||
bw.putBits(8, 4) // yawRateConfidence: unavailable(8)
|
||||
|
||||
// ---- LowFrequencyContainer CHOICE ---- extension(0) -> basicVehicleContainerLowFrequency
|
||||
bw.putBits(0, 1)
|
||||
@@ -217,9 +230,34 @@ object CamUperCodec {
|
||||
br.getBits(20) // altitudeValue
|
||||
br.getBits(4) // altitudeConfidence
|
||||
|
||||
// HighFrequencyContainer is an extensible CHOICE: extension bit, then a 1-bit index
|
||||
// selecting basicVehicleContainerHighFrequency(0) or rsuContainerHighFrequency(1).
|
||||
val highFreqExt = br.getBitsInt(1)
|
||||
val highFreqIndex = br.getBitsInt(1)
|
||||
if (highFreqExt != 0 || highFreqIndex != 0) return null // extension, or rsuContainerHighFrequency
|
||||
if (highFreqExt != 0) return null // an alternative added in a later revision
|
||||
|
||||
if (highFreqIndex == HF_CONTAINER_RSU) {
|
||||
// An RSU's CAM carries no kinematics at all - RSUContainerHighFrequency holds only an
|
||||
// optional protected-zone list. Everything meaningful (position, stationType) has
|
||||
// already been read from the basicContainer above, so return that rather than
|
||||
// dropping the message: an RSU is exactly the station a rider wants to see, and the
|
||||
// roadside unit at this bench sends CAM and SPATEM from the same station id.
|
||||
//
|
||||
// speed/heading are reported as zero because the model has no "unknown" for them.
|
||||
// That is safe only because RSU CAMs are kept out of UseCaseDetectionEngine - a
|
||||
// permanently stationary "vehicle" would otherwise trip the stopped-vehicle use case
|
||||
// forever. See CamUseCaseRepository.handleCamFromSerial.
|
||||
return Cam(
|
||||
stationId = stationId,
|
||||
stationType = stationType,
|
||||
latitude = latitude,
|
||||
longitude = longitude,
|
||||
speedMps = 0.0,
|
||||
headingDeg = 0.0,
|
||||
timestamp = receivedAtEpochMs,
|
||||
isOwn = false,
|
||||
)
|
||||
}
|
||||
|
||||
// Optional-presence bitmap for BasicVehicleContainerHighFrequency's 7 trailing OPTIONAL
|
||||
// fields: accelerationControl, lanePosition, steeringWheelAngle, lateralAcceleration,
|
||||
@@ -266,7 +304,7 @@ object CamUperCodec {
|
||||
br.getBits(2) // curvatureCalculationMode root index
|
||||
|
||||
val yawRateRaw = br.getBitsInt(16) + (-32766)
|
||||
br.getBits(3) // yawRateConfidence
|
||||
br.getBits(4) // yawRateConfidence: 9 enumerands -> 4 bits, see the note in encode()
|
||||
val yawRateDps = if (yawRateRaw == YAW_RATE_UNAVAILABLE) null else yawRateRaw / 100.0
|
||||
|
||||
// Everything after yawRate is deliberately left unread: the 7 optional high-frequency
|
||||
|
||||
@@ -0,0 +1,199 @@
|
||||
package com.hawhamburg.micr0bu.domain.asn1
|
||||
|
||||
import com.hawhamburg.micr0bu.domain.spat.IntersectionSignalState
|
||||
import com.hawhamburg.micr0bu.domain.spat.SignalMovement
|
||||
import com.hawhamburg.micr0bu.domain.spat.SignalPhase
|
||||
import com.hawhamburg.micr0bu.domain.spat.SignalPhaseEvent
|
||||
import com.hawhamburg.micr0bu.domain.spat.SpatEvent
|
||||
|
||||
/**
|
||||
* ASN.1 UPER **decoder** for SPATEM (ETSI TS 103 301 / SAE J2735 DSRC), for messages received over
|
||||
* the air on the ESP32-C5 path.
|
||||
*
|
||||
* Decode-only: this project never transmits SPATEM, that is an RSU's job.
|
||||
*
|
||||
* ## Field widths
|
||||
* Every width is taken from the ETSI ASN.1 modules in the `C-ITS-Parser` checkout
|
||||
* (`dsrc_2_2_1.asn`, `cdd_2_2_1.asn`):
|
||||
*
|
||||
* - `MinuteOfTheYear` (0..527040) = 20 bits
|
||||
* - `DSecond` (0..65535) = 16 bits
|
||||
* - `MsgCount` (0..127) = 7 bits
|
||||
* - `SignalGroupID` / `LaneID` / `LaneConnectionID` (0..255) = 8 bits
|
||||
* - `IntersectionID` / `RoadRegulatorID` (0..65535) = 16 bits
|
||||
* - `TimeMark` (0..36001) = 16 bits
|
||||
* - `TimeIntervalConfidence` (0..15) = 4 bits
|
||||
* - `ZoneLength` (0..10000) = 14 bits
|
||||
* - `IntersectionStatusObject` BIT STRING SIZE(16) = 16 bits
|
||||
*
|
||||
* A `SEQUENCE (SIZE(lo..hi)) OF` writes its count in `ceil(log2(hi-lo+1))` bits holding `n-lo`:
|
||||
* intersections 1..32 gives 5 bits, movements 1..255 gives 8, events 1..16 gives 4.
|
||||
*
|
||||
* Two traps worth naming, both of which silently shift every later field:
|
||||
* - `TimeChangeDetails` is **not** extensible, so it has 5 optional bits and no extension bit,
|
||||
* unlike almost every other SEQUENCE here, which all carry one.
|
||||
* - `maneuverAssistList` cannot be skipped when present. It is variable-length, so the only way
|
||||
* to reach the next movement is to walk it, even though nothing here consumes it.
|
||||
*
|
||||
* ## Verification
|
||||
* The layout was validated by replaying **79,042 real SPATEMs**, the entire 2026-03-18 drive
|
||||
* (7+ RSUs) plus the live bench trigger, through a port of this decoder and comparing every field
|
||||
* against `asn1tools` decoding the same bytes with the real ETSI modules. All 79,042 matched
|
||||
* exactly, with no message hitting an unsupported branch.
|
||||
*
|
||||
* Returns null rather than guessing whenever an extension bit is set or an unsupported optional
|
||||
* appears: a dropped SPATEM is recoverable (they repeat at ~2 Hz), a misread one shows a driver
|
||||
* the wrong light.
|
||||
*/
|
||||
object SpatemUperCodec {
|
||||
|
||||
private const val MESSAGE_ID_SPATEM = 4
|
||||
private const val PROTOCOL_VERSION = 2
|
||||
|
||||
/** Guards against a malformed length field turning into a long decode loop. */
|
||||
private const val MAX_INTERSECTIONS = 32
|
||||
private const val MAX_MOVEMENTS = 255
|
||||
|
||||
fun decode(bytes: ByteArray, receivedAtEpochMs: Long, rssiDbm: Int? = null): SpatEvent? = try {
|
||||
decodeOrNull(bytes, receivedAtEpochMs, rssiDbm)
|
||||
} catch (e: IndexOutOfBoundsException) {
|
||||
null // truncated frame
|
||||
}
|
||||
|
||||
private fun decodeOrNull(bytes: ByteArray, receivedAtEpochMs: Long, rssiDbm: Int?): SpatEvent? {
|
||||
val br = BitReader(bytes)
|
||||
|
||||
// ---- ItsPduHeader ---- not extensible, no optionals, so no preamble.
|
||||
if (br.getBitsInt(8) != PROTOCOL_VERSION) return null
|
||||
if (br.getBitsInt(8) != MESSAGE_ID_SPATEM) return null
|
||||
val stationId = br.getBits(32)
|
||||
|
||||
// ---- SPAT ---- extensible: extension bit, then timeStamp/name/regional optional bits.
|
||||
if (br.getBitsInt(1) != 0) return null
|
||||
val hasTimeStamp = br.getBitsInt(1) == 1
|
||||
val hasName = br.getBitsInt(1) == 1
|
||||
val hasRegional = br.getBitsInt(1) == 1
|
||||
val minuteOfYear = if (hasTimeStamp) br.getBitsInt(20) else null
|
||||
// DescriptiveName is a variable-length IA5String. Nothing in 79k real messages uses it,
|
||||
// and guessing its length would desynchronise everything after it.
|
||||
if (hasName) return null
|
||||
|
||||
val intersectionCount = br.getBitsInt(5) + 1
|
||||
if (intersectionCount > MAX_INTERSECTIONS) return null
|
||||
val intersections = ArrayList<IntersectionSignalState>(intersectionCount)
|
||||
repeat(intersectionCount) {
|
||||
intersections.add(readIntersection(br) ?: return null)
|
||||
}
|
||||
if (hasRegional) return null
|
||||
|
||||
return SpatEvent(
|
||||
stationId = stationId,
|
||||
minuteOfYear = minuteOfYear,
|
||||
intersections = intersections,
|
||||
rssiDbm = rssiDbm,
|
||||
timestamp = receivedAtEpochMs,
|
||||
)
|
||||
}
|
||||
|
||||
private fun readIntersection(br: BitReader): IntersectionSignalState? {
|
||||
if (br.getBitsInt(1) != 0) return null // IntersectionState extension
|
||||
val hasName = br.getBitsInt(1) == 1
|
||||
val hasMoy = br.getBitsInt(1) == 1
|
||||
val hasTimeStamp = br.getBitsInt(1) == 1
|
||||
val hasEnabledLanes = br.getBitsInt(1) == 1
|
||||
val hasManeuvers = br.getBitsInt(1) == 1
|
||||
val hasRegional = br.getBitsInt(1) == 1
|
||||
if (hasName) return null
|
||||
|
||||
// IntersectionReferenceID - not extensible, one optional bit for region.
|
||||
val region = if (br.getBitsInt(1) == 1) br.getBitsInt(16) else null
|
||||
val id = br.getBitsInt(16)
|
||||
val revision = br.getBitsInt(7)
|
||||
br.getBits(16) // IntersectionStatusObject - not surfaced yet
|
||||
val moy = if (hasMoy) br.getBitsInt(20) else null
|
||||
val timeStampMs = if (hasTimeStamp) br.getBitsInt(16) else null
|
||||
if (hasEnabledLanes) {
|
||||
repeat(br.getBitsInt(4) + 1) { br.getBits(8) } // EnabledLaneList SIZE(1..16) OF LaneID
|
||||
}
|
||||
|
||||
val movementCount = br.getBitsInt(8) + 1
|
||||
if (movementCount > MAX_MOVEMENTS) return null
|
||||
val movements = ArrayList<SignalMovement>(movementCount)
|
||||
repeat(movementCount) {
|
||||
movements.add(readMovement(br) ?: return null)
|
||||
}
|
||||
|
||||
if (hasManeuvers && !skipManeuverAssistList(br)) return null
|
||||
if (hasRegional) return null
|
||||
|
||||
return IntersectionSignalState(region, id, revision, moy, timeStampMs, movements)
|
||||
}
|
||||
|
||||
private fun readMovement(br: BitReader): SignalMovement? {
|
||||
if (br.getBitsInt(1) != 0) return null // MovementState extension
|
||||
val hasName = br.getBitsInt(1) == 1
|
||||
val hasManeuvers = br.getBitsInt(1) == 1
|
||||
val hasRegional = br.getBitsInt(1) == 1
|
||||
if (hasName) return null
|
||||
|
||||
val signalGroup = br.getBitsInt(8)
|
||||
val eventCount = br.getBitsInt(4) + 1
|
||||
val events = ArrayList<SignalPhaseEvent>(eventCount)
|
||||
repeat(eventCount) {
|
||||
if (br.getBitsInt(1) != 0) return null // MovementEvent extension
|
||||
val hasTiming = br.getBitsInt(1) == 1
|
||||
val hasSpeeds = br.getBitsInt(1) == 1
|
||||
val hasRegionalEvent = br.getBitsInt(1) == 1
|
||||
val phase = SignalPhase.fromWire(br.getBitsInt(4)) ?: return null
|
||||
|
||||
var minEnd: Int? = null
|
||||
var maxEnd: Int? = null
|
||||
var likely: Int? = null
|
||||
if (hasTiming) {
|
||||
// TimeChangeDetails: NOT extensible - 5 optional bits, no extension bit.
|
||||
val hasStart = br.getBitsInt(1) == 1
|
||||
val hasMaxEnd = br.getBitsInt(1) == 1
|
||||
val hasLikely = br.getBitsInt(1) == 1
|
||||
val hasConfidence = br.getBitsInt(1) == 1
|
||||
val hasNext = br.getBitsInt(1) == 1
|
||||
if (hasStart) br.getBits(16)
|
||||
minEnd = br.getBitsInt(16)
|
||||
if (hasMaxEnd) maxEnd = br.getBitsInt(16)
|
||||
if (hasLikely) likely = br.getBitsInt(16)
|
||||
if (hasConfidence) br.getBits(4)
|
||||
if (hasNext) br.getBits(16)
|
||||
}
|
||||
if (hasSpeeds || hasRegionalEvent) return null
|
||||
events.add(SignalPhaseEvent(phase, minEnd, maxEnd, likely))
|
||||
}
|
||||
|
||||
if (hasManeuvers && !skipManeuverAssistList(br)) return null
|
||||
if (hasRegional) return null
|
||||
return SignalMovement(signalGroup, events)
|
||||
}
|
||||
|
||||
/**
|
||||
* Walks a `ManeuverAssistList` without keeping it. Nothing consumes queue lengths yet, but the
|
||||
* list is variable-length, so it has to be parsed to find where the next field starts.
|
||||
*
|
||||
* Returns false if it contains something this decoder cannot size, in which case the whole
|
||||
* message must be abandoned - the bit position is no longer trustworthy.
|
||||
*/
|
||||
private fun skipManeuverAssistList(br: BitReader): Boolean {
|
||||
repeat(br.getBitsInt(4) + 1) { // SIZE(1..16)
|
||||
if (br.getBitsInt(1) != 0) return false // ConnectionManeuverAssist extension
|
||||
val hasQueue = br.getBitsInt(1) == 1
|
||||
val hasStorage = br.getBitsInt(1) == 1
|
||||
val hasWaitOnStop = br.getBitsInt(1) == 1
|
||||
val hasPedBicycle = br.getBitsInt(1) == 1
|
||||
val hasRegional = br.getBitsInt(1) == 1
|
||||
br.getBits(8) // connectionID
|
||||
if (hasQueue) br.getBits(14) // ZoneLength (0..10000)
|
||||
if (hasStorage) br.getBits(14)
|
||||
if (hasWaitOnStop) br.getBits(1) // BOOLEAN
|
||||
if (hasPedBicycle) br.getBits(1) // BOOLEAN
|
||||
if (hasRegional) return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
}
|
||||
@@ -8,6 +8,12 @@ package com.hawhamburg.micr0bu.domain.cam
|
||||
object StationType {
|
||||
const val CYCLIST = 2
|
||||
const val PASSENGER_CAR = 5
|
||||
|
||||
/**
|
||||
* Roadside infrastructure. Note the gap: the enumeration runs 0..11 then jumps to 15, so this
|
||||
* is 15 and not 12 - mapping by list position mislabels every RSU.
|
||||
*/
|
||||
const val ROAD_SIDE_UNIT = 15
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -47,8 +47,12 @@ internal object JsonFieldReader {
|
||||
if (lat != null && lon != null) return lat to lon
|
||||
}
|
||||
|
||||
val geoJson = obj.optJSONObject("position")
|
||||
val coords = geoJson?.optJSONArray("coordinates")
|
||||
// A GeoJSON Point may be nested under "position", or `obj` may BE the Point itself - the
|
||||
// Use Case API's DENM sends `"eventPosition": {"type":"Point","coordinates":[...]}`, so the
|
||||
// caller passes eventPosition in directly. Missing this second case meant every DENM was
|
||||
// rejected for having no position.
|
||||
val coords = obj.optJSONObject("position")?.optJSONArray("coordinates")
|
||||
?: obj.optJSONArray("coordinates")
|
||||
if (coords != null && coords.length() >= 2) {
|
||||
// GeoJSON coordinate order is [longitude, latitude, altitude?]
|
||||
val lon = coords.optDouble(0, Double.NaN)
|
||||
|
||||
@@ -11,10 +11,12 @@ import org.json.JSONObject
|
||||
* more (validity duration, relevance area, traffic direction, trace paths); none of it is used
|
||||
* yet, and inventing a fuller model before there's a consumer for it would just be guesswork.
|
||||
*
|
||||
* **Availability:** DENM reaches the app only on the CiT One path, via the Use Case API's
|
||||
* `v2x-uca/output/json/denm` topic. The ESP32-C5 path receives none — the firmware's
|
||||
* `gn_unwrap.c` accepts BTP-B destination port 2001 (CAM) only and drops port 2002 (DENM) before
|
||||
* anything is forwarded over the serial link. See that file's header comment.
|
||||
* **Availability:** both hardware paths. On the CiT One path DENM arrives as processed JSON on
|
||||
* the Use Case API's `v2x-uca/output/json/denm` topic ([DenmParser]); on the ESP32-C5 path it is
|
||||
* decoded from over-the-air GeoBroadcast traffic on BTP-B port 2002
|
||||
* ([com.hawhamburg.micr0bu.domain.asn1.DenmUperCodec]). Fields sourced from the GeoNetworking
|
||||
* header - [relevanceRadiusM], [rssiDbm] - exist only on the ESP32-C5 path, since the Use Case
|
||||
* API never exposes the GN layer.
|
||||
*/
|
||||
data class DenmEvent(
|
||||
/** Originating station ID — `actionID.originatingStationID`, not the radio source. */
|
||||
@@ -23,8 +25,8 @@ data class DenmEvent(
|
||||
/**
|
||||
* `actionID.sequenceNumber`. Together with [stationId] this is ETSI's real event identity:
|
||||
* repetitions of one hazard reuse it, and under GeoBroadcast several stations may relay the
|
||||
* same DENM, so this pair is what dedup must key on. Null on the MQTT path when the Use Case
|
||||
* API doesn't supply it.
|
||||
* same DENM, so this pair is what dedup must key on. Both transports supply it - the Use Case
|
||||
* API as a `sequenceNumber` key - so the fallback below is for malformed payloads only.
|
||||
*/
|
||||
val sequenceNumber: Int? = null,
|
||||
|
||||
@@ -64,9 +66,10 @@ data class DenmEvent(
|
||||
val timestamp: Long,
|
||||
) {
|
||||
/**
|
||||
* Stable identity for map/list dedup. Prefers ETSI's actionID (`stationId` + `sequenceNumber`)
|
||||
* where available; falls back to station + cause on the MQTT path, which doesn't reliably
|
||||
* expose a sequence number.
|
||||
* Stable identity for map/list dedup. Prefers ETSI's actionID (`stationId` + `sequenceNumber`),
|
||||
* which both transports carry; falls back to station + cause only when a payload omits the
|
||||
* sequence number. The fallback is weaker than it looks: a terminating DENM carries no
|
||||
* SituationContainer, so its cause is null and it would NOT collide with the event it ends.
|
||||
*/
|
||||
val dedupKey: String
|
||||
get() = if (sequenceNumber != null) "$stationId/$sequenceNumber"
|
||||
@@ -77,11 +80,12 @@ data class DenmEvent(
|
||||
* Parses the processed DENM JSON published by the consider it Use Case API on
|
||||
* `v2x-uca/output/json/denm`.
|
||||
*
|
||||
* Same field-name tolerance approach as [com.hawhamburg.micr0bu.domain.cam.CamParser] — confirmed
|
||||
* spellings first, plausible alternatives as fallbacks via [JsonFieldReader] — because the exact
|
||||
* schema hasn't been pinned against real OBU payloads yet. Returns null rather than a
|
||||
* half-populated event when position is missing: a DENM with no position is useless to a map and
|
||||
* worse than absent on a hazard display.
|
||||
* Field names follow `CI-CiT-MQTT_API_Documentation-v6-20250221.pdf` section 2.2.4 / listing 2.6,
|
||||
* which is the contract for this topic; `DenmParserMqttTest` pins this parser to that worked
|
||||
* example. Alternative spellings are still accepted via [JsonFieldReader] as fallbacks.
|
||||
*
|
||||
* Returns null rather than a half-populated event when position is missing: a DENM with no
|
||||
* position is useless to a map and worse than absent on a hazard display.
|
||||
*/
|
||||
object DenmParser {
|
||||
|
||||
@@ -92,7 +96,10 @@ object DenmParser {
|
||||
*/
|
||||
private val CAUSE_CODE_BY_NAME = mapOf(
|
||||
"trafficCondition" to 1, "accident" to 2, "roadworks" to 3, "impassability" to 5,
|
||||
"adverseWeatherCondition_Adhesion" to 6, "aquaplanning" to 7,
|
||||
"adverseWeatherCondition_Adhesion" to 6,
|
||||
// Three n's: that is how ETSI's CauseCodeType spells it, and the API follows.
|
||||
// The correctly-spelled variant is accepted too, in case that is ever fixed.
|
||||
"aquaplannning" to 7, "aquaplanning" to 7,
|
||||
"hazardousLocation_SurfaceCondition" to 9, "hazardousLocation_ObstacleOnTheRoad" to 10,
|
||||
"hazardousLocation_AnimalOnTheRoad" to 11, "humanPresenceOnTheRoad" to 12,
|
||||
"wrongWayDriving" to 14, "rescueAndRecoveryWorkInProgress" to 15,
|
||||
@@ -116,6 +123,29 @@ object DenmParser {
|
||||
*/
|
||||
fun causeCodeName(causeCode: Int?): String? = causeCode?.let { NAME_BY_CAUSE_CODE[it] }
|
||||
|
||||
/**
|
||||
* The Use Case API's `stationType` string enum mapped to its ITS-G5 integer, per StationType in
|
||||
* the ETSI CDD. Note roadSideUnit is **15**, not 12 - the enumeration has a gap after tram(11),
|
||||
* so mapping by list position would silently mislabel every RSU.
|
||||
*/
|
||||
private val STATION_TYPE_BY_NAME = mapOf(
|
||||
"unknown" to 0, "pedestrian" to 1, "cyclist" to 2, "moped" to 3, "motorcycle" to 4,
|
||||
"passengerCar" to 5, "bus" to 6, "lightTruck" to 7, "heavyTruck" to 8, "trailer" to 9,
|
||||
"specialVehicles" to 10, "tram" to 11, "roadSideUnit" to 15,
|
||||
)
|
||||
|
||||
/**
|
||||
* Parses the API's RFC3339 timestamps ("2021-05-11T12:01:02+00:00") to epoch millis. The air
|
||||
* path carries a binary TimestampIts instead, so the two transports arrive here in completely
|
||||
* different formats and both end up as epoch ms on [DenmEvent].
|
||||
*/
|
||||
private fun parseRfc3339(value: String?): Long? {
|
||||
if (value.isNullOrBlank()) return null
|
||||
return runCatching { java.time.OffsetDateTime.parse(value).toInstant().toEpochMilli() }
|
||||
.recoverCatching { java.time.Instant.parse(value).toEpochMilli() }
|
||||
.getOrNull()
|
||||
}
|
||||
|
||||
fun parse(json: String, timestamp: Long = System.currentTimeMillis()): DenmEvent? {
|
||||
val obj = runCatching { JSONObject(json) }.getOrNull() ?: return null
|
||||
|
||||
@@ -127,9 +157,16 @@ object DenmParser {
|
||||
?.let { JsonFieldReader.firstLatLon(it) }
|
||||
?: return null
|
||||
|
||||
val stationId = JsonFieldReader.firstLong(obj, "stationId", "stationID", "station_id")
|
||||
// "originatingStationId" is what the Use Case API actually sends (API doc listing 2.6);
|
||||
// without it every DENM from the CiT One path was dropped here, before anything else in
|
||||
// this function ran. The other spellings are kept as fallbacks.
|
||||
val stationId = JsonFieldReader.firstLong(
|
||||
obj, "originatingStationId", "originatingStationID", "stationId", "stationID", "station_id",
|
||||
)
|
||||
?: obj.optJSONObject("management")?.let {
|
||||
JsonFieldReader.firstLong(it, "stationId", "stationID", "station_id")
|
||||
JsonFieldReader.firstLong(
|
||||
it, "originatingStationId", "originatingStationID", "stationId", "stationID", "station_id",
|
||||
)
|
||||
}
|
||||
?: return null
|
||||
|
||||
@@ -144,12 +181,25 @@ object DenmParser {
|
||||
val subCauseCode = JsonFieldReader.firstInt(obj, "subCauseCode", "sub_cause_code", "subCause")
|
||||
?: situation?.let { JsonFieldReader.firstInt(it, "subCauseCode", "sub_cause_code", "subCause") }
|
||||
|
||||
// "Key is present, if the DENM is cancelled or negated" (API doc 2.2.4) - so presence is
|
||||
// the signal, not the value. An explicit `false` is still honoured in case the API ever
|
||||
// starts always emitting the key.
|
||||
val isTermination = obj.has("termination") && obj.optBoolean("termination", true)
|
||||
|
||||
val stationType = JsonFieldReader.firstInt(obj, "stationType", "station_type")
|
||||
?: STATION_TYPE_BY_NAME[obj.optString("stationType").takeIf { it.isNotBlank() }]
|
||||
|
||||
return DenmEvent(
|
||||
stationId = stationId,
|
||||
sequenceNumber = JsonFieldReader.firstInt(obj, "sequenceNumber", "sequence_number"),
|
||||
latitude = lat,
|
||||
longitude = lon,
|
||||
causeCode = causeCode,
|
||||
subCauseCode = subCauseCode,
|
||||
stationType = stationType,
|
||||
isTermination = isTermination,
|
||||
detectionTimeMs = parseRfc3339(obj.optString("detectionTime").takeIf { it.isNotBlank() })
|
||||
?: JsonFieldReader.firstLong(obj, "detectionTime", "detection_time"),
|
||||
timestamp = timestamp,
|
||||
)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,141 @@
|
||||
package com.hawhamburg.micr0bu.domain.spat
|
||||
|
||||
/**
|
||||
* Signal phase and timing for one or more intersections, decoded from a SPATEM heard over the air.
|
||||
*
|
||||
* A SPATEM is the live counterpart to MAPEM's static geometry: MAPEM says where the lanes are,
|
||||
* SPATEM says what the lights are doing right now. The two join on [IntersectionSignalState.key].
|
||||
* Only SPATEM is decoded today - without MAPEM there is no lane geometry, so a signal group is
|
||||
* shown as a bare number rather than "the left-turn lane you are in".
|
||||
*
|
||||
* Repetition is ~2 Hz per intersection, so consumers should key on [IntersectionSignalState.key]
|
||||
* and keep the latest rather than accumulating a log.
|
||||
*/
|
||||
data class SpatEvent(
|
||||
/** Originating RSU's station ID, from the ItsPduHeader. */
|
||||
val stationId: Long,
|
||||
|
||||
/** `SPAT.timeStamp`, minute of the year, when present. */
|
||||
val minuteOfYear: Int?,
|
||||
|
||||
val intersections: List<IntersectionSignalState>,
|
||||
|
||||
/** Received signal strength, dBm - ESP32-C5 path only. */
|
||||
val rssiDbm: Int? = null,
|
||||
|
||||
/** Wall-clock ms this SPATEM was received. */
|
||||
val timestamp: Long,
|
||||
)
|
||||
|
||||
/** One intersection's current signal state. */
|
||||
data class IntersectionSignalState(
|
||||
/** `RoadRegulatorID`, when the sender qualifies its intersection id with one. */
|
||||
val region: Int?,
|
||||
/** `IntersectionID` - only unique *within* [region]. */
|
||||
val id: Int,
|
||||
/** `MsgCount`, bumped when the intersection's MAP geometry changes. */
|
||||
val revision: Int,
|
||||
/** Minute of the year this state refers to, when present. */
|
||||
val moy: Int?,
|
||||
/** `DSecond` - milliseconds within [moy]'s minute, when present. */
|
||||
val timeStampMs: Int?,
|
||||
val movements: List<SignalMovement>,
|
||||
) {
|
||||
/**
|
||||
* Identity for dedup and for joining against MAPEM. `IntersectionID` alone is NOT unique -
|
||||
* it is only unique within a `RoadRegulatorID`, and the recorded drive contains the same id
|
||||
* under different regions - so the region must be part of the key.
|
||||
*/
|
||||
val key: String get() = "${region ?: -1}/$id"
|
||||
}
|
||||
|
||||
/** The signal state of one signal group (one movement through the intersection). */
|
||||
data class SignalMovement(
|
||||
/** `SignalGroupID` - the number MAPEM's lane connections refer to. */
|
||||
val signalGroup: Int,
|
||||
/**
|
||||
* Predicted phases, in order. The first entry is the state now; later entries are the
|
||||
* upcoming sequence, which is what makes a countdown possible.
|
||||
*/
|
||||
val events: List<SignalPhaseEvent>,
|
||||
) {
|
||||
val current: SignalPhaseEvent? get() = events.firstOrNull()
|
||||
}
|
||||
|
||||
data class SignalPhaseEvent(
|
||||
val phase: SignalPhase,
|
||||
/**
|
||||
* `TimeMark`: tenths of a second within the current or next UTC hour, so it wraps hourly.
|
||||
* 36001 means "unknown". Use [secondsUntil] rather than comparing these directly.
|
||||
*/
|
||||
val minEndTimeDs: Int?,
|
||||
val maxEndTimeDs: Int?,
|
||||
val likelyTimeDs: Int?,
|
||||
) {
|
||||
/**
|
||||
* Seconds from [nowEpochMs] until [minEndTimeDs], or null if unknown.
|
||||
*
|
||||
* TimeMark counts tenths of a second from the top of the hour and wraps, so a mark that looks
|
||||
* like it is in the past is really in the next hour - hence the wrap correction. Without it, a
|
||||
* countdown reads as a large negative number for the seconds either side of the hour.
|
||||
*/
|
||||
fun secondsUntil(nowEpochMs: Long): Double? {
|
||||
val mark = minEndTimeDs ?: return null
|
||||
if (mark >= UNKNOWN_TIME_MARK) return null
|
||||
val msIntoHour = nowEpochMs % 3_600_000L
|
||||
var deltaMs = mark * 100L - msIntoHour
|
||||
if (deltaMs < -HALF_HOUR_MS) deltaMs += 3_600_000L // mark is in the next hour
|
||||
return deltaMs / 1000.0
|
||||
}
|
||||
|
||||
private companion object {
|
||||
const val UNKNOWN_TIME_MARK = 36001
|
||||
const val HALF_HOUR_MS = 1_800_000L
|
||||
}
|
||||
}
|
||||
|
||||
/** `MovementPhaseState` (ETSI/SAE J2735), in enumeration order - the ordinal IS the wire value. */
|
||||
enum class SignalPhase {
|
||||
UNAVAILABLE,
|
||||
DARK,
|
||||
STOP_THEN_PROCEED,
|
||||
STOP_AND_REMAIN,
|
||||
PRE_MOVEMENT,
|
||||
PERMISSIVE_MOVEMENT_ALLOWED,
|
||||
PROTECTED_MOVEMENT_ALLOWED,
|
||||
PERMISSIVE_CLEARANCE,
|
||||
PROTECTED_CLEARANCE,
|
||||
CAUTION_CONFLICTING_TRAFFIC;
|
||||
|
||||
/** True for the two "you may go" states. */
|
||||
val isGo: Boolean
|
||||
get() = this == PERMISSIVE_MOVEMENT_ALLOWED || this == PROTECTED_MOVEMENT_ALLOWED
|
||||
|
||||
/** True for the two "you must stop" states. */
|
||||
val isStop: Boolean
|
||||
get() = this == STOP_AND_REMAIN || this == STOP_THEN_PROCEED
|
||||
|
||||
/** True while the light is changing - amber, or red-amber before green. */
|
||||
val isTransition: Boolean
|
||||
get() = this == PRE_MOVEMENT || this == PERMISSIVE_CLEARANCE || this == PROTECTED_CLEARANCE
|
||||
|
||||
companion object {
|
||||
fun fromWire(value: Int): SignalPhase? = entries.getOrNull(value)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* One intersection's signal state as the UI consumes it: the decoded state plus who sent it and
|
||||
* when, flattened out of the [SpatEvent] that carried it.
|
||||
*
|
||||
* A single SPATEM may describe several intersections, and the same intersection may be heard from
|
||||
* more than one RSU, so the UI keys on the intersection rather than on the message.
|
||||
*/
|
||||
data class SpatIntersection(
|
||||
val state: IntersectionSignalState,
|
||||
val stationId: Long,
|
||||
val rssiDbm: Int?,
|
||||
val timestamp: Long,
|
||||
) {
|
||||
val key: String get() = state.key
|
||||
}
|
||||
@@ -8,6 +8,7 @@ import androidx.compose.foundation.clickable
|
||||
import androidx.compose.foundation.layout.Arrangement
|
||||
import androidx.compose.foundation.layout.Box
|
||||
import androidx.compose.foundation.layout.Column
|
||||
import androidx.compose.foundation.layout.FlowRow
|
||||
import androidx.compose.foundation.layout.Row
|
||||
import androidx.compose.foundation.layout.Spacer
|
||||
import androidx.compose.foundation.layout.fillMaxSize
|
||||
@@ -84,8 +85,12 @@ import java.text.SimpleDateFormat
|
||||
import java.util.Date
|
||||
import java.util.Locale
|
||||
|
||||
/** List (raw topics) vs Map (V2X live map, Section 13) toggle for [TopicListPane]. */
|
||||
private enum class TopicViewMode { LIST, MAP }
|
||||
/**
|
||||
* View toggle for [TopicListPane]: decoded CAM/DENM traffic (LIST), the raw MQTT topic list
|
||||
* (TOPICS, CiT One only - there is no broker on the ESP32-C5 path), or the V2X live map
|
||||
* (MAP, Section 13).
|
||||
*/
|
||||
private enum class TopicViewMode { LIST, TOPICS, MAP }
|
||||
|
||||
private val timeFormat = SimpleDateFormat("HH:mm:ss.SSS", Locale.US)
|
||||
|
||||
@@ -130,6 +135,7 @@ fun MqttTopicViewerScreen(
|
||||
val camSendFailures by viewModel.camSendFailures.collectAsState()
|
||||
val espLinkStatus by viewModel.espLinkStatus.collectAsState()
|
||||
val denmEvents by viewModel.denmEvents.collectAsState()
|
||||
val spatIntersections by viewModel.spatIntersections.collectAsState()
|
||||
|
||||
// Sort: sys/ topics first (heartbeat/health), then alphabetical
|
||||
val sortedTopics = topicMessages.keys.sortedWith(
|
||||
@@ -226,6 +232,7 @@ fun MqttTopicViewerScreen(
|
||||
showCamPinger = isEsp32,
|
||||
isEsp32 = isEsp32,
|
||||
denmEvents = denmEvents,
|
||||
spatIntersections = spatIntersections,
|
||||
usbSerialState = usbSerialState,
|
||||
camPingerActive = camPingerActive,
|
||||
camPingerSentCount = camPingerSentCount,
|
||||
@@ -265,6 +272,7 @@ private fun TopicListPane(
|
||||
showCamPinger: Boolean = false,
|
||||
isEsp32: Boolean = false,
|
||||
denmEvents: List<com.hawhamburg.micr0bu.domain.denm.DenmEvent> = emptyList(),
|
||||
spatIntersections: List<com.hawhamburg.micr0bu.domain.spat.SpatIntersection> = emptyList(),
|
||||
usbSerialState: UsbSerialState = UsbSerialState.DISCONNECTED,
|
||||
camPingerActive: Boolean = false,
|
||||
camPingerSentCount: Int = 0,
|
||||
@@ -321,30 +329,41 @@ private fun TopicListPane(
|
||||
HorizontalDivider(color = MaterialTheme.colorScheme.outline.copy(alpha = 0.25f))
|
||||
}
|
||||
|
||||
// ── List / Map toggle — the raw topic list stays available either way (Section 13
|
||||
// asks for the map "in addition to", not instead of, the topic list). ──────────────
|
||||
// ── List / Topics / Map toggle ────────────────────────────────────────────────────
|
||||
// Decoded traffic is the default on BOTH hardware paths: what a tester wants to see is
|
||||
// the road users and hazards, not the transport that carried them. The raw MQTT topic
|
||||
// list stays one tap away on the CiT One path (Section 13 asks for the map "in addition
|
||||
// to", not instead of, the topic list). It is hidden on the ESP32-C5 path, where there is
|
||||
// no broker and `topics` is permanently empty.
|
||||
Row(
|
||||
modifier = Modifier.fillMaxWidth().padding(horizontal = 12.dp, vertical = 6.dp),
|
||||
horizontalArrangement = Arrangement.spacedBy(8.dp),
|
||||
) {
|
||||
OutlinedButton(
|
||||
onClick = { viewMode = TopicViewMode.LIST },
|
||||
colors = ButtonDefaults.outlinedButtonColors(
|
||||
containerColor = if (viewMode == TopicViewMode.LIST) MaterialTheme.colorScheme.primaryContainer else Color.Transparent,
|
||||
contentColor = if (viewMode == TopicViewMode.LIST) MaterialTheme.colorScheme.onPrimaryContainer else MaterialTheme.colorScheme.onSurface,
|
||||
),
|
||||
) { Text(stringResource(R.string.mqtt_view_list)) }
|
||||
OutlinedButton(
|
||||
onClick = { viewMode = TopicViewMode.MAP },
|
||||
colors = ButtonDefaults.outlinedButtonColors(
|
||||
containerColor = if (viewMode == TopicViewMode.MAP) MaterialTheme.colorScheme.primaryContainer else Color.Transparent,
|
||||
contentColor = if (viewMode == TopicViewMode.MAP) MaterialTheme.colorScheme.onPrimaryContainer else MaterialTheme.colorScheme.onSurface,
|
||||
),
|
||||
) { Text(stringResource(R.string.mqtt_view_map)) }
|
||||
ViewModeButton(
|
||||
label = stringResource(R.string.mqtt_view_list),
|
||||
selected = viewMode == TopicViewMode.LIST,
|
||||
) { viewMode = TopicViewMode.LIST }
|
||||
|
||||
if (!isEsp32) {
|
||||
ViewModeButton(
|
||||
label = stringResource(R.string.mqtt_view_topics),
|
||||
selected = viewMode == TopicViewMode.TOPICS,
|
||||
) { viewMode = TopicViewMode.TOPICS }
|
||||
}
|
||||
|
||||
ViewModeButton(
|
||||
label = stringResource(R.string.mqtt_view_map),
|
||||
selected = viewMode == TopicViewMode.MAP,
|
||||
) { viewMode = TopicViewMode.MAP }
|
||||
}
|
||||
|
||||
// ── Topic rows / received CAMs / live map ─────────────────────────────
|
||||
if (viewMode == TopicViewMode.MAP) {
|
||||
// ── Decoded traffic / raw topics / live map ───────────────────────────
|
||||
// TOPICS can still be the saved selection from a CiT One session after switching hardware
|
||||
// to the ESP32-C5, where that button no longer exists - fall back to the decoded list
|
||||
// rather than stranding the user on a pane they can't navigate away from.
|
||||
val shownMode = if (viewMode == TopicViewMode.TOPICS && isEsp32) TopicViewMode.LIST else viewMode
|
||||
|
||||
if (shownMode == TopicViewMode.MAP) {
|
||||
V2xLiveMapView(
|
||||
own = ownCamPosition,
|
||||
remotes = remoteCamPositions,
|
||||
@@ -352,15 +371,13 @@ private fun TopicListPane(
|
||||
denms = denmEvents,
|
||||
modifier = Modifier.fillMaxSize(),
|
||||
)
|
||||
} else if (isEsp32) {
|
||||
// The MQTT topic list is meaningless on this path - there is no broker, so `topics`
|
||||
// is permanently empty and the list would read as "nothing is happening" even while
|
||||
// CAMs stream in over the serial link. Show the decoded traffic instead.
|
||||
} else if (shownMode == TopicViewMode.LIST) {
|
||||
ReceivedCamPane(
|
||||
own = ownCamPosition,
|
||||
remotes = remoteCamPositions,
|
||||
alerts = useCaseAlerts,
|
||||
denms = denmEvents,
|
||||
spats = spatIntersections,
|
||||
modifier = Modifier.fillMaxSize(),
|
||||
)
|
||||
} else if (topics.isEmpty()) {
|
||||
@@ -421,9 +438,10 @@ private fun ReceivedCamPane(
|
||||
remotes: Map<Long, com.hawhamburg.micr0bu.domain.cam.Cam>,
|
||||
alerts: List<UseCaseAlert>,
|
||||
denms: List<com.hawhamburg.micr0bu.domain.denm.DenmEvent>,
|
||||
spats: List<com.hawhamburg.micr0bu.domain.spat.SpatIntersection>,
|
||||
modifier: Modifier = Modifier,
|
||||
) {
|
||||
if (remotes.isEmpty() && denms.isEmpty()) {
|
||||
if (remotes.isEmpty() && denms.isEmpty() && spats.isEmpty()) {
|
||||
Box(modifier = modifier, contentAlignment = Alignment.Center) {
|
||||
Column(horizontalAlignment = Alignment.CenterHorizontally) {
|
||||
Text(
|
||||
@@ -485,6 +503,14 @@ private fun ReceivedCamPane(
|
||||
}
|
||||
}
|
||||
|
||||
if (spats.isNotEmpty()) {
|
||||
item { PaneSectionHeader(stringResource(R.string.v2x_spat_rx_count, spats.size)) }
|
||||
items(spats, key = { "spat/" + it.key }) { spat ->
|
||||
ReceivedSpatRow(spat)
|
||||
HorizontalDivider(color = MaterialTheme.colorScheme.outline.copy(alpha = 0.25f))
|
||||
}
|
||||
}
|
||||
|
||||
item {
|
||||
PaneSectionHeader(
|
||||
if (rows.isEmpty()) stringResource(R.string.v2x_cam_rx_none_stations)
|
||||
@@ -498,6 +524,82 @@ private fun ReceivedCamPane(
|
||||
}
|
||||
}
|
||||
|
||||
@Composable
|
||||
private fun ViewModeButton(label: String, selected: Boolean, onClick: () -> Unit) {
|
||||
OutlinedButton(
|
||||
onClick = onClick,
|
||||
colors = ButtonDefaults.outlinedButtonColors(
|
||||
containerColor = if (selected) MaterialTheme.colorScheme.primaryContainer else Color.Transparent,
|
||||
contentColor = if (selected) MaterialTheme.colorScheme.onPrimaryContainer else MaterialTheme.colorScheme.onSurface,
|
||||
),
|
||||
) { Text(label) }
|
||||
}
|
||||
|
||||
/**
|
||||
* One signalised intersection: every signal group's current phase, with a countdown where the RSU
|
||||
* supplies one.
|
||||
*
|
||||
* Signal groups are shown as bare numbers because that is genuinely all we know - mapping a group
|
||||
* to "your lane" needs MAPEM geometry, which nothing on the air is currently sending. Inventing a
|
||||
* friendlier label would imply knowledge the app does not have.
|
||||
*/
|
||||
@OptIn(androidx.compose.foundation.layout.ExperimentalLayoutApi::class)
|
||||
@Composable
|
||||
private fun ReceivedSpatRow(spat: com.hawhamburg.micr0bu.domain.spat.SpatIntersection) {
|
||||
val now = System.currentTimeMillis()
|
||||
Column(
|
||||
modifier = Modifier
|
||||
.fillMaxWidth()
|
||||
.padding(horizontal = 16.dp, vertical = 10.dp),
|
||||
) {
|
||||
Text(
|
||||
text = stringResource(R.string.v2x_spat_rx_title, spat.state.key, spat.stationId),
|
||||
style = MaterialTheme.typography.bodyMedium,
|
||||
fontWeight = FontWeight.SemiBold,
|
||||
color = MaterialTheme.colorScheme.primary,
|
||||
)
|
||||
Spacer(Modifier.height(4.dp))
|
||||
|
||||
// Wraps rather than scrolls: a busy intersection has 20+ groups and a horizontal
|
||||
// scroller inside a vertical list is awkward to drive one-handed on a bike mount.
|
||||
FlowRow(horizontalArrangement = Arrangement.spacedBy(6.dp)) {
|
||||
spat.state.movements.forEach { movement ->
|
||||
val phase = movement.current?.phase
|
||||
val tint = when {
|
||||
phase == null -> MaterialTheme.colorScheme.onSurfaceVariant
|
||||
phase.isGo -> ConnectedGreen
|
||||
phase.isStop -> WarningRed
|
||||
phase.isTransition -> AwarenessAmber
|
||||
else -> DisconnectedGray
|
||||
}
|
||||
val countdown = movement.current?.secondsUntil(now)?.takeIf { it in 0.0..99.0 }
|
||||
Text(
|
||||
text = stringResource(R.string.v2x_spat_group, movement.signalGroup) +
|
||||
(countdown?.let { " " + stringResource(R.string.v2x_spat_countdown, it) } ?: ""),
|
||||
style = MaterialTheme.typography.bodySmall,
|
||||
fontFamily = FontFamily.Monospace,
|
||||
color = tint,
|
||||
modifier = Modifier
|
||||
.padding(vertical = 2.dp)
|
||||
.clip(RoundedCornerShape(4.dp))
|
||||
.background(tint.copy(alpha = 0.15f))
|
||||
.padding(horizontal = 6.dp, vertical = 2.dp),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
spat.rssiDbm?.let {
|
||||
Spacer(Modifier.height(3.dp))
|
||||
Text(
|
||||
text = stringResource(R.string.v2x_cam_rx_rssi, it),
|
||||
style = MaterialTheme.typography.bodySmall,
|
||||
fontFamily = FontFamily.Monospace,
|
||||
color = MaterialTheme.colorScheme.onSurfaceVariant,
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Composable
|
||||
private fun PaneSectionHeader(text: String) {
|
||||
Column {
|
||||
|
||||
@@ -17,13 +17,17 @@ import com.hawhamburg.micr0bu.data.transport.UsbSerialState
|
||||
import com.hawhamburg.micr0bu.data.transport.UsbSerialTransport
|
||||
import com.hawhamburg.micr0bu.domain.denm.DenmEvent
|
||||
import com.hawhamburg.micr0bu.domain.denm.DenmParser
|
||||
import com.hawhamburg.micr0bu.domain.spat.SpatIntersection
|
||||
import com.hawhamburg.micr0bu.domain.denm.DenmUseCase
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseAlert
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseType
|
||||
import com.hawhamburg.micr0bu.service.CamPinger
|
||||
import dagger.hilt.android.lifecycle.HiltViewModel
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.combine
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.flow
|
||||
import kotlinx.coroutines.flow.runningFold
|
||||
import kotlinx.coroutines.flow.SharingStarted
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
@@ -155,17 +159,77 @@ class MqttViewModel @Inject constructor(
|
||||
camUseCaseRepository.airDenm
|
||||
.runningFold(emptyMap<String, DenmEvent>()) { acc, denm -> acc + (denm.dedupKey to denm) }
|
||||
.map { it.values.toList() },
|
||||
) { fromMqtt, fromAir ->
|
||||
// Expiry has to be driven by a clock, not by arrivals. Both upstream flows only re-emit
|
||||
// when a DENM arrives, so a sender that simply stops transmitting - drives away, loses
|
||||
// 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, _ ->
|
||||
val now = System.currentTimeMillis()
|
||||
(fromMqtt + fromAir)
|
||||
.filterNot { it.isTermination } // the hazard is over - stop drawing it
|
||||
.associateBy { it.dedupKey } // last write wins = most recent per hazard
|
||||
.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".
|
||||
.filter { now - it.timestamp <= DENM_TTL_MS }
|
||||
.sortedByDescending { it.timestamp }
|
||||
}.stateIn(viewModelScope, SharingStarted.Eagerly, emptyList())
|
||||
|
||||
/**
|
||||
* 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.
|
||||
*
|
||||
* 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
|
||||
.runningFold(emptyMap<String, SpatIntersection>()) { acc, spat ->
|
||||
acc + spat.intersections.associate { i ->
|
||||
i.key to SpatIntersection(i, spat.stationId, spat.rssiDbm, spat.timestamp)
|
||||
}
|
||||
},
|
||||
tickerFlow(SPAT_EXPIRY_TICK_MS),
|
||||
) { byKey, _ ->
|
||||
val now = System.currentTimeMillis()
|
||||
byKey.values
|
||||
.filter { now - it.timestamp <= SPAT_TTL_MS }
|
||||
.sortedByDescending { it.timestamp }
|
||||
}.stateIn(viewModelScope, SharingStarted.Eagerly, emptyList())
|
||||
|
||||
/** Emits immediately, then every [periodMs], purely to re-trigger a time-dependent combine. */
|
||||
private fun tickerFlow(periodMs: Long): Flow<Long> = flow {
|
||||
while (true) {
|
||||
emit(System.currentTimeMillis())
|
||||
delay(periodMs)
|
||||
}
|
||||
}
|
||||
|
||||
private companion object {
|
||||
/** Use Case API topic carrying received DENMs (CiT One path only). */
|
||||
const val DENM_RX_TOPIC = "v2x-uca/output/json/denm"
|
||||
|
||||
/**
|
||||
* How long a hazard stays listed after its last repetition. A DENM has no "still here"
|
||||
* guarantee beyond the sender repeating it, and its own validityDuration is not decoded
|
||||
* yet, so silence is the only expiry signal available.
|
||||
*/
|
||||
const val DENM_TTL_MS = 60_000L
|
||||
|
||||
/** How often the list is re-evaluated for expiry. Sets the worst-case lateness of a fade. */
|
||||
const val DENM_EXPIRY_TICK_MS = 5_000L
|
||||
|
||||
/**
|
||||
* SPATEM repeats at ~2 Hz, so 15 s of silence is ~30 missed repetitions: the RSU is out of
|
||||
* range. Much shorter than the DENM window because a stale traffic light is more
|
||||
* misleading than a stale hazard - a light that stopped updating is not "still green".
|
||||
*/
|
||||
const val SPAT_TTL_MS = 15_000L
|
||||
const val SPAT_EXPIRY_TICK_MS = 2_000L
|
||||
}
|
||||
|
||||
// ── DENM transmission ─────────────────────────────────────────────────────
|
||||
|
||||
Reference in New Issue
Block a user