From 0b08a83d4b6d4d1ebeea98a0b6e94f46bb3ff33e Mon Sep 17 00:00:00 2001 From: fufesou Date: Thu, 27 Aug 2026 13:01:11 +0800 Subject: [PATCH] fix(file-transfer): improve large directory loading (#15830) * fix(file-transfer): improve large directory loading Signed-off-by: fufesou * fix(file-transfer): avoid failing newer directory reads Track each remote directory request by its registered completer and only remove the task when it still matches, preventing stale failures from affecting newer requests for the same path. Signed-off-by: fufesou * fix(file-transfer): handle slow directory listings safely Signed-off-by: fufesou * fix(file transfer): correlate directory responses with requests Signed-off-by: fufesou * fix(file transfer): prevent automatic directory responses from matching requests Signed-off-by: fufesou * fix(file-transfer): handle large remote directory listings reliably - build file rows lazily - register remote reads before sending requests - handle Home paths, stale responses, errors, and timeouts - serialize same-path reads with different hidden-file options Signed-off-by: fufesou * fix(file transfer): reduce diffs Signed-off-by: fufesou * fix: build Signed-off-by: fufesou * fix: invalidate pending dir reads on reconnect Signed-off-by: fufesou * test(file-transfer): cover remote directory read lifecycle Signed-off-by: fufesou --------- Signed-off-by: fufesou --- .../lib/desktop/pages/file_manager_page.dart | 5 +- flutter/lib/models/file_model.dart | 165 +++++++--- flutter/test/file_model_test.dart | 281 ++++++++++++++++++ src/ui_cm_interface.rs | 34 ++- 4 files changed, 440 insertions(+), 45 deletions(-) create mode 100644 flutter/test/file_model_test.dart diff --git a/flutter/lib/desktop/pages/file_manager_page.dart b/flutter/lib/desktop/pages/file_manager_page.dart index e1130fdaa..17674c268 100644 --- a/flutter/lib/desktop/pages/file_manager_page.dart +++ b/flutter/lib/desktop/pages/file_manager_page.dart @@ -1126,6 +1126,7 @@ class _FileManagerViewState extends State { return element.name.contains(_searchText.value); }).toList(growable: false) : entries; + // Keep rows lazy so large directories only build visible list items. final rows = filteredEntries.map((entry) { final sizeStr = entry.isFile ? readableFileSize(entry.size.toDouble()) : ""; @@ -1308,7 +1309,7 @@ class _FileManagerViewState extends State { ], ))), ); - }).toList(growable: false); + }); return Column( children: [ @@ -1324,7 +1325,7 @@ class _FileManagerViewState extends State { controller: scrollController, itemExtent: kDesktopFileTransferRowHeight, itemBuilder: (context, index) { - return rows[index]; + return rows.elementAt(index); }, itemCount: rows.length, ), diff --git a/flutter/lib/models/file_model.dart b/flutter/lib/models/file_model.dart index 94f0fcb7b..22bf1eab6 100644 --- a/flutter/lib/models/file_model.dart +++ b/flutter/lib/models/file_model.dart @@ -46,6 +46,12 @@ class JobID { typedef GetSessionID = SessionID Function(); typedef GetDialogManager = OverlayDialogManager? Function(); +typedef ReadRemoteDirectory = Future Function( + SessionID sessionId, String path, bool includeHidden); + +const _kRemoteReadDirTimeout = Duration(seconds: 30); +const _kRemoteSessionChangedError = + 'Remote directory read cancelled because the session changed'; class FileModel { final WeakReference parent; @@ -84,6 +90,7 @@ class FileModel { } Future onReady() async { + fileFetcher.beginRemoteSession(); await evtLoop.onReady(); if (!isWeb) await localController.onReady(); await remoteController.onReady(); @@ -133,7 +140,11 @@ class FileModel { final id = int.tryParse(evt['id']?.toString() ?? ''); if (id != null) { final err = evt['err']?.toString() ?? 'Unknown error'; - fileFetcher.tryCompleteRecursiveTaskWithError(id, err); + if (id == 0) { + fileFetcher.tryCompleteRemoteTaskWithError(err); + } else { + fileFetcher.tryCompleteRecursiveTaskWithError(id, err); + } } // Always call jobController.jobError(evt) to ensure all error events are processed, // even if the event does not have a valid job ID. This allows for generic error handling @@ -350,6 +361,8 @@ class FileController { final history = RxList.empty(growable: true); final sortBy = SortBy.name.obs; var sortAscending = true; + // Incremented for each navigation; only the latest generation applies results. + int _directoryRequestGeneration = 0; final JobController jobController; final WeakReference rootState; @@ -484,12 +497,19 @@ class FileController { path = "$path\\"; } } + final requestGeneration = ++_directoryRequestGeneration; try { final fd = await fileFetcher.fetchDirectory(path, isLocal, showHidden); + if (requestGeneration != _directoryRequestGeneration) { + return true; + } fd.format(isWindows, sort: sortBy.value); directory.value = fd; return true; } catch (e) { + if (requestGeneration != _directoryRequestGeneration) { + return true; + } debugPrint("Failed to openDirectory $path: $e"); return false; } @@ -541,6 +561,7 @@ class FileController { void initDirAndHome(Map evt) { try { final fd = FileDirectory.fromJson(jsonDecode(evt['value'])); + final isHomeResponse = fileFetcher.isLikelyRemoteHomeResponse(fd.path); fd.format(options.value.isWindows, sort: sortBy.value); if (fd.id > 0) { final jobIndex = jobController.getJob(fd.id); @@ -556,10 +577,12 @@ class FileController { debugPrint("update receive details: ${fd.path}"); jobController.jobTable.refresh(); } - } else if (options.value.home.isEmpty) { + } else if (options.value.home.isEmpty && isHomeResponse) { options.value.home = fd.path; debugPrint("init remote home: ${fd.path}"); - directory.value = fd; + if (_directoryRequestGeneration == 0) { + directory.value = fd; + } } } catch (e) { debugPrint("initDirAndHome err=$e"); @@ -1362,16 +1385,78 @@ class JobResultListener { } } +class _RemoteReadTask { + final bool includeHidden; + final Completer completer = Completer(); + final Completer released = Completer(); + late final Timer timer; + + _RemoteReadTask(this.includeHidden); +} + class FileFetcher { // Map> localTasks = {}; // now we only use read local dir sync - Map> remoteTasks = {}; + final Map _remoteReadTasks = {}; Map>> remoteEmptyDirsTasks = {}; Map> readRecursiveTasks = {}; + int _remoteSessionGeneration = 0; final GetSessionID getSessionID; + final ReadRemoteDirectory _readRemoteDirectory; SessionID get sessionId => getSessionID(); - FileFetcher(this.getSessionID); + FileFetcher(this.getSessionID, {ReadRemoteDirectory? readRemoteDirectory}) + : _readRemoteDirectory = readRemoteDirectory ?? + ((sessionId, path, includeHidden) => bind.sessionReadRemoteDir( + sessionId: sessionId, + path: path, + includeHidden: includeHidden)); + + bool hasPendingRemoteRead(String path) => _remoteReadTasks.containsKey(path); + + bool isLikelyRemoteHomeResponse(String path) => + _remoteReadTasks.isEmpty || + (_remoteReadTasks.length == 1 && + hasPendingRemoteRead("") && + !hasPendingRemoteRead(path)); + + void beginRemoteSession() { + _remoteSessionGeneration++; + final pendingTasks = _remoteReadTasks.entries.toList(growable: false); + for (final entry in pendingTasks) { + final task = entry.value; + if (!_removeRemoteReadTask(entry.key, task)) continue; + task.completer.completeError(StateError(_kRemoteSessionChangedError)); + } + } + + _RemoteReadTask _registerRemoteReadTask(String path, bool includeHidden) { + if (hasPendingRemoteRead(path)) { + throw "Failed to registerReadTask, already have same read job"; + } + final task = _RemoteReadTask(includeHidden); + _remoteReadTasks[path] = task; + task.timer = Timer(_kRemoteReadDirTimeout, () { + if (!_removeRemoteReadTask(path, task)) return; + task.completer.completeError("Failed to read dir, timeout"); + }); + return task; + } + + bool _removeRemoteReadTask(String path, _RemoteReadTask task) { + if (!identical(_remoteReadTasks[path], task)) return false; + _remoteReadTasks.remove(path); + task.timer.cancel(); + task.released.complete(); + return true; + } + + bool _completeRemoteReadTask(String path, FileDirectory directory) { + final task = _remoteReadTasks[path]; + if (task == null || !_removeRemoteReadTask(path, task)) return false; + task.completer.complete(directory); + return true; + } Future> registerReadEmptyDirsTask( bool isLocal, String path) { @@ -1391,23 +1476,6 @@ class FileFetcher { return c.future; } - Future registerReadTask(bool isLocal, String path) { - // final jobs = isLocal?localJobs:remoteJobs; // maybe we will use read local dir async later - final tasks = remoteTasks; // bypass now - if (tasks.containsKey(path)) { - throw "Failed to registerReadTask, already have same read job"; - } - final c = Completer(); - tasks[path] = c; - - Timer(Duration(seconds: 2), () { - tasks.remove(path); - if (c.isCompleted) return; - c.completeError("Failed to read dir, timeout"); - }); - return c.future; - } - Future registerReadRecursiveTask(int actID) { final tasks = readRecursiveTasks; if (tasks.containsKey(actID)) { @@ -1445,27 +1513,37 @@ class FileFetcher { tryCompleteTask(String? msg, String? isLocalStr) { if (msg == null || isLocalStr == null) return; - late final Map> tasks; try { final fd = FileDirectory.fromJson(jsonDecode(msg)); if (fd.id > 0) { // fd.id > 0 is result for read recursive - // to-do later,will be better if every fetch use ID,so that there will only one task map for read and recursive read - tasks = readRecursiveTasks; - final completer = tasks.remove(fd.id); - completer?.complete(fd); - } else if (fd.path.isNotEmpty) { - // result for normal read dir - // final jobs = isLocal?localJobs:remoteJobs; // maybe we will use read local dir async later - tasks = remoteTasks; // bypass now - final completer = tasks.remove(fd.path); + final completer = readRecursiveTasks.remove(fd.id); completer?.complete(fd); + return; + } + if (isLocalStr == "false" && fd.path.isNotEmpty) { + if (_completeRemoteReadTask(fd.path, fd)) { + return; + } + // A Home request uses an empty path but returns its resolved path. + if (isLikelyRemoteHomeResponse(fd.path)) { + _completeRemoteReadTask("", fd); + } } } catch (e) { debugPrint("tryCompleteJob err: $e"); } } + bool tryCompleteRemoteTaskWithError(String error) { + if (_remoteReadTasks.length != 1) return false; + final entry = _remoteReadTasks.entries.single; + final task = entry.value; + if (!_removeRemoteReadTask(entry.key, task)) return false; + task.completer.completeError(error); + return true; + } + // Complete a pending recursive read task with an error. // See FileModel.handleJobError() for why this is necessary. void tryCompleteRecursiveTaskWithError(int id, String error) { @@ -1506,9 +1584,26 @@ class FileFetcher { final fd = FileDirectory.fromJson(jsonDecode(res)); return fd; } else { - await bind.sessionReadRemoteDir( - sessionId: sessionId, path: path, includeHidden: showHidden); - return registerReadTask(isLocal, path); + final remoteSessionGeneration = _remoteSessionGeneration; + final pendingTask = _remoteReadTasks[path]; + if (pendingTask != null) { + if (pendingTask.includeHidden == showHidden) { + return pendingTask.completer.future; + } + await pendingTask.released.future; + if (remoteSessionGeneration != _remoteSessionGeneration) { + throw StateError(_kRemoteSessionChangedError); + } + return fetchDirectory(path, isLocal, showHidden); + } + final task = _registerRemoteReadTask(path, showHidden); + unawaited(Future.sync( + () => _readRemoteDirectory(sessionId, path, showHidden)) + .catchError((Object error, StackTrace stackTrace) { + if (!_removeRemoteReadTask(path, task)) return; + task.completer.completeError(error, stackTrace); + })); + return task.completer.future; } } catch (e) { return Future.error(e); diff --git a/flutter/test/file_model_test.dart b/flutter/test/file_model_test.dart new file mode 100644 index 000000000..9455f2cab --- /dev/null +++ b/flutter/test/file_model_test.dart @@ -0,0 +1,281 @@ +import 'dart:async'; +import 'dart:convert'; + +import 'package:flutter_hbb/models/file_model.dart'; +import 'package:flutter_hbb/models/model.dart'; +import 'package:flutter_test/flutter_test.dart'; +import 'package:uuid/uuid.dart'; + +final _sessionId = UuidValue('00000000-0000-0000-0000-000000000000'); + +class _FakeFFI implements FFI { + @override + String id = 'test-peer'; + @override + UuidValue get sessionId => _sessionId; + @override + late final FfiModel ffiModel = FfiModel(WeakReference(this)); + + @override + dynamic noSuchMethod(Invocation invocation) => super.noSuchMethod(invocation); +} + +FileController _createController(FileFetcher fileFetcher) { + final ffi = _FakeFFI(); + return FileController( + isLocal: false, + getSessionID: () => _sessionId, + rootState: WeakReference(ffi), + jobController: JobController(() => _sessionId, () => null), + fileFetcher: fileFetcher, + getOtherSideDirectoryData: () => + DirectoryData(FileDirectory(), DirectoryOptions()), + ); +} + +FileDirectory _directory(String path) => FileDirectory()..path = path; + +String _directoryJson(String path) => jsonEncode({ + 'id': 0, + 'path': path, + 'entries': [], + }); + +class _SentRead { + final String path; + final bool includeHidden; + + const _SentRead(this.path, this.includeHidden); +} + +void main() { + test('a fast remote response is matched after registration', () async { + late final FileFetcher fileFetcher; + fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, path, __) { + fileFetcher.tryCompleteTask(_directoryJson(path), 'false'); + return Future.value(); + }, + ); + final directory = await fileFetcher.fetchDirectory('/fast', false, false); + expect(directory.path, '/fast'); + }); + + test('a send failure fails and removes its registered task', () async { + final failure = StateError('send failed'); + final fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, __, ___) => Future.error(failure), + ); + await expectLater( + fileFetcher.fetchDirectory('/failed', false, false), + throwsA(same(failure)), + ); + expect(fileFetcher.hasPendingRemoteRead('/failed'), isFalse); + }); + + test('a resolved Home path completes the sole empty-path request', () async { + final sent = <_SentRead>[]; + final fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, path, includeHidden) async { + sent.add(_SentRead(path, includeHidden)); + }, + ); + final controller = _createController(fileFetcher); + controller.directory.value = _directory('/initial'); + + final home = controller.openDirectory(''); + await Future.delayed(Duration.zero); + final response = _directoryJson('/home/user'); + controller.initDirAndHome({'value': response}); + expect(controller.homePath, '/home/user'); + expect(controller.directory.value.path, '/initial'); + fileFetcher.tryCompleteTask(response, 'false'); + expect(await home, isTrue); + expect(controller.directory.value.path, '/home/user'); + expect(sent.single.path, isEmpty); + }); + + test('an automatic response initializes Home without a pending request', () { + final controller = _createController(FileFetcher(() => _sessionId)); + + controller.initDirAndHome({'value': _directoryJson('/home/user')}); + + expect(controller.homePath, '/home/user'); + expect(controller.directory.value.path, '/home/user'); + }); + + test('an exact path response is not taken by a pending Home request', + () async { + final sent = <_SentRead>[]; + final fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, path, includeHidden) async { + sent.add(_SentRead(path, includeHidden)); + }, + ); + final home = fileFetcher.fetchDirectory('', false, false); + final regular = fileFetcher.fetchDirectory('/regular', false, false); + await Future.delayed(Duration.zero); + var homeCompleted = false; + home.then((_) => homeCompleted = true); + + fileFetcher.tryCompleteTask(_directoryJson('/unmatched'), 'false'); + await Future.delayed(Duration.zero); + expect(homeCompleted, isFalse); + + fileFetcher.tryCompleteTask(_directoryJson('/regular'), 'false'); + expect((await regular).path, '/regular'); + await Future.delayed(Duration.zero); + expect(homeCompleted, isFalse); + + fileFetcher.tryCompleteTask(_directoryJson('/home/user'), 'false'); + expect((await home).path, '/home/user'); + expect(sent.map((request) => request.path), ['', '/regular']); + }); + + test('a read error completes the sole pending request', () async { + final fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, __, ___) async {}, + ); + final request = fileFetcher.fetchDirectory('/denied', false, false); + await Future.delayed(Duration.zero); + final expectation = expectLater(request, throwsA('permission denied')); + + fileFetcher.tryCompleteRemoteTaskWithError('permission denied'); + + await expectation; + expect(fileFetcher.hasPendingRemoteRead('/denied'), isFalse); + }); + + test('same-path requests share the pending read', () async { + final sent = <_SentRead>[]; + late final FileFetcher fileFetcher; + fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, path, includeHidden) async { + sent.add(_SentRead(path, includeHidden)); + }, + ); + final controller = _createController(fileFetcher); + controller.directory.value = _directory('/initial'); + + final first = controller.openDirectory('/same'); + final waiting = controller.openDirectory('/same'); + + await Future.delayed(Duration.zero); + expect(sent.map((request) => request.path), ['/same']); + fileFetcher.tryCompleteTask(_directoryJson('/same'), 'false'); + expect(await first, isTrue); + await Future.delayed(Duration.zero); + + expect(sent.map((request) => request.path), ['/same']); + expect(await waiting, isTrue); + expect(controller.directory.value.path, '/same'); + }); + + test('same-path requests with different hidden options are serialized', + () async { + final sent = <_SentRead>[]; + final fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, path, includeHidden) async { + sent.add(_SentRead(path, includeHidden)); + }, + ); + + final first = fileFetcher.fetchDirectory('/same', false, false); + final second = fileFetcher.fetchDirectory('/same', false, true); + + await Future.delayed(Duration.zero); + expect(sent.map((request) => request.includeHidden), [false]); + + fileFetcher.tryCompleteTask(_directoryJson('/same'), 'false'); + expect((await first).path, '/same'); + await Future.delayed(Duration.zero); + expect(sent.map((request) => request.includeHidden), [false, true]); + + fileFetcher.tryCompleteTask(_directoryJson('/same'), 'false'); + expect((await second).path, '/same'); + }); + + test('session invalidation cancels active and waiting reads', () async { + final sent = <_SentRead>[]; + final fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, path, includeHidden) async { + sent.add(_SentRead(path, includeHidden)); + }, + ); + final first = fileFetcher.fetchDirectory('/same', false, false); + final waiting = fileFetcher.fetchDirectory('/same', false, true); + await Future.delayed(Duration.zero); + final firstError = expectLater(first, throwsA(isA())); + final waitingError = expectLater(waiting, throwsA(isA())); + + fileFetcher.beginRemoteSession(); + + await firstError; + await Future.delayed(Duration.zero); + expect(sent.map((request) => request.includeHidden), [false]); + await waitingError; + expect(fileFetcher.hasPendingRemoteRead('/same'), isFalse); + final replacement = fileFetcher.fetchDirectory('/same', false, true); + await Future.delayed(Duration.zero); + expect(sent.map((request) => request.includeHidden), [false, true]); + fileFetcher.tryCompleteTask(_directoryJson('/same'), 'false'); + expect((await replacement).path, '/same'); + }); + + test('a late dispatch failure cannot remove a replacement task', () async { + final dispatches = >[]; + final fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, __, ___) { + final dispatch = Completer(); + dispatches.add(dispatch); + return dispatch.future; + }, + ); + final first = fileFetcher.fetchDirectory('/same', false, false); + await Future.delayed(Duration.zero); + final firstError = expectLater(first, throwsA(isA())); + fileFetcher.beginRemoteSession(); + await firstError; + + final replacement = fileFetcher.fetchDirectory('/same', false, false); + await Future.delayed(Duration.zero); + expect(dispatches, hasLength(2)); + dispatches.first.completeError(StateError('late dispatch failure')); + await Future.delayed(Duration.zero); + + expect(fileFetcher.hasPendingRemoteRead('/same'), isTrue); + fileFetcher.tryCompleteTask(_directoryJson('/same'), 'false'); + expect((await replacement).path, '/same'); + dispatches.last.complete(); + await Future.delayed(Duration.zero); + }); + + test('navigation ignores stale directory responses', () async { + final fileFetcher = FileFetcher( + () => _sessionId, + readRemoteDirectory: (_, __, ___) async {}, + ); + final controller = _createController(fileFetcher); + controller.directory.value = _directory('/initial'); + + final stale = controller.openDirectory('/stale'); + final latest = controller.openDirectory('/latest'); + await Future.delayed(Duration.zero); + + fileFetcher.tryCompleteTask(_directoryJson('/latest'), 'false'); + expect(await latest, isTrue); + fileFetcher.tryCompleteTask(_directoryJson('/stale'), 'false'); + expect(await stale, isTrue); + + expect(controller.directory.value.path, '/latest'); + }); +} diff --git a/src/ui_cm_interface.rs b/src/ui_cm_interface.rs index 1474ce093..5e13ef82b 100644 --- a/src/ui_cm_interface.rs +++ b/src/ui_cm_interface.rs @@ -1546,13 +1546,19 @@ async fn read_dir(dir: &str, include_hidden: bool, tx: &UnboundedSender) { fs::get_path(dir) } }; - if let Ok(Ok(fd)) = spawn_blocking(move || fs::read_dir(&path, include_hidden)).await { - let mut msg_out = Message::new(); - let mut file_response = FileResponse::new(); - file_response.set_dir(fd); - msg_out.set_file_response(file_response); - send_raw(msg_out, tx); - } + let result = spawn_blocking(move || fs::read_dir(&path, include_hidden)).await; + let msg_out = match result { + Ok(Ok(fd)) => { + let mut msg_out = Message::new(); + let mut file_response = FileResponse::new(); + file_response.set_dir(fd); + msg_out.set_file_response(file_response); + msg_out + } + Ok(Err(err)) => fs::new_error(0, err, -1), + Err(err) => fs::new_error(0, err, -1), + }; + send_raw(msg_out, tx); } #[cfg(not(any(target_os = "ios")))] @@ -1750,7 +1756,7 @@ mod tests { #[test] #[cfg(not(any(target_os = "ios")))] - fn read_dir_success() { + fn read_dir_reports_success_and_error() { let rt = Runtime::new().unwrap(); rt.block_on(async { let (tx, mut rx) = unbounded_channel(); @@ -1773,6 +1779,18 @@ mod tests { _ => panic!("unexpected data"), } let _ = fs::remove_dir_all(&dir); + + super::read_dir(&dir.to_string_lossy(), false, &tx).await; + + match rx.recv().await.unwrap() { + Data::RawMessage(bytes) => { + let mut msg = Message::new(); + msg.merge_from_bytes(&bytes).unwrap(); + assert_eq!(msg.file_response().error().id, 0); + assert!(!msg.file_response().error().error.is_empty()); + } + _ => panic!("unexpected data"), + } }); }