| OLD | NEW |
| 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 Loading... |
| 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 |
| OLD | NEW |