| 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 class _BaseDataInputStream { | 5 class _BaseDataInputStream { |
| 6 abstract int available(); | 6 abstract int available(); |
| 7 | 7 |
| 8 List<int> read([int len]) { | 8 List<int> read([int len]) { |
| 9 if (_closeCallbackCalled) return null; | 9 if (_closeCallbackCalled) return null; |
| 10 int bytesToRead = available(); | 10 int bytesToRead = available(); |
| (...skipping 18 matching lines...) Expand all Loading... |
| 29 if (len < 0) throw new StreamException("Illegal length $len"); | 29 if (len < 0) throw new StreamException("Illegal length $len"); |
| 30 int bytesToRead = Math.min(len, available()); | 30 int bytesToRead = Math.min(len, available()); |
| 31 return _readInto(buffer, offset, bytesToRead); | 31 return _readInto(buffer, offset, bytesToRead); |
| 32 } | 32 } |
| 33 | 33 |
| 34 void pipe(OutputStream output, [bool close = true]) { | 34 void pipe(OutputStream output, [bool close = true]) { |
| 35 _pipe(this, output, close: close); | 35 _pipe(this, output, close: close); |
| 36 } | 36 } |
| 37 | 37 |
| 38 void close() { | 38 void close() { |
| 39 if (_scheduledDataCallback != null) { | 39 _cancelScheduledDataCallback(); |
| 40 _scheduledDataCallback.cancel(); | |
| 41 } | |
| 42 _close(); | 40 _close(); |
| 43 _checkScheduleCallbacks(); | 41 _checkScheduleCallbacks(); |
| 44 } | 42 } |
| 45 | 43 |
| 46 bool get closed() => _closeCallbackCalled; | 44 bool get closed() => _closeCallbackCalled; |
| 47 | 45 |
| 48 void set dataHandler(void callback()) { | 46 void set dataHandler(void callback()) { |
| 49 _clientDataHandler = callback; | 47 _clientDataHandler = callback; |
| 50 _checkScheduleCallbacks(); | 48 _checkScheduleCallbacks(); |
| 51 } | 49 } |
| 52 | 50 |
| 53 void set closeHandler(void callback()) { | 51 void set closeHandler(void callback()) { |
| 54 _clientCloseHandler = callback; | 52 _clientCloseHandler = callback; |
| 55 _checkScheduleCallbacks(); | 53 _checkScheduleCallbacks(); |
| 56 } | 54 } |
| 57 | 55 |
| 58 void set errorHandler(void callback()) { | 56 void set errorHandler(void callback()) { |
| 59 _clientErrorHandler = callback; | 57 _clientErrorHandler = callback; |
| 60 } | 58 } |
| 61 | 59 |
| 62 abstract List<int> _read(int bytesToRead); | 60 abstract List<int> _read(int bytesToRead); |
| 63 | 61 |
| 62 void _dataReceived() { |
| 63 // More data has been received asynchronously. Perform the data |
| 64 // handler callback now. |
| 65 _cancelScheduledDataCallback(); |
| 66 if (_clientDataHandler !== null) { |
| 67 _clientDataHandler(); |
| 68 } |
| 69 _checkScheduleCallbacks(); |
| 70 } |
| 71 |
| 72 void _closeReceived() { |
| 73 // Close indication has been received asynchronously. Perform the |
| 74 // close callback now if all data has been delivered. |
| 75 _streamMarkedClosed = true; |
| 76 if (available() == 0) { |
| 77 if (_clientCloseHandler !== null) { |
| 78 _clientCloseHandler(); |
| 79 _closeCallbackCalled = true; |
| 80 } |
| 81 } else { |
| 82 _checkScheduleCallbacks(); |
| 83 } |
| 84 } |
| 85 |
| 86 void _cancelScheduledDataCallback() { |
| 87 if (_scheduledDataCallback != null) { |
| 88 _scheduledDataCallback.cancel(); |
| 89 _scheduledDataCallback = null; |
| 90 } |
| 91 } |
| 92 |
| 64 void _checkScheduleCallbacks() { | 93 void _checkScheduleCallbacks() { |
| 65 void issueDataCallback(Timer timer) { | 94 void issueDataCallback(Timer timer) { |
| 66 _scheduledDataCallback = null; | 95 _scheduledDataCallback = null; |
| 67 if (_clientDataHandler !== null) { | 96 if (_clientDataHandler !== null) { |
| 68 _clientDataHandler(); | 97 _clientDataHandler(); |
| 69 _checkScheduleCallbacks(); | 98 _checkScheduleCallbacks(); |
| 70 } | 99 } |
| 71 } | 100 } |
| 72 | 101 |
| 73 void issueCloseCallback(Timer timer) { | 102 void issueCloseCallback(Timer timer) { |
| 74 _scheduledCloseCallback = null; | 103 _scheduledCloseCallback = null; |
| 75 if (_clientCloseHandler !== null) _clientCloseHandler(); | 104 if (_clientCloseHandler !== null) _clientCloseHandler(); |
| 76 } | 105 } |
| 77 | 106 |
| 78 // Schedule data callback if there is more data to read. Schedule | 107 // Schedule data callback if there is more data to read. Schedule |
| 79 // close callback once when all data has been read. Only schedule | 108 // close callback once when all data has been read. Only schedule |
| 80 // a new callback if the previous one has actually been called. | 109 // a new callback if the previous one has actually been called. |
| 81 if (!_closeCallbackCalled) { | 110 if (!_closeCallbackCalled) { |
| 82 if (available() > 0) { | 111 if (available() > 0) { |
| 83 if (_scheduledDataCallback == null) { | 112 if (_scheduledDataCallback == null) { |
| 84 _scheduledDataCallback = new Timer(issueDataCallback, 0); | 113 _scheduledDataCallback = new Timer(issueDataCallback, 0); |
| 85 } | 114 } |
| 86 } else if (_streamMarkedClosed && !_closeCallbackCalled) { | 115 } else if (_streamMarkedClosed && !_closeCallbackCalled) { |
| 87 if (_scheduledDataCallback != null) { | 116 _cancelScheduledDataCallback(); |
| 88 _scheduledDataCallback.cancel(); | |
| 89 } | |
| 90 _close(); | 117 _close(); |
| 91 _scheduledCloseCallback = new Timer(issueCloseCallback, 0); | 118 _scheduledCloseCallback = new Timer(issueCloseCallback, 0); |
| 92 _closeCallbackCalled = true; | 119 _closeCallbackCalled = true; |
| 93 } | 120 } |
| 94 } | 121 } |
| 95 } | 122 } |
| 96 | 123 |
| 97 // When this is set to true the stream is marked closed. When a | 124 // When this is set to true the stream is marked closed. When a |
| 98 // stream is marked closed no more data can arrive and the value | 125 // stream is marked closed no more data can arrive and the value |
| 99 // from available is now all remaining data. If this is true and the | 126 // from available is now all remaining data. If this is true and the |
| (...skipping 41 matching lines...) Expand 10 before | Expand all | Expand 10 after Loading... |
| 141 input.dataHandler = pipeDataHandler; | 168 input.dataHandler = pipeDataHandler; |
| 142 output.noPendingWriteHandler = null; | 169 output.noPendingWriteHandler = null; |
| 143 }; | 170 }; |
| 144 | 171 |
| 145 _inputCloseHandler = input._clientCloseHandler; | 172 _inputCloseHandler = input._clientCloseHandler; |
| 146 input.dataHandler = pipeDataHandler; | 173 input.dataHandler = pipeDataHandler; |
| 147 input.closeHandler = pipeCloseHandler; | 174 input.closeHandler = pipeCloseHandler; |
| 148 output.noPendingWriteHandler = null; | 175 output.noPendingWriteHandler = null; |
| 149 } | 176 } |
| 150 | 177 |
| OLD | NEW |