SelectMany.cs 63 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705
  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.Collections.Generic;
  5. using System.Reactive.Disposables;
  6. using System.Threading;
  7. using System.Threading.Tasks;
  8. namespace System.Reactive.Linq.ObservableImpl
  9. {
  10. internal static class SelectMany<TSource, TCollection, TResult>
  11. {
  12. internal sealed class ObservableSelector : Producer<TResult, ObservableSelector._>
  13. {
  14. private readonly IObservable<TSource> _source;
  15. private readonly Func<TSource, IObservable<TCollection>> _collectionSelector;
  16. private readonly Func<TSource, TCollection, TResult> _resultSelector;
  17. public ObservableSelector(IObservable<TSource> source, Func<TSource, IObservable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  18. {
  19. _source = source;
  20. _collectionSelector = collectionSelector;
  21. _resultSelector = resultSelector;
  22. }
  23. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  24. protected override void Run(_ sink) => sink.Run(_source);
  25. internal sealed class _ : Sink<TSource, TResult>
  26. {
  27. private readonly object _gate = new object();
  28. private readonly SingleAssignmentDisposable _sourceSubscription = new SingleAssignmentDisposable();
  29. private readonly CompositeDisposable _group = new CompositeDisposable();
  30. private readonly Func<TSource, IObservable<TCollection>> _collectionSelector;
  31. private readonly Func<TSource, TCollection, TResult> _resultSelector;
  32. public _(ObservableSelector parent, IObserver<TResult> observer)
  33. : base(observer)
  34. {
  35. _collectionSelector = parent._collectionSelector;
  36. _resultSelector = parent._resultSelector;
  37. _group.Add(_sourceSubscription);
  38. }
  39. private bool _isStopped;
  40. public void Run(IObservable<TSource> source)
  41. {
  42. _isStopped = false;
  43. _sourceSubscription.Disposable = source.SubscribeSafe(this);
  44. SetUpstream(_group);
  45. }
  46. public override void OnNext(TSource value)
  47. {
  48. var collection = default(IObservable<TCollection>);
  49. try
  50. {
  51. collection = _collectionSelector(value);
  52. }
  53. catch (Exception ex)
  54. {
  55. lock (_gate)
  56. {
  57. ForwardOnError(ex);
  58. }
  59. return;
  60. }
  61. var innerSubscription = new SingleAssignmentDisposable();
  62. _group.Add(innerSubscription);
  63. innerSubscription.Disposable = collection.SubscribeSafe(new InnerObserver(this, value, innerSubscription));
  64. }
  65. public override void OnError(Exception error)
  66. {
  67. lock (_gate)
  68. {
  69. ForwardOnError(error);
  70. }
  71. }
  72. public override void OnCompleted()
  73. {
  74. _isStopped = true;
  75. if (_group.Count == 1)
  76. {
  77. //
  78. // Notice there can be a race between OnCompleted of the source and any
  79. // of the inner sequences, where both see _group.Count == 1, and one is
  80. // waiting for the lock. There won't be a double OnCompleted observation
  81. // though, because the call to Dispose silences the observer by swapping
  82. // in a NopObserver<T>.
  83. //
  84. lock (_gate)
  85. {
  86. ForwardOnCompleted();
  87. }
  88. }
  89. else
  90. {
  91. _sourceSubscription.Dispose();
  92. }
  93. }
  94. private sealed class InnerObserver : IObserver<TCollection>
  95. {
  96. private readonly _ _parent;
  97. private readonly TSource _value;
  98. private readonly IDisposable _self;
  99. public InnerObserver(_ parent, TSource value, IDisposable self)
  100. {
  101. _parent = parent;
  102. _value = value;
  103. _self = self;
  104. }
  105. public void OnNext(TCollection value)
  106. {
  107. var res = default(TResult);
  108. try
  109. {
  110. res = _parent._resultSelector(_value, value);
  111. }
  112. catch (Exception ex)
  113. {
  114. lock (_parent._gate)
  115. {
  116. _parent.ForwardOnError(ex);
  117. }
  118. return;
  119. }
  120. lock (_parent._gate)
  121. _parent.ForwardOnNext(res);
  122. }
  123. public void OnError(Exception error)
  124. {
  125. lock (_parent._gate)
  126. {
  127. _parent.ForwardOnError(error);
  128. }
  129. }
  130. public void OnCompleted()
  131. {
  132. _parent._group.Remove(_self);
  133. if (_parent._isStopped && _parent._group.Count == 1)
  134. {
  135. //
  136. // Notice there can be a race between OnCompleted of the source and any
  137. // of the inner sequences, where both see _group.Count == 1, and one is
  138. // waiting for the lock. There won't be a double OnCompleted observation
  139. // though, because the call to Dispose silences the observer by swapping
  140. // in a NopObserver<T>.
  141. //
  142. lock (_parent._gate)
  143. {
  144. _parent.ForwardOnCompleted();
  145. }
  146. }
  147. }
  148. }
  149. }
  150. }
  151. internal sealed class ObservableSelectorIndexed : Producer<TResult, ObservableSelectorIndexed._>
  152. {
  153. private readonly IObservable<TSource> _source;
  154. private readonly Func<TSource, int, IObservable<TCollection>> _collectionSelector;
  155. private readonly Func<TSource, int, TCollection, int, TResult> _resultSelector;
  156. public ObservableSelectorIndexed(IObservable<TSource> source, Func<TSource, int, IObservable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  157. {
  158. _source = source;
  159. _collectionSelector = collectionSelector;
  160. _resultSelector = resultSelector;
  161. }
  162. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  163. protected override void Run(_ sink) => sink.Run(_source);
  164. internal sealed class _ : Sink<TSource, TResult>
  165. {
  166. private readonly object _gate = new object();
  167. private readonly SingleAssignmentDisposable _sourceSubscription = new SingleAssignmentDisposable();
  168. private readonly CompositeDisposable _group = new CompositeDisposable();
  169. private readonly Func<TSource, int, IObservable<TCollection>> _collectionSelector;
  170. private readonly Func<TSource, int, TCollection, int, TResult> _resultSelector;
  171. public _(ObservableSelectorIndexed parent, IObserver<TResult> observer)
  172. : base(observer)
  173. {
  174. _collectionSelector = parent._collectionSelector;
  175. _resultSelector = parent._resultSelector;
  176. _group.Add(_sourceSubscription);
  177. }
  178. private bool _isStopped;
  179. private int _index;
  180. public void Run(IObservable<TSource> source)
  181. {
  182. _isStopped = false;
  183. _sourceSubscription.Disposable = source.SubscribeSafe(this);
  184. SetUpstream(_group);
  185. }
  186. public override void OnNext(TSource value)
  187. {
  188. var index = checked(_index++);
  189. var collection = default(IObservable<TCollection>);
  190. try
  191. {
  192. collection = _collectionSelector(value, index);
  193. }
  194. catch (Exception ex)
  195. {
  196. lock (_gate)
  197. {
  198. ForwardOnError(ex);
  199. }
  200. return;
  201. }
  202. var innerSubscription = new SingleAssignmentDisposable();
  203. _group.Add(innerSubscription);
  204. innerSubscription.Disposable = collection.SubscribeSafe(new InnerObserver(this, value, index, innerSubscription));
  205. }
  206. public override void OnError(Exception error)
  207. {
  208. lock (_gate)
  209. {
  210. ForwardOnError(error);
  211. }
  212. }
  213. public override void OnCompleted()
  214. {
  215. _isStopped = true;
  216. if (_group.Count == 1)
  217. {
  218. //
  219. // Notice there can be a race between OnCompleted of the source and any
  220. // of the inner sequences, where both see _group.Count == 1, and one is
  221. // waiting for the lock. There won't be a double OnCompleted observation
  222. // though, because the call to Dispose silences the observer by swapping
  223. // in a NopObserver<T>.
  224. //
  225. lock (_gate)
  226. {
  227. ForwardOnCompleted();
  228. }
  229. }
  230. else
  231. {
  232. _sourceSubscription.Dispose();
  233. }
  234. }
  235. private sealed class InnerObserver : IObserver<TCollection>
  236. {
  237. private readonly _ _parent;
  238. private readonly TSource _value;
  239. private readonly int _valueIndex;
  240. private readonly IDisposable _self;
  241. public InnerObserver(_ parent, TSource value, int index, IDisposable self)
  242. {
  243. _parent = parent;
  244. _value = value;
  245. _valueIndex = index;
  246. _self = self;
  247. }
  248. private int _index;
  249. public void OnNext(TCollection value)
  250. {
  251. var res = default(TResult);
  252. try
  253. {
  254. res = _parent._resultSelector(_value, _valueIndex, value, checked(_index++));
  255. }
  256. catch (Exception ex)
  257. {
  258. lock (_parent._gate)
  259. {
  260. _parent.ForwardOnError(ex);
  261. }
  262. return;
  263. }
  264. lock (_parent._gate)
  265. _parent.ForwardOnNext(res);
  266. }
  267. public void OnError(Exception error)
  268. {
  269. lock (_parent._gate)
  270. {
  271. _parent.ForwardOnError(error);
  272. }
  273. }
  274. public void OnCompleted()
  275. {
  276. _parent._group.Remove(_self);
  277. if (_parent._isStopped && _parent._group.Count == 1)
  278. {
  279. //
  280. // Notice there can be a race between OnCompleted of the source and any
  281. // of the inner sequences, where both see _group.Count == 1, and one is
  282. // waiting for the lock. There won't be a double OnCompleted observation
  283. // though, because the call to Dispose silences the observer by swapping
  284. // in a NopObserver<T>.
  285. //
  286. lock (_parent._gate)
  287. {
  288. _parent.ForwardOnCompleted();
  289. }
  290. }
  291. }
  292. }
  293. }
  294. }
  295. internal sealed class EnumerableSelector : Producer<TResult, EnumerableSelector._>
  296. {
  297. private readonly IObservable<TSource> _source;
  298. private readonly Func<TSource, IEnumerable<TCollection>> _collectionSelector;
  299. private readonly Func<TSource, TCollection, TResult> _resultSelector;
  300. public EnumerableSelector(IObservable<TSource> source, Func<TSource, IEnumerable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  301. {
  302. _source = source;
  303. _collectionSelector = collectionSelector;
  304. _resultSelector = resultSelector;
  305. }
  306. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  307. protected override void Run(_ sink) => sink.Run(_source);
  308. internal sealed class _ : Sink<TSource, TResult>
  309. {
  310. private readonly Func<TSource, IEnumerable<TCollection>> _collectionSelector;
  311. private readonly Func<TSource, TCollection, TResult> _resultSelector;
  312. public _(EnumerableSelector parent, IObserver<TResult> observer)
  313. : base(observer)
  314. {
  315. _collectionSelector = parent._collectionSelector;
  316. _resultSelector = parent._resultSelector;
  317. }
  318. public override void OnNext(TSource value)
  319. {
  320. var xs = default(IEnumerable<TCollection>);
  321. try
  322. {
  323. xs = _collectionSelector(value);
  324. }
  325. catch (Exception exception)
  326. {
  327. ForwardOnError(exception);
  328. return;
  329. }
  330. var e = default(IEnumerator<TCollection>);
  331. try
  332. {
  333. e = xs.GetEnumerator();
  334. }
  335. catch (Exception exception)
  336. {
  337. ForwardOnError(exception);
  338. return;
  339. }
  340. try
  341. {
  342. var hasNext = true;
  343. while (hasNext)
  344. {
  345. hasNext = false;
  346. var current = default(TResult);
  347. try
  348. {
  349. hasNext = e.MoveNext();
  350. if (hasNext)
  351. current = _resultSelector(value, e.Current);
  352. }
  353. catch (Exception exception)
  354. {
  355. ForwardOnError(exception);
  356. return;
  357. }
  358. if (hasNext)
  359. ForwardOnNext(current);
  360. }
  361. }
  362. finally
  363. {
  364. if (e != null)
  365. e.Dispose();
  366. }
  367. }
  368. }
  369. }
  370. internal sealed class EnumerableSelectorIndexed : Producer<TResult, EnumerableSelectorIndexed._>
  371. {
  372. private readonly IObservable<TSource> _source;
  373. private readonly Func<TSource, int, IEnumerable<TCollection>> _collectionSelector;
  374. private readonly Func<TSource, int, TCollection, int, TResult> _resultSelector;
  375. public EnumerableSelectorIndexed(IObservable<TSource> source, Func<TSource, int, IEnumerable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  376. {
  377. _source = source;
  378. _collectionSelector = collectionSelector;
  379. _resultSelector = resultSelector;
  380. }
  381. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  382. protected override void Run(_ sink) => sink.Run(_source);
  383. internal sealed class _ : Sink<TSource, TResult>
  384. {
  385. private readonly Func<TSource, int, IEnumerable<TCollection>> _collectionSelector;
  386. private readonly Func<TSource, int, TCollection, int, TResult> _resultSelector;
  387. public _(EnumerableSelectorIndexed parent, IObserver<TResult> observer)
  388. : base(observer)
  389. {
  390. _collectionSelector = parent._collectionSelector;
  391. _resultSelector = parent._resultSelector;
  392. }
  393. private int _index;
  394. public override void OnNext(TSource value)
  395. {
  396. var index = checked(_index++);
  397. var xs = default(IEnumerable<TCollection>);
  398. try
  399. {
  400. xs = _collectionSelector(value, index);
  401. }
  402. catch (Exception exception)
  403. {
  404. ForwardOnError(exception);
  405. return;
  406. }
  407. var e = default(IEnumerator<TCollection>);
  408. try
  409. {
  410. e = xs.GetEnumerator();
  411. }
  412. catch (Exception exception)
  413. {
  414. ForwardOnError(exception);
  415. return;
  416. }
  417. try
  418. {
  419. var eIndex = 0;
  420. var hasNext = true;
  421. while (hasNext)
  422. {
  423. hasNext = false;
  424. var current = default(TResult);
  425. try
  426. {
  427. hasNext = e.MoveNext();
  428. if (hasNext)
  429. current = _resultSelector(value, index, e.Current, checked(eIndex++));
  430. }
  431. catch (Exception exception)
  432. {
  433. ForwardOnError(exception);
  434. return;
  435. }
  436. if (hasNext)
  437. ForwardOnNext(current);
  438. }
  439. }
  440. finally
  441. {
  442. if (e != null)
  443. e.Dispose();
  444. }
  445. }
  446. }
  447. }
  448. internal sealed class TaskSelector : Producer<TResult, TaskSelector._>
  449. {
  450. private readonly IObservable<TSource> _source;
  451. private readonly Func<TSource, CancellationToken, Task<TCollection>> _collectionSelector;
  452. private readonly Func<TSource, TCollection, TResult> _resultSelector;
  453. public TaskSelector(IObservable<TSource> source, Func<TSource, CancellationToken, Task<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  454. {
  455. _source = source;
  456. _collectionSelector = collectionSelector;
  457. _resultSelector = resultSelector;
  458. }
  459. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  460. protected override void Run(_ sink) => sink.Run(_source);
  461. internal sealed class _ : Sink<TSource, TResult>
  462. {
  463. private readonly object _gate = new object();
  464. private readonly CancellationDisposable _cancel = new CancellationDisposable();
  465. private readonly Func<TSource, CancellationToken, Task<TCollection>> _collectionSelector;
  466. private readonly Func<TSource, TCollection, TResult> _resultSelector;
  467. public _(TaskSelector parent, IObserver<TResult> observer)
  468. : base(observer)
  469. {
  470. _collectionSelector = parent._collectionSelector;
  471. _resultSelector = parent._resultSelector;
  472. }
  473. private volatile int _count;
  474. public void Run(IObservable<TSource> source)
  475. {
  476. _count = 1;
  477. SetUpstream(StableCompositeDisposable.Create(source.SubscribeSafe(this), _cancel));
  478. }
  479. public override void OnNext(TSource value)
  480. {
  481. var task = default(Task<TCollection>);
  482. try
  483. {
  484. Interlocked.Increment(ref _count);
  485. task = _collectionSelector(value, _cancel.Token);
  486. }
  487. catch (Exception ex)
  488. {
  489. lock (_gate)
  490. {
  491. ForwardOnError(ex);
  492. }
  493. return;
  494. }
  495. if (task.IsCompleted)
  496. {
  497. OnCompletedTask(value, task);
  498. }
  499. else
  500. {
  501. AttachContinuation(value, task);
  502. }
  503. }
  504. private void AttachContinuation(TSource value, Task<TCollection> task)
  505. {
  506. //
  507. // Separate method to avoid closure in synchronous completion case.
  508. //
  509. task.ContinueWith(t => OnCompletedTask(value, t));
  510. }
  511. private void OnCompletedTask(TSource value, Task<TCollection> task)
  512. {
  513. switch (task.Status)
  514. {
  515. case TaskStatus.RanToCompletion:
  516. {
  517. var res = default(TResult);
  518. try
  519. {
  520. res = _resultSelector(value, task.Result);
  521. }
  522. catch (Exception ex)
  523. {
  524. lock (_gate)
  525. {
  526. ForwardOnError(ex);
  527. }
  528. return;
  529. }
  530. lock (_gate)
  531. ForwardOnNext(res);
  532. OnCompleted();
  533. }
  534. break;
  535. case TaskStatus.Faulted:
  536. {
  537. lock (_gate)
  538. {
  539. ForwardOnError(task.Exception.InnerException);
  540. }
  541. }
  542. break;
  543. case TaskStatus.Canceled:
  544. {
  545. if (!_cancel.IsDisposed)
  546. {
  547. lock (_gate)
  548. {
  549. ForwardOnError(new TaskCanceledException(task));
  550. }
  551. }
  552. }
  553. break;
  554. }
  555. }
  556. public override void OnError(Exception error)
  557. {
  558. lock (_gate)
  559. {
  560. ForwardOnError(error);
  561. }
  562. }
  563. public override void OnCompleted()
  564. {
  565. if (Interlocked.Decrement(ref _count) == 0)
  566. {
  567. lock (_gate)
  568. {
  569. ForwardOnCompleted();
  570. }
  571. }
  572. }
  573. }
  574. }
  575. internal sealed class TaskSelectorIndexed : Producer<TResult, TaskSelectorIndexed._>
  576. {
  577. private readonly IObservable<TSource> _source;
  578. private readonly Func<TSource, int, CancellationToken, Task<TCollection>> _collectionSelector;
  579. private readonly Func<TSource, int, TCollection, TResult> _resultSelector;
  580. public TaskSelectorIndexed(IObservable<TSource> source, Func<TSource, int, CancellationToken, Task<TCollection>> collectionSelector, Func<TSource, int, TCollection, TResult> resultSelector)
  581. {
  582. _source = source;
  583. _collectionSelector = collectionSelector;
  584. _resultSelector = resultSelector;
  585. }
  586. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  587. protected override void Run(_ sink) => sink.Run(_source);
  588. internal sealed class _ : Sink<TSource, TResult>
  589. {
  590. private readonly object _gate = new object();
  591. private readonly CancellationDisposable _cancel = new CancellationDisposable();
  592. private readonly Func<TSource, int, CancellationToken, Task<TCollection>> _collectionSelector;
  593. private readonly Func<TSource, int, TCollection, TResult> _resultSelector;
  594. public _(TaskSelectorIndexed parent, IObserver<TResult> observer)
  595. : base(observer)
  596. {
  597. _collectionSelector = parent._collectionSelector;
  598. _resultSelector = parent._resultSelector;
  599. }
  600. private volatile int _count;
  601. private int _index;
  602. public void Run(IObservable<TSource> source)
  603. {
  604. _count = 1;
  605. SetUpstream(StableCompositeDisposable.Create(source.SubscribeSafe(this), _cancel));
  606. }
  607. public override void OnNext(TSource value)
  608. {
  609. var index = checked(_index++);
  610. var task = default(Task<TCollection>);
  611. try
  612. {
  613. Interlocked.Increment(ref _count);
  614. task = _collectionSelector(value, index, _cancel.Token);
  615. }
  616. catch (Exception ex)
  617. {
  618. lock (_gate)
  619. {
  620. ForwardOnError(ex);
  621. }
  622. return;
  623. }
  624. if (task.IsCompleted)
  625. {
  626. OnCompletedTask(value, index, task);
  627. }
  628. else
  629. {
  630. AttachContinuation(value, index, task);
  631. }
  632. }
  633. private void AttachContinuation(TSource value, int index, Task<TCollection> task)
  634. {
  635. //
  636. // Separate method to avoid closure in synchronous completion case.
  637. //
  638. task.ContinueWith(t => OnCompletedTask(value, index, t));
  639. }
  640. private void OnCompletedTask(TSource value, int index, Task<TCollection> task)
  641. {
  642. switch (task.Status)
  643. {
  644. case TaskStatus.RanToCompletion:
  645. {
  646. var res = default(TResult);
  647. try
  648. {
  649. res = _resultSelector(value, index, task.Result);
  650. }
  651. catch (Exception ex)
  652. {
  653. lock (_gate)
  654. {
  655. ForwardOnError(ex);
  656. }
  657. return;
  658. }
  659. lock (_gate)
  660. ForwardOnNext(res);
  661. OnCompleted();
  662. }
  663. break;
  664. case TaskStatus.Faulted:
  665. {
  666. lock (_gate)
  667. {
  668. ForwardOnError(task.Exception.InnerException);
  669. }
  670. }
  671. break;
  672. case TaskStatus.Canceled:
  673. {
  674. if (!_cancel.IsDisposed)
  675. {
  676. lock (_gate)
  677. {
  678. ForwardOnError(new TaskCanceledException(task));
  679. }
  680. }
  681. }
  682. break;
  683. }
  684. }
  685. public override void OnError(Exception error)
  686. {
  687. lock (_gate)
  688. {
  689. ForwardOnError(error);
  690. }
  691. }
  692. public override void OnCompleted()
  693. {
  694. if (Interlocked.Decrement(ref _count) == 0)
  695. {
  696. lock (_gate)
  697. {
  698. ForwardOnCompleted();
  699. }
  700. }
  701. }
  702. }
  703. }
  704. }
  705. internal static class SelectMany<TSource, TResult>
  706. {
  707. internal class ObservableSelector : Producer<TResult, ObservableSelector._>
  708. {
  709. protected readonly IObservable<TSource> _source;
  710. protected readonly Func<TSource, IObservable<TResult>> _selector;
  711. public ObservableSelector(IObservable<TSource> source, Func<TSource, IObservable<TResult>> selector)
  712. {
  713. _source = source;
  714. _selector = selector;
  715. }
  716. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  717. protected override void Run(_ sink) => sink.Run(_source);
  718. internal class _ : Sink<TSource, TResult>
  719. {
  720. protected readonly object _gate = new object();
  721. private readonly SingleAssignmentDisposable _sourceSubscription = new SingleAssignmentDisposable();
  722. private readonly CompositeDisposable _group = new CompositeDisposable();
  723. private readonly Func<TSource, IObservable<TResult>> _selector;
  724. public _(ObservableSelector parent, IObserver<TResult> observer)
  725. : base(observer)
  726. {
  727. _selector = parent._selector;
  728. _group.Add(_sourceSubscription);
  729. }
  730. private bool _isStopped;
  731. public void Run(IObservable<TSource> source)
  732. {
  733. _isStopped = false;
  734. _sourceSubscription.Disposable = source.SubscribeSafe(this);
  735. SetUpstream(_group);
  736. }
  737. public override void OnNext(TSource value)
  738. {
  739. var inner = default(IObservable<TResult>);
  740. try
  741. {
  742. inner = _selector(value);
  743. }
  744. catch (Exception ex)
  745. {
  746. lock (_gate)
  747. {
  748. ForwardOnError(ex);
  749. }
  750. return;
  751. }
  752. SubscribeInner(inner);
  753. }
  754. public override void OnError(Exception error)
  755. {
  756. lock (_gate)
  757. {
  758. ForwardOnError(error);
  759. }
  760. }
  761. public override void OnCompleted()
  762. {
  763. Final();
  764. }
  765. protected void Final()
  766. {
  767. _isStopped = true;
  768. if (_group.Count == 1)
  769. {
  770. //
  771. // Notice there can be a race between OnCompleted of the source and any
  772. // of the inner sequences, where both see _group.Count == 1, and one is
  773. // waiting for the lock. There won't be a double OnCompleted observation
  774. // though, because the call to Dispose silences the observer by swapping
  775. // in a NopObserver<T>.
  776. //
  777. lock (_gate)
  778. {
  779. ForwardOnCompleted();
  780. }
  781. }
  782. else
  783. {
  784. _sourceSubscription.Dispose();
  785. }
  786. }
  787. protected void SubscribeInner(IObservable<TResult> inner)
  788. {
  789. var innerSubscription = new SingleAssignmentDisposable();
  790. _group.Add(innerSubscription);
  791. innerSubscription.Disposable = inner.SubscribeSafe(new InnerObserver(this, innerSubscription));
  792. }
  793. private sealed class InnerObserver : IObserver<TResult>
  794. {
  795. private readonly _ _parent;
  796. private readonly IDisposable _self;
  797. public InnerObserver(_ parent, IDisposable self)
  798. {
  799. _parent = parent;
  800. _self = self;
  801. }
  802. public void OnNext(TResult value)
  803. {
  804. lock (_parent._gate)
  805. _parent.ForwardOnNext(value);
  806. }
  807. public void OnError(Exception error)
  808. {
  809. lock (_parent._gate)
  810. {
  811. _parent.ForwardOnError(error);
  812. }
  813. }
  814. public void OnCompleted()
  815. {
  816. _parent._group.Remove(_self);
  817. if (_parent._isStopped && _parent._group.Count == 1)
  818. {
  819. //
  820. // Notice there can be a race between OnCompleted of the source and any
  821. // of the inner sequences, where both see _group.Count == 1, and one is
  822. // waiting for the lock. There won't be a double OnCompleted observation
  823. // though, because the call to Dispose silences the observer by swapping
  824. // in a NopObserver<T>.
  825. //
  826. lock (_parent._gate)
  827. {
  828. _parent.ForwardOnCompleted();
  829. }
  830. }
  831. }
  832. }
  833. }
  834. }
  835. internal sealed class ObservableSelectors : ObservableSelector
  836. {
  837. private readonly Func<Exception, IObservable<TResult>> _selectorOnError;
  838. private readonly Func<IObservable<TResult>> _selectorOnCompleted;
  839. public ObservableSelectors(IObservable<TSource> source, Func<TSource, IObservable<TResult>> selector, Func<Exception, IObservable<TResult>> selectorOnError, Func<IObservable<TResult>> selectorOnCompleted)
  840. : base(source, selector)
  841. {
  842. _selectorOnError = selectorOnError;
  843. _selectorOnCompleted = selectorOnCompleted;
  844. }
  845. protected override ObservableSelector._ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  846. new internal sealed class _ : ObservableSelector._
  847. {
  848. private readonly Func<Exception, IObservable<TResult>> _selectorOnError;
  849. private readonly Func<IObservable<TResult>> _selectorOnCompleted;
  850. public _(ObservableSelectors parent, IObserver<TResult> observer)
  851. : base(parent, observer)
  852. {
  853. _selectorOnError = parent._selectorOnError;
  854. _selectorOnCompleted = parent._selectorOnCompleted;
  855. }
  856. public override void OnError(Exception error)
  857. {
  858. if (_selectorOnError != null)
  859. {
  860. var inner = default(IObservable<TResult>);
  861. try
  862. {
  863. inner = _selectorOnError(error);
  864. }
  865. catch (Exception ex)
  866. {
  867. lock (_gate)
  868. {
  869. ForwardOnError(ex);
  870. }
  871. return;
  872. }
  873. SubscribeInner(inner);
  874. Final();
  875. }
  876. else
  877. {
  878. base.OnError(error);
  879. }
  880. }
  881. public override void OnCompleted()
  882. {
  883. if (_selectorOnCompleted != null)
  884. {
  885. var inner = default(IObservable<TResult>);
  886. try
  887. {
  888. inner = _selectorOnCompleted();
  889. }
  890. catch (Exception ex)
  891. {
  892. lock (_gate)
  893. {
  894. ForwardOnError(ex);
  895. }
  896. return;
  897. }
  898. SubscribeInner(inner);
  899. }
  900. Final();
  901. }
  902. }
  903. }
  904. internal class ObservableSelectorIndexed : Producer<TResult, ObservableSelectorIndexed._>
  905. {
  906. protected readonly IObservable<TSource> _source;
  907. protected readonly Func<TSource, int, IObservable<TResult>> _selector;
  908. public ObservableSelectorIndexed(IObservable<TSource> source, Func<TSource, int, IObservable<TResult>> selector)
  909. {
  910. _source = source;
  911. _selector = selector;
  912. }
  913. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  914. protected override void Run(_ sink) => sink.Run(_source);
  915. internal class _ : Sink<TSource, TResult>
  916. {
  917. private readonly object _gate = new object();
  918. private readonly SingleAssignmentDisposable _sourceSubscription = new SingleAssignmentDisposable();
  919. private readonly CompositeDisposable _group = new CompositeDisposable();
  920. protected readonly Func<TSource, int, IObservable<TResult>> _selector;
  921. public _(ObservableSelectorIndexed parent, IObserver<TResult> observer)
  922. : base(observer)
  923. {
  924. _selector = parent._selector;
  925. _group.Add(_sourceSubscription);
  926. }
  927. private bool _isStopped;
  928. private int _index;
  929. public void Run(IObservable<TSource> source)
  930. {
  931. _isStopped = false;
  932. _sourceSubscription.Disposable = source.SubscribeSafe(this);
  933. SetUpstream(_group);
  934. }
  935. public override void OnNext(TSource value)
  936. {
  937. var inner = default(IObservable<TResult>);
  938. try
  939. {
  940. inner = _selector(value, checked(_index++));
  941. }
  942. catch (Exception ex)
  943. {
  944. lock (_gate)
  945. {
  946. ForwardOnError(ex);
  947. }
  948. return;
  949. }
  950. SubscribeInner(inner);
  951. }
  952. public override void OnError(Exception error)
  953. {
  954. lock (_gate)
  955. {
  956. ForwardOnError(error);
  957. }
  958. }
  959. public override void OnCompleted()
  960. {
  961. Final();
  962. }
  963. protected void Final()
  964. {
  965. _isStopped = true;
  966. if (_group.Count == 1)
  967. {
  968. //
  969. // Notice there can be a race between OnCompleted of the source and any
  970. // of the inner sequences, where both see _group.Count == 1, and one is
  971. // waiting for the lock. There won't be a double OnCompleted observation
  972. // though, because the call to Dispose silences the observer by swapping
  973. // in a NopObserver<T>.
  974. //
  975. lock (_gate)
  976. {
  977. ForwardOnCompleted();
  978. }
  979. }
  980. else
  981. {
  982. _sourceSubscription.Dispose();
  983. }
  984. }
  985. protected void SubscribeInner(IObservable<TResult> inner)
  986. {
  987. var innerSubscription = new SingleAssignmentDisposable();
  988. _group.Add(innerSubscription);
  989. innerSubscription.Disposable = inner.SubscribeSafe(new InnerObserver(this, innerSubscription));
  990. }
  991. private sealed class InnerObserver : IObserver<TResult>
  992. {
  993. private readonly _ _parent;
  994. private readonly IDisposable _self;
  995. public InnerObserver(_ parent, IDisposable self)
  996. {
  997. _parent = parent;
  998. _self = self;
  999. }
  1000. public void OnNext(TResult value)
  1001. {
  1002. lock (_parent._gate)
  1003. _parent.ForwardOnNext(value);
  1004. }
  1005. public void OnError(Exception error)
  1006. {
  1007. lock (_parent._gate)
  1008. {
  1009. _parent.ForwardOnError(error);
  1010. }
  1011. }
  1012. public void OnCompleted()
  1013. {
  1014. _parent._group.Remove(_self);
  1015. if (_parent._isStopped && _parent._group.Count == 1)
  1016. {
  1017. //
  1018. // Notice there can be a race between OnCompleted of the source and any
  1019. // of the inner sequences, where both see _group.Count == 1, and one is
  1020. // waiting for the lock. There won't be a double OnCompleted observation
  1021. // though, because the call to Dispose silences the observer by swapping
  1022. // in a NopObserver<T>.
  1023. //
  1024. lock (_parent._gate)
  1025. {
  1026. _parent.ForwardOnCompleted();
  1027. }
  1028. }
  1029. }
  1030. }
  1031. }
  1032. }
  1033. internal sealed class ObservableSelectorsIndexed : ObservableSelectorIndexed
  1034. {
  1035. private readonly Func<Exception, IObservable<TResult>> _selectorOnError;
  1036. private readonly Func<IObservable<TResult>> _selectorOnCompleted;
  1037. public ObservableSelectorsIndexed(IObservable<TSource> source, Func<TSource, int, IObservable<TResult>> selector, Func<Exception, IObservable<TResult>> selectorOnError, Func<IObservable<TResult>> selectorOnCompleted)
  1038. : base(source, selector)
  1039. {
  1040. _selectorOnError = selectorOnError;
  1041. _selectorOnCompleted = selectorOnCompleted;
  1042. }
  1043. protected override ObservableSelectorIndexed._ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  1044. new internal sealed class _ : ObservableSelectorIndexed._
  1045. {
  1046. private readonly object _gate = new object();
  1047. private readonly SingleAssignmentDisposable _sourceSubscription = new SingleAssignmentDisposable();
  1048. private readonly CompositeDisposable _group = new CompositeDisposable();
  1049. private readonly Func<Exception, IObservable<TResult>> _selectorOnError;
  1050. private readonly Func<IObservable<TResult>> _selectorOnCompleted;
  1051. public _(ObservableSelectorsIndexed parent, IObserver<TResult> observer)
  1052. : base(parent, observer)
  1053. {
  1054. _selectorOnError = parent._selectorOnError;
  1055. _selectorOnCompleted = parent._selectorOnCompleted;
  1056. _group.Add(_sourceSubscription);
  1057. }
  1058. public override void OnError(Exception error)
  1059. {
  1060. if (_selectorOnError != null)
  1061. {
  1062. var inner = default(IObservable<TResult>);
  1063. try
  1064. {
  1065. inner = _selectorOnError(error);
  1066. }
  1067. catch (Exception ex)
  1068. {
  1069. lock (_gate)
  1070. {
  1071. ForwardOnError(ex);
  1072. }
  1073. return;
  1074. }
  1075. SubscribeInner(inner);
  1076. Final();
  1077. }
  1078. else
  1079. {
  1080. base.OnError(error);
  1081. }
  1082. }
  1083. public override void OnCompleted()
  1084. {
  1085. if (_selectorOnCompleted != null)
  1086. {
  1087. var inner = default(IObservable<TResult>);
  1088. try
  1089. {
  1090. inner = _selectorOnCompleted();
  1091. }
  1092. catch (Exception ex)
  1093. {
  1094. lock (_gate)
  1095. {
  1096. ForwardOnError(ex);
  1097. }
  1098. return;
  1099. }
  1100. SubscribeInner(inner);
  1101. }
  1102. Final();
  1103. }
  1104. }
  1105. }
  1106. internal sealed class EnumerableSelector : Producer<TResult, EnumerableSelector._>
  1107. {
  1108. private readonly IObservable<TSource> _source;
  1109. private readonly Func<TSource, IEnumerable<TResult>> _selector;
  1110. public EnumerableSelector(IObservable<TSource> source, Func<TSource, IEnumerable<TResult>> selector)
  1111. {
  1112. _source = source;
  1113. _selector = selector;
  1114. }
  1115. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  1116. protected override void Run(_ sink) => sink.Run(_source);
  1117. internal sealed class _ : Sink<TSource, TResult>
  1118. {
  1119. private readonly Func<TSource, IEnumerable<TResult>> _selector;
  1120. public _(EnumerableSelector parent, IObserver<TResult> observer)
  1121. : base(observer)
  1122. {
  1123. _selector = parent._selector;
  1124. }
  1125. public override void OnNext(TSource value)
  1126. {
  1127. var xs = default(IEnumerable<TResult>);
  1128. try
  1129. {
  1130. xs = _selector(value);
  1131. }
  1132. catch (Exception exception)
  1133. {
  1134. ForwardOnError(exception);
  1135. return;
  1136. }
  1137. var e = default(IEnumerator<TResult>);
  1138. try
  1139. {
  1140. e = xs.GetEnumerator();
  1141. }
  1142. catch (Exception exception)
  1143. {
  1144. ForwardOnError(exception);
  1145. return;
  1146. }
  1147. try
  1148. {
  1149. var hasNext = true;
  1150. while (hasNext)
  1151. {
  1152. hasNext = false;
  1153. var current = default(TResult);
  1154. try
  1155. {
  1156. hasNext = e.MoveNext();
  1157. if (hasNext)
  1158. current = e.Current;
  1159. }
  1160. catch (Exception exception)
  1161. {
  1162. ForwardOnError(exception);
  1163. return;
  1164. }
  1165. if (hasNext)
  1166. ForwardOnNext(current);
  1167. }
  1168. }
  1169. finally
  1170. {
  1171. if (e != null)
  1172. e.Dispose();
  1173. }
  1174. }
  1175. }
  1176. }
  1177. internal sealed class EnumerableSelectorIndexed : Producer<TResult, EnumerableSelectorIndexed._>
  1178. {
  1179. private readonly IObservable<TSource> _source;
  1180. private readonly Func<TSource, int, IEnumerable<TResult>> _selector;
  1181. public EnumerableSelectorIndexed(IObservable<TSource> source, Func<TSource, int, IEnumerable<TResult>> selector)
  1182. {
  1183. _source = source;
  1184. _selector = selector;
  1185. }
  1186. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  1187. protected override void Run(_ sink) => sink.Run(_source);
  1188. internal sealed class _ : Sink<TSource, TResult>
  1189. {
  1190. private readonly Func<TSource, int, IEnumerable<TResult>> _selector;
  1191. public _(EnumerableSelectorIndexed parent, IObserver<TResult> observer)
  1192. : base(observer)
  1193. {
  1194. _selector = parent._selector;
  1195. }
  1196. private int _index;
  1197. public override void OnNext(TSource value)
  1198. {
  1199. var xs = default(IEnumerable<TResult>);
  1200. try
  1201. {
  1202. xs = _selector(value, checked(_index++));
  1203. }
  1204. catch (Exception exception)
  1205. {
  1206. ForwardOnError(exception);
  1207. return;
  1208. }
  1209. var e = default(IEnumerator<TResult>);
  1210. try
  1211. {
  1212. e = xs.GetEnumerator();
  1213. }
  1214. catch (Exception exception)
  1215. {
  1216. ForwardOnError(exception);
  1217. return;
  1218. }
  1219. try
  1220. {
  1221. var hasNext = true;
  1222. while (hasNext)
  1223. {
  1224. hasNext = false;
  1225. var current = default(TResult);
  1226. try
  1227. {
  1228. hasNext = e.MoveNext();
  1229. if (hasNext)
  1230. current = e.Current;
  1231. }
  1232. catch (Exception exception)
  1233. {
  1234. ForwardOnError(exception);
  1235. return;
  1236. }
  1237. if (hasNext)
  1238. ForwardOnNext(current);
  1239. }
  1240. }
  1241. finally
  1242. {
  1243. if (e != null)
  1244. e.Dispose();
  1245. }
  1246. }
  1247. }
  1248. }
  1249. internal sealed class TaskSelector : Producer<TResult, TaskSelector._>
  1250. {
  1251. private readonly IObservable<TSource> _source;
  1252. private readonly Func<TSource, CancellationToken, Task<TResult>> _selector;
  1253. public TaskSelector(IObservable<TSource> source, Func<TSource, CancellationToken, Task<TResult>> selector)
  1254. {
  1255. _source = source;
  1256. _selector = selector;
  1257. }
  1258. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  1259. protected override void Run(_ sink) => sink.Run(_source);
  1260. internal sealed class _ : Sink<TSource, TResult>
  1261. {
  1262. private readonly object _gate = new object();
  1263. private readonly CancellationDisposable _cancel = new CancellationDisposable();
  1264. private readonly Func<TSource, CancellationToken, Task<TResult>> _selector;
  1265. public _(TaskSelector parent, IObserver<TResult> observer)
  1266. : base(observer)
  1267. {
  1268. _selector = parent._selector;
  1269. }
  1270. private volatile int _count;
  1271. public void Run(IObservable<TSource> source)
  1272. {
  1273. _count = 1;
  1274. SetUpstream(StableCompositeDisposable.Create(source.SubscribeSafe(this), _cancel));
  1275. }
  1276. public override void OnNext(TSource value)
  1277. {
  1278. var task = default(Task<TResult>);
  1279. try
  1280. {
  1281. Interlocked.Increment(ref _count);
  1282. task = _selector(value, _cancel.Token);
  1283. }
  1284. catch (Exception ex)
  1285. {
  1286. lock (_gate)
  1287. {
  1288. ForwardOnError(ex);
  1289. }
  1290. return;
  1291. }
  1292. if (task.IsCompleted)
  1293. {
  1294. OnCompletedTask(task);
  1295. }
  1296. else
  1297. {
  1298. task.ContinueWith(OnCompletedTask);
  1299. }
  1300. }
  1301. private void OnCompletedTask(Task<TResult> task)
  1302. {
  1303. switch (task.Status)
  1304. {
  1305. case TaskStatus.RanToCompletion:
  1306. {
  1307. lock (_gate)
  1308. ForwardOnNext(task.Result);
  1309. OnCompleted();
  1310. }
  1311. break;
  1312. case TaskStatus.Faulted:
  1313. {
  1314. lock (_gate)
  1315. {
  1316. ForwardOnError(task.Exception.InnerException);
  1317. }
  1318. }
  1319. break;
  1320. case TaskStatus.Canceled:
  1321. {
  1322. if (!_cancel.IsDisposed)
  1323. {
  1324. lock (_gate)
  1325. {
  1326. ForwardOnError(new TaskCanceledException(task));
  1327. }
  1328. }
  1329. }
  1330. break;
  1331. }
  1332. }
  1333. public override void OnError(Exception error)
  1334. {
  1335. lock (_gate)
  1336. {
  1337. ForwardOnError(error);
  1338. }
  1339. }
  1340. public override void OnCompleted()
  1341. {
  1342. if (Interlocked.Decrement(ref _count) == 0)
  1343. {
  1344. lock (_gate)
  1345. {
  1346. ForwardOnCompleted();
  1347. }
  1348. }
  1349. }
  1350. }
  1351. }
  1352. internal sealed class TaskSelectorIndexed : Producer<TResult, TaskSelectorIndexed._>
  1353. {
  1354. private readonly IObservable<TSource> _source;
  1355. private readonly Func<TSource, int, CancellationToken, Task<TResult>> _selector;
  1356. public TaskSelectorIndexed(IObservable<TSource> source, Func<TSource, int, CancellationToken, Task<TResult>> selector)
  1357. {
  1358. _source = source;
  1359. _selector = selector;
  1360. }
  1361. protected override _ CreateSink(IObserver<TResult> observer) => new _(this, observer);
  1362. protected override void Run(_ sink) => sink.Run(_source);
  1363. internal sealed class _ : Sink<TSource, TResult>
  1364. {
  1365. private readonly object _gate = new object();
  1366. private readonly CancellationDisposable _cancel = new CancellationDisposable();
  1367. private readonly Func<TSource, int, CancellationToken, Task<TResult>> _selector;
  1368. public _(TaskSelectorIndexed parent, IObserver<TResult> observer)
  1369. : base(observer)
  1370. {
  1371. _selector = parent._selector;
  1372. }
  1373. private volatile int _count;
  1374. private int _index;
  1375. public void Run(IObservable<TSource> source)
  1376. {
  1377. _count = 1;
  1378. SetUpstream(StableCompositeDisposable.Create(source.SubscribeSafe(this), _cancel));
  1379. }
  1380. public override void OnNext(TSource value)
  1381. {
  1382. var task = default(Task<TResult>);
  1383. try
  1384. {
  1385. Interlocked.Increment(ref _count);
  1386. task = _selector(value, checked(_index++), _cancel.Token);
  1387. }
  1388. catch (Exception ex)
  1389. {
  1390. lock (_gate)
  1391. {
  1392. ForwardOnError(ex);
  1393. }
  1394. return;
  1395. }
  1396. if (task.IsCompleted)
  1397. {
  1398. OnCompletedTask(task);
  1399. }
  1400. else
  1401. {
  1402. task.ContinueWith(OnCompletedTask);
  1403. }
  1404. }
  1405. private void OnCompletedTask(Task<TResult> task)
  1406. {
  1407. switch (task.Status)
  1408. {
  1409. case TaskStatus.RanToCompletion:
  1410. {
  1411. lock (_gate)
  1412. ForwardOnNext(task.Result);
  1413. OnCompleted();
  1414. }
  1415. break;
  1416. case TaskStatus.Faulted:
  1417. {
  1418. lock (_gate)
  1419. {
  1420. ForwardOnError(task.Exception.InnerException);
  1421. }
  1422. }
  1423. break;
  1424. case TaskStatus.Canceled:
  1425. {
  1426. if (!_cancel.IsDisposed)
  1427. {
  1428. lock (_gate)
  1429. {
  1430. ForwardOnError(new TaskCanceledException(task));
  1431. }
  1432. }
  1433. }
  1434. break;
  1435. }
  1436. }
  1437. public override void OnError(Exception error)
  1438. {
  1439. lock (_gate)
  1440. {
  1441. ForwardOnError(error);
  1442. }
  1443. }
  1444. public override void OnCompleted()
  1445. {
  1446. if (Interlocked.Decrement(ref _count) == 0)
  1447. {
  1448. lock (_gate)
  1449. {
  1450. ForwardOnCompleted();
  1451. }
  1452. }
  1453. }
  1454. }
  1455. }
  1456. }
  1457. }