Ensure ready future completes (with error) in the case the watcher fails and closes. (dart-lang/watcher#123)
See dart-lang/watcher#115.
diff --git a/pkgs/watcher/CHANGELOG.md b/pkgs/watcher/CHANGELOG.md
index d987225..af54045 100644
--- a/pkgs/watcher/CHANGELOG.md
+++ b/pkgs/watcher/CHANGELOG.md
@@ -1,6 +1,7 @@
# 1.0.2-dev
- Require Dart SDK >= 2.14
+- Ensure `DirectoryWatcher.ready` completes even when errors occur that close the watcher.
# 1.0.1
diff --git a/pkgs/watcher/lib/src/directory_watcher/linux.dart b/pkgs/watcher/lib/src/directory_watcher/linux.dart
index 1bf5efd..dafaeb1 100644
--- a/pkgs/watcher/lib/src/directory_watcher/linux.dart
+++ b/pkgs/watcher/lib/src/directory_watcher/linux.dart
@@ -82,20 +82,32 @@
// Batch the inotify changes together so that we can dedup events.
var innerStream = _nativeEvents.stream.batchEvents();
- _listen(innerStream, _onBatch, onError: _eventsController.addError);
-
- _listen(Directory(path).list(recursive: true), (FileSystemEntity entity) {
- if (entity is Directory) {
- _watchSubdir(entity.path);
- } else {
- _files.add(entity.path);
+ _listen(innerStream, _onBatch,
+ onError: (Object error, StackTrace stackTrace) {
+ // Guarantee that ready always completes.
+ if (!isReady) {
+ _readyCompleter.complete();
}
- }, onError: (Object error, StackTrace stackTrace) {
_eventsController.addError(error, stackTrace);
- close();
- }, onDone: () {
- _readyCompleter.complete();
- }, cancelOnError: true);
+ });
+
+ _listen(
+ Directory(path).list(recursive: true),
+ (FileSystemEntity entity) {
+ if (entity is Directory) {
+ _watchSubdir(entity.path);
+ } else {
+ _files.add(entity.path);
+ }
+ },
+ onError: _emitError,
+ onDone: () {
+ if (!isReady) {
+ _readyCompleter.complete();
+ }
+ },
+ cancelOnError: true,
+ );
}
@override
@@ -195,15 +207,15 @@
// contents,
if (files.contains(path) && _files.contains(path)) continue;
for (var file in _files.remove(path)) {
- _emit(ChangeType.REMOVE, file);
+ _emitEvent(ChangeType.REMOVE, file);
}
}
for (var file in files) {
if (_files.contains(file)) {
- _emit(ChangeType.MODIFY, file);
+ _emitEvent(ChangeType.MODIFY, file);
} else {
- _emit(ChangeType.ADD, file);
+ _emitEvent(ChangeType.ADD, file);
_files.add(file);
}
}
@@ -221,15 +233,14 @@
_watchSubdir(entity.path);
} else {
_files.add(entity.path);
- _emit(ChangeType.ADD, entity.path);
+ _emitEvent(ChangeType.ADD, entity.path);
}
}, onError: (Object error, StackTrace stackTrace) {
// Ignore an exception caused by the dir not existing. It's fine if it
// was added and then quickly removed.
if (error is FileSystemException) return;
- _eventsController.addError(error, stackTrace);
- close();
+ _emitError(error, stackTrace);
}, cancelOnError: true);
}
@@ -242,7 +253,7 @@
// caused by a MOVE, we need to manually emit events.
if (isReady) {
for (var file in _files.paths) {
- _emit(ChangeType.REMOVE, file);
+ _emitEvent(ChangeType.REMOVE, file);
}
}
@@ -251,12 +262,22 @@
/// Emits a [WatchEvent] with [type] and [path] if this watcher is in a state
/// to emit events.
- void _emit(ChangeType type, String path) {
+ void _emitEvent(ChangeType type, String path) {
if (!isReady) return;
if (_eventsController.isClosed) return;
_eventsController.add(WatchEvent(type, path));
}
+ /// Emit an error, then close the watcher.
+ void _emitError(Object error, StackTrace stackTrace) {
+ // Guarantee that ready always completes.
+ if (!isReady) {
+ _readyCompleter.complete();
+ }
+ _eventsController.addError(error, stackTrace);
+ close();
+ }
+
/// Like [Stream.listen], but automatically adds the subscription to
/// [_subscriptions] so that it can be canceled when [close] is called.
void _listen<T>(Stream<T> stream, void Function(T) onData,
diff --git a/pkgs/watcher/lib/src/directory_watcher/mac_os.dart b/pkgs/watcher/lib/src/directory_watcher/mac_os.dart
index 2046ce0..12648c8 100644
--- a/pkgs/watcher/lib/src/directory_watcher/mac_os.dart
+++ b/pkgs/watcher/lib/src/directory_watcher/mac_os.dart
@@ -87,8 +87,11 @@
//
// If we do receive a batch of events, [_onBatch] will ensure that these
// futures don't fire and that the directory is re-listed.
- Future.wait([_listDir(), _waitForBogusEvents()])
- .then((_) => _readyCompleter.complete());
+ Future.wait([_listDir(), _waitForBogusEvents()]).then((_) {
+ if (!isReady) {
+ _readyCompleter.complete();
+ }
+ });
}
@override
@@ -115,7 +118,11 @@
// Cancel the timer because bogus events only occur in the first batch, so
// we can fire [ready] as soon as we're done listing the directory.
_bogusEventTimer.cancel();
- _listDir().then((_) => _readyCompleter.complete());
+ _listDir().then((_) {
+ if (!isReady) {
+ _readyCompleter.complete();
+ }
+ });
return;
}
@@ -396,6 +403,10 @@
/// Emit an error, then close the watcher.
void _emitError(Object error, StackTrace stackTrace) {
+ // Guarantee that ready always completes.
+ if (!isReady) {
+ _readyCompleter.complete();
+ }
_eventsController.addError(error, stackTrace);
close();
}
diff --git a/pkgs/watcher/lib/src/directory_watcher/polling.dart b/pkgs/watcher/lib/src/directory_watcher/polling.dart
index 2a43937..3dcec05 100644
--- a/pkgs/watcher/lib/src/directory_watcher/polling.dart
+++ b/pkgs/watcher/lib/src/directory_watcher/polling.dart
@@ -43,11 +43,11 @@
final _events = StreamController<WatchEvent>.broadcast();
@override
- bool get isReady => _ready.isCompleted;
+ bool get isReady => _readyCompleter.isCompleted;
@override
- Future<void> get ready => _ready.future;
- final _ready = Completer<void>();
+ Future<void> get ready => _readyCompleter.future;
+ final _readyCompleter = Completer<void>();
/// The amount of time the watcher pauses between successive polls of the
/// directory contents.
@@ -119,6 +119,10 @@
if (entity is! File) return;
_filesToProcess.add(entity.path);
}, onError: (Object error, StackTrace stackTrace) {
+ // Guarantee that ready always completes.
+ if (!isReady) {
+ _readyCompleter.complete();
+ }
if (!isDirectoryNotFoundException(error)) {
// It's some unknown error. Pipe it over to the event stream so the
// user can see it.
@@ -177,7 +181,7 @@
_lastModifieds.remove(removed);
}
- if (!isReady) _ready.complete();
+ if (!isReady) _readyCompleter.complete();
// Wait and then poll again.
await Future.delayed(_pollingDelay);
diff --git a/pkgs/watcher/lib/src/directory_watcher/windows.dart b/pkgs/watcher/lib/src/directory_watcher/windows.dart
index 09b4b36..710caf5 100644
--- a/pkgs/watcher/lib/src/directory_watcher/windows.dart
+++ b/pkgs/watcher/lib/src/directory_watcher/windows.dart
@@ -92,7 +92,9 @@
_listDir().then((_) {
_startWatch();
_startParentWatcher();
- _readyCompleter.complete();
+ if (!isReady) {
+ _readyCompleter.complete();
+ }
});
}
@@ -427,6 +429,10 @@
/// Emit an error, then close the watcher.
void _emitError(Object error, StackTrace stackTrace) {
+ // Guarantee that ready always completes.
+ if (!isReady) {
+ _readyCompleter.complete();
+ }
_eventsController.addError(error, stackTrace);
close();
}
diff --git a/pkgs/watcher/lib/src/utils.dart b/pkgs/watcher/lib/src/utils.dart
index ecf4e10..c2e71b3 100644
--- a/pkgs/watcher/lib/src/utils.dart
+++ b/pkgs/watcher/lib/src/utils.dart
@@ -3,8 +3,8 @@
// BSD-style license that can be found in the LICENSE file.
import 'dart:async';
-import 'dart:io';
import 'dart:collection';
+import 'dart:io';
/// Returns `true` if [error] is a [FileSystemException] for a missing
/// directory.
diff --git a/pkgs/watcher/lib/watcher.dart b/pkgs/watcher/lib/watcher.dart
index 22e0d6e..12a5369 100644
--- a/pkgs/watcher/lib/watcher.dart
+++ b/pkgs/watcher/lib/watcher.dart
@@ -42,6 +42,9 @@
///
/// If the watcher is already monitoring, this returns an already complete
/// future.
+ ///
+ /// This future always completes successfully as errors are provided through
+ /// the [events] stream.
Future get ready;
/// Creates a new [DirectoryWatcher] or [FileWatcher] monitoring [path],
diff --git a/pkgs/watcher/test/ready/shared.dart b/pkgs/watcher/test/ready/shared.dart
index b6b3cce..6578c52 100644
--- a/pkgs/watcher/test/ready/shared.dart
+++ b/pkgs/watcher/test/ready/shared.dart
@@ -61,4 +61,17 @@
// Should be back to not ready.
expect(watcher.ready, doesNotComplete);
});
+
+ test('completes even if directory does not exist', () async {
+ var watcher = createWatcher(path: 'does/not/exist');
+
+ // Subscribe to the events (else ready will never fire).
+ var subscription = watcher.events.listen((event) {}, onError: (error) {});
+
+ // Expect ready still completes.
+ await watcher.ready;
+
+ // Now unsubscribe.
+ await subscription.cancel();
+ });
}