QueryLanguage.StandardSequenceOperators.cs 57 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458
  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 ObservableImpl;
  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 + OrderBy +
  673. public virtual IOrderedObservable<TSource> OrderBy<TSource, TKey>(IObservable<TSource> source, Func<TSource, TKey> keySelector)
  674. {
  675. return new OrderBy<TSource, TKey>(source, keySelector, comparer: null, descending: false);
  676. }
  677. public virtual IOrderedObservable<TSource> OrderBy<TSource, TKey>(IObservable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
  678. {
  679. return new OrderBy<TSource, TKey>(source, keySelector, comparer, descending: false);
  680. }
  681. public virtual IOrderedObservable<TSource> OrderBy<TSource, TOther>(IObservable<TSource> source, Func<TSource, IObservable<TOther>> timeSelector)
  682. {
  683. return new OrderBy<TSource, TOther>(source, timeSelector, descending: false);
  684. }
  685. #endregion
  686. #region + OrderByDescending +
  687. public virtual IOrderedObservable<TSource> OrderByDescending<TSource, TKey>(IObservable<TSource> source, Func<TSource, TKey> keySelector)
  688. {
  689. return new OrderBy<TSource, TKey>(source, keySelector, comparer: null, descending: true);
  690. }
  691. public virtual IOrderedObservable<TSource> OrderByDescending<TSource, TKey>(IObservable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
  692. {
  693. return new OrderBy<TSource, TKey>(source, keySelector, comparer, descending: true);
  694. }
  695. public virtual IOrderedObservable<TSource> OrderByDescending<TSource, TOther>(IObservable<TSource> source, Func<TSource, IObservable<TOther>> timeSelector)
  696. {
  697. return new OrderBy<TSource, TOther>(source, timeSelector, descending: true);
  698. }
  699. #endregion
  700. #region + Select +
  701. public virtual IObservable<TResult> Select<TSource, TResult>(IObservable<TSource> source, Func<TSource, TResult> selector)
  702. {
  703. #if !NO_PERF
  704. var select = source as Select<TSource>;
  705. if (select != null)
  706. return select.Omega(selector);
  707. return new Select<TSource, TResult>(source, selector);
  708. #else
  709. var s = source as SelectObservable<TSource>;
  710. if (s != null)
  711. return s.Select(selector);
  712. return new SelectObservable<TSource, TResult>(source, selector);
  713. #endif
  714. }
  715. #if NO_PERF
  716. abstract class SelectObservable<TResult> : ObservableBase<TResult>
  717. {
  718. public abstract IObservable<TResult2> Select<TResult2>(Func<TResult, TResult2> selector);
  719. }
  720. class SelectObservable<TSource, TResult> : SelectObservable<TResult>
  721. {
  722. private readonly IObservable<TSource> _source;
  723. private readonly Func<TSource, TResult> _selector;
  724. public SelectObservable(IObservable<TSource> source, Func<TSource, TResult> selector)
  725. {
  726. _source = source;
  727. _selector = selector;
  728. }
  729. protected override IDisposable SubscribeCore(IObserver<TResult> observer)
  730. {
  731. return _source.Subscribe(new Observer(observer, _selector));
  732. }
  733. public override IObservable<TResult2> Select<TResult2>(Func<TResult, TResult2> selector)
  734. {
  735. return new SelectObservable<TSource, TResult2>(_source, x => selector(_selector(x)));
  736. }
  737. class Observer : ObserverBase<TSource>
  738. {
  739. private readonly IObserver<TResult> _observer;
  740. private readonly Func<TSource, TResult> _selector;
  741. public Observer(IObserver<TResult> observer, Func<TSource, TResult> selector)
  742. {
  743. _observer = observer;
  744. _selector = selector;
  745. }
  746. protected override void OnNextCore(TSource value)
  747. {
  748. TResult result;
  749. try
  750. {
  751. result = _selector(value);
  752. }
  753. catch (Exception exception)
  754. {
  755. _observer.OnError(exception);
  756. return;
  757. }
  758. _observer.OnNext(result);
  759. }
  760. protected override void OnErrorCore(Exception error)
  761. {
  762. _observer.OnError(error);
  763. }
  764. protected override void OnCompletedCore()
  765. {
  766. _observer.OnCompleted();
  767. }
  768. }
  769. }
  770. #endif
  771. public virtual IObservable<TResult> Select<TSource, TResult>(IObservable<TSource> source, Func<TSource, int, TResult> selector)
  772. {
  773. #if !NO_PERF
  774. return new Select<TSource, TResult>(source, selector);
  775. #else
  776. return Defer(() =>
  777. {
  778. var index = 0;
  779. return source.Select(x => selector(x, checked(index++)));
  780. });
  781. #endif
  782. }
  783. #endregion
  784. #region + SelectMany +
  785. public virtual IObservable<TOther> SelectMany<TSource, TOther>(IObservable<TSource> source, IObservable<TOther> other)
  786. {
  787. return SelectMany_<TSource, TOther>(source, _ => other);
  788. }
  789. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TResult>> selector)
  790. {
  791. return SelectMany_<TSource, TResult>(source, selector);
  792. }
  793. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, int, IObservable<TResult>> selector)
  794. {
  795. return SelectMany_<TSource, TResult>(source, selector);
  796. }
  797. #if !NO_TPL
  798. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, Task<TResult>> selector)
  799. {
  800. #if !NO_PERF
  801. return new SelectMany<TSource, TResult>(source, (x, token) => selector(x));
  802. #else
  803. return SelectMany_<TSource, TResult>(source, x => selector(x).ToObservable());
  804. #endif
  805. }
  806. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, CancellationToken, Task<TResult>> selector)
  807. {
  808. #if !NO_PERF
  809. return new SelectMany<TSource, TResult>(source, selector);
  810. #else
  811. return SelectMany_<TSource, TResult>(source, x => FromAsync(ct => selector(x, ct)));
  812. #endif
  813. }
  814. #endif
  815. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  816. {
  817. return SelectMany_<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  818. }
  819. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, int, IObservable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  820. {
  821. return SelectMany_<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  822. }
  823. #if !NO_TPL
  824. public virtual IObservable<TResult> SelectMany<TSource, TTaskResult, TResult>(IObservable<TSource> source, Func<TSource, Task<TTaskResult>> taskSelector, Func<TSource, TTaskResult, TResult> resultSelector)
  825. {
  826. #if !NO_PERF
  827. return new SelectMany<TSource, TTaskResult, TResult>(source, (x, token) => taskSelector(x), resultSelector);
  828. #else
  829. return SelectMany_<TSource, TTaskResult, TResult>(source, x => taskSelector(x).ToObservable(), resultSelector);
  830. #endif
  831. }
  832. public virtual IObservable<TResult> SelectMany<TSource, TTaskResult, TResult>(IObservable<TSource> source, Func<TSource, CancellationToken, Task<TTaskResult>> taskSelector, Func<TSource, TTaskResult, TResult> resultSelector)
  833. {
  834. #if !NO_PERF
  835. return new SelectMany<TSource, TTaskResult, TResult>(source, taskSelector, resultSelector);
  836. #else
  837. return SelectMany_<TSource, TTaskResult, TResult>(source, x => FromAsync(ct => taskSelector(x, ct)), resultSelector);
  838. #endif
  839. }
  840. #endif
  841. private static IObservable<TResult> SelectMany_<TSource, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TResult>> selector)
  842. {
  843. #if !NO_PERF
  844. return new SelectMany<TSource, TResult>(source, selector);
  845. #else
  846. return source.Select(selector).Merge();
  847. #endif
  848. }
  849. private static IObservable<TResult> SelectMany_<TSource, TResult>(IObservable<TSource> source, Func<TSource, int, IObservable<TResult>> selector)
  850. {
  851. #if !NO_PERF
  852. return new SelectMany<TSource, TResult>(source, selector);
  853. #else
  854. return source.Select(selector).Merge();
  855. #endif
  856. }
  857. private static IObservable<TResult> SelectMany_<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  858. {
  859. #if !NO_PERF
  860. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  861. #else
  862. return SelectMany_<TSource, TResult>(source, x => collectionSelector(x).Select(y => resultSelector(x, y)));
  863. #endif
  864. }
  865. private static IObservable<TResult> SelectMany_<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, int, IObservable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  866. {
  867. #if !NO_PERF
  868. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  869. #else
  870. return SelectMany_<TSource, TResult>(source, x => collectionSelector(x).Select(y => resultSelector(x, y)));
  871. #endif
  872. }
  873. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, IObservable<TResult>> onNext, Func<Exception, IObservable<TResult>> onError, Func<IObservable<TResult>> onCompleted)
  874. {
  875. #if !NO_PERF
  876. return new SelectMany<TSource, TResult>(source, onNext, onError, onCompleted);
  877. #else
  878. return source.Materialize().SelectMany(notification =>
  879. {
  880. if (notification.Kind == NotificationKind.OnNext)
  881. return onNext(notification.Value);
  882. else if (notification.Kind == NotificationKind.OnError)
  883. return onError(notification.Exception);
  884. else
  885. return onCompleted();
  886. });
  887. #endif
  888. }
  889. 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)
  890. {
  891. #if !NO_PERF
  892. return new SelectMany<TSource, TResult>(source, onNext, onError, onCompleted);
  893. #else
  894. return source.Materialize().SelectMany(notification =>
  895. {
  896. if (notification.Kind == NotificationKind.OnNext)
  897. return onNext(notification.Value);
  898. else if (notification.Kind == NotificationKind.OnError)
  899. return onError(notification.Exception);
  900. else
  901. return onCompleted();
  902. });
  903. #endif
  904. }
  905. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, IEnumerable<TResult>> selector)
  906. {
  907. #if !NO_PERF
  908. return new SelectMany<TSource, TResult>(source, selector);
  909. #else
  910. return SelectMany_<TSource, TResult, TResult>(source, selector, (_, x) => x);
  911. #endif
  912. }
  913. public virtual IObservable<TResult> SelectMany<TSource, TResult>(IObservable<TSource> source, Func<TSource, int, IEnumerable<TResult>> selector)
  914. {
  915. #if !NO_PERF
  916. return new SelectMany<TSource, TResult>(source, selector);
  917. #else
  918. return SelectMany_<TSource, TResult, TResult>(source, selector, (_, x) => x);
  919. #endif
  920. }
  921. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IEnumerable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  922. {
  923. return SelectMany_<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  924. }
  925. private static IObservable<TResult> SelectMany_<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, IEnumerable<TCollection>> collectionSelector, Func<TSource, TCollection, TResult> resultSelector)
  926. {
  927. #if !NO_PERF
  928. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  929. #else
  930. return new AnonymousObservable<TResult>(observer =>
  931. source.Subscribe(
  932. x =>
  933. {
  934. var xs = default(IEnumerable<TCollection>);
  935. try
  936. {
  937. xs = collectionSelector(x);
  938. }
  939. catch (Exception exception)
  940. {
  941. observer.OnError(exception);
  942. return;
  943. }
  944. var e = xs.GetEnumerator();
  945. try
  946. {
  947. var hasNext = true;
  948. while (hasNext)
  949. {
  950. hasNext = false;
  951. var current = default(TResult);
  952. try
  953. {
  954. hasNext = e.MoveNext();
  955. if (hasNext)
  956. current = resultSelector(x, e.Current);
  957. }
  958. catch (Exception exception)
  959. {
  960. observer.OnError(exception);
  961. return;
  962. }
  963. if (hasNext)
  964. observer.OnNext(current);
  965. }
  966. }
  967. finally
  968. {
  969. if (e != null)
  970. e.Dispose();
  971. }
  972. },
  973. observer.OnError,
  974. observer.OnCompleted
  975. )
  976. );
  977. #endif
  978. }
  979. public virtual IObservable<TResult> SelectMany<TSource, TCollection, TResult>(IObservable<TSource> source, Func<TSource, int, IEnumerable<TCollection>> collectionSelector, Func<TSource, int, TCollection, int, TResult> resultSelector)
  980. {
  981. return new SelectMany<TSource, TCollection, TResult>(source, collectionSelector, resultSelector);
  982. }
  983. #endregion
  984. #region + Skip +
  985. public virtual IObservable<TSource> Skip<TSource>(IObservable<TSource> source, int count)
  986. {
  987. #if !NO_PERF
  988. var skip = source as Skip<TSource>;
  989. if (skip != null && skip._scheduler == null)
  990. return skip.Omega(count);
  991. return new Skip<TSource>(source, count);
  992. #else
  993. return new AnonymousObservable<TSource>(observer =>
  994. {
  995. var remaining = count;
  996. return source.Subscribe(
  997. x =>
  998. {
  999. if (remaining <= 0)
  1000. observer.OnNext(x);
  1001. else
  1002. remaining--;
  1003. },
  1004. observer.OnError,
  1005. observer.OnCompleted);
  1006. });
  1007. #endif
  1008. }
  1009. #endregion
  1010. #region + SkipWhile +
  1011. public virtual IObservable<TSource> SkipWhile<TSource>(IObservable<TSource> source, Func<TSource, bool> predicate)
  1012. {
  1013. #if !NO_PERF
  1014. return new SkipWhile<TSource>(source, predicate);
  1015. #else
  1016. return SkipWhile_(source, (x, i) => predicate(x));
  1017. #endif
  1018. }
  1019. public virtual IObservable<TSource> SkipWhile<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  1020. {
  1021. #if !NO_PERF
  1022. return new SkipWhile<TSource>(source, predicate);
  1023. #else
  1024. return SkipWhile_(source, predicate);
  1025. #endif
  1026. }
  1027. #if NO_PERF
  1028. private static IObservable<TSource> SkipWhile_<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  1029. {
  1030. return new AnonymousObservable<TSource>(observer =>
  1031. {
  1032. var running = false;
  1033. var i = 0;
  1034. return source.Subscribe(
  1035. x =>
  1036. {
  1037. if (!running)
  1038. try
  1039. {
  1040. running = !predicate(x, checked(i++));
  1041. }
  1042. catch (Exception exception)
  1043. {
  1044. observer.OnError(exception);
  1045. return;
  1046. }
  1047. if (running)
  1048. observer.OnNext(x);
  1049. },
  1050. observer.OnError,
  1051. observer.OnCompleted);
  1052. });
  1053. }
  1054. #endif
  1055. #endregion
  1056. #region + Take +
  1057. public virtual IObservable<TSource> Take<TSource>(IObservable<TSource> source, int count)
  1058. {
  1059. if (count == 0)
  1060. return Empty<TSource>();
  1061. return Take_(source, count);
  1062. }
  1063. public virtual IObservable<TSource> Take<TSource>(IObservable<TSource> source, int count, IScheduler scheduler)
  1064. {
  1065. if (count == 0)
  1066. return Empty<TSource>(scheduler);
  1067. return Take_(source, count);
  1068. }
  1069. #if !NO_PERF
  1070. private static IObservable<TSource> Take_<TSource>(IObservable<TSource> source, int count)
  1071. {
  1072. var take = source as Take<TSource>;
  1073. if (take != null && take._scheduler == null)
  1074. return take.Omega(count);
  1075. return new Take<TSource>(source, count);
  1076. }
  1077. #else
  1078. private static IObservable<TSource> Take_<TSource>(IObservable<TSource> source, int count)
  1079. {
  1080. return new AnonymousObservable<TSource>(observer =>
  1081. {
  1082. var remaining = count;
  1083. return source.Subscribe(
  1084. x =>
  1085. {
  1086. if (remaining > 0)
  1087. {
  1088. --remaining;
  1089. observer.OnNext(x);
  1090. if (remaining == 0)
  1091. observer.OnCompleted();
  1092. }
  1093. },
  1094. observer.OnError,
  1095. observer.OnCompleted);
  1096. });
  1097. }
  1098. #endif
  1099. #endregion
  1100. #region + TakeWhile +
  1101. public virtual IObservable<TSource> TakeWhile<TSource>(IObservable<TSource> source, Func<TSource, bool> predicate)
  1102. {
  1103. #if !NO_PERF
  1104. return new TakeWhile<TSource>(source, predicate);
  1105. #else
  1106. return TakeWhile_(source, (x, i) => predicate(x));
  1107. #endif
  1108. }
  1109. public virtual IObservable<TSource> TakeWhile<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  1110. {
  1111. #if !NO_PERF
  1112. return new TakeWhile<TSource>(source, predicate);
  1113. #else
  1114. return TakeWhile_(source, predicate);
  1115. #endif
  1116. }
  1117. #if NO_PERF
  1118. private static IObservable<TSource> TakeWhile_<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  1119. {
  1120. return new AnonymousObservable<TSource>(observer =>
  1121. {
  1122. var running = true;
  1123. var i = 0;
  1124. return source.Subscribe(
  1125. x =>
  1126. {
  1127. if (running)
  1128. {
  1129. try
  1130. {
  1131. running = predicate(x, checked(i++));
  1132. }
  1133. catch (Exception exception)
  1134. {
  1135. observer.OnError(exception);
  1136. return;
  1137. }
  1138. if (running)
  1139. observer.OnNext(x);
  1140. else
  1141. observer.OnCompleted();
  1142. }
  1143. },
  1144. observer.OnError,
  1145. observer.OnCompleted);
  1146. });
  1147. }
  1148. #endif
  1149. #endregion
  1150. #region + ThenBy +
  1151. public virtual IOrderedObservable<TSource> ThenBy<TSource, TKey>(IOrderedObservable<TSource> source, Func<TSource, TKey> keySelector)
  1152. {
  1153. return source.CreateOrderedObservable(keySelector, comparer: null, descending: false);
  1154. }
  1155. public virtual IOrderedObservable<TSource> ThenBy<TSource, TKey>(IOrderedObservable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
  1156. {
  1157. return source.CreateOrderedObservable(keySelector, comparer, descending: false);
  1158. }
  1159. public virtual IOrderedObservable<TSource> ThenBy<TSource, TOther>(IOrderedObservable<TSource> source, Func<TSource, IObservable<TOther>> timeSelector)
  1160. {
  1161. return source.CreateOrderedObservable(timeSelector, descending: false);
  1162. }
  1163. #endregion
  1164. #region + ThenByDescending +
  1165. public virtual IOrderedObservable<TSource> ThenByDescending<TSource, TKey>(IOrderedObservable<TSource> source, Func<TSource, TKey> keySelector)
  1166. {
  1167. return source.CreateOrderedObservable(keySelector, comparer: null, descending: true);
  1168. }
  1169. public virtual IOrderedObservable<TSource> ThenByDescending<TSource, TKey>(IOrderedObservable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
  1170. {
  1171. return source.CreateOrderedObservable(keySelector, comparer, descending: true);
  1172. }
  1173. public virtual IOrderedObservable<TSource> ThenByDescending<TSource, TOther>(IOrderedObservable<TSource> source, Func<TSource, IObservable<TOther>> timeSelector)
  1174. {
  1175. return source.CreateOrderedObservable(timeSelector, descending: true);
  1176. }
  1177. #endregion
  1178. #region + Where +
  1179. public virtual IObservable<TSource> Where<TSource>(IObservable<TSource> source, Func<TSource, bool> predicate)
  1180. {
  1181. #if !NO_PERF
  1182. var where = source as Where<TSource>;
  1183. if (where != null)
  1184. return where.Omega(predicate);
  1185. return new Where<TSource>(source, predicate);
  1186. #else
  1187. var w = source as WhereObservable<TSource>;
  1188. if (w != null)
  1189. return w.Where(predicate);
  1190. return new WhereObservable<TSource>(source, predicate);
  1191. #endif
  1192. }
  1193. #if NO_PERF
  1194. class WhereObservable<TSource> : ObservableBase<TSource>
  1195. {
  1196. private readonly IObservable<TSource> _source;
  1197. private readonly Func<TSource, bool> _predicate;
  1198. public WhereObservable(IObservable<TSource> source, Func<TSource, bool> predicate)
  1199. {
  1200. _source = source;
  1201. _predicate = predicate;
  1202. }
  1203. protected override IDisposable SubscribeCore(IObserver<TSource> observer)
  1204. {
  1205. return _source.Subscribe(new Observer(observer, _predicate));
  1206. }
  1207. public IObservable<TSource> Where(Func<TSource, bool> predicate)
  1208. {
  1209. return new WhereObservable<TSource>(_source, x => _predicate(x) && predicate(x));
  1210. }
  1211. class Observer : ObserverBase<TSource>
  1212. {
  1213. private readonly IObserver<TSource> _observer;
  1214. private readonly Func<TSource, bool> _predicate;
  1215. public Observer(IObserver<TSource> observer, Func<TSource, bool> predicate)
  1216. {
  1217. _observer = observer;
  1218. _predicate = predicate;
  1219. }
  1220. protected override void OnNextCore(TSource value)
  1221. {
  1222. bool shouldRun;
  1223. try
  1224. {
  1225. shouldRun = _predicate(value);
  1226. }
  1227. catch (Exception exception)
  1228. {
  1229. _observer.OnError(exception);
  1230. return;
  1231. }
  1232. if (shouldRun)
  1233. _observer.OnNext(value);
  1234. }
  1235. protected override void OnErrorCore(Exception error)
  1236. {
  1237. _observer.OnError(error);
  1238. }
  1239. protected override void OnCompletedCore()
  1240. {
  1241. _observer.OnCompleted();
  1242. }
  1243. }
  1244. }
  1245. #endif
  1246. public virtual IObservable<TSource> Where<TSource>(IObservable<TSource> source, Func<TSource, int, bool> predicate)
  1247. {
  1248. #if !NO_PERF
  1249. return new Where<TSource>(source, predicate);
  1250. #else
  1251. return Defer(() =>
  1252. {
  1253. var index = 0;
  1254. return source.Where(x => predicate(x, checked(index++)));
  1255. });
  1256. #endif
  1257. }
  1258. #endregion
  1259. }
  1260. }