Skip to content
Merged
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
1 change: 1 addition & 0 deletions app/build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,7 @@ dependencies {
implementation(libs.kotlinx.serialization.json)

implementation(libs.lib.snapcast.android)
implementation(libs.lib.shairport.android)
implementation("com.github.capullo-tech.lib-librespot-android:librespot-android:0.2.0")
implementation(libs.slf4j.handroid)

Expand Down
2 changes: 2 additions & 0 deletions app/src/main/AndroidManifest.xml
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@

<uses-permission android:name="android.permission.ACCESS_WIFI_STATE" />
<uses-permission android:name="android.permission.CHANGE_WIFI_STATE" />
<!-- MulticastLock for shairport-sync's in-process mDNS responder -->
<uses-permission android:name="android.permission.CHANGE_WIFI_MULTICAST_STATE" />
<uses-permission android:name="android.permission.NEARBY_WIFI_DEVICES"
android:usesPermissionFlags="neverForLocation"
tools:targetApi="tiramisu" />
Expand Down
2 changes: 1 addition & 1 deletion app/src/main/java/tech/capullo/radio/RadioNavHost.kt
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,7 @@ fun RadioCapulloNavHost() {
schemeChoice = SchemeChoice.GREEN,
) {
EspotiSessionLoadingScreen(
onPlayerReady = {
onProceedToBroadcast = {
backStack.add(Broadcast)
},
)
Expand Down
216 changes: 216 additions & 0 deletions app/src/main/java/tech/capullo/radio/airplay/AirplayMetadataParser.kt
Original file line number Diff line number Diff line change
@@ -0,0 +1,216 @@
package tech.capullo.radio.airplay

import kotlinx.serialization.json.JsonArray
import kotlinx.serialization.json.JsonPrimitive
import tech.capullo.radio.snapcast.ArtData
import tech.capullo.radio.snapcast.StreamMetadata
import tech.capullo.radio.snapcast.StreamProperties

/**
* Parses shairport-sync's metadata pipe stream (emitted when built with
* CONFIG_METADATA) into Snapcast [StreamProperties]. Pure / no Android deps so
* it can be unit-tested; base64 decoding is injected.
*
* shairport writes newline-terminated items (see rtsp.c `metadata_process`):
*
* <item><type>HEX</type><code>HEX</code><length>N</length>
* <data encoding="base64">
* BASE64</data></item>
*
* where the `<data>` block is absent for zero-length items, the base64 has no
* internal newlines, and HEX is the big-endian uint32 of a 4-char code. Codes
* arrive under two types: `core` (DAAP metadata) and `ssnc` (shairport-sync
* control). Feed raw pipe text in arbitrary chunks; complete items are parsed
* and a cumulative snapshot is returned.
*/
class AirplayMetadataParser(private val decodeBase64: (String) -> ByteArray) {

private val buffer = StringBuilder()

// Cumulative player state, updated as items arrive.
private var title: String? = null
private var artist: String? = null
private var album: String? = null
private var durationSec: Float? = null
private var positionSec: Float = 0f
private var artDataBase64: String? = null
private var artExtension: String = "jpg"
private var volumePercent: Int? = null
private var muted: Boolean = false

// null = nothing has played yet (reported as "stopped").
private var playing: Boolean? = null

// DACP remote-control credentials from the connected sender (ssnc acre/daid).
private var activeRemote: String? = null
private var dacpId: String? = null

/** DACP credentials, non-null only once both acre and daid have arrived. */
fun credentials(): DacpCredentials? {
val ar = activeRemote
val id = dacpId
return if (ar != null && id != null) DacpCredentials(id, ar) else null
}

/**
* Feeds a chunk of raw pipe text. Returns a cumulative [StreamProperties]
* snapshot if at least one complete item was parsed, else null.
*/
fun feed(text: String): StreamProperties? {
buffer.append(text)
var lastEnd = 0
var changed = false
for (match in ITEM_REGEX.findAll(buffer)) {
processItem(match)
lastEnd = match.range.last + 1
changed = true
}
if (lastEnd > 0) buffer.delete(0, lastEnd)
return if (changed) snapshot() else null
}

/** Clears any partial item left over (e.g. after the writer closed). */
fun resetBuffer() = buffer.setLength(0)

fun snapshot(): StreamProperties {
val controllable = credentials() != null
val hasMetadata = title != null || artist != null || album != null ||
artDataBase64 != null || durationSec != null
val metadata = if (hasMetadata) {
StreamMetadata(
album = album,
artist = artist?.let { JsonArray(listOf(JsonPrimitive(it))) },
track = null,
title = title,
duration = durationSec,
artUrl = null,
artData = artDataBase64?.let { ArtData(data = it, extension = artExtension) },
)
} else {
null
}
return StreamProperties(
playbackStatus = when (playing) {
true -> "playing"
false -> "paused"
null -> "stopped"
},
loopStatus = "none",
shuffle = false,
volume = volumePercent,
mute = muted,
rate = 1.0f,
position = positionSec,
// Controls are advertised only once we hold DACP credentials, i.e.
// the sender is reachable for transport commands. Seek is unsupported.
canControl = controllable,
canGoNext = controllable,
canGoPrevious = controllable,
canPause = controllable,
canPlay = controllable,
canSeek = false,
metadata = metadata,
)
}

private fun processItem(match: MatchResult) {
val type = fourCc(match.groupValues[1])
val code = fourCc(match.groupValues[2])
val data = match.groups[4]?.value

when (type) {
"core" -> when (code) {
"minm" -> title = data?.let { decodeText(it) }
"asar" -> artist = data?.let { decodeText(it) }
"asal" -> album = data?.let { decodeText(it) }
"astm" -> data?.let { durationSec = decodeBigEndianInt(it) / 1000f }
}

"ssnc" -> when (code) {
"PICT" -> if (data != null && data.isNotBlank()) {
artDataBase64 = data
artExtension = detectImageExtension(data)
}

// Begin / resume -> playing; flush (pause) / end (stop) -> not.
"pbeg", "prsm" -> playing = true

"pfls", "pend" -> playing = false

"prgr" -> data?.let { parseProgress(decodeText(it)) }

"pvol" -> data?.let { parseVolume(decodeText(it)) }

// DACP remote-control credentials for driving the sender.
"acre" -> activeRemote = data?.let { decodeText(it) }

"daid" -> dacpId = data?.let { decodeText(it) }
}
}
}

private fun decodeText(base64: String): String =
String(decodeBase64(base64), Charsets.UTF_8).trim()

private fun decodeBigEndianInt(base64: String): Int {
val bytes = decodeBase64(base64)
var value = 0
for (b in bytes) value = (value shl 8) or (b.toInt() and 0xFF)
return value
}

// "start/current/end" in RTP frames (44100 Hz). Yields position & duration.
private fun parseProgress(text: String) {
val parts = text.split('/')
if (parts.size != 3) return
val start = parts[0].trim().toLongOrNull() ?: return
val current = parts[1].trim().toLongOrNull() ?: return
val end = parts[2].trim().toLongOrNull() ?: return
positionSec = ((current - start).coerceAtLeast(0)) / RTP_RATE
if (end > start) durationSec = (end - start) / RTP_RATE
}

// "airplay_volume,volume,lowest,highest"; airplay_volume is -30..0 dB, or
// <= -144 for mute. Mapped linearly to 0..100 for display only.
private fun parseVolume(text: String) {
val airplayVolume = text.split(',').firstOrNull()?.trim()?.toFloatOrNull() ?: return
if (airplayVolume <= MUTE_THRESHOLD_DB) {
muted = true
} else {
muted = false
volumePercent = (((airplayVolume + VOLUME_RANGE_DB) / VOLUME_RANGE_DB) * 100f)
.toInt().coerceIn(0, 100)
}
}

// Detect from the base64 prefix of the image's magic bytes:
// PNG 0x89 'P' 'N' 'G' -> "iVBOR"; JPEG 0xFF 0xD8 0xFF -> "/9j/".
private fun detectImageExtension(base64: String): String = when {
base64.startsWith("iVBOR") -> "png"
base64.startsWith("/9j/") -> "jpg"
else -> "jpg"
}

private fun fourCc(hex: String): String {
val v = hex.toLong(16)
return charArrayOf(
((v shr 24) and 0xFF).toInt().toChar(),
((v shr 16) and 0xFF).toInt().toChar(),
((v shr 8) and 0xFF).toInt().toChar(),
(v and 0xFF).toInt().toChar(),
).concatToString()
}

companion object {
private const val RTP_RATE = 44100f
private const val MUTE_THRESHOLD_DB = -144f
private const val VOLUME_RANGE_DB = 30f

private val ITEM_REGEX = Regex(
"<item><type>([0-9a-fA-F]+)</type><code>([0-9a-fA-F]+)</code>" +
"<length>(\\d+)</length>" +
"(?:\\n<data encoding=\"base64\">\\n(.*?)</data>)?</item>",
RegexOption.DOT_MATCHES_ALL,
)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
package tech.capullo.radio.airplay

import android.util.Base64
import android.util.Log
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.currentCoroutineContext
import kotlinx.coroutines.ensureActive
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.isActive
import tech.capullo.radio.data.RadioRepository
import tech.capullo.radio.snapcast.StreamProperties
import java.io.FileInputStream
import java.io.InputStreamReader
import javax.inject.Inject

/**
* Reads shairport-sync's metadata FIFO and exposes the parsed AirPlay
* [StreamProperties] (title/artist/album/cover-art + play state) as a
* [StateFlow] for the Airplay Snapcast control bridge to push to clients.
*
* Like [AirplayProcess]/SnapserverProcess this is a supervised run loop: opening
* the FIFO O_RDONLY blocks until shairport-sync opens the write end (same
* rendezvous as the audio pipe), and a returned [run] (writer closed -> EOF)
* lets the service's supervisor reopen it for the next session.
*/
class AirplayMetadataReader @Inject constructor(private val radioRepository: RadioRepository) {
private val parser = AirplayMetadataParser { Base64.decode(it, Base64.DEFAULT) }

private val _properties = MutableStateFlow(parser.snapshot())
val properties: StateFlow<StreamProperties> = _properties.asStateFlow()

private val _credentials = MutableStateFlow(parser.credentials())
val credentials: StateFlow<DacpCredentials?> = _credentials.asStateFlow()

suspend fun run() = coroutineScope {
val pipePath = radioRepository.getAirplayMetadataPipeFilepath()
if (pipePath == null) {
Log.e(TAG, "No AirPlay metadata PIPE; not reading shairport metadata")
return@coroutineScope
}

// Drop any partial item buffered from a previous (closed) session.
parser.resetBuffer()

// Blocks until shairport-sync opens the write end.
FileInputStream(pipePath).use { fis ->
val reader = InputStreamReader(fis, Charsets.ISO_8859_1)
val buf = CharArray(BUFFER_SIZE)
var read: Int
while (reader.read(buf).also { read = it } != -1) {
currentCoroutineContext().ensureActive()
if (!isActive) break
parser.feed(String(buf, 0, read))?.let {
_properties.value = it
_credentials.value = parser.credentials()
}
}
}
}

companion object {
private val TAG = AirplayMetadataReader::class.java.simpleName
private const val BUFFER_SIZE = 4096

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

While this is likely sufficient for metadata, if a shairport-sync metadata burst exceeds 4KB, reader.read(buf) will chunk it. Ensure the AirplayMetadataParser logic remains resilient to items being split across these 4KB boundaries (your current buffer implementation seems to handle this, but it is worth verifying with a stress test).

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The unit test handlesItemSplitAcrossFeeds() in AirplayMetadataParserTest explicitly verifies this: it splits a metadata item at an arbitrary midpoint across two feed() calls and asserts the title is still parsed correctly.

}
}
64 changes: 64 additions & 0 deletions app/src/main/java/tech/capullo/radio/airplay/AirplayProcess.kt
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
package tech.capullo.radio.airplay

import android.util.Log
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.ensureActive
import tech.capullo.radio.data.RadioRepository
import java.io.BufferedReader
import java.io.InputStreamReader
import javax.inject.Inject

/**
* Runs shairport-sync (packaged by lib-shairport-android as an exec-able
* libshairport.so) as a classic AirPlay receiver writing 44100:16:2 PCM to
* the AirPlay FIFO, where snapserver picks it up as a stream source.
*/
class AirplayProcess @Inject constructor(private val radioRepository: RadioRepository) {

private val nativeLibDir = radioRepository.getNativeLibDirPath()

suspend fun start() = coroutineScope {
val pipeFilepath = radioRepository.getAirplayPipeFilepath()
if (pipeFilepath == null) {
Log.e(TAG, "Could not create the AirPlay PIPE, not starting shairport-sync")
return@coroutineScope
}
// Idempotent like the audio FIFO; AirplayMetadataReader opens the same
// path for reading. Null (mkfifo failure) just omits the metadata block.
val metadataPipeFilepath = radioRepository.getAirplayMetadataPipeFilepath()
val confFile = radioRepository.getShairportConfPath(pipeFilepath, metadataPipeFilepath)

val pb = ProcessBuilder()
.command(
"$nativeLibDir/libshairport.so",
"-c",
confFile,
)
.redirectErrorStream(true)
// Android hides MAC addresses from apps; shairport-sync's Android port
// derives its AirPlay device ID from this variable instead.
pb.environment()["SPS_DEVICE_ID"] = radioRepository.getAirplayDeviceId()

val process = pb.start()
// Cleanup lives in finally so cancellation propagates to the
// supervisor (which distinguishes a clean stop from a crash); a normal
// native exit returns and lets the supervisor restart.
try {
val bufferedReader = BufferedReader(
InputStreamReader(process.inputStream),
)
var line: String?
while (bufferedReader.readLine().also { line = it } != null) {
ensureActive()
Log.d(TAG, "shairport-sync: ${line!!}")
}
} finally {
process.destroy()
process.waitFor()

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

waitFor() is a blocking call. If a process hangs in a zombie state, this could potentially block the coroutine indefinitely. Consider using withTimeout(TIMEOUT_MS) { withContext(Dispatchers.IO) { process.waitFor() } } to ensure the supervisor can eventually recover even if the native process refuses to exit gracefully.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The same pattern should be applied to SnapclientProcess and SnapserverProcess for consistency.

}
}

companion object {
private val TAG = AirplayProcess::class.java.simpleName
}
}
Loading
Loading