1
0

QueryLanguage.StandardSequenceOperators.cs 47 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265
  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. #if !NO_TPL
  732. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, Task<TResult>> selector)
  733. {
  734. return SelectMany_<TSource, TResult>(source, x => selector(x).ToObservable());
  735. }
  736. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, CancellationToken, Task<TResult>> selector)
  737. {
  738. return SelectMany_<TSource, TResult>(source, x => FromAsync(ct => selector(x, ct)));
  739. }
  740. #endif
  741. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  742. {
  743. return SelectMany_<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  744. }
  745. #if !NO_TPL
  746. public virtual IObservable<TResult> SelectMany<TSource, TTaskResult, TResult>(IObservable<TSource> source, Func<TSource, Task<TTaskResult>> taskSelector, Func<TSource, TTaskResult, TResult> resultSelector)
  747. {
  748. return SelectMany_<TSource, TTaskResult, TResult>(source, x => taskSelector(x).ToObservable(), resultSelector);
  749. }
  750. public virtual IObservable<TResult> SelectMany<TSource, TTaskResult, TResult>(IObservable<TSource> source, Func<TSource, CancellationToken, Task<TTaskResult>> taskSelector, Func<TSource, TTaskResult, TResult> resultSelector)
  751. {
  752. return SelectMany_<TSource, TTaskResult, TResult>(source, x => FromAsync(ct => taskSelector(x, ct)), resultSelector);
  753. }
  754. #endif
  755. private static IObservable<TResult> SelectMany_<TSource, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TResult>> selector)
  756. {
  757. #if !NO_PERF
  758. return new SelectMany<TSource, TResult>(source, selector);
  759. #else
  760. return source.Select(selector).Merge();
  761. #endif
  762. }
  763. private static IObservable<TResult> SelectMany_<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  764. {
  765. #if !NO_PERF
  766. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  767. #else
  768. return SelectMany_<TSource, TResult>(source, x => collectionSelector(x).Select(y => resultSelector(x, y)));
  769. #endif
  770. }
  771. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TResult>> onNext, Func<Exception, IObservable<TResult>> onError, Func<IObservable<TResult>> onCompleted)
  772. {
  773. #if !NO_PERF
  774. return new SelectMany<TSource, TResult>(source, onNext, onError, onCompleted);
  775. #else
  776. return source.Materialize().SelectMany(notification =>
  777. {
  778. if (notification.Kind == NotificationKind.OnNext)
  779. return onNext(notification.Value);
  780. else if (notification.Kind == NotificationKind.OnError)
  781. return onError(notification.Exception);
  782. else
  783. return onCompleted();
  784. });
  785. #endif
  786. }
  787. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, IEnumerable<TResult>> selector)
  788. {
  789. #if !NO_PERF
  790. return new SelectMany<TSource, TResult>(source, selector);
  791. #else
  792. return SelectMany_<TSource, TResult, TResult>(source, selector, (_, x) => x);
  793. #endif
  794. }
  795. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IEnumerable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  796. {
  797. return SelectMany_<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  798. }
  799. private static IObservable<TResult> SelectMany_<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IEnumerable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  800. {
  801. #if !NO_PERF
  802. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  803. #else
  804. return new AnonymousObservable<TResult>(observer =>
  805. source.Subscribe(
  806. x =>
  807. {
  808. var xs = default(IEnumerable<TCollection>);
  809. try
  810. {
  811. xs = collectionSelector(x);
  812. }
  813. catch (Exception exception)
  814. {
  815. observer.OnError(exception);
  816. return;
  817. }
  818. var e = xs.GetEnumerator();
  819. try
  820. {
  821. var hasNext = true;
  822. while (hasNext)
  823. {
  824. hasNext = false;
  825. var current = default(TResult);
  826. try
  827. {
  828. hasNext = e.MoveNext();
  829. if (hasNext)
  830. current = resultSelector(x, e.Current);
  831. }
  832. catch (Exception exception)
  833. {
  834. observer.OnError(exception);
  835. return;
  836. }
  837. if (hasNext)
  838. observer.OnNext(current);
  839. }
  840. }
  841. finally
  842. {
  843. if (e != null)
  844. e.Dispose();
  845. }
  846. },
  847. observer.OnError,
  848. observer.OnCompleted
  849. )
  850. );
  851. #endif
  852. }
  853. #endregion
  854. #region + Skip +
  855. public virtual IObservable<TSource> Skip<TSource>(IObservable<TSource> source, int count)
  856. {
  857. #if !NO_PERF
  858. var skip = source as Skip<TSource>;
  859. if (skip != null && skip._scheduler == null)
  860. return skip.Ω(count);
  861. return new Skip<TSource>(source, count);
  862. #else
  863. return new AnonymousObservable<TSource>(observer =>
  864. {
  865. var remaining = count;
  866. return source.Subscribe(
  867. x =>
  868. {
  869. if (remaining <= 0)
  870. observer.OnNext(x);
  871. else
  872. remaining--;
  873. },
  874. observer.OnError,
  875. observer.OnCompleted);
  876. });
  877. #endif
  878. }
  879. #endregion
  880. #region + SkipWhile +
  881. public virtual IObservable<TSource> SkipWhile<TSource>(IObservable<TSource> source, Func<TSource, bool> predicate)
  882. {
  883. #if !NO_PERF
  884. return new SkipWhile<TSource>(source, predicate);
  885. #else
  886. return SkipWhile_(source, (x, i) => predicate(x));
  887. #endif
  888. }
  889. public virtual IObservable<TSource> SkipWhile<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  890. {
  891. #if !NO_PERF
  892. return new SkipWhile<TSource>(source, predicate);
  893. #else
  894. return SkipWhile_(source, predicate);
  895. #endif
  896. }
  897. #if NO_PERF
  898. private static IObservable<TSource> SkipWhile_<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  899. {
  900. return new AnonymousObservable<TSource>(observer =>
  901. {
  902. var running = false;
  903. var i = 0;
  904. return source.Subscribe(
  905. x =>
  906. {
  907. if (!running)
  908. try
  909. {
  910. running = !predicate(x, checked(i++));
  911. }
  912. catch (Exception exception)
  913. {
  914. observer.OnError(exception);
  915. return;
  916. }
  917. if (running)
  918. observer.OnNext(x);
  919. },
  920. observer.OnError,
  921. observer.OnCompleted);
  922. });
  923. }
  924. #endif
  925. #endregion
  926. #region + Take +
  927. public virtual IObservable<TSource> Take<TSource>(IObservable<TSource> source, int count)
  928. {
  929. if (count == 0)
  930. return Empty<TSource>();
  931. return Take_(source, count);
  932. }
  933. public virtual IObservable<TSource> Take<TSource>(IObservable<TSource> source, int count, IScheduler scheduler)
  934. {
  935. if (count == 0)
  936. return Empty<TSource>(scheduler);
  937. return Take_(source, count);
  938. }
  939. #if !NO_PERF
  940. private static IObservable<TSource> Take_<TSource>(IObservable<TSource> source, int count)
  941. {
  942. var take = source as Take<TSource>;
  943. if (take != null && take._scheduler == null)
  944. return take.Ω(count);
  945. return new Take<TSource>(source, count);
  946. }
  947. #else
  948. private static IObservable<TSource> Take_<TSource>(IObservable<TSource> source, int count)
  949. {
  950. return new AnonymousObservable<TSource>(observer =>
  951. {
  952. var remaining = count;
  953. return source.Subscribe(
  954. x =>
  955. {
  956. if (remaining > 0)
  957. {
  958. --remaining;
  959. observer.OnNext(x);
  960. if (remaining == 0)
  961. observer.OnCompleted();
  962. }
  963. },
  964. observer.OnError,
  965. observer.OnCompleted);
  966. });
  967. }
  968. #endif
  969. #endregion
  970. #region + TakeWhile +
  971. public virtual IObservable<TSource> TakeWhile<TSource>(IObservable<TSource> source, Func<TSource, bool> predicate)
  972. {
  973. #if !NO_PERF
  974. return new TakeWhile<TSource>(source, predicate);
  975. #else
  976. return TakeWhile_(source, (x, i) => predicate(x));
  977. #endif
  978. }
  979. public virtual IObservable<TSource> TakeWhile<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  980. {
  981. #if !NO_PERF
  982. return new TakeWhile<TSource>(source, predicate);
  983. #else
  984. return TakeWhile_(source, predicate);
  985. #endif
  986. }
  987. #if NO_PERF
  988. private static IObservable<TSource> TakeWhile_<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  989. {
  990. return new AnonymousObservable<TSource>(observer =>
  991. {
  992. var running = true;
  993. var i = 0;
  994. return source.Subscribe(
  995. x =>
  996. {
  997. if (running)
  998. {
  999. try
  1000. {
  1001. running = predicate(x, checked(i++));
  1002. }
  1003. catch (Exception exception)
  1004. {
  1005. observer.OnError(exception);
  1006. return;
  1007. }
  1008. if (running)
  1009. observer.OnNext(x);
  1010. else
  1011. observer.OnCompleted();
  1012. }
  1013. },
  1014. observer.OnError,
  1015. observer.OnCompleted);
  1016. });
  1017. }
  1018. #endif
  1019. #endregion
  1020. #region + Where +
  1021. public virtual IObservable<TSource> Where<TSource>(IObservable<TSource> source, Func<TSource, bool> predicate)
  1022. {
  1023. #if !NO_PERF
  1024. var where = source as Where<TSource>;
  1025. if (where != null)
  1026. return where.Ω(predicate);
  1027. return new Where<TSource>(source, predicate);
  1028. #else
  1029. var w = source as WhereObservable<TSource>;
  1030. if (w != null)
  1031. return w.Where(predicate);
  1032. return new WhereObservable<TSource>(source, predicate);
  1033. #endif
  1034. }
  1035. #if NO_PERF
  1036. class WhereObservable<TSource> : ObservableBase<TSource>
  1037. {
  1038. private readonly IObservable<TSource> _source;
  1039. private readonly Func<TSource, bool> _predicate;
  1040. public WhereObservable(IObservable<TSource> source, Func<TSource, bool> predicate)
  1041. {
  1042. _source = source;
  1043. _predicate = predicate;
  1044. }
  1045. protected override IDisposable SubscribeCore(IObserver<TSource> observer)
  1046. {
  1047. return _source.Subscribe(new Observer(observer, _predicate));
  1048. }
  1049. public IObservable<TSource> Where(Func<TSource, bool> predicate)
  1050. {
  1051. return new WhereObservable<TSource>(_source, x => _predicate(x) && predicate(x));
  1052. }
  1053. class Observer : ObserverBase<TSource>
  1054. {
  1055. private readonly IObserver<TSource> _observer;
  1056. private readonly Func<TSource, bool> _predicate;
  1057. public Observer(IObserver<TSource> observer, Func<TSource, bool> predicate)
  1058. {
  1059. _observer = observer;
  1060. _predicate = predicate;
  1061. }
  1062. protected override void OnNextCore(TSource value)
  1063. {
  1064. bool shouldRun;
  1065. try
  1066. {
  1067. shouldRun = _predicate(value);
  1068. }
  1069. catch (Exception exception)
  1070. {
  1071. _observer.OnError(exception);
  1072. return;
  1073. }
  1074. if (shouldRun)
  1075. _observer.OnNext(value);
  1076. }
  1077. protected override void OnErrorCore(Exception error)
  1078. {
  1079. _observer.OnError(error);
  1080. }
  1081. protected override void OnCompletedCore()
  1082. {
  1083. _observer.OnCompleted();
  1084. }
  1085. }
  1086. }
  1087. #endif
  1088. public virtual IObservable<TSource> Where<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  1089. {
  1090. #if !NO_PERF
  1091. return new Where<TSource>(source, predicate);
  1092. #else
  1093. return Defer(() =>
  1094. {
  1095. var index = 0;
  1096. return source.Where(x => predicate(x, checked(index++)));
  1097. });
  1098. #endif
  1099. }
  1100. #endregion
  1101. }
  1102. }