feat(backend): hold library scans in place while paused

Path and SAF scans wait at their next file checkpoint while paused and keep their position in memory. Cancel, backend shutdown and released request leases still end a paused scan.
This commit is contained in:
zarzet committed 2026-10-01 21:00:42 +07:00
1 parent 72b689b1f8
commit fe36f08918
8 files changed
+206 -12

No files matched your search

@@ -203,6 +203,8 @@ internal interface CoreBackend {
fun scanLibraryFolderIncrementalFromSnapshot(folder: String, snapshot: String): String
fun getLibraryScanProgress(): String
fun cancelLibraryScan()
fun pauseLibraryScan()
fun resumeLibraryScan()
fun parseCueSheet(path: String, audioDirectory: String): String
fun parseCueSheetWithResolvedAudio(path: String, audioPath: String): String
fun scanCueForLibrary(path: String, audioDirectory: String, virtualPrefix: String, modTime: Long, cacheKey: String): String
@@ -101,6 +101,8 @@ class MainActivity: FlutterFragmentActivity() {
"scanSafTreeIncrementalFromSnapshot",
"getLibraryScanProgress",
"cancelLibraryScan",
"pauseLibraryScan",
"resumeLibraryScan",
"parseCueSheet",
"pickSafTree",
"safExists",
@@ -167,6 +169,7 @@ class MainActivity: FlutterFragmentActivity() {
private val playbackLeaseLock = Any()
private val playbackLeases = LinkedHashMap<String, ParcelFileDescriptor>()
@Volatile internal var safScanCancel = false
@Volatile internal var safScanPaused = false
@Volatile internal var safScanActive = false
private val safTreeLauncher = registerForActivityResult(
ActivityResultContracts.StartActivityForResult()
@@ -2018,10 +2021,25 @@ class MainActivity: FlutterFragmentActivity() {
"cancelLibraryScan" -> {
withContext(Dispatchers.IO) {
safScanCancel = true
safScanPaused = false
coreBackend.cancelLibraryScan()
}
result.success(null)
}
"pauseLibraryScan" -> {
withContext(Dispatchers.IO) {
safScanPaused = true
coreBackend.pauseLibraryScan()
}
result.success(null)
}
"resumeLibraryScan" -> {
withContext(Dispatchers.IO) {
safScanPaused = false
coreBackend.resumeLibraryScan()
}
result.success(null)
}
"readAudioMetadata" -> {
val filePath = call.argument<String>("file_path") ?: ""
val response = withContext(Dispatchers.IO) {
@@ -590,6 +590,18 @@ private fun rememberCueDirectoryListing(
// 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
private const val SAF_SCAN_PAUSE_POLL_MS = 100L
/**
* Holds a SAF scan at its checkpoint while paused, keeping its position in
* memory. Returns true when the scan should stop because it was cancelled.
*/
internal fun MainActivity.safScanStopRequested(): Boolean {
while (safScanPaused && !safScanCancel) {
Thread.sleep(SAF_SCAN_PAUSE_POLL_MS)
}
return safScanCancel
}
/**
* Prefetches the listings of directories already waiting in a breadth-first
@@ -892,7 +904,7 @@ internal fun MainActivity.scanSafTree(
val lister = SafTreeLister(this)
while (queue.isNotEmpty()) {
if (safScanCancel) {
if (safScanStopRequested()) {
return cancelledResult()
}
@@ -916,7 +928,7 @@ internal fun MainActivity.scanSafTree(
rememberCueDirectoryListing(dir, children, safChildLookupCache)
for (child in children) {
if (safScanCancel) {
if (safScanStopRequested()) {
return cancelledResult()
}
@@ -1036,7 +1048,7 @@ internal fun MainActivity.scanSafTree(
val parentDir = cue.parentDir
val cueUri = cueDoc.uri.toString()
val cueAlreadyIndexed = checkpointed(cueUri, cue.lastModified)
if (safScanCancel) {
if (safScanStopRequested()) {
ndjsonWriter?.close()
spill?.abandon()
return cancelledResult()
@@ -1148,7 +1160,7 @@ internal fun MainActivity.scanSafTree(
// Skip resumable and CUE entries before parallel reads.
for (audio in audioFiles) {
val doc = audio.doc
if (safScanCancel) {
if (safScanStopRequested()) {
ndjsonWriter?.close()
spill?.abandon()
return cancelledResult()
@@ -1176,7 +1188,7 @@ internal fun MainActivity.scanSafTree(
// pool is bounded by provider round trips rather than copy memory.
val completed = runSafReadsInOrder(
pendingAudio,
cancelled = { safScanCancel },
cancelled = { safScanStopRequested() },
task = { audio ->
val doc = audio.doc
val stableUri = doc.uri.toString()
@@ -1326,7 +1338,7 @@ internal fun MainActivity.scanSafTreeIncremental(
val lister = SafTreeLister(this)
while (queue.isNotEmpty()) {
if (safScanCancel) {
if (safScanStopRequested()) {
updateSafScanProgress { it.isComplete = true }
val result = JSONObject()
result.put("files", JSONArray())
@@ -1357,7 +1369,7 @@ internal fun MainActivity.scanSafTreeIncremental(
rememberCueDirectoryListing(dir, children, safChildLookupCache)
for (child in children) {
if (safScanCancel) {
if (safScanStopRequested()) {
updateSafScanProgress { it.isComplete = true }
val result = JSONObject()
result.put("files", JSONArray())
@@ -1461,7 +1473,7 @@ internal fun MainActivity.scanSafTreeIncremental(
val cueReferencedAudioUris = mutableSetOf<String>()
for ((cueDoc, parentDir, cueName, cueLastModified) in cueFilesToScan) {
if (safScanCancel) {
if (safScanStopRequested()) {
updateSafScanProgress { it.isComplete = true }
spill.abandon()
val result = JSONObject()
@@ -1620,7 +1632,7 @@ internal fun MainActivity.scanSafTreeIncremental(
val pendingAudio = mutableListOf<ChangedAudio>()
for (audio in audioFiles) {
if (safScanCancel) return cancelledIncrementalResult()
if (safScanStopRequested()) return cancelledIncrementalResult()
if (cueReferencedAudioUris.contains(audio.doc.uri.toString())) {
scanned++
reportProcessed()
@@ -1631,7 +1643,7 @@ internal fun MainActivity.scanSafTreeIncremental(
val completed = runSafReadsInOrder(
pendingAudio,
cancelled = { safScanCancel },
cancelled = { safScanStopRequested() },
task = { audio ->
val ext = audio.name.substringAfterLast('.', "").lowercase(Locale.ROOT)
val fallbackExt = if (ext.isNotBlank()) ".${ext}" else null
@@ -335,6 +335,10 @@ internal object RustCoreBackend : CoreBackend {
override fun cancelLibraryScan() { synchronized(this) { manager }?.cancelLibraryScan() }
override fun pauseLibraryScan() { synchronized(this) { manager }?.pauseLibraryScan() }
override fun resumeLibraryScan() { synchronized(this) { manager }?.resumeLibraryScan() }
override fun parseCueSheet(path: String, audioDirectory: String): String {
val cue = File(path).canonicalFile
val audio = if (audioDirectory.isEmpty()) cue.parentFile!! else File(audioDirectory).canonicalFile
+10 -1
View File
@@ -263,7 +263,8 @@ import UniformTypeIdentifiers
"startBackgroundWork", "updateBackgroundWork", "stopBackgroundWork",
"pickIosDirectory", "startAccessingIosBookmark", "stopAccessingIosBookmark", "downloadCoverToFile", "releaseMemory", "releaseMemoryUnderPressure",
"setLibraryCoverCacheDir", "scanLibraryFolderToNDJSONFile", "scanLibraryFolderIncremental",
"getLibraryScanProgress", "cancelLibraryScan", "parseCueSheet", "extractCoverToFile",
"getLibraryScanProgress", "cancelLibraryScan", "pauseLibraryScan", "resumeLibraryScan",
"parseCueSheet", "extractCoverToFile",
"rewriteSplitArtistTags", "writeM4AFreeformTags", "ensureAC4Config", "writeAC4Metadata", "reEnrichFile",
"checkHiResAuthenticity"]
if call.method == "setScreenAwake" {
@@ -616,6 +617,14 @@ import UniformTypeIdentifiers
try coreBackend.cancelLibraryScan()
return nil
case "pauseLibraryScan":
try coreBackend.pauseLibraryScan()
return nil
case "resumeLibraryScan":
try coreBackend.resumeLibraryScan()
return nil
case "startAccessingIosBookmark":
guard
+16
View File
@@ -177,6 +177,8 @@ protocol CoreBackend {
func scanLibraryFolderIncremental(folder: String, existing: String) throws -> String
func getLibraryScanProgress() throws -> String
func cancelLibraryScan() throws
func pauseLibraryScan() throws
func resumeLibraryScan() throws
func parseCueSheet(path: String, audioDirectory: String) throws -> String
func openDownloadDirectory(path: String) throws -> CoreDirectoryScope
func createTemporaryMediaFile(prefix: String, suffix: String) throws -> URL
@@ -378,6 +380,20 @@ final class RustCoreBackend: CoreBackend {
try current?.cancelLibraryScan()
}
func pauseLibraryScan() throws {
ownerLock.lock()
let current = manager
ownerLock.unlock()
try current?.pauseLibraryScan()
}
func resumeLibraryScan() throws {
ownerLock.lock()
let current = manager
ownerLock.unlock()
try current?.resumeLibraryScan()
}
func parseCueSheet(path: String, audioDirectory: String) throws -> String {
let cue = URL(fileURLWithPath: path).resolvingSymlinksInPath().standardizedFileURL
let audio = audioDirectory.isEmpty ? cue.deletingLastPathComponent() : URL(fileURLWithPath: audioDirectory).resolvingSymlinksInPath().standardizedFileURL
@@ -8,16 +8,23 @@ use std::collections::{BTreeMap, BTreeSet};
use std::io::{BufRead, BufReader, BufWriter, Read, Write};
use std::path::Path;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex, TryLockError, mpsc};
use std::sync::{Arc, Condvar, Mutex, TryLockError, mpsc};
use std::thread;
use std::time::Duration;
type Check<'a> = &'a (dyn Fn() -> Result<(), String> + Sync);
/// How often a paused scan re-checks backend shutdown and its request lease.
const PAUSE_POLL: Duration = Duration::from_millis(250);
#[derive(Default)]
pub(in crate::backend) struct ScanState {
current: Mutex<Option<Arc<Run>>>,
owner: Mutex<()>,
// Pausing is scan-wide rather than per run so a pause requested just
// before the native scan starts still holds it at its first checkpoint.
paused: Mutex<bool>,
resumed: Condvar,
}
#[derive(Default)]
@@ -72,6 +79,51 @@ impl Backend {
{
run.cancelled.store(true, Ordering::Release);
}
// Wake paused workers so they observe the cancellation.
self.set_library_scan_paused(false);
Ok(())
}
/// Holds the scan at its next checkpoint. Workers keep their position in
/// memory, so [`Self::resume_library_scan`] continues with the next file.
pub fn pause_library_scan(&self) -> Result<(), String> {
let _operation = self.enter()?;
self.set_library_scan_paused(true);
Ok(())
}
pub fn resume_library_scan(&self) -> Result<(), String> {
let _operation = self.enter()?;
self.set_library_scan_paused(false);
Ok(())
}
fn set_library_scan_paused(&self, paused: bool) {
*self
.library_scan
.paused
.lock()
.expect("library scan pause lock") = paused;
self.library_scan.resumed.notify_all();
}
fn wait_while_library_scan_paused(&self, run: &Run, check: Check<'_>) -> Result<(), String> {
let state = &self.library_scan;
let mut paused = state.paused.lock().expect("library scan pause lock");
while *paused && !run.cancelled.load(Ordering::Acquire) {
paused = state
.resumed
.wait_timeout(paused, PAUSE_POLL)
.expect("library scan pause lock")
.0;
if *paused {
drop(paused);
// Shutdown and released request leases still end a paused scan.
self.check()?;
check()?;
paused = state.paused.lock().expect("library scan pause lock");
}
}
Ok(())
}
@@ -321,6 +373,9 @@ impl Backend {
previous.cancelled.store(true, Ordering::Release);
}
let check = || {
// Wait before taking the failure lock so paused workers do not
// serialize on it.
self.wait_while_library_scan_paused(&run, check)?;
let mut failure = run.failure.lock().expect("scan failure lock");
if let Some(error) = failure.as_ref() {
return Err(error.clone());
@@ -683,4 +738,70 @@ mod tests {
"native media path must not contain symlinks or special files"
);
}
fn paused_library(root: &Path) -> (Backend, String) {
let library = fs::canonicalize(root).unwrap().join("library");
fs::create_dir_all(&library).unwrap();
for index in 0..3 {
fs::write(library.join(format!("{index}.mp3")), b"not audio").unwrap();
}
let backend = backend(root, &library);
backend.pause_library_scan().unwrap();
(backend, library.to_string_lossy().into_owned())
}
fn wait_for_run(backend: &Backend) {
while backend
.library_scan
.current
.lock()
.expect("library scan lock")
.is_none()
{
thread::sleep(Duration::from_millis(5));
}
}
#[test]
fn paused_scan_holds_its_position_until_resumed() {
let root = tempfile::tempdir().unwrap();
let (backend, folder) = paused_library(root.path());
thread::scope(|scope| {
let scan =
scope.spawn(|| backend.scan_library_folder_incremental(&folder, "{}", &|| Ok(())));
wait_for_run(&backend);
thread::sleep(Duration::from_millis(300));
assert!(!scan.is_finished());
let progress = backend.get_library_scan_progress().unwrap();
assert_eq!(progress["scanned_files"], 0);
assert_eq!(progress["is_complete"], false);
backend.resume_library_scan().unwrap();
let result = scan.join().unwrap().unwrap();
assert_eq!(result["totalFiles"], 3);
});
let progress = backend.get_library_scan_progress().unwrap();
assert_eq!(progress["scanned_files"], 3);
assert_eq!(progress["is_complete"], true);
}
#[test]
fn cancel_releases_a_paused_scan() {
let root = tempfile::tempdir().unwrap();
let (backend, folder) = paused_library(root.path());
thread::scope(|scope| {
let scan =
scope.spawn(|| backend.scan_library_folder_incremental(&folder, "{}", &|| Ok(())));
wait_for_run(&backend);
backend.cancel_library_scan().unwrap();
assert_eq!(scan.join().unwrap().unwrap_err(), "scan cancelled");
});
// Cancelling also clears the pause, so the next scan is not held.
let result = backend
.scan_library_folder_incremental(&folder, "{}", &|| Ok(()))
.unwrap();
assert_eq!(result["totalFiles"], 3);
}
}
+12
View File
@@ -145,6 +145,18 @@ impl ExtensionManager {
.map_err(ExtensionManagerError::Operation)
}
pub fn pause_library_scan(&self) -> Result<(), ExtensionManagerError> {
self.inner
.pause_library_scan()
.map_err(ExtensionManagerError::Operation)
}
pub fn resume_library_scan(&self) -> Result<(), ExtensionManagerError> {
self.inner
.resume_library_scan()
.map_err(ExtensionManagerError::Operation)
}
pub fn rewrite_split_artist_tags(
&self,
path: String,