| OLD | NEW |
| (Empty) | |
| 1 // Copyright (c) 2013, the Dart project authors. Please see the AUTHORS file |
| 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. |
| 4 |
| 5 library scheduled_test.scheduled_stream_test; |
| 6 |
| 7 import 'dart:async'; |
| 8 |
| 9 import 'package:scheduled_test/scheduled_stream.dart'; |
| 10 import 'package:unittest/unittest.dart'; |
| 11 import 'package:scheduled_test/src/utils.dart'; |
| 12 |
| 13 void main() { |
| 14 group("a completed stream with no elements", () { |
| 15 var stream; |
| 16 setUp(() { |
| 17 var controller = new StreamController()..close(); |
| 18 stream = new ScheduledStream(controller.stream); |
| 19 |
| 20 // Wait for [stream] to register the wrapped stream's close event. |
| 21 return pumpEventQueue(); |
| 22 }); |
| 23 |
| 24 test("throws an error for [next]", () { |
| 25 expect(stream.next(), throwsStateError); |
| 26 }); |
| 27 |
| 28 test("returns false for [hasNext]", () { |
| 29 expect(stream.hasNext, completion(isFalse)); |
| 30 }); |
| 31 |
| 32 test("emittedValues is empty", () { |
| 33 expect(stream.emittedValues, isEmpty); |
| 34 }); |
| 35 |
| 36 test("allValues is empty", () { |
| 37 expect(stream.allValues, isEmpty); |
| 38 }); |
| 39 |
| 40 test("close() does nothing", () { |
| 41 // This just shouldn't throw any exceptions. |
| 42 stream.close(); |
| 43 }); |
| 44 |
| 45 test("fork() returns a closed stream", () { |
| 46 var fork = stream.fork(); |
| 47 expect(fork.emittedValues, isEmpty); |
| 48 expect(fork.allValues, isEmpty); |
| 49 expect(fork.next(), throwsStateError); |
| 50 expect(fork.hasNext, completion(isFalse)); |
| 51 }); |
| 52 }); |
| 53 |
| 54 group("a stream that completes later with no elements", () { |
| 55 var controller; |
| 56 var stream; |
| 57 setUp(() { |
| 58 controller = new StreamController(); |
| 59 stream = new ScheduledStream(controller.stream); |
| 60 }); |
| 61 |
| 62 test("throws an error for [next]", () { |
| 63 expect(stream.next(), throwsStateError); |
| 64 return new Future(controller.close); |
| 65 }); |
| 66 |
| 67 test("returns false for [hasNext]", () { |
| 68 expect(stream.hasNext, completion(isFalse)); |
| 69 return new Future(controller.close); |
| 70 }); |
| 71 |
| 72 test("emittedValues is empty", () { |
| 73 expect(stream.emittedValues, isEmpty); |
| 74 }); |
| 75 |
| 76 test("allValues is empty", () { |
| 77 expect(stream.allValues, isEmpty); |
| 78 }); |
| 79 |
| 80 test("fork() returns a stream that closes when the controller is closed", |
| 81 () { |
| 82 var fork = stream.fork(); |
| 83 expect(fork.emittedValues, isEmpty); |
| 84 expect(fork.allValues, isEmpty); |
| 85 |
| 86 var nextComplete = false; |
| 87 expect(fork.next().whenComplete(() { |
| 88 nextComplete = true; |
| 89 }), throwsStateError); |
| 90 |
| 91 var hasNextComplete = false; |
| 92 expect(fork.hasNext.whenComplete(() { |
| 93 hasNextComplete = true; |
| 94 }), completion(isFalse)); |
| 95 |
| 96 // Pump the event queue to give [next] and [hasNext] a chance to |
| 97 // (incorrectly) fire. |
| 98 return pumpEventQueue().then((_) { |
| 99 expect(nextComplete, isFalse); |
| 100 expect(hasNextComplete, isFalse); |
| 101 |
| 102 controller.close(); |
| 103 }); |
| 104 }); |
| 105 }); |
| 106 |
| 107 test("forking and then closing a stream closes the fork", () { |
| 108 var stream = new ScheduledStream(new StreamController().stream); |
| 109 var fork = stream.fork(); |
| 110 expect(fork.emittedValues, isEmpty); |
| 111 expect(fork.allValues, isEmpty); |
| 112 |
| 113 var nextComplete = false; |
| 114 expect(fork.next().whenComplete(() { |
| 115 nextComplete = true; |
| 116 }), throwsStateError); |
| 117 |
| 118 var hasNextComplete = false; |
| 119 expect(fork.hasNext.whenComplete(() { |
| 120 hasNextComplete = true; |
| 121 }), completion(isFalse)); |
| 122 |
| 123 // Pump the event queue to give [next] and [hasNext] a chance to |
| 124 // (incorrectly) fire. |
| 125 return pumpEventQueue().then((_) { |
| 126 expect(nextComplete, isFalse); |
| 127 expect(hasNextComplete, isFalse); |
| 128 |
| 129 stream.close(); |
| 130 }); |
| 131 }); |
| 132 |
| 133 group("a completed stream with several values", () { |
| 134 var stream; |
| 135 setUp(() { |
| 136 var controller = new StreamController<int>() |
| 137 ..add(1)..add(2)..add(3)..close(); |
| 138 stream = new ScheduledStream<int>(controller.stream); |
| 139 |
| 140 return pumpEventQueue(); |
| 141 }); |
| 142 |
| 143 test("next() returns each value then throws an error", () { |
| 144 return stream.next().then((value) { |
| 145 expect(value, equals(1)); |
| 146 return stream.next(); |
| 147 }).then((value) { |
| 148 expect(value, equals(2)); |
| 149 return stream.next(); |
| 150 }).then((value) { |
| 151 expect(value, equals(3)); |
| 152 expect(stream.next(), throwsStateError); |
| 153 }); |
| 154 }); |
| 155 |
| 156 test("parallel next() calls are disallowed", () { |
| 157 expect(stream.next(), completion(equals(1))); |
| 158 expect(stream.next(), throwsStateError); |
| 159 }); |
| 160 |
| 161 test("parallel hasNext calls are allowed", () { |
| 162 expect(stream.hasNext, completion(isTrue)); |
| 163 expect(stream.hasNext, completion(isTrue)); |
| 164 }); |
| 165 |
| 166 test("hasNext returns true until there are no more values", () { |
| 167 return stream.hasNext.then((hasNext) { |
| 168 expect(hasNext, isTrue); |
| 169 return stream.next(); |
| 170 }).then((_) => stream.hasNext).then((hasNext) { |
| 171 expect(hasNext, isTrue); |
| 172 return stream.next(); |
| 173 }).then((_) => stream.hasNext).then((hasNext) { |
| 174 expect(hasNext, isTrue); |
| 175 return stream.next(); |
| 176 }).then((_) => expect(stream.hasNext, completion(isFalse))); |
| 177 }); |
| 178 |
| 179 test("emittedValues returns the values that have been emitted", () { |
| 180 expect(stream.emittedValues, isEmpty); |
| 181 |
| 182 return stream.next().then((_) { |
| 183 expect(stream.emittedValues, equals([1])); |
| 184 return stream.next(); |
| 185 }).then((_) { |
| 186 expect(stream.emittedValues, equals([1, 2])); |
| 187 return stream.next(); |
| 188 }).then((_) { |
| 189 expect(stream.emittedValues, equals([1, 2, 3])); |
| 190 }); |
| 191 }); |
| 192 |
| 193 test("allValues returns all values that the inner stream emitted", () { |
| 194 expect(stream.allValues, equals([1, 2, 3])); |
| 195 }); |
| 196 |
| 197 test("closing the stream means it doesn't emit additional events", () { |
| 198 return stream.next().then((_) { |
| 199 stream.close(); |
| 200 expect(stream.next(), throwsStateError); |
| 201 }); |
| 202 }); |
| 203 |
| 204 test("a fork created before any values are emitted emits all values", () { |
| 205 var fork = stream.fork(); |
| 206 return fork.next().then((value) { |
| 207 expect(value, equals(1)); |
| 208 return fork.next(); |
| 209 }).then((value) { |
| 210 expect(value, equals(2)); |
| 211 return fork.next(); |
| 212 }).then((value) { |
| 213 expect(value, equals(3)); |
| 214 expect(fork.next(), throwsStateError); |
| 215 }); |
| 216 }); |
| 217 |
| 218 test("a fork created after some values are emitted emits remaining values", |
| 219 () { |
| 220 var fork; |
| 221 return stream.next().then((_) { |
| 222 fork = stream.fork(); |
| 223 return fork.next(); |
| 224 }).then((value) { |
| 225 expect(value, equals(2)); |
| 226 return fork.next(); |
| 227 }).then((value) { |
| 228 expect(value, equals(3)); |
| 229 expect(fork.next(), throwsStateError); |
| 230 }); |
| 231 }); |
| 232 |
| 233 test("a fork doesn't push forward its parent stream", () { |
| 234 var fork = stream.fork(); |
| 235 return fork.next().then((_) { |
| 236 expect(stream.next(), completion(equals(1))); |
| 237 }); |
| 238 }); |
| 239 |
| 240 test("closing a fork doesn't close its parent stream", () { |
| 241 var fork = stream.fork(); |
| 242 fork.close(); |
| 243 expect(stream.next(), completion(equals(1))); |
| 244 }); |
| 245 |
| 246 test("closing a stream closes its forks immediately", () { |
| 247 var fork = stream.fork(); |
| 248 return stream.next().then((_) { |
| 249 stream.close(); |
| 250 expect(fork.next(), throwsStateError); |
| 251 }); |
| 252 }); |
| 253 }); |
| 254 |
| 255 group("a stream with several values added asynchronously", () { |
| 256 var stream; |
| 257 setUp(() { |
| 258 var controller = new StreamController<int>(); |
| 259 stream = new ScheduledStream<int>(controller.stream); |
| 260 |
| 261 pumpEventQueue().then((_) { |
| 262 controller.add(1); |
| 263 return pumpEventQueue(); |
| 264 }).then((_) { |
| 265 controller.add(2); |
| 266 return pumpEventQueue(); |
| 267 }).then((_) { |
| 268 controller.add(3); |
| 269 return pumpEventQueue(); |
| 270 }).then((_) { |
| 271 controller.close(); |
| 272 }); |
| 273 }); |
| 274 |
| 275 test("next() returns each value then throws an error", () { |
| 276 return stream.next().then((value) { |
| 277 expect(value, equals(1)); |
| 278 return stream.next(); |
| 279 }).then((value) { |
| 280 expect(value, equals(2)); |
| 281 return stream.next(); |
| 282 }).then((value) { |
| 283 expect(value, equals(3)); |
| 284 expect(stream.next(), throwsStateError); |
| 285 }); |
| 286 }); |
| 287 |
| 288 test("parallel next() calls are disallowed", () { |
| 289 expect(stream.next(), completion(equals(1))); |
| 290 expect(stream.next(), throwsStateError); |
| 291 }); |
| 292 |
| 293 test("parallel hasNext calls are allowed", () { |
| 294 expect(stream.hasNext, completion(isTrue)); |
| 295 expect(stream.hasNext, completion(isTrue)); |
| 296 }); |
| 297 |
| 298 test("hasNext returns true until there are no more values", () { |
| 299 return stream.hasNext.then((hasNext) { |
| 300 expect(hasNext, isTrue); |
| 301 return stream.next(); |
| 302 }).then((_) => stream.hasNext).then((hasNext) { |
| 303 expect(hasNext, isTrue); |
| 304 return stream.next(); |
| 305 }).then((_) => stream.hasNext).then((hasNext) { |
| 306 expect(hasNext, isTrue); |
| 307 return stream.next(); |
| 308 }).then((_) => expect(stream.hasNext, completion(isFalse))); |
| 309 }); |
| 310 |
| 311 test("emittedValues returns the values that have been emitted", () { |
| 312 expect(stream.emittedValues, isEmpty); |
| 313 |
| 314 return stream.next().then((_) { |
| 315 expect(stream.emittedValues, equals([1])); |
| 316 return stream.next(); |
| 317 }).then((_) { |
| 318 expect(stream.emittedValues, equals([1, 2])); |
| 319 return stream.next(); |
| 320 }).then((_) { |
| 321 expect(stream.emittedValues, equals([1, 2, 3])); |
| 322 }); |
| 323 }); |
| 324 |
| 325 test("allValues returns all values that the inner stream emitted", () { |
| 326 expect(stream.allValues, isEmpty); |
| 327 |
| 328 return stream.next().then((_) { |
| 329 expect(stream.allValues, equals([1])); |
| 330 return stream.next(); |
| 331 }).then((_) { |
| 332 expect(stream.allValues, equals([1, 2])); |
| 333 return stream.next(); |
| 334 }).then((_) { |
| 335 expect(stream.allValues, equals([1, 2, 3])); |
| 336 }); |
| 337 }); |
| 338 |
| 339 test("closing the stream means it doesn't emit additional events", () { |
| 340 return stream.next().then((_) { |
| 341 stream.close(); |
| 342 expect(stream.next(), throwsStateError); |
| 343 }); |
| 344 }); |
| 345 |
| 346 test("a fork created before any values are emitted emits all values", () { |
| 347 var fork = stream.fork(); |
| 348 return fork.next().then((value) { |
| 349 expect(value, equals(1)); |
| 350 return fork.next(); |
| 351 }).then((value) { |
| 352 expect(value, equals(2)); |
| 353 return fork.next(); |
| 354 }).then((value) { |
| 355 expect(value, equals(3)); |
| 356 expect(fork.next(), throwsStateError); |
| 357 }); |
| 358 }); |
| 359 |
| 360 test("a fork created after some values are emitted emits remaining values", |
| 361 () { |
| 362 var fork; |
| 363 return stream.next().then((_) { |
| 364 fork = stream.fork(); |
| 365 return fork.next(); |
| 366 }).then((value) { |
| 367 expect(value, equals(2)); |
| 368 return fork.next(); |
| 369 }).then((value) { |
| 370 expect(value, equals(3)); |
| 371 expect(fork.next(), throwsStateError); |
| 372 }); |
| 373 }); |
| 374 |
| 375 test("a fork doesn't push forward its parent stream", () { |
| 376 var fork = stream.fork(); |
| 377 return fork.next().then((_) { |
| 378 expect(stream.next(), completion(equals(1))); |
| 379 }); |
| 380 }); |
| 381 |
| 382 test("closing a fork doesn't close its parent stream", () { |
| 383 var fork = stream.fork(); |
| 384 fork.close(); |
| 385 expect(stream.next(), completion(equals(1))); |
| 386 }); |
| 387 |
| 388 test("closing a stream closes its forks immediately", () { |
| 389 var fork = stream.fork(); |
| 390 return stream.next().then((_) { |
| 391 stream.close(); |
| 392 expect(fork.next(), throwsStateError); |
| 393 }); |
| 394 }); |
| 395 }); |
| 396 } |
| OLD | NEW |