diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..15d0e52 --- /dev/null +++ b/.gitignore @@ -0,0 +1,19 @@ +*.iml +.gradle/ +/local.properties +/.idea/caches +/.idea/libraries +/.idea/modules.xml +/.idea/workspace.xml +/.idea/navEditor.xml +/.idea/assetWizardSettings.xml +.DS_Store +build/ +captures/ +.externalNativeBuild +.cxx +local.properties +*.hprof +*.jks +*.keystore +secrets.properties diff --git a/.idea/.gitignore b/.idea/.gitignore new file mode 100644 index 0000000..26d3352 --- /dev/null +++ b/.idea/.gitignore @@ -0,0 +1,3 @@ +# Default ignored files +/shelf/ +/workspace.xml diff --git a/.idea/AndroidProjectSystem.xml b/.idea/AndroidProjectSystem.xml new file mode 100644 index 0000000..4a53bee --- /dev/null +++ b/.idea/AndroidProjectSystem.xml @@ -0,0 +1,6 @@ + + + + + \ No newline at end of file diff --git a/.idea/appInsightsSettings.xml b/.idea/appInsightsSettings.xml new file mode 100644 index 0000000..7f73dc8 --- /dev/null +++ b/.idea/appInsightsSettings.xml @@ -0,0 +1,6 @@ + + + + + \ No newline at end of file diff --git a/.idea/compiler.xml b/.idea/compiler.xml new file mode 100644 index 0000000..b86273d --- /dev/null +++ b/.idea/compiler.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/.idea/deploymentTargetSelector.xml b/.idea/deploymentTargetSelector.xml new file mode 100644 index 0000000..c07d1d1 --- /dev/null +++ b/.idea/deploymentTargetSelector.xml @@ -0,0 +1,31 @@ + + + + + + + + + \ No newline at end of file diff --git a/.idea/deviceManager.xml b/.idea/deviceManager.xml new file mode 100644 index 0000000..91f9558 --- /dev/null +++ b/.idea/deviceManager.xml @@ -0,0 +1,13 @@ + + + + + + \ No newline at end of file diff --git a/.idea/gradle.xml b/.idea/gradle.xml new file mode 100644 index 0000000..639c779 --- /dev/null +++ b/.idea/gradle.xml @@ -0,0 +1,19 @@ + + + + + + + \ No newline at end of file diff --git a/.idea/inspectionProfiles/Project_Default.xml b/.idea/inspectionProfiles/Project_Default.xml new file mode 100644 index 0000000..f0c6ad0 --- /dev/null +++ b/.idea/inspectionProfiles/Project_Default.xml @@ -0,0 +1,50 @@ + + + + \ No newline at end of file diff --git a/.idea/markdown.xml b/.idea/markdown.xml new file mode 100644 index 0000000..c61ea33 --- /dev/null +++ b/.idea/markdown.xml @@ -0,0 +1,8 @@ + + + + + + \ No newline at end of file diff --git a/.idea/migrations.xml b/.idea/migrations.xml new file mode 100644 index 0000000..f8051a6 --- /dev/null +++ b/.idea/migrations.xml @@ -0,0 +1,10 @@ + + + + + + \ No newline at end of file diff --git a/.idea/misc.xml b/.idea/misc.xml new file mode 100644 index 0000000..b2c751a --- /dev/null +++ b/.idea/misc.xml @@ -0,0 +1,9 @@ + + + + + + + + \ No newline at end of file diff --git a/.idea/runConfigurations.xml b/.idea/runConfigurations.xml new file mode 100644 index 0000000..16660f1 --- /dev/null +++ b/.idea/runConfigurations.xml @@ -0,0 +1,17 @@ + + + + + + \ No newline at end of file diff --git a/.idea/studiobot.xml b/.idea/studiobot.xml new file mode 100644 index 0000000..9298202 --- /dev/null +++ b/.idea/studiobot.xml @@ -0,0 +1,6 @@ + + + + + \ No newline at end of file diff --git a/.idea/vcs.xml b/.idea/vcs.xml new file mode 100644 index 0000000..94a25f7 --- /dev/null +++ b/.idea/vcs.xml @@ -0,0 +1,6 @@ + + + + + + \ No newline at end of file diff --git a/V2X_MicrOBU_Android_App_Requirements_v10.docx b/V2X_MicrOBU_Android_App_Requirements_v10.docx new file mode 100644 index 0000000..b2aa5dd Binary files /dev/null and b/V2X_MicrOBU_Android_App_Requirements_v10.docx differ diff --git a/app/src/main/java/com/hawhamburg/micr0bu/data/cam/CamUseCaseRepository.kt b/app/src/main/java/com/hawhamburg/micr0bu/data/cam/CamUseCaseRepository.kt new file mode 100644 index 0000000..5a5d8ec --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/data/cam/CamUseCaseRepository.kt @@ -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(null) + /** The ego OBU's own station ID, learned from `v2x/rx/obu_gnss`. Null until known. */ + val ownStationId: StateFlow = _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> = prefs.enabledMapFlow.stateIn( + scope, SharingStarted.Eagerly, UseCaseType.entries.associateWith { true }, + ) + + /** All currently active alerts, regardless of per-use-case enablement. */ + val allAlerts: StateFlow> = engine.currentAlerts + + /** Alerts filtered to only the use cases the user has enabled — what the UI should show. */ + val enabledAlerts: StateFlow> = 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) + } + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/UseCaseAlertPreferences.kt b/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/UseCaseAlertPreferences.kt new file mode 100644 index 0000000..1327364 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/data/mqtt/UseCaseAlertPreferences.kt @@ -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> = 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 } + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/Cam.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/Cam.kt new file mode 100644 index 0000000..404fd62 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/Cam.kt @@ -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, +) diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/CamParser.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/CamParser.kt new file mode 100644 index 0000000..477baf2 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/CamParser.kt @@ -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, + ) + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/JsonFieldReader.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/JsonFieldReader.kt new file mode 100644 index 0000000..d2e9a8c --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/JsonFieldReader.kt @@ -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? { + 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 + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/ObuGnssParser.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/ObuGnssParser.kt new file mode 100644 index 0000000..a1e7e09 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/cam/ObuGnssParser.kt @@ -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? { + 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 + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/denm/DenmUseCase.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/denm/DenmUseCase.kt new file mode 100644 index 0000000..b948885 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/denm/DenmUseCase.kt @@ -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":{}}""" diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/AlertLevel.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/AlertLevel.kt new file mode 100644 index 0000000..8e9015f --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/AlertLevel.kt @@ -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, +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/GeoMath.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/GeoMath.kt new file mode 100644 index 0000000..32c3e17 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/GeoMath.kt @@ -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 { + 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) + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseAlert.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseAlert.kt new file mode 100644 index 0000000..316d702 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseAlert.kt @@ -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, +) diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseDetectionConfig.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseDetectionConfig.kt new file mode 100644 index 0000000..0cc6383 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseDetectionConfig.kt @@ -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 +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseDetectionEngine.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseDetectionEngine.kt new file mode 100644 index 0000000..c3d47e8 --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseDetectionEngine.kt @@ -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>() + + // 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, UseCaseAlert>() + + private val _currentAlerts = MutableStateFlow>(emptyList()) + + /** Currently active alerts across all remote road users, most severe first. */ + val currentAlerts: StateFlow> = _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, 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, 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().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>() // 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 { it.alertLevel.ordinal * -1 } + .thenBy { it.timeToConflictSec } + ) + } +} diff --git a/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseType.kt b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseType.kt new file mode 100644 index 0000000..8b65d8c --- /dev/null +++ b/app/src/main/java/com/hawhamburg/micr0bu/domain/usecase/UseCaseType.kt @@ -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"), +}