Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -146,11 +146,15 @@ class AACPManager(private val device: AppleDevice) {
}

} else if (bytesRead == -1) {
Log.i("AirPodsService", "socket closed (bytesRead = -1)")
// isConnected can stay true after the remote end closes the channel, so without breaking
// this spins on read() returning -1 until something else closes the socket
Log.i(TAG, "socket closed (bytesRead = -1), stopping read loop")
break
}
} catch (e: Exception) {
Log.i(TAG, "Error reading data, we have probably disconnected.")
e.printStackTrace()
break
}
}
}
Expand Down Expand Up @@ -1181,7 +1185,16 @@ class AACPManager(private val device: AppleDevice) {
val payload = data.command.payload.toByteArray()
val timestamp = Clock.System.now()
if (payload.size == 18) {
val heartRate = payload[1].toInt()
// unsigned: a signed read turns anything >= 128 bpm negative and drops it
val heartRate = payload[1].toInt() and 0xFF

// bit 0 of the last byte is set for the first few samples after the sensor starts (payload[2],
// which looks like a confidence value, is also very low then); those readings are unreliable,
// e.g. 169 bpm at rest, so skip them and stay in WAITING until the sensor settles
if (payload[17].toInt() and 0x01 != 0) {
Log.d(TAG, "skipping heart rate sample while the sensor is acquiring: ${payload.toHexString()}")
return
}

// same as healthconnect's datatype. 300 isn't possible anyway, but whatever
if (heartRate !in 1..300) {
Expand All @@ -1202,7 +1215,7 @@ class AACPManager(private val device: AppleDevice) {

Log.i(
TAG,
"hr: $heartRate bpm"
"hr: $heartRate bpm, payload: ${payload.toHexString()}"
)

val heartRateSample = HeartRateSample(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asSharedFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.getAndUpdate
import kotlinx.coroutines.flow.update
import kotlinx.coroutines.launch
import me.kavishdevar.librepods.bluetooth.MacAddress
Expand Down Expand Up @@ -124,8 +125,13 @@ class AppleDevice(
}

override fun connect(): Boolean {
_connectionState.update {
ConnectionState.CONNECTING
// the service (ACL/UUID broadcasts) and the UI can both call connect(); only let one of them open a socket
val previousState = _connectionState.getAndUpdate {
if (it == ConnectionState.CONNECTING || it == ConnectionState.CONNECTED) it else ConnectionState.CONNECTING
}
if (previousState == ConnectionState.CONNECTING || previousState == ConnectionState.CONNECTED) {
Log.d(TAG, "connect() ignored, already $previousState")
return previousState == ConnectionState.CONNECTED
}

val success = aacp.connect() // && att.connect()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,13 @@ import android.content.Context
import android.content.Intent
import android.content.IntentFilter
import android.content.pm.PackageManager
import android.media.AudioDeviceInfo
import android.media.AudioManager
import android.os.BatteryManager
import android.os.Binder
import android.os.Build
import android.os.IBinder
import android.os.Looper
import android.os.ParcelUuid
import android.os.ext.SdkExtensions
import android.provider.Settings
Expand All @@ -44,6 +46,7 @@ import kotlinx.coroutines.Job
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.update
import kotlinx.coroutines.launch
import me.kavishdevar.librepods.LibrePodsApplication
Expand Down Expand Up @@ -74,11 +77,13 @@ import me.kavishdevar.librepods.utils.MediaController
import me.kavishdevar.librepods.utils.calculateLevel
import me.kavishdevar.librepods.utils.redactMac
import java.time.ZoneOffset
import java.util.concurrent.ConcurrentHashMap
import kotlin.time.Duration
import kotlin.time.Duration.Companion.seconds
import kotlin.time.toJavaInstant

private const val TAG = "LibrePodsService"
private val A2DP_RECONNECT_TIMEOUT = 15.seconds

@SuppressLint("MissingPermission")
class LibrePodsService: Service() {
Expand All @@ -92,6 +97,7 @@ class LibrePodsService: Service() {
val devices = _devices.asStateFlow()

private val deviceJobs = mutableMapOf<MacAddress, MutableList<Job>>()
private val devicesBeingConnected: MutableSet<MacAddress> = ConcurrentHashMap.newKeySet()

val irkMap = mutableMapOf<MacAddress, ByteArray>()
val rpasByPublicMac = mutableMapOf<MacAddress, MutableSet<MacAddress>>()
Expand Down Expand Up @@ -272,59 +278,77 @@ class LibrePodsService: Service() {
return
}

// ACL_CONNECTED and ACTION_UUID can both arrive for the same connection; set the device up once at a time
if (!devicesBeingConnected.add(device.macAddress)) {
Log.d(TAG, "Device already being connected: ${bluetoothDevice.address}")
return
}

when (device) {
is AppleDevice -> CoroutineScope(Dispatchers.IO).launch {
Log.i(TAG, "Loading device ${device.macAddress.toRedactedString()} from db")
try {
Log.i(TAG, "Loading device ${device.macAddress.toRedactedString()} from db")

appleRepository.load(device.macAddress)?.let { entity ->
val cache = entity.cache
Log.i(
TAG,
"Loaded cached state for device ${device.macAddress.toRedactedString()}: $cache"
)
val settings = entity.settings
Log.i(
TAG,
"Loaded settings for device ${device.macAddress.toRedactedString()}: $settings"
)
val metadata = entity.metadata
Log.i(
TAG,
"Loaded metadata for device ${device.macAddress.toRedactedString()}: $metadata"
)

device.loadInitialState(
state = AppleState().copy(
capabilities = cache.capabilities,
magicKeys = cache.magicKeys,
controlStates = cache.controlStates,
),
settings = settings,
metadata = metadata
)
}

// createDevice() already registered observers for this device; cancel them so each state change is handled once
deviceJobs[MacAddress(bluetoothDevice.address)]?.forEach { it.cancel() }
deviceJobs[MacAddress(bluetoothDevice.address)] = mutableListOf()

deviceJobs[MacAddress(bluetoothDevice.address)]?.add(observeAppleState(device))
deviceJobs[MacAddress(bluetoothDevice.address)]?.add(observeAppleSettings(device))
deviceJobs[MacAddress(bluetoothDevice.address)]?.add(observeAppleMetadata(device))
deviceJobs[MacAddress(bluetoothDevice.address)]?.add(observeAppleMicrophoneFrames(device))

// connect only after the cached state is loaded and the observers are registered, otherwise
// loadInitialState() can overwrite what the first packets set and the observers miss those changes.
// This also keeps the blocking socket connect off the main thread (this runs from a broadcast receiver).
device.connect()
// connect() returns right away if another caller (e.g. the device list) is already connecting,
// so wait for that attempt to finish before using the connection
val connected = device.connectionState.first { it != ConnectionState.CONNECTING } == ConnectionState.CONNECTED

appleRepository.load(device.macAddress)?.let { entity ->
val cache = entity.cache
Log.i(
TAG,
"Loaded cached state for device ${device.macAddress.toRedactedString()}: $cache"
)
val settings = entity.settings
Log.i(
TAG,
"Loaded settings for device ${device.macAddress.toRedactedString()}: $settings"
)
val metadata = entity.metadata
Log.i(
TAG,
"Loaded metadata for device ${device.macAddress.toRedactedString()}: $metadata"
"Device ${if (connected) "connected" else "failed to connect"}: ${device.macAddress.toRedactedString()} (${device.javaClass.simpleName})"
)

device.loadInitialState(
state = AppleState().copy(
capabilities = cache.capabilities,
magicKeys = cache.magicKeys,
controlStates = cache.controlStates,
),
settings = settings,
metadata = metadata
)
_devices.update { it + (device.macAddress to device) }

if (device.settings.value.hrmAlertEnabled) {
if (connected && device.settings.value.hrmAlertEnabled) {
device.startHr()
}
} finally {
devicesBeingConnected.remove(device.macAddress)
}

deviceJobs[MacAddress(bluetoothDevice.address)] = mutableListOf()

deviceJobs[MacAddress(bluetoothDevice.address)]?.add(observeAppleState(device))
deviceJobs[MacAddress(bluetoothDevice.address)]?.add(observeAppleSettings(device))
deviceJobs[MacAddress(bluetoothDevice.address)]?.add(observeAppleMetadata(device))
deviceJobs[MacAddress(bluetoothDevice.address)]?.add(observeAppleMicrophoneFrames(device))
}
}

device.connect()

Log.i(
TAG,
"Device connected: ${device.macAddress.toRedactedString()} (${device.javaClass.simpleName})"
)

_devices.update { it + (device.macAddress to device) }
}

private fun onDeviceDisconnected(mac: MacAddress) {
Expand Down Expand Up @@ -676,7 +700,9 @@ class LibrePodsService: Service() {

Log.d(TAG, "updating island window")
if (islandWindow?.isVisible == true) {
islandWindow?.updateBattery(state.battery)
CoroutineScope(Dispatchers.Main).launch {
islandWindow?.updateBattery(state.battery)
}
}

Log.d(TAG, "updating notification")
Expand Down Expand Up @@ -919,6 +945,14 @@ class LibrePodsService: Service() {
reversed: Boolean = false,
otherDeviceName: String? = null
) {
// the island is a window, so it has to be added from the main thread (state observers run on IO)
if (Looper.myLooper() != Looper.getMainLooper()) {
CoroutineScope(Dispatchers.Main).launch {
showIsland(device, type, reversed, otherDeviceName)
}
return
}

Log.d(TAG, "Showing island window")

val state = device.state.value
Expand Down Expand Up @@ -1221,7 +1255,9 @@ class LibrePodsService: Service() {
}

if (new == EarPresence.NONE && islandWindow?.isVisible == true) {
islandWindow?.close()
CoroutineScope(Dispatchers.Main).launch {
islandWindow?.close()
}
}

var justEnabledA2dp = false
Expand All @@ -1233,20 +1269,38 @@ class LibrePodsService: Service() {
"User put in at least one component, enabling audio for device ${device.macAddress.toRedactedString()}"
)
device.enableAudio()
device.connectA2dp()
// both profiles: disconnectAudio() drops A2DP and the headset profile when all components are taken out
device.connectAudio()
justEnabledA2dp = true

device.waitForA2dpConnection(this) {
MediaController.sendPlay()
MediaController.iPausedTheMedia = false
// the heart rate sensor only streams while worn; a request made while the buds were in the case
// (e.g. on connect) doesn't start delivering once they are put in, so request it again now
if (device is AppleDevice && device.settings.value.hrmAlertEnabled) {
device.startHr()
}

if (MediaController.getMusicActive()) {
MediaController.userPlayedTheMedia = true
}
if (new == EarPresence.PARTIAL) {

if (isA2dpAudioConnected(device.macAddress)) {
MediaController.sendPlay()
MediaController.iPausedTheMedia = false
} else {
// resuming before A2DP is up starts playback on the phone speaker, so wait for it
val receiver = device.waitForA2dpConnection(this) {
MediaController.sendPlay()
MediaController.iPausedTheMedia = false
}
// if A2DP doesn't connect, don't leave the receiver around to resume playback much later
CoroutineScope(Dispatchers.Main).launch {
delay(A2DP_RECONNECT_TIMEOUT)
try {
unregisterReceiver(receiver)
} catch (_: IllegalArgumentException) {
// already unregistered itself after A2DP connected
}
}
}
}

Expand Down Expand Up @@ -1283,6 +1337,11 @@ class LibrePodsService: Service() {
}
}

private fun isA2dpAudioConnected(macAddress: MacAddress): Boolean =
getSystemService(AudioManager::class.java)
.getDevices(AudioManager.GET_DEVICES_OUTPUTS)
.any { it.type == AudioDeviceInfo.TYPE_BLUETOOTH_A2DP && it.address.equals(macAddress.value, ignoreCase = true) }

private fun processHeartRateSample(heartRateSample: HeartRateSample, interval: Duration, alertThreshold: Int) {
CoroutineScope(Dispatchers.IO).launch {
Log.d(TAG, "inserting to local db")
Expand Down