QueryLanguage.StandardSequenceOperators.cs 54 KB

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