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

Unified Diff: sdk/lib/io/string_transformer.dart

Issue 12313127: Fix subscription handling in stream decoder (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Added missing file Created 7 years, 10 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/io/string_transformer.dart
diff --git a/sdk/lib/io/string_transformer.dart b/sdk/lib/io/string_transformer.dart
index 2b34622781a2fb5df3f88e79539d9bf35b57578a..8c0574e19ee1b40742904ab899d1f3729e11abb8 100644
--- a/sdk/lib/io/string_transformer.dart
+++ b/sdk/lib/io/string_transformer.dart
@@ -153,131 +153,81 @@ List<int> _encodeString(String string, [Encoding encoding = Encoding.UTF_8]) {
}
-class LineTransformer implements StreamTransformer<String, String> {
+class LineTransformer extends StreamEventTransformer<String, String> {
const int _LF = 10;
const int _CR = 13;
final StringBuffer _buffer = new StringBuffer();
-
- StreamSubscription<String> _subscription;
- StreamController<String> _controller;
String _carry;
- Stream<String> bind(Stream<String> stream) {
- _controller = new StreamController<String>(
- onPauseStateChange: _pauseChanged,
- onSubscriptionStateChange: _subscriptionChanged);
-
- void handle(String data, bool isClosing) {
- if (_carry != null) {
- data = _carry.concat(data);
- _carry = null;
- }
- int startPos = 0;
- int pos = 0;
- while (pos < data.length) {
- int skip = 0;
- int char = data.codeUnitAt(pos);
- if (char == _LF) {
- skip = 1;
- } else if (char == _CR) {
- skip = 1;
- if (pos + 1 < data.length) {
- if (data.codeUnitAt(pos + 1) == _LF) {
- skip = 2;
- }
- } else if (!isClosing) {
- _carry = data.substring(startPos);
- return;
+ void handle(String data, StreamSink<String> sink, bool isClosing) {
floitsch 2013/02/26 14:47:26 make it private.
Søren Gjesse 2013/02/26 15:23:33 Done.
+ if (_carry != null) {
+ data = _carry.concat(data);
+ _carry = null;
+ }
+ int startPos = 0;
+ int pos = 0;
+ while (pos < data.length) {
+ int skip = 0;
+ int char = data.codeUnitAt(pos);
+ if (char == _LF) {
+ skip = 1;
+ } else if (char == _CR) {
+ skip = 1;
+ if (pos + 1 < data.length) {
+ if (data.codeUnitAt(pos + 1) == _LF) {
+ skip = 2;
}
- }
- if (skip > 0) {
- _buffer.add(data.substring(startPos, pos));
- _controller.add(_buffer.toString());
- _buffer.clear();
- startPos = pos = pos + skip;
- } else {
- pos++;
+ } else if (!isClosing) {
+ _carry = data.substring(startPos);
+ return;
}
}
- if (pos != startPos) {
- // Add remaining
+ if (skip > 0) {
_buffer.add(data.substring(startPos, pos));
- }
- if (isClosing && !_buffer.isEmpty) {
- _controller.add(_buffer.toString());
+ sink.add(_buffer.toString());
_buffer.clear();
+ startPos = pos = pos + skip;
+ } else {
+ pos++;
}
}
-
- _subscription = stream.listen(
- (data) => handle(data, false),
- onDone: () {
- // Handle remaining data (mainly _carry).
- handle("", true);
- _controller.close();
- },
- onError: _controller.signalError);
- return _controller.stream;
+ if (pos != startPos) {
+ // Add remaining
+ _buffer.add(data.substring(startPos, pos));
+ }
+ if (isClosing && !_buffer.isEmpty) {
+ sink.add(_buffer.toString());
+ _buffer.clear();
+ }
}
- void _pauseChanged() {
- if (_controller.isPaused) {
- _subscription.pause();
- } else {
- _subscription.resume();
- }
+ void handleData(String data, StreamSink<String> sink) {
+ handle(data, sink, false);
}
- void _subscriptionChanged() {
- if (!_controller.hasSubscribers) {
- _subscription.cancel();
- }
+ void handleDone(StreamSink<String> sink) {
+ handle("", sink, true);
floitsch 2013/02/26 14:47:26 sink.close()
Søren Gjesse 2013/02/26 15:23:33 Thanks.
}
}
-class _SingleByteDecoder implements StreamTransformer<List<int>, String> {
- StreamSubscription<List<int>> _subscription;
- StreamController<String> _controller;
+class _SingleByteDecoder extends StreamEventTransformer<List<int>, String> {
final int _replacementChar;
_SingleByteDecoder(this._replacementChar);
- Stream<String> bind(Stream<List<int>> stream) {
- _controller = new StreamController<String>(
- onPauseStateChange: _pauseChanged,
- onSubscriptionStateChange: _subscriptionChanged);
- _subscription = stream.listen(
- (data) {
- var buffer = new List<int>.fixedLength(data.length);
- for (int i = 0; i < data.length; i++) {
- int char = _decodeByte(data[i]);
- if (char < 0) char = _replacementChar;
- buffer[i] = char;
- }
- _controller.add(new String.fromCharCodes(buffer));
- },
- onDone: _controller.close,
- onError: _controller.signalError);
- return _controller.stream;
- }
-
- int _decodeByte(int byte);
-
- void _pauseChanged() {
- if (_controller.isPaused) {
- _subscription.pause();
- } else {
- _subscription.resume();
+ void handleData(List<int> data, StreamSink<String> sink) {
+ var buffer = new List<int>.fixedLength(data.length);
+ for (int i = 0; i < data.length; i++) {
+ int char = _decodeByte(data[i]);
+ if (char < 0) char = _replacementChar;
+ buffer[i] = char;
}
+ sink.add(new String.fromCharCodes(buffer));
}
- void _subscriptionChanged() {
- if (!_controller.hasSubscribers) {
- _subscription.cancel();
- }
- }
+ int _decodeByte(int byte);
}
@@ -299,46 +249,18 @@ class _Latin1Decoder extends _SingleByteDecoder {
}
-class _SingleByteEncoder implements StreamTransformer<String, List<int>> {
- StreamSubscription<String> _subscription;
- StreamController<List<int>> _controller;
-
- Stream<List<int>> bind(Stream<String> stream) {
- _controller = new StreamController<List<int>>(
- onPauseStateChange: _pauseChanged,
- onSubscriptionStateChange: _subscriptionChanged);
- _subscription = stream.listen(
- (string) {
- var bytes = _encode(string);
- if (bytes == null) {
- _controller.signalError(new FormatException(
- "Invalid character for encoding"));
- _controller.close();
- _subscription.cancel();
- } else {
- _controller.add(bytes);
- }
- },
- onDone: _controller.close,
- onError: _controller.signalError);
- return _controller.stream;
- }
-
- List<int> _encode(String string);
-
- void _pauseChanged() {
- if (_controller.isPaused) {
- _subscription.pause();
+class _SingleByteEncoder extends StreamEventTransformer<String, List<int>> {
+ void handleData(String data, StreamSink<List<int>> sink) {
+ var bytes = _encode(data);
+ if (bytes == null) {
+ throw new FormatException("Invalid character for encoding");
+ sink.close();
} else {
- _subscription.resume();
+ sink.add(bytes);
}
}
- void _subscriptionChanged() {
- if (!_controller.hasSubscribers) {
- _subscription.cancel();
- }
- }
+ List<int> _encode(String string);
}
@@ -379,36 +301,10 @@ class _WindowsCodePageEncoder extends _SingleByteEncoder {
// Utility class for decoding Windows current code page data delivered
// as a stream of bytes.
-class _WindowsCodePageDecoder implements StreamTransformer<List<int>, String> {
- StreamSubscription<List<int>> _subscription;
- StreamController<String> _controller;
-
- Stream<String> bind(Stream<List<int>> stream) {
- _controller = new StreamController<String>(
- onPauseStateChange: _pauseChanged,
- onSubscriptionStateChange: _subscriptionChanged);
- _subscription = stream.listen(
- (data) {
- _controller.add(_decodeBytes(data));
- },
- onDone: _controller.close,
- onError: _controller.signalError);
- return _controller.stream;
+class _WindowsCodePageDecoder extends StreamEventTransformer<List<int>, String> {
+ void handleData(List<int> data, StreamSink<String> sink) {
+ sink.add(_decodeBytes(data));
}
external static String _decodeBytes(List<int> bytes);
-
- void _pauseChanged() {
- if (_controller.isPaused) {
- _subscription.pause();
- } else {
- _subscription.resume();
- }
- }
-
- void _subscriptionChanged() {
- if (!_controller.hasSubscribers) {
- _subscription.cancel();
- }
- }
}

Powered by Google App Engine
This is Rietveld 408576698