| 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 |