Chromium Code Reviews
chromiumcodereview-hr@appspot.gserviceaccount.com (chromiumcodereview-hr) | Please choose your nickname with Settings | Help | Chromium Project | Gerrit Changes | Sign out
(299)

Unified Diff: sdk/lib/async/stream_pipe.dart

Issue 11862008: Make Streams also cosider a thrown AsyncError a rethrow. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Fix indentation. Created 7 years, 11 months ago
Use n/p to move between diff chunks; N/P to move between comments. Draft comments are only viewable by you.
Jump to:
View side-by-side diff with in-line comments
Download patch
Index: sdk/lib/async/stream_pipe.dart
diff --git a/sdk/lib/async/stream_pipe.dart b/sdk/lib/async/stream_pipe.dart
index ccf8a8ad7b236177b186da890a2ad756ad38b51e..747ba003502f8c119670a5aa7cb4ad5edd5868c6 100644
--- a/sdk/lib/async/stream_pipe.dart
+++ b/sdk/lib/async/stream_pipe.dart
@@ -69,13 +69,19 @@ class _ForwardingMultiStream<S, T> extends _MultiStreamImpl<T> {
void _handleDone() {
_close();
}
+
+ AsyncError _asyncError(Object error, Object stackTrace, [AsyncError cause]) {
+ if (error is AsyncError) return error;
+ if (cause == null) return new AsyncError(error, stackTrace);
+ return new AsyncError.withCause(error, stackTrace, cause);
+ }
}
abstract class _ForwardingTransformer<S, T> extends _ForwardingMultiStream<S, T>
implements StreamTransformer<S, T> {
Stream<T> bind(Stream<S> source) {
- assert(_source == null);
+ if (_source != null) throw new StateError("Already bound to source.");
_source = source;
if (_hasSubscribers) {
_subscribeToSource();
@@ -101,7 +107,7 @@ class WhereStream<T> extends _ForwardingTransformer<T, T> {
try {
satisfies = _test(inputEvent);
} catch (e, s) {
- _signalError(new AsyncError(e, s));
+ _signalError(_asyncError(e, s));
return;
}
if (satisfies) {
@@ -127,7 +133,7 @@ class MapStream<S, T> extends _ForwardingTransformer<S, T> {
try {
outputEvent = _transform(inputEvent);
} catch (e, s) {
- _signalError(new AsyncError(e, s));
+ _signalError(_asyncError(e, s));
return;
}
_add(outputEvent);
@@ -151,13 +157,14 @@ class ExpandStream<S, T> extends _ForwardingTransformer<S, T> {
} catch (e, s) {
// If either _expand or iterating the generated iterator throws,
// we abort the iteration.
- _signalError(new AsyncError(e, s));
+ _signalError(_asyncError(e, s));
}
}
}
typedef AsyncError _ErrorTransformation(AsyncError error);
+typedef bool _ErrorTest(error);
/**
* A stream pipe that converts or disposes error events
@@ -165,18 +172,31 @@ typedef AsyncError _ErrorTransformation(AsyncError error);
*/
class HandleErrorStream<T> extends _ForwardingTransformer<T, T> {
final _ErrorTransformation _transform;
floitsch 2013/01/11 13:23:33 Type is not right anymore. It's a void function no
Lasse Reichstein Nielsen 2013/01/14 08:33:30 Done.
+ final _ErrorTest _test;
- HandleErrorStream(AsyncError transform(AsyncError event))
- : this._transform = transform;
+ HandleErrorStream(void transform(AsyncError event), bool test(error))
+ : this._transform = transform, this._test = test;
void _handleError(AsyncError error) {
- try {
- error = _transform(error);
- if (error == null) return;
- } catch (e, s) {
- error = new AsyncError.withCause(e, s, error);
+ bool matches = true;
+ if (_test != null) {
+ try {
+ matches = _test(error.error);
+ } catch (e, s) {
+ _signalError(_asyncError(e, s, error));
+ return;
+ }
+ }
+ if (matches) {
+ try {
+ _transform(error);
+ } catch (e, s) {
+ _signalError(_asyncError(e, s, error));
+ return;
+ }
+ } else {
+ _signalError(error);
}
- _signalError(error);
}
}
@@ -214,7 +234,7 @@ class PipeStream<S, T> extends _ForwardingTransformer<S, T> {
try {
return _onData(data, _sink);
} catch (e, s) {
- _signalError(new AsyncError(e, s));
+ _signalError(_asyncError(e, s));
}
}
@@ -222,7 +242,7 @@ class PipeStream<S, T> extends _ForwardingTransformer<S, T> {
try {
_onError(error, _sink);
} catch (e, s) {
- _signalError(new AsyncError.withCause(e, s, error));
+ _signalError(_asyncError(e, s, error));
}
}
@@ -230,7 +250,7 @@ class PipeStream<S, T> extends _ForwardingTransformer<S, T> {
try {
_onDone(_sink);
} catch (e, s) {
- _signalError(new AsyncError(e, s));
+ _signalError(_asyncError(e, s));
}
}
@@ -280,7 +300,7 @@ class TransformStream<S, T> extends _ForwardingTransformer<S, T> {
try {
return _transform.handleData(data, _sink);
} catch (e, s) {
- _controller.signalError(new AsyncError(e, s));
+ _controller.signalError(_asyncError(e, s));
}
}
@@ -288,7 +308,7 @@ class TransformStream<S, T> extends _ForwardingTransformer<S, T> {
try {
_transform.handleError(error, _sink);
} catch (e, s) {
- _controller.signalError(new AsyncError.withCause(e, s, error));
+ _controller.signalError(_asyncError(e, s, error));
}
}
@@ -296,7 +316,7 @@ class TransformStream<S, T> extends _ForwardingTransformer<S, T> {
try {
_transform.handleDone(_sink);
} catch (e, s) {
- _controller.signalError(new AsyncError(e, s));
+ _controller.signalError(_asyncError(e, s));
}
}
}
@@ -367,7 +387,7 @@ class TakeWhileStream<T> extends _ForwardingTransformer<T, T> {
try {
satisfies = _test(inputEvent);
} catch (e, s) {
- _signalError(new AsyncError(e, s));
+ _signalError(_asyncError(e, s));
// The test didn't say true. Didn't say false either, but we stop anyway.
_close();
return;
@@ -414,7 +434,7 @@ class SkipWhileStream<T> extends _ForwardingTransformer<T, T> {
try {
satisfies = _test(inputEvent);
} catch (e, s) {
- _signalError(new AsyncError(e, s));
+ _signalError(_asyncError(e, s));
// A failure to return a boolean is considered "not matching".
_hasFailed = true;
return;
@@ -450,7 +470,7 @@ class DistinctStream<T> extends _ForwardingTransformer<T, T> {
isEqual = _equals(_previous, inputEvent);
}
} catch (e, s) {
- _signalError(new AsyncError(e, s));
+ _signalError(_asyncError(e, s));
return null;
}
if (!isEqual) {

Powered by Google App Engine
This is Rietveld 408576698