2026-08-04 17:58:07 +02:00
|
|
|
package com.hawhamburg.micr0bu.service
|
|
|
|
|
|
|
|
|
|
import android.content.Context
|
|
|
|
|
import com.hawhamburg.micr0bu.data.GnssReading
|
2026-09-10 14:47:30 +02:00
|
|
|
import com.hawhamburg.micr0bu.data.GnssTimeSource
|
2026-08-04 17:58:07 +02:00
|
|
|
import com.hawhamburg.micr0bu.data.SensorRepository
|
2026-09-10 14:47:30 +02:00
|
|
|
import com.hawhamburg.micr0bu.data.cam.PseudonymManager
|
2026-08-04 17:58:07 +02:00
|
|
|
import com.hawhamburg.micr0bu.data.mqtt.ObuHardwarePreferences
|
2026-09-23 17:28:05 +02:00
|
|
|
import com.hawhamburg.micr0bu.data.transport.Esp32Link
|
2026-09-10 14:47:30 +02:00
|
|
|
import com.hawhamburg.micr0bu.data.transport.GnPositionVector
|
2026-08-04 17:58:07 +02:00
|
|
|
import com.hawhamburg.micr0bu.data.transport.ObuHardware
|
2026-09-23 17:28:05 +02:00
|
|
|
import com.hawhamburg.micr0bu.data.transport.OutgoingIts
|
|
|
|
|
import com.hawhamburg.micr0bu.data.transport.OutgoingMessage
|
2026-08-04 17:58:07 +02:00
|
|
|
import com.hawhamburg.micr0bu.domain.asn1.RealAsn1UperCodec
|
2026-09-23 17:28:05 +02:00
|
|
|
import com.hawhamburg.micr0bu.domain.asn1.VamContent
|
|
|
|
|
import com.hawhamburg.micr0bu.domain.asn1.VamUperCodec
|
2026-08-04 17:58:07 +02:00
|
|
|
import com.hawhamburg.micr0bu.domain.cam.CamTransmitConfig
|
|
|
|
|
import com.hawhamburg.micr0bu.domain.cam.PhoneCamBuilder
|
|
|
|
|
import com.hawhamburg.micr0bu.domain.usecase.GeoMath
|
2026-09-23 17:28:05 +02:00
|
|
|
import com.hawhamburg.micr0bu.domain.vam.VamGenerationRules
|
2026-08-04 17:58:07 +02:00
|
|
|
import dagger.hilt.android.qualifiers.ApplicationContext
|
|
|
|
|
import kotlinx.coroutines.CoroutineScope
|
|
|
|
|
import kotlinx.coroutines.Dispatchers
|
|
|
|
|
import kotlinx.coroutines.Job
|
|
|
|
|
import kotlinx.coroutines.SupervisorJob
|
|
|
|
|
import kotlinx.coroutines.coroutineScope
|
|
|
|
|
import kotlinx.coroutines.delay
|
|
|
|
|
import kotlinx.coroutines.flow.collectLatest
|
|
|
|
|
import kotlinx.coroutines.launch
|
|
|
|
|
import javax.inject.Inject
|
|
|
|
|
import javax.inject.Singleton
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Real CAM transmit loop for the ESP32-C5 hardware path (Phase 03, Section 13) — the phone-side
|
|
|
|
|
* counterpart to the OBU's old autonomous beacon, now driven from here since the ESP32-C5 has
|
|
|
|
|
* no onboard CAM generation at all (see `obu-firmware/main/main.c`'s rewritten TX path, which
|
|
|
|
|
* is purely receive-and-transmit-on-serial-arrival with no timer of its own).
|
|
|
|
|
*
|
|
|
|
|
* Started/stopped by [TripRecordingService] around an active recording session — per the
|
|
|
|
|
* Section 13 spec, CAM transmission only runs while recording, matching the CiT One path's
|
|
|
|
|
* behavior of "no traffic until there's a trip to correlate it with." Internally also gated on
|
|
|
|
|
* [ObuHardwarePreferences] currently reporting [ObuHardware.ESP32_C5] — on the CiT One path
|
|
|
|
|
* this loop stays parked (via [kotlinx.coroutines.flow.collectLatest] on the hardware
|
|
|
|
|
* preference) and never sends anything.
|
|
|
|
|
*
|
|
|
|
|
* Rate policy: [CamTransmitConfig.baseRateHz] (1 Hz) everywhere, bumped to
|
|
|
|
|
* [CamTransmitConfig.elevatedRateHz] inside a [com.hawhamburg.micr0bu.domain.cam.CamGeofence] or
|
|
|
|
|
* for [ELEVATED_HOLD_MS] after an external event trigger (harsh braking/turning/stopping — see
|
|
|
|
|
* [onDetectedEvent], called by [TripRecordingService] from the same
|
|
|
|
|
* [com.hawhamburg.micr0bu.domain.detection.EventDetector] stream that already drives trip event
|
|
|
|
|
* logging). Both rate figures are placeholders pending real-world tuning, per
|
|
|
|
|
* [CamTransmitConfig]'s own disclaimer.
|
2026-09-23 17:28:05 +02:00
|
|
|
*
|
|
|
|
|
* ## CAM or VAM
|
|
|
|
|
* Settings chooses what goes out ([OutgoingMessage]). CAM follows the rate policy above. VAM is
|
|
|
|
|
* checked every [VAM_TICK_MS] against the generation rules of TS 103 300-3 clause 6.4
|
|
|
|
|
* ([VamGenerationRules]) and sent when one fires, from the same GNSS fix, pseudonym and position
|
|
|
|
|
* vector a CAM would use. Whether either is signed is the "Sign outgoing messages" setting; the
|
|
|
|
|
* micrOBU does the signing ([Esp32Link]).
|
2026-08-04 17:58:07 +02:00
|
|
|
*/
|
|
|
|
|
@Singleton
|
|
|
|
|
class CamTransmitLoop @Inject constructor(
|
|
|
|
|
@ApplicationContext private val context: Context,
|
|
|
|
|
private val obuHardwarePrefs: ObuHardwarePreferences,
|
2026-09-23 17:28:05 +02:00
|
|
|
private val esp32Link: Esp32Link,
|
2026-08-04 17:58:07 +02:00
|
|
|
private val codec: RealAsn1UperCodec,
|
2026-09-10 14:47:30 +02:00
|
|
|
private val pseudonymManager: PseudonymManager,
|
2026-08-04 17:58:07 +02:00
|
|
|
) {
|
|
|
|
|
private val config = CamTransmitConfig()
|
|
|
|
|
private val sensorRepository = SensorRepository(context)
|
|
|
|
|
|
|
|
|
|
private var job: Job? = null
|
|
|
|
|
private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Default)
|
|
|
|
|
|
|
|
|
|
@Volatile private var latestGnss: GnssReading? = null
|
|
|
|
|
@Volatile private var latestGyroZ: Float? = null
|
|
|
|
|
@Volatile private var elevatedUntilMs: Long = 0L
|
|
|
|
|
|
2026-08-10 14:09:18 +02:00
|
|
|
/** Previous GNSS fix, kept only to derive along-track acceleration — see [longitudinalAccel]. */
|
|
|
|
|
@Volatile private var previousGnss: GnssReading? = null
|
|
|
|
|
|
2026-09-23 17:28:05 +02:00
|
|
|
@Volatile private var outgoing: OutgoingMessage = OutgoingMessage.CAM
|
|
|
|
|
@Volatile private var signOutgoing: Boolean = true
|
|
|
|
|
private val vamRules = VamGenerationRules()
|
|
|
|
|
|
2026-08-04 17:58:07 +02:00
|
|
|
/**
|
|
|
|
|
* Call when a braking/turning/stopping event fires during an active trip — bumps the CAM
|
|
|
|
|
* rate to [CamTransmitConfig.elevatedRateHz] for [ELEVATED_HOLD_MS] so nearby stations get
|
|
|
|
|
* denser updates through the maneuver, not just at the instant it was detected.
|
|
|
|
|
*/
|
|
|
|
|
fun onDetectedEvent() {
|
|
|
|
|
elevatedUntilMs = System.currentTimeMillis() + ELEVATED_HOLD_MS
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
/**
|
|
|
|
|
* Starts the loop for the duration of a recording session. Internally stays idle (no
|
|
|
|
|
* transmission) unless/until the ESP32-C5 is the selected OBU hardware, and automatically
|
|
|
|
|
* pauses/resumes if that selection changes mid-trip.
|
|
|
|
|
*/
|
|
|
|
|
fun start() {
|
|
|
|
|
if (job?.isActive == true) return
|
|
|
|
|
elevatedUntilMs = 0L
|
2026-08-10 14:09:18 +02:00
|
|
|
previousGnss = null
|
2026-09-23 17:28:05 +02:00
|
|
|
vamRules.reset()
|
2026-08-04 17:58:07 +02:00
|
|
|
job = scope.launch {
|
|
|
|
|
obuHardwarePrefs.obuHardwareFlow.collectLatest { hardware ->
|
|
|
|
|
if (hardware != ObuHardware.ESP32_C5) return@collectLatest
|
|
|
|
|
runTransmitLoop()
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
fun stop() {
|
|
|
|
|
job?.cancel()
|
|
|
|
|
job = null
|
|
|
|
|
elevatedUntilMs = 0L
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private suspend fun runTransmitLoop() = coroutineScope {
|
|
|
|
|
launch { sensorRepository.gnssFlow().collect { latestGnss = it } }
|
|
|
|
|
launch { sensorRepository.gyroscopeFlow().collect { latestGyroZ = it.z } }
|
2026-09-23 17:28:05 +02:00
|
|
|
launch { obuHardwarePrefs.signOutgoingFlow.collect { signOutgoing = it } }
|
|
|
|
|
launch {
|
|
|
|
|
obuHardwarePrefs.outgoingMessageFlow.collect {
|
|
|
|
|
if (it != outgoing) vamRules.reset()
|
|
|
|
|
outgoing = it
|
|
|
|
|
}
|
|
|
|
|
}
|
2026-08-04 17:58:07 +02:00
|
|
|
|
|
|
|
|
while (true) {
|
|
|
|
|
val gnss = latestGnss
|
|
|
|
|
if (gnss != null) {
|
2026-09-23 17:28:05 +02:00
|
|
|
when (outgoing) {
|
|
|
|
|
OutgoingMessage.CAM -> sendCam(gnss)
|
|
|
|
|
OutgoingMessage.VAM -> sendVamIfDue(gnss)
|
|
|
|
|
}
|
2026-08-04 17:58:07 +02:00
|
|
|
}
|
2026-09-23 17:28:05 +02:00
|
|
|
delay(if (outgoing == OutgoingMessage.VAM) VAM_TICK_MS else (1000.0 / currentRateHz(gnss)).toLong())
|
2026-08-04 17:58:07 +02:00
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-23 17:28:05 +02:00
|
|
|
private suspend fun sendCam(gnss: GnssReading) {
|
|
|
|
|
// Asked for per CAM rather than once per trip: that is what lets a pseudonym
|
|
|
|
|
// rotation fall cleanly between two frames instead of inside one.
|
|
|
|
|
val pseudonym = pseudonymManager.current()
|
|
|
|
|
// Stamped on GNSS time rather than the phone clock; see GnssTimeSource. Only the
|
|
|
|
|
// outgoing CAM is: acceleration below still differences wall-clock samples.
|
|
|
|
|
val fix = gnss.copy(timestamp = GnssTimeSource.correct(gnss.timestamp))
|
|
|
|
|
val cam = PhoneCamBuilder.build(fix, latestGyroZ, pseudonym.stationId, longitudinalAccel(gnss))
|
|
|
|
|
val bytes = codec.encodeCam(cam)
|
|
|
|
|
esp32Link.send(OutgoingIts(OutgoingMessage.CAM, bytes,
|
|
|
|
|
GnPositionVector.fromCam(cam, gnss.accuracyM, pseudonym.mac), gnss.accuracyM, signOutgoing))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
private suspend fun sendVamIfDue(gnss: GnssReading) {
|
|
|
|
|
val now = System.currentTimeMillis()
|
|
|
|
|
val kinematics = VamGenerationRules.Kinematics(gnss.latitude, gnss.longitude,
|
|
|
|
|
gnss.speedMs.toDouble(), gnss.bearingDeg.toDouble())
|
|
|
|
|
if (!vamRules.due(now, kinematics)) return
|
|
|
|
|
val pseudonym = pseudonymManager.current()
|
|
|
|
|
val fix = gnss.copy(timestamp = GnssTimeSource.correct(gnss.timestamp))
|
|
|
|
|
// The CAM view of this fix is built only for its position vector, so the GN header of a VAM
|
|
|
|
|
// follows exactly the rules a CAM's does. It is not transmitted.
|
|
|
|
|
val cam = PhoneCamBuilder.build(fix, latestGyroZ, pseudonym.stationId, longitudinalAccel(gnss))
|
|
|
|
|
val withLowFrequency = vamRules.includeLowFrequency(now)
|
|
|
|
|
val bytes = VamUperCodec.encode(VamContent(
|
|
|
|
|
stationId = pseudonym.stationId,
|
|
|
|
|
timestamp = cam.timestamp,
|
|
|
|
|
latitude = cam.latitude,
|
|
|
|
|
longitude = cam.longitude,
|
|
|
|
|
accuracyM = gnss.accuracyM,
|
|
|
|
|
speedMps = cam.speedMps,
|
|
|
|
|
headingDeg = cam.headingDeg,
|
|
|
|
|
accelerationMps2 = cam.accelerationMps2,
|
|
|
|
|
includeLowFrequency = withLowFrequency,
|
|
|
|
|
))
|
|
|
|
|
val handedOver = esp32Link.send(OutgoingIts(OutgoingMessage.VAM, bytes,
|
|
|
|
|
GnPositionVector.fromCam(cam, gnss.accuracyM, pseudonym.mac), gnss.accuracyM, signOutgoing))
|
|
|
|
|
if (handedOver) vamRules.onSent(now, kinematics, withLowFrequency)
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-10 14:09:18 +02:00
|
|
|
/**
|
|
|
|
|
* Along-track acceleration in m/s², from the change in GNSS speed since the previous fix.
|
|
|
|
|
*
|
|
|
|
|
* Deliberately not from the accelerometer: CAM's `longitudinalAcceleration` is acceleration
|
|
|
|
|
* along the direction of travel, while the raw accelerometer reads in the device frame with
|
|
|
|
|
* gravity included — extracting the along-track component from it needs a full orientation
|
|
|
|
|
* estimate, which this path doesn't have (the detection engine sidesteps the same problem by
|
|
|
|
|
* working on orientation-independent magnitudes, which is not what CAM wants here).
|
|
|
|
|
*
|
|
|
|
|
* Returns null — encoded as ASN.1 `unavailable` — when there's no usable previous fix, when
|
|
|
|
|
* the gap is too short to divide by safely, or when it's long enough that the two samples
|
|
|
|
|
* aren't really consecutive. Better an honest "unavailable" than a fabricated number a
|
|
|
|
|
* receiving vehicle might brake on.
|
|
|
|
|
*/
|
|
|
|
|
private fun longitudinalAccel(gnss: GnssReading): Double? {
|
|
|
|
|
val prev = previousGnss
|
|
|
|
|
previousGnss = gnss
|
|
|
|
|
if (prev == null) return null
|
|
|
|
|
|
|
|
|
|
val dtSec = (gnss.timestamp - prev.timestamp) / 1000.0
|
|
|
|
|
if (dtSec < MIN_ACCEL_DT_SEC || dtSec > MAX_ACCEL_DT_SEC) return null
|
|
|
|
|
|
|
|
|
|
val dv = gnss.speedMs.toDouble() - prev.speedMs.toDouble()
|
|
|
|
|
return dv / dtSec
|
|
|
|
|
}
|
|
|
|
|
|
2026-08-04 17:58:07 +02:00
|
|
|
private fun currentRateHz(gnss: GnssReading?): Double {
|
|
|
|
|
val now = System.currentTimeMillis()
|
|
|
|
|
val inGeofence = gnss != null && config.geofences.any { fence ->
|
|
|
|
|
GeoMath.haversineMeters(gnss.latitude, gnss.longitude, fence.latitude, fence.longitude) <= fence.radiusM
|
|
|
|
|
}
|
|
|
|
|
val eventBoosted = now < elevatedUntilMs
|
|
|
|
|
return if (inGeofence || eventBoosted) config.elevatedRateHz else config.baseRateHz
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
companion object {
|
|
|
|
|
private const val ELEVATED_HOLD_MS = 5_000L
|
2026-08-10 14:09:18 +02:00
|
|
|
|
2026-09-23 17:28:05 +02:00
|
|
|
/** How often VAM generation rules are checked: T_GenVamMin, TS 103 300-3 Table 16. */
|
|
|
|
|
private const val VAM_TICK_MS = VamGenerationRules.T_GEN_VAM_MIN_MS
|
|
|
|
|
|
2026-08-10 14:09:18 +02:00
|
|
|
/** Below this gap, GNSS speed noise divided by a tiny dt produces absurd accelerations. */
|
|
|
|
|
private const val MIN_ACCEL_DT_SEC = 0.2
|
|
|
|
|
|
|
|
|
|
/** Above this gap the two fixes aren't consecutive enough to call the result acceleration. */
|
|
|
|
|
private const val MAX_ACCEL_DT_SEC = 3.0
|
2026-08-04 17:58:07 +02:00
|
|
|
}
|
|
|
|
|
}
|