initial commit
This commit is contained in:
@@ -0,0 +1,192 @@
|
||||
package com.hawhamburg.micr0bu.data.cam
|
||||
|
||||
import android.content.Context
|
||||
import com.hawhamburg.micr0bu.data.GnssReading
|
||||
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.UseCaseAlertPreferences
|
||||
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.usecase.UseCaseAlert
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseDetectionEngine
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseType
|
||||
import dagger.hilt.android.qualifiers.ApplicationContext
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.SharingStarted
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
import kotlinx.coroutines.flow.combine
|
||||
import kotlinx.coroutines.flow.stateIn
|
||||
import kotlinx.coroutines.launch
|
||||
import javax.inject.Inject
|
||||
import javax.inject.Singleton
|
||||
|
||||
private const val CAM_TOPIC = "v2x-uca/output/json/cam"
|
||||
private const val OBU_GNSS_TOPIC = "v2x/rx/obu_gnss"
|
||||
private const val PRUNE_INTERVAL_MS = 1_000L
|
||||
|
||||
// If no v2x/rx/obu_gnss update has arrived within this window, the ego state is considered
|
||||
// stale enough that a fresh phone GNSS fix (if available) is preferred over it — see
|
||||
// [handlePhoneGnss]. obu_gnss updates at ~4 Hz per the requirements doc, so 2.5 s is several
|
||||
// missed updates, not just normal jitter between samples.
|
||||
private const val OBU_GNSS_STALE_MS = 2_500L
|
||||
|
||||
/**
|
||||
* 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
|
||||
* (requirements doc Section 10.2 / 10.4).
|
||||
*
|
||||
* **Ego state sourcing:** `v2x/rx/obu_gnss` is the primary source for the ego bike's own
|
||||
* position/speed/heading/yaw rate (~4 Hz, includes yaw rate). If it goes stale (cable
|
||||
* unplugged, OBU hiccup, etc.), phone GNSS (`FusedLocationProviderClient` via
|
||||
* [SensorRepository]) is used as a fallback so the engine keeps running — at the cost of yaw
|
||||
* rate, which the phone alone doesn't provide. The CAM topic's own low-rate "own" entry is
|
||||
* also fed in as a third fallback. [UseCaseDetectionEngine.onOwnCam] always keeps whichever
|
||||
* update is freshest and ignores out-of-order ones, so no explicit priority juggling is needed
|
||||
* beyond "only use phone GNSS when obu_gnss is stale".
|
||||
*
|
||||
* 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.
|
||||
*
|
||||
* No DENM is generated or consumed anywhere in this class.
|
||||
*/
|
||||
@Singleton
|
||||
class CamUseCaseRepository @Inject constructor(
|
||||
private val mqttRepository: MqttRepository,
|
||||
private val prefs: UseCaseAlertPreferences,
|
||||
@ApplicationContext private val context: Context,
|
||||
) {
|
||||
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
|
||||
private val engine = UseCaseDetectionEngine()
|
||||
private val sensorRepository = SensorRepository(context)
|
||||
|
||||
private val _ownStationId = MutableStateFlow<Long?>(null)
|
||||
/** The ego OBU's own station ID, learned from `v2x/rx/obu_gnss`. Null until known. */
|
||||
val ownStationId: StateFlow<Long?> = _ownStationId.asStateFlow()
|
||||
|
||||
@Volatile private var lastOwnStationType: Int = StationType.CYCLIST
|
||||
@Volatile private var lastObuGnssTimestamp: 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 },
|
||||
)
|
||||
|
||||
/** All currently active alerts, regardless of per-use-case enablement. */
|
||||
val allAlerts: StateFlow<List<UseCaseAlert>> = engine.currentAlerts
|
||||
|
||||
/** Alerts filtered to only the use cases the user has enabled — what the UI should show. */
|
||||
val enabledAlerts: StateFlow<List<UseCaseAlert>> = combine(engine.currentAlerts, enabledMap) { alerts, enabled ->
|
||||
alerts.filter { enabled[it.useCase] != false }
|
||||
}.stateIn(scope, SharingStarted.Eagerly, emptyList())
|
||||
|
||||
init {
|
||||
scope.launch {
|
||||
mqttRepository.messages.collect { msg ->
|
||||
when (msg.topic) {
|
||||
OBU_GNSS_TOPIC -> handleObuGnss(msg.payload, msg.timestamp)
|
||||
CAM_TOPIC -> handleCam(msg.payload, msg.timestamp)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 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
|
||||
// enough — re-subscribe periodically until it succeeds, and again if it ever ends.
|
||||
scope.launch {
|
||||
while (true) {
|
||||
runCatching {
|
||||
sensorRepository.gnssFlow().collect { reading -> handlePhoneGnss(reading) }
|
||||
}
|
||||
delay(5_000L)
|
||||
}
|
||||
}
|
||||
|
||||
// Clear all state on disconnect — stale remote CAMs from a previous session
|
||||
// shouldn't linger as alerts after the OBU link drops.
|
||||
scope.launch {
|
||||
mqttRepository.connectionState.collect { state ->
|
||||
if (state == MqttConnectionState.DISCONNECTED) {
|
||||
engine.reset()
|
||||
_ownStationId.value = null
|
||||
lastObuGnssTimestamp = 0L
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Periodic staleness sweep so alerts clear once a remote road user goes out of
|
||||
// range / stops transmitting, even without a new CAM arriving to trigger re-evaluation.
|
||||
scope.launch {
|
||||
while (true) {
|
||||
delay(PRUNE_INTERVAL_MS)
|
||||
engine.pruneStale(System.currentTimeMillis())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fun setUseCaseEnabled(type: UseCaseType, enabled: Boolean) {
|
||||
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
|
||||
|
||||
/**
|
||||
* Primary ego state source: `v2x/rx/obu_gnss`, ~4 Hz, carries position/speed/heading/yaw
|
||||
* rate/stationType directly (Section 6 / audit context) — a richer and higher-rate source
|
||||
* than the CAM topic's own low-rate entry.
|
||||
*/
|
||||
private fun handleObuGnss(payload: String, timestamp: Long) {
|
||||
val ego = ObuGnssParser.parseEgo(payload, fallbackStationId = _ownStationId.value ?: 0L, timestamp = timestamp)
|
||||
?: return
|
||||
lastObuGnssTimestamp = timestamp
|
||||
if (ego.stationId != 0L) _ownStationId.value = ego.stationId
|
||||
lastOwnStationType = ego.stationType
|
||||
engine.onOwnCam(ego)
|
||||
}
|
||||
|
||||
/**
|
||||
* Fallback ego state source: phone GNSS via [SensorRepository], used only once obu_gnss
|
||||
* has gone stale for longer than [OBU_GNSS_STALE_MS] — "more accurate/reliable" in this
|
||||
* context means "still updating" when the OBU feed isn't. Yaw rate is unavailable from
|
||||
* phone GNSS alone, so alerts relying on it fall back to heading-trend from history
|
||||
* (see [UseCaseDetectionEngine]) while this source is active.
|
||||
*/
|
||||
private fun handlePhoneGnss(reading: GnssReading) {
|
||||
val now = System.currentTimeMillis()
|
||||
if (now - lastObuGnssTimestamp <= OBU_GNSS_STALE_MS) return // obu_gnss is fresh enough — prefer it
|
||||
|
||||
val ego = Cam(
|
||||
stationId = _ownStationId.value ?: 0L,
|
||||
stationType = lastOwnStationType,
|
||||
latitude = reading.latitude,
|
||||
longitude = reading.longitude,
|
||||
speedMps = reading.speedMs.toDouble(),
|
||||
headingDeg = reading.bearingDeg.toDouble(),
|
||||
yawRateDps = null, // phone alone has no yaw rate; engine falls back to heading-trend
|
||||
timestamp = reading.timestamp,
|
||||
isOwn = true,
|
||||
)
|
||||
engine.onOwnCam(ego)
|
||||
}
|
||||
|
||||
private fun handleCam(payload: String, timestamp: Long) {
|
||||
val cam = CamParser.parse(payload, _ownStationId.value, timestamp) ?: return
|
||||
if (cam.isOwn) {
|
||||
// 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.
|
||||
engine.onOwnCam(cam)
|
||||
} else {
|
||||
engine.onRemoteCam(cam)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,35 @@
|
||||
package com.hawhamburg.micr0bu.data.mqtt
|
||||
|
||||
import android.content.Context
|
||||
import androidx.datastore.preferences.core.booleanPreferencesKey
|
||||
import androidx.datastore.preferences.core.edit
|
||||
import androidx.datastore.preferences.preferencesDataStore
|
||||
import com.hawhamburg.micr0bu.domain.usecase.UseCaseType
|
||||
import dagger.hilt.android.qualifiers.ApplicationContext
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.map
|
||||
import javax.inject.Inject
|
||||
import javax.inject.Singleton
|
||||
|
||||
private val Context.useCaseAlertDataStore by preferencesDataStore(name = "use_case_alert_prefs")
|
||||
|
||||
/**
|
||||
* Per-use-case enable/disable toggles for the CAM-based Use Case Alert panel
|
||||
* (Settings > Use Case Alerts, requirements doc Section 2.2 / 10.4). All in-scope use cases
|
||||
* default to enabled.
|
||||
*/
|
||||
@Singleton
|
||||
class UseCaseAlertPreferences @Inject constructor(
|
||||
@ApplicationContext private val context: Context,
|
||||
) {
|
||||
private fun keyFor(type: UseCaseType) = booleanPreferencesKey("uc_enabled_${type.name}")
|
||||
|
||||
/** Map of use case -> enabled, defaulting to true for any use case not yet persisted. */
|
||||
val enabledMapFlow: Flow<Map<UseCaseType, Boolean>> = context.useCaseAlertDataStore.data.map { prefs ->
|
||||
UseCaseType.entries.associateWith { type -> prefs[keyFor(type)] ?: true }
|
||||
}
|
||||
|
||||
suspend fun setEnabled(type: UseCaseType, enabled: Boolean) {
|
||||
context.useCaseAlertDataStore.edit { prefs -> prefs[keyFor(type)] = enabled }
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
package com.hawhamburg.micr0bu.domain.cam
|
||||
|
||||
/**
|
||||
* ETSI EN 302 637-2 stationType values relevant to this project (Table 1).
|
||||
* Only the values the app actively checks for are enumerated; the raw int is always
|
||||
* preserved on [Cam.stationType] regardless.
|
||||
*/
|
||||
object StationType {
|
||||
const val CYCLIST = 2
|
||||
const val PASSENGER_CAR = 5
|
||||
}
|
||||
|
||||
/**
|
||||
* A single Cooperative Awareness Message, parsed from the consider it Use Case API's
|
||||
* processed JSON on `v2x-uca/output/json/cam` (Section 10.2 of the requirements doc).
|
||||
*
|
||||
* This is a pure domain model — no Android or Room imports — so [com.hawhamburg.micr0bu.domain.usecase.UseCaseDetectionEngine]
|
||||
* stays fully unit-testable, matching the pattern already used by
|
||||
* [com.hawhamburg.micr0bu.domain.detection.EventDetector] / [com.hawhamburg.micr0bu.domain.detection.DetectedEvent].
|
||||
*
|
||||
* CAM is broadcast autonomously and periodically by every V2X station (own OBU and any
|
||||
* nearby remote road users) — no application trigger is required, unlike DENM.
|
||||
*/
|
||||
data class Cam(
|
||||
/** Originating V2X station ID. Used to tell own CAMs from remote ones. */
|
||||
val stationId: Long,
|
||||
|
||||
/** Raw ETSI EN 302 637-2 stationType (see [StationType] for values this project cares about). */
|
||||
val stationType: Int,
|
||||
|
||||
/** Reference position latitude/longitude (WGS84 degrees). */
|
||||
val latitude: Double,
|
||||
val longitude: Double,
|
||||
|
||||
/** Ground speed, m/s. */
|
||||
val speedMps: Double,
|
||||
|
||||
/** Heading, degrees clockwise from true north, [0, 360). */
|
||||
val headingDeg: Double,
|
||||
|
||||
/**
|
||||
* Yaw rate, degrees per second, where available (positive = turning clockwise).
|
||||
* Carried by both `v2x/rx/obu_gnss` (ego, ~4 Hz) and `v2x-uca/output/json/cam` (remote).
|
||||
* A much more direct "is this road user turning" signal than heading-delta heuristics —
|
||||
* see [com.hawhamburg.micr0bu.domain.usecase.UseCaseDetectionEngine].
|
||||
*/
|
||||
val yawRateDps: Double? = null,
|
||||
|
||||
/** Optional drive direction flag (forward/backward), where available. */
|
||||
val driveDirection: Int? = null,
|
||||
|
||||
/** Optional vehicle dimensions, metres, where available (CAMv2 / extended fields). */
|
||||
val vehicleLengthM: Double? = null,
|
||||
val vehicleWidthM: Double? = null,
|
||||
|
||||
/** Optional longitudinal acceleration control field, m/s², where available. */
|
||||
val accelerationMps2: Double? = null,
|
||||
|
||||
/** Wall-clock ms this CAM was received/processed. */
|
||||
val timestamp: Long,
|
||||
|
||||
/** True if this CAM originated from the ego micrOBU itself, false if a remote road user. */
|
||||
val isOwn: Boolean,
|
||||
)
|
||||
@@ -0,0 +1,64 @@
|
||||
package com.hawhamburg.micr0bu.domain.cam
|
||||
|
||||
import org.json.JSONObject
|
||||
|
||||
/**
|
||||
* Parses the processed CAM JSON published by the consider it Use Case API on
|
||||
* `v2x-uca/output/json/cam` (CI-CiT-MQTT_API_Documentation-v6) into a [Cam].
|
||||
*
|
||||
* **Field-name tolerance:** confirmed field names are tried first (`stationId`, `heading_deg`,
|
||||
* `speed_mps`, `yawRate_dps`, GeoJSON `position`), with several other plausible spellings/
|
||||
* nestings as fallbacks — see [JsonFieldReader]. Once the exact schema is fully confirmed
|
||||
* against real OBU payloads (Phase 02 bench test, Section 10.3), the fallbacks that never hit
|
||||
* can be deleted.
|
||||
*/
|
||||
object CamParser {
|
||||
|
||||
/**
|
||||
* @param json raw MQTT payload string from `v2x-uca/output/json/cam`.
|
||||
* @param ownStationId the ego OBU's own station ID, if known yet (from `v2x/rx/obu_gnss` —
|
||||
* see [ObuGnssParser] / `CamUseCaseRepository`). May be null before the first
|
||||
* obu_gnss message arrives, in which case [Cam.isOwn] falls back to any explicit
|
||||
* own/self flag in the payload, or false.
|
||||
* @param timestamp wall-clock ms to stamp the parsed [Cam] with.
|
||||
* @return the parsed [Cam], or null if the payload isn't a recognisable CAM message.
|
||||
*/
|
||||
fun parse(json: String, ownStationId: Long?, timestamp: Long = System.currentTimeMillis()): Cam? {
|
||||
val obj = runCatching { JSONObject(json) }.getOrNull() ?: return null
|
||||
|
||||
val stationId = JsonFieldReader.firstLong(obj, "stationId", "stationID", "station_id") ?: return null
|
||||
val stationType = JsonFieldReader.firstInt(obj, "stationType", "station_type") ?: return null
|
||||
|
||||
val (lat, lon) = JsonFieldReader.firstLatLon(obj) ?: return null
|
||||
|
||||
val speed = JsonFieldReader.firstDouble(obj, "speed_mps", "speed", "speedMps") ?: 0.0
|
||||
val heading = JsonFieldReader.firstDouble(obj, "heading_deg", "heading", "headingDeg") ?: 0.0
|
||||
val yawRate = JsonFieldReader.firstDouble(obj, "yawRate_dps", "yawRateDps", "yaw_rate_dps", "yawRate")
|
||||
|
||||
val driveDirection = JsonFieldReader.firstInt(obj, "driveDirection", "drive_direction")
|
||||
val vehicleLength = JsonFieldReader.firstDouble(obj, "vehicleLength", "vehicle_length")
|
||||
val vehicleWidth = JsonFieldReader.firstDouble(obj, "vehicleWidth", "vehicle_width")
|
||||
val accel = JsonFieldReader.firstDouble(obj, "accelerationControl", "longitudinalAcceleration", "acceleration_mps2")
|
||||
|
||||
// Own/remote: prefer an explicit flag if the API provides one; otherwise compare
|
||||
// against the ego station ID learned from v2x/rx/obu_gnss.
|
||||
val explicitOwn = obj.optBoolean("own", obj.optBoolean("isOwn", false))
|
||||
val isOwn = explicitOwn || (ownStationId != null && ownStationId == stationId)
|
||||
|
||||
return Cam(
|
||||
stationId = stationId,
|
||||
stationType = stationType,
|
||||
latitude = lat,
|
||||
longitude = lon,
|
||||
speedMps = speed,
|
||||
headingDeg = JsonFieldReader.normaliseHeadingDeg(heading),
|
||||
yawRateDps = yawRate,
|
||||
driveDirection = driveDirection,
|
||||
vehicleLengthM = vehicleLength,
|
||||
vehicleWidthM = vehicleWidth,
|
||||
accelerationMps2 = accel,
|
||||
timestamp = timestamp,
|
||||
isOwn = isOwn,
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,80 @@
|
||||
package com.hawhamburg.micr0bu.domain.cam
|
||||
|
||||
import org.json.JSONObject
|
||||
|
||||
/**
|
||||
* Shared defensive JSON field readers used by [CamParser] and [ObuGnssParser].
|
||||
*
|
||||
* Both `v2x-uca/output/json/cam` and `v2x/rx/obu_gnss` are processed JSON from the consider it
|
||||
* Use Case API, and the exact key spelling wasn't available when this was written — so every
|
||||
* field is looked up by trying several plausible spellings (confirmed ones first, e.g.
|
||||
* `heading_deg` / `speed_mps` / `yawRate_dps`), and unwraps either a bare numeric value or an
|
||||
* ETSI-style `{"value": ...}` wrapper.
|
||||
*/
|
||||
internal object JsonFieldReader {
|
||||
|
||||
/** Reads the first present key as a Long, unwrapping a `{"value": n}` object if needed. */
|
||||
fun firstLong(obj: JSONObject, vararg keys: String): Long? {
|
||||
for (key in keys) {
|
||||
if (!obj.has(key)) continue
|
||||
unwrapValue(obj.opt(key))?.let { return it.toLong() }
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
fun firstInt(obj: JSONObject, vararg keys: String): Int? =
|
||||
firstLong(obj, *keys)?.toInt()
|
||||
|
||||
fun firstDouble(obj: JSONObject, vararg keys: String): Double? {
|
||||
for (key in keys) {
|
||||
if (!obj.has(key)) continue
|
||||
unwrapValue(obj.opt(key))?.let { return it }
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
/**
|
||||
* Reads a WGS84 lat/lon pair, trying (in order): a `referencePosition`/`reference_position`
|
||||
* sub-object with `lat`/`lon` keys, a GeoJSON-style `position` sub-object with
|
||||
* `{"type":"Point","coordinates":[lon,lat,alt]}` (the shape the codebase's own prior CAM-TX
|
||||
* scaffold used), or bare `lat`/`lon` keys on [obj] itself.
|
||||
*/
|
||||
fun firstLatLon(obj: JSONObject): Pair<Double, Double>? {
|
||||
val posObj = obj.optJSONObject("referencePosition") ?: obj.optJSONObject("reference_position")
|
||||
if (posObj != null) {
|
||||
val lat = firstDouble(posObj, "lat", "latitude")
|
||||
val lon = firstDouble(posObj, "lon", "lng", "longitude")
|
||||
if (lat != null && lon != null) return lat to lon
|
||||
}
|
||||
|
||||
val geoJson = obj.optJSONObject("position")
|
||||
val coords = geoJson?.optJSONArray("coordinates")
|
||||
if (coords != null && coords.length() >= 2) {
|
||||
// GeoJSON coordinate order is [longitude, latitude, altitude?]
|
||||
val lon = coords.optDouble(0, Double.NaN)
|
||||
val lat = coords.optDouble(1, Double.NaN)
|
||||
if (!lat.isNaN() && !lon.isNaN()) return lat to lon
|
||||
}
|
||||
|
||||
val lat = firstDouble(obj, "lat", "latitude")
|
||||
val lon = firstDouble(obj, "lon", "lng", "longitude")
|
||||
if (lat != null && lon != null) return lat to lon
|
||||
|
||||
return null
|
||||
}
|
||||
|
||||
/** Normalises a heading/bearing in degrees to [0, 360). */
|
||||
fun normaliseHeadingDeg(deg: Double): Double {
|
||||
var h = deg % 360.0
|
||||
if (h < 0) h += 360.0
|
||||
return h
|
||||
}
|
||||
|
||||
/** Unwraps a raw Number, a numeric String, or a `{"value": n}` object into a Double. */
|
||||
private fun unwrapValue(v: Any?): Double? = when (v) {
|
||||
is Number -> v.toDouble()
|
||||
is String -> v.toDoubleOrNull()
|
||||
is JSONObject -> if (v.has("value")) unwrapValue(v.opt("value")) else null
|
||||
else -> null
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
package com.hawhamburg.micr0bu.domain.cam
|
||||
|
||||
import org.json.JSONObject
|
||||
|
||||
/**
|
||||
* Parses the ego OBU's own kinematic state from `v2x/rx/obu_gnss`.
|
||||
*
|
||||
* This topic updates at ~4 Hz and carries the ego bike's own position, speed, heading, yaw
|
||||
* rate, and stationType — a higher-rate, more complete source for the ego side of
|
||||
* [com.hawhamburg.micr0bu.domain.usecase.UseCaseDetectionEngine]'s evaluation than the lower
|
||||
* rate (1-10 Hz) "own" entry that also appears on `v2x-uca/output/json/cam`. See
|
||||
* `CamUseCaseRepository`, which prefers this source and falls back to the CAM-topic "own"
|
||||
* entry (or phone GNSS) only when this one goes stale.
|
||||
*
|
||||
* Same field-name-tolerance caveat as [CamParser] — see [JsonFieldReader].
|
||||
*/
|
||||
object ObuGnssParser {
|
||||
|
||||
/**
|
||||
* @param json raw MQTT payload string from `v2x/rx/obu_gnss`.
|
||||
* @param fallbackStationId used if the payload doesn't carry its own station ID (rare —
|
||||
* own_info normally does, but keeps this robust to a partial payload).
|
||||
* @return the ego [Cam] (always [Cam.isOwn] == true), or null if position is missing.
|
||||
*/
|
||||
fun parseEgo(json: String, fallbackStationId: Long = 0L, timestamp: Long = System.currentTimeMillis()): Cam? {
|
||||
val root = runCatching { JSONObject(json) }.getOrNull() ?: return null
|
||||
val obj = root.optJSONObject("own_info") ?: root
|
||||
|
||||
val stationId = JsonFieldReader.firstLong(obj, "stationID", "stationId", "station_id") ?: fallbackStationId
|
||||
val stationType = JsonFieldReader.firstInt(obj, "stationType", "station_type") ?: StationType.CYCLIST
|
||||
|
||||
val (lat, lon) = JsonFieldReader.firstLatLon(obj) ?: return null
|
||||
|
||||
val speed = JsonFieldReader.firstDouble(obj, "speed_mps", "speed", "speedMps") ?: 0.0
|
||||
val heading = JsonFieldReader.firstDouble(obj, "heading_deg", "heading", "headingDeg") ?: 0.0
|
||||
val yawRate = JsonFieldReader.firstDouble(obj, "yawRate_dps", "yawRateDps", "yaw_rate_dps", "yawRate")
|
||||
|
||||
return Cam(
|
||||
stationId = stationId,
|
||||
stationType = stationType,
|
||||
latitude = lat,
|
||||
longitude = lon,
|
||||
speedMps = speed,
|
||||
headingDeg = JsonFieldReader.normaliseHeadingDeg(heading),
|
||||
yawRateDps = yawRate,
|
||||
timestamp = timestamp,
|
||||
isOwn = true,
|
||||
)
|
||||
}
|
||||
|
||||
/** Just the stationID/stationType, for when only identity (not a full fix) is needed. */
|
||||
fun parseOwnIdentity(json: String): Pair<Long, Int>? {
|
||||
val root = runCatching { JSONObject(json) }.getOrNull() ?: return null
|
||||
val obj = root.optJSONObject("own_info") ?: return null
|
||||
val stationId = JsonFieldReader.firstLong(obj, "stationID", "stationId", "station_id") ?: return null
|
||||
val stationType = JsonFieldReader.firstInt(obj, "stationType", "station_type") ?: return null
|
||||
return stationId to stationType
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,29 @@
|
||||
package com.hawhamburg.micr0bu.domain.denm
|
||||
|
||||
/**
|
||||
* MQTT topic for the consider it Use Case API control messages.
|
||||
*
|
||||
* **Manual/antenna-test path only.** This project's CAM-based use case architecture
|
||||
* (`com.hawhamburg.micr0bu.domain.usecase`) never generates or consumes DENM — see
|
||||
* Section 0.2 / 10.5 of the requirements doc. The DENM trigger below is kept solely as a
|
||||
* manual test tool: it lets a tester fire an OBU→RSU DENM broadcast on demand to verify the
|
||||
* antennas/ITS-G5 link are actually working end-to-end, independent of any detected event.
|
||||
*/
|
||||
const val DENM_CTRL_TOPIC = "v2x-uca/input/denmtrg"
|
||||
|
||||
/**
|
||||
* DENM use cases triggerable via [DENM_CTRL_TOPIC].
|
||||
*
|
||||
* Only one use case may be active on the OBU at a time (enforced in MqttRepository).
|
||||
*/
|
||||
enum class DenmUseCase(val id: String) {
|
||||
/** Electronic Emergency Brake Light (CC99/1). */
|
||||
EEBL("c2c-eebl"),
|
||||
|
||||
/** Aftermarket Stationary Recovery Vehicle (CC94/0). Default manual test case. */
|
||||
STATIONARY("hln-sv"),
|
||||
}
|
||||
|
||||
/** Builds a uca-denmctrl JSON control payload. */
|
||||
fun buildDenmPayload(useCase: String, active: Boolean): String =
|
||||
"""{"type":"uca-denmctrl","active":$active,"usecase":"$useCase","params":{}}"""
|
||||
@@ -0,0 +1,21 @@
|
||||
package com.hawhamburg.micr0bu.domain.usecase
|
||||
|
||||
/**
|
||||
* The C2C-CC three-tier alert level model (requirements doc Section 5.5).
|
||||
*
|
||||
* Governs both HMI presentation (colour, audio/haptic intensity) and the detection engine's
|
||||
* time-to-conflict thresholds. All in-scope use cases are designed to reach at least
|
||||
* [AWARENESS] under current BSP1-compliant CAM accuracy; [WARNING] additionally requires
|
||||
* lane-level positioning accuracy where achievable. The engine supports both tiers from the
|
||||
* start so the HMI upgrades automatically as positioning accuracy improves.
|
||||
*/
|
||||
enum class AlertLevel {
|
||||
/** Danger present but not imminent. Example threshold: > 10-30 s to conflict. */
|
||||
INFO,
|
||||
|
||||
/** Achievable without lane-level accuracy. Example threshold: > 3-15 s to conflict. */
|
||||
AWARENESS,
|
||||
|
||||
/** Imminent danger. Example threshold: <= 5 s to conflict. */
|
||||
WARNING,
|
||||
}
|
||||
@@ -0,0 +1,128 @@
|
||||
package com.hawhamburg.micr0bu.domain.usecase
|
||||
|
||||
import kotlin.math.abs
|
||||
import kotlin.math.atan2
|
||||
import kotlin.math.cos
|
||||
import kotlin.math.sin
|
||||
import kotlin.math.sqrt
|
||||
|
||||
/**
|
||||
* Small geodesy + 2D kinematics helpers used by [UseCaseDetectionEngine].
|
||||
*
|
||||
* Pure Kotlin, no Android dependencies, fully unit-testable — same pattern as
|
||||
* [com.hawhamburg.micr0bu.domain.detection.RunningStats].
|
||||
*/
|
||||
internal object GeoMath {
|
||||
|
||||
private const val EARTH_RADIUS_M = 6_371_000.0
|
||||
|
||||
/** A local-tangent-plane 2D vector: x = east (m), y = north (m), or an east/north velocity (m/s). */
|
||||
data class Vector2(val x: Double, val y: Double) {
|
||||
operator fun minus(other: Vector2) = Vector2(x - other.x, y - other.y)
|
||||
operator fun plus(other: Vector2) = Vector2(x + other.x, y + other.y)
|
||||
operator fun times(scalar: Double) = Vector2(x * scalar, y * scalar)
|
||||
fun dot(other: Vector2) = x * other.x + y * other.y
|
||||
fun length() = sqrt(x * x + y * y)
|
||||
}
|
||||
|
||||
/** Great-circle distance between two WGS84 points, metres (haversine). */
|
||||
fun haversineMeters(lat1: Double, lon1: Double, lat2: Double, lon2: Double): Double {
|
||||
val phi1 = Math.toRadians(lat1)
|
||||
val phi2 = Math.toRadians(lat2)
|
||||
val dPhi = Math.toRadians(lat2 - lat1)
|
||||
val dLambda = Math.toRadians(lon2 - lon1)
|
||||
val a = sin(dPhi / 2).let { it * it } +
|
||||
cos(phi1) * cos(phi2) * sin(dLambda / 2).let { it * it }
|
||||
val c = 2 * atan2(sqrt(a), sqrt(1 - a))
|
||||
return EARTH_RADIUS_M * c
|
||||
}
|
||||
|
||||
/** Initial bearing from point 1 to point 2, degrees clockwise from true north, [0, 360). */
|
||||
fun initialBearingDeg(lat1: Double, lon1: Double, lat2: Double, lon2: Double): Double {
|
||||
val phi1 = Math.toRadians(lat1)
|
||||
val phi2 = Math.toRadians(lat2)
|
||||
val dLambda = Math.toRadians(lon2 - lon1)
|
||||
val y = sin(dLambda) * cos(phi2)
|
||||
val x = cos(phi1) * sin(phi2) - sin(phi1) * cos(phi2) * cos(dLambda)
|
||||
return normalizeAngle(Math.toDegrees(atan2(y, x)))
|
||||
}
|
||||
|
||||
/** Normalises an angle in degrees to [0, 360). */
|
||||
fun normalizeAngle(deg: Double): Double {
|
||||
var a = deg % 360.0
|
||||
if (a < 0) a += 360.0
|
||||
return a
|
||||
}
|
||||
|
||||
/** Smallest absolute angular difference between two bearings/headings, in [0, 180]. */
|
||||
fun angleDiffDeg(a: Double, b: Double): Double {
|
||||
val diff = abs(normalizeAngle(a) - normalizeAngle(b)) % 360.0
|
||||
return if (diff > 180.0) 360.0 - diff else diff
|
||||
}
|
||||
|
||||
/**
|
||||
* Signed angular difference `to - from`, in (-180, 180]. Positive means [to] is clockwise
|
||||
* of [from] (e.g. if [from] is a heading and [to] is a bearing, positive = target is to the
|
||||
* right). Used to test whether a yaw-rate direction is rotating a heading *toward* a bearing.
|
||||
*/
|
||||
fun angleDiffSigned(from: Double, to: Double): Double {
|
||||
var diff = (normalizeAngle(to) - normalizeAngle(from)) % 360.0
|
||||
if (diff > 180.0) diff -= 360.0
|
||||
if (diff <= -180.0) diff += 360.0
|
||||
return diff
|
||||
}
|
||||
|
||||
/**
|
||||
* Projects a point [lat]/[lon] onto a local east/north tangent plane centred at
|
||||
* [lat0]/[lon0], in metres. Equirectangular approximation — accurate for the short
|
||||
* (sub-kilometre) ranges relevant to V2X.
|
||||
*/
|
||||
fun toLocalMeters(lat0: Double, lon0: Double, lat: Double, lon: Double): Vector2 {
|
||||
val phi0 = Math.toRadians(lat0)
|
||||
val dLat = Math.toRadians(lat - lat0)
|
||||
val dLon = Math.toRadians(lon - lon0)
|
||||
val north = dLat * EARTH_RADIUS_M
|
||||
val east = dLon * EARTH_RADIUS_M * cos(phi0)
|
||||
return Vector2(east, north)
|
||||
}
|
||||
|
||||
/** East/north velocity vector (m/s) from speed (m/s) and heading (deg clockwise from north). */
|
||||
fun velocityVector(speedMps: Double, headingDeg: Double): Vector2 {
|
||||
val rad = Math.toRadians(headingDeg)
|
||||
return Vector2(x = speedMps * sin(rad), y = speedMps * cos(rad))
|
||||
}
|
||||
|
||||
/**
|
||||
* Closest point of approach between two converging tracks, given the relative position
|
||||
* ([relPos] = other - self, metres) and relative velocity ([relVel] = other's velocity
|
||||
* minus self's velocity, m/s).
|
||||
*
|
||||
* @return time to CPA (s, clamped to [0, maxHorizonSec]) and the separation distance (m)
|
||||
* at that time. If the tracks are not closing (or [relVel] is ~stationary), time
|
||||
* is reported as 0 and distance as the current separation.
|
||||
*/
|
||||
fun closestPointOfApproach(
|
||||
relPos: Vector2,
|
||||
relVel: Vector2,
|
||||
maxHorizonSec: Double,
|
||||
): Pair<Double, Double> {
|
||||
val vSq = relVel.dot(relVel)
|
||||
if (vSq < 1e-6) return 0.0 to relPos.length()
|
||||
|
||||
val tRaw = -relPos.dot(relVel) / vSq
|
||||
val t = tRaw.coerceIn(0.0, maxHorizonSec)
|
||||
val posAtT = relPos + relVel * t
|
||||
return t to posAtT.length()
|
||||
}
|
||||
|
||||
/**
|
||||
* Rate of closure (m/s, positive = closing) of [relVel] (other's velocity minus self's)
|
||||
* along the current line of sight [relPos] (other - self, metres).
|
||||
*/
|
||||
fun closingSpeed(relPos: Vector2, relVel: Vector2): Double {
|
||||
val dist = relPos.length()
|
||||
if (dist < 1e-6) return 0.0
|
||||
val unit = Vector2(relPos.x / dist, relPos.y / dist)
|
||||
return -unit.dot(relVel)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
package com.hawhamburg.micr0bu.domain.usecase
|
||||
|
||||
/**
|
||||
* A currently-active use case alert produced by [UseCaseDetectionEngine], evaluating one
|
||||
* remote CAM against the ego OBU's own most recent CAM.
|
||||
*
|
||||
* Pure domain model — no Android imports.
|
||||
*/
|
||||
data class UseCaseAlert(
|
||||
val useCase: UseCaseType,
|
||||
val alertLevel: AlertLevel,
|
||||
|
||||
/** stationID of the remote road user contributing to this alert. */
|
||||
val remoteStationId: Long,
|
||||
|
||||
/** Raw ETSI stationType of the remote road user (see [com.hawhamburg.micr0bu.domain.cam.StationType]). */
|
||||
val remoteStationType: Int,
|
||||
|
||||
/** Great-circle distance between ego and remote, metres. */
|
||||
val distanceMeters: Double,
|
||||
|
||||
/** Rate of closure along the line of sight, m/s. Positive = closing. */
|
||||
val closingSpeedMps: Double,
|
||||
|
||||
/** Estimated time to closest point of approach / conflict point, seconds. */
|
||||
val timeToConflictSec: Double,
|
||||
|
||||
/**
|
||||
* Which signal(s) contributed to this alert, for debugging/verification — e.g.
|
||||
* "yaw-rate", "heading-trend", "static-heuristic". Not shown as the primary UI text, but
|
||||
* useful to confirm the engine is actually using per-station history/yaw rate rather than
|
||||
* a single-sample heuristic.
|
||||
*/
|
||||
val signalNote: String,
|
||||
|
||||
/** Wall-clock ms this alert was last (re)computed. */
|
||||
val lastUpdated: Long,
|
||||
)
|
||||
@@ -0,0 +1,111 @@
|
||||
package com.hawhamburg.micr0bu.domain.usecase
|
||||
|
||||
/**
|
||||
* All CAM-based use case detection thresholds in one place.
|
||||
*
|
||||
* As with [com.hawhamburg.micr0bu.domain.detection.DetectionConfig], these are **initial
|
||||
* engineering estimates** — the C2C-CC White Paper's exact geometric/kinematic definitions
|
||||
* were not available when this was written, only the plain-language use case descriptions in
|
||||
* Section 1.2 of the requirements doc. Validate and tune against real CAM traffic from the
|
||||
* Phase 02 bench test (Section 10.3) and, ultimately, real test-intersection data.
|
||||
*
|
||||
* Notably, baseline CAM carries no turn-signal/indicator field, so RTW-B/LTW-B are
|
||||
* approximated here via parallel-heading + lateral-sector + closing-distance heuristics
|
||||
* rather than true turn-intent detection. Precision should improve once CAMv2
|
||||
* exterior-lights/indicator fields are consumed.
|
||||
*/
|
||||
data class UseCaseDetectionConfig(
|
||||
|
||||
// ── Distance gates ────────────────────────────────────────────────────────
|
||||
/** Max distance (m) considered for intersection-scale use cases (IMA-B, IMA-S, RTW-B, LTW-B). */
|
||||
val intersectionRadiusM: Double = 40.0,
|
||||
|
||||
/** Max distance (m) considered for same-direction/rural use cases (SMVA/BCW-B). */
|
||||
val roadwayRadiusM: Double = 120.0,
|
||||
|
||||
/** Beyond this distance (m), skip evaluation entirely (cheap early-out). */
|
||||
val maxConsiderationRadiusM: Double = 150.0,
|
||||
|
||||
/** Closest-point-of-approach distance (m) below which paths are considered "in conflict". */
|
||||
val conflictRadiusM: Double = 5.0,
|
||||
|
||||
/** Slightly wider CPA tolerance (m) for the turn-warning use cases, which model an
|
||||
* as-yet-unexecuted turn rather than the vehicle's current heading. */
|
||||
val turnConflictRadiusM: Double = 7.5,
|
||||
|
||||
// ── Heading / bearing gates ──────────────────────────────────────────────
|
||||
/** Heading delta (deg) range considered "crossing paths" (near-perpendicular). */
|
||||
val crossingHeadingMinDeg: Double = 50.0,
|
||||
val crossingHeadingMaxDeg: Double = 130.0,
|
||||
|
||||
/** Heading delta (deg) at or below which two road users are considered travelling parallel. */
|
||||
val parallelHeadingMaxDeg: Double = 30.0,
|
||||
|
||||
/** Lateral-sector half-width (deg either side of dead-ahead/dead-astern) used for the
|
||||
* same-direction closing check in SMVA/BCW-B. */
|
||||
val sameDirectionSectorDeg: Double = 20.0,
|
||||
|
||||
// ── Speed gates ───────────────────────────────────────────────────────────
|
||||
/** Below this speed (m/s), a road user is considered stationary/standstill. */
|
||||
val standstillSpeedThresholdMps: Double = 0.5,
|
||||
|
||||
/** Minimum speed (m/s) for a road user to count as "moving" for crossing use cases. */
|
||||
val minMovingSpeedMps: Double = 1.0,
|
||||
|
||||
/** Minimum positive closing speed (m/s) along the line of sight to count as "approaching". */
|
||||
val closingSpeedMinMps: Double = 0.3,
|
||||
|
||||
/** Minimum speed differential (m/s) between ego and remote for SMVA/BCW-B to trigger. */
|
||||
val speedDifferentialMinMps: Double = 3.0,
|
||||
|
||||
// ── Time-to-conflict → alert level thresholds (Section 5.5) ──────────────
|
||||
/** TTC (s) at or below which the alert escalates to WARNING. */
|
||||
val ttcWarningSec: Double = 5.0,
|
||||
|
||||
/** TTC (s) at or below which the alert is at least AWARENESS. */
|
||||
val ttcAwarenessSec: Double = 15.0,
|
||||
|
||||
/** TTC (s) at or below which the alert is at least INFO. Beyond this, no alert is raised. */
|
||||
val ttcInfoSec: Double = 30.0,
|
||||
|
||||
/** Cap (s) on the closest-point-of-approach time projection — road users converging
|
||||
* further out than this are not yet considered for an alert. */
|
||||
val maxProjectionHorizonSec: Double = 30.0,
|
||||
|
||||
// ── Staleness ─────────────────────────────────────────────────────────────
|
||||
/** Drop a remote CAM (and any alerts derived from it) if no update arrives within this
|
||||
* many ms — CAMs typically arrive at 1-10 Hz, so a multi-second gap means the remote
|
||||
* road user is out of range or the link dropped. */
|
||||
val staleRemoteMs: Long = 3_000L,
|
||||
|
||||
// ── Per-station history / trend (not reacting to a single message in isolation) ───
|
||||
/** Max number of past CAMs retained per remote station for trend computation. */
|
||||
val historyMaxSamples: Int = 8,
|
||||
|
||||
/** Max age (ms) of a history sample before it's dropped from the trend window. */
|
||||
val historyMaxAgeMs: Long = 5_000L,
|
||||
|
||||
/** Minimum samples in a station's history before trend (heading-rate/distance-rate) is
|
||||
* trusted; below this, classification falls back to the instantaneous-geometry heuristics
|
||||
* only. */
|
||||
val minHistorySamplesForTrend: Int = 2,
|
||||
|
||||
// ── Turn-toward-ego detection (yaw rate / heading-trend based) ────────────
|
||||
/** Minimum |yaw rate| (deg/s, from the CAM field if present, else derived from heading
|
||||
* history) to consider a remote road user "actively turning". */
|
||||
val turnYawRateThresholdDegPerSec: Double = 8.0,
|
||||
|
||||
/** Max |signed angle from the remote's heading to the ego's bearing| (deg) for the ego to
|
||||
* count as roughly "in front of" the turning remote — beyond this the remote is turning
|
||||
* away from/behind the ego, not toward it. */
|
||||
val turnTowardEgoMaxAngleDeg: Double = 90.0,
|
||||
)
|
||||
|
||||
/** Buckets a time-to-conflict estimate into an [AlertLevel], or null if beyond all thresholds. */
|
||||
fun UseCaseDetectionConfig.alertLevelForTtc(ttcSec: Double): AlertLevel? = when {
|
||||
ttcSec < 0 -> null
|
||||
ttcSec <= ttcWarningSec -> AlertLevel.WARNING
|
||||
ttcSec <= ttcAwarenessSec -> AlertLevel.AWARENESS
|
||||
ttcSec <= ttcInfoSec -> AlertLevel.INFO
|
||||
else -> null
|
||||
}
|
||||
@@ -0,0 +1,278 @@
|
||||
package com.hawhamburg.micr0bu.domain.usecase
|
||||
|
||||
import com.hawhamburg.micr0bu.domain.cam.Cam
|
||||
import kotlinx.coroutines.flow.MutableStateFlow
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.asStateFlow
|
||||
import kotlin.math.abs
|
||||
import kotlin.math.sign
|
||||
|
||||
/**
|
||||
* CAM-based Use Case Detection Engine (requirements doc Section 4.2 / 10.2).
|
||||
*
|
||||
* Continuously correlates the ego bike's own latest state (from `v2x/rx/obu_gnss`, ~4 Hz —
|
||||
* see `CamUseCaseRepository`) with a short **history** of CAMs received per remote road user
|
||||
* to evaluate the geometric/kinematic conditions for each in-scope C2C-CC bike safety use case
|
||||
* ([UseCaseType]), and raises alerts per the three-tier model ([AlertLevel], Section 5.5).
|
||||
*
|
||||
* Reacting to a single CAM in isolation is noisy (GNSS jitter, one-off heading glitches), so
|
||||
* each remote station's recent samples are kept in [remoteHistory] and used to derive a
|
||||
* short-term trend (heading-rate, distance-rate) that corroborates or substitutes for
|
||||
* instantaneous fields — notably yaw rate, which not all CAM sources carry.
|
||||
*
|
||||
* **No DENM is generated or consumed here.** The OBU's own autonomous CAM broadcast, received
|
||||
* by the other road user's OBU, already carries the ego vehicle's presence — this engine only
|
||||
* needs to *consume* CAM (and the ego's own obu_gnss state) to decide when to surface a warning
|
||||
* locally.
|
||||
*
|
||||
* Pure Kotlin, no Android imports — unit-testable with synthetic CAM input exactly like
|
||||
* [com.hawhamburg.micr0bu.domain.detection.EventDetector].
|
||||
*
|
||||
* Thread-safety: this class is not synchronized. Callers (see `CamUseCaseRepository`) should
|
||||
* confine calls to a single coroutine/thread, e.g. by collecting from a single Flow.
|
||||
*/
|
||||
class UseCaseDetectionEngine(private val config: UseCaseDetectionConfig = UseCaseDetectionConfig()) {
|
||||
|
||||
private var ownCam: Cam? = null
|
||||
|
||||
// Short history per remote station ID — oldest first, bounded by size and age
|
||||
// (config.historyMaxSamples / historyMaxAgeMs). This is what lets the engine compute
|
||||
// trends (heading-rate, distance-rate) instead of reacting to one message in isolation.
|
||||
private val remoteHistory = mutableMapOf<Long, ArrayDeque<Cam>>()
|
||||
|
||||
// Currently active alerts, keyed by (remote station ID, use case) so one remote road user
|
||||
// can concurrently contribute to more than one use case (e.g. IMA-B and SMVA/BCW-B).
|
||||
private val activeAlerts = mutableMapOf<Pair<Long, UseCaseType>, UseCaseAlert>()
|
||||
|
||||
private val _currentAlerts = MutableStateFlow<List<UseCaseAlert>>(emptyList())
|
||||
|
||||
/** Currently active alerts across all remote road users, most severe first. */
|
||||
val currentAlerts: StateFlow<List<UseCaseAlert>> = _currentAlerts.asStateFlow()
|
||||
|
||||
/**
|
||||
* Feed the ego bike's own most recent state (from `v2x/rx/obu_gnss`, or a fallback source
|
||||
* — see `CamUseCaseRepository`). Out-of-order/late updates are ignored. Re-evaluates all
|
||||
* tracked remote stations against the new state.
|
||||
*/
|
||||
fun onOwnCam(cam: Cam) {
|
||||
val current = ownCam
|
||||
if (current != null && cam.timestamp < current.timestamp) return // stale/out-of-order
|
||||
ownCam = cam
|
||||
remoteHistory.keys.toList().forEach { id -> remoteHistory[id]?.lastOrNull()?.let { evaluate(it) } }
|
||||
publish()
|
||||
}
|
||||
|
||||
/**
|
||||
* Feed a CAM from a remote road user. Pushed into that station's history, then evaluated
|
||||
* immediately against the last ego state, if any.
|
||||
*/
|
||||
fun onRemoteCam(cam: Cam) {
|
||||
val history = remoteHistory.getOrPut(cam.stationId) { ArrayDeque() }
|
||||
history.addLast(cam)
|
||||
trimHistory(history, cam.timestamp)
|
||||
evaluate(cam)
|
||||
publish()
|
||||
}
|
||||
|
||||
private fun trimHistory(history: ArrayDeque<Cam>, nowMs: Long) {
|
||||
while (history.size > config.historyMaxSamples) history.removeFirst()
|
||||
while (history.isNotEmpty() && nowMs - history.first().timestamp > config.historyMaxAgeMs) {
|
||||
history.removeFirst()
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Drops remote CAMs (and any alerts derived from them) that haven't been updated within
|
||||
* [UseCaseDetectionConfig.staleRemoteMs]. Call periodically (e.g. once per second) from a
|
||||
* ticker — this is what makes an alert "clear automatically once the geometry resolves" or
|
||||
* the remote road user goes out of range (Section 10.3 test procedure).
|
||||
*/
|
||||
fun pruneStale(nowMs: Long) {
|
||||
val staleIds = remoteHistory.filterValues { history ->
|
||||
val last = history.lastOrNull() ?: return@filterValues true
|
||||
nowMs - last.timestamp > config.staleRemoteMs
|
||||
}.keys
|
||||
if (staleIds.isEmpty()) return
|
||||
staleIds.forEach { id ->
|
||||
remoteHistory.remove(id)
|
||||
activeAlerts.keys.filter { it.first == id }.forEach { activeAlerts.remove(it) }
|
||||
}
|
||||
publish()
|
||||
}
|
||||
|
||||
/** Resets all state (e.g. on disconnect). */
|
||||
fun reset() {
|
||||
ownCam = null
|
||||
remoteHistory.clear()
|
||||
activeAlerts.clear()
|
||||
publish()
|
||||
}
|
||||
|
||||
// ── Trend (per-station history → deltas) ─────────────────────────────────
|
||||
|
||||
/** Trend derived from a remote station's history: how its heading/distance are changing. */
|
||||
private data class Trend(
|
||||
val headingRateDegPerSec: Double?,
|
||||
val distanceTrendMps: Double?,
|
||||
val sampleCount: Int,
|
||||
)
|
||||
|
||||
private fun computeTrend(history: ArrayDeque<Cam>, own: Cam): Trend {
|
||||
if (history.size < config.minHistorySamplesForTrend) return Trend(null, null, history.size)
|
||||
|
||||
val oldest = history.first()
|
||||
val newest = history.last()
|
||||
val dtSec = (newest.timestamp - oldest.timestamp) / 1000.0
|
||||
if (dtSec < 0.1) return Trend(null, null, history.size) // too little time elapsed to trust a rate
|
||||
|
||||
val headingRate = GeoMath.angleDiffSigned(oldest.headingDeg, newest.headingDeg) / dtSec
|
||||
|
||||
val distOld = GeoMath.haversineMeters(own.latitude, own.longitude, oldest.latitude, oldest.longitude)
|
||||
val distNew = GeoMath.haversineMeters(own.latitude, own.longitude, newest.latitude, newest.longitude)
|
||||
val distanceTrend = (distNew - distOld) / dtSec
|
||||
|
||||
return Trend(headingRate, distanceTrend, history.size)
|
||||
}
|
||||
|
||||
// ── Evaluation ────────────────────────────────────────────────────────────
|
||||
|
||||
private fun evaluate(remote: Cam) {
|
||||
val own = ownCam
|
||||
if (own == null) {
|
||||
// No ego state yet — nothing to correlate against.
|
||||
UseCaseType.entries.forEach { activeAlerts.remove(remote.stationId to it) }
|
||||
return
|
||||
}
|
||||
|
||||
val distance = GeoMath.haversineMeters(own.latitude, own.longitude, remote.latitude, remote.longitude)
|
||||
if (distance > config.maxConsiderationRadiusM) {
|
||||
UseCaseType.entries.forEach { activeAlerts.remove(remote.stationId to it) }
|
||||
return
|
||||
}
|
||||
|
||||
val history = remoteHistory[remote.stationId] ?: ArrayDeque<Cam>().also { it.addLast(remote) }
|
||||
val trend = computeTrend(history, own)
|
||||
|
||||
val relPos = GeoMath.toLocalMeters(own.latitude, own.longitude, remote.latitude, remote.longitude)
|
||||
val ownVel = GeoMath.velocityVector(own.speedMps, own.headingDeg)
|
||||
val remoteVel = GeoMath.velocityVector(remote.speedMps, remote.headingDeg)
|
||||
val relVel = remoteVel - ownVel
|
||||
|
||||
val (tCpa, dCpa) = GeoMath.closestPointOfApproach(relPos, relVel, config.maxProjectionHorizonSec)
|
||||
val closingSpeed = GeoMath.closingSpeed(relPos, relVel)
|
||||
val headingDelta = GeoMath.angleDiffDeg(own.headingDeg, remote.headingDeg)
|
||||
|
||||
// Bearing to the remote, relative to the ego's own heading: 0 = dead ahead,
|
||||
// 90 = right side, 180 = dead astern, 270 = left side.
|
||||
val bearingToRemote = GeoMath.initialBearingDeg(own.latitude, own.longitude, remote.latitude, remote.longitude)
|
||||
val relBearing = GeoMath.normalizeAngle(bearingToRemote - own.headingDeg)
|
||||
|
||||
// Corroborate instantaneous closing speed with the distance trend when we have enough
|
||||
// history to trust it; otherwise fall back to the single-sample closing speed alone.
|
||||
val isActuallyClosing = when {
|
||||
trend.distanceTrendMps != null -> trend.distanceTrendMps < 0.0
|
||||
else -> closingSpeed >= config.closingSpeedMinMps
|
||||
}
|
||||
|
||||
// ── Turn-toward-ego detection: prefer the CAM's own yaw rate field; fall back to the
|
||||
// heading-rate derived from this station's history when the field isn't present. ──
|
||||
val effectiveYawRateDps = remote.yawRateDps ?: trend.headingRateDegPerSec
|
||||
val yawSignalSource = when {
|
||||
remote.yawRateDps != null -> "yaw-rate"
|
||||
trend.headingRateDegPerSec != null -> "heading-trend"
|
||||
else -> null
|
||||
}
|
||||
val bearingFromRemoteToEgo = GeoMath.initialBearingDeg(remote.latitude, remote.longitude, own.latitude, own.longitude)
|
||||
val angleRemoteHeadingToEgo = GeoMath.angleDiffSigned(remote.headingDeg, bearingFromRemoteToEgo)
|
||||
val turningTowardEgo = effectiveYawRateDps != null &&
|
||||
abs(effectiveYawRateDps) >= config.turnYawRateThresholdDegPerSec &&
|
||||
abs(angleRemoteHeadingToEgo) <= config.turnTowardEgoMaxAngleDeg &&
|
||||
sign(effectiveYawRateDps) == sign(angleRemoteHeadingToEgo)
|
||||
|
||||
val results = mutableMapOf<UseCaseType, Pair<Double, String>>() // use case -> (ttc, signalNote)
|
||||
|
||||
// ── IMA-B: crossing paths near an intersection (primary use case, Section 1.1) ──
|
||||
if (headingDelta in config.crossingHeadingMinDeg..config.crossingHeadingMaxDeg &&
|
||||
distance <= config.intersectionRadiusM &&
|
||||
own.speedMps >= config.minMovingSpeedMps &&
|
||||
remote.speedMps >= config.minMovingSpeedMps &&
|
||||
dCpa <= config.conflictRadiusM &&
|
||||
tCpa > 0.0
|
||||
) {
|
||||
results[UseCaseType.IMA_B] = tCpa to "crossing-geometry"
|
||||
}
|
||||
|
||||
// ── IMA-S: standstill remote vehicle, ego approaching an intersection ───────────
|
||||
if (remote.speedMps < config.standstillSpeedThresholdMps &&
|
||||
own.speedMps >= config.minMovingSpeedMps &&
|
||||
distance <= config.intersectionRadiusM &&
|
||||
isActuallyClosing
|
||||
) {
|
||||
val ttc = if (tCpa > 0.0) tCpa else distance / closingSpeed.coerceAtLeast(config.closingSpeedMinMps)
|
||||
val note = if (trend.distanceTrendMps != null) "distance-trend" else "static-heuristic"
|
||||
results[UseCaseType.IMA_S] = ttc to note
|
||||
}
|
||||
|
||||
// ── RTW-B / LTW-B: turning toward the ego (yaw-rate/heading-trend) with a parallel-
|
||||
// heading + lateral-sector + closing-distance fallback for when no yaw signal exists. ──
|
||||
val turnSideGatesPass = distance <= config.intersectionRadiusM && isActuallyClosing
|
||||
if (turnSideGatesPass) {
|
||||
val staticHeuristicPasses = headingDelta <= config.parallelHeadingMaxDeg && dCpa <= config.turnConflictRadiusM
|
||||
val ttcFallback = { tCpa.takeIf { it > 0.0 } ?: (distance / closingSpeed.coerceAtLeast(config.closingSpeedMinMps)) }
|
||||
|
||||
if (relBearing in 0.0..90.0 && (turningTowardEgo || staticHeuristicPasses)) {
|
||||
val note = if (turningTowardEgo) (yawSignalSource ?: "static-heuristic") else "static-heuristic"
|
||||
results[UseCaseType.RTW_B] = ttcFallback() to note
|
||||
}
|
||||
if (relBearing in 270.0..360.0 && (turningTowardEgo || staticHeuristicPasses)) {
|
||||
val note = if (turningTowardEgo) (yawSignalSource ?: "static-heuristic") else "static-heuristic"
|
||||
results[UseCaseType.LTW_B] = ttcFallback() to note
|
||||
}
|
||||
}
|
||||
|
||||
// ── SMVA / BCW-B: same-direction closing speed (rural / lateral accident pattern) ──
|
||||
val sameDirectionSector = relBearing <= config.sameDirectionSectorDeg ||
|
||||
relBearing >= 360.0 - config.sameDirectionSectorDeg ||
|
||||
(relBearing - 180.0) in -config.sameDirectionSectorDeg..config.sameDirectionSectorDeg
|
||||
if (headingDelta <= config.parallelHeadingMaxDeg &&
|
||||
distance <= config.roadwayRadiusM &&
|
||||
sameDirectionSector &&
|
||||
abs(own.speedMps - remote.speedMps) >= config.speedDifferentialMinMps &&
|
||||
isActuallyClosing
|
||||
) {
|
||||
val ttc = if (tCpa > 0.0) tCpa else distance / closingSpeed.coerceAtLeast(config.closingSpeedMinMps)
|
||||
val note = if (trend.distanceTrendMps != null) "distance-trend" else "static-heuristic"
|
||||
results[UseCaseType.SMVA_BCW_B] = ttc to note
|
||||
}
|
||||
|
||||
// ── Publish / clear per use case ─────────────────────────────────────────────
|
||||
UseCaseType.entries.forEach { type ->
|
||||
val key = remote.stationId to type
|
||||
val result = results[type]
|
||||
val level = result?.let { config.alertLevelForTtc(it.first) }
|
||||
if (result != null && level != null) {
|
||||
activeAlerts[key] = UseCaseAlert(
|
||||
useCase = type,
|
||||
alertLevel = level,
|
||||
remoteStationId = remote.stationId,
|
||||
remoteStationType = remote.stationType,
|
||||
distanceMeters = distance,
|
||||
closingSpeedMps = closingSpeed,
|
||||
timeToConflictSec = result.first,
|
||||
signalNote = result.second,
|
||||
lastUpdated = remote.timestamp,
|
||||
)
|
||||
} else {
|
||||
activeAlerts.remove(key)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun publish() {
|
||||
_currentAlerts.value = activeAlerts.values
|
||||
.sortedWith(
|
||||
compareBy<UseCaseAlert> { it.alertLevel.ordinal * -1 }
|
||||
.thenBy { it.timeToConflictSec }
|
||||
)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
package com.hawhamburg.micr0bu.domain.usecase
|
||||
|
||||
/**
|
||||
* The C2C-CC Bicycle Safety Use Cases (White Paper C2CCC_WP_2324, v1.0, 2026-05-07) that are
|
||||
* in scope for this project's CAM-only architecture (Section 1.2 of the requirements doc).
|
||||
*
|
||||
* All five are achievable from CAM/CAMv2 exchange alone (position, speed, heading, station
|
||||
* type) — no DENM generation or consumption is involved.
|
||||
*
|
||||
* Deliberately **not** included: Bike Accident Warning (BAW, UC_BIKE_00010) — it depends on
|
||||
* DENM generation by a fallen cyclist's OBU, which this project does not implement. See
|
||||
* `com.hawhamburg.micr0bu.domain.denm` for the (unrelated, manual/test-only) DENM trigger.
|
||||
*/
|
||||
enum class UseCaseType(
|
||||
/** Short label as used in the C2C-CC white paper and this project's requirements doc. */
|
||||
val label: String,
|
||||
/** C2C-CC use case ID(s). */
|
||||
val c2cId: String,
|
||||
) {
|
||||
/** Intersection Movement Assist for bikes — primary Use Case 1 (Section 1.1). */
|
||||
IMA_B("IMA-B", "UC_BIKE_00001/2"),
|
||||
|
||||
/** Intersection Movement Assist with Standstill Vehicle. */
|
||||
IMA_S("IMA-S", "UC_BIKE_00003"),
|
||||
|
||||
/** Right-turn Warning for bike. */
|
||||
RTW_B("RTW-B", "UC_BIKE_00004/5"),
|
||||
|
||||
/** Left-turn Warning for bike. */
|
||||
LTW_B("LTW-B", "UC_BIKE_00006/7"),
|
||||
|
||||
/** Slow Moving Vehicle Ahead / Backward Collision Warning for bike. */
|
||||
SMVA_BCW_B("SMVA/BCW-B", "UC_BIKE_00008/9"),
|
||||
}
|
||||
Reference in New Issue
Block a user