ForEachAsyncTest.cs 18 KB

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