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
Original file line number Diff line number Diff line change
Expand Up @@ -449,15 +449,19 @@ class StreamFileUploadProcessor(
Log.w("StreamFileUploadProcessor", "⚠️ No saved server can serve ${item.title} — will retry")
return false
}
// Header auth on top of the query token, like playback and the artwork backfill: newer ABS
// versions reject query-string tokens (401) and only accept the Authorization header.
// Header auth, like playback and the artwork backfill: the Jellyfin URL carries no token, and
// newer ABS versions reject query-string tokens (401). Custom headers are sanitized like the
// download's: a persisted illegal name/value throws on addHeader and would wedge the pipe.
val headers = ExternalServiceUtils.serviceTypeFor(resource.providerName)
?.let { ExternalServiceUtils.playbackHeaders(it, server.token, server.customHeaders) }
?.let { ExternalServiceUtils.playbackHeaders(it, server.token, ExternalServiceUtils.sanitizeCustomHeaders(server.customHeaders)) }

return try {
val getRequest = okhttp3.Request.Builder().url(sourceUrl)
headers?.forEach { (k, v) -> getRequest.addHeader(k, v) }
client.newCall(getRequest.build()).execute().use { response ->
// Headers ride only the hops that stay on the media server (a redirect elsewhere must not get
// them); newBuilder() shares the pipe client's connection pool.
val sourceClient = headers?.let {
client.newBuilder().addNetworkInterceptor(ExternalServiceUtils.originPinnedHeaders(sourceUrl, it)).build()
} ?: client
sourceClient.newCall(okhttp3.Request.Builder().url(sourceUrl).build()).execute().use { response ->
val body = response.body
if (!response.isSuccessful || body == null) {
Log.e("StreamFileUploadProcessor", "❌ Source GET failed (${response.code}) for ${item.title}")
Expand Down Expand Up @@ -553,7 +557,19 @@ class StreamFileUploadProcessor(
}
}

class DownloadFileProcessor(private val context: Context) : TaskProcessor {
class DownloadFileProcessor(
private val context: Context,
// Through the repository, not the DAO: stored credentials are encrypted at rest, and media-server
// downloads authenticate with this token. Overridable so tests can swap the Keystore cipher.
private val serverRepository: ExternalServerRepository =
ExternalServerRepository(AppDatabase.getDatabase(context).externalServerDao()),
) : TaskProcessor {
companion object {
// One base client for every download (and retry); per-server variants derive via newBuilder(),
// which shares this client's connection pool and dispatcher threads.
private val baseHttpClient by lazy { okhttp3.OkHttpClient() }
}

override suspend fun process(task: SyncTaskEntity): Boolean {
val gson = Gson()
val payloadType = object : TypeToken<Map<String, Any?>>() {}.type
Expand Down Expand Up @@ -585,61 +601,66 @@ class DownloadFileProcessor(private val context: Context) : TaskProcessor {
// Ensure parent directories exist
destFile.parentFile?.mkdirs()

val client = okhttp3.OkHttpClient()
val request = okhttp3.Request.Builder().url(remoteURL).build()

return try {
val response = client.newCall(request).execute()
if (!response.isSuccessful) {
Log.e("DownloadFileProcessor", "❌ Download failed: ${response.code}")
return false
}

val body = response.body ?: return false
val contentLength = body.contentLength()
// Refuse up front when the file can't fit with headroom to spare: a download that fills the
// disk takes the database down with it. The engine holds downloads until storage recovers.
if (contentLength > 0 && !StorageMonitor.hasRoomFor(context, contentLength)) {
StorageMonitor.noteTransferDoesNotFit(context, contentLength)
Log.w("DownloadFileProcessor", "⛔ Not enough storage for $relativePath ($contentLength bytes)")
return false
}
var bytesRead = 0L
var cancelled = false
// Media-server headers ride only the hops that stay on that server: OkHttp would carry custom
// headers (often Cloudflare Access secrets) across a redirect to another host.
val client = mediaServerHeaders(taskId, remoteURL)?.let {
baseHttpClient.newBuilder().addNetworkInterceptor(ExternalServiceUtils.originPinnedHeaders(remoteURL, it)).build()
} ?: baseHttpClient
// `use` closes the response on every path: the early returns below (error status, no room) would
// otherwise leak the connection on each retry.
client.newCall(okhttp3.Request.Builder().url(remoteURL).build()).execute().use { response ->
if (!response.isSuccessful) {
Log.e("DownloadFileProcessor", "❌ Download failed: ${response.code}")
return false
}

body.byteStream().use { input: java.io.InputStream ->
FileOutputStream(destFile).use { output: FileOutputStream ->
val buffer = ByteArray(8 * 1024)
var read: Int
while (input.read(buffer).also { read = it } != -1) {
// Cooperative cancellation: abort mid-stream if the user cancelled this download.
if (SyncStatusManager.isCancelRequested(taskId)) {
cancelled = true
break
}
output.write(buffer, 0, read)
bytesRead += read
if (contentLength > 0) {
val progress = bytesRead.toDouble() / contentLength
SyncStatusManager.updateTaskProgress(taskId, progress)
val body = response.body ?: return false
val contentLength = body.contentLength()
// Refuse up front when the file can't fit with headroom to spare: a download that fills the
// disk takes the database down with it. The engine holds downloads until storage recovers.
if (contentLength > 0 && !StorageMonitor.hasRoomFor(context, contentLength)) {
StorageMonitor.noteTransferDoesNotFit(context, contentLength)
Log.w("DownloadFileProcessor", "⛔ Not enough storage for $relativePath ($contentLength bytes)")
return false
}
var bytesRead = 0L
var cancelled = false

body.byteStream().use { input: java.io.InputStream ->
FileOutputStream(destFile).use { output: FileOutputStream ->
val buffer = ByteArray(8 * 1024)
var read: Int
while (input.read(buffer).also { read = it } != -1) {
// Cooperative cancellation: abort mid-stream if the user cancelled this download.
if (SyncStatusManager.isCancelRequested(taskId)) {
cancelled = true
break
}
output.write(buffer, 0, read)
bytesRead += read
if (contentLength > 0) {
val progress = bytesRead.toDouble() / contentLength
SyncStatusManager.updateTaskProgress(taskId, progress)
}
}
output.flush()
}
output.flush()
}
}

SyncStatusManager.clearTaskProgress(taskId)
if (cancelled) {
Log.d("DownloadFileProcessor", "🚫 Download cancelled: $relativePath")
if (destFile.exists()) destFile.delete()
// Leave the cancel flag SET on purpose: TaskConcurrencyManager reads it on this false
// return to make the task terminal (delete, no retry) and then clears it. Clearing here
// would let the failure path re-queue the task and silently re-download it to completion.
return false
SyncStatusManager.clearTaskProgress(taskId)
if (cancelled) {
Log.d("DownloadFileProcessor", "🚫 Download cancelled: $relativePath")
if (destFile.exists()) destFile.delete()
// Leave the cancel flag SET on purpose: TaskConcurrencyManager reads it on this false
// return to make the task terminal (delete, no retry) and then clears it. Clearing here
// would let the failure path re-queue the task and silently re-download it to completion.
return false
}
SyncStatusManager.clearCancel(taskId)
Log.d("DownloadFileProcessor", "✅ Download complete: $relativePath")
true
}
SyncStatusManager.clearCancel(taskId)
Log.d("DownloadFileProcessor", "✅ Download complete: $relativePath")
true
} catch (e: Exception) {
StorageMonitor.reportFailure(context, e) // ENOSPC mid-write: the storage state holds further downloads
Log.e("DownloadFileProcessor", "💥 Exception during download: ${e.message}", e)
Expand All @@ -653,6 +674,15 @@ class DownloadFileProcessor(private val context: Context) : TaskProcessor {
}
}

// Resolved per run, not stored in the payload: tokens stay out of the task table, and a re-auth's
// fresh token applies to an already-queued download. The resource pick mirrors externalStreamUrlFor.
private suspend fun mediaServerHeaders(uuid: String, url: String): Map<String, String>? {
val resource = AppDatabase.getDatabase(context).libraryDao().getExternalResourcesForBookSync(uuid)
.find { it.syncStatus == ExternalResourceEntity.STATUS_STREAM || it.syncStatus == ExternalResourceEntity.STATUS_DOWNLOADED }
?: return null
return ExternalServiceUtils.downloadHeadersFor(serverRepository, resource, url)
}

override fun canHandle(jobType: String): Boolean {
return jobType == SyncTaskFactory.JOB_DOWNLOAD_FILE
}
Expand Down Expand Up @@ -996,7 +1026,10 @@ class SetExternalResourceToDownloadProcessor : TaskProcessor {
}

class ExternalUpdateProcessor(
private val context: Context
private val context: Context,
// Through the repository, not the DAO (see process()). Overridable so tests can swap the Keystore cipher.
private val serverRepository: ExternalServerRepository =
ExternalServerRepository(AppDatabase.getDatabase(context).externalServerDao()),
) : TaskProcessor {
private val gson = Gson()

Expand All @@ -1008,7 +1041,9 @@ class ExternalUpdateProcessor(

private fun buildApiClient(sanitizedUrl: String, customHeaders: Map<String, String>?): retrofit2.Retrofit {
val okHttpClientBuilder = baseHttpClient.newBuilder()
customHeaders?.forEach { (key, value) ->
// Sanitized like JellyfinService's client: a custom `Authorization` entry would replace the
// provider's own auth header, and an illegal name/value throws at request time.
ExternalServiceUtils.sanitizeCustomHeaders(customHeaders)?.forEach { (key, value) ->
okHttpClientBuilder.addInterceptor { chain ->
val request = chain.request().newBuilder().header(key, value).build()
chain.proceed(request)
Expand All @@ -1034,34 +1069,6 @@ class ExternalUpdateProcessor(
return permanent
}

private fun getDeviceId(): String {
return try {
if (!com.tortugapower.audiobookplayer.core.CoreContext.isInitialized()) return "BookPlayerAndroidID"
val appCtx = com.tortugapower.audiobookplayer.core.CoreContext.appContext
val prefs = appCtx.getSharedPreferences("jellyfin_prefs", Context.MODE_PRIVATE)
var id = prefs.getString("device_id", null)
if (id == null) {
id = java.util.UUID.randomUUID().toString()
prefs.edit().putString("device_id", id).apply()
}
id
} catch (e: Exception) {
"BookPlayerAndroidID"
}
}

private fun getJellyfinAuthHeader(token: String? = null): String {
val device = "Android"
val deviceId = getDeviceId()
val client = "BookPlayer"
val version = "1.0.0"
var header = "MediaBrowser Client=\"$client\", Device=\"$device\", DeviceId=\"$deviceId\", Version=\"$version\""
if (token != null) {
header += ", Token=\"$token\""
}
return header
}

override suspend fun process(task: SyncTaskEntity): Boolean {
val payloadType = object : TypeToken<Map<String, Any?>>() {}.type
val payload: Map<String, Any?> = gson.fromJson(task.payload, payloadType)
Expand All @@ -1075,13 +1082,10 @@ class ExternalUpdateProcessor(
val percentCompleted = (payload["percentCompleted"] as? Double) ?: 0.0
val isFinished = (payload["isFinished"] as? Boolean) ?: false

val db = AppDatabase.getDatabase(context)

// Resolve through THE shared resolver (stable-id contract + decrypted credentials): the
// old inline rowid lookup read the DAO directly, so the token below was ciphertext and
// the provider rejected it with 401; it also stopped matching once hostIds became
// GUIDs/URL keys, silently discarding every progress push.
val serverRepository = com.tortugapower.audiobookplayer.repository.ExternalServerRepository(db.externalServerDao())
val server = ExternalServiceUtils.serverForResource(
serverRepository,
com.tortugapower.audiobookplayer.database.entities.ExternalResourceEntity(
Expand Down Expand Up @@ -1121,7 +1125,7 @@ class ExternalUpdateProcessor(
val api = buildApiClient(sanitizedUrl, customHeaders)
.create(com.tortugapower.audiobookplayer.network.services.JellyfinApi::class.java)

val authHeader = getJellyfinAuthHeader(token)
val authHeader = com.tortugapower.audiobookplayer.network.services.JellyfinService.getAuthHeader(token)
val response = api.updateUserData(authHeader, providerId, requestBody)
handleResponse(providerName, response)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@ import com.tortugapower.audiobookplayer.database.entities.ExternalServerEntity
import com.tortugapower.audiobookplayer.database.entities.ExternalServiceType
import com.tortugapower.audiobookplayer.repository.ExternalServerRepository
import kotlinx.coroutines.flow.first
import okhttp3.HttpUrl.Companion.toHttpUrlOrNull
import okhttp3.Interceptor

object ExternalServiceUtils {
fun sanitizeUrl(url: String): String {
Expand Down Expand Up @@ -136,16 +138,56 @@ object ExternalServiceUtils {
}

/**
* The provider's direct-download URL for [resource] on [server] (query-token auth, so it needs no
* extra headers), or null for an unknown provider. Pure counterpart of the URL rebuild in
* `resolveStreamingUrl`, also used to GET the source file for the stream-to-cloud pipe.
* The provider's direct-download URL for [resource] on [server], or null for an unknown provider.
* Pure counterpart of the URL rebuild in `resolveStreamingUrl`, also used to GET the source file for
* the stream-to-cloud pipe. The Jellyfin URL carries no token — Jellyfin 12 ignores `api_key`, and a
* URL token leaks into logs and the task table — so every request for it needs the provider's header
* auth: playback via PlaybackManager's host registry, the pipe and downloads via [downloadHeadersFor].
* ABS keeps its `token` query param; its consumers send the Bearer header on top of it.
*/
fun downloadUrlFor(server: ExternalServerEntity, resource: ExternalResourceEntity): String? {
val path = when (serviceTypeFor(resource.providerName)) {
ExternalServiceType.JELLYFIN -> "Items/${resource.providerId}/Download?api_key=${server.token ?: ""}"
ExternalServiceType.JELLYFIN -> "Items/${resource.providerId}/Download"
ExternalServiceType.AUDIOBOOKSHELF -> "api/items/${resource.providerId}/download?token=${server.token ?: ""}"
null -> return null
}
return "${sanitizeUrl(server.url)}$path"
}

/**
* The headers a download of [url] must carry when it comes from the saved server behind [resource]:
* the provider's Authorization header plus the user's custom headers, like playback and the pipe.
* The query token alone isn't enough — Jellyfin 12 rejects it (401), as do newer ABS versions.
* Null when [url] is anywhere else: a BookPlayer-cloud presigned URL must go out bare, since S3
* rejects a request that carries a second auth mechanism.
*/
suspend fun downloadHeadersFor(
servers: ExternalServerRepository,
resource: ExternalResourceEntity,
url: String,
): Map<String, String>? {
val server = serverForResource(servers, resource) ?: return null
if (!url.startsWith(sanitizeUrl(server.url))) return null
val type = serviceTypeFor(resource.providerName) ?: return null
return playbackHeaders(type, server.token, sanitizeCustomHeaders(server.customHeaders))
}

/**
* A network interceptor that adds [headers] to each hop of a request only while it stays on [url]'s
* origin (scheme, host and port — the rule OkHttp applies to `Authorization` on redirects). OkHttp
* keeps every other header across a cross-host redirect, and custom headers are often Cloudflare
* Access secrets; playback pins its headers to the server's host the same way.
*/
fun originPinnedHeaders(url: String, headers: Map<String, String>): Interceptor {
val origin = url.toHttpUrlOrNull()
return Interceptor { chain ->
val request = chain.request()
val sameOrigin = origin != null && request.url.scheme == origin.scheme &&
request.url.host == origin.host && request.url.port == origin.port
if (!sameOrigin) return@Interceptor chain.proceed(request)
val pinned = request.newBuilder()
headers.forEach { (name, value) -> pinned.header(name, value) }
chain.proceed(pinned.build())
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -707,7 +707,7 @@ object PlaybackManager {
* Best-effort audio file extension for picking the chapter parser, from most to least reliable:
* the item's `relativePath`, then its `originalFileName` (set for external items whose relativePath
* is null, e.g. AudiobookShelf), then the remote URL's last path segment with any query/fragment
* stripped. A streaming URL like `Items/<id>/Download?api_key=...` yields no extension → we fall
* stripped. A streaming URL like `Items/<id>/Download` yields no extension → we fall
* through rather than mis-detecting. Pure (no Android APIs) so it's unit-tested. Lowercased, no dot.
*/
internal fun audioExtensionFor(item: LibraryItemEntity, url: String): String {
Expand Down
Loading
Loading