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

Side by Side Diff: sdk/lib/async/stream_pipe.dart

Issue 12082047: Remove transformEvents and make StreamEventTransformer extend StreamTransformer. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Addressed review comments. 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 unified diff | Download patch | Annotate | Revision Log
« no previous file with comments | « sdk/lib/async/stream.dart ('k') | tests/lib/async/stream_controller_test.dart » ('j') | no next file with comments »
Toggle Intra-line Diffs ('i') | Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
OLDNEW
1 // Copyright (c) 2012, the Dart project authors. Please see the AUTHORS file 1 // Copyright (c) 2012, the Dart project authors. Please see the AUTHORS file
2 // for details. All rights reserved. Use of this source code is governed by a 2 // for details. All rights reserved. Use of this source code is governed by a
3 // BSD-style license that can be found in the LICENSE file. 3 // BSD-style license that can be found in the LICENSE file.
4 4
5 part of dart.async; 5 part of dart.async;
6 6
7 /** Utility function to create an [AsyncError] if [error] isn't one already. */ 7 /** Utility function to create an [AsyncError] if [error] isn't one already. */
8 AsyncError _asyncError(Object error, Object stackTrace, [AsyncError cause]) { 8 AsyncError _asyncError(Object error, Object stackTrace, [AsyncError cause]) {
9 if (error is AsyncError) return error; 9 if (error is AsyncError) return error;
10 if (cause == null) return new AsyncError(error, stackTrace); 10 if (cause == null) return new AsyncError(error, stackTrace);
(...skipping 465 matching lines...) Expand 10 before | Expand all | Expand 10 after
476 void _defaultHandleError(AsyncError error, StreamSink sink) { 476 void _defaultHandleError(AsyncError error, StreamSink sink) {
477 sink.signalError(error); 477 sink.signalError(error);
478 } 478 }
479 479
480 /** Default done handler forwards done. */ 480 /** Default done handler forwards done. */
481 void _defaultHandleDone(StreamSink sink) { 481 void _defaultHandleDone(StreamSink sink) {
482 sink.close(); 482 sink.close();
483 } 483 }
484 484
485 485
486 /**
487 * A stream transformer that intercepts all events and can generate any event as
488 * output.
489 *
490 * Each incoming event on the source stream is passed to the corresponding
491 * provided event handler, along with a [StreamSink] linked to the output
492 * Stream.
493 * The handler can then decide exactly which events to send to the output.
494 */
495 class _StreamTransformerImpl<S, T> implements StreamTransformer<S, T> {
496 final _TransformDataHandler<S, T> _onData;
497 final _TransformErrorHandler<T> _onError;
498 final _TransformDoneHandler<T> _onDone;
499 StreamSink<T> _sink;
500
501 _StreamTransformerImpl(void onData(S data, StreamSink<T> sink),
502 void onError(AsyncError data, StreamSink<T> sink),
503 void onDone(StreamSink<T> sink))
504 : this._onData = (onData == null ? _defaultHandleData : onData),
505 this._onError = (onError == null ? _defaultHandleError : onError),
506 this._onDone = (onDone == null ? _defaultHandleDone : onDone);
507
508 Stream<T> bind(Stream<S> source) {
509 Stream<T> stream = new _SingleStreamImpl<T>();
510 // Cache a Sink object to avoid creating a new one for each event.
511 _sink = new _StreamImplSink(stream);
512 source.listen(_handleData, onError: _handleError, onDone: _handleDone);
513 return stream;
514 }
515
516 void _handleData(S data) {
517 try {
518 _onData(data, _sink);
519 } catch (e, s) {
520 _sink.signalError(_asyncError(e, s));
521 }
522 }
523
524 void _handleError(AsyncError error) {
525 try {
526 _onError(error, _sink);
527 } catch (e, s) {
528 _sink.signalError(_asyncError(e, s, error));
529 }
530 }
531
532 void _handleDone() {
533 try {
534 _onDone(_sink);
535 } catch (e, s) {
536 _sink.signalError(_asyncError(e, s));
537 }
538 }
539 }
540
541 /** Creates a [StreamSink] from a [_StreamImpl]'s input methods. */ 486 /** Creates a [StreamSink] from a [_StreamImpl]'s input methods. */
542 class _StreamImplSink<T> implements StreamSink<T> { 487 class _StreamImplSink<T> implements StreamSink<T> {
543 _StreamImpl<T> _target; 488 _StreamImpl<T> _target;
544 _StreamImplSink(this._target); 489 _StreamImplSink(this._target);
545 void add(T data) { _target._add(data); } 490 void add(T data) { _target._add(data); }
546 void signalError(AsyncError error) { _target._signalError(error); } 491 void signalError(AsyncError error) { _target._signalError(error); }
547 void close() { _target._close(); } 492 void close() { _target._close(); }
548 } 493 }
549 494
550
551 /** 495 /**
552 * A stream transformer that intercepts all events and can generate any event as 496 * A [StreamTransformer] that modifies stream events.
553 * output.
554 * 497 *
555 * Each incoming event on the source stream is passed to the corresponding 498 * This class is used by [StreamTransformer]'s factory constructor.
556 * provided event handler, along with a [StreamSink] linked to the output 499 * It is actually an [StreamEventTransformer] where the functions used to
557 * Stream. 500 * modify the events are passed as constructor arguments.
558 * The handler can then decide exactly which events to send to the output. 501 *
502 * If an argument is omitted, it acts as the default method from
503 * [StreamEventTransformer].
559 */ 504 */
560 class _StreamEventTransformerImpl<S, T> 505 class _StreamTransformerImpl<S, T> extends StreamEventTransformer<S, T> {
561 implements StreamEventTransformer<S, T> {
562 // TODO(ahe): Restore type when feature is implemented in dart2js 506 // TODO(ahe): Restore type when feature is implemented in dart2js
563 // checked mode. http://dartbug.com/7733 507 // checked mode. http://dartbug.com/7733
564 final Function /*_TransformDataHandler<S, T>*/ _handleData; 508 final Function /*_TransformDataHandler<S, T>*/ _handleData;
565 final _TransformErrorHandler<T> _handleError; 509 final _TransformErrorHandler<T> _handleError;
566 final _TransformDoneHandler<T> _handleDone; 510 final _TransformDoneHandler<T> _handleDone;
567 511
568 _StreamEventTransformerImpl(void onData(S data, StreamSink<T> sink), 512 _StreamTransformerImpl(void handleData(S data, StreamSink<T> sink),
569 void onError(AsyncError data, StreamSink<T> sink), 513 void handleError(AsyncError data, StreamSink<T> sink),
570 void onDone(StreamSink<T> sink)) 514 void handleDone(StreamSink<T> sink))
571 : this._handleData = (onData == null ? _defaultHandleData : onData), 515 : this._handleData = (handleData == null ? _defaultHandleData
572 this._handleError = (onError == null ? _defaultHandleError : onError), 516 : handleData),
573 this._handleDone = (onDone == null ? _defaultHandleDone : onDone); 517 this._handleError = (handleError == null ? _defaultHandleError
518 : handleError),
519 this._handleDone = (handleDone == null ? _defaultHandleDone
520 : handleDone);
574 521
575 void handleData(S data, StreamSink<T> sink) { 522 void handleData(S data, StreamSink<T> sink) {
576 try { 523 _handleData(data, sink);
577 _handleData(data, sink);
578 } catch (e, s) {
579 sink.signalError(_asyncError(e, s));
580 }
581 } 524 }
582 525
583 void handleError(AsyncError error, StreamSink<T> sink) { 526 void handleError(AsyncError error, StreamSink<T> sink) {
584 try { 527 _handleError(error, sink);
585 _handleError(error, sink);
586 } catch (e, s) {
587 sink.signalError(_asyncError(e, s, error));
588 }
589 } 528 }
590 529
591 void handleDone(StreamSink<T> sink) { 530 void handleDone(StreamSink<T> sink) {
592 try { 531 _handleDone(sink);
593 _handleDone(sink);
594 } catch (e, s) {
595 sink.signalError(_asyncError(e, s));
596 }
597 } 532 }
598 } 533 }
534
OLDNEW
« no previous file with comments | « sdk/lib/async/stream.dart ('k') | tests/lib/async/stream_controller_test.dart » ('j') | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698