diff --git a/android/app/src/main/kotlin/com/zarz/spotiflac/MainActivitySafIo.kt b/android/app/src/main/kotlin/com/zarz/spotiflac/MainActivitySafIo.kt index de9aec80..ec73ff4c 100644 --- a/android/app/src/main/kotlin/com/zarz/spotiflac/MainActivitySafIo.kt +++ b/android/app/src/main/kotlin/com/zarz/spotiflac/MainActivitySafIo.kt @@ -122,19 +122,27 @@ internal fun MainActivity.extFromMimeType(mime: String?): String { } } -internal fun MainActivity.copyUriToTemp(uri: Uri, fallbackExt: String? = null): String? { +/** [displayName] skips a provider name query when the caller already listed it. */ +internal fun MainActivity.copyUriToTemp( + uri: Uri, + fallbackExt: String? = null, + displayName: String? = null, +): String? { var tempFile: File? = null var success = false try { - val mime = try { contentResolver.getType(uri) } catch (_: Exception) { null } val nameHint = ( - try { DocumentFile.fromSingleUri(this, uri)?.name } catch (_: Exception) { null } + displayName + ?: try { DocumentFile.fromSingleUri(this, uri)?.name } catch (_: Exception) { null } ?: uri.lastPathSegment ?: "" ).lowercase(Locale.ROOT) val extFromName = extFromFileName(nameHint) - val extFromMime = extFromMimeType(mime) + // The MIME type only matters when the name has no extension. + val extFromMime = if (extFromName.isNotBlank()) "" else extFromMimeType( + try { contentResolver.getType(uri) } catch (_: Exception) { null }, + ) val ext = if (extFromName.isNotBlank()) extFromName else if (extFromMime.isNotBlank()) extFromMime else (fallbackExt ?: "") val suffix = ext.ifBlank { ".tmp" } tempFile = coreBackend.createTemporaryMediaFile(this, "saf_", suffix) diff --git a/android/app/src/main/kotlin/com/zarz/spotiflac/MainActivitySafScan.kt b/android/app/src/main/kotlin/com/zarz/spotiflac/MainActivitySafScan.kt index c71d8176..0ef31ecf 100644 --- a/android/app/src/main/kotlin/com/zarz/spotiflac/MainActivitySafScan.kt +++ b/android/app/src/main/kotlin/com/zarz/spotiflac/MainActivitySafScan.kt @@ -40,8 +40,12 @@ import java.io.IOException import java.io.RandomAccessFile import java.security.MessageDigest import java.util.Locale +import java.util.concurrent.ExecutionException import java.util.concurrent.Executors import java.util.concurrent.Future +import java.util.concurrent.LinkedBlockingQueue +import java.util.concurrent.ThreadPoolExecutor +import java.util.concurrent.TimeUnit // SAF library-scan subsystem: tree walking, incremental diff, CUE resolution, // and scan-progress state shared with the progress stream in MainActivity. @@ -558,18 +562,124 @@ internal fun MainActivity.getSafChildFileLookup( cache: MutableMap>, ): Map { val dirKey = dir.uri.toString() - return cache.getOrPut(dirKey) { - buildMap { - for (child in listSafChildrenOrThrow(dir)) { - if (child.isDirectory) continue - val childName = child.name?.trim().orEmpty() - if (childName.isBlank()) continue - put(childName.lowercase(Locale.ROOT), child.doc) - } - } + return cache.getOrPut(dirKey) { safChildFileLookup(listSafChildrenOrThrow(dir)) } + } + +private fun safChildFileLookup(children: List): Map = + buildMap { + for (child in children) { + if (child.isDirectory) continue + val childName = child.name?.trim().orEmpty() + if (childName.isBlank()) continue + put(childName.lowercase(Locale.ROOT), child.doc) } } +/** Keeps a traversal listing for CUE sibling lookup so the directory is not + * queried again; on network providers every listing is a round trip. */ +private fun rememberCueDirectoryListing( + dir: DocumentFile, + children: List, + cache: MutableMap>, +) { + if (children.any { !it.isDirectory && it.name?.endsWith(".cue", ignoreCase = true) == true }) { + cache[dir.uri.toString()] = safChildFileLookup(children) + } +} + +// Provider round trips dominate on SD cards, USB drives and network shares. +private const val SAF_LIST_CONCURRENCY = 4 +private const val SAF_READ_WORKERS = 6 + +/** + * Prefetches the listings of directories already waiting in a breadth-first + * queue. The caller still dequeues, deduplicates and records errors serially, + * so traversal order and results match a one-directory-at-a-time walk. + * Idle daemon threads exit on their own, so abandoned walks leak nothing. + */ +private class SafTreeLister(private val context: Context) { + private val executor = ThreadPoolExecutor( + SAF_LIST_CONCURRENCY, + SAF_LIST_CONCURRENCY, + 1, + TimeUnit.SECONDS, + LinkedBlockingQueue(), + ) { runnable -> Thread(runnable, "saf-scan-list").apply { isDaemon = true } } + .apply { allowCoreThreadTimeOut(true) } + private val prefetched = mutableMapOf>>>() + + private fun submit(dir: DocumentFile) = executor.submit>> { + try { + Result.success(context.listSafChildrenOrThrow(dir)) + } catch (e: Exception) { + Result.failure(e) + } + } + + fun list( + dir: DocumentFile, + queue: Collection>, + visited: Set, + ): Result> { + val uri = dir.uri.toString() + val current = prefetched.remove(uri) ?: submit(dir) + for ((next, _) in queue) { + if (prefetched.size >= SAF_LIST_CONCURRENCY * 2) break + val nextUri = next.uri.toString() + if (nextUri == uri || nextUri in visited || nextUri in prefetched) continue + prefetched[nextUri] = submit(next) + } + return try { + current.get() + } catch (e: ExecutionException) { + throw e.cause ?: e + } + } +} + +/** + * Runs [task] on a bounded pool and passes each result to [consume] in input + * order. A sliding window keeps workers busy past one slow file. Returns + * false when cancelled. + */ +private fun runSafReadsInOrder( + items: List, + cancelled: () -> Boolean, + task: (T) -> R, + consume: (T, Result) -> Unit, +): Boolean { + if (items.isEmpty()) return true + val workers = minOf(SAF_READ_WORKERS, items.size) + val executor = Executors.newFixedThreadPool(workers) + try { + val window = workers * 4 + val inFlight = ArrayDeque>>(window) + var next = 0 + while (next < items.size || inFlight.isNotEmpty()) { + if (cancelled()) { + inFlight.forEach { it.second.cancel(true) } + return false + } + while (inFlight.size < window && next < items.size) { + val item = items[next++] + inFlight.addLast(item to executor.submit { task(item) }) + } + val (item, future) = inFlight.removeFirst() + val result = try { + Result.success(future.get()) + } catch (e: ExecutionException) { + Result.failure(e.cause ?: e) + } catch (e: Exception) { + Result.failure(e) + } + consume(item, result) + } + return true + } finally { + executor.shutdownNow() + } +} + internal fun MainActivity.resolveCueAudioSibling( parentDir: DocumentFile, cueName: String, @@ -779,6 +889,7 @@ internal fun MainActivity.scanSafTree( val queue: ArrayDeque> = ArrayDeque() queue.add(root to "") + val lister = SafTreeLister(this) while (queue.isNotEmpty()) { if (safScanCancel) { @@ -791,17 +902,18 @@ internal fun MainActivity.scanSafTree( continue } - val children = try { - listSafChildrenOrThrow(dir) - } catch (e: Exception) { + val listing = lister.list(dir, queue, visitedDirUris) + val children = listing.getOrNull() + if (children == null) { traversalErrors++ updateSafScanProgress { it.errorCount = traversalErrors } android.util.Log.w( "SpotiFLAC", - "SAF scan: failed listing directory $dirUri: ${e.message}", + "SAF scan: failed listing directory $dirUri: ${listing.exceptionOrNull()?.message}", ) continue } + rememberCueDirectoryListing(dir, children, safChildLookupCache) for (child in children) { if (safScanCancel) { @@ -936,7 +1048,7 @@ internal fun MainActivity.scanSafTree( var tempCuePath: String? = null var tempAudioPath: String? = null try { - tempCuePath = copyUriToTemp(cueDoc.uri, ".cue") + tempCuePath = copyUriToTemp(cueDoc.uri, ".cue", cueName) if (tempCuePath == null) { errors++ android.util.Log.w("SpotiFLAC", "SAF scan: failed to copy CUE ${cueDoc.uri}") @@ -1060,11 +1172,12 @@ internal fun MainActivity.scanSafTree( pendingAudio.add(audio) } - // Bound parallelism to limit SAF full-copy memory and I/O. - val workerCount = 2 - val executor = Executors.newFixedThreadPool(workerCount) - fun submitAudio(audio: SafAudioEntry): Future = - executor.submit { + // Descriptor reads stream from the provider without copying, so the + // pool is bounded by provider round trips rather than copy memory. + val completed = runSafReadsInOrder( + pendingAudio, + cancelled = { safScanCancel }, + task = { audio -> val doc = audio.doc val stableUri = doc.uri.toString() val name = audio.name @@ -1087,62 +1200,43 @@ internal fun MainActivity.scanSafTree( null } SafAudioScanOutcome(stableUri, name, lastModified, metadata) - } - try { - // A sliding window keeps both workers busy past one slow file, - // unlike fixed batches that waited for their slowest entry. At most - // workerCount * 2 reads are queued or running, and rows are still - // consumed in input order. - val window = workerCount * 2 - val inFlight = ArrayDeque>(window) - var nextAudio = 0 - while (nextAudio < pendingAudio.size || inFlight.isNotEmpty()) { - if (safScanCancel) { - inFlight.forEach { it.cancel(true) } - ndjsonWriter?.close() - spill?.abandon() - return cancelledResult() - } - while (inFlight.size < window && nextAudio < pendingAudio.size) { - inFlight.addLast(submitAudio(pendingAudio[nextAudio++])) - } - - val outcome = try { - inFlight.removeFirst().get() - } catch (e: Exception) { + }, + ) { _, result -> + val outcome = result.getOrNull() + if (outcome == null) { + errors++ + } else { + updateSafScanProgress { it.currentFile = outcome.name } + val metadataObj = outcome.metadata + if (metadataObj == null) { errors++ - null - } - if (outcome != null) { - updateSafScanProgress { it.currentFile = outcome.name } - val metadataObj = outcome.metadata - if (metadataObj == null) { + } else { + try { + metadataObj.put("id", buildStableLibraryId(outcome.uri)) + metadataObj.put("filePath", outcome.uri) + metadataObj.put("fileModTime", outcome.lastModified) + putResult(metadataObj) + // Flush before recording the checkpoint to avoid losing the row. + ndjsonWriter?.flush() + recordCheckpoint(outcome.uri, outcome.lastModified) + } catch (_: Exception) { errors++ - } else { - try { - metadataObj.put("id", buildStableLibraryId(outcome.uri)) - metadataObj.put("filePath", outcome.uri) - metadataObj.put("fileModTime", outcome.lastModified) - putResult(metadataObj) - // Flush before recording the checkpoint to avoid losing the row. - ndjsonWriter?.flush() - recordCheckpoint(outcome.uri, outcome.lastModified) - } catch (_: Exception) { - errors++ - } } } - - scanned++ - val pct = scanned.toDouble() / totalItems.toDouble() * 100.0 - updateSafScanProgress { - it.scannedFiles = scanned - it.errorCount = errors - it.progressPct = pct - } } - } finally { - executor.shutdownNow() + + scanned++ + val pct = scanned.toDouble() / totalItems.toDouble() * 100.0 + updateSafScanProgress { + it.scannedFiles = scanned + it.errorCount = errors + it.progressPct = pct + } + } + if (!completed) { + ndjsonWriter?.close() + spill?.abandon() + return cancelledResult() } updateSafScanProgress { @@ -1201,9 +1295,18 @@ internal fun MainActivity.scanSafTreeIncremental( } val supportedAudioExt = libraryScanAudioExtensions - val audioFiles = mutableListOf>() - val cueFilesToScan = mutableListOf>() - val unchangedCueFiles = mutableListOf>() + // Names and timestamps come from the directory listing; querying each + // document again costs a provider round trip per changed file. + data class ChangedAudio(val doc: DocumentFile, val name: String, val lastModified: Long) + data class CueEntry( + val doc: DocumentFile, + val parentDir: DocumentFile, + val name: String, + val lastModified: Long, + ) + val audioFiles = mutableListOf() + val cueFilesToScan = mutableListOf() + val unchangedCueFiles = mutableListOf() val currentUris = mutableSetOf() val visitedDirUris = mutableSetOf() val safChildLookupCache = mutableMapOf>() @@ -1220,6 +1323,7 @@ internal fun MainActivity.scanSafTreeIncremental( val queue: ArrayDeque> = ArrayDeque() queue.add(root to "") + val lister = SafTreeLister(this) while (queue.isNotEmpty()) { if (safScanCancel) { @@ -1239,17 +1343,18 @@ internal fun MainActivity.scanSafTreeIncremental( continue } - val children = try { - listSafChildrenOrThrow(dir) - } catch (e: Exception) { + val listing = lister.list(dir, queue, visitedDirUris) + val children = listing.getOrNull() + if (children == null) { traversalErrors++ updateSafScanProgress { it.errorCount = traversalErrors } android.util.Log.w( "SpotiFLAC", - "SAF incremental scan: failed listing directory $dirUri: ${e.message}", + "SAF incremental scan: failed listing directory $dirUri: ${listing.exceptionOrNull()?.message}", ) continue } + rememberCueDirectoryListing(dir, children, safChildLookupCache) for (child in children) { if (safScanCancel) { @@ -1286,19 +1391,19 @@ internal fun MainActivity.scanSafTreeIncremental( val existingModified = virtualPaths?.firstOrNull()?.let { existingFiles[it] } if (existingModified != null && existingModified == lastModified) { - unchangedCueFiles.add(child.doc to dir) + unchangedCueFiles.add(CueEntry(child.doc, dir, name, lastModified)) for (vp in virtualPaths) { currentUris.add(vp) } } else { - cueFilesToScan.add(Triple(child.doc, dir, lastModified)) + cueFilesToScan.add(CueEntry(child.doc, dir, name, lastModified)) } } else if (ext.isNotBlank() && supportedAudioExt.contains(".$ext")) { val existingModified = existingFiles[uriStr] val lastModified = child.lastModified if (existingModified == null || existingModified != lastModified) { - audioFiles.add(Triple(child.doc, path, lastModified)) + audioFiles.add(ChangedAudio(child.doc, name, lastModified)) } } } @@ -1355,7 +1460,7 @@ internal fun MainActivity.scanSafTreeIncremental( val cueReferencedAudioUris = mutableSetOf() - for ((cueDoc, parentDir, cueLastModified) in cueFilesToScan) { + for ((cueDoc, parentDir, cueName, cueLastModified) in cueFilesToScan) { if (safScanCancel) { updateSafScanProgress { it.isComplete = true } spill.abandon() @@ -1368,13 +1473,12 @@ internal fun MainActivity.scanSafTreeIncremental( return result.toString() } - val cueName = try { cueDoc.name ?: "" } catch (_: Exception) { "" } updateSafScanProgress { it.currentFile = cueName } var tempCuePath: String? = null var tempAudioPath: String? = null try { - tempCuePath = copyUriToTemp(cueDoc.uri, ".cue") + tempCuePath = copyUriToTemp(cueDoc.uri, ".cue", cueName) if (tempCuePath == null) { errors++ android.util.Log.w("SpotiFLAC", "SAF incremental scan: failed to copy CUE ${cueDoc.uri}") @@ -1465,13 +1569,12 @@ internal fun MainActivity.scanSafTreeIncremental( } } - for ((cueDoc, parentDir) in unchangedCueFiles) { + for ((cueDoc, parentDir, cueName) in unchangedCueFiles) { var tempCue: String? = null try { - tempCue = copyUriToTemp(cueDoc.uri, ".cue") + tempCue = copyUriToTemp(cueDoc.uri, ".cue", cueName) if (tempCue != null) { val audioFileName = extractCueAudioFileName(tempCue) - val cueName = try { cueDoc.name ?: "" } catch (_: Exception) { "" } val audioDoc = resolveCueAudioSibling( parentDir = parentDir, cueName = cueName, @@ -1489,65 +1592,19 @@ internal fun MainActivity.scanSafTreeIncremental( } } - for ((doc, _, lastModified) in audioFiles) { - if (safScanCancel) { - updateSafScanProgress { it.isComplete = true } - spill.abandon() - val result = JSONObject() - result.put("files", JSONArray()) - result.put("removedUris", JSONArray()) - result.put("skippedCount", skippedCount) - result.put("totalFiles", totalFiles) - result.put("cancelled", true) - return result.toString() - } + fun cancelledIncrementalResult(): String { + updateSafScanProgress { it.isComplete = true } + spill.abandon() + val result = JSONObject() + result.put("files", JSONArray()) + result.put("removedUris", JSONArray()) + result.put("skippedCount", skippedCount) + result.put("totalFiles", totalFiles) + result.put("cancelled", true) + return result.toString() + } - if (cueReferencedAudioUris.contains(doc.uri.toString())) { - scanned++ - val processed = skippedCount + scanned - val pct = if (totalFiles > 0) { - processed.toDouble() / totalFiles.toDouble() * 100.0 - } else { - 100.0 - } - updateSafScanProgress { - it.scannedFiles = processed - it.progressPct = pct - } - continue - } - - val name = try { doc.name ?: "" } catch (_: Exception) { "" } - updateSafScanProgress { - it.currentFile = name - } - - val ext = name.substringAfterLast('.', "").lowercase(Locale.ROOT) - val fallbackExt = if (ext.isNotBlank()) ".${ext}" else null - val safeLastModified = try { doc.lastModified() } catch (_: Exception) { lastModified } - val stableUri = doc.uri.toString() - val coverCacheKey = buildLibraryCoverCacheKey(stableUri, safeLastModified) - val metadataObj = readAudioMetadataFromUri( - doc.uri, - name, - fallbackExt, - coverCacheKey, - ) - if (metadataObj == null) { - errors++ - } else { - try { - metadataObj.put("id", buildStableLibraryId(stableUri)) - metadataObj.put("filePath", stableUri) - metadataObj.put("fileModTime", safeLastModified) - metadataObj.put("lastModified", safeLastModified) - putFile(metadataObj) - } catch (_: Exception) { - errors++ - } - } - - scanned++ + fun reportProcessed() { val processed = skippedCount + scanned val pct = if (totalFiles > 0) { processed.toDouble() / totalFiles.toDouble() * 100.0 @@ -1561,6 +1618,53 @@ internal fun MainActivity.scanSafTreeIncremental( } } + val pendingAudio = mutableListOf() + for (audio in audioFiles) { + if (safScanCancel) return cancelledIncrementalResult() + if (cueReferencedAudioUris.contains(audio.doc.uri.toString())) { + scanned++ + reportProcessed() + } else { + pendingAudio.add(audio) + } + } + + val completed = runSafReadsInOrder( + pendingAudio, + cancelled = { safScanCancel }, + task = { audio -> + val ext = audio.name.substringAfterLast('.', "").lowercase(Locale.ROOT) + val fallbackExt = if (ext.isNotBlank()) ".${ext}" else null + readAudioMetadataFromUri( + audio.doc.uri, + audio.name, + fallbackExt, + buildLibraryCoverCacheKey(audio.doc.uri.toString(), audio.lastModified), + ) + }, + ) { audio, result -> + updateSafScanProgress { it.currentFile = audio.name } + // A failed read aborts the incremental scan, as it did when serial. + val metadataObj = result.getOrThrow() + if (metadataObj == null) { + errors++ + } else { + try { + val stableUri = audio.doc.uri.toString() + metadataObj.put("id", buildStableLibraryId(stableUri)) + metadataObj.put("filePath", stableUri) + metadataObj.put("fileModTime", audio.lastModified) + metadataObj.put("lastModified", audio.lastModified) + putFile(metadataObj) + } catch (_: Exception) { + errors++ + } + } + scanned++ + reportProcessed() + } + if (!completed) return cancelledIncrementalResult() + val finalRemovedUris = existingFiles.keys.filter { !currentUris.contains(it) } updateSafScanProgress {