SelectMany.cs 32 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910
  1. using System;
  2. using System.Collections.Generic;
  3. namespace UniRx.Operators
  4. {
  5. internal class SelectManyObservable<TSource, TResult> : OperatorObservableBase<TResult>
  6. {
  7. readonly IObservable<TSource> source;
  8. readonly Func<TSource, IObservable<TResult>> selector;
  9. readonly Func<TSource, int, IObservable<TResult>> selectorWithIndex;
  10. readonly Func<TSource, IEnumerable<TResult>> selectorEnumerable;
  11. readonly Func<TSource, int, IEnumerable<TResult>> selectorEnumerableWithIndex;
  12. public SelectManyObservable(IObservable<TSource> source, Func<TSource, IObservable<TResult>> selector)
  13. : base(source.IsRequiredSubscribeOnCurrentThread())
  14. {
  15. this.source = source;
  16. this.selector = selector;
  17. }
  18. public SelectManyObservable(IObservable<TSource> source, Func<TSource, int, IObservable<TResult>> selector)
  19. : base(source.IsRequiredSubscribeOnCurrentThread())
  20. {
  21. this.source = source;
  22. this.selectorWithIndex = selector;
  23. }
  24. public SelectManyObservable(IObservable<TSource> source, Func<TSource, IEnumerable<TResult>> selector)
  25. : base(source.IsRequiredSubscribeOnCurrentThread())
  26. {
  27. this.source = source;
  28. this.selectorEnumerable = selector;
  29. }
  30. public SelectManyObservable(IObservable<TSource> source, Func<TSource, int, IEnumerable<TResult>> selector)
  31. : base(source.IsRequiredSubscribeOnCurrentThread())
  32. {
  33. this.source = source;
  34. this.selectorEnumerableWithIndex = selector;
  35. }
  36. protected override IDisposable SubscribeCore(IObserver<TResult> observer, IDisposable cancel)
  37. {
  38. if (this.selector != null)
  39. {
  40. return new SelectManyOuterObserver(this, observer, cancel).Run();
  41. }
  42. else if (this.selectorWithIndex != null)
  43. {
  44. return new SelectManyObserverWithIndex(this, observer, cancel).Run();
  45. }
  46. else if (this.selectorEnumerable != null)
  47. {
  48. return new SelectManyEnumerableObserver(this, observer, cancel).Run();
  49. }
  50. else if (this.selectorEnumerableWithIndex != null)
  51. {
  52. return new SelectManyEnumerableObserverWithIndex(this, observer, cancel).Run();
  53. }
  54. else
  55. {
  56. throw new InvalidOperationException();
  57. }
  58. }
  59. class SelectManyOuterObserver : OperatorObserverBase<TSource, TResult>
  60. {
  61. readonly SelectManyObservable<TSource, TResult> parent;
  62. CompositeDisposable collectionDisposable;
  63. SingleAssignmentDisposable sourceDisposable;
  64. object gate = new object();
  65. bool isStopped = false;
  66. public SelectManyOuterObserver(SelectManyObservable<TSource, TResult> parent, IObserver<TResult> observer, IDisposable cancel) : base(observer, cancel)
  67. {
  68. this.parent = parent;
  69. }
  70. public IDisposable Run()
  71. {
  72. collectionDisposable = new CompositeDisposable();
  73. sourceDisposable = new SingleAssignmentDisposable();
  74. collectionDisposable.Add(sourceDisposable);
  75. sourceDisposable.Disposable = parent.source.Subscribe(this);
  76. return collectionDisposable;
  77. }
  78. public override void OnNext(TSource value)
  79. {
  80. IObservable<TResult> nextObservable;
  81. try
  82. {
  83. nextObservable = parent.selector(value);
  84. }
  85. catch (Exception ex)
  86. {
  87. try { observer.OnError(ex); } finally { Dispose(); };
  88. return;
  89. }
  90. var disposable = new SingleAssignmentDisposable();
  91. collectionDisposable.Add(disposable);
  92. var collectionObserver = new SelectMany(this, disposable);
  93. disposable.Disposable = nextObservable.Subscribe(collectionObserver);
  94. }
  95. public override void OnError(Exception error)
  96. {
  97. lock (gate)
  98. {
  99. try { observer.OnError(error); } finally { Dispose(); };
  100. }
  101. }
  102. public override void OnCompleted()
  103. {
  104. isStopped = true;
  105. if (collectionDisposable.Count == 1)
  106. {
  107. lock (gate)
  108. {
  109. try { observer.OnCompleted(); } finally { Dispose(); };
  110. }
  111. }
  112. else
  113. {
  114. sourceDisposable.Dispose();
  115. }
  116. }
  117. class SelectMany : OperatorObserverBase<TResult, TResult>
  118. {
  119. readonly SelectManyOuterObserver parent;
  120. readonly IDisposable cancel;
  121. public SelectMany(SelectManyOuterObserver parent, IDisposable cancel)
  122. : base(parent.observer, cancel)
  123. {
  124. this.parent = parent;
  125. this.cancel = cancel;
  126. }
  127. public override void OnNext(TResult value)
  128. {
  129. lock (parent.gate)
  130. {
  131. base.observer.OnNext(value);
  132. }
  133. }
  134. public override void OnError(Exception error)
  135. {
  136. lock (parent.gate)
  137. {
  138. try { observer.OnError(error); } finally { Dispose(); };
  139. }
  140. }
  141. public override void OnCompleted()
  142. {
  143. parent.collectionDisposable.Remove(cancel);
  144. if (parent.isStopped && parent.collectionDisposable.Count == 1)
  145. {
  146. lock (parent.gate)
  147. {
  148. try { observer.OnCompleted(); } finally { Dispose(); };
  149. }
  150. }
  151. }
  152. }
  153. }
  154. class SelectManyObserverWithIndex : OperatorObserverBase<TSource, TResult>
  155. {
  156. readonly SelectManyObservable<TSource, TResult> parent;
  157. CompositeDisposable collectionDisposable;
  158. int index = 0;
  159. object gate = new object();
  160. bool isStopped = false;
  161. SingleAssignmentDisposable sourceDisposable;
  162. public SelectManyObserverWithIndex(SelectManyObservable<TSource, TResult> parent, IObserver<TResult> observer, IDisposable cancel) : base(observer, cancel)
  163. {
  164. this.parent = parent;
  165. }
  166. public IDisposable Run()
  167. {
  168. collectionDisposable = new CompositeDisposable();
  169. sourceDisposable = new SingleAssignmentDisposable();
  170. collectionDisposable.Add(sourceDisposable);
  171. sourceDisposable.Disposable = parent.source.Subscribe(this);
  172. return collectionDisposable;
  173. }
  174. public override void OnNext(TSource value)
  175. {
  176. IObservable<TResult> nextObservable;
  177. try
  178. {
  179. nextObservable = parent.selectorWithIndex(value, index++);
  180. }
  181. catch (Exception ex)
  182. {
  183. try { observer.OnError(ex); } finally { Dispose(); };
  184. return;
  185. }
  186. var disposable = new SingleAssignmentDisposable();
  187. collectionDisposable.Add(disposable);
  188. var collectionObserver = new SelectMany(this, disposable);
  189. disposable.Disposable = nextObservable.Subscribe(collectionObserver);
  190. }
  191. public override void OnError(Exception error)
  192. {
  193. lock (gate)
  194. {
  195. try { observer.OnError(error); } finally { Dispose(); };
  196. }
  197. }
  198. public override void OnCompleted()
  199. {
  200. isStopped = true;
  201. if (collectionDisposable.Count == 1)
  202. {
  203. lock (gate)
  204. {
  205. try { observer.OnCompleted(); } finally { Dispose(); };
  206. }
  207. }
  208. else
  209. {
  210. sourceDisposable.Dispose();
  211. }
  212. }
  213. class SelectMany : OperatorObserverBase<TResult, TResult>
  214. {
  215. readonly SelectManyObserverWithIndex parent;
  216. readonly IDisposable cancel;
  217. public SelectMany(SelectManyObserverWithIndex parent, IDisposable cancel)
  218. : base(parent.observer, cancel)
  219. {
  220. this.parent = parent;
  221. this.cancel = cancel;
  222. }
  223. public override void OnNext(TResult value)
  224. {
  225. lock (parent.gate)
  226. {
  227. base.observer.OnNext(value);
  228. }
  229. }
  230. public override void OnError(Exception error)
  231. {
  232. lock (parent.gate)
  233. {
  234. try { observer.OnError(error); } finally { Dispose(); };
  235. }
  236. }
  237. public override void OnCompleted()
  238. {
  239. parent.collectionDisposable.Remove(cancel);
  240. if (parent.isStopped && parent.collectionDisposable.Count == 1)
  241. {
  242. lock (parent.gate)
  243. {
  244. try { observer.OnCompleted(); } finally { Dispose(); };
  245. }
  246. }
  247. }
  248. }
  249. }
  250. class SelectManyEnumerableObserver : OperatorObserverBase<TSource, TResult>
  251. {
  252. readonly SelectManyObservable<TSource, TResult> parent;
  253. public SelectManyEnumerableObserver(SelectManyObservable<TSource, TResult> parent, IObserver<TResult> observer, IDisposable cancel) : base(observer, cancel)
  254. {
  255. this.parent = parent;
  256. }
  257. public IDisposable Run()
  258. {
  259. return parent.source.Subscribe(this);
  260. }
  261. public override void OnNext(TSource value)
  262. {
  263. IEnumerable<TResult> nextEnumerable;
  264. try
  265. {
  266. nextEnumerable = parent.selectorEnumerable(value);
  267. }
  268. catch (Exception ex)
  269. {
  270. try { observer.OnError(ex); } finally { Dispose(); };
  271. return;
  272. }
  273. var e = nextEnumerable.GetEnumerator();
  274. try
  275. {
  276. var hasNext = true;
  277. while (hasNext)
  278. {
  279. hasNext = false;
  280. var current = default(TResult);
  281. try
  282. {
  283. hasNext = e.MoveNext();
  284. if (hasNext)
  285. {
  286. current = e.Current;
  287. }
  288. }
  289. catch (Exception exception)
  290. {
  291. try { observer.OnError(exception); } finally { Dispose(); }
  292. return;
  293. }
  294. if (hasNext)
  295. {
  296. observer.OnNext(current);
  297. }
  298. }
  299. }
  300. finally
  301. {
  302. if (e != null)
  303. {
  304. e.Dispose();
  305. }
  306. }
  307. }
  308. public override void OnError(Exception error)
  309. {
  310. try { observer.OnError(error); } finally { Dispose(); }
  311. }
  312. public override void OnCompleted()
  313. {
  314. try { observer.OnCompleted(); } finally { Dispose(); }
  315. }
  316. }
  317. class SelectManyEnumerableObserverWithIndex : OperatorObserverBase<TSource, TResult>
  318. {
  319. readonly SelectManyObservable<TSource, TResult> parent;
  320. int index = 0;
  321. public SelectManyEnumerableObserverWithIndex(SelectManyObservable<TSource, TResult> parent, IObserver<TResult> observer, IDisposable cancel) : base(observer, cancel)
  322. {
  323. this.parent = parent;
  324. }
  325. public IDisposable Run()
  326. {
  327. return parent.source.Subscribe(this);
  328. }
  329. public override void OnNext(TSource value)
  330. {
  331. IEnumerable<TResult> nextEnumerable;
  332. try
  333. {
  334. nextEnumerable = parent.selectorEnumerableWithIndex(value, index++);
  335. }
  336. catch (Exception ex)
  337. {
  338. OnError(ex);
  339. return;
  340. }
  341. var e = nextEnumerable.GetEnumerator();
  342. try
  343. {
  344. var hasNext = true;
  345. while (hasNext)
  346. {
  347. hasNext = false;
  348. var current = default(TResult);
  349. try
  350. {
  351. hasNext = e.MoveNext();
  352. if (hasNext)
  353. {
  354. current = e.Current;
  355. }
  356. }
  357. catch (Exception exception)
  358. {
  359. try { observer.OnError(exception); } finally { Dispose(); }
  360. return;
  361. }
  362. if (hasNext)
  363. {
  364. observer.OnNext(current);
  365. }
  366. }
  367. }
  368. finally
  369. {
  370. if (e != null)
  371. {
  372. e.Dispose();
  373. }
  374. }
  375. }
  376. public override void OnError(Exception error)
  377. {
  378. try { observer.OnError(error); } finally { Dispose(); }
  379. }
  380. public override void OnCompleted()
  381. {
  382. try { observer.OnCompleted(); } finally { Dispose(); }
  383. }
  384. }
  385. }
  386. // with resultSelector
  387. internal class SelectManyObservable<TSource, TCollection, TResult> : OperatorObservableBase<TResult>
  388. {
  389. readonly IObservable<TSource> source;
  390. readonly Func<TSource, IObservable<TCollection>> collectionSelector;
  391. readonly Func<TSource, int, IObservable<TCollection>> collectionSelectorWithIndex;
  392. readonly Func<TSource, IEnumerable<TCollection>> collectionSelectorEnumerable;
  393. readonly Func<TSource, int, IEnumerable<TCollection>> collectionSelectorEnumerableWithIndex;
  394. readonly Func<TSource, TCollection, TResult> resultSelector;
  395. readonly Func<TSource, int, TCollection, int, TResult> resultSelectorWithIndex;
  396. public SelectManyObservable(IObservable<TSource> source, Func<TSource, IObservable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  397. : base(source.IsRequiredSubscribeOnCurrentThread())
  398. {
  399. this.source = source;
  400. this.collectionSelector = collectionSelector;
  401. this.resultSelector = resultSelector;
  402. }
  403. public SelectManyObservable(IObservable<TSource> source, Func<TSource, int, IObservable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  404. : base(source.IsRequiredSubscribeOnCurrentThread())
  405. {
  406. this.source = source;
  407. this.collectionSelectorWithIndex = collectionSelector;
  408. this.resultSelectorWithIndex = resultSelector;
  409. }
  410. public SelectManyObservable(IObservable<TSource> source, Func<TSource, IEnumerable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  411. : base(source.IsRequiredSubscribeOnCurrentThread())
  412. {
  413. this.source = source;
  414. this.collectionSelectorEnumerable = collectionSelector;
  415. this.resultSelector = resultSelector;
  416. }
  417. public SelectManyObservable(IObservable<TSource> source, Func<TSource, int, IEnumerable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  418. : base(source.IsRequiredSubscribeOnCurrentThread())
  419. {
  420. this.source = source;
  421. this.collectionSelectorEnumerableWithIndex = collectionSelector;
  422. this.resultSelectorWithIndex = resultSelector;
  423. }
  424. protected override IDisposable SubscribeCore(IObserver<TResult> observer, IDisposable cancel)
  425. {
  426. if (collectionSelector != null)
  427. {
  428. return new SelectManyOuterObserver(this, observer, cancel).Run();
  429. }
  430. else if (collectionSelectorWithIndex != null)
  431. {
  432. return new SelectManyObserverWithIndex(this, observer, cancel).Run();
  433. }
  434. else if (collectionSelectorEnumerable != null)
  435. {
  436. return new SelectManyEnumerableObserver(this, observer, cancel).Run();
  437. }
  438. else if (collectionSelectorEnumerableWithIndex != null)
  439. {
  440. return new SelectManyEnumerableObserverWithIndex(this, observer, cancel).Run();
  441. }
  442. else
  443. {
  444. throw new InvalidOperationException();
  445. }
  446. }
  447. class SelectManyOuterObserver : OperatorObserverBase<TSource, TResult>
  448. {
  449. readonly SelectManyObservable<TSource, TCollection, TResult> parent;
  450. CompositeDisposable collectionDisposable;
  451. object gate = new object();
  452. bool isStopped = false;
  453. SingleAssignmentDisposable sourceDisposable;
  454. public SelectManyOuterObserver(SelectManyObservable<TSource, TCollection, TResult> parent, IObserver<TResult> observer, IDisposable cancel) : base(observer, cancel)
  455. {
  456. this.parent = parent;
  457. }
  458. public IDisposable Run()
  459. {
  460. collectionDisposable = new CompositeDisposable();
  461. sourceDisposable = new SingleAssignmentDisposable();
  462. collectionDisposable.Add(sourceDisposable);
  463. sourceDisposable.Disposable = parent.source.Subscribe(this);
  464. return collectionDisposable;
  465. }
  466. public override void OnNext(TSource value)
  467. {
  468. IObservable<TCollection> nextObservable;
  469. try
  470. {
  471. nextObservable = parent.collectionSelector(value);
  472. }
  473. catch (Exception ex)
  474. {
  475. OnError(ex);
  476. return;
  477. }
  478. var disposable = new SingleAssignmentDisposable();
  479. collectionDisposable.Add(disposable);
  480. var collectionObserver = new SelectMany(this, value, disposable);
  481. disposable.Disposable = nextObservable.Subscribe(collectionObserver);
  482. }
  483. public override void OnError(Exception error)
  484. {
  485. lock (gate)
  486. {
  487. try { observer.OnError(error); } finally { Dispose(); };
  488. }
  489. }
  490. public override void OnCompleted()
  491. {
  492. isStopped = true;
  493. if (collectionDisposable.Count == 1)
  494. {
  495. lock (gate)
  496. {
  497. try { observer.OnCompleted(); } finally { Dispose(); };
  498. }
  499. }
  500. else
  501. {
  502. sourceDisposable.Dispose();
  503. }
  504. }
  505. class SelectMany : OperatorObserverBase<TCollection, TResult>
  506. {
  507. readonly SelectManyOuterObserver parent;
  508. readonly TSource sourceValue;
  509. readonly IDisposable cancel;
  510. public SelectMany(SelectManyOuterObserver parent, TSource value, IDisposable cancel)
  511. : base(parent.observer, cancel)
  512. {
  513. this.parent = parent;
  514. this.sourceValue = value;
  515. this.cancel = cancel;
  516. }
  517. public override void OnNext(TCollection value)
  518. {
  519. TResult resultValue;
  520. try
  521. {
  522. resultValue = parent.parent.resultSelector(sourceValue, value);
  523. }
  524. catch (Exception ex)
  525. {
  526. OnError(ex);
  527. return;
  528. }
  529. lock (parent.gate)
  530. {
  531. base.observer.OnNext(resultValue);
  532. }
  533. }
  534. public override void OnError(Exception error)
  535. {
  536. lock (parent.gate)
  537. {
  538. try { observer.OnError(error); } finally { Dispose(); };
  539. }
  540. }
  541. public override void OnCompleted()
  542. {
  543. parent.collectionDisposable.Remove(cancel);
  544. if (parent.isStopped && parent.collectionDisposable.Count == 1)
  545. {
  546. lock (parent.gate)
  547. {
  548. try { observer.OnCompleted(); } finally { Dispose(); };
  549. }
  550. }
  551. }
  552. }
  553. }
  554. class SelectManyObserverWithIndex : OperatorObserverBase<TSource, TResult>
  555. {
  556. readonly SelectManyObservable<TSource, TCollection, TResult> parent;
  557. CompositeDisposable collectionDisposable;
  558. object gate = new object();
  559. bool isStopped = false;
  560. SingleAssignmentDisposable sourceDisposable;
  561. int index = 0;
  562. public SelectManyObserverWithIndex(SelectManyObservable<TSource, TCollection, TResult> parent, IObserver<TResult> observer, IDisposable cancel) : base(observer, cancel)
  563. {
  564. this.parent = parent;
  565. }
  566. public IDisposable Run()
  567. {
  568. collectionDisposable = new CompositeDisposable();
  569. sourceDisposable = new SingleAssignmentDisposable();
  570. collectionDisposable.Add(sourceDisposable);
  571. sourceDisposable.Disposable = parent.source.Subscribe(this);
  572. return collectionDisposable;
  573. }
  574. public override void OnNext(TSource value)
  575. {
  576. var i = index++;
  577. IObservable<TCollection> nextObservable;
  578. try
  579. {
  580. nextObservable = parent.collectionSelectorWithIndex(value, i);
  581. }
  582. catch (Exception ex)
  583. {
  584. OnError(ex);
  585. return;
  586. }
  587. var disposable = new SingleAssignmentDisposable();
  588. collectionDisposable.Add(disposable);
  589. var collectionObserver = new SelectManyObserver(this, value, i, disposable);
  590. disposable.Disposable = nextObservable.Subscribe(collectionObserver);
  591. }
  592. public override void OnError(Exception error)
  593. {
  594. lock (gate)
  595. {
  596. try { observer.OnError(error); } finally { Dispose(); };
  597. }
  598. }
  599. public override void OnCompleted()
  600. {
  601. isStopped = true;
  602. if (collectionDisposable.Count == 1)
  603. {
  604. lock (gate)
  605. {
  606. try { observer.OnCompleted(); } finally { Dispose(); };
  607. }
  608. }
  609. else
  610. {
  611. sourceDisposable.Dispose();
  612. }
  613. }
  614. class SelectManyObserver : OperatorObserverBase<TCollection, TResult>
  615. {
  616. readonly SelectManyObserverWithIndex parent;
  617. readonly TSource sourceValue;
  618. readonly int sourceIndex;
  619. readonly IDisposable cancel;
  620. int index;
  621. public SelectManyObserver(SelectManyObserverWithIndex parent, TSource value, int index, IDisposable cancel)
  622. : base(parent.observer, cancel)
  623. {
  624. this.parent = parent;
  625. this.sourceValue = value;
  626. this.sourceIndex = index;
  627. this.cancel = cancel;
  628. }
  629. public override void OnNext(TCollection value)
  630. {
  631. TResult resultValue;
  632. try
  633. {
  634. resultValue = parent.parent.resultSelectorWithIndex(sourceValue, sourceIndex, value, index++);
  635. }
  636. catch (Exception ex)
  637. {
  638. try { observer.OnError(ex); } finally { Dispose(); };
  639. return;
  640. }
  641. lock (parent.gate)
  642. {
  643. base.observer.OnNext(resultValue);
  644. }
  645. }
  646. public override void OnError(Exception error)
  647. {
  648. lock (parent.gate)
  649. {
  650. try { observer.OnError(error); } finally { Dispose(); };
  651. }
  652. }
  653. public override void OnCompleted()
  654. {
  655. parent.collectionDisposable.Remove(cancel);
  656. if (parent.isStopped && parent.collectionDisposable.Count == 1)
  657. {
  658. lock (parent.gate)
  659. {
  660. try { observer.OnCompleted(); } finally { Dispose(); };
  661. }
  662. }
  663. }
  664. }
  665. }
  666. class SelectManyEnumerableObserver : OperatorObserverBase<TSource, TResult>
  667. {
  668. readonly SelectManyObservable<TSource, TCollection, TResult> parent;
  669. public SelectManyEnumerableObserver(SelectManyObservable<TSource, TCollection, TResult> parent, IObserver<TResult> observer, IDisposable cancel) : base(observer, cancel)
  670. {
  671. this.parent = parent;
  672. }
  673. public IDisposable Run()
  674. {
  675. return parent.source.Subscribe(this);
  676. }
  677. public override void OnNext(TSource value)
  678. {
  679. IEnumerable<TCollection> nextEnumerable;
  680. try
  681. {
  682. nextEnumerable = parent.collectionSelectorEnumerable(value);
  683. }
  684. catch (Exception ex)
  685. {
  686. try { observer.OnError(ex); } finally { Dispose(); };
  687. return;
  688. }
  689. var e = nextEnumerable.GetEnumerator();
  690. try
  691. {
  692. var hasNext = true;
  693. while (hasNext)
  694. {
  695. hasNext = false;
  696. var current = default(TResult);
  697. try
  698. {
  699. hasNext = e.MoveNext();
  700. if (hasNext)
  701. {
  702. current = parent.resultSelector(value, e.Current);
  703. }
  704. }
  705. catch (Exception exception)
  706. {
  707. try { observer.OnError(exception); } finally { Dispose(); }
  708. return;
  709. }
  710. if (hasNext)
  711. {
  712. observer.OnNext(current);
  713. }
  714. }
  715. }
  716. finally
  717. {
  718. if (e != null)
  719. {
  720. e.Dispose();
  721. }
  722. }
  723. }
  724. public override void OnError(Exception error)
  725. {
  726. try { observer.OnError(error); } finally { Dispose(); }
  727. }
  728. public override void OnCompleted()
  729. {
  730. try { observer.OnCompleted(); } finally { Dispose(); }
  731. }
  732. }
  733. class SelectManyEnumerableObserverWithIndex : OperatorObserverBase<TSource, TResult>
  734. {
  735. readonly SelectManyObservable<TSource, TCollection, TResult> parent;
  736. int index = 0;
  737. public SelectManyEnumerableObserverWithIndex(SelectManyObservable<TSource, TCollection, TResult> parent, IObserver<TResult> observer, IDisposable cancel) : base(observer, cancel)
  738. {
  739. this.parent = parent;
  740. }
  741. public IDisposable Run()
  742. {
  743. return parent.source.Subscribe(this);
  744. }
  745. public override void OnNext(TSource value)
  746. {
  747. var i = index++;
  748. IEnumerable<TCollection> nextEnumerable;
  749. try
  750. {
  751. nextEnumerable = parent.collectionSelectorEnumerableWithIndex(value, i);
  752. }
  753. catch (Exception ex)
  754. {
  755. try { observer.OnError(ex); } finally { Dispose(); };
  756. return;
  757. }
  758. var e = nextEnumerable.GetEnumerator();
  759. try
  760. {
  761. var sequenceI = 0;
  762. var hasNext = true;
  763. while (hasNext)
  764. {
  765. hasNext = false;
  766. var current = default(TResult);
  767. try
  768. {
  769. hasNext = e.MoveNext();
  770. if (hasNext)
  771. {
  772. current = parent.resultSelectorWithIndex(value, i, e.Current, sequenceI++);
  773. }
  774. }
  775. catch (Exception exception)
  776. {
  777. try { observer.OnError(exception); } finally { Dispose(); }
  778. return;
  779. }
  780. if (hasNext)
  781. {
  782. observer.OnNext(current);
  783. }
  784. }
  785. }
  786. finally
  787. {
  788. if (e != null)
  789. {
  790. e.Dispose();
  791. }
  792. }
  793. }
  794. public override void OnError(Exception error)
  795. {
  796. try { observer.OnError(error); } finally { Dispose(); }
  797. }
  798. public override void OnCompleted()
  799. {
  800. try { observer.OnCompleted(); } finally { Dispose(); }
  801. }
  802. }
  803. }
  804. }