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"), + } }); }