| 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 class _BroadcastStream<T> extends _ControllerStream<T> { | 7 class _BroadcastStream<T> extends _ControllerStream<T> { |
| 8 _BroadcastStream(_StreamControllerLifecycle controller) : super(controller); | 8 _BroadcastStream(_StreamControllerLifecycle controller) : super(controller); |
| 9 | 9 |
| 10 bool get isBroadcast => true; | 10 bool get isBroadcast => true; |
| (...skipping 221 matching lines...) Expand 10 before | Expand all | Expand 10 after Loading... |
| 232 return new StateError("Cannot add new events while doing an addStream"); | 232 return new StateError("Cannot add new events while doing an addStream"); |
| 233 } | 233 } |
| 234 | 234 |
| 235 void add(T data) { | 235 void add(T data) { |
| 236 if (!_mayAddEvent) throw _addEventError(); | 236 if (!_mayAddEvent) throw _addEventError(); |
| 237 _sendData(data); | 237 _sendData(data); |
| 238 } | 238 } |
| 239 | 239 |
| 240 void addError(Object error, [StackTrace stackTrace]) { | 240 void addError(Object error, [StackTrace stackTrace]) { |
| 241 if (!_mayAddEvent) throw _addEventError(); | 241 if (!_mayAddEvent) throw _addEventError(); |
| 242 AsyncError replacement = Zone.current.errorCallback(error, stackTrace); |
| 243 if (replacement != null) { |
| 244 error = replacement.error; |
| 245 stackTrace = replacement.stackTrace; |
| 246 } |
| 242 _sendError(error, stackTrace); | 247 _sendError(error, stackTrace); |
| 243 } | 248 } |
| 244 | 249 |
| 245 Future close() { | 250 Future close() { |
| 246 if (isClosed) { | 251 if (isClosed) { |
| 247 assert(_doneFuture != null); | 252 assert(_doneFuture != null); |
| 248 return _doneFuture; | 253 return _doneFuture; |
| 249 } | 254 } |
| 250 if (!_mayAddEvent) throw _addEventError(); | 255 if (!_mayAddEvent) throw _addEventError(); |
| 251 _state |= _STATE_CLOSED; | 256 _state |= _STATE_CLOSED; |
| (...skipping 10 matching lines...) Expand all Loading... |
| 262 _addStreamState = new _AddStreamState(this, stream, cancelOnError); | 267 _addStreamState = new _AddStreamState(this, stream, cancelOnError); |
| 263 return _addStreamState.addStreamFuture; | 268 return _addStreamState.addStreamFuture; |
| 264 } | 269 } |
| 265 | 270 |
| 266 // _EventSink interface, called from AddStreamState. | 271 // _EventSink interface, called from AddStreamState. |
| 267 void _add(T data) { | 272 void _add(T data) { |
| 268 _sendData(data); | 273 _sendData(data); |
| 269 } | 274 } |
| 270 | 275 |
| 271 void _addError(Object error, StackTrace stackTrace) { | 276 void _addError(Object error, StackTrace stackTrace) { |
| 272 assert(_isAddingStream); | |
| 273 _sendError(error, stackTrace); | 277 _sendError(error, stackTrace); |
| 274 } | 278 } |
| 275 | 279 |
| 276 void _close() { | 280 void _close() { |
| 277 assert(_isAddingStream); | 281 assert(_isAddingStream); |
| 278 _AddStreamState addState = _addStreamState; | 282 _AddStreamState addState = _addStreamState; |
| 279 _addStreamState = null; | 283 _addStreamState = null; |
| 280 _state &= ~_STATE_ADDSTREAM; | 284 _state &= ~_STATE_ADDSTREAM; |
| 281 addState.complete(); | 285 addState.complete(); |
| 282 } | 286 } |
| (...skipping 169 matching lines...) Expand 10 before | Expand all | Expand 10 after Loading... |
| 452 while (_hasPending) { | 456 while (_hasPending) { |
| 453 _pending.handleNext(this); | 457 _pending.handleNext(this); |
| 454 } | 458 } |
| 455 } | 459 } |
| 456 | 460 |
| 457 void addError(Object error, [StackTrace stackTrace]) { | 461 void addError(Object error, [StackTrace stackTrace]) { |
| 458 if (!isClosed && _isFiring) { | 462 if (!isClosed && _isFiring) { |
| 459 _addPendingEvent(new _DelayedError(error, stackTrace)); | 463 _addPendingEvent(new _DelayedError(error, stackTrace)); |
| 460 return; | 464 return; |
| 461 } | 465 } |
| 462 super.addError(error, stackTrace); | 466 if (!_mayAddEvent) throw _addEventError(); |
| 467 _sendError(error, stackTrace); |
| 463 while (_hasPending) { | 468 while (_hasPending) { |
| 464 _pending.handleNext(this); | 469 _pending.handleNext(this); |
| 465 } | 470 } |
| 466 } | 471 } |
| 467 | 472 |
| 468 Future close() { | 473 Future close() { |
| 469 if (!isClosed && _isFiring) { | 474 if (!isClosed && _isFiring) { |
| 470 _addPendingEvent(const _DelayedDone()); | 475 _addPendingEvent(const _DelayedDone()); |
| 471 _state |= _BroadcastStreamController._STATE_CLOSED; | 476 _state |= _BroadcastStreamController._STATE_CLOSED; |
| 472 return super.done; | 477 return super.done; |
| (...skipping 24 matching lines...) Expand all Loading... |
| 497 _pauseCount++; | 502 _pauseCount++; |
| 498 } | 503 } |
| 499 void resume() { _resume(null); } | 504 void resume() { _resume(null); } |
| 500 void _resume(_) { | 505 void _resume(_) { |
| 501 if (_pauseCount > 0) _pauseCount--; | 506 if (_pauseCount > 0) _pauseCount--; |
| 502 } | 507 } |
| 503 Future cancel() { return new _Future.immediate(null); } | 508 Future cancel() { return new _Future.immediate(null); } |
| 504 bool get isPaused => _pauseCount > 0; | 509 bool get isPaused => _pauseCount > 0; |
| 505 Future asFuture([Object value]) => new _Future(); | 510 Future asFuture([Object value]) => new _Future(); |
| 506 } | 511 } |
| OLD | NEW |