perf(library): reduce scan allocations and keep SAF workers busy

This commit is contained in:
zarzet committed 2026-09-30 23:39:22 +07:00
1 parent 1d3f7a0d44
commit 9051f37cf5
2 files changed
+223 -96

No files matched your search

@@ -41,6 +41,7 @@ import java.io.RandomAccessFile
import java.security.MessageDigest
import java.util.Locale
import java.util.concurrent.Executors
import java.util.concurrent.Future
// SAF library-scan subsystem: tree walking, incremental diff, CUE resolution,
// and scan-progress state shared with the progress stream in MainActivity.
@@ -1062,82 +1063,82 @@ internal fun MainActivity.scanSafTree(
// Bound parallelism to limit SAF full-copy memory and I/O.
val workerCount = 2
val executor = Executors.newFixedThreadPool(workerCount)
fun submitAudio(audio: SafAudioEntry): Future<SafAudioScanOutcome> =
executor.submit<SafAudioScanOutcome> {
val doc = audio.doc
val stableUri = doc.uri.toString()
val name = audio.name
val lastModified = audio.lastModified
val ext = name.substringAfterLast('.', "").lowercase(Locale.ROOT)
val fallbackExt = if (ext.isNotBlank()) ".${ext}" else null
val coverCacheKey = buildLibraryCoverCacheKey(stableUri, lastModified)
val metadata = try {
readAudioMetadataFromUri(
doc.uri,
name,
fallbackExt,
coverCacheKey,
)
} catch (e: Exception) {
android.util.Log.w(
"SpotiFLAC",
"SAF scan: metadata read failed for $stableUri: ${e.message}",
)
null
}
SafAudioScanOutcome(stableUri, name, lastModified, metadata)
}
try {
for (batch in pendingAudio.chunked(workerCount * 2)) {
// 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<Future<SafAudioScanOutcome>>(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 futures = batch.map { audio ->
executor.submit<SafAudioScanOutcome> {
val doc = audio.doc
val stableUri = doc.uri.toString()
val name = audio.name
val lastModified = audio.lastModified
val ext = name.substringAfterLast('.', "").lowercase(Locale.ROOT)
val fallbackExt = if (ext.isNotBlank()) ".${ext}" else null
val coverCacheKey = buildLibraryCoverCacheKey(stableUri, lastModified)
val metadata = try {
readAudioMetadataFromUri(
doc.uri,
name,
fallbackExt,
coverCacheKey,
)
} catch (e: Exception) {
android.util.Log.w(
"SpotiFLAC",
"SAF scan: metadata read failed for $stableUri: ${e.message}",
)
null
val outcome = try {
inFlight.removeFirst().get()
} catch (e: Exception) {
errors++
null
}
if (outcome != null) {
updateSafScanProgress { it.currentFile = outcome.name }
val metadataObj = outcome.metadata
if (metadataObj == null) {
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++
}
SafAudioScanOutcome(stableUri, name, lastModified, metadata)
}
}
for (future in futures) {
if (safScanCancel) {
futures.forEach { it.cancel(true) }
ndjsonWriter?.close()
spill?.abandon()
return cancelledResult()
}
val outcome = try {
future.get()
} catch (e: Exception) {
errors++
null
}
if (outcome != null) {
updateSafScanProgress { it.currentFile = outcome.name }
val metadataObj = outcome.metadata
if (metadataObj == null) {
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
}
scanned++
val pct = scanned.toDouble() / totalItems.toDouble() * 100.0
updateSafScanProgress {
it.scannedFiles = scanned
it.errorCount = errors
it.progressPct = pct
}
}
} finally {