Index: sdk/lib/io/io_sink.dart |
diff --git a/sdk/lib/io/io_sink.dart b/sdk/lib/io/io_sink.dart |
index a41b5ef959205fa7fde847e9903c0875338757bf..0c2752244994d403db9773487ce45177bf22ba35 100644 |
--- a/sdk/lib/io/io_sink.dart |
+++ b/sdk/lib/io/io_sink.dart |
@@ -12,7 +12,7 @@ part of dart.io; |
* |
* When the [IOSink] is bound to a stream (through [addStream]) any call |
* to the [IOSink] will throw a [StateError]. When the [addStream] completes, |
- * the [IOSink] will again accept all method calls. |
+ * the [IOSink] will again be open for all calls. |
* |
* If data is added to the [IOSink] after the sink is closed, the data will be |
* ignored. Use the [done] future to be notified when the [IOSink] is closed. |
@@ -124,8 +124,138 @@ abstract class IOSink implements StreamSink<List<int>>, StringSink { |
Future get done; |
} |
+class _StreamSinkImpl<T> implements StreamSink<T> { |
+ final StreamConsumer<T> _target; |
+ Completer _doneCompleter = new Completer(); |
+ Future _doneFuture; |
+ StreamController<T> _controllerInstance; |
+ Completer _controllerCompleter; |
+ bool _isClosed = false; |
+ bool _isBound = false; |
+ bool _hasError = false; |
-class _IOSinkImpl extends StreamSinkAdapter<List<int>> implements IOSink { |
+ _StreamSinkImpl(this._target) { |
+ _doneFuture = _doneCompleter.future; |
+ } |
+ |
+ void add(T data) { |
+ if (_isClosed) return; |
+ _controller.add(data); |
+ } |
+ |
+ void addError(error, [StackTrace stackTrace]) => |
+ _controller.addError(error, stackTrace); |
+ |
+ Future addStream(Stream<T> stream) { |
+ if (_isBound) { |
+ throw new StateError("StreamSink is already bound to a stream"); |
+ } |
+ _isBound = true; |
+ if (_hasError) return done; |
+ // Wait for any sync operations to complete. |
+ Future targetAddStream() { |
+ return _target.addStream(stream) |
+ .whenComplete(() { |
+ _isBound = false; |
+ }); |
+ } |
+ if (_controllerInstance == null) return targetAddStream(); |
+ var future = _controllerCompleter.future; |
+ _controllerInstance.close(); |
+ return future.then((_) => targetAddStream()); |
+ } |
+ |
+ Future flush() { |
+ if (_isBound) { |
+ throw new StateError("StreamSink is bound to a stream"); |
+ } |
+ if (_controllerInstance == null) return new Future.value(this); |
+ // Adding an empty stream-controller will return a future that will complete |
+ // when all data is done. |
+ _isBound = true; |
+ var future = _controllerCompleter.future; |
+ _controllerInstance.close(); |
+ return future.whenComplete(() { |
+ _isBound = false; |
+ }); |
+ } |
+ |
+ Future close() { |
+ if (_isBound) { |
+ throw new StateError("StreamSink is bound to a stream"); |
+ } |
+ if (!_isClosed) { |
+ _isClosed = true; |
+ if (_controllerInstance != null) { |
+ _controllerInstance.close(); |
+ } else { |
+ _closeTarget(); |
+ } |
+ } |
+ return done; |
+ } |
+ |
+ void _closeTarget() { |
+ _target.close() |
+ .then((value) => _completeDone(value: value), |
+ onError: (error) => _completeDone(error: error)); |
+ } |
+ |
+ Future get done => _doneFuture; |
+ |
+ void _completeDone({value, error}) { |
+ if (_doneCompleter == null) return; |
+ if (error == null) { |
+ _doneCompleter.complete(value); |
+ } else { |
+ _hasError = true; |
+ _doneCompleter.completeError(error); |
+ } |
+ _doneCompleter = null; |
+ } |
+ |
+ StreamController<T> get _controller { |
+ if (_isBound) { |
+ throw new StateError("StreamSink is bound to a stream"); |
+ } |
+ if (_isClosed) { |
+ throw new StateError("StreamSink is closed"); |
+ } |
+ if (_controllerInstance == null) { |
+ _controllerInstance = new StreamController<T>(sync: true); |
+ _controllerCompleter = new Completer(); |
+ _target.addStream(_controller.stream) |
+ .then( |
+ (_) { |
+ if (_isBound) { |
+ // A new stream takes over - forward values to that stream. |
+ _controllerCompleter.complete(this); |
+ _controllerCompleter = null; |
+ _controllerInstance = null; |
+ } else { |
+ // No new stream, .close was called. Close _target. |
+ _closeTarget(); |
+ } |
+ }, |
+ onError: (error) { |
+ if (_isBound) { |
+ // A new stream takes over - forward errors to that stream. |
+ _controllerCompleter.completeError(error); |
+ _controllerCompleter = null; |
+ _controllerInstance = null; |
+ } else { |
+ // No new stream. No need to close target, as it have already |
+ // failed. |
+ _completeDone(error: error); |
+ } |
+ }); |
+ } |
+ return _controllerInstance; |
+ } |
+} |
+ |
+ |
+class _IOSinkImpl extends _StreamSinkImpl<List<int>> implements IOSink { |
Encoding _encoding; |
bool _encodingMutable = true; |