1
0

QueryLanguage.StandardSequenceOperators.cs 51 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340
  1. // Copyright (c) Microsoft Open Technologies, Inc. All rights reserved. See License.txt in the project root for license information.
  2. using System.Collections.Generic;
  3. using System.Linq;
  4. using System.Reactive.Concurrency;
  5. using System.Reactive.Disposables;
  6. using System.Reactive.Subjects;
  7. #if !NO_TPL
  8. using System.Reactive.Threading.Tasks;
  9. using System.Threading;
  10. using System.Threading.Tasks;
  11. #endif
  12. namespace System.Reactive.Linq
  13. {
  14. #if !NO_PERF
  15. using Observαble;
  16. #endif
  17. internal partial class QueryLanguage
  18. {
  19. #region + Cast +
  20. public virtual IObservable<TResult> Cast<TResult>(IObservable<object> source)
  21. {
  22. #if !NO_PERF
  23. return new Cast<object, TResult>(source);
  24. #else
  25. return source.Select(x => (TResult)x);
  26. #endif
  27. }
  28. #endregion
  29. #region + DefaultIfEmpty +
  30. public virtual IObservable<TSource> DefaultIfEmpty<TSource>(IObservable<TSource> source)
  31. {
  32. #if !NO_PERF
  33. return new DefaultIfEmpty<TSource>(source, default(TSource));
  34. #else
  35. return DefaultIfEmpty_(source, default(TSource));
  36. #endif
  37. }
  38. public virtual IObservable<TSource> DefaultIfEmpty<TSource>(IObservable<TSource> source, TSource defaultValue)
  39. {
  40. #if !NO_PERF
  41. return new DefaultIfEmpty<TSource>(source, defaultValue);
  42. #else
  43. return DefaultIfEmpty_(source, defaultValue);
  44. #endif
  45. }
  46. #if NO_PERF
  47. private static IObservable<TSource> DefaultIfEmpty_<TSource>(IObservable<TSource> source, TSource defaultValue)
  48. {
  49. return new AnonymousObservable<TSource>(observer =>
  50. {
  51. var found = false;
  52. return source.Subscribe(
  53. x =>
  54. {
  55. found = true;
  56. observer.OnNext(x);
  57. },
  58. observer.OnError,
  59. () =>
  60. {
  61. if (!found)
  62. observer.OnNext(defaultValue);
  63. observer.OnCompleted();
  64. }
  65. );
  66. });
  67. }
  68. #endif
  69. #endregion
  70. #region + Distinct +
  71. public virtual IObservable<TSource> Distinct<TSource>(IObservable<TSource> source)
  72. {
  73. #if !NO_PERF
  74. return new Distinct<TSource, TSource>(source, x => x, EqualityComparer<TSource>.Default);
  75. #else
  76. return Distinct_(source, x => x, EqualityComparer<TSource>.Default);
  77. #endif
  78. }
  79. public virtual IObservable<TSource> Distinct<TSource>(IObservable<TSource> source, IEqualityComparer<TSource> comparer)
  80. {
  81. #if !NO_PERF
  82. return new Distinct<TSource, TSource>(source, x => x, comparer);
  83. #else
  84. return Distinct_(source, x => x, comparer);
  85. #endif
  86. }
  87. public virtual IObservable<TSource> Distinct<TSource, TKey>(IObservable<TSource> source, Func<TSource, TKey> keySelector)
  88. {
  89. #if !NO_PERF
  90. return new Distinct<TSource, TKey>(source, keySelector, EqualityComparer<TKey>.Default);
  91. #else
  92. return Distinct_(source, keySelector, EqualityComparer<TKey>.Default);
  93. #endif
  94. }
  95. public virtual IObservable<TSource> Distinct<TSource, TKey>(IObservable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
  96. {
  97. #if !NO_PERF
  98. return new Distinct<TSource, TKey>(source, keySelector, comparer);
  99. #else
  100. return Distinct_(source, keySelector, comparer);
  101. #endif
  102. }
  103. #if NO_PERF
  104. private static IObservable<TSource> Distinct_<TSource, TKey>(IObservable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
  105. {
  106. return new AnonymousObservable<TSource>(observer =>
  107. {
  108. var hashSet = new HashSet<TKey>(comparer);
  109. return source.Subscribe(
  110. x =>
  111. {
  112. var key = default(TKey);
  113. var hasAdded = false;
  114. try
  115. {
  116. key = keySelector(x);
  117. hasAdded = hashSet.Add(key);
  118. }
  119. catch (Exception exception)
  120. {
  121. observer.OnError(exception);
  122. return;
  123. }
  124. if (hasAdded)
  125. observer.OnNext(x);
  126. },
  127. observer.OnError,
  128. observer.OnCompleted
  129. );
  130. });
  131. }
  132. #endif
  133. #endregion
  134. #region + GroupBy +
  135. public virtual IObservable<IGroupedObservable<TKey, TElement>> GroupBy<TSource, TKey, TElement>(IObservable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector)
  136. {
  137. return GroupBy_<TSource, TKey, TElement>(source, keySelector, elementSelector, EqualityComparer<TKey>.Default);
  138. }
  139. public virtual IObservable<IGroupedObservable<TKey, TSource>> GroupBy<TSource, TKey>(IObservable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
  140. {
  141. return GroupBy_<TSource, TKey, TSource>(source, keySelector, x => x, comparer);
  142. }
  143. public virtual IObservable<IGroupedObservable<TKey, TSource>> GroupBy<TSource, TKey>(IObservable<TSource> source, Func<TSource, TKey> keySelector)
  144. {
  145. return GroupBy_<TSource, TKey, TSource>(source, keySelector, x => x, EqualityComparer<TKey>.Default);
  146. }
  147. public virtual IObservable<IGroupedObservable<TKey, TElement>> GroupBy<TSource, TKey, TElement>(IObservable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, IEqualityComparer<TKey> comparer)
  148. {
  149. return GroupBy_<TSource, TKey, TElement>(source, keySelector, elementSelector, comparer);
  150. }
  151. private static IObservable<IGroupedObservable<TKey, TElement>> GroupBy_<TSource, TKey, TElement>(IObservable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, IEqualityComparer<TKey> comparer)
  152. {
  153. #if !NO_PERF
  154. return new GroupBy<TSource, TKey, TElement>(source, keySelector, elementSelector, comparer);
  155. #else
  156. return GroupByUntil_<TSource, TKey, TElement, Unit>(source, keySelector, elementSelector, _ => Observable.Never<Unit>(), comparer);
  157. #endif
  158. }
  159. #endregion
  160. #region + GroupByUntil +
  161. public virtual IObservable<IGroupedObservable<TKey, TElement>> GroupByUntil<TSource, TKey, TElement, TDuration>(IObservable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, Func<IGroupedObservable<TKey, TElement>, IObservable<TDuration>> durationSelector, IEqualityComparer<TKey> comparer)
  162. {
  163. return GroupByUntil_<TSource, TKey, TElement, TDuration>(source, keySelector, elementSelector, durationSelector, comparer);
  164. }
  165. public virtual IObservable<IGroupedObservable<TKey, TElement>> GroupByUntil<TSource, TKey, TElement, TDuration>(IObservable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, Func<IGroupedObservable<TKey, TElement>, IObservable<TDuration>> durationSelector)
  166. {
  167. return GroupByUntil_<TSource, TKey, TElement, TDuration>(source, keySelector, elementSelector, durationSelector, EqualityComparer<TKey>.Default);
  168. }
  169. public virtual IObservable<IGroupedObservable<TKey, TSource>> GroupByUntil<TSource, TKey, TDuration>(IObservable<TSource> source, Func<TSource, TKey> keySelector, Func<IGroupedObservable<TKey, TSource>, IObservable<TDuration>> durationSelector, IEqualityComparer<TKey> comparer)
  170. {
  171. return GroupByUntil_<TSource, TKey, TSource, TDuration>(source, keySelector, x => x, durationSelector, comparer);
  172. }
  173. public virtual IObservable<IGroupedObservable<TKey, TSource>> GroupByUntil<TSource, TKey, TDuration>(IObservable<TSource> source, Func<TSource, TKey> keySelector, Func<IGroupedObservable<TKey, TSource>, IObservable<TDuration>> durationSelector)
  174. {
  175. return GroupByUntil_<TSource, TKey, TSource, TDuration>(source, keySelector, x => x, durationSelector, EqualityComparer<TKey>.Default);
  176. }
  177. private static IObservable<IGroupedObservable<TKey, TElement>> GroupByUntil_<TSource, TKey, TElement, TDuration>(IObservable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, Func<IGroupedObservable<TKey, TElement>, IObservable<TDuration>> durationSelector, IEqualityComparer<TKey> comparer)
  178. {
  179. #if !NO_PERF
  180. return new GroupByUntil<TSource, TKey, TElement, TDuration>(source, keySelector, elementSelector, durationSelector, comparer);
  181. #else
  182. return new AnonymousObservable<IGroupedObservable<TKey, TElement>>(observer =>
  183. {
  184. var map = new Dictionary<TKey, ISubject<TElement>>(comparer);
  185. var groupDisposable = new CompositeDisposable();
  186. var refCountDisposable = new RefCountDisposable(groupDisposable);
  187. groupDisposable.Add(source.Subscribe(x =>
  188. {
  189. var key = default(TKey);
  190. try
  191. {
  192. key = keySelector(x);
  193. }
  194. catch (Exception exception)
  195. {
  196. lock (map)
  197. foreach (var w in map.Values.ToArray())
  198. w.OnError(exception);
  199. observer.OnError(exception);
  200. return;
  201. }
  202. var fireNewMapEntry = false;
  203. var writer = default(ISubject<TElement>);
  204. try
  205. {
  206. lock (map)
  207. {
  208. if (!map.TryGetValue(key, out writer))
  209. {
  210. writer = new Subject<TElement>();
  211. map.Add(key, writer);
  212. fireNewMapEntry = true;
  213. }
  214. }
  215. }
  216. catch (Exception exception)
  217. {
  218. lock (map)
  219. {
  220. foreach (var w in map.Values.ToArray())
  221. w.OnError(exception);
  222. }
  223. observer.OnError(exception);
  224. return;
  225. }
  226. if (fireNewMapEntry)
  227. {
  228. var group = new GroupedObservable<TKey, TElement>(key, writer, refCountDisposable);
  229. var durationGroup = new GroupedObservable<TKey, TElement>(key, writer);
  230. var duration = default(IObservable<TDuration>);
  231. try
  232. {
  233. duration = durationSelector(durationGroup);
  234. }
  235. catch (Exception exception)
  236. {
  237. foreach (var w in map.Values.ToArray())
  238. w.OnError(exception);
  239. observer.OnError(exception);
  240. return;
  241. }
  242. observer.OnNext(group);
  243. var md = new SingleAssignmentDisposable();
  244. groupDisposable.Add(md);
  245. Action expire = () =>
  246. {
  247. lock (map)
  248. {
  249. if (map.Remove(key))
  250. writer.OnCompleted();
  251. }
  252. groupDisposable.Remove(md);
  253. };
  254. md.Disposable = duration.Take(1).Subscribe(
  255. _ => { },
  256. exception =>
  257. {
  258. lock (map)
  259. foreach (var o in map.Values.ToArray())
  260. o.OnError(exception);
  261. observer.OnError(exception);
  262. },
  263. expire);
  264. }
  265. var element = default(TElement);
  266. try
  267. {
  268. element = elementSelector(x);
  269. }
  270. catch (Exception exception)
  271. {
  272. lock (map)
  273. foreach (var w in map.Values.ToArray())
  274. w.OnError(exception);
  275. observer.OnError(exception);
  276. return;
  277. }
  278. writer.OnNext(element);
  279. },
  280. e =>
  281. {
  282. lock (map)
  283. foreach (var w in map.Values.ToArray())
  284. w.OnError(e);
  285. observer.OnError(e);
  286. },
  287. () =>
  288. {
  289. lock (map)
  290. foreach (var w in map.Values.ToArray())
  291. w.OnCompleted();
  292. observer.OnCompleted();
  293. }));
  294. return refCountDisposable;
  295. });
  296. #endif
  297. }
  298. #endregion
  299. #region + GroupJoin +
  300. public virtual IObservable<TResult> GroupJoin<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(IObservable<TLeft> left, IObservable<TRight> right, Func<TLeft, IObservable<TLeftDuration>> leftDurationSelector, Func<TRight, IObservable<TRightDuration>> rightDurationSelector, Func<TLeft, IObservable<TRight>, TResult> resultSelector)
  301. {
  302. return GroupJoin_<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(left, right, leftDurationSelector, rightDurationSelector, resultSelector);
  303. }
  304. private static IObservable<TResult> GroupJoin_<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(IObservable<TLeft> left, IObservable<TRight> right, Func<TLeft, IObservable<TLeftDuration>> leftDurationSelector, Func<TRight, IObservable<TRightDuration>> rightDurationSelector, Func<TLeft, IObservable<TRight>, TResult> resultSelector)
  305. {
  306. #if !NO_PERF
  307. return new GroupJoin<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(left, right, leftDurationSelector, rightDurationSelector, resultSelector);
  308. #else
  309. return new AnonymousObservable<TResult>(observer =>
  310. {
  311. var gate = new object();
  312. var group = new CompositeDisposable();
  313. var r = new RefCountDisposable(group);
  314. var leftMap = new Dictionary<int, IObserver<TRight>>();
  315. var rightMap = new Dictionary<int, TRight>();
  316. var leftID = 0;
  317. var rightID = 0;
  318. group.Add(left.Subscribe(
  319. value =>
  320. {
  321. var s = new Subject<TRight>();
  322. var id = 0;
  323. lock (gate)
  324. {
  325. id = leftID++;
  326. leftMap.Add(id, s);
  327. }
  328. lock (gate)
  329. {
  330. var result = default(TResult);
  331. try
  332. {
  333. result = resultSelector(value, s.AddRef(r));
  334. }
  335. catch (Exception exception)
  336. {
  337. foreach (var o in leftMap.Values.ToArray())
  338. o.OnError(exception);
  339. observer.OnError(exception);
  340. return;
  341. }
  342. observer.OnNext(result);
  343. foreach (var rightValue in rightMap.Values.ToArray())
  344. {
  345. s.OnNext(rightValue);
  346. }
  347. }
  348. var md = new SingleAssignmentDisposable();
  349. group.Add(md);
  350. Action expire = () =>
  351. {
  352. lock (gate)
  353. if (leftMap.Remove(id))
  354. s.OnCompleted();
  355. group.Remove(md);
  356. };
  357. var duration = default(IObservable<TLeftDuration>);
  358. try
  359. {
  360. duration = leftDurationSelector(value);
  361. }
  362. catch (Exception exception)
  363. {
  364. lock (gate)
  365. {
  366. foreach (var o in leftMap.Values.ToArray())
  367. o.OnError(exception);
  368. observer.OnError(exception);
  369. }
  370. return;
  371. }
  372. md.Disposable = duration.Take(1).Subscribe(
  373. _ => { },
  374. exception =>
  375. {
  376. lock (gate)
  377. {
  378. foreach (var o in leftMap.Values.ToArray())
  379. o.OnError(exception);
  380. observer.OnError(exception);
  381. }
  382. },
  383. expire);
  384. },
  385. exception =>
  386. {
  387. lock (gate)
  388. {
  389. foreach (var o in leftMap.Values.ToArray())
  390. o.OnError(exception);
  391. observer.OnError(exception);
  392. }
  393. },
  394. () =>
  395. {
  396. lock (gate)
  397. observer.OnCompleted();
  398. }));
  399. group.Add(right.Subscribe(
  400. value =>
  401. {
  402. var id = 0;
  403. lock (gate)
  404. {
  405. id = rightID++;
  406. rightMap.Add(id, value);
  407. }
  408. var md = new SingleAssignmentDisposable();
  409. group.Add(md);
  410. Action expire = () =>
  411. {
  412. lock (gate)
  413. rightMap.Remove(id);
  414. group.Remove(md);
  415. };
  416. var duration = default(IObservable<TRightDuration>);
  417. try
  418. {
  419. duration = rightDurationSelector(value);
  420. }
  421. catch (Exception exception)
  422. {
  423. lock (gate)
  424. {
  425. foreach (var o in leftMap.Values.ToArray())
  426. o.OnError(exception);
  427. observer.OnError(exception);
  428. }
  429. return;
  430. }
  431. md.Disposable = duration.Take(1).Subscribe(
  432. _ => { },
  433. exception =>
  434. {
  435. lock (gate)
  436. {
  437. foreach (var o in leftMap.Values.ToArray())
  438. o.OnError(exception);
  439. observer.OnError(exception);
  440. }
  441. },
  442. expire);
  443. lock (gate)
  444. {
  445. foreach (var o in leftMap.Values.ToArray())
  446. o.OnNext(value);
  447. }
  448. },
  449. exception =>
  450. {
  451. lock (gate)
  452. {
  453. foreach (var o in leftMap.Values.ToArray())
  454. o.OnError(exception);
  455. observer.OnError(exception);
  456. }
  457. }));
  458. return r;
  459. });
  460. #endif
  461. }
  462. #endregion
  463. #region + Join +
  464. public virtual IObservable<TResult> Join<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(IObservable<TLeft> left, IObservable<TRight> right, Func<TLeft, IObservable<TLeftDuration>> leftDurationSelector, Func<TRight, IObservable<TRightDuration>> rightDurationSelector, Func<TLeft, TRight, TResult> resultSelector)
  465. {
  466. return Join_<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(left, right, leftDurationSelector, rightDurationSelector, resultSelector);
  467. }
  468. private static IObservable<TResult> Join_<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(IObservable<TLeft> left, IObservable<TRight> right, Func<TLeft, IObservable<TLeftDuration>> leftDurationSelector, Func<TRight, IObservable<TRightDuration>> rightDurationSelector, Func<TLeft, TRight, TResult> resultSelector)
  469. {
  470. #if !NO_PERF
  471. return new Join<TLeft, TRight, TLeftDuration, TRightDuration, TResult>(left, right, leftDurationSelector, rightDurationSelector, resultSelector);
  472. #else
  473. return new AnonymousObservable<TResult>(observer =>
  474. {
  475. var gate = new object();
  476. var leftDone = false;
  477. var rightDone = false;
  478. var group = new CompositeDisposable();
  479. var leftMap = new Dictionary<int, TLeft>();
  480. var rightMap = new Dictionary<int, TRight>();
  481. var leftID = 0;
  482. var rightID = 0;
  483. group.Add(left.Subscribe(
  484. value =>
  485. {
  486. var id = 0;
  487. lock (gate)
  488. {
  489. id = leftID++;
  490. leftMap.Add(id, value);
  491. }
  492. var md = new SingleAssignmentDisposable();
  493. group.Add(md);
  494. Action expire = () =>
  495. {
  496. lock (gate)
  497. {
  498. if (leftMap.Remove(id) && leftMap.Count == 0 && leftDone)
  499. observer.OnCompleted();
  500. }
  501. group.Remove(md);
  502. };
  503. var duration = default(IObservable<TLeftDuration>);
  504. try
  505. {
  506. duration = leftDurationSelector(value);
  507. }
  508. catch (Exception exception)
  509. {
  510. observer.OnError(exception);
  511. return;
  512. }
  513. md.Disposable = duration.Take(1).Subscribe(
  514. _ => { },
  515. error =>
  516. {
  517. lock (gate)
  518. observer.OnError(error);
  519. },
  520. expire);
  521. lock (gate)
  522. {
  523. foreach (var rightValue in rightMap.Values.ToArray())
  524. {
  525. var result = default(TResult);
  526. try
  527. {
  528. result = resultSelector(value, rightValue);
  529. }
  530. catch (Exception exception)
  531. {
  532. observer.OnError(exception);
  533. return;
  534. }
  535. observer.OnNext(result);
  536. }
  537. }
  538. },
  539. error =>
  540. {
  541. lock (gate)
  542. observer.OnError(error);
  543. },
  544. () =>
  545. {
  546. lock (gate)
  547. {
  548. leftDone = true;
  549. if (rightDone || leftMap.Count == 0)
  550. observer.OnCompleted();
  551. }
  552. }));
  553. group.Add(right.Subscribe(
  554. value =>
  555. {
  556. var id = 0;
  557. lock (gate)
  558. {
  559. id = rightID++;
  560. rightMap.Add(id, value);
  561. }
  562. var md = new SingleAssignmentDisposable();
  563. group.Add(md);
  564. Action expire = () =>
  565. {
  566. lock (gate)
  567. {
  568. if (rightMap.Remove(id) && rightMap.Count == 0 && rightDone)
  569. observer.OnCompleted();
  570. }
  571. group.Remove(md);
  572. };
  573. var duration = default(IObservable<TRightDuration>);
  574. try
  575. {
  576. duration = rightDurationSelector(value);
  577. }
  578. catch (Exception exception)
  579. {
  580. observer.OnError(exception);
  581. return;
  582. }
  583. md.Disposable = duration.Take(1).Subscribe(
  584. _ => { },
  585. error =>
  586. {
  587. lock (gate)
  588. observer.OnError(error);
  589. },
  590. expire);
  591. lock (gate)
  592. {
  593. foreach (var leftValue in leftMap.Values.ToArray())
  594. {
  595. var result = default(TResult);
  596. try
  597. {
  598. result = resultSelector(leftValue, value);
  599. }
  600. catch (Exception exception)
  601. {
  602. observer.OnError(exception);
  603. return;
  604. }
  605. observer.OnNext(result);
  606. }
  607. }
  608. },
  609. error =>
  610. {
  611. lock (gate)
  612. observer.OnError(error);
  613. },
  614. () =>
  615. {
  616. lock (gate)
  617. {
  618. rightDone = true;
  619. if (leftDone || rightMap.Count == 0)
  620. observer.OnCompleted();
  621. }
  622. }));
  623. return group;
  624. });
  625. #endif
  626. }
  627. #endregion
  628. #region + OfType +
  629. public virtual IObservable<TResult> OfType<TResult>(IObservable<object> source)
  630. {
  631. #if !NO_PERF
  632. return new OfType<object, TResult>(source);
  633. #else
  634. return source.Where(x => x is TResult).Cast<TResult>();
  635. #endif
  636. }
  637. #endregion
  638. #region + Select +
  639. public virtual IObservable<TResult> Select<TSource, TResult>(IObservable<TSource> source, Func<TSource, TResult> selector)
  640. {
  641. #if !NO_PERF
  642. var select = source as Select<TSource>;
  643. if (select != null)
  644. return select.Ω(selector);
  645. return new Select<TSource, TResult>(source, selector);
  646. #else
  647. var s = source as SelectObservable<TSource>;
  648. if (s != null)
  649. return s.Select(selector);
  650. return new SelectObservable<TSource, TResult>(source, selector);
  651. #endif
  652. }
  653. #if NO_PERF
  654. abstract class SelectObservable<TResult> : ObservableBase<TResult>
  655. {
  656. public abstract IObservable<TResult2> Select<TResult2>(Func<TResult, TResult2> selector);
  657. }
  658. class SelectObservable<TSource, TResult> : SelectObservable<TResult>
  659. {
  660. private readonly IObservable<TSource> _source;
  661. private readonly Func<TSource, TResult> _selector;
  662. public SelectObservable(IObservable<TSource> source, Func<TSource, TResult> selector)
  663. {
  664. _source = source;
  665. _selector = selector;
  666. }
  667. protected override IDisposable SubscribeCore(IObserver<TResult> observer)
  668. {
  669. return _source.Subscribe(new Observer(observer, _selector));
  670. }
  671. public override IObservable<TResult2> Select<TResult2>(Func<TResult, TResult2> selector)
  672. {
  673. return new SelectObservable<TSource, TResult2>(_source, x => selector(_selector(x)));
  674. }
  675. class Observer : ObserverBase<TSource>
  676. {
  677. private readonly IObserver<TResult> _observer;
  678. private readonly Func<TSource, TResult> _selector;
  679. public Observer(IObserver<TResult> observer, Func<TSource, TResult> selector)
  680. {
  681. _observer = observer;
  682. _selector = selector;
  683. }
  684. protected override void OnNextCore(TSource value)
  685. {
  686. TResult result;
  687. try
  688. {
  689. result = _selector(value);
  690. }
  691. catch (Exception exception)
  692. {
  693. _observer.OnError(exception);
  694. return;
  695. }
  696. _observer.OnNext(result);
  697. }
  698. protected override void OnErrorCore(Exception error)
  699. {
  700. _observer.OnError(error);
  701. }
  702. protected override void OnCompletedCore()
  703. {
  704. _observer.OnCompleted();
  705. }
  706. }
  707. }
  708. #endif
  709. public virtual IObservable<TResult> Select<TSource, TResult>(IObservable<TSource> source, Func<TSource, int, TResult> selector)
  710. {
  711. #if !NO_PERF
  712. return new Select<TSource, TResult>(source, selector);
  713. #else
  714. return Defer(() =>
  715. {
  716. var index = 0;
  717. return source.Select(x => selector(x, checked(index++)));
  718. });
  719. #endif
  720. }
  721. #endregion
  722. #region + SelectMany +
  723. public virtual IObservable<TOther> SelectMany<TSource, TOther>(IObservable<TSource> source, IObservable<TOther> other)
  724. {
  725. return SelectMany_<TSource, TOther>(source, _ => other);
  726. }
  727. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TResult>> selector)
  728. {
  729. return SelectMany_<TSource, TResult>(source, selector);
  730. }
  731. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, int, IObservable<TResult>> selector)
  732. {
  733. return SelectMany_<TSource, TResult>(source, selector);
  734. }
  735. #if !NO_TPL
  736. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, Task<TResult>> selector)
  737. {
  738. #if !NO_PERF
  739. return new SelectMany<TSource, TResult>(source, (x, token) => selector(x));
  740. #else
  741. return SelectMany_<TSource, TResult>(source, x => selector(x).ToObservable());
  742. #endif
  743. }
  744. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, CancellationToken, Task<TResult>> selector)
  745. {
  746. #if !NO_PERF
  747. return new SelectMany<TSource, TResult>(source, selector);
  748. #else
  749. return SelectMany_<TSource, TResult>(source, x => FromAsync(ct => selector(x, ct)));
  750. #endif
  751. }
  752. #endif
  753. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  754. {
  755. return SelectMany_<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  756. }
  757. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, int, IObservable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  758. {
  759. return SelectMany_<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  760. }
  761. #if !NO_TPL
  762. public virtual IObservable<TResult> SelectMany<TSource, TTaskResult, TResult>(IObservable<TSource> source, Func<TSource, Task<TTaskResult>> taskSelector, Func<TSource, TTaskResult, TResult> resultSelector)
  763. {
  764. #if !NO_PERF
  765. return new SelectMany<TSource, TTaskResult, TResult>(source, (x, token) => taskSelector(x), resultSelector);
  766. #else
  767. return SelectMany_<TSource, TTaskResult, TResult>(source, x => taskSelector(x).ToObservable(), resultSelector);
  768. #endif
  769. }
  770. public virtual IObservable<TResult> SelectMany<TSource, TTaskResult, TResult>(IObservable<TSource> source, Func<TSource, CancellationToken, Task<TTaskResult>> taskSelector, Func<TSource, TTaskResult, TResult> resultSelector)
  771. {
  772. #if !NO_PERF
  773. return new SelectMany<TSource, TTaskResult, TResult>(source, taskSelector, resultSelector);
  774. #else
  775. return SelectMany_<TSource, TTaskResult, TResult>(source, x => FromAsync(ct => taskSelector(x, ct)), resultSelector);
  776. #endif
  777. }
  778. #endif
  779. private static IObservable<TResult> SelectMany_<TSource, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TResult>> selector)
  780. {
  781. #if !NO_PERF
  782. return new SelectMany<TSource, TResult>(source, selector);
  783. #else
  784. return source.Select(selector).Merge();
  785. #endif
  786. }
  787. private static IObservable<TResult> SelectMany_<TSource, TResult>(IObservable<TSource> source, Func<TSource, int, IObservable<TResult>> selector)
  788. {
  789. #if !NO_PERF
  790. return new SelectMany<TSource, TResult>(source, selector);
  791. #else
  792. return source.Select(selector).Merge();
  793. #endif
  794. }
  795. private static IObservable<TResult> SelectMany_<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  796. {
  797. #if !NO_PERF
  798. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  799. #else
  800. return SelectMany_<TSource, TResult>(source, x => collectionSelector(x).Select(y => resultSelector(x, y)));
  801. #endif
  802. }
  803. private static IObservable<TResult> SelectMany_<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, int, IObservable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  804. {
  805. #if !NO_PERF
  806. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  807. #else
  808. return SelectMany_<TSource, TResult>(source, x => collectionSelector(x).Select(y => resultSelector(x, y)));
  809. #endif
  810. }
  811. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TResult>> onNext, Func<Exception, IObservable<TResult>> onError, Func<IObservable<TResult>> onCompleted)
  812. {
  813. #if !NO_PERF
  814. return new SelectMany<TSource, TResult>(source, onNext, onError, onCompleted);
  815. #else
  816. return source.Materialize().SelectMany(notification =>
  817. {
  818. if (notification.Kind == NotificationKind.OnNext)
  819. return onNext(notification.Value);
  820. else if (notification.Kind == NotificationKind.OnError)
  821. return onError(notification.Exception);
  822. else
  823. return onCompleted();
  824. });
  825. #endif
  826. }
  827. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, int, IObservable<TResult>> onNext, Func<Exception, int, IObservable<TResult>> onError, Func<int, IObservable<TResult>> onCompleted)
  828. {
  829. #if !NO_PERF
  830. return new SelectMany<TSource, TResult>(source, onNext, onError, onCompleted);
  831. #else
  832. return source.Materialize().SelectMany(notification =>
  833. {
  834. if (notification.Kind == NotificationKind.OnNext)
  835. return onNext(notification.Value);
  836. else if (notification.Kind == NotificationKind.OnError)
  837. return onError(notification.Exception);
  838. else
  839. return onCompleted();
  840. });
  841. #endif
  842. }
  843. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, IEnumerable<TResult>> selector)
  844. {
  845. #if !NO_PERF
  846. return new SelectMany<TSource, TResult>(source, selector);
  847. #else
  848. return SelectMany_<TSource, TResult, TResult>(source, selector, (_, x) => x);
  849. #endif
  850. }
  851. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, int, IEnumerable<TResult>> selector)
  852. {
  853. #if !NO_PERF
  854. return new SelectMany<TSource, TResult>(source, selector);
  855. #else
  856. return SelectMany_<TSource, TResult, TResult>(source, selector, (_, x) => x);
  857. #endif
  858. }
  859. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IEnumerable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  860. {
  861. return SelectMany_<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  862. }
  863. private static IObservable<TResult> SelectMany_<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IEnumerable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  864. {
  865. #if !NO_PERF
  866. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  867. #else
  868. return new AnonymousObservable<TResult>(observer =>
  869. source.Subscribe(
  870. x =>
  871. {
  872. var xs = default(IEnumerable<TCollection>);
  873. try
  874. {
  875. xs = collectionSelector(x);
  876. }
  877. catch (Exception exception)
  878. {
  879. observer.OnError(exception);
  880. return;
  881. }
  882. var e = xs.GetEnumerator();
  883. try
  884. {
  885. var hasNext = true;
  886. while (hasNext)
  887. {
  888. hasNext = false;
  889. var current = default(TResult);
  890. try
  891. {
  892. hasNext = e.MoveNext();
  893. if (hasNext)
  894. current = resultSelector(x, e.Current);
  895. }
  896. catch (Exception exception)
  897. {
  898. observer.OnError(exception);
  899. return;
  900. }
  901. if (hasNext)
  902. observer.OnNext(current);
  903. }
  904. }
  905. finally
  906. {
  907. if (e != null)
  908. e.Dispose();
  909. }
  910. },
  911. observer.OnError,
  912. observer.OnCompleted
  913. )
  914. );
  915. #endif
  916. }
  917. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, int, IEnumerable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  918. {
  919. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  920. }
  921. #endregion
  922. #region + Skip +
  923. public virtual IObservable<TSource> Skip<TSource>(IObservable<TSource> source, int count)
  924. {
  925. #if !NO_PERF
  926. var skip = source as Skip<TSource>;
  927. if (skip != null && skip._scheduler == null)
  928. return skip.Ω(count);
  929. return new Skip<TSource>(source, count);
  930. #else
  931. return new AnonymousObservable<TSource>(observer =>
  932. {
  933. var remaining = count;
  934. return source.Subscribe(
  935. x =>
  936. {
  937. if (remaining <= 0)
  938. observer.OnNext(x);
  939. else
  940. remaining--;
  941. },
  942. observer.OnError,
  943. observer.OnCompleted);
  944. });
  945. #endif
  946. }
  947. #endregion
  948. #region + SkipWhile +
  949. public virtual IObservable<TSource> SkipWhile<TSource>(IObservable<TSource> source, Func<TSource, bool> predicate)
  950. {
  951. #if !NO_PERF
  952. return new SkipWhile<TSource>(source, predicate);
  953. #else
  954. return SkipWhile_(source, (x, i) => predicate(x));
  955. #endif
  956. }
  957. public virtual IObservable<TSource> SkipWhile<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  958. {
  959. #if !NO_PERF
  960. return new SkipWhile<TSource>(source, predicate);
  961. #else
  962. return SkipWhile_(source, predicate);
  963. #endif
  964. }
  965. #if NO_PERF
  966. private static IObservable<TSource> SkipWhile_<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  967. {
  968. return new AnonymousObservable<TSource>(observer =>
  969. {
  970. var running = false;
  971. var i = 0;
  972. return source.Subscribe(
  973. x =>
  974. {
  975. if (!running)
  976. try
  977. {
  978. running = !predicate(x, checked(i++));
  979. }
  980. catch (Exception exception)
  981. {
  982. observer.OnError(exception);
  983. return;
  984. }
  985. if (running)
  986. observer.OnNext(x);
  987. },
  988. observer.OnError,
  989. observer.OnCompleted);
  990. });
  991. }
  992. #endif
  993. #endregion
  994. #region + Take +
  995. public virtual IObservable<TSource> Take<TSource>(IObservable<TSource> source, int count)
  996. {
  997. if (count == 0)
  998. return Empty<TSource>();
  999. return Take_(source, count);
  1000. }
  1001. public virtual IObservable<TSource> Take<TSource>(IObservable<TSource> source, int count, IScheduler scheduler)
  1002. {
  1003. if (count == 0)
  1004. return Empty<TSource>(scheduler);
  1005. return Take_(source, count);
  1006. }
  1007. #if !NO_PERF
  1008. private static IObservable<TSource> Take_<TSource>(IObservable<TSource> source, int count)
  1009. {
  1010. var take = source as Take<TSource>;
  1011. if (take != null && take._scheduler == null)
  1012. return take.Ω(count);
  1013. return new Take<TSource>(source, count);
  1014. }
  1015. #else
  1016. private static IObservable<TSource> Take_<TSource>(IObservable<TSource> source, int count)
  1017. {
  1018. return new AnonymousObservable<TSource>(observer =>
  1019. {
  1020. var remaining = count;
  1021. return source.Subscribe(
  1022. x =>
  1023. {
  1024. if (remaining > 0)
  1025. {
  1026. --remaining;
  1027. observer.OnNext(x);
  1028. if (remaining == 0)
  1029. observer.OnCompleted();
  1030. }
  1031. },
  1032. observer.OnError,
  1033. observer.OnCompleted);
  1034. });
  1035. }
  1036. #endif
  1037. #endregion
  1038. #region + TakeWhile +
  1039. public virtual IObservable<TSource> TakeWhile<TSource>(IObservable<TSource> source, Func<TSource, bool> predicate)
  1040. {
  1041. #if !NO_PERF
  1042. return new TakeWhile<TSource>(source, predicate);
  1043. #else
  1044. return TakeWhile_(source, (x, i) => predicate(x));
  1045. #endif
  1046. }
  1047. public virtual IObservable<TSource> TakeWhile<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  1048. {
  1049. #if !NO_PERF
  1050. return new TakeWhile<TSource>(source, predicate);
  1051. #else
  1052. return TakeWhile_(source, predicate);
  1053. #endif
  1054. }
  1055. #if NO_PERF
  1056. private static IObservable<TSource> TakeWhile_<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  1057. {
  1058. return new AnonymousObservable<TSource>(observer =>
  1059. {
  1060. var running = true;
  1061. var i = 0;
  1062. return source.Subscribe(
  1063. x =>
  1064. {
  1065. if (running)
  1066. {
  1067. try
  1068. {
  1069. running = predicate(x, checked(i++));
  1070. }
  1071. catch (Exception exception)
  1072. {
  1073. observer.OnError(exception);
  1074. return;
  1075. }
  1076. if (running)
  1077. observer.OnNext(x);
  1078. else
  1079. observer.OnCompleted();
  1080. }
  1081. },
  1082. observer.OnError,
  1083. observer.OnCompleted);
  1084. });
  1085. }
  1086. #endif
  1087. #endregion
  1088. #region + Where +
  1089. public virtual IObservable<TSource> Where<TSource>(IObservable<TSource> source, Func<TSource, bool> predicate)
  1090. {
  1091. #if !NO_PERF
  1092. var where = source as Where<TSource>;
  1093. if (where != null)
  1094. return where.Ω(predicate);
  1095. return new Where<TSource>(source, predicate);
  1096. #else
  1097. var w = source as WhereObservable<TSource>;
  1098. if (w != null)
  1099. return w.Where(predicate);
  1100. return new WhereObservable<TSource>(source, predicate);
  1101. #endif
  1102. }
  1103. #if NO_PERF
  1104. class WhereObservable<TSource> : ObservableBase<TSource>
  1105. {
  1106. private readonly IObservable<TSource> _source;
  1107. private readonly Func<TSource, bool> _predicate;
  1108. public WhereObservable(IObservable<TSource> source, Func<TSource, bool> predicate)
  1109. {
  1110. _source = source;
  1111. _predicate = predicate;
  1112. }
  1113. protected override IDisposable SubscribeCore(IObserver<TSource> observer)
  1114. {
  1115. return _source.Subscribe(new Observer(observer, _predicate));
  1116. }
  1117. public IObservable<TSource> Where(Func<TSource, bool> predicate)
  1118. {
  1119. return new WhereObservable<TSource>(_source, x => _predicate(x) && predicate(x));
  1120. }
  1121. class Observer : ObserverBase<TSource>
  1122. {
  1123. private readonly IObserver<TSource> _observer;
  1124. private readonly Func<TSource, bool> _predicate;
  1125. public Observer(IObserver<TSource> observer, Func<TSource, bool> predicate)
  1126. {
  1127. _observer = observer;
  1128. _predicate = predicate;
  1129. }
  1130. protected override void OnNextCore(TSource value)
  1131. {
  1132. bool shouldRun;
  1133. try
  1134. {
  1135. shouldRun = _predicate(value);
  1136. }
  1137. catch (Exception exception)
  1138. {
  1139. _observer.OnError(exception);
  1140. return;
  1141. }
  1142. if (shouldRun)
  1143. _observer.OnNext(value);
  1144. }
  1145. protected override void OnErrorCore(Exception error)
  1146. {
  1147. _observer.OnError(error);
  1148. }
  1149. protected override void OnCompletedCore()
  1150. {
  1151. _observer.OnCompleted();
  1152. }
  1153. }
  1154. }
  1155. #endif
  1156. public virtual IObservable<TSource> Where<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  1157. {
  1158. #if !NO_PERF
  1159. return new Where<TSource>(source, predicate);
  1160. #else
  1161. return Defer(() =>
  1162. {
  1163. var index = 0;
  1164. return source.Where(x => predicate(x, checked(index++)));
  1165. });
  1166. #endif
  1167. }
  1168. #endregion
  1169. }
  1170. }