Chromium Code Reviews
chromiumcodereview-hr@appspot.gserviceaccount.com (chromiumcodereview-hr) | Please choose your nickname with Settings | Help | Chromium Project | Gerrit Changes | Sign out
(265)

Side by Side Diff: tests/lib/async/stream_transformation_broadcast_test.dart

Issue 920373003: Fix behavior when listening multiple times to a broadcast take/skip stream. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Fix typo Created 5 years, 10 months ago
Use n/p to move between diff chunks; N/P to move between comments. Draft comments are only viewable by you.
Jump to:
View unified diff | Download patch | Annotate | Revision Log
« no previous file with comments | « sdk/lib/async/stream_pipe.dart ('k') | no next file » | no next file with comments »
Toggle Intra-line Diffs ('i') | Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
OLDNEW
1 // Copyright (c) 2013, the Dart project authors. Please see the AUTHORS file 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 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 // Test that transformations like `map` and `where` preserve broadcast flag. 5 // Test that transformations like `map` and `where` preserve broadcast flag.
6 library stream_join_test; 6 library stream_join_test;
7 7
8 import 'dart:async'; 8 import 'dart:async';
9 import 'event_helper.dart'; 9 import 'event_helper.dart';
10 import 'package:unittest/unittest.dart'; 10 import 'package:unittest/unittest.dart';
11 import "package:expect/expect.dart"; 11 import "package:expect/expect.dart";
12 12
13 main() {
14 testStream("singlesub", () => new StreamController(), (c) => c.stream);
15 testStream("broadcast", () => new StreamController.broadcast(),
16 (c) => c.stream);
17 testStream("asBroadcast", () => new StreamController(),
18 (c) => c.stream.asBroadcastStream());
19 testStream("broadcast.asBroadcast", () => new StreamController.broadcast(),
20 (c) => c.stream.asBroadcastStream());
21 }
22
13 void testStream(String name, 23 void testStream(String name,
14 StreamController create(), 24 StreamController create(),
15 Stream getStream(controller)) { 25 Stream getStream(controller)) {
16 test("$name-map", () { 26 test("$name-map", () {
17 var c = create(); 27 var c = create();
18 var s = getStream(c); 28 var s = getStream(c);
19 Stream newStream = s.map((x) => x + 1); 29 Stream newStream = s.map((x) => x + 1);
20 Expect.equals(s.isBroadcast, newStream.isBroadcast); 30 Expect.equals(s.isBroadcast, newStream.isBroadcast);
21 newStream.single.then(expectAsync((v) { 31 newStream.single.then(expectAsync((v) {
22 Expect.equals(43, v); 32 Expect.equals(43, v);
(...skipping 141 matching lines...) Expand 10 before | Expand all | Expand 10 after
164 var c = create(); 174 var c = create();
165 var s = getStream(c); 175 var s = getStream(c);
166 Stream newStream = s.asyncExpand((x) => new Stream.fromIterable([x + 1])); 176 Stream newStream = s.asyncExpand((x) => new Stream.fromIterable([x + 1]));
167 Expect.equals(s.isBroadcast, newStream.isBroadcast); 177 Expect.equals(s.isBroadcast, newStream.isBroadcast);
168 newStream.single.then(expectAsync((v) { 178 newStream.single.then(expectAsync((v) {
169 Expect.equals(43, v); 179 Expect.equals(43, v);
170 })); 180 }));
171 c.add(42); 181 c.add(42);
172 c.close(); 182 c.close();
173 }); 183 });
184
185 // The following tests are only on broadcast streams, they require listening
186 // more than once.
187 if (name.startsWith("singlesub")) return;
188
189 test("$name-skip-multilisten", () {
190 if (name.startsWith("singlesub") ||
191 name.startsWith("asBroadcast")) return;
192 var c = create();
193 var s = getStream(c);
194 Stream newStream = s.skip(5);
195 // Listen immediately, to ensure that an asBroadcast stream is started.
196 var sub = newStream.listen((_){});
197 int i = 0;
198 var expect1 = 11;
199 var expect2 = 21;
200 var handler2 = expectAsync((v) {
201 expect(v, expect2);
202 expect2++;
203 }, count: 5);
204 var handler1 = expectAsync((v) {
205 expect(v, expect1);
206 expect1++;
207 }, count: 15);
208 var loop;
209 loop = expectAsync(() {
210 i++;
211 c.add(i);
212 if (i == 5) {
213 scheduleMicrotask(() {
214 newStream.listen(handler1);
215 });
216 }
217 if (i == 15) {
218 scheduleMicrotask(() {
219 newStream.listen(handler2);
220 });
221 }
222 if (i < 25) {
223 scheduleMicrotask(loop);
224 } else {
225 sub.cancel();
226 c.close();
227 }
228 }, count: 25);
229 scheduleMicrotask(loop);
230 });
231
232 test("$name-take-multilisten", () {
233 var c = create();
234 var s = getStream(c);
235 Stream newStream = s.take(10);
236 // Listen immediately, to ensure that an asBroadcast stream is started.
237 var sub = newStream.listen((_){});
238 int i = 0;
239 var expect1 = 6;
240 var expect2 = 11;
241 var handler2 = expectAsync((v) {
242 expect(v, expect2);
243 expect(v <= 20, isTrue);
244 expect2++;
245 }, count: 10);
246 var handler1 = expectAsync((v) {
247 expect(v, expect1);
248 expect(v <= 15, isTrue);
249 expect1++;
250 }, count: 10);
251 var loop;
252 loop = expectAsync(() {
253 i++;
254 c.add(i);
255 if (i == 5) {
256 scheduleMicrotask(() {
257 newStream.listen(handler1);
258 });
259 }
260 if (i == 10) {
261 scheduleMicrotask(() {
262 newStream.listen(handler2);
263 });
264 }
265 if (i < 25) {
266 scheduleMicrotask(loop);
267 } else {
268 sub.cancel();
269 c.close();
270 }
271 }, count: 25);
272 scheduleMicrotask(loop);
273 });
274
275 test("$name-skipWhile-multilisten", () {
276 if (name.startsWith("singlesub") ||
277 name.startsWith("asBroadcast")) return;
278 var c = create();
279 var s = getStream(c);
280 Stream newStream = s.skipWhile((x) => (x % 10) != 1);
281 // Listen immediately, to ensure that an asBroadcast stream is started.
282 var sub = newStream.listen((_){});
283 int i = 0;
284 var expect1 = 11;
285 var expect2 = 21;
286 var handler2 = expectAsync((v) {
287 expect(v, expect2);
288 expect2++;
289 }, count: 5);
290 var handler1 = expectAsync((v) {
291 expect(v, expect1);
292 expect1++;
293 }, count: 15);
294 var loop;
295 loop = expectAsync(() {
296 i++;
297 c.add(i);
298 if (i == 5) {
299 scheduleMicrotask(() {
300 newStream.listen(handler1);
301 });
302 }
303 if (i == 15) {
304 scheduleMicrotask(() {
305 newStream.listen(handler2);
306 });
307 }
308 if (i < 25) {
309 scheduleMicrotask(loop);
310 } else {
311 sub.cancel();
312 c.close();
313 }
314 }, count: 25);
315 scheduleMicrotask(loop);
316 });
317
318 test("$name-takeWhile-multilisten", () {
319 var c = create();
320 var s = getStream(c);
321 Stream newStream = s.takeWhile((x) => (x % 10) != 5);
322 // Listen immediately, to ensure that an asBroadcast stream is started.
323 var sub = newStream.listen((_){});
324 int i = 0;
325 // Non-overlapping ranges means the test must not remember its first
326 // failure.
327 var expect1 = 6;
328 var expect2 = 16;
329 var handler2 = expectAsync((v) {
330 expect(v, expect2);
331 expect(v <= 25, isTrue);
332 expect2++;
333 }, count: 9);
334 var handler1 = expectAsync((v) {
335 expect(v, expect1);
336 expect(v <= 15, isTrue);
337 expect1++;
338 }, count: 9);
339 var loop;
340 loop = expectAsync(() {
341 i++;
342 c.add(i);
343 if (i == 5) {
344 scheduleMicrotask(() {
345 newStream.listen(handler1);
346 });
347 }
348 if (i == 15) {
349 scheduleMicrotask(() {
350 newStream.listen(handler2);
351 });
352 }
353 if (i < 25) {
354 scheduleMicrotask(loop);
355 } else {
356 sub.cancel();
357 c.close();
358 }
359 }, count: 25);
360 scheduleMicrotask(loop);
361 });
174 } 362 }
175
176 main() {
177 testStream("singlesub", () => new StreamController(), (c) => c.stream);
178 testStream("broadcast", () => new StreamController.broadcast(),
179 (c) => c.stream);
180 testStream("asBroadcast", () => new StreamController(),
181 (c) => c.stream.asBroadcastStream());
182 testStream("broadcast.asBroadcast", () => new StreamController.broadcast(),
183 (c) => c.stream.asBroadcastStream());
184 }
OLDNEW
« no previous file with comments | « sdk/lib/async/stream_pipe.dart ('k') | no next file » | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698