ForEachAsyncTest.cs 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596
  1. // Licensed to the .NET Foundation under one or more agreements.
  2. // The .NET Foundation licenses this file to you under the Apache 2.0 License.
  3. // See the LICENSE file in the project root for more information.
  4. using System;
  5. using System.Collections.Generic;
  6. using System.Reactive.Concurrency;
  7. using System.Reactive.Disposables;
  8. using System.Reactive.Linq;
  9. using System.Threading;
  10. using System.Threading.Tasks;
  11. using Microsoft.Reactive.Testing;
  12. using Xunit;
  13. namespace ReactiveTests.Tests
  14. {
  15. public class ForEachAsyncTest : ReactiveTest
  16. {
  17. [Fact]
  18. public void ForEachAsync_ArgumentChecking()
  19. {
  20. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.ForEachAsync(default(IObservable<int>), x => { }));
  21. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.ForEachAsync(Observable.Never<int>(), default(Action<int>)));
  22. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.ForEachAsync(default(IObservable<int>), x => { }, CancellationToken.None));
  23. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.ForEachAsync(Observable.Never<int>(), default(Action<int>), CancellationToken.None));
  24. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.ForEachAsync(default(IObservable<int>), (x, i) => { }));
  25. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.ForEachAsync(Observable.Never<int>(), default(Action<int, int>)));
  26. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.ForEachAsync(default(IObservable<int>), (x, i) => { }, CancellationToken.None));
  27. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.ForEachAsync(Observable.Never<int>(), default(Action<int, int>), CancellationToken.None));
  28. }
  29. [Fact]
  30. public void ForEachAsync_Never()
  31. {
  32. var scheduler = new TestScheduler();
  33. var xs = scheduler.CreateHotObservable(
  34. OnNext(100, 1),
  35. OnNext(200, 2),
  36. OnNext(300, 3),
  37. OnNext(400, 4),
  38. OnNext(500, 5)
  39. );
  40. var task = default(Task);
  41. var cts = new CancellationTokenSource();
  42. var list = new List<Recorded<int>>();
  43. scheduler.ScheduleAbsolute(150, () => task = xs.ForEachAsync(x => list.Add(new Recorded<int>(scheduler.Clock, x)), cts.Token));
  44. scheduler.Start();
  45. xs.Subscriptions.AssertEqual(
  46. Subscribe(150)
  47. );
  48. list.AssertEqual(
  49. new Recorded<int>(200, 2),
  50. new Recorded<int>(300, 3),
  51. new Recorded<int>(400, 4),
  52. new Recorded<int>(500, 5)
  53. );
  54. Assert.Equal(TaskStatus.WaitingForActivation, task.Status);
  55. }
  56. [Fact]
  57. public void ForEachAsync_Completed()
  58. {
  59. var scheduler = new TestScheduler();
  60. var xs = scheduler.CreateHotObservable(
  61. OnNext(100, 1),
  62. OnNext(200, 2),
  63. OnNext(300, 3),
  64. OnNext(400, 4),
  65. OnNext(500, 5),
  66. OnCompleted<int>(600)
  67. );
  68. var task = default(Task);
  69. var cts = new CancellationTokenSource();
  70. var list = new List<Recorded<int>>();
  71. scheduler.ScheduleAbsolute(150, () => task = xs.ForEachAsync(x => list.Add(new Recorded<int>(scheduler.Clock, x)), cts.Token));
  72. scheduler.Start();
  73. xs.Subscriptions.AssertEqual(
  74. Subscribe(150, 600)
  75. );
  76. list.AssertEqual(
  77. new Recorded<int>(200, 2),
  78. new Recorded<int>(300, 3),
  79. new Recorded<int>(400, 4),
  80. new Recorded<int>(500, 5)
  81. );
  82. Assert.Equal(TaskStatus.RanToCompletion, task.Status);
  83. }
  84. [Fact]
  85. public void ForEachAsync_Error()
  86. {
  87. var scheduler = new TestScheduler();
  88. var exception = new Exception();
  89. var xs = scheduler.CreateHotObservable(
  90. OnNext(100, 1),
  91. OnNext(200, 2),
  92. OnNext(300, 3),
  93. OnNext(400, 4),
  94. OnNext(500, 5),
  95. OnError<int>(600, exception)
  96. );
  97. var task = default(Task);
  98. var cts = new CancellationTokenSource();
  99. var list = new List<Recorded<int>>();
  100. scheduler.ScheduleAbsolute(150, () => task = xs.ForEachAsync(x => list.Add(new Recorded<int>(scheduler.Clock, x)), cts.Token));
  101. scheduler.Start();
  102. xs.Subscriptions.AssertEqual(
  103. Subscribe(150, 600)
  104. );
  105. list.AssertEqual(
  106. new Recorded<int>(200, 2),
  107. new Recorded<int>(300, 3),
  108. new Recorded<int>(400, 4),
  109. new Recorded<int>(500, 5)
  110. );
  111. Assert.Equal(TaskStatus.Faulted, task.Status);
  112. Assert.Same(exception, task.Exception.InnerException);
  113. }
  114. [Fact]
  115. public void ForEachAsync_Throw()
  116. {
  117. var scheduler = new TestScheduler();
  118. var exception = new Exception();
  119. var xs = scheduler.CreateHotObservable(
  120. OnNext(100, 1),
  121. OnNext(200, 2),
  122. OnNext(300, 3),
  123. OnNext(400, 4),
  124. OnNext(500, 5),
  125. OnCompleted<int>(600)
  126. );
  127. var task = default(Task);
  128. var cts = new CancellationTokenSource();
  129. var list = new List<Recorded<int>>();
  130. scheduler.ScheduleAbsolute(150, () => task = xs.ForEachAsync(x =>
  131. {
  132. if (scheduler.Clock > 400)
  133. {
  134. throw exception;
  135. }
  136. list.Add(new Recorded<int>(scheduler.Clock, x));
  137. }, cts.Token));
  138. scheduler.Start();
  139. xs.Subscriptions.AssertEqual(
  140. Subscribe(150, 500)
  141. );
  142. list.AssertEqual(
  143. new Recorded<int>(200, 2),
  144. new Recorded<int>(300, 3),
  145. new Recorded<int>(400, 4)
  146. );
  147. Assert.Equal(TaskStatus.Faulted, task.Status);
  148. Assert.Same(exception, task.Exception.InnerException);
  149. }
  150. [Fact]
  151. public void ForEachAsync_CancelDuring()
  152. {
  153. var scheduler = new TestScheduler();
  154. var xs = scheduler.CreateHotObservable(
  155. OnNext(100, 1),
  156. OnNext(200, 2),
  157. OnNext(300, 3),
  158. OnNext(400, 4),
  159. OnNext(500, 5),
  160. OnCompleted<int>(600)
  161. );
  162. var task = default(Task);
  163. var cts = new CancellationTokenSource();
  164. var list = new List<Recorded<int>>();
  165. scheduler.ScheduleAbsolute(150, () => task = xs.ForEachAsync(x => list.Add(new Recorded<int>(scheduler.Clock, x)), cts.Token));
  166. scheduler.ScheduleAbsolute(350, () => cts.Cancel());
  167. scheduler.Start();
  168. xs.Subscriptions.AssertEqual(
  169. Subscribe(150, 350)
  170. );
  171. list.AssertEqual(
  172. new Recorded<int>(200, 2),
  173. new Recorded<int>(300, 3)
  174. );
  175. Assert.Equal(TaskStatus.Canceled, task.Status);
  176. }
  177. [Fact]
  178. public void ForEachAsync_CancelBefore()
  179. {
  180. var scheduler = new TestScheduler();
  181. var xs = scheduler.CreateHotObservable(
  182. OnNext(100, 1),
  183. OnNext(200, 2),
  184. OnNext(300, 3),
  185. OnNext(400, 4),
  186. OnNext(500, 5),
  187. OnCompleted<int>(600)
  188. );
  189. var task = default(Task);
  190. var cts = new CancellationTokenSource();
  191. var list = new List<Recorded<int>>();
  192. cts.Cancel();
  193. scheduler.ScheduleAbsolute(150, () => task = xs.ForEachAsync(x => list.Add(new Recorded<int>(scheduler.Clock, x)), cts.Token));
  194. scheduler.Start();
  195. xs.Subscriptions.AssertEqual(
  196. );
  197. list.AssertEqual(
  198. );
  199. Assert.Equal(TaskStatus.Canceled, task.Status);
  200. }
  201. [Fact]
  202. public void ForEachAsync_CancelAfter()
  203. {
  204. var scheduler = new TestScheduler();
  205. var xs = scheduler.CreateHotObservable(
  206. OnNext(100, 1),
  207. OnNext(200, 2),
  208. OnNext(300, 3),
  209. OnNext(400, 4),
  210. OnNext(500, 5),
  211. OnCompleted<int>(600)
  212. );
  213. var task = default(Task);
  214. var cts = new CancellationTokenSource();
  215. var list = new List<Recorded<int>>();
  216. scheduler.ScheduleAbsolute(150, () => task = xs.ForEachAsync(x => list.Add(new Recorded<int>(scheduler.Clock, x)), cts.Token));
  217. scheduler.ScheduleAbsolute(700, () => cts.Cancel());
  218. scheduler.Start();
  219. xs.Subscriptions.AssertEqual(
  220. Subscribe(150, 600)
  221. );
  222. list.AssertEqual(
  223. new Recorded<int>(200, 2),
  224. new Recorded<int>(300, 3),
  225. new Recorded<int>(400, 4),
  226. new Recorded<int>(500, 5)
  227. );
  228. Assert.Equal(TaskStatus.RanToCompletion, task.Status);
  229. }
  230. [Fact]
  231. public void ForEachAsync_Default()
  232. {
  233. var list = new List<int>();
  234. Observable.Range(1, 10).ForEachAsync(list.Add).Wait();
  235. list.AssertEqual(1, 2, 3, 4, 5, 6, 7, 8, 9, 10);
  236. }
  237. [Fact]
  238. public void ForEachAsync_Index()
  239. {
  240. var list = new List<int>();
  241. Observable.Range(3, 5).ForEachAsync((x, i) => list.Add(x * i)).Wait();
  242. list.AssertEqual(3 * 0, 4 * 1, 5 * 2, 6 * 3, 7 * 4);
  243. }
  244. [Fact]
  245. public void ForEachAsync_Default_Cancel()
  246. {
  247. var N = 10;
  248. for (var n = 0; n < N; n++)
  249. {
  250. var cts = new CancellationTokenSource();
  251. var done = false;
  252. var xs = Observable.Create<int>(observer =>
  253. {
  254. return new CompositeDisposable(
  255. Observable.Repeat(42, Scheduler.Default).Subscribe(observer),
  256. Disposable.Create(() => done = true)
  257. );
  258. });
  259. var lst = new List<int>();
  260. var t = xs.ForEachAsync(
  261. x =>
  262. {
  263. lock (lst)
  264. {
  265. lst.Add(x);
  266. }
  267. },
  268. cts.Token
  269. );
  270. while (true)
  271. {
  272. lock (lst)
  273. {
  274. if (lst.Count >= 10)
  275. {
  276. break;
  277. }
  278. }
  279. }
  280. cts.Cancel();
  281. while (!t.IsCompleted)
  282. {
  283. ;
  284. }
  285. for (var i = 0; i < 10; i++)
  286. {
  287. Assert.Equal(42, lst[i]);
  288. }
  289. Assert.True(done);
  290. Assert.True(t.IsCanceled);
  291. }
  292. }
  293. [Fact]
  294. public void ForEachAsync_Index_Cancel()
  295. {
  296. var N = 10;
  297. for (var n = 0; n < N; n++)
  298. {
  299. var cts = new CancellationTokenSource();
  300. var done = false;
  301. var xs = Observable.Create<int>(observer =>
  302. {
  303. return new CompositeDisposable(
  304. Observable.Repeat(42, Scheduler.Default).Subscribe(observer),
  305. Disposable.Create(() => done = true)
  306. );
  307. });
  308. var lst = new List<int>();
  309. var t = xs.ForEachAsync(
  310. (x, i) =>
  311. {
  312. lock (lst)
  313. {
  314. lst.Add(x * i);
  315. }
  316. },
  317. cts.Token
  318. );
  319. while (true)
  320. {
  321. lock (lst)
  322. {
  323. if (lst.Count >= 10)
  324. {
  325. break;
  326. }
  327. }
  328. }
  329. cts.Cancel();
  330. while (!t.IsCompleted)
  331. {
  332. ;
  333. }
  334. for (var i = 0; i < 10; i++)
  335. {
  336. Assert.Equal(i * 42, lst[i]);
  337. }
  338. Assert.True(done);
  339. Assert.True(t.IsCanceled);
  340. }
  341. }
  342. [Fact]
  343. public void ForEachAsync_DisposeThrows1()
  344. {
  345. var cts = new CancellationTokenSource();
  346. var ex = new Exception();
  347. var xs = Observable.Create<int>(observer =>
  348. {
  349. return new CompositeDisposable(
  350. Observable.Range(0, 10, Scheduler.CurrentThread).Subscribe(observer),
  351. Disposable.Create(() => { throw ex; })
  352. );
  353. });
  354. var lst = new List<int>();
  355. var t = xs.ForEachAsync(lst.Add, cts.Token);
  356. //
  357. // Unfortunately, this doesn't throw for CurrentThread scheduling. The
  358. // subscription completes prior to assignment of the disposable, so we
  359. // succeed calling the TrySetResult method for the OnCompleted handler
  360. // prior to observing the exception of the Dispose operation, which is
  361. // surfacing upon assignment to the SingleAssignmentDisposable. As a
  362. // result, the exception evaporates.
  363. //
  364. // It'd be a breaking change at this point to rethrow the exception in
  365. // that case, so we're merely asserting regressions here.
  366. //
  367. try
  368. {
  369. t.Wait();
  370. }
  371. catch
  372. {
  373. Assert.True(false);
  374. }
  375. }
  376. [Fact]
  377. public void ForEachAsync_DisposeThrows2()
  378. {
  379. var cts = new CancellationTokenSource();
  380. var ex = new Exception();
  381. var xs = Observable.Create<int>(observer =>
  382. {
  383. return new CompositeDisposable(
  384. Observable.Range(0, 10, Scheduler.CurrentThread).Subscribe(observer),
  385. Disposable.Create(() => { throw ex; })
  386. );
  387. });
  388. var lst = new List<int>();
  389. var t = default(Task);
  390. Scheduler.CurrentThread.Schedule(() =>
  391. {
  392. t = xs.ForEachAsync(lst.Add, cts.Token);
  393. });
  394. //
  395. // If the trampoline of the CurrentThread has been installed higher
  396. // up the stack, the assignment of the subscription's disposable to
  397. // the SingleAssignmentDisposable can complete prior to the Dispose
  398. // method being called from the OnCompleted handler. In this case,
  399. // the OnCompleted handler's invocation of Dispose will cause the
  400. // exception to occur, and it bubbles out through TrySetException.
  401. //
  402. try
  403. {
  404. t.Wait();
  405. }
  406. catch (AggregateException err)
  407. {
  408. Assert.Equal(1, err.InnerExceptions.Count);
  409. Assert.Same(ex, err.InnerExceptions[0]);
  410. }
  411. }
  412. #if !NO_THREAD
  413. [Fact]
  414. [Trait("SkipCI", "true")]
  415. public void ForEachAsync_DisposeThrows()
  416. {
  417. //
  418. // Unfortunately, this test is non-deterministic due to the race
  419. // conditions described above in the tests using the CurrentThread
  420. // scheduler. The exception can come out through the OnCompleted
  421. // handler but can equally well get swallowed if the main thread
  422. // hasn't reached the assignment of the disposable yet, causing
  423. // the OnCompleted handler to win the race. The user can deal with
  424. // this by hooking an exception handler to the scheduler, so we
  425. // assert this behavior here.
  426. //
  427. // It'd be a breaking change at this point to change rethrowing
  428. // behavior, so we're merely asserting regressions here.
  429. //
  430. var hasCaughtEscapingException = 0;
  431. var cts = new CancellationTokenSource();
  432. var ex = new Exception();
  433. var s = Scheduler.Default.Catch<Exception>(err =>
  434. {
  435. Volatile.Write(ref hasCaughtEscapingException, 1);
  436. return ex == err;
  437. });
  438. while (Volatile.Read(ref hasCaughtEscapingException) == 0)
  439. {
  440. var xs = Observable.Create<int>(observer =>
  441. {
  442. return new CompositeDisposable(
  443. Observable.Range(0, 10, s).Subscribe(observer),
  444. Disposable.Create(() => { throw ex; })
  445. );
  446. });
  447. var lst = new List<int>();
  448. var t = xs.ForEachAsync(lst.Add, cts.Token);
  449. try
  450. {
  451. t.Wait();
  452. }
  453. catch (AggregateException err)
  454. {
  455. Assert.Equal(1, err.InnerExceptions.Count);
  456. Assert.Same(ex, err.InnerExceptions[0]);
  457. }
  458. }
  459. }
  460. [Fact]
  461. public void ForEachAsync_SubscribeThrows()
  462. {
  463. var ex = new Exception();
  464. var x = 42;
  465. var xs = Observable.Create<int>(observer =>
  466. {
  467. if (x == 42)
  468. {
  469. throw ex;
  470. }
  471. return Disposable.Empty;
  472. });
  473. var t = xs.ForEachAsync(_ => { });
  474. try
  475. {
  476. t.Wait();
  477. Assert.True(false);
  478. }
  479. catch (AggregateException err)
  480. {
  481. Assert.Equal(1, err.InnerExceptions.Count);
  482. Assert.Same(ex, err.InnerExceptions[0]);
  483. }
  484. }
  485. #endif
  486. }
  487. }