mirror of
https://github.com/d4rken-org/capod.git
synced 2026-09-14 18:26:11 -04:00
Compare commits
19
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a91ef44686 | ||
|
|
e3bef4e74b | ||
|
|
358bf5e244 | ||
|
|
fba2d3a59d | ||
|
|
0abba4ece2 | ||
|
|
e9b0694b67 | ||
|
|
49bbf1bc85 | ||
|
|
f5b6a5c605 | ||
|
|
fccc6cedd5 | ||
|
|
a03766448e | ||
|
|
f0c9d318ff | ||
|
|
5c5cc6f28e | ||
|
|
246b96bfb4 | ||
|
|
352e020b54 | ||
|
|
72f5690ea6 | ||
|
|
1ed1199574 | ||
|
|
fd4a43de96 | ||
|
|
e406fd4cc1 | ||
|
|
1e555cc712 |
@@ -1,11 +1,11 @@
|
||||
<?xml version="1.0" encoding="utf-8"?>
|
||||
<resources>
|
||||
<string name="foss_upgrade_donate_label">تبرع</string>
|
||||
<string name="foss_upgrade_alreadydonated_label">أنا متبرّع فعلًا</string>
|
||||
<string name="foss_upgrade_no_money_label">أنفق كل أموالي على AirPods</string>
|
||||
<string name="upgrade_foss_preamble">CAPod FOSS مجاني ومفتوح المصدر. إذا وجدته مفيدًا، ففكّر في دعم تطويره للمساعدة في استمرار المشروع.</string>
|
||||
<string name="upgrade_foss_sponsor_action">دعم المشروع</string>
|
||||
<string name="upgrade_foss_sponsor_subtitle">لا إعلانات. لا تتبع. لا قيود من Google Play.</string>
|
||||
<string name="upgrade_foss_sponsor_returned_early">هل عدت فعلًا؟ دعمكم يُبقي CAPod حيًا.</string>
|
||||
<string name="upgrade_badge_label">البرمجيات الحرة (FOSS)</string>
|
||||
<string name="foss_upgrade_donate_label">تبرُّع</string>
|
||||
<string name="foss_upgrade_alreadydonated_label">لقد تبرّعتُ بالفعل</string>
|
||||
<string name="foss_upgrade_no_money_label">أنفق كل أموالي على سمّاعات AirPods</string>
|
||||
<string name="upgrade_foss_preamble">برنامج CAPod FOSS مجاني ومفتوح المصدر. إذا وجدتهُ مفيدًا، ففكّر في دعم تطويره للمساعدة في استمرار المشروع.</string>
|
||||
<string name="upgrade_foss_sponsor_action">دعم التطوير</string>
|
||||
<string name="upgrade_foss_sponsor_subtitle">لا إعلانات. لا تتبُّع. لا قيود من جوجل بلاي.</string>
|
||||
<string name="upgrade_foss_sponsor_returned_early">هل عدت بالفعل؟ دعمك يُبقي CAPod حيًّا.</string>
|
||||
<string name="upgrade_badge_label">البرمجيات الحرة والمفتوحة المصدر</string>
|
||||
</resources>
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
<string name="foss_upgrade_no_money_label">Es iztērēju visu savu naudu AirPods austiņām</string>
|
||||
<string name="upgrade_foss_preamble">CAPod FOSS ir bezmaksas un atkļējkoda. Ja tā ir noderba, apsveriet attīstības sponsēšanu, lai palīdzētu turpināt projektu.</string>
|
||||
<string name="upgrade_foss_sponsor_action">Sponsorēt attīstību</string>
|
||||
<string name="upgrade_foss_sponsor_subtitle">Bez reklamām. Bez seklōšanas. Bez Google Play saistībām.</string>
|
||||
<string name="upgrade_foss_sponsor_returned_early">Jau atpakaļ? Jūsu atbalsts uztur CAPod pie dzīvības.</string>
|
||||
<string name="upgrade_foss_sponsor_subtitle">Bez reklamām. Bez izsekošanas. Bez Google Play saistībām.</string>
|
||||
<string name="upgrade_foss_sponsor_returned_early">Jau atpakaļ? Tavs atbalsts uztur CAPod esamību.</string>
|
||||
<string name="upgrade_badge_label">FOSS</string>
|
||||
</resources>
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
<?xml version="1.0" encoding="utf-8"?>
|
||||
<resources>
|
||||
<string name="upgrades_gplay_unavailable_error">خدمات Google Play غير متوفّرة.</string>
|
||||
<string name="upgrades_no_purchases_found_check_account">لم يتم العثور على أي مشتريات. هل تستخدم الحساب الصحيح؟</string>
|
||||
<string name="upgrades_gplay_billing_error_label">خطأ في Google Play</string>
|
||||
<string name="upgrades_gplay_billing_error_description">حدث خطأ في Google Play. يُرجى المحاولة لاحقًا أو إعادة تشغيل هاتفك.\n\nخطأ: %s</string>
|
||||
<string name="upgrades_gplay_billing_result_error_label">خطأ في فوترة Google Play</string>
|
||||
<string name="upgrades_gplay_billing_result_error_description">حدث خطأ أثناء طلب تفاصيل عملية الشراء من Google Play. يُرجى مسح ذاكرة التخزين المؤقّت لـ Google Play وإعادة تشغيل هاتفك.\n\nخطأ %s</string>
|
||||
<string name="upgrade_restore_action">استعادة المشتريات</string>
|
||||
<string name="upgrades_gplay_unavailable_error">خدمات جوجل بلاي غير متوفّرة.</string>
|
||||
<string name="upgrades_no_purchases_found_check_account">لم يُعثَر على أيّ مشتريات. هل تستخدم الحساب الصحيح؟</string>
|
||||
<string name="upgrades_gplay_billing_error_label">خطأ في جوجل بلاي</string>
|
||||
<string name="upgrades_gplay_billing_error_description">حدث خطأ في سوق جوجل بلاي. يُرجى المحاولة لاحقًا أو إعادة تشغيل هاتفك.\n\nالخطأ: %s</string>
|
||||
<string name="upgrades_gplay_billing_result_error_label">خطأ في فوترة جوجل بلاي</string>
|
||||
<string name="upgrades_gplay_billing_result_error_description">حدث خطأ أثناء طلب تفاصيل عملية الشراء من سوق جوجل بلاي. يُرجى مسح ذاكرة التخزين المؤقّت لسوق جوجل بلاي وإعادة تشغيل هاتفك.\n\nالخطأ %s</string>
|
||||
<string name="upgrade_restore_action">استعادة عملية الشراء</string>
|
||||
<string name="upgrade_badge_label">النسخة الاحترافية</string>
|
||||
</resources>
|
||||
|
||||
@@ -1,11 +1,11 @@
|
||||
<?xml version="1.0" encoding="utf-8"?>
|
||||
<resources>
|
||||
<string name="upgrades_gplay_unavailable_error">Google Play pakalpojumi nav pieejami.</string>
|
||||
<string name="upgrades_no_purchases_found_check_account">Nav atrasts neviens pirkums. Vai izmantojat pareizo kontu?</string>
|
||||
<string name="upgrades_no_purchases_found_check_account">Nav atrasts neviens pirkums. Vai izmanto pareizo kontu?</string>
|
||||
<string name="upgrades_gplay_billing_error_label">Google Play kļūda</string>
|
||||
<string name="upgrades_gplay_billing_error_description">Google Play radās kļūda. Lūdzu, mēģiniet vēlreiz vēlāk vai pārstartējiet tālruni.\n\nKļūda: %s</string>
|
||||
<string name="upgrades_gplay_billing_error_description">Google Play radās kļūda. Lūgums mēģināt vēlreiz vēlāk vai pārsāknēt tālruni.\n\nKļūda: %s</string>
|
||||
<string name="upgrades_gplay_billing_result_error_label">Google Play norēķinu kļūda</string>
|
||||
<string name="upgrades_gplay_billing_result_error_description">Radās kļūda, pieprasot no Google Play jūsu pirkuma informāciju. Notīriet Google Play kešatmiņu un pārstartējiet tālruni.\n\nKļūda %s</string>
|
||||
<string name="upgrades_gplay_billing_result_error_description">Radās kļūda, pieprasot no Google Play pirkuma informāciju. Jāiztīra Google Play kešatmiņa un jāpārsāknē tālrunis.\n\nKļūda %s</string>
|
||||
<string name="upgrade_restore_action">Atjaunot pirkumu</string>
|
||||
<string name="upgrade_badge_label">Pro</string>
|
||||
</resources>
|
||||
|
||||
@@ -26,8 +26,14 @@ class MonitorModeResolver @Inject constructor(
|
||||
profilesRepo.profiles,
|
||||
nudgeCapabilityStore.availability,
|
||||
) { profiles, nudge ->
|
||||
val primary = profiles.firstOrNull() ?: return@combine MonitorMode.MANUAL
|
||||
primary.requiredMode(nudge)
|
||||
// The mode is driven by the first *addressed* profile, not just the first profile. An
|
||||
// address-less profile can neither host an AAP session (AapAutoConnect requires an address)
|
||||
// nor auto-connect, so it can only ever resolve to MANUAL. Letting such a profile sit at the
|
||||
// top of the list and force MANUAL would silently stop the foreground service that keeps the
|
||||
// process — and thus background AAP/stem controls — alive for an addressed profile below it.
|
||||
// No addressed profile at all => genuinely nothing to monitor => MANUAL.
|
||||
profiles.firstOrNull { !it.address.isNullOrBlank() }?.requiredMode(nudge)
|
||||
?: MonitorMode.MANUAL
|
||||
}
|
||||
.distinctUntilChanged()
|
||||
.onEach { log(TAG, VERBOSE) { "effectiveMode = $it" } }
|
||||
|
||||
@@ -9,20 +9,22 @@ import eu.darken.capod.common.flow.setupCommonEventHandlers
|
||||
import eu.darken.capod.monitor.core.ble.BlePodMonitor
|
||||
import eu.darken.capod.pods.core.apple.PodModel
|
||||
import eu.darken.capod.pods.core.apple.aap.AapConnectionManager
|
||||
import eu.darken.capod.pods.core.apple.aap.AapDisconnectEvent
|
||||
import eu.darken.capod.pods.core.apple.aap.AapPodState
|
||||
import eu.darken.capod.profiles.core.AppleDeviceProfile
|
||||
import eu.darken.capod.profiles.core.DeviceProfilesRepo
|
||||
import kotlinx.coroutines.coroutineScope
|
||||
import kotlinx.coroutines.delay
|
||||
import kotlinx.coroutines.flow.Flow
|
||||
import kotlinx.coroutines.flow.channelFlow
|
||||
import kotlinx.coroutines.flow.combine
|
||||
import kotlinx.coroutines.flow.distinctUntilChanged
|
||||
import kotlinx.coroutines.flow.first
|
||||
import kotlinx.coroutines.flow.map
|
||||
import kotlinx.coroutines.flow.mapLatest
|
||||
import kotlinx.coroutines.flow.merge
|
||||
import kotlinx.coroutines.flow.onEach
|
||||
import kotlinx.coroutines.launch
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
import javax.inject.Inject
|
||||
import javax.inject.Singleton
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
@@ -37,6 +39,21 @@ class AapAutoConnect @Inject constructor(
|
||||
private val activeReconnects = java.util.Collections.synchronizedSet(mutableSetOf<String>())
|
||||
private val processedModelCorrections = java.util.Collections.synchronizedSet(mutableSetOf<String>())
|
||||
|
||||
/**
|
||||
* Per-address count of consecutive connection attempts that never reached READY (failed socket
|
||||
* connects and handshake-timeout disconnects). Drives the escalating reconnect cooldown so a
|
||||
* device that keeps stalling its handshake in a crowded RF environment doesn't churn a tight
|
||||
* reconnect loop. Reset to 0 once a session works (a was-ready disconnect) or the device leaves.
|
||||
*/
|
||||
private val preReadyFailures = ConcurrentHashMap<String, Int>()
|
||||
|
||||
/** Cooldown (ms) to wait before the next attempt, given how many consecutive pre-READY failures we've seen. */
|
||||
private fun cooldownFor(address: String): Long {
|
||||
val failures = preReadyFailures[address] ?: 0
|
||||
if (failures <= 0) return 0L
|
||||
return HANDSHAKE_BACKOFF[minOf(failures - 1, HANDSHAKE_BACKOFF.lastIndex)]
|
||||
}
|
||||
|
||||
fun monitor(): Flow<Unit> = merge(
|
||||
initialConnect(),
|
||||
reconnectOnDisconnect(),
|
||||
@@ -74,104 +91,150 @@ class AapAutoConnect @Inject constructor(
|
||||
return
|
||||
}
|
||||
|
||||
// Escalating cooldown carried over from prior failed attempts (socket failures or handshake
|
||||
// stalls). 0 on the first attempt so a healthy device connects without delay.
|
||||
val cooldown = cooldownFor(address)
|
||||
if (cooldown > 0) {
|
||||
log(TAG) { "AAP initial connect cooldown ${cooldown}ms for $address (failures=${preReadyFailures[address]})" }
|
||||
delay(cooldown)
|
||||
}
|
||||
|
||||
log(TAG) { "AAP connecting to $address (${profile.label})" }
|
||||
try {
|
||||
aapManager.connect(address, bonded.internal!!, profile.model)
|
||||
log(TAG) { "AAP connected to $address" }
|
||||
return
|
||||
} catch (e: Exception) {
|
||||
log(TAG, WARN) { "AAP initial connect failed for $address: ${e.message}" }
|
||||
|
||||
for ((attempt, delayMs) in RETRY_DELAYS.withIndex()) {
|
||||
delay(delayMs)
|
||||
|
||||
// Bail out if classic BT disconnected
|
||||
// Bail out if classic BT disconnected — device left, clear the penalty
|
||||
val currentConnected = bluetoothManager.connectedDevices.first().map { it.address }.toSet()
|
||||
if (address !in currentConnected) {
|
||||
log(TAG) { "AAP initial retry: $address no longer classically connected, stopping" }
|
||||
break
|
||||
preReadyFailures.remove(address)
|
||||
return
|
||||
}
|
||||
|
||||
val retryState = aapManager.allStates.value[address]
|
||||
if (retryState != null && retryState.connectionState != AapPodState.ConnectionState.DISCONNECTED) {
|
||||
log(TAG) { "AAP initial retry: $address already reconnected, stopping" }
|
||||
break
|
||||
return
|
||||
}
|
||||
|
||||
try {
|
||||
log(TAG) { "AAP initial retry ${attempt + 1} for $address after ${delayMs}ms" }
|
||||
aapManager.connect(address, bonded.internal!!, profile.model)
|
||||
log(TAG) { "AAP connected to $address on retry ${attempt + 1}" }
|
||||
break
|
||||
return
|
||||
} catch (retryException: Exception) {
|
||||
log(TAG, WARN) { "AAP initial retry ${attempt + 1} failed for $address: ${retryException.message}" }
|
||||
}
|
||||
}
|
||||
|
||||
// Every attempt this episode failed at the socket level — escalate the next round's cooldown.
|
||||
preReadyFailures.merge(address, 1, Int::plus)
|
||||
}
|
||||
}
|
||||
|
||||
private fun reconnectOnDisconnect(): Flow<Unit> = aapManager.disconnectEvents
|
||||
.onEach { address ->
|
||||
if (!activeReconnects.add(address)) {
|
||||
log(TAG, VERBOSE) { "AAP reconnect already in progress for $address, skipping" }
|
||||
return@onEach
|
||||
}
|
||||
|
||||
try {
|
||||
for ((attempt, delayMs) in RETRY_DELAYS.withIndex()) {
|
||||
delay(delayMs)
|
||||
|
||||
// Check if still profiled
|
||||
val profile = profilesRepo.profiles.first()
|
||||
.firstOrNull { it.address == address }
|
||||
if (profile == null) {
|
||||
log(TAG) { "AAP reconnect: $address no longer profiled, stopping" }
|
||||
break
|
||||
}
|
||||
|
||||
// Check if still bonded
|
||||
val bonded = bluetoothManager.bondedDevices().first()
|
||||
.firstOrNull { it.address == address }
|
||||
if (bonded == null) {
|
||||
log(TAG) { "AAP reconnect: $address no longer bonded, stopping" }
|
||||
break
|
||||
}
|
||||
|
||||
// Check if still classically connected
|
||||
val currentConnected = bluetoothManager.connectedDevices.first().map { it.address }.toSet()
|
||||
if (address !in currentConnected) {
|
||||
log(TAG) { "AAP reconnect: $address no longer classically connected, stopping" }
|
||||
break
|
||||
}
|
||||
|
||||
// Check if still visible in BLE
|
||||
val bleDevices = blePodMonitor.devices.first()
|
||||
if (bleDevices.none { it.meta?.profile?.address == address }) {
|
||||
log(TAG) { "AAP reconnect: $address no longer visible in BLE, stopping" }
|
||||
break
|
||||
}
|
||||
|
||||
// Check if already reconnected (e.g., by initialConnect)
|
||||
val currentState = aapManager.allStates.value[address]
|
||||
if (currentState != null && currentState.connectionState != AapPodState.ConnectionState.DISCONNECTED) {
|
||||
log(TAG) { "AAP reconnect: $address already reconnected" }
|
||||
break
|
||||
}
|
||||
|
||||
try {
|
||||
log(TAG) { "AAP reconnect attempt ${attempt + 1} for $address in ${delayMs}ms" }
|
||||
aapManager.connect(address, bonded.internal!!, profile.model)
|
||||
log(TAG) { "AAP reconnected to $address" }
|
||||
break
|
||||
} catch (e: Exception) {
|
||||
log(TAG, WARN) { "AAP reconnect attempt ${attempt + 1} failed for $address: ${e.message}" }
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
activeReconnects.remove(address)
|
||||
}
|
||||
// Process each disconnect on its own child coroutine so a long per-address backoff never blocks
|
||||
// another device's reconnect (the collector would otherwise serialise events). channelFlow gives
|
||||
// us a scope to launch into; it intentionally emits nothing — it stays subscribed for its lifetime.
|
||||
private fun reconnectOnDisconnect(): Flow<Unit> = channelFlow<Unit> {
|
||||
aapManager.disconnectEvents.collect { event ->
|
||||
launch { handleReconnect(event) }
|
||||
}
|
||||
.map { } // SharedFlow<BluetoothAddress> → Flow<Unit>
|
||||
.setupCommonEventHandlers(TAG) { "reconnect" }
|
||||
}.setupCommonEventHandlers(TAG) { "reconnect" }
|
||||
|
||||
private suspend fun handleReconnect(event: AapDisconnectEvent) {
|
||||
val address = event.address
|
||||
|
||||
// A session that actually worked resets the penalty (prompt reconnect). One that never reached
|
||||
// READY is a failed handshake/short session — escalate so a persistent stall backs off.
|
||||
if (event.wasEverReady) {
|
||||
preReadyFailures.remove(address)
|
||||
} else {
|
||||
preReadyFailures.merge(address, 1, Int::plus)
|
||||
}
|
||||
|
||||
if (!activeReconnects.add(address)) {
|
||||
log(TAG, VERBOSE) { "AAP reconnect already in progress for $address, skipping" }
|
||||
return
|
||||
}
|
||||
|
||||
try {
|
||||
// Escalating cooldown before the normal retry cadence when handshakes keep stalling.
|
||||
val cooldown = cooldownFor(address)
|
||||
if (cooldown > 0) {
|
||||
log(TAG) { "AAP reconnect cooldown ${cooldown}ms for $address (failures=${preReadyFailures[address]})" }
|
||||
delay(cooldown)
|
||||
}
|
||||
|
||||
for ((attempt, delayMs) in RETRY_DELAYS.withIndex()) {
|
||||
delay(delayMs)
|
||||
|
||||
// Check if still profiled
|
||||
val profile = profilesRepo.profiles.first()
|
||||
.firstOrNull { it.address == address }
|
||||
if (profile == null) {
|
||||
log(TAG) { "AAP reconnect: $address no longer profiled, stopping" }
|
||||
preReadyFailures.remove(address)
|
||||
return
|
||||
}
|
||||
|
||||
// Check if still bonded
|
||||
val bonded = bluetoothManager.bondedDevices().first()
|
||||
.firstOrNull { it.address == address }
|
||||
if (bonded == null) {
|
||||
log(TAG) { "AAP reconnect: $address no longer bonded, stopping" }
|
||||
preReadyFailures.remove(address)
|
||||
return
|
||||
}
|
||||
|
||||
// Check if still classically connected
|
||||
val currentConnected = bluetoothManager.connectedDevices.first().map { it.address }.toSet()
|
||||
if (address !in currentConnected) {
|
||||
log(TAG) { "AAP reconnect: $address no longer classically connected, stopping" }
|
||||
preReadyFailures.remove(address)
|
||||
return
|
||||
}
|
||||
|
||||
// Check if still visible in BLE
|
||||
val bleDevices = blePodMonitor.devices.first()
|
||||
if (bleDevices.none { it.meta?.profile?.address == address }) {
|
||||
log(TAG) { "AAP reconnect: $address no longer visible in BLE, stopping" }
|
||||
preReadyFailures.remove(address)
|
||||
return
|
||||
}
|
||||
|
||||
// Check if already reconnected (e.g., by initialConnect)
|
||||
val currentState = aapManager.allStates.value[address]
|
||||
if (currentState != null && currentState.connectionState != AapPodState.ConnectionState.DISCONNECTED) {
|
||||
log(TAG) { "AAP reconnect: $address already reconnected" }
|
||||
return
|
||||
}
|
||||
|
||||
try {
|
||||
log(TAG) { "AAP reconnect attempt ${attempt + 1} for $address in ${delayMs}ms" }
|
||||
aapManager.connect(address, bonded.internal!!, profile.model)
|
||||
log(TAG) { "AAP reconnected to $address" }
|
||||
return
|
||||
} catch (e: Exception) {
|
||||
log(TAG, WARN) { "AAP reconnect attempt ${attempt + 1} failed for $address: ${e.message}" }
|
||||
}
|
||||
}
|
||||
|
||||
// Every attempt threw at the socket level (the device is still present but won't accept a
|
||||
// socket) — escalate so we don't hammer it. Mirrors connectWithRetries. Handshake stalls
|
||||
// don't reach here: their reconnect socket succeeds and we return above; they escalate via
|
||||
// the never-ready disconnect event instead.
|
||||
preReadyFailures.merge(address, 1, Int::plus)
|
||||
} finally {
|
||||
activeReconnects.remove(address)
|
||||
}
|
||||
}
|
||||
|
||||
private fun correctModelOnDeviceInfo(): Flow<Unit> = aapManager.allStates
|
||||
.map { states -> states.mapValues { (_, state) -> state.deviceInfo?.modelNumber } }
|
||||
@@ -224,5 +287,12 @@ class AapAutoConnect @Inject constructor(
|
||||
companion object {
|
||||
private val TAG = logTag("Monitor", "AapAutoConnect")
|
||||
internal val RETRY_DELAYS = longArrayOf(3_000, 3_000, 3_000, 5_000, 5_000, 10_000, 10_000)
|
||||
|
||||
/**
|
||||
* Escalating cooldown (ms) applied before a reconnect/connect attempt, indexed by
|
||||
* consecutive pre-READY failures minus one. Caps at 60s so a device that perpetually stalls
|
||||
* its handshake in crowded RF settles into a slow ~minute cadence instead of a tight loop.
|
||||
*/
|
||||
internal val HANDSHAKE_BACKOFF = longArrayOf(5_000, 15_000, 30_000, 60_000)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -54,8 +54,8 @@ class AapConnectionManager @Inject constructor(
|
||||
private val _allStates = MutableStateFlow<Map<BluetoothAddress, AapPodState>>(emptyMap())
|
||||
val allStates: StateFlow<Map<BluetoothAddress, AapPodState>> = _allStates.asStateFlow()
|
||||
|
||||
private val _disconnectEvents = MutableSharedFlow<BluetoothAddress>(extraBufferCapacity = 16)
|
||||
val disconnectEvents: SharedFlow<BluetoothAddress> = _disconnectEvents.asSharedFlow()
|
||||
private val _disconnectEvents = MutableSharedFlow<AapDisconnectEvent>(extraBufferCapacity = 16)
|
||||
val disconnectEvents: SharedFlow<AapDisconnectEvent> = _disconnectEvents.asSharedFlow()
|
||||
|
||||
/** Emits when a connection receives private keys (IRK/ENC) from the device. */
|
||||
private val _keysReceived = MutableSharedFlow<Pair<BluetoothAddress, KeyExchangeResult>>(extraBufferCapacity = 16)
|
||||
@@ -75,10 +75,10 @@ class AapConnectionManager @Inject constructor(
|
||||
val sleepEvents: SharedFlow<BluetoothAddress> = _sleepEvents.asSharedFlow()
|
||||
|
||||
/**
|
||||
* Emits when a connected device reports a Conversational Awareness speaking transition
|
||||
* (START/STOP, AAP command 0x4B). Paired with the origin address so the conversation
|
||||
* reaction can gate on the primary device. Only the known 0x01/0x04 markers reach here —
|
||||
* the engine drops unknown raw values (see [AapSessionEngine]).
|
||||
* Emits when a connected device reports a Conversational Awareness transition (START / RESUME /
|
||||
* HOLD / STOP, AAP command 0x4B). Paired with the origin address so the conversation reaction
|
||||
* can gate on the primary device. Every well-formed frame is classified and emitted (unknown
|
||||
* status bytes map to HOLD); only structurally malformed frames are dropped (see [AapSessionEngine]).
|
||||
*/
|
||||
private val _conversationalAwarenessEvents =
|
||||
MutableSharedFlow<Pair<BluetoothAddress, ConversationAwarenessEvent>>(extraBufferCapacity = 16)
|
||||
@@ -188,7 +188,7 @@ class AapConnectionManager @Inject constructor(
|
||||
}
|
||||
|
||||
if (!wasIntentional) {
|
||||
_disconnectEvents.tryEmit(address)
|
||||
_disconnectEvents.tryEmit(AapDisconnectEvent(address, connection.wasEverReady))
|
||||
}
|
||||
|
||||
// End this collector coroutine (also cancels child key-forwarding coroutine)
|
||||
@@ -217,3 +217,13 @@ class AapConnectionManager @Inject constructor(
|
||||
connection.send(command)
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Emitted on an unintentional AAP disconnect. [wasEverReady] tells the reconnect layer whether this
|
||||
* session ever completed the handshake (reached READY) — a never-ready disconnect is a failed
|
||||
* handshake/short session and feeds the escalating reconnect backoff; a was-ready drop reconnects promptly.
|
||||
*/
|
||||
data class AapDisconnectEvent(
|
||||
val address: BluetoothAddress,
|
||||
val wasEverReady: Boolean,
|
||||
)
|
||||
|
||||
@@ -23,12 +23,14 @@ import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.Job
|
||||
import kotlinx.coroutines.flow.SharedFlow
|
||||
import kotlinx.coroutines.flow.StateFlow
|
||||
import kotlinx.coroutines.flow.first
|
||||
import kotlinx.coroutines.isActive
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.coroutines.withContext
|
||||
import kotlinx.coroutines.withTimeout
|
||||
import kotlinx.coroutines.withTimeoutOrNull
|
||||
import java.io.IOException
|
||||
import java.util.concurrent.atomic.AtomicBoolean
|
||||
import kotlin.time.Duration
|
||||
@@ -46,10 +48,14 @@ internal class AapConnection(
|
||||
private val socketFactory: L2capSocketFactory,
|
||||
timeSource: TimeSource,
|
||||
private val connectTimeout: Duration = DEFAULT_CONNECT_TIMEOUT,
|
||||
private val handshakeTimeout: Duration = DEFAULT_HANDSHAKE_TIMEOUT,
|
||||
) {
|
||||
|
||||
private val engine = AapSessionEngine(profile, timeSource)
|
||||
|
||||
/** True once this connection's session reached READY at least once. Survives the engine reset. */
|
||||
val wasEverReady: Boolean get() = engine.wasEverReady
|
||||
|
||||
val state: StateFlow<AapPodState> get() = engine.state
|
||||
val keysReceived: SharedFlow<KeyExchangeResult> get() = engine.keysReceived
|
||||
val stemPressEvents: SharedFlow<StemPressEvent> get() = engine.stemPressEvents
|
||||
@@ -60,6 +66,7 @@ internal class AapConnection(
|
||||
|
||||
private var socket: BluetoothSocket? = null
|
||||
private var readerJob: Job? = null
|
||||
private val disconnected = AtomicBoolean(false)
|
||||
private val writeMutex = Mutex()
|
||||
private val framer = AapFramer()
|
||||
|
||||
@@ -109,6 +116,27 @@ internal class AapConnection(
|
||||
|
||||
// Launch read loop in the provided scope — connect() returns immediately
|
||||
readerJob = scope.launch(Dispatchers.IO) { readLoop(sock) }
|
||||
|
||||
// Handshake watchdog: the socket can connect and the handshake be sent, but if the peer
|
||||
// never replies the engine sits in HANDSHAKING forever (blocking read, no deadline).
|
||||
// Bound it: if neither READY nor DISCONNECTED is reached in time, tear the socket down so
|
||||
// the reconnect path can recover. READY / an earlier DISCONNECTED end the wait with no action.
|
||||
scope.launch {
|
||||
val settled = withTimeoutOrNull(handshakeTimeout) {
|
||||
state.first {
|
||||
it.connectionState == AapPodState.ConnectionState.READY ||
|
||||
it.connectionState == AapPodState.ConnectionState.DISCONNECTED
|
||||
}
|
||||
}
|
||||
// Re-check after the timeout: READY may have landed in the boundary race between
|
||||
// withTimeoutOrNull cancelling and us getting here — don't tear down a live session.
|
||||
if (settled == null && state.value.connectionState != AapPodState.ConnectionState.READY) {
|
||||
log(TAG, Logging.Priority.WARN) {
|
||||
"Handshake timed out after $handshakeTimeout for ${device.address} — disconnecting"
|
||||
}
|
||||
disconnect()
|
||||
}
|
||||
}
|
||||
} catch (e: Exception) {
|
||||
log(TAG, Logging.Priority.ERROR) { "Connection failed: $e" }
|
||||
cleanupSocket()
|
||||
@@ -118,6 +146,9 @@ internal class AapConnection(
|
||||
}
|
||||
|
||||
suspend fun disconnect() = withContext(Dispatchers.IO) {
|
||||
// Idempotent: the handshake watchdog and the manager's DISCONNECTED observer can both reach
|
||||
// here for the same dying session. Run the teardown exactly once.
|
||||
if (!disconnected.compareAndSet(false, true)) return@withContext
|
||||
log(TAG, Logging.Priority.INFO) { "Disconnecting" }
|
||||
readerJob?.cancel()
|
||||
readerJob = null
|
||||
@@ -261,6 +292,13 @@ internal class AapConnection(
|
||||
companion object {
|
||||
private const val PSM = 0x1001
|
||||
internal val DEFAULT_CONNECT_TIMEOUT = 5.seconds
|
||||
|
||||
/**
|
||||
* Upper bound on the post-handshake wait for the first sign of life (READY). Normal
|
||||
* handshakes complete in well under 2s; 10s tolerates a congested 2.4 GHz band before
|
||||
* giving up so the reconnect path can recover instead of wedging in HANDSHAKING forever.
|
||||
*/
|
||||
internal val DEFAULT_HANDSHAKE_TIMEOUT = 10.seconds
|
||||
private val TAG = logTag("AAP", "Connection")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -86,6 +86,15 @@ internal class AapSessionEngine(
|
||||
private var activeSendRaw: (suspend (AapCommand) -> Unit)? = null
|
||||
private var runtimeState = EngineRuntimeState()
|
||||
|
||||
/**
|
||||
* True once this session has reached [AapPodState.ConnectionState.READY] at least once.
|
||||
* Deliberately NOT cleared by [reset] — consumers read it at disconnect time to tell a
|
||||
* working session that dropped apart from one that never completed the handshake (drives
|
||||
* the reconnect backoff). Cleared only when a fresh session starts.
|
||||
*/
|
||||
private var _wasEverReady: Boolean = false
|
||||
val wasEverReady: Boolean get() = _wasEverReady
|
||||
|
||||
fun start(scope: CoroutineScope) {
|
||||
dispatch(AapEngineEvent.SessionStarted(scope))
|
||||
}
|
||||
@@ -130,6 +139,7 @@ internal class AapSessionEngine(
|
||||
when (event) {
|
||||
is AapEngineEvent.SessionStarted -> {
|
||||
scope = event.scope
|
||||
_wasEverReady = false
|
||||
runtimeState = runtimeState.copy(handshakeResponseReceived = false)
|
||||
_state.value = _state.value.copy(connectionState = AapPodState.ConnectionState.CONNECTING)
|
||||
}
|
||||
@@ -159,11 +169,14 @@ internal class AapSessionEngine(
|
||||
private fun handleConnectResponse(packet: AapPacket.ConnectResponse) {
|
||||
if (packet.status != 0) {
|
||||
log(TAG, ERROR) {
|
||||
"ConnectResponse failed: status=0x${"%04X".format(packet.status)} " +
|
||||
"ConnectResponse rejected (status=0x${"%04X".format(packet.status)}) — tearing down: " +
|
||||
"major=${packet.major} minor=${packet.minor} " +
|
||||
"features=0x${"%016X".format(packet.features.toLong())}"
|
||||
}
|
||||
_state.value = _state.value.copy(connectResponseStatus = packet.status)
|
||||
// Hard protocol rejection: don't wait out the handshake watchdog — disconnect now so
|
||||
// the reconnect path can recover (or back off). Nested dispatch mirrors the existing
|
||||
// inbound-update re-dispatch pattern.
|
||||
dispatch(AapEngineEvent.ResetRequested)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -192,6 +205,7 @@ internal class AapSessionEngine(
|
||||
_state.value.connectionState == AapPodState.ConnectionState.HANDSHAKING
|
||||
) {
|
||||
runtimeState = runtimeState.copy(handshakeResponseReceived = true)
|
||||
_wasEverReady = true
|
||||
_state.value = _state.value.copy(connectionState = AapPodState.ConnectionState.READY)
|
||||
log(TAG) { "Connection READY" }
|
||||
}
|
||||
@@ -309,8 +323,8 @@ internal class AapSessionEngine(
|
||||
}
|
||||
|
||||
// Re-emit every (well-formed) Conversational Awareness frame as a classified event for the
|
||||
// conversation reaction: START / STOP / HOLD (keep-alive). The decoder already dropped
|
||||
// malformed frames (rawValue stays null only in that case). The raw payload remains logged.
|
||||
// conversation reaction: START / RESUME / HOLD / STOP. The decoder already dropped malformed
|
||||
// frames (rawValue stays null only in that case). The raw payload remains logged.
|
||||
if (value is AapSetting.ConversationalAwarenessState) {
|
||||
value.rawValue?.let { _conversationalAwarenessEvents.tryEmit(ConversationAwarenessEvent.fromStatus(it)) }
|
||||
}
|
||||
|
||||
@@ -103,9 +103,10 @@ sealed class AapSetting {
|
||||
* Push-only from device — reports speaking detection state (command 0x4B).
|
||||
*
|
||||
* [rawValue] is the status byte: the last byte of the 4-byte `02 00 01 XX` form (or the single
|
||||
* byte of the legacy form), preserved so consumers can classify it. [speaking] is `true` only
|
||||
* for the speaking-onset statuses (`1`, `2`); every other value (`0`, `3`, `4`, `5`, `0x0B`, …)
|
||||
* is `false`. START/STOP/HOLD classification for the reaction lives in [ConversationAwarenessEvent].
|
||||
* byte of the legacy form), preserved so consumers can classify it. [speaking] is `true` while
|
||||
* the wearer is talking — onset (`1`, `2`) or resumed after a pause (`5`); every other value
|
||||
* (`0`, `3`, `4`, `0x0B`, …) is `false`. START/RESUME/HOLD/STOP classification for the reaction
|
||||
* lives in [ConversationAwarenessEvent].
|
||||
*/
|
||||
data class ConversationalAwarenessState(
|
||||
val speaking: Boolean,
|
||||
|
||||
+25
-15
@@ -3,36 +3,46 @@ package eu.darken.capod.pods.core.apple.aap.protocol
|
||||
/**
|
||||
* Classified Conversational Awareness signal derived from the status byte of a `0x4B` frame.
|
||||
*
|
||||
* Status-byte mapping (from live AirPods Pro 3 + Pro 2 USB-C captures — both models share one
|
||||
* firmware train and a byte-identical protocol — plus the librepods project):
|
||||
* - `1`, `2` → [START] (wearer started / is speaking → engage the reaction)
|
||||
* - `5`, `6`, `8`, `9` → [STOP] (wearer stopped → disengage). All four confirmed live on Pro 2
|
||||
* USB-C fw `…6814`; which one terminates a given flurry varies with how speech ended.
|
||||
* - any other value (`0`, `3`, `4`, `7`, `0x0B`, … and anything unrecognised) → [HOLD]: a
|
||||
* transitional wind-down frame (`7` was only discovered on fw `…6814` — the set is open-ended,
|
||||
* so unknown values are deliberately classified as HOLD rather than guessed at).
|
||||
* Status-byte mapping, derived from a labelled capture set across AirPods Pro 3 and Pro 2 USB-C
|
||||
* (`protocol-research/conversationalawareness/` — both models share one firmware train and emit
|
||||
* byte-identical sequences). Each scenario was captured with a known action (normal talking,
|
||||
* bursty talking, single-pod, volume-up abort, pod removal/case):
|
||||
* - `1`, `2` → [START] (conversation onset / wind-up → engage the reaction). A conversation always
|
||||
* opens `1` then `2`; re-onset only ever happens after a terminal.
|
||||
* - `5` → [RESUME] (speech resumed after a pause; the wind-down was aborted, CA stays engaged). In
|
||||
* bursty speech the pod cycles `3,5,3,5,…`; `5` is NEVER a terminal. Misreading `5` as a stop was
|
||||
* the root cause of the premature-resume bug — see [ConversationReaction].
|
||||
* - `6`, `8`, `9` → [STOP] (conversation ended → disengage). The usual terminal is the `8`→`9`
|
||||
* pair; `6` is a standalone terminal seen live on Pro 2 USB-C fw `…6814`. Pod removal / case-close
|
||||
* also emit `8`,`9` (sometimes with no prior `1`,`2`).
|
||||
* - any other value (`3` pause, `4` / `0x0B` wind-down, `7` abort, and anything unrecognised)
|
||||
* → [HOLD]: a transitional "possible/real wind-down" frame. The real wind-down runs `3→0x0B→4`
|
||||
* then the `8,9` terminal; `7` precedes an aborted terminal. Unknown values are deliberately
|
||||
* classified as HOLD (arm the safety fuse, never resume immediately) rather than guessed at.
|
||||
*
|
||||
* The pod sends NO frames during active speech — it stays engaged (and silent) for as long as it
|
||||
* hears nearby voices, 20-30s+ observed. So frame-silence must NOT be read as "speaking ended".
|
||||
* Conversely, any non-START frame means the wind-down has begun: a short flurry of transitional
|
||||
* and terminal frames (e.g. `3,0xB,4,8,9` or `3,5,7,8,9`). With only ONE pod worn (other in
|
||||
* case/disconnected) the terminal is deterministically dropped — the flurry ends on a transitional
|
||||
* `4` (#608; reproduced on Pro 3 and Pro 2 alike) — so [ConversationReaction] treats a HOLD as
|
||||
* "terminal imminent" and arms a short fuse, with a long stale backstop for a fully-dropped flurry.
|
||||
* A flurry's trailing frames may arrive after its terminal; they are ignored once disengaged.
|
||||
* A HOLD means the wind-down may have begun; with only ONE pod worn the `8,9` terminal is
|
||||
* deterministically dropped — the flurry ends on `4` (#608, reproduced on Pro 3 and Pro 2) — so
|
||||
* [ConversationReaction] treats a HOLD as "terminal imminent" and arms a short fuse, with a long
|
||||
* stale backstop for a fully-dropped flurry. A RESUME (`5`) cancels that fuse and re-arms the long
|
||||
* backstop; trailing frames after a terminal are ignored once disengaged.
|
||||
*/
|
||||
enum class ConversationAwarenessEvent {
|
||||
START,
|
||||
RESUME,
|
||||
HOLD,
|
||||
STOP,
|
||||
;
|
||||
|
||||
companion object {
|
||||
val SPEAKING_STATUSES = setOf(1, 2)
|
||||
val STOPPED_STATUSES = setOf(5, 6, 8, 9)
|
||||
val RESUME_STATUSES = setOf(5)
|
||||
val STOPPED_STATUSES = setOf(6, 8, 9)
|
||||
|
||||
fun fromStatus(status: Int): ConversationAwarenessEvent = when (status) {
|
||||
in SPEAKING_STATUSES -> START
|
||||
in RESUME_STATUSES -> RESUME
|
||||
in STOPPED_STATUSES -> STOP
|
||||
else -> HOLD
|
||||
}
|
||||
|
||||
+5
-2
@@ -197,8 +197,11 @@ class DefaultAapDeviceProfile(
|
||||
p[3].toInt() and 0xFF
|
||||
else -> return null
|
||||
}
|
||||
// speaking = the "started/active speaking" statuses (1,2); see ConversationAwarenessEvent.
|
||||
val speaking = status in ConversationAwarenessEvent.SPEAKING_STATUSES
|
||||
// speaking = the wearer is currently speaking: onset (1,2) or resumed after a pause (5).
|
||||
// See ConversationAwarenessEvent. (Consumers read rawValue for full classification; this
|
||||
// bool just exposes "is talking now", so a resume must count as speaking too.)
|
||||
val speaking = status in ConversationAwarenessEvent.SPEAKING_STATUSES ||
|
||||
status in ConversationAwarenessEvent.RESUME_STATUSES
|
||||
return AapSetting.ConversationalAwarenessState::class to
|
||||
AapSetting.ConversationalAwarenessState(speaking, rawValue = status)
|
||||
}
|
||||
|
||||
@@ -16,6 +16,12 @@ data class KnownDevice(
|
||||
val seenCounter: Int,
|
||||
val history: List<ApplePods>,
|
||||
val lastCaseBattery: Float?,
|
||||
/**
|
||||
* Profile id this device has ever IRK-resolved to, or null if it has only ever been seen as a
|
||||
* keyless device. Durable — kept even when [history] is trimmed — so the weak history-recovery
|
||||
* paths can reliably refuse to hand an IRK-backed identity to a foreign same-model frame.
|
||||
*/
|
||||
val boundProfileId: String? = null,
|
||||
) {
|
||||
val lastPayload: ProximityPayload
|
||||
get() = history.last().payload
|
||||
|
||||
@@ -76,7 +76,8 @@ class PodHistoryRepo @Inject constructor(
|
||||
val caseIgnored = if (current is DualApplePods) {
|
||||
val target = current.getCaseMatchMarkings()
|
||||
knownDevices.values.filter { known ->
|
||||
known.history.filterIsInstance<DualApplePods>().any { it.getCaseMatchMarkings() == target }
|
||||
known.history.filterIsInstance<DualApplePods>().any { it.getCaseMatchMarkings() == target } &&
|
||||
known.isIrkConsistentWith(current)
|
||||
}
|
||||
} else {
|
||||
emptySet()
|
||||
@@ -143,7 +144,7 @@ class PodHistoryRepo @Inject constructor(
|
||||
if (recognizedDevice == null) {
|
||||
val currentMarkers = payload.getFuzzyIdentifier()
|
||||
recognizedDevice = knownDevices.values
|
||||
.firstOrNull { it.lastPayload.getFuzzyIdentifier() == currentMarkers }
|
||||
.firstOrNull { it.lastPayload.getFuzzyIdentifier() == currentMarkers && it.isIrkConsistentWith(current) }
|
||||
?.also { log(TAG, DEBUG) { "search1: Similarity match for device=${current.model}" } }
|
||||
}
|
||||
|
||||
@@ -154,6 +155,24 @@ class PodHistoryRepo @Inject constructor(
|
||||
return recognizedDevice
|
||||
}
|
||||
|
||||
/**
|
||||
* The weak history paths (fuzzy markers, case-ignored markings) match on low-entropy broadcast
|
||||
* data — model + battery nibbles + colour — which a stranger's same-model AirPods at a similar
|
||||
* battery satisfies. Those paths must not let a foreign frame inherit the identity of one of OUR
|
||||
* keyed (IRK-backed) devices, or the foreign snapshot poisons that device's cache slot and the
|
||||
* dashboard card glitches (Symptom B).
|
||||
*
|
||||
* A [KnownDevice] is "IRK-backed" if it has ever resolved against a profile's IRK
|
||||
* ([KnownDevice.boundProfileId], latched durably so history trimming can't erode it). Such a
|
||||
* device may only be claimed via the strong address/IRK paths, or by a current frame that
|
||||
* resolved to the SAME profile. Keyless devices (never IRK-backed) are unaffected — fuzzy
|
||||
* recovery still works for them exactly as before.
|
||||
*/
|
||||
private fun KnownDevice.isIrkConsistentWith(current: ApplePods): Boolean {
|
||||
val bound = boundProfileId ?: return true // not IRK-backed -> keyless, weak match allowed
|
||||
return current.meta.profile?.id == bound
|
||||
}
|
||||
|
||||
fun updateHistory(device: ApplePods) {
|
||||
val existing = knownDevices[device.identifier]
|
||||
|
||||
@@ -162,12 +181,16 @@ class PodHistoryRepo @Inject constructor(
|
||||
.mapNotNull { it.batteryCasePercent }
|
||||
.lastOrNull()
|
||||
|
||||
// Latch the IRK-resolved profile id (never cleared by an unbound frame), so it survives history trimming.
|
||||
val boundProfileId = device.meta.profile?.id?.takeIf { device.meta.isIRKMatch }
|
||||
|
||||
knownDevices[device.identifier] = if (existing != null) {
|
||||
val history = existing.history.plus(device)
|
||||
existing.copy(
|
||||
seenCounter = existing.seenCounter + 1,
|
||||
history = history,
|
||||
lastCaseBattery = history.determineLatestCaseBattery() ?: existing.lastCaseBattery
|
||||
lastCaseBattery = history.determineLatestCaseBattery() ?: existing.lastCaseBattery,
|
||||
boundProfileId = boundProfileId ?: existing.boundProfileId,
|
||||
)
|
||||
} else {
|
||||
log(TAG, DEBUG) { "Creating new history for ${device.logSummary()}" }
|
||||
@@ -177,7 +200,8 @@ class PodHistoryRepo @Inject constructor(
|
||||
seenFirstAt = device.seenFirstAt,
|
||||
seenCounter = 1,
|
||||
history = history,
|
||||
lastCaseBattery = history.determineLatestCaseBattery()
|
||||
lastCaseBattery = history.determineLatestCaseBattery(),
|
||||
boundProfileId = boundProfileId,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
+232
-51
@@ -12,6 +12,8 @@ import eu.darken.capod.monitor.core.DeviceMonitor
|
||||
import eu.darken.capod.monitor.core.PodDevice
|
||||
import eu.darken.capod.monitor.core.primaryDevice
|
||||
import eu.darken.capod.pods.core.apple.aap.AapConnectionManager
|
||||
import eu.darken.capod.pods.core.apple.aap.AapPodState
|
||||
import eu.darken.capod.pods.core.apple.aap.protocol.AapSetting
|
||||
import eu.darken.capod.pods.core.apple.aap.protocol.ConversationAwarenessEvent
|
||||
import eu.darken.capod.profiles.core.ReactionConfig
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
@@ -30,7 +32,6 @@ import kotlinx.coroutines.sync.withLock
|
||||
import kotlinx.coroutines.withContext
|
||||
import javax.inject.Inject
|
||||
import javax.inject.Singleton
|
||||
import kotlin.time.Duration
|
||||
import kotlin.time.Duration.Companion.milliseconds
|
||||
import kotlin.time.Duration.Companion.minutes
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
@@ -40,17 +41,29 @@ import kotlin.time.Duration.Companion.seconds
|
||||
* volume or pausing, per the primary device's [ReactionConfig.conversationAction], and reverts when
|
||||
* speaking stops. On Android the pod firmware does not duck audio itself, so CAPod performs it.
|
||||
*
|
||||
* Disengage is driven by the pod's explicit end-of-speech frame ([ConversationAwarenessEvent.STOP]).
|
||||
* The pod sends NO frames during active speech — it stays engaged for as long as it hears nearby
|
||||
* voices (observed: CA held engaged 21-32s with zero `0x4B` frames, indefinitely against ambient
|
||||
* noise). So frame-silence must NOT be read as "speaking ended"; after a START only the long
|
||||
* [STALE_TIMEOUT] backstop applies, and a link drop is handled by the owner-disconnect revert.
|
||||
* Disengage is driven by the pod's explicit end-of-speech frame ([ConversationAwarenessEvent.STOP],
|
||||
* status `8,9`). The pod sends NO frames during active speech — it stays engaged for as long as it
|
||||
* hears nearby voices (observed: CA held engaged 21-32s with zero `0x4B` frames, indefinitely
|
||||
* against ambient noise). So frame-silence must NOT be read as "speaking ended"; after a START only
|
||||
* the long [STALE_TIMEOUT] backstop applies, and a link drop is handled by the owner-disconnect revert.
|
||||
*
|
||||
* Because the pod is silent while engaged, any non-START frame is evidence the wind-down has begun.
|
||||
* The wind-down is a short flurry of transitional ([ConversationAwarenessEvent.HOLD]) and terminal
|
||||
* frames — but the terminal is sometimes dropped entirely (fw `…6589` emitted `3,0xB,4` then nothing,
|
||||
* stranding the volume low — #608). So a HOLD frame re-arms the timer with the short
|
||||
* [WIND_DOWN_TIMEOUT] fuse: if no terminal follows, speech is treated as ended anyway.
|
||||
* A [ConversationAwarenessEvent.HOLD] frame (`3` pause, `0x0B`/`4` wind-down, `7` abort) is evidence
|
||||
* the wind-down may have begun. The real wind-down is a short flurry `3→0x0B→4` then the `8,9`
|
||||
* terminal — but the terminal is sometimes dropped entirely (single-pod wear: fw `…6589`/`…6861`
|
||||
* emitted `3,0xB,4` then nothing, stranding the volume low — #608). So a HOLD frame re-arms the
|
||||
* timer with the short [WIND_DOWN_TIMEOUT] fuse: if no terminal (or resume) follows, speech is
|
||||
* treated as ended anyway.
|
||||
*
|
||||
* A [ConversationAwarenessEvent.RESUME] frame (`5`) means speech resumed after a pause — in bursty
|
||||
* talking the pod cycles `3,5,3,5,…` while CA stays engaged. It is NOT a terminal: it cancels any
|
||||
* armed wind-down fuse and re-arms the long backstop, so media stays paused/ducked through the whole
|
||||
* conversation. (Treating `5` as a stop was the premature-resume + stuck bug.)
|
||||
*
|
||||
* Removing/inserting a pod makes the firmware re-key CA — it emits a terminal then a fresh onset
|
||||
* around the physical change even though the conversation continues. A terminal arriving within
|
||||
* [EAR_TRANSITION_WINDOW] of an AAP ear-detection change is therefore treated as a re-key artifact:
|
||||
* it defers to the wind-down fuse instead of disengaging, so a brief pod removal mid-conversation
|
||||
* neither strands a pause nor causes a duck→restore→re-duck volume blip.
|
||||
*
|
||||
* State is a single global slot (media volume / playback is system-wide, not per-device) guarded by
|
||||
* a [Mutex] — events, AAP-state-removal, the stale timer, and monitor completion all mutate it.
|
||||
@@ -78,16 +91,42 @@ class ConversationReaction @Inject constructor(
|
||||
val at: Long,
|
||||
)
|
||||
|
||||
/** What the single pending [disengageJob] is currently waiting for. */
|
||||
private enum class TimerPhase {
|
||||
/** Long backstop while only START/RESUME seen — inferred end if it fires. */
|
||||
STALE_BACKSTOP,
|
||||
|
||||
/** Short fuse armed by a wind-down HOLD — terminal imminent (or dropped, #608). */
|
||||
WIND_DOWN_FUSE,
|
||||
|
||||
/** Brief wait on an uncorroborated cold terminal to see if a pod-removal ear change follows. */
|
||||
STOP_SETTLE,
|
||||
}
|
||||
|
||||
private val mutex = Mutex()
|
||||
private var active: Active? = null
|
||||
private var staleJob: Job? = null
|
||||
private var disengageJob: Job? = null
|
||||
private var timerPhase: TimerPhase? = null
|
||||
// Bumped on every (re)arm. A timer job captures its generation and only acts if still current —
|
||||
// guards against a stale job that already passed its delay (so cancel() no longer stops it) being
|
||||
// re-armed to the SAME phase while it blocks on the mutex. This is a runtime concurrency guard: the
|
||||
// single-threaded virtual-time test harness serializes the timer callback and event handlers, so a
|
||||
// re-arm there always cancels the prior job before it runs — the race is not reproducible in a unit
|
||||
// test (a test would exercise cancel(), not this check). Left deliberately uncovered, not an oversight.
|
||||
private var timerGeneration = 0L
|
||||
private var idCounter = 0L
|
||||
|
||||
// Per-owner AAP ear-detection tracking, so a CA terminal that lands right after a pod
|
||||
// removal/insertion can be recognised as a re-key artifact rather than a real conversation end.
|
||||
private val lastEarDetection = mutableMapOf<BluetoothAddress, AapSetting.EarDetection?>()
|
||||
private val earTransitionAt = mutableMapOf<BluetoothAddress, Long>()
|
||||
|
||||
fun monitor(): Flow<Unit> = merge(
|
||||
aapManager.conversationalAwarenessEvents.onEach { (address, event) -> onEvent(address, event) },
|
||||
// Reverts a stranded duck when the owning device leaves the AAP state map for any reason
|
||||
// (intentional or not) — disconnectEvents only fires for unintentional drops.
|
||||
aapManager.allStates.onEach { states -> onActiveDevicesChanged(states.keys) },
|
||||
// Tracks ear-detection transitions (for terminal-artifact suppression) and reverts a stranded
|
||||
// duck when the owning device leaves the AAP state map for any reason (intentional or not) —
|
||||
// disconnectEvents only fires for unintentional drops.
|
||||
aapManager.allStates.onEach { states -> onStatesUpdated(states) },
|
||||
)
|
||||
.map { }
|
||||
// Service stop / scope cancellation: undo any active duck so we don't leave volume lowered.
|
||||
@@ -96,6 +135,7 @@ class ConversationReaction @Inject constructor(
|
||||
|
||||
private suspend fun onEvent(address: BluetoothAddress, event: ConversationAwarenessEvent) = when (event) {
|
||||
ConversationAwarenessEvent.START -> onSpeakingStart(address)
|
||||
ConversationAwarenessEvent.RESUME -> onSpeakingResume(address)
|
||||
ConversationAwarenessEvent.HOLD -> onSpeakingHold(address)
|
||||
ConversationAwarenessEvent.STOP -> onSpeakingStop(address)
|
||||
}
|
||||
@@ -113,13 +153,17 @@ class ConversationReaction @Inject constructor(
|
||||
val current = active
|
||||
if (current != null && current.owner == address) {
|
||||
// Duplicate START for the same speaker — don't re-act, just keep the session alive.
|
||||
// A START also cancels a pending wind-down fuse: the wearer is speaking again.
|
||||
armDisengageTimer(current, STALE_TIMEOUT)
|
||||
// A START also cancels a pending wind-down fuse / settle: the wearer is speaking again.
|
||||
armTimer(current, TimerPhase.STALE_BACKSTOP)
|
||||
log(TAG) { "START from $address — already active ($action), keep-alive" }
|
||||
return
|
||||
}
|
||||
// A different device started speaking while we were active — undo the old one first.
|
||||
if (current != null) revert(current, "superseded by $address")
|
||||
// A different device started speaking while we were active — undo the old one first and
|
||||
// cancel its pending timer (the no-op engage paths below would otherwise leak it).
|
||||
if (current != null) {
|
||||
revert(current, "superseded by $address")
|
||||
clearActive()
|
||||
}
|
||||
|
||||
when (action) {
|
||||
ConversationAction.PAUSE -> {
|
||||
@@ -127,7 +171,7 @@ class ConversationReaction @Inject constructor(
|
||||
if (paused) {
|
||||
val record = Active(nextId(), address, Kind.Paused, timeSource.elapsedRealtime())
|
||||
active = record
|
||||
armDisengageTimer(record, STALE_TIMEOUT)
|
||||
armTimer(record, TimerPhase.STALE_BACKSTOP)
|
||||
log(TAG, INFO) { "START on $address → paused media" }
|
||||
} else {
|
||||
active = null
|
||||
@@ -150,7 +194,7 @@ class ConversationReaction @Inject constructor(
|
||||
timeSource.elapsedRealtime(),
|
||||
)
|
||||
active = record
|
||||
armDisengageTimer(record, STALE_TIMEOUT)
|
||||
armTimer(record, TimerPhase.STALE_BACKSTOP)
|
||||
log(TAG, INFO) { "START on $address → ducked volume ${duck.priorVolume}→${duck.appliedVolume}" }
|
||||
} else {
|
||||
active = null
|
||||
@@ -164,14 +208,29 @@ class ConversationReaction @Inject constructor(
|
||||
}
|
||||
|
||||
/**
|
||||
* Transitional wind-down frame. The pod is silent during active speech, so this frame means the
|
||||
* wind-down has begun and a terminal frame is imminent — but firmware sometimes drops it (#608).
|
||||
* Arm the short fuse: if no terminal (or fresh START) follows, disengage anyway.
|
||||
* Speech resumed (status `5`) after a pause — the wind-down was aborted, the wearer is talking
|
||||
* again. Cancel any armed wind-down fuse and re-arm the long [STALE_TIMEOUT] backstop, exactly
|
||||
* like a duplicate-START keep-alive, so media stays paused/ducked through the conversation. Does
|
||||
* NOT engage from scratch: a RESUME with no active session is ignored (a conversation always
|
||||
* opens with a START, and a stray `5` should never start a pause/duck on its own).
|
||||
*/
|
||||
private suspend fun onSpeakingResume(address: BluetoothAddress) = mutex.withLock {
|
||||
val current = active ?: return
|
||||
if (current.owner != address) return
|
||||
armTimer(current, TimerPhase.STALE_BACKSTOP)
|
||||
log(TAG) { "RESUME from $address — speech resumed, keep-alive" }
|
||||
}
|
||||
|
||||
/**
|
||||
* Transitional wind-down frame (`3` pause, `0x0B`/`4` wind-down, `7` abort). The pod is silent
|
||||
* during active speech, so this frame means the wind-down may have begun and a terminal frame is
|
||||
* imminent — but firmware sometimes drops it (#608, single-pod wear). Arm the short fuse: if no
|
||||
* terminal (or a fresh START / RESUME) follows, disengage anyway.
|
||||
*/
|
||||
private suspend fun onSpeakingHold(address: BluetoothAddress) = mutex.withLock {
|
||||
val current = active ?: return
|
||||
if (current.owner != address) return
|
||||
armDisengageTimer(current, WIND_DOWN_TIMEOUT)
|
||||
armTimer(current, TimerPhase.WIND_DOWN_FUSE)
|
||||
}
|
||||
|
||||
private suspend fun onSpeakingStop(address: BluetoothAddress) {
|
||||
@@ -182,13 +241,42 @@ class ConversationReaction @Inject constructor(
|
||||
log(TAG) { "STOP from $address ignored — owner is ${current.owner}" }
|
||||
return
|
||||
}
|
||||
clearActive()
|
||||
disengage(current, primary, "STOP on $address")
|
||||
// The 0x06 ear-detection frame may have arrived but not yet been processed by the
|
||||
// allStates collector — fold the latest snapshot in before deciding (in-app ordering race).
|
||||
refreshEarTransition(address)
|
||||
|
||||
// Pulling/inserting a pod makes the firmware re-key CA — a terminal (8,9) then a fresh
|
||||
// onset (1,2) around the physical change, even though the conversation continues.
|
||||
when {
|
||||
// Ear change already seen → re-key artifact. Defer to the fuse; a re-onset cancels it.
|
||||
isRecentEarTransition(address) -> {
|
||||
armTimer(current, TimerPhase.WIND_DOWN_FUSE)
|
||||
log(TAG) { "STOP from $address during ear transition — deferring to fuse" }
|
||||
}
|
||||
// Corroborated by a preceding wind-down (3,0xB,4 / abort 7): a real end → resume now.
|
||||
timerPhase == TimerPhase.WIND_DOWN_FUSE -> {
|
||||
clearActive()
|
||||
disengage(current, primary, "STOP on $address", applyAgeGuard = false)
|
||||
}
|
||||
// Cold terminal with no wind-down and no ear change yet. On primary-pod removal the
|
||||
// firmware sends the terminal tens of ms BEFORE the ear-detection frame (measured ~30ms
|
||||
// on Pro 3), so wait a short settle to see if an ear change lands before treating it as
|
||||
// a real end.
|
||||
else -> {
|
||||
armTimer(current, TimerPhase.STOP_SETTLE)
|
||||
log(TAG) { "STOP from $address — uncorroborated, settling" }
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Graceful disengage (STOP event or stale timeout). Must be called under [mutex]. */
|
||||
private suspend fun disengage(record: Active, primary: PodDevice?, reason: String) {
|
||||
/**
|
||||
* Graceful disengage (STOP event or timer expiry). Must be called under [mutex]. [applyAgeGuard]
|
||||
* gates resume on [PAUSE_RESUME_WINDOW]: `true` only for the inferred stale-backstop end (we lost
|
||||
* track and timed out — surprise-resuming long after is worse than a stranded pause); `false` for
|
||||
* an explicit terminal or the wind-down fuse, where a recent/real end signal makes resume correct.
|
||||
*/
|
||||
private suspend fun disengage(record: Active, primary: PodDevice?, reason: String, applyAgeGuard: Boolean) {
|
||||
when (val kind = record.kind) {
|
||||
is Kind.Paused -> {
|
||||
// Resume the pause WE caused, regardless of the current action setting. Gating on
|
||||
@@ -197,11 +285,11 @@ class ConversationReaction @Inject constructor(
|
||||
// surprising behaviour. The remaining guards are about real device/playback state.
|
||||
val age = timeSource.elapsedRealtime() - record.at
|
||||
when {
|
||||
age.milliseconds > PAUSE_RESUME_WINDOW ->
|
||||
applyAgeGuard && age.milliseconds > PAUSE_RESUME_WINDOW ->
|
||||
log(TAG) { "$reason — resume skipped (stale, ${age}ms)" }
|
||||
primary?.address != record.owner ->
|
||||
log(TAG) { "$reason — resume skipped (primary switched)" }
|
||||
primary.isBeingWorn == false ->
|
||||
wornForResume(primary) == false ->
|
||||
log(TAG) { "$reason — resume skipped (not worn)" }
|
||||
mediaControl.isPlaying ->
|
||||
log(TAG) { "$reason — resume skipped (already playing)" }
|
||||
@@ -216,14 +304,35 @@ class ConversationReaction @Inject constructor(
|
||||
}
|
||||
}
|
||||
|
||||
private suspend fun onActiveDevicesChanged(addresses: Set<BluetoothAddress>) = mutex.withLock {
|
||||
private suspend fun onStatesUpdated(states: Map<BluetoothAddress, AapPodState>) = mutex.withLock {
|
||||
// Record when each device's AAP ear-detection last changed (pod removed/inserted/cased).
|
||||
for ((address, state) in states) recordEarIfChanged(address, state.aapEarDetection)
|
||||
lastEarDetection.keys.retainAll(states.keys)
|
||||
earTransitionAt.keys.retainAll(states.keys)
|
||||
|
||||
val current = active ?: return
|
||||
if (current.owner !in addresses) {
|
||||
if (current.owner !in states.keys) {
|
||||
clearActive()
|
||||
revert(current, "owner ${current.owner} gone")
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Must be called under [mutex]. Notes an AAP ear-detection transition for [address]. The very
|
||||
* first observation of a device only seeds the baseline — it is not counted as a transition.
|
||||
*/
|
||||
private fun recordEarIfChanged(address: BluetoothAddress, ear: AapSetting.EarDetection?) {
|
||||
if (lastEarDetection.containsKey(address) && lastEarDetection[address] != ear) {
|
||||
earTransitionAt[address] = timeSource.elapsedRealtime()
|
||||
}
|
||||
lastEarDetection[address] = ear
|
||||
}
|
||||
|
||||
/** Must be called under [mutex]. Folds the latest allStates snapshot in for [address]. */
|
||||
private fun refreshEarTransition(address: BluetoothAddress) {
|
||||
recordEarIfChanged(address, aapManager.allStates.value[address]?.aapEarDetection)
|
||||
}
|
||||
|
||||
private suspend fun onMonitorCompleted() = mutex.withLock {
|
||||
val current = active ?: return
|
||||
clearActive()
|
||||
@@ -249,41 +358,86 @@ class ConversationReaction @Inject constructor(
|
||||
mediaControl.restoreMusicVolume(kind.priorVolume)
|
||||
}
|
||||
|
||||
/** Must be called under [mutex]. Clears the active slot and cancels its stale timer. */
|
||||
/** Must be called under [mutex]. Clears the active slot and cancels its pending timer. */
|
||||
private fun clearActive() {
|
||||
staleJob?.cancel()
|
||||
staleJob = null
|
||||
disengageJob?.cancel()
|
||||
disengageJob = null
|
||||
timerPhase = null
|
||||
active = null
|
||||
}
|
||||
|
||||
/**
|
||||
* Must be called under [mutex]. (Re)arms the disengage timer for [record] — every frame picks
|
||||
* the fuse matching its meaning: START → [STALE_TIMEOUT] (active speech, frames cease for its
|
||||
* whole duration), HOLD → [WIND_DOWN_TIMEOUT] (wind-down begun, terminal imminent). On expiry
|
||||
* the session is force-ended. Identity-checked so a late timer can't disengage a newer session.
|
||||
* Must be called under [mutex]. (Re)arms the single disengage timer for [record] in [phase],
|
||||
* replacing whatever was pending. On expiry [onTimerExpired] decides per phase. Guarded on both
|
||||
* [Active.id] and [phase] so a late timer can't act on a newer session or a re-armed phase.
|
||||
*/
|
||||
private fun armDisengageTimer(record: Active, timeout: Duration) {
|
||||
staleJob?.cancel()
|
||||
staleJob = appScope.launch {
|
||||
private fun armTimer(record: Active, phase: TimerPhase) {
|
||||
disengageJob?.cancel()
|
||||
timerPhase = phase
|
||||
val generation = ++timerGeneration
|
||||
val timeout = when (phase) {
|
||||
TimerPhase.STALE_BACKSTOP -> STALE_TIMEOUT
|
||||
TimerPhase.WIND_DOWN_FUSE -> WIND_DOWN_TIMEOUT
|
||||
TimerPhase.STOP_SETTLE -> STOP_SETTLE_DELAY
|
||||
}
|
||||
disengageJob = appScope.launch {
|
||||
delay(timeout)
|
||||
val primary = deviceMonitor.primaryDevice().first()
|
||||
mutex.withLock {
|
||||
if (active?.id == record.id) {
|
||||
val current = active!!
|
||||
clearActive()
|
||||
disengage(current, primary, "no terminal frame within $timeout")
|
||||
}
|
||||
// generation guards against a stale job re-armed to the same phase for the same record.
|
||||
if (active?.id == record.id && timerGeneration == generation) onTimerExpired(record, phase, primary)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Must be called under [mutex] from inside the firing timer job. Detaches the slot WITHOUT
|
||||
* cancelling (we ARE that job — self-cancel could abort a following suspension), then acts:
|
||||
* - STOP_SETTLE + an ear change appeared meanwhile → re-key artifact, hand off to the fuse.
|
||||
* - otherwise force-end the session. The age guard applies only to the inferred STALE backstop;
|
||||
* the wind-down fuse and a settled cold terminal carry real/recent end evidence → resume regardless.
|
||||
*/
|
||||
private suspend fun onTimerExpired(record: Active, phase: TimerPhase, primary: PodDevice?) {
|
||||
disengageJob = null
|
||||
timerPhase = null
|
||||
if (phase == TimerPhase.STOP_SETTLE && isRecentEarTransition(record.owner)) {
|
||||
armTimer(record, TimerPhase.WIND_DOWN_FUSE)
|
||||
log(TAG) { "STOP settle elapsed for ${record.owner} — ear transition appeared, deferring to fuse" }
|
||||
return
|
||||
}
|
||||
active = null
|
||||
val reason = when (phase) {
|
||||
TimerPhase.STALE_BACKSTOP -> "no terminal within backstop"
|
||||
TimerPhase.WIND_DOWN_FUSE -> "no terminal within wind-down fuse"
|
||||
TimerPhase.STOP_SETTLE -> "cold terminal settled"
|
||||
}
|
||||
disengage(record, primary, reason, applyAgeGuard = phase == TimerPhase.STALE_BACKSTOP)
|
||||
}
|
||||
|
||||
/**
|
||||
* Whether the pod is worn enough to resume into, honoring One-Pod Mode (mirrors
|
||||
* [eu.darken.capod.reaction.core.autoconnect.AutoConnect]): with One-Pod Mode on, a single pod
|
||||
* in ear counts as worn, falling back to the both-pods reading only when the single-pod signal is
|
||||
* unknown. Returns `null` when ear state is genuinely unknown — the caller treats only an explicit
|
||||
* `false` as not-worn, so an unknown stays lenient (resumes).
|
||||
*/
|
||||
private fun wornForResume(device: PodDevice): Boolean? =
|
||||
if (device.reactions.onePodMode) device.isEitherPodInEar ?: device.isBeingWorn
|
||||
else device.isBeingWorn
|
||||
|
||||
/** Must be called under [mutex]. True if [address] had an AAP ear-detection change very recently. */
|
||||
private fun isRecentEarTransition(address: BluetoothAddress): Boolean {
|
||||
val at = earTransitionAt[address] ?: return false
|
||||
return (timeSource.elapsedRealtime() - at).milliseconds <= EAR_TRANSITION_WINDOW
|
||||
}
|
||||
|
||||
private fun nextId(): Long = ++idCounter
|
||||
|
||||
companion object {
|
||||
private val TAG = logTag("Reaction", "Conversation")
|
||||
|
||||
/**
|
||||
* Backstop fuse while engaged with no wind-down evidence yet (only START frames seen).
|
||||
* Backstop fuse while engaged with no wind-down evidence yet (START/RESUME frames seen).
|
||||
* Must stay LONG: the pod sends zero frames during active speech and stays engaged against
|
||||
* ambient noise, so a short timeout here resumes media mid-conversation (the original 12s
|
||||
* value did exactly that). Only recovers a session whose entire wind-down flurry was lost
|
||||
@@ -299,11 +453,38 @@ class ConversationReaction @Inject constructor(
|
||||
* pod deterministically drops the terminal (#608, reproduced on Pro 3 and Pro 2) — this
|
||||
* fuse disengages instead of stranding the volume low for [STALE_TIMEOUT]. Must be ≥ ~5s:
|
||||
* gaps up to 2.8s were observed between consecutive wind-down frames, and each HOLD re-arms
|
||||
* this fuse. A fresh START re-arms the long fuse (speaking resumed).
|
||||
* this fuse. A fresh START or RESUME (`5`, speech resumed) re-arms the long fuse instead.
|
||||
*/
|
||||
private val WIND_DOWN_TIMEOUT = 6.seconds
|
||||
|
||||
/** A pause older than this no longer auto-resumes — unexpected late playback is worse. */
|
||||
/**
|
||||
* For an *inferred* end only (the [STALE_TIMEOUT] backstop fired — we lost track of the
|
||||
* conversation), a pause older than this no longer auto-resumes: surprise playback long after
|
||||
* is worse than a stranded pause. An explicit terminal frame, or the short wind-down fuse,
|
||||
* resumes regardless of age — those carry real/recent evidence the conversation just ended.
|
||||
*/
|
||||
private val PAUSE_RESUME_WINDOW = 2.minutes
|
||||
|
||||
/**
|
||||
* A CA terminal landing within this window of an AAP ear-detection change is treated as a pod
|
||||
* removal/insertion re-key artifact (firmware emits 8,9 then 1,2 around the physical change),
|
||||
* not a real end. The artifact terminal and the ear change land within tens of ms of each other
|
||||
* (measured ~30ms on Pro 3, either side depending on which pod); a real end is seconds away.
|
||||
* Kept well under that gap so genuine terminals aren't deferred to the fuse. (The ear-detection
|
||||
* 0x06 frames are emitted on removal even when the EarDetectionEnabled setting is off, so this
|
||||
* suppression does not depend on that setting being enabled.)
|
||||
*/
|
||||
private val EAR_TRANSITION_WINDOW = 2.seconds
|
||||
|
||||
/**
|
||||
* How long to wait on an *uncorroborated cold* terminal (no preceding wind-down HOLD, no ear
|
||||
* change yet) before treating it as a real end. Covers the reverse frame-ordering on
|
||||
* primary-pod removal, where the firmware sends the terminal tens of ms BEFORE the ear-detection
|
||||
* frame (measured ~30ms on Pro 3) — too early for the backward-looking check. Only cold
|
||||
* terminals settle; a normal end
|
||||
* (`…3,0xB,4,8`) and an abort (`7,8`) are already corroborated by the fuse and resume instantly,
|
||||
* so this adds no latency to genuine conversation ends.
|
||||
*/
|
||||
private val STOP_SETTLE_DELAY = 250.milliseconds
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,7 +12,7 @@
|
||||
<clip>
|
||||
<shape android:shape="rectangle">
|
||||
<corners android:radius="4dp" />
|
||||
<solid android:color="?android:attr/colorAccent" />
|
||||
<solid android:color="@color/brand_primary" />
|
||||
</shape>
|
||||
</clip>
|
||||
</item>
|
||||
|
||||
@@ -24,8 +24,8 @@
|
||||
android:id="@+id/pod_left_icon"
|
||||
android:layout_width="24dp"
|
||||
android:layout_height="24dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/device_airpods_gen1_left" />
|
||||
|
||||
<ProgressBar
|
||||
@@ -58,8 +58,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_power_24" />
|
||||
|
||||
<ImageView
|
||||
@@ -67,8 +67,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_hearing_24" />
|
||||
|
||||
</LinearLayout>
|
||||
@@ -88,8 +88,8 @@
|
||||
android:id="@+id/pod_case_icon"
|
||||
android:layout_width="24dp"
|
||||
android:layout_height="24dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/device_airpods_gen1_case" />
|
||||
|
||||
<ProgressBar
|
||||
@@ -122,8 +122,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_power_24" />
|
||||
|
||||
</LinearLayout>
|
||||
@@ -142,8 +142,8 @@
|
||||
android:id="@+id/pod_right_icon"
|
||||
android:layout_width="24dp"
|
||||
android:layout_height="24dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/device_airpods_gen1_right" />
|
||||
|
||||
<ProgressBar
|
||||
@@ -176,8 +176,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_power_24" />
|
||||
|
||||
<ImageView
|
||||
@@ -185,8 +185,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_hearing_24" />
|
||||
|
||||
</LinearLayout>
|
||||
|
||||
@@ -38,6 +38,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_power_24" />
|
||||
|
||||
<ImageView
|
||||
@@ -46,6 +48,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_hearing_24" />
|
||||
</LinearLayout>
|
||||
|
||||
@@ -76,6 +80,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_power_24" />
|
||||
</LinearLayout>
|
||||
|
||||
@@ -105,6 +111,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_power_24" />
|
||||
|
||||
<ImageView
|
||||
@@ -113,6 +121,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_hearing_24" />
|
||||
</LinearLayout>
|
||||
|
||||
|
||||
@@ -24,8 +24,8 @@
|
||||
android:id="@+id/headphones_icon"
|
||||
android:layout_width="24dp"
|
||||
android:layout_height="24dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/twotone_headphones_24" />
|
||||
|
||||
<TextView
|
||||
@@ -71,8 +71,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_power_24" />
|
||||
|
||||
<ImageView
|
||||
@@ -80,8 +80,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tintMode="multiply"
|
||||
android:tint="?android:attr/colorAccent"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_hearing_24" />
|
||||
|
||||
</LinearLayout>
|
||||
|
||||
@@ -30,6 +30,8 @@
|
||||
android:id="@+id/headphones_battery_icon"
|
||||
style="@style/PodInfoItemIcon.Notification"
|
||||
android:layout_marginStart="12dp"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_battery_unknown_24" />
|
||||
|
||||
<TextView
|
||||
@@ -46,6 +48,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_power_24" />
|
||||
|
||||
<ImageView
|
||||
@@ -54,6 +58,8 @@
|
||||
android:layout_width="16dp"
|
||||
android:layout_height="16dp"
|
||||
android:layout_marginStart="2dp"
|
||||
android:tint="@color/notification_icon_tint"
|
||||
android:tintMode="src_in"
|
||||
android:src="@drawable/ic_baseline_hearing_24" />
|
||||
|
||||
</LinearLayout>
|
||||
|
||||
@@ -6,10 +6,10 @@
|
||||
<string name="general_copy_action">نسخ</string>
|
||||
<string name="general_thank_you_label">شكرًا لك</string>
|
||||
<string name="general_upgrade_action">ترقية</string>
|
||||
<string name="common_feature_requires_pro_msg">رقّي للفتح.</string>
|
||||
<string name="common_upgrade_required_label">يتطلب الترقية</string>
|
||||
<string name="general_donate_action">تبرع</string>
|
||||
<string name="general_check_action">تحقق</string>
|
||||
<string name="common_feature_requires_pro_msg">رقِّ للفتح.</string>
|
||||
<string name="common_upgrade_required_label">يتطلب ترقية</string>
|
||||
<string name="general_donate_action">تبرُّع</string>
|
||||
<string name="general_check_action">فحص</string>
|
||||
<string name="general_close_action">إغلاق</string>
|
||||
<string name="general_edit_action">تعديل</string>
|
||||
<string name="general_save_action">حفظ</string>
|
||||
|
||||
@@ -514,8 +514,8 @@
|
||||
<string name="reaction_sleep_notification_text">%1$s oznámil, že jste usnuli, takže vaše hudba byla pozastavena. Toto můžete vypnout v nastavení zařízení.</string>
|
||||
<string name="reaction_charged_channel_label">Oznámení o nabití</string>
|
||||
<string name="reaction_charged_notification_title">Nabíjení dokončeno</string>
|
||||
<string name="reaction_charged_notification_text_full">%1$s je plně nabitý.</string>
|
||||
<string name="reaction_charged_notification_text_partial">%1$s byl nabit na %2$d%%.</string>
|
||||
<string name="reaction_charged_notification_text_full">%1$s je plně nabito.</string>
|
||||
<string name="reaction_charged_notification_text_partial">%1$s nabito na %2$d%%.</string>
|
||||
<string name="device_settings_rename_label">Přejmenovat</string>
|
||||
<string name="device_settings_rename_hint">Název zařízení</string>
|
||||
<string name="device_settings_rename_confirm">Přejmenovat</string>
|
||||
|
||||
@@ -480,7 +480,7 @@
|
||||
<!-- New device settings -->
|
||||
<string name="device_settings_category_general_label">Allgemein</string>
|
||||
<string name="device_settings_microphone_mode_label">Mikrofon</string>
|
||||
<string name="device_settings_microphone_mode_description">Welches Ohrstöpsel als Mikrofon verwendet wird</string>
|
||||
<string name="device_settings_microphone_mode_description">Welcher Ohrstöpsel als Mikrofon verwendet wird</string>
|
||||
<string name="device_settings_microphone_mode_auto">Automatisch</string>
|
||||
<string name="device_settings_microphone_mode_right">Rechts</string>
|
||||
<string name="device_settings_microphone_mode_left">Links</string>
|
||||
|
||||
@@ -1,34 +1,34 @@
|
||||
<?xml version="1.0" encoding="utf-8"?>
|
||||
<resources xmlns:tools="http://schemas.android.com/tools" tools:ignore="MissingTranslation">
|
||||
<string name="general_share_action">Dalīties</string>
|
||||
<string name="general_share_action">Kopīgot</string>
|
||||
<string name="general_done_action">Gatavs</string>
|
||||
<string name="general_cancel_action">Atcelt</string>
|
||||
<string name="general_copy_action">Kopēt</string>
|
||||
<string name="general_copy_action">Ievietot starpliktuvē</string>
|
||||
<string name="general_thank_you_label">Paldies</string>
|
||||
<string name="general_upgrade_action">Jaunināt</string>
|
||||
<string name="common_feature_requires_pro_msg">Jauniniet, lai atbloķētu.</string>
|
||||
<string name="common_upgrade_required_label">Nepieciešams jauninājums</string>
|
||||
<string name="general_upgrade_action">Uzlabot</string>
|
||||
<string name="common_feature_requires_pro_msg">Jāuzlabo, lai atslēgtu.</string>
|
||||
<string name="common_upgrade_required_label">Nepieciešama uzlabošana</string>
|
||||
<string name="general_donate_action">Ziedot</string>
|
||||
<string name="general_check_action">Pārbaudīt</string>
|
||||
<string name="general_close_action">Aizvērt</string>
|
||||
<string name="general_edit_action">Rediģēt</string>
|
||||
<string name="general_edit_action">Labot</string>
|
||||
<string name="general_save_action">Saglabāt</string>
|
||||
<string name="general_guide_action">Ceļvedis</string>
|
||||
<string name="general_continue_action">Turpināt</string>
|
||||
<string name="general_show_action">Rādīt</string>
|
||||
<string name="general_hide_action">Paslēpt</string>
|
||||
<string name="general_example_label">Piem.: %s</string>
|
||||
<string name="upgrade_capod_label">Jaunināt CAPod</string>
|
||||
<string name="upgrade_capod_description">Iegūstiet papildu funkcijas un atbalstiet izstrādātāju.</string>
|
||||
<string name="upgrade_preamble">CAPod ir izstrādājusi viena persona. Jaunināšana atbloķē papildu funkcijas un palīdz uzturēt lietotni pie dzīvības.</string>
|
||||
<string name="upgrade_benefit_themes">Tēmu pielāgošana</string>
|
||||
<string name="upgrade_benefit_autoplay">Automātiska atskaņošana un pauze</string>
|
||||
<string name="upgrade_benefit_popups">Uznirstošs logs, atverot futrāli un savienojoties</string>
|
||||
<string name="upgrade_capod_label">Uzlabot CAPod</string>
|
||||
<string name="upgrade_capod_description">Iegūsti papildu iespējas un atbalsti izstrādātāju!</string>
|
||||
<string name="upgrade_preamble">CAPod ir izstrādājis viens cilvēks. Uzlabošana atslēdz papildu iespējas un palīdz uzturēt lietotnes esamību.</string>
|
||||
<string name="upgrade_benefit_themes">Izskata pielāgošana</string>
|
||||
<string name="upgrade_benefit_autoplay">Automātiska atskaņošana un apturēšana</string>
|
||||
<string name="upgrade_benefit_popups">Uznirstošs logs pēc futrāļa atvēršanas un savienošanās</string>
|
||||
<string name="upgrade_benefit_widgets">Sākuma ekrāna logrīki</string>
|
||||
<string name="upgrade_benefit_device_settings">Papildu ierīces iestatījumi</string>
|
||||
<string name="upgrade_benefit_device_controls">Ierīces vadīklēšana — troksnis, kātieni un vairāk</string>
|
||||
<string name="upgrade_benefit_support">Atbalstiet izstrādātāju</string>
|
||||
<string name="upgrade_benefit_disclaimer">Funkciju pieejamība atkarīga no jūsu ausţu un ierīces.</string>
|
||||
<string name="upgrade_benefit_device_controls">Ierīces vadīklas — troksnis, kātieni un vairāk</string>
|
||||
<string name="upgrade_benefit_support">Izstrādātāja atbalstīšana</string>
|
||||
<string name="upgrade_benefit_disclaimer">Iespēju pieejamība atkarīga no austiņam un ierīces.</string>
|
||||
<string name="upgrade_screen_options_description">Vienādas funkcijas, atšķirīgas cenas. Abonements ietver bezmaksas izmēģinājumu un var tikt atcelts jebkurā laikā.</string>
|
||||
<string name="upgrade_screen_subscription_trial_action">Sākt bezmaksas izmēģinājumu</string>
|
||||
<string name="upgrade_screen_subscription_action">Abonēt gadā</string>
|
||||
|
||||
@@ -0,0 +1,4 @@
|
||||
<?xml version="1.0" encoding="utf-8"?>
|
||||
<resources>
|
||||
<color name="notification_icon_tint">#EEEEEE</color>
|
||||
</resources>
|
||||
@@ -60,4 +60,6 @@
|
||||
<color name="brand_primary">#3f7aff</color>
|
||||
<color name="brand_secondary">#ffd83c</color>
|
||||
<color name="brand_tertiary">#41eb8e</color>
|
||||
|
||||
<color name="notification_icon_tint">#757575</color>
|
||||
</resources>
|
||||
@@ -13,15 +13,13 @@
|
||||
|
||||
<style name="Notification.Container">
|
||||
<!-- <item name="android:background">?android:attr/colorBackground</item>-->
|
||||
<item name="android:theme">@style/Theme.Material3.DynamicColors.DayNight</item>
|
||||
</style>
|
||||
|
||||
<style name="PodInfoItemIcon.Notification" parent="TextAppearance.Compat.Notification.Title">
|
||||
<item name="android:layout_height">20dp</item>
|
||||
<item name="android:layout_width">20dp</item>
|
||||
<item name="android:tintMode">multiply</item>
|
||||
<item name="android:tint">?android:attr/colorAccent</item>
|
||||
<item name="tint">?android:attr/colorAccent</item>
|
||||
<item name="android:tintMode">src_in</item>
|
||||
<item name="android:tint">@color/notification_icon_tint</item>
|
||||
</style>
|
||||
|
||||
<style name="PodInfoItemIcon.Notification.Large" parent="PodInfoItemIcon.Notification">
|
||||
@@ -32,7 +30,6 @@
|
||||
<style name="PodInfoItemText.Notification" parent="TextAppearance.Compat.Notification.Title">
|
||||
<item name="android:singleLine">true</item>
|
||||
<item name="android:ellipsize">end</item>
|
||||
<item name="android:textColor">?android:attr/textColorPrimary</item>
|
||||
</style>
|
||||
|
||||
<style name="PodInfoItemIcon" parent="TextAppearance.Material3.BodyMedium">
|
||||
|
||||
@@ -100,7 +100,7 @@ class MonitorModeResolverTest : BaseTest() {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `multi-profile primary-only - primary case 2 + secondary case 3 - AUTOMATIC`() = runTest {
|
||||
fun `first addressed profile controls mode - addressed primary case 2 + addressed secondary case 3 - AUTOMATIC`() = runTest {
|
||||
profilesFlow.value = listOf(
|
||||
profile(id = "primary", autoConnect = false),
|
||||
profile(id = "secondary", autoConnect = true),
|
||||
@@ -111,25 +111,59 @@ class MonitorModeResolverTest : BaseTest() {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `multi-profile primary-only - primary case 4 + secondary case 2 - MANUAL`() = runTest {
|
||||
fun `unpaired primary is skipped when an addressed secondary exists - AUTOMATIC`() = runTest {
|
||||
profilesFlow.value = listOf(
|
||||
profile(id = "primary", address = null),
|
||||
profile(id = "secondary", address = "AA:BB:CC:DD:EE:FF"),
|
||||
)
|
||||
|
||||
resolver.effectiveMode.first() shouldBe MonitorMode.AUTOMATIC
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `unpaired primary is skipped - addressed autoConnect secondary drives ALWAYS`() = runTest {
|
||||
profilesFlow.value = listOf(
|
||||
profile(id = "primary", address = null),
|
||||
profile(id = "secondary", address = "AA:BB:CC:DD:EE:FF", autoConnect = true),
|
||||
)
|
||||
nudgeFlow.value = NudgeAvailability.AVAILABLE
|
||||
|
||||
resolver.effectiveMode.first() shouldBe MonitorMode.ALWAYS
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `blank-address primary is skipped when an addressed secondary exists`() = runTest {
|
||||
profilesFlow.value = listOf(
|
||||
profile(id = "primary", address = " "),
|
||||
profile(id = "secondary", address = "AA:BB:CC:DD:EE:FF", autoConnect = true),
|
||||
)
|
||||
nudgeFlow.value = NudgeAvailability.AVAILABLE
|
||||
|
||||
resolver.effectiveMode.first() shouldBe MonitorMode.ALWAYS
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `all profiles unpaired - MANUAL`() = runTest {
|
||||
profilesFlow.value = listOf(
|
||||
profile(id = "a", address = null, autoConnect = true),
|
||||
profile(id = "b", address = "", autoConnect = true),
|
||||
)
|
||||
nudgeFlow.value = NudgeAvailability.AVAILABLE
|
||||
|
||||
resolver.effectiveMode.first() shouldBe MonitorMode.MANUAL
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `list reorder makes the secondary primary - mode flips`() = runTest {
|
||||
val a = profile(id = "a", autoConnect = true)
|
||||
val b = profile(id = "b", autoConnect = false)
|
||||
fun `reorder among addressed profiles still flips the mode, unpaired stays ignored`() = runTest {
|
||||
val unpaired = profile(id = "unpaired", address = null)
|
||||
val addressedAuto = profile(id = "addressedAuto", autoConnect = true)
|
||||
val addressedManual = profile(id = "addressedManual", autoConnect = false)
|
||||
nudgeFlow.value = NudgeAvailability.AVAILABLE
|
||||
|
||||
profilesFlow.value = listOf(a, b)
|
||||
profilesFlow.value = listOf(unpaired, addressedAuto, addressedManual)
|
||||
resolver.effectiveMode.first() shouldBe MonitorMode.ALWAYS
|
||||
|
||||
profilesFlow.value = listOf(b, a)
|
||||
profilesFlow.value = listOf(unpaired, addressedManual, addressedAuto)
|
||||
resolver.effectiveMode.first() shouldBe MonitorMode.AUTOMATIC
|
||||
}
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@ import eu.darken.capod.common.bluetooth.BluetoothManager2
|
||||
import eu.darken.capod.monitor.core.ble.BlePodMonitor
|
||||
import eu.darken.capod.pods.core.apple.PodModel
|
||||
import eu.darken.capod.pods.core.apple.aap.AapConnectionManager
|
||||
import eu.darken.capod.pods.core.apple.aap.AapDisconnectEvent
|
||||
import eu.darken.capod.pods.core.apple.aap.AapPodState
|
||||
import eu.darken.capod.pods.core.apple.aap.protocol.AapDeviceInfo
|
||||
import eu.darken.capod.pods.core.apple.ble.BlePodSnapshot
|
||||
@@ -22,7 +23,9 @@ import kotlinx.coroutines.flow.flowOf
|
||||
import kotlinx.coroutines.flow.toList
|
||||
import kotlinx.coroutines.launch
|
||||
import kotlinx.coroutines.test.UnconfinedTestDispatcher
|
||||
import kotlinx.coroutines.test.advanceTimeBy
|
||||
import kotlinx.coroutines.test.advanceUntilIdle
|
||||
import kotlinx.coroutines.test.runCurrent
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import org.junit.jupiter.api.BeforeEach
|
||||
import org.junit.jupiter.api.Nested
|
||||
@@ -40,7 +43,7 @@ class AapAutoConnectTest : BaseTest() {
|
||||
|
||||
private lateinit var profilesFlow: MutableStateFlow<List<DeviceProfile>>
|
||||
private lateinit var connectedDevicesFlow: MutableStateFlow<List<BluetoothDevice2>>
|
||||
private lateinit var disconnectEventsFlow: MutableSharedFlow<String>
|
||||
private lateinit var disconnectEventsFlow: MutableSharedFlow<AapDisconnectEvent>
|
||||
private lateinit var allStatesFlow: MutableStateFlow<Map<String, AapPodState>>
|
||||
|
||||
private val testAddress = "AA:BB:CC:DD:EE:FF"
|
||||
@@ -351,7 +354,7 @@ class AapAutoConnectTest : BaseTest() {
|
||||
// Clear profiles, then disconnect
|
||||
profilesFlow.value = emptyList()
|
||||
allStatesFlow.value = emptyMap()
|
||||
disconnectEventsFlow.tryEmit(testAddress)
|
||||
disconnectEventsFlow.tryEmit(AapDisconnectEvent(testAddress, wasEverReady = true))
|
||||
advanceUntilIdle()
|
||||
|
||||
// The reconnect loop should not call connect since profile is gone
|
||||
@@ -369,7 +372,7 @@ class AapAutoConnectTest : BaseTest() {
|
||||
// Remove bonded device, then disconnect
|
||||
every { bluetoothManager.bondedDevices() } returns flowOf(emptySet())
|
||||
allStatesFlow.value = emptyMap()
|
||||
disconnectEventsFlow.tryEmit(testAddress)
|
||||
disconnectEventsFlow.tryEmit(AapDisconnectEvent(testAddress, wasEverReady = true))
|
||||
advanceUntilIdle()
|
||||
|
||||
// Reconnect should not call connect since not bonded
|
||||
@@ -387,7 +390,7 @@ class AapAutoConnectTest : BaseTest() {
|
||||
// Remove classic BT connection, then disconnect
|
||||
connectedDevicesFlow.value = emptyList()
|
||||
allStatesFlow.value = emptyMap()
|
||||
disconnectEventsFlow.tryEmit(testAddress)
|
||||
disconnectEventsFlow.tryEmit(AapDisconnectEvent(testAddress, wasEverReady = true))
|
||||
advanceUntilIdle()
|
||||
|
||||
// Reconnect should not call connect since not classically connected
|
||||
@@ -405,7 +408,7 @@ class AapAutoConnectTest : BaseTest() {
|
||||
// Remove from BLE scans, then disconnect
|
||||
every { blePodMonitor.devices } returns flowOf(emptyList())
|
||||
allStatesFlow.value = emptyMap()
|
||||
disconnectEventsFlow.tryEmit(testAddress)
|
||||
disconnectEventsFlow.tryEmit(AapDisconnectEvent(testAddress, wasEverReady = true))
|
||||
advanceUntilIdle()
|
||||
|
||||
// Reconnect should not call connect since not visible in BLE
|
||||
@@ -422,7 +425,7 @@ class AapAutoConnectTest : BaseTest() {
|
||||
|
||||
// Keep as READY, emit disconnect event
|
||||
// allStatesFlow still shows READY → reconnect should skip
|
||||
disconnectEventsFlow.tryEmit(testAddress)
|
||||
disconnectEventsFlow.tryEmit(AapDisconnectEvent(testAddress, wasEverReady = true))
|
||||
advanceUntilIdle()
|
||||
|
||||
// Should not attempt connect — already connected
|
||||
@@ -430,6 +433,44 @@ class AapAutoConnectTest : BaseTest() {
|
||||
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `was-ready disconnect reconnects without extra cooldown`() = runTest(testDispatcher) {
|
||||
val autoConnect = createAutoConnect()
|
||||
val job = setupForReconnect(autoConnect)
|
||||
advanceUntilIdle()
|
||||
allStatesFlow.value = emptyMap()
|
||||
|
||||
// A session that reached READY drops: only the normal first retry delay (3s), no backoff cooldown.
|
||||
disconnectEventsFlow.tryEmit(AapDisconnectEvent(testAddress, wasEverReady = true))
|
||||
advanceTimeBy(3_100)
|
||||
runCurrent()
|
||||
|
||||
coVerify(exactly = 1) { aapManager.connect(testAddress, any(), any()) }
|
||||
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `never-ready disconnect backs off before reconnecting`() = runTest(testDispatcher) {
|
||||
val autoConnect = createAutoConnect()
|
||||
val job = setupForReconnect(autoConnect)
|
||||
advanceUntilIdle()
|
||||
allStatesFlow.value = emptyMap()
|
||||
|
||||
// A session that never reached READY (failed handshake): first failure adds a 5s cooldown
|
||||
// on top of the 3s retry delay, so nothing should connect within the 3.1s a was-ready drop would.
|
||||
disconnectEventsFlow.tryEmit(AapDisconnectEvent(testAddress, wasEverReady = false))
|
||||
advanceTimeBy(3_100)
|
||||
runCurrent()
|
||||
coVerify(exactly = 0) { aapManager.connect(testAddress, any(), any()) }
|
||||
|
||||
// After the 5s cooldown + 3s retry delay elapse, the reconnect fires.
|
||||
advanceUntilIdle()
|
||||
coVerify(exactly = 1) { aapManager.connect(testAddress, any(), any()) }
|
||||
|
||||
job.cancel()
|
||||
}
|
||||
}
|
||||
|
||||
@Nested
|
||||
|
||||
@@ -12,8 +12,15 @@ import eu.darken.capod.pods.core.apple.aap.protocol.AapSetting
|
||||
import io.kotest.assertions.throwables.shouldThrow
|
||||
import io.kotest.matchers.maps.shouldBeEmpty
|
||||
import io.kotest.matchers.shouldBe
|
||||
import io.mockk.Runs
|
||||
import io.mockk.every
|
||||
import io.mockk.just
|
||||
import io.mockk.mockk
|
||||
import kotlinx.coroutines.CoroutineScope
|
||||
import kotlinx.coroutines.Dispatchers
|
||||
import kotlinx.coroutines.SupervisorJob
|
||||
import kotlinx.coroutines.cancel
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.test.TestScope
|
||||
import kotlinx.coroutines.TimeoutCancellationException
|
||||
import kotlinx.coroutines.test.UnconfinedTestDispatcher
|
||||
@@ -27,10 +34,12 @@ import testhelpers.TestTimeSource
|
||||
import java.io.ByteArrayInputStream
|
||||
import java.io.ByteArrayOutputStream
|
||||
import java.io.IOException
|
||||
import java.io.InputStream
|
||||
import java.util.concurrent.CountDownLatch
|
||||
import java.util.concurrent.TimeUnit
|
||||
import java.util.concurrent.atomic.AtomicInteger
|
||||
import kotlin.time.Duration.Companion.milliseconds
|
||||
import kotlin.time.Duration.Companion.seconds
|
||||
|
||||
class AapConnectionManagerTest : BaseTest() {
|
||||
|
||||
@@ -135,6 +144,65 @@ class AapConnectionManagerTest : BaseTest() {
|
||||
advanceUntilIdle()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `handshake timeout disconnects a silent session`() = runBlocking {
|
||||
// Socket connects and the handshake is sent, but the peer never replies: read() blocks
|
||||
// until the socket is closed. The handshake watchdog must time out and tear the session
|
||||
// down. This exercises real dispatchers + a real socket close, so it runs on REAL time
|
||||
// with generous bounds — mixing runTest's virtual time with Dispatchers.IO made it flaky.
|
||||
val readBlocked = CountDownLatch(1)
|
||||
val closeCalls = AtomicInteger(0)
|
||||
val blockingInput = object : InputStream() {
|
||||
// Block well past the test's own wait so the watchdog (not a read self-timeout) is
|
||||
// always what ends the read.
|
||||
override fun read(): Int {
|
||||
readBlocked.await(30, TimeUnit.SECONDS)
|
||||
return -1
|
||||
}
|
||||
|
||||
override fun read(b: ByteArray): Int {
|
||||
readBlocked.await(30, TimeUnit.SECONDS)
|
||||
return -1
|
||||
}
|
||||
}
|
||||
val silentSocket = mockk<BluetoothSocket>(relaxed = true) {
|
||||
every { connect() } just Runs
|
||||
every { outputStream } returns ByteArrayOutputStream()
|
||||
every { inputStream } returns blockingInput
|
||||
every { close() } answers {
|
||||
closeCalls.incrementAndGet()
|
||||
readBlocked.countDown()
|
||||
}
|
||||
}
|
||||
every { socketFactory.createSocket(any(), any()) } returns silentSocket
|
||||
val connection = AapConnection(
|
||||
device = testDevice,
|
||||
profile = AapDeviceProfile.forModel(PodModel.AIRPODS_PRO3),
|
||||
socketFactory = socketFactory,
|
||||
timeSource = timeSource,
|
||||
connectTimeout = 1.seconds,
|
||||
handshakeTimeout = 200.milliseconds,
|
||||
)
|
||||
|
||||
val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
||||
try {
|
||||
connection.connect(scope)
|
||||
|
||||
// The 200ms watchdog tears the silent session down: disconnect() resets the engine to
|
||||
// DISCONNECTED and THEN closes the socket (the close mock increments before counting
|
||||
// the latch down). Awaiting that latch is the unambiguous "watchdog fired" signal and
|
||||
// sidesteps any reset-vs-close ordering race. Generous real-time bound for loaded CI.
|
||||
readBlocked.await(15, TimeUnit.SECONDS) shouldBe true
|
||||
connection.state.value.connectionState shouldBe AapPodState.ConnectionState.DISCONNECTED
|
||||
// close() may run once (watchdog) or twice (watchdog + read-loop finally race on
|
||||
// cleanupSocket before socket is nulled) — both are correct, so assert "at least once".
|
||||
(closeCalls.get() >= 1) shouldBe true
|
||||
} finally {
|
||||
readBlocked.countDown()
|
||||
scope.cancel()
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `remote disconnect cleans up allStates`() = testScope.runTest {
|
||||
// Empty inputStream → readLoop gets -1 immediately → DISCONNECTED
|
||||
|
||||
+2
@@ -412,6 +412,8 @@ class DefaultAapDeviceProfileTest : BaseAapSessionTest() {
|
||||
@Nested
|
||||
inner class ConversationAwarenessStateTests {
|
||||
@Test fun `speaking start`() { decodeSetting<AapSetting.ConversationalAwarenessState>(aapMessage("04 00 04 00 4B 00 01")).speaking shouldBe true }
|
||||
@Test fun `speaking resume`() { decodeSetting<AapSetting.ConversationalAwarenessState>(aapMessage("04 00 04 00 4B 00 05")).speaking shouldBe true }
|
||||
@Test fun `speaking pause`() { decodeSetting<AapSetting.ConversationalAwarenessState>(aapMessage("04 00 04 00 4B 00 03")).speaking shouldBe false }
|
||||
@Test fun `speaking stop`() { decodeSetting<AapSetting.ConversationalAwarenessState>(aapMessage("04 00 04 00 4B 00 04")).speaking shouldBe false }
|
||||
@Test fun `empty payload returns null`() { profile.decodeSetting(aapMessage("04 00 04 00 4B 00")).shouldBeNull() }
|
||||
}
|
||||
|
||||
+51
-7
@@ -114,6 +114,38 @@ class AapSessionEngineTest : BaseTest() {
|
||||
engine.state.value.connectionState shouldBe AapPodState.ConnectionState.DISCONNECTED
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `wasEverReady is false until READY then true`() {
|
||||
val engine = createEngine()
|
||||
val scope = TestScope(UnconfinedTestDispatcher())
|
||||
|
||||
engine.start(scope)
|
||||
engine.wasEverReady shouldBe false
|
||||
engine.onHandshakeSent()
|
||||
engine.wasEverReady shouldBe false
|
||||
|
||||
// First non-CONTROL message during HANDSHAKING → READY
|
||||
engine.processMessage(dummyMessage(commandType = 0x0002))
|
||||
engine.wasEverReady shouldBe true
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `wasEverReady survives reset`() {
|
||||
val engine = createEngine()
|
||||
val scope = TestScope(UnconfinedTestDispatcher())
|
||||
|
||||
engine.start(scope)
|
||||
engine.onHandshakeSent()
|
||||
engine.processMessage(dummyMessage(commandType = 0x0002))
|
||||
engine.wasEverReady shouldBe true
|
||||
|
||||
engine.reset()
|
||||
|
||||
engine.state.value.connectionState shouldBe AapPodState.ConnectionState.DISCONNECTED
|
||||
// Consumers read this at disconnect time — reset must not clear it.
|
||||
engine.wasEverReady shouldBe true
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `reset is idempotent`() {
|
||||
val engine = createEngine()
|
||||
@@ -220,12 +252,13 @@ class AapSessionEngineTest : BaseTest() {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `Connect Response with non-zero status does not store features`() {
|
||||
fun `Connect Response with non-zero status tears down the session`() {
|
||||
val engine = createEngine()
|
||||
val scope = TestScope(UnconfinedTestDispatcher())
|
||||
|
||||
engine.start(scope)
|
||||
engine.onHandshakeSent()
|
||||
engine.state.value.connectionState shouldBe AapPodState.ConnectionState.HANDSHAKING
|
||||
|
||||
val packet = AapPacket.ConnectResponse(
|
||||
raw = ByteArray(18),
|
||||
@@ -237,9 +270,9 @@ class AapSessionEngineTest : BaseTest() {
|
||||
)
|
||||
engine.processConnectResponse(packet)
|
||||
|
||||
engine.state.value.connectResponseStatus shouldBe 0x0001
|
||||
// Hard protocol rejection: fast-fail to DISCONNECTED instead of waiting out the watchdog.
|
||||
engine.state.value.connectionState shouldBe AapPodState.ConnectionState.DISCONNECTED
|
||||
engine.state.value.negotiatedFeatures.shouldBeNull()
|
||||
engine.state.value.connectionState shouldBe AapPodState.ConnectionState.HANDSHAKING
|
||||
}
|
||||
}
|
||||
|
||||
@@ -529,9 +562,17 @@ class AapSessionEngineTest : BaseTest() {
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `terminal statuses 5, 6, 8, 9 emit STOP`() = runTest(UnconfinedTestDispatcher()) {
|
||||
// 5 is the terminal wind-down value on fw …6861 (never reaches 6/8/9); 6/8/9 on fw …6503.
|
||||
firstEventFor(5) shouldBe ConversationAwarenessEvent.STOP
|
||||
fun `status 5 emits RESUME`() = runTest(UnconfinedTestDispatcher()) {
|
||||
// 5 = speech resumed after a pause (bursty talking cycles 3,5,3,5,…); NOT a terminal.
|
||||
// Labelled captures across Pro 3 + Pro 2 USB-C show 5 only ever inside an active
|
||||
// conversation, never ending one. Misreading it as STOP was the premature-resume bug.
|
||||
firstEventFor(5) shouldBe ConversationAwarenessEvent.RESUME
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `terminal statuses 6, 8 and 9 emit STOP`() = runTest(UnconfinedTestDispatcher()) {
|
||||
// The usual terminal is the 8→9 pair (both pods; also emitted by pod removal / case-close).
|
||||
// 6 is a standalone terminal seen live on Pro 2 USB-C fw …6814.
|
||||
firstEventFor(6) shouldBe ConversationAwarenessEvent.STOP
|
||||
firstEventFor(8) shouldBe ConversationAwarenessEvent.STOP
|
||||
firstEventFor(9) shouldBe ConversationAwarenessEvent.STOP
|
||||
@@ -539,11 +580,14 @@ class AapSessionEngineTest : BaseTest() {
|
||||
|
||||
@Test
|
||||
fun `transitional and unknown statuses emit HOLD (stay engaged)`() = runTest(UnconfinedTestDispatcher()) {
|
||||
// Must never disengage on these — only an explicit terminal STOP does.
|
||||
// Must never resume immediately on these — they arm the short wind-down fuse instead.
|
||||
// 3 = pause, 0x0B/4 = wind-down, 7 = abort; any unknown value defaults to HOLD.
|
||||
firstEventFor(3) shouldBe ConversationAwarenessEvent.HOLD
|
||||
firstEventFor(4) shouldBe ConversationAwarenessEvent.HOLD
|
||||
firstEventFor(0x0B) shouldBe ConversationAwarenessEvent.HOLD
|
||||
firstEventFor(7) shouldBe ConversationAwarenessEvent.HOLD
|
||||
firstEventFor(0) shouldBe ConversationAwarenessEvent.HOLD
|
||||
firstEventFor(0xFF) shouldBe ConversationAwarenessEvent.HOLD
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+90
@@ -0,0 +1,90 @@
|
||||
package eu.darken.capod.pods.core.apple.ble.history
|
||||
|
||||
import eu.darken.capod.common.fromHex
|
||||
import eu.darken.capod.pods.core.apple.PodModel
|
||||
import eu.darken.capod.pods.core.apple.ble.devices.BaseBlePodsTest
|
||||
import eu.darken.capod.pods.core.apple.ble.devices.airpods.AirPodsPro2Usbc
|
||||
import eu.darken.capod.profiles.core.AppleDeviceProfile
|
||||
import io.kotest.matchers.shouldBe
|
||||
import io.kotest.matchers.shouldNotBe
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import org.junit.jupiter.api.Test
|
||||
|
||||
/**
|
||||
* Reproduces Symptom B (the dashboard-card glitch): a FOREIGN same-model AirPods, whose BLE address
|
||||
* does not resolve against our IRK, gets recovered as OUR device's stable identity purely via the
|
||||
* fuzzy/case-ignored history fallback ([PodHistoryRepo.baseSearch]). The foreign snapshot carries
|
||||
* `meta.profile == null` (IRK failed) but adopts our identifier, so it overwrites our BLE cache slot
|
||||
* keyed by that identifier — and the merged device flaps to `ble = null` (cached-only) until our next
|
||||
* real frame re-binds.
|
||||
*
|
||||
* The two frames use the SAME advertisement payload (so the low-entropy fuzzy markers — model bits,
|
||||
* pod/case battery nibbles, colour — are identical) but DIFFERENT BLE addresses, mimicking another
|
||||
* person's same-model AirPods at a similar battery level in a crowded environment.
|
||||
*/
|
||||
class PodHistoryFuzzyCollisionTest : BaseBlePodsTest() {
|
||||
|
||||
// IRK + a resolvable RPA pair lifted from RPACheckerTest.
|
||||
private val irkHex = "79-04-65-1E-E2-CC-D9-26-F2-6E-20-EE-3E-CC-DE-79"
|
||||
private val myRpa = "5A:16:2B:91:D1:CD" // resolves against irkHex
|
||||
private val foreignAddress = "77:49:4C:D8:25:0C" // does NOT resolve against irkHex
|
||||
|
||||
// A valid AirPods Pro 2 (USB-C) proximity advertisement (model 0x2420), reused for both frames.
|
||||
private val payload = "07 19 01 24 20 0B 99 8F 11 00 04 BD A7 3B FF 2D 8A 3C AF 9B 1A 7C 74 B7 A9 D1 C3"
|
||||
|
||||
@Test
|
||||
fun `foreign same-model frame must not adopt a keyed device's identity`() = runTest {
|
||||
profileList.add(
|
||||
AppleDeviceProfile(
|
||||
label = "Mine",
|
||||
model = PodModel.AIRPODS_PRO2_USBC,
|
||||
identityKey = irkHex.fromHex(),
|
||||
address = "AA:BB:CC:DD:EE:FF",
|
||||
)
|
||||
)
|
||||
|
||||
// 1) Our own frame: address resolves against our IRK -> IRK-bound, registers history.
|
||||
lateinit var mine: AirPodsPro2Usbc
|
||||
create<AirPodsPro2Usbc>(payload, address = myRpa) { mine = this }
|
||||
mine.meta.isIRKMatch shouldBe true
|
||||
mine.meta.profile shouldNotBe null
|
||||
|
||||
// 2) A foreign device's frame: same payload (identical fuzzy markers), different address.
|
||||
lateinit var foreign: AirPodsPro2Usbc
|
||||
create<AirPodsPro2Usbc>(payload, address = foreignAddress) { foreign = this }
|
||||
|
||||
// The foreign frame is correctly NOT IRK-bound...
|
||||
foreign.meta.isIRKMatch shouldBe false
|
||||
foreign.meta.profile shouldBe null
|
||||
|
||||
// ...and therefore must NOT inherit our stable identity. Before the gate this FAILED: the
|
||||
// case-ignored/fuzzy fallback recovered our KnownDevice for the foreign frame, so it adopted
|
||||
// our identifier and poisoned our cache slot -> the merged device glitched to ble=null.
|
||||
foreign.identifier shouldNotBe mine.identifier
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `our own keyed device keeps its identity across an RPA rotation`() = runTest {
|
||||
// Guard: narrowing the weak paths must not fragment our own identity — across an RPA
|
||||
// rotation our device is still recovered via the strong IRK path.
|
||||
profileList.add(
|
||||
AppleDeviceProfile(
|
||||
label = "Mine",
|
||||
model = PodModel.AIRPODS_PRO2_USBC,
|
||||
identityKey = irkHex.fromHex(),
|
||||
address = "AA:BB:CC:DD:EE:FF",
|
||||
)
|
||||
)
|
||||
val myRotatedRpa = "45:23:51:E3:40:6E" // also resolves against irkHex (computed)
|
||||
|
||||
lateinit var first: AirPodsPro2Usbc
|
||||
create<AirPodsPro2Usbc>(payload, address = myRpa) { first = this }
|
||||
first.meta.isIRKMatch shouldBe true
|
||||
|
||||
lateinit var second: AirPodsPro2Usbc
|
||||
create<AirPodsPro2Usbc>(payload, address = myRotatedRpa) { second = this }
|
||||
second.meta.isIRKMatch shouldBe true
|
||||
// Same physical device across rotation -> same stable identity (recovered via IRK).
|
||||
second.identifier shouldBe first.identifier
|
||||
}
|
||||
}
|
||||
+466
-10
@@ -6,6 +6,7 @@ import eu.darken.capod.monitor.core.DeviceMonitor
|
||||
import eu.darken.capod.monitor.core.PodDevice
|
||||
import eu.darken.capod.pods.core.apple.aap.AapConnectionManager
|
||||
import eu.darken.capod.pods.core.apple.aap.AapPodState
|
||||
import eu.darken.capod.pods.core.apple.aap.protocol.AapSetting
|
||||
import eu.darken.capod.pods.core.apple.aap.protocol.ConversationAwarenessEvent
|
||||
import eu.darken.capod.profiles.core.ReactionConfig
|
||||
import io.mockk.coEvery
|
||||
@@ -31,11 +32,16 @@ class ConversationReactionTest : BaseTest() {
|
||||
private val primaryAddress: BluetoothAddress = "AA:BB:CC:DD:EE:FF"
|
||||
private val otherAddress: BluetoothAddress = "11:22:33:44:55:66"
|
||||
|
||||
private val inEar = AapSetting.EarDetection.PodPlacement.IN_EAR
|
||||
private val outOfEar = AapSetting.EarDetection.PodPlacement.NOT_IN_EAR
|
||||
|
||||
// Mirrors of ConversationReaction's private timing constants. STALE_TIMEOUT: long backstop
|
||||
// while only START frames were seen. WIND_DOWN_TIMEOUT: short fuse armed by a HOLD frame
|
||||
// (wind-down begun, terminal imminent — but sometimes dropped, #608).
|
||||
// (wind-down begun, terminal imminent — but sometimes dropped, #608). STOP_SETTLE_DELAY: brief
|
||||
// wait on an uncorroborated cold terminal (no preceding HOLD) before treating it as a real end.
|
||||
private val staleTimeoutMs = 5L * 60 * 1000
|
||||
private val windDownTimeoutMs = 6_000L
|
||||
private val stopSettleMs = 250L
|
||||
|
||||
private lateinit var eventsFlow: MutableSharedFlow<Pair<BluetoothAddress, ConversationAwarenessEvent>>
|
||||
private lateinit var statesFlow: MutableStateFlow<Map<BluetoothAddress, AapPodState>>
|
||||
@@ -50,14 +56,18 @@ class ConversationReactionTest : BaseTest() {
|
||||
action: ConversationAction,
|
||||
reduction: Int = 50,
|
||||
worn: Boolean = true,
|
||||
onePodMode: Boolean = false,
|
||||
eitherInEar: Boolean? = worn,
|
||||
): PodDevice = mockk(relaxed = true) {
|
||||
every { profileId } returns address
|
||||
every { this@mockk.address } returns address
|
||||
every { reactions } returns ReactionConfig(
|
||||
conversationAction = action,
|
||||
conversationVolumeReduction = reduction,
|
||||
onePodMode = onePodMode,
|
||||
)
|
||||
every { isBeingWorn } returns worn
|
||||
every { isEitherPodInEar } returns eitherInEar
|
||||
}
|
||||
|
||||
@BeforeEach
|
||||
@@ -95,6 +105,38 @@ class ConversationReactionTest : BaseTest() {
|
||||
runCurrent()
|
||||
}
|
||||
|
||||
/**
|
||||
* Seeds the AAP ear-detection snapshot WITHOUT draining the collector. Call before [launchReaction]
|
||||
* to establish a realistic worn baseline so the first pull registers as a real in→out transition
|
||||
* (the very first observation only seeds the baseline, it is not counted as a transition).
|
||||
*/
|
||||
private fun setEarDetectionSnapshot(primary: AapSetting.EarDetection.PodPlacement, secondary: AapSetting.EarDetection.PodPlacement) {
|
||||
statesFlow.value = mapOf(
|
||||
primaryAddress to mockk(relaxed = true) {
|
||||
every { aapEarDetection } returns AapSetting.EarDetection(primary, secondary)
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
/** Simulates an AAP ear-detection transition (pod removed/inserted) for the primary device. */
|
||||
private fun TestScope.changeEarDetection(primary: AapSetting.EarDetection.PodPlacement, secondary: AapSetting.EarDetection.PodPlacement) {
|
||||
setEarDetectionSnapshot(primary, secondary)
|
||||
runCurrent()
|
||||
}
|
||||
|
||||
/**
|
||||
* Advances BOTH the coroutine scheduler (drives the delay-based timers: STOP_SETTLE, the wind-down
|
||||
* fuse, the stale backstop) AND the [TestTimeSource] that backs `isRecentEarTransition`'s
|
||||
* `elapsedRealtime()` reads. Advancing only the coroutine clock leaves the recorded ear-transition
|
||||
* age frozen at 0, which makes settle / EAR_TRANSITION_WINDOW assertions pass regardless of the real
|
||||
* durations. The wall clock is advanced first so a timer firing mid-advance sees the full elapsed age.
|
||||
*/
|
||||
private fun TestScope.advanceBoth(ms: Long) {
|
||||
timeSource.advanceBy(java.time.Duration.ofMillis(ms))
|
||||
advanceTimeBy(ms)
|
||||
runCurrent()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `LOWER_VOLUME start ducks, stop restores`() = runTest(UnconfinedTestDispatcher()) {
|
||||
val job = launchReaction()
|
||||
@@ -103,7 +145,9 @@ class ConversationReactionTest : BaseTest() {
|
||||
verify(exactly = 1) { mediaControl.duckMusicVolume(50) }
|
||||
verify(exactly = 0) { mediaControl.restoreMusicVolume(any()) }
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // cold terminal → settles briefly
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
verify(exactly = 1) { mediaControl.restoreMusicVolume(10) }
|
||||
job.cancel()
|
||||
}
|
||||
@@ -116,7 +160,9 @@ class ConversationReactionTest : BaseTest() {
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // cold terminal → settles briefly
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
|
||||
verify(exactly = 1) { mediaControl.restoreMusicVolume(10) }
|
||||
job.cancel()
|
||||
@@ -245,7 +291,9 @@ class ConversationReactionTest : BaseTest() {
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
coVerify(exactly = 1) { mediaControl.sendPause(false) }
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // cold terminal → settles briefly
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
@@ -257,11 +305,88 @@ class ConversationReactionTest : BaseTest() {
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `PAUSE one-pod mode resumes via wind-down fuse when the single worn pod drops the terminal`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// #608 single-pod wear: both-pods isBeingWorn is false, but One-Pod Mode means the single
|
||||
// in-ear pod counts as worn, so the fuse-driven resume must fire. Was skipped as "not worn"
|
||||
// because the guard read isBeingWorn (both pods) instead of honoring One-Pod Mode.
|
||||
devicesFlow.value = listOf(
|
||||
mockPodDevice(
|
||||
primaryAddress,
|
||||
ConversationAction.PAUSE,
|
||||
worn = false,
|
||||
onePodMode = true,
|
||||
eitherInEar = true,
|
||||
),
|
||||
)
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
coVerify(exactly = 1) { mediaControl.sendPause(false) }
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // wind-down begins, terminal dropped
|
||||
|
||||
advanceTimeBy(windDownTimeoutMs + 500)
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `PAUSE one-pod mode stays paused via fuse when no pod is worn`() = runTest(UnconfinedTestDispatcher()) {
|
||||
// One-Pod Mode on, but neither pod in ear (both isBeingWorn and isEitherPodInEar false) — the
|
||||
// single-pod leniency must NOT resume into pods that are out.
|
||||
devicesFlow.value = listOf(
|
||||
mockPodDevice(
|
||||
primaryAddress,
|
||||
ConversationAction.PAUSE,
|
||||
worn = false,
|
||||
onePodMode = true,
|
||||
eitherInEar = false,
|
||||
),
|
||||
)
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD)
|
||||
advanceTimeBy(windDownTimeoutMs + 500)
|
||||
runCurrent()
|
||||
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `PAUSE one-pod mode resumes on terminal after wind-down with a single worn pod`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// The timerPhase==WIND_DOWN_FUSE branch disengages immediately on the explicit terminal
|
||||
// (no STOP_SETTLE). Verify the One-Pod Mode worn check applies on that shared path too.
|
||||
devicesFlow.value = listOf(
|
||||
mockPodDevice(
|
||||
primaryAddress,
|
||||
ConversationAction.PAUSE,
|
||||
worn = false,
|
||||
onePodMode = true,
|
||||
eitherInEar = true,
|
||||
),
|
||||
)
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // wind-down → fuse armed
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // real terminal → immediate disengage
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `PAUSE stop does not resume when something is already playing`() = runTest(UnconfinedTestDispatcher()) {
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
@@ -270,6 +395,8 @@ class ConversationReactionTest : BaseTest() {
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
every { mediaControl.isPlaying } returns true // user/app restarted playback during the talk
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
@@ -285,7 +412,9 @@ class ConversationReactionTest : BaseTest() {
|
||||
|
||||
// User changes the action mid-conversation; we must still undo the pause WE caused.
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.LOWER_VOLUME))
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // cold terminal → settles briefly
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
@@ -338,9 +467,9 @@ class ConversationReactionTest : BaseTest() {
|
||||
@Test
|
||||
fun `PAUSE stays paused through frame silence, resumes only on explicit STOP, then re-engages`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// Regression for the fw …6861 bug: the pod sends an onset, then NO frames for ~20s while
|
||||
// the wearer keeps talking, then a terminal STOP. The old 12s stale timeout resumed media
|
||||
// mid-speech; the backstop must not, and a fresh talk must re-arm.
|
||||
// The pod sends NO frames during continuous speech (29s silent gaps observed), then a
|
||||
// terminal STOP. The old 12s stale timeout resumed media mid-speech; the long backstop
|
||||
// must not, and a fresh talk must re-arm.
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
@@ -352,7 +481,9 @@ class ConversationReactionTest : BaseTest() {
|
||||
runCurrent()
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() } // NOT resumed mid-speech
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // wearer stopped → pod's terminal frame
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // cold terminal → settles briefly
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START) // a fresh talk re-arms
|
||||
@@ -377,6 +508,283 @@ class ConversationReactionTest : BaseTest() {
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `RESUME keeps media paused through a bursty conversation, resumes only at the real terminal`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// Regression for the fw …6861 status-5 bug. Bursty talking emits 1,2 then 3,5 (pause,
|
||||
// resume) pairs while CA stays engaged, ending with the real wind-down 3,0xB,4,8,9.
|
||||
// Status 5 was misclassified as a terminal STOP, so media resumed on the first burst
|
||||
// pause and — with no fresh 1/2 onset mid-conversation — never paused again. RESUME must
|
||||
// keep media paused until the genuine terminal.
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START) // 1
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START) // 2
|
||||
coVerify(exactly = 1) { mediaControl.sendPause(false) }
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 3 pause
|
||||
emit(primaryAddress, ConversationAwarenessEvent.RESUME) // 5 resume
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 3 pause
|
||||
emit(primaryAddress, ConversationAwarenessEvent.RESUME) // 5 resume
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() } // stayed paused through the bursts
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 3
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 0x0B
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 4
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // 8 terminal
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `RESUME cancels the wind-down fuse`() = runTest(UnconfinedTestDispatcher()) {
|
||||
// A pause (3) arms the short fuse; a resume (5) must cancel it and switch back to the long
|
||||
// backstop — otherwise media resumes ~6s into renewed speech.
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 3 — wind-down fuse armed
|
||||
advanceTimeBy(windDownTimeoutMs * 2 / 3)
|
||||
runCurrent()
|
||||
emit(primaryAddress, ConversationAwarenessEvent.RESUME) // 5 — speech resumed, cancel fuse
|
||||
advanceTimeBy(windDownTimeoutMs * 2) // well past the original fuse
|
||||
runCurrent()
|
||||
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `wind-down after a RESUME still disengages via the fuse (dropped terminal)`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// After a resume re-arms the long backstop, a later genuine wind-down (0xB,4 with the
|
||||
// 8,9 terminal dropped — single-pod) must still disengage via the short fuse. Proves the
|
||||
// RESUME keep-alive doesn't permanently disable #608 recovery.
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 3 pause
|
||||
emit(primaryAddress, ConversationAwarenessEvent.RESUME) // 5 resume → long backstop
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 3 pause again
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 0x0B wind-down
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 4 — terminal dropped
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() }
|
||||
|
||||
advanceTimeBy(windDownTimeoutMs + 500)
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `PAUSE resumes on the terminal even after a conversation longer than the resume window`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// A bursty conversation that runs past PAUSE_RESUME_WINDOW: the explicit STOP must still
|
||||
// resume. The age guard applies only to the inferred stale backstop, not to a real
|
||||
// terminal. Advance BOTH clocks so the window is genuinely exercised (TestTimeSource
|
||||
// drives `age`).
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
coVerify(exactly = 1) { mediaControl.sendPause(false) }
|
||||
|
||||
repeat(3) {
|
||||
timeSource.advanceBy(java.time.Duration.ofSeconds(60))
|
||||
advanceTimeBy(60_000) // < STALE_TIMEOUT, so the backstop never fires
|
||||
runCurrent()
|
||||
emit(primaryAddress, ConversationAwarenessEvent.RESUME)
|
||||
}
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() } // 3 min in, still paused
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // cold terminal → settles briefly
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `PAUSE resumes on a direct STOP after long frame silence`() = runTest(UnconfinedTestDispatcher()) {
|
||||
// Continuous speech sends no frames; then a direct terminal (no preceding wind-down) arrives
|
||||
// past PAUSE_RESUME_WINDOW but before STALE_TIMEOUT. An explicit STOP is positive evidence CA
|
||||
// ended now, so it must resume regardless of engage-age. (Fails if the age guard is applied to
|
||||
// explicit terminals.)
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
coVerify(exactly = 1) { mediaControl.sendPause(false) }
|
||||
|
||||
timeSource.advanceBy(java.time.Duration.ofSeconds(150)) // > 2-min window, < 5-min backstop
|
||||
advanceTimeBy(150_000)
|
||||
runCurrent()
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() } // still paused, no frames yet
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // cold terminal → settles briefly
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `PAUSE does not resume when the stale backstop fires past the resume window`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// The inferred end (5-min backstop, no terminal ever) keeps the age guard: a long-stale
|
||||
// pause must NOT surprise-resume.
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
coVerify(exactly = 1) { mediaControl.sendPause(false) }
|
||||
|
||||
timeSource.advanceBy(java.time.Duration.ofMinutes(5))
|
||||
advanceTimeBy(5L * 60 * 1000 + 500) // STALE_TIMEOUT fires
|
||||
runCurrent()
|
||||
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() } // inferred stale end → not resumed
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `RESUME without a prior start is a no-op`() = runTest(UnconfinedTestDispatcher()) {
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.RESUME) // stray 5, nothing engaged
|
||||
|
||||
coVerify(exactly = 0) { mediaControl.sendPause(any()) }
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `LOWER_VOLUME RESUME keeps the volume ducked`() = runTest(UnconfinedTestDispatcher()) {
|
||||
// Same status-5 bug seen on the default action: it restored volume on the first 3→5 pause.
|
||||
val job = launchReaction() // devicesFlow default = LOWER_VOLUME
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
verify(exactly = 1) { mediaControl.duckMusicVolume(any()) }
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 3
|
||||
emit(primaryAddress, ConversationAwarenessEvent.RESUME) // 5
|
||||
verify(exactly = 0) { mediaControl.restoreMusicVolume(any()) }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `PAUSE pod removal mid-conversation does not strand media paused`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// Pulling a pod re-keys CA (terminal then fresh onset) while the conversation continues.
|
||||
// The terminal lands right after the ear-detection change → defer to the fuse, the re-onset
|
||||
// cancels it, media stays paused and the session alive. The real end later resumes.
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
coVerify(exactly = 1) { mediaControl.sendPause(false) }
|
||||
|
||||
changeEarDetection(inEar, outOfEar) // pod pulled
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // re-key terminal
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START) // re-onset, still talking
|
||||
advanceTimeBy(windDownTimeoutMs * 2) // fuse would have fired — but was cancelled
|
||||
runCurrent()
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() } // still paused, not stranded-resumed early
|
||||
|
||||
timeSource.advanceBy(java.time.Duration.ofSeconds(3)) // past EAR_TRANSITION_WINDOW
|
||||
advanceTimeBy(3_000)
|
||||
runCurrent()
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // genuine end, no recent ear change
|
||||
advanceTimeBy(stopSettleMs + 50) // cold terminal → settles briefly
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() } // resumes
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `PAUSE primary-pod removal where the terminal precedes the ear change is caught by the settle`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// Ground truth (on-device, AirPods Pro 3, primary/left pod): the doubled re-key terminal
|
||||
// (8,9) arrives ~30ms BEFORE the 0x06 ear frame, so the backward-looking check misses it and
|
||||
// onSpeakingStop takes the cold "settling" branch. The ear change then lands DURING the
|
||||
// 250ms settle, and at settle expiry the transition is recent → deferred to the fuse, not a
|
||||
// resume. advanceBoth is required: isRecentEarTransition reads TestTimeSource, so advancing
|
||||
// only coroutine time would freeze the transition age at 0 and pass regardless of durations.
|
||||
setEarDetectionSnapshot(inEar, inEar) // both pods worn baseline
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
coVerify(exactly = 1) { mediaControl.sendPause(false) }
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // terminal first → cold, settling
|
||||
advanceBoth(30)
|
||||
changeEarDetection(outOfEar, inEar) // ear change lands during the settle
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // second terminal frame
|
||||
advanceBoth(stopSettleMs) // settle expires; transition is recent
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() } // deferred to fuse, NOT resumed
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START) // re-onset on remaining pod
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
advanceBoth(windDownTimeoutMs * 2) // fuse would fire — re-onset cancelled it
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() } // still paused, session alive
|
||||
|
||||
// Genuine end, past EAR_TRANSITION_WINDOW so the terminal is fuse-corroborated (immediate).
|
||||
advanceBoth(2_100)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() } // resumes exactly once, only here
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `LOWER_VOLUME pod removal mid-conversation does not blip the volume`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// The re-key terminal must not restore (which the immediate re-onset would then re-duck —
|
||||
// an audible 5→10→5 jump). Defer to the fuse; the re-onset keeps it ducked.
|
||||
val job = launchReaction() // default LOWER_VOLUME
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
verify(exactly = 1) { mediaControl.duckMusicVolume(any()) }
|
||||
|
||||
changeEarDetection(inEar, outOfEar) // pod pulled
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // re-key terminal
|
||||
verify(exactly = 0) { mediaControl.restoreMusicVolume(any()) } // no restore → no blip
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START) // re-onset
|
||||
verify(exactly = 1) { mediaControl.duckMusicVolume(any()) } // not re-ducked (keep-alive only)
|
||||
verify(exactly = 0) { mediaControl.restoreMusicVolume(any()) }
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `terminal during ear transition still disengages via the fuse if no re-onset follows`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// Removed a pod AND the conversation actually ended (no re-onset): the deferred terminal's
|
||||
// fuse must still disengage within seconds rather than waiting for the long backstop.
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
coVerify(exactly = 1) { mediaControl.sendPause(false) }
|
||||
|
||||
changeEarDetection(inEar, outOfEar)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // deferred to fuse
|
||||
coVerify(exactly = 0) { mediaControl.sendPlay() }
|
||||
|
||||
advanceTimeBy(windDownTimeoutMs + 500)
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() } // fuse disengaged
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `STOP from a non-owner does not disengage the active owner`() = runTest(UnconfinedTestDispatcher()) {
|
||||
devicesFlow.value = listOf(mockPodDevice(primaryAddress, ConversationAction.PAUSE))
|
||||
@@ -395,7 +803,9 @@ class ConversationReactionTest : BaseTest() {
|
||||
val job = launchReaction()
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // cold terminal → settles briefly
|
||||
advanceTimeBy(stopSettleMs + 50)
|
||||
runCurrent()
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() }
|
||||
|
||||
advanceTimeBy(staleTimeoutMs + 500) // backstop would fire if STOP hadn't cancelled it
|
||||
@@ -403,4 +813,50 @@ class ConversationReactionTest : BaseTest() {
|
||||
coVerify(exactly = 1) { mediaControl.sendPlay() } // still only once
|
||||
job.cancel()
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `LOWER_VOLUME secondary-pod removal and reinsert holds the duck, restores only at the real end`() =
|
||||
runTest(UnconfinedTestDispatcher()) {
|
||||
// Ground truth (on-device, AirPods Pro 3, secondary/right pod): the 0x06 ear change lands
|
||||
// ~30ms BEFORE the doubled re-key terminal (8,9), so the backward isRecentEarTransition
|
||||
// check defers each terminal to the fuse; the doubled re-onset (1,2) keeps the session
|
||||
// alive. The duck must be restored only at the genuine wind-down end — never during the
|
||||
// pull or the reinsert (otherwise the immediate re-onset re-ducks → an audible 5→10→5 jump).
|
||||
setEarDetectionSnapshot(inEar, inEar) // both pods worn baseline
|
||||
val job = launchReaction() // default LOWER_VOLUME
|
||||
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
verify(exactly = 1) { mediaControl.duckMusicVolume(any()) }
|
||||
|
||||
// Pull: ear-out first, then the doubled terminal, then the doubled re-onset.
|
||||
changeEarDetection(inEar, outOfEar)
|
||||
advanceBoth(30)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
verify(exactly = 0) { mediaControl.restoreMusicVolume(any()) }
|
||||
|
||||
// Reinsert ~1.5s later: same ear-before-terminal ordering.
|
||||
advanceBoth(1_500)
|
||||
changeEarDetection(inEar, inEar)
|
||||
advanceBoth(30)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.START)
|
||||
verify(exactly = 0) { mediaControl.restoreMusicVolume(any()) }
|
||||
|
||||
// Genuine end, past EAR_TRANSITION_WINDOW so the terminal is fuse-corroborated (immediate).
|
||||
advanceBoth(2_100)
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 3
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 0x0B
|
||||
emit(primaryAddress, ConversationAwarenessEvent.HOLD) // 4
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // 8
|
||||
emit(primaryAddress, ConversationAwarenessEvent.STOP) // 9
|
||||
verify(exactly = 1) { mediaControl.restoreMusicVolume(any()) } // restored exactly once, here
|
||||
verify(exactly = 1) { mediaControl.duckMusicVolume(any()) } // engaged exactly once total
|
||||
job.cancel()
|
||||
}
|
||||
}
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
### Updated by tools/release/bump.sh ###
|
||||
project.versioning.major=5
|
||||
project.versioning.minor=1
|
||||
project.versioning.patch=8
|
||||
project.versioning.patch=10
|
||||
project.versioning.build=0
|
||||
project.versioning.type=rc
|
||||
#############################
|
||||
|
||||
Reference in New Issue
Block a user