AsyncEnumerable.Single.cs 92 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363236423652366236723682369237023712372237323742375237623772378237923802381238223832384238523862387238823892390239123922393239423952396239723982399240024012402240324042405240624072408240924102411241224132414241524162417241824192420242124222423242424252426242724282429243024312432243324342435243624372438243924402441244224432444244524462447244824492450245124522453245424552456245724582459246024612462246324642465246624672468246924702471247224732474247524762477247824792480248124822483248424852486248724882489249024912492249324942495249624972498249925002501
  1. // Copyright (c) Microsoft Open Technologies, Inc. All rights reserved. See License.txt in the project root for license information.
  2. using System;
  3. using System.Collections.Generic;
  4. using System.Linq;
  5. using System.Threading;
  6. using System.Threading.Tasks;
  7. namespace System.Linq
  8. {
  9. public static partial class AsyncEnumerable
  10. {
  11. public static IAsyncEnumerable<TResult> Select<TSource, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, TResult> selector)
  12. {
  13. if (source == null)
  14. throw new ArgumentNullException("source");
  15. if (selector == null)
  16. throw new ArgumentNullException("selector");
  17. return Create(() =>
  18. {
  19. var e = source.GetEnumerator();
  20. var current = default(TResult);
  21. var cts = new CancellationTokenDisposable();
  22. var d = new CompositeDisposable(cts, e);
  23. return Create(
  24. (ct, tcs) =>
  25. {
  26. e.MoveNext(cts.Token).ContinueWith(t =>
  27. {
  28. t.Handle(tcs, res =>
  29. {
  30. if (res)
  31. {
  32. try
  33. {
  34. current = selector(e.Current);
  35. }
  36. catch (Exception ex)
  37. {
  38. tcs.TrySetException(ex);
  39. return;
  40. }
  41. tcs.TrySetResult(true);
  42. }
  43. else
  44. {
  45. tcs.TrySetResult(false);
  46. }
  47. });
  48. });
  49. return tcs.Task.UsingEnumerator(e);
  50. },
  51. () => current,
  52. d.Dispose
  53. );
  54. });
  55. }
  56. public static IAsyncEnumerable<TResult> Select<TSource, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, int, TResult> selector)
  57. {
  58. if (source == null)
  59. throw new ArgumentNullException("source");
  60. if (selector == null)
  61. throw new ArgumentNullException("selector");
  62. return Create(() =>
  63. {
  64. var e = source.GetEnumerator();
  65. var current = default(TResult);
  66. var index = 0;
  67. var cts = new CancellationTokenDisposable();
  68. var d = new CompositeDisposable(cts, e);
  69. return Create(
  70. (ct, tcs) =>
  71. {
  72. e.MoveNext(cts.Token).ContinueWith(t =>
  73. {
  74. t.Handle(tcs, res =>
  75. {
  76. if (res)
  77. {
  78. try
  79. {
  80. current = selector(e.Current, index++);
  81. }
  82. catch (Exception ex)
  83. {
  84. tcs.TrySetException(ex);
  85. return;
  86. }
  87. tcs.TrySetResult(true);
  88. }
  89. else
  90. {
  91. tcs.TrySetResult(false);
  92. }
  93. });
  94. });
  95. return tcs.Task.UsingEnumerator(e);
  96. },
  97. () => current,
  98. d.Dispose
  99. );
  100. });
  101. }
  102. public static IAsyncEnumerable<TSource> AsAsyncEnumerable<TSource>(this IAsyncEnumerable<TSource> source)
  103. {
  104. if (source == null)
  105. throw new ArgumentNullException("source");
  106. return source.Select(x => x);
  107. }
  108. public static IAsyncEnumerable<TSource> Where<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, bool> predicate)
  109. {
  110. if (source == null)
  111. throw new ArgumentNullException("source");
  112. if (predicate == null)
  113. throw new ArgumentNullException("predicate");
  114. return Create(() =>
  115. {
  116. var e = source.GetEnumerator();
  117. var cts = new CancellationTokenDisposable();
  118. var d = new CompositeDisposable(cts, e);
  119. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  120. f = (tcs, ct) =>
  121. {
  122. e.MoveNext(ct).ContinueWith(t =>
  123. {
  124. t.Handle(tcs, res =>
  125. {
  126. if (res)
  127. {
  128. var b = false;
  129. try
  130. {
  131. b = predicate(e.Current);
  132. }
  133. catch (Exception ex)
  134. {
  135. tcs.TrySetException(ex);
  136. return;
  137. }
  138. if (b)
  139. tcs.TrySetResult(true);
  140. else
  141. f(tcs, ct);
  142. }
  143. else
  144. tcs.TrySetResult(false);
  145. });
  146. });
  147. };
  148. return Create(
  149. (ct, tcs) =>
  150. {
  151. f(tcs, cts.Token);
  152. return tcs.Task.UsingEnumerator(e);
  153. },
  154. () => e.Current,
  155. d.Dispose
  156. );
  157. });
  158. }
  159. public static IAsyncEnumerable<TSource> Where<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, int, bool> predicate)
  160. {
  161. if (source == null)
  162. throw new ArgumentNullException("source");
  163. if (predicate == null)
  164. throw new ArgumentNullException("predicate");
  165. return Create(() =>
  166. {
  167. var e = source.GetEnumerator();
  168. var index = 0;
  169. var cts = new CancellationTokenDisposable();
  170. var d = new CompositeDisposable(cts, e);
  171. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  172. f = (tcs, ct) =>
  173. {
  174. e.MoveNext(ct).ContinueWith(t =>
  175. {
  176. t.Handle(tcs, res =>
  177. {
  178. if (res)
  179. {
  180. var b = false;
  181. try
  182. {
  183. b = predicate(e.Current, index++);
  184. }
  185. catch (Exception ex)
  186. {
  187. tcs.TrySetException(ex);
  188. return;
  189. }
  190. if (b)
  191. tcs.TrySetResult(true);
  192. else
  193. f(tcs, ct);
  194. }
  195. else
  196. tcs.TrySetResult(false);
  197. });
  198. });
  199. };
  200. return Create(
  201. (ct, tcs) =>
  202. {
  203. f(tcs, cts.Token);
  204. return tcs.Task.UsingEnumerator(e);
  205. },
  206. () => e.Current,
  207. d.Dispose
  208. );
  209. });
  210. }
  211. public static IAsyncEnumerable<TResult> SelectMany<TSource, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, IAsyncEnumerable<TResult>> selector)
  212. {
  213. if (source == null)
  214. throw new ArgumentNullException("source");
  215. if (selector == null)
  216. throw new ArgumentNullException("selector");
  217. return Create(() =>
  218. {
  219. // A lock seems inevitable. Disposal of the outer enumerator and completion of
  220. // MoveNext of the inner enumerator can happen concurrently.
  221. var syncRoot = new object();
  222. var e = source.GetEnumerator();
  223. var ie = default(IAsyncEnumerator<TResult>);
  224. var disposeIe = new Action(() =>
  225. {
  226. IAsyncEnumerator<TResult> localIe;
  227. lock (syncRoot)
  228. {
  229. localIe = ie;
  230. ie = null;
  231. }
  232. if (localIe != null)
  233. localIe.Dispose();
  234. });
  235. var cts = new CancellationTokenDisposable();
  236. var d = new CompositeDisposable(cts, new Disposable(disposeIe), e);
  237. var outer = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  238. var inner = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  239. inner = (tcs, ct) =>
  240. {
  241. ie.MoveNext(ct).ContinueWith(t =>
  242. {
  243. t.Handle(tcs, res =>
  244. {
  245. if (res)
  246. {
  247. tcs.TrySetResult(true);
  248. }
  249. else
  250. {
  251. disposeIe();
  252. outer(tcs, ct);
  253. }
  254. });
  255. });
  256. };
  257. outer = (tcs, ct) =>
  258. {
  259. e.MoveNext(ct).ContinueWith(t =>
  260. {
  261. t.Handle(tcs, res =>
  262. {
  263. if (res)
  264. {
  265. try
  266. {
  267. ie = selector(e.Current).GetEnumerator();
  268. inner(tcs, ct);
  269. }
  270. catch (Exception ex)
  271. {
  272. tcs.TrySetException(ex);
  273. }
  274. }
  275. else
  276. tcs.TrySetResult(false);
  277. });
  278. });
  279. };
  280. return Create(
  281. (ct, tcs) =>
  282. {
  283. if (ie == null)
  284. outer(tcs, cts.Token);
  285. else
  286. inner(tcs, cts.Token);
  287. return tcs.Task.UsingEnumerator(e);
  288. },
  289. () => ie.Current,
  290. d.Dispose
  291. );
  292. });
  293. }
  294. public static IAsyncEnumerable<TResult> SelectMany<TSource, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, int, IAsyncEnumerable<TResult>> selector)
  295. {
  296. if (source == null)
  297. throw new ArgumentNullException("source");
  298. if (selector == null)
  299. throw new ArgumentNullException("selector");
  300. return Create(() =>
  301. {
  302. // A lock seems inevitable. Disposal of the outer enumerator and completion of
  303. // MoveNext of the inner enumerator can happen concurrently.
  304. var syncRoot = new object();
  305. var e = source.GetEnumerator();
  306. var ie = default(IAsyncEnumerator<TResult>);
  307. var disposeIe = new Action(() =>
  308. {
  309. IAsyncEnumerator<TResult> localIe;
  310. lock (syncRoot)
  311. {
  312. localIe = ie;
  313. ie = null;
  314. }
  315. if (localIe != null)
  316. localIe.Dispose();
  317. });
  318. var index = 0;
  319. var cts = new CancellationTokenDisposable();
  320. var d = new CompositeDisposable(cts, new Disposable(disposeIe), e);
  321. var outer = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  322. var inner = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  323. inner = (tcs, ct) =>
  324. {
  325. ie.MoveNext(ct).ContinueWith(t =>
  326. {
  327. t.Handle(tcs, res =>
  328. {
  329. if (res)
  330. {
  331. tcs.TrySetResult(true);
  332. }
  333. else
  334. {
  335. disposeIe();
  336. outer(tcs, ct);
  337. }
  338. });
  339. });
  340. };
  341. outer = (tcs, ct) =>
  342. {
  343. e.MoveNext(ct).ContinueWith(t =>
  344. {
  345. t.Handle(tcs, res =>
  346. {
  347. if (res)
  348. {
  349. try
  350. {
  351. ie = selector(e.Current, index++).GetEnumerator();
  352. inner(tcs, ct);
  353. }
  354. catch (Exception ex)
  355. {
  356. tcs.TrySetException(ex);
  357. }
  358. }
  359. else
  360. tcs.TrySetResult(false);
  361. });
  362. });
  363. };
  364. return Create(
  365. (ct, tcs) =>
  366. {
  367. if (ie == null)
  368. outer(tcs, cts.Token);
  369. else
  370. inner(tcs, cts.Token);
  371. return tcs.Task.UsingEnumerator(e);
  372. },
  373. () => ie.Current,
  374. d.Dispose
  375. );
  376. });
  377. }
  378. public static IAsyncEnumerable<TResult> SelectMany<TSource, TCollection, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, IAsyncEnumerable<TCollection>> selector, Func<TSource, TCollection, TResult> resultSelector)
  379. {
  380. if (source == null)
  381. throw new ArgumentNullException("source");
  382. if (selector == null)
  383. throw new ArgumentNullException("selector");
  384. if (resultSelector == null)
  385. throw new ArgumentNullException("resultSelector");
  386. return source.SelectMany(x => selector(x).Select(y => resultSelector(x, y)));
  387. }
  388. public static IAsyncEnumerable<TResult> SelectMany<TSource, TCollection, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, int, IAsyncEnumerable<TCollection>> selector, Func<TSource, TCollection, TResult> resultSelector)
  389. {
  390. if (source == null)
  391. throw new ArgumentNullException("source");
  392. if (selector == null)
  393. throw new ArgumentNullException("selector");
  394. if (resultSelector == null)
  395. throw new ArgumentNullException("resultSelector");
  396. return source.SelectMany((x, i) => selector(x, i).Select(y => resultSelector(x, y)));
  397. }
  398. public static IAsyncEnumerable<TType> OfType<TType>(this IAsyncEnumerable<object> source)
  399. {
  400. if (source == null)
  401. throw new ArgumentNullException("source");
  402. return source.Where(x => x is TType).Cast<TType>();
  403. }
  404. public static IAsyncEnumerable<TResult> Cast<TResult>(this IAsyncEnumerable<object> source)
  405. {
  406. if (source == null)
  407. throw new ArgumentNullException("source");
  408. return source.Select(x => (TResult)x);
  409. }
  410. public static IAsyncEnumerable<TSource> Take<TSource>(this IAsyncEnumerable<TSource> source, int count)
  411. {
  412. if (source == null)
  413. throw new ArgumentNullException("source");
  414. if (count < 0)
  415. throw new ArgumentOutOfRangeException("count");
  416. return Create(() =>
  417. {
  418. var e = source.GetEnumerator();
  419. var n = count;
  420. var cts = new CancellationTokenDisposable();
  421. var d = new CompositeDisposable(cts, e);
  422. return Create(
  423. (ct, tcs) =>
  424. {
  425. if (n == 0)
  426. return TaskExt.Return(false, cts.Token);
  427. e.MoveNext(cts.Token).ContinueWith(t =>
  428. {
  429. t.Handle(tcs, res =>
  430. {
  431. --n;
  432. tcs.TrySetResult(res);
  433. });
  434. });
  435. return tcs.Task.UsingEnumerator(e);
  436. },
  437. () => e.Current,
  438. d.Dispose
  439. );
  440. });
  441. }
  442. public static IAsyncEnumerable<TSource> TakeWhile<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, bool> predicate)
  443. {
  444. if (source == null)
  445. throw new ArgumentNullException("source");
  446. if (predicate == null)
  447. throw new ArgumentNullException("predicate");
  448. return Create(() =>
  449. {
  450. var e = source.GetEnumerator();
  451. var cts = new CancellationTokenDisposable();
  452. var d = new CompositeDisposable(cts, e);
  453. return Create(
  454. (ct, tcs) =>
  455. {
  456. e.MoveNext(cts.Token).ContinueWith(t =>
  457. {
  458. t.Handle(tcs, res =>
  459. {
  460. if (res)
  461. {
  462. var b = false;
  463. try
  464. {
  465. b = predicate(e.Current);
  466. }
  467. catch (Exception ex)
  468. {
  469. tcs.TrySetException(ex);
  470. return;
  471. }
  472. if (b)
  473. {
  474. tcs.TrySetResult(true);
  475. return;
  476. }
  477. }
  478. tcs.TrySetResult(false);
  479. });
  480. });
  481. return tcs.Task.UsingEnumerator(e);
  482. },
  483. () => e.Current,
  484. d.Dispose
  485. );
  486. });
  487. }
  488. public static IAsyncEnumerable<TSource> TakeWhile<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, int, bool> predicate)
  489. {
  490. if (source == null)
  491. throw new ArgumentNullException("source");
  492. if (predicate == null)
  493. throw new ArgumentNullException("predicate");
  494. return Create(() =>
  495. {
  496. var e = source.GetEnumerator();
  497. var index = 0;
  498. var cts = new CancellationTokenDisposable();
  499. var d = new CompositeDisposable(cts, e);
  500. return Create(
  501. (ct, tcs) =>
  502. {
  503. e.MoveNext(cts.Token).ContinueWith(t =>
  504. {
  505. t.Handle(tcs, res =>
  506. {
  507. if (res)
  508. {
  509. var b = false;
  510. try
  511. {
  512. b = predicate(e.Current, index++);
  513. }
  514. catch (Exception ex)
  515. {
  516. tcs.TrySetException(ex);
  517. return;
  518. }
  519. if (b)
  520. {
  521. tcs.TrySetResult(true);
  522. return;
  523. }
  524. }
  525. tcs.TrySetResult(false);
  526. });
  527. });
  528. return tcs.Task.UsingEnumerator(e);
  529. },
  530. () => e.Current,
  531. d.Dispose
  532. );
  533. });
  534. }
  535. public static IAsyncEnumerable<TSource> Skip<TSource>(this IAsyncEnumerable<TSource> source, int count)
  536. {
  537. if (source == null)
  538. throw new ArgumentNullException("source");
  539. if (count < 0)
  540. throw new ArgumentOutOfRangeException("count");
  541. return Create(() =>
  542. {
  543. var e = source.GetEnumerator();
  544. var n = count;
  545. var cts = new CancellationTokenDisposable();
  546. var d = new CompositeDisposable(cts, e);
  547. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  548. f = (tcs, ct) =>
  549. {
  550. if (n == 0)
  551. e.MoveNext(ct).ContinueWith(t =>
  552. {
  553. t.Handle(tcs, x => tcs.TrySetResult(x));
  554. });
  555. else
  556. {
  557. --n;
  558. e.MoveNext(ct).ContinueWith(t =>
  559. {
  560. t.Handle(tcs, res =>
  561. {
  562. if (!res)
  563. tcs.TrySetResult(false);
  564. else
  565. f(tcs, ct);
  566. });
  567. });
  568. }
  569. };
  570. return Create(
  571. (ct, tcs) =>
  572. {
  573. f(tcs, cts.Token);
  574. return tcs.Task.UsingEnumerator(e);
  575. },
  576. () => e.Current,
  577. d.Dispose
  578. );
  579. });
  580. }
  581. public static IAsyncEnumerable<TSource> SkipWhile<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, bool> predicate)
  582. {
  583. if (source == null)
  584. throw new ArgumentNullException("source");
  585. if (predicate == null)
  586. throw new ArgumentNullException("predicate");
  587. return Create(() =>
  588. {
  589. var e = source.GetEnumerator();
  590. var skipping = true;
  591. var cts = new CancellationTokenDisposable();
  592. var d = new CompositeDisposable(cts, e);
  593. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  594. f = (tcs, ct) =>
  595. {
  596. if (skipping)
  597. e.MoveNext(ct).ContinueWith(t =>
  598. {
  599. t.Handle(tcs, res =>
  600. {
  601. if (res)
  602. {
  603. var result = false;
  604. try
  605. {
  606. result = predicate(e.Current);
  607. }
  608. catch (Exception ex)
  609. {
  610. tcs.TrySetException(ex);
  611. return;
  612. }
  613. if (result)
  614. f(tcs, ct);
  615. else
  616. {
  617. skipping = false;
  618. tcs.TrySetResult(true);
  619. }
  620. }
  621. else
  622. tcs.TrySetResult(false);
  623. });
  624. });
  625. else
  626. e.MoveNext(ct).ContinueWith(t =>
  627. {
  628. t.Handle(tcs, x => tcs.TrySetResult(x));
  629. });
  630. };
  631. return Create(
  632. (ct, tcs) =>
  633. {
  634. f(tcs, cts.Token);
  635. return tcs.Task.UsingEnumerator(e);
  636. },
  637. () => e.Current,
  638. d.Dispose
  639. );
  640. });
  641. }
  642. public static IAsyncEnumerable<TSource> SkipWhile<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, int, bool> predicate)
  643. {
  644. if (source == null)
  645. throw new ArgumentNullException("source");
  646. if (predicate == null)
  647. throw new ArgumentNullException("predicate");
  648. return Create(() =>
  649. {
  650. var e = source.GetEnumerator();
  651. var skipping = true;
  652. var index = 0;
  653. var cts = new CancellationTokenDisposable();
  654. var d = new CompositeDisposable(cts, e);
  655. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  656. f = (tcs, ct) =>
  657. {
  658. if (skipping)
  659. e.MoveNext(ct).ContinueWith(t =>
  660. {
  661. t.Handle(tcs, res =>
  662. {
  663. if (res)
  664. {
  665. var result = false;
  666. try
  667. {
  668. result = predicate(e.Current, index++);
  669. }
  670. catch (Exception ex)
  671. {
  672. tcs.TrySetException(ex);
  673. return;
  674. }
  675. if (result)
  676. f(tcs, ct);
  677. else
  678. {
  679. skipping = false;
  680. tcs.TrySetResult(true);
  681. }
  682. }
  683. else
  684. tcs.TrySetResult(false);
  685. });
  686. });
  687. else
  688. e.MoveNext(ct).ContinueWith(t =>
  689. {
  690. t.Handle(tcs, x => tcs.TrySetResult(x));
  691. });
  692. };
  693. return Create(
  694. (ct, tcs) =>
  695. {
  696. f(tcs, cts.Token);
  697. return tcs.Task.UsingEnumerator(e);
  698. },
  699. () => e.Current,
  700. d.Dispose
  701. );
  702. });
  703. }
  704. public static IAsyncEnumerable<TSource> DefaultIfEmpty<TSource>(this IAsyncEnumerable<TSource> source, TSource defaultValue)
  705. {
  706. if (source == null)
  707. throw new ArgumentNullException("source");
  708. return Create(() =>
  709. {
  710. var done = false;
  711. var hasElements = false;
  712. var e = source.GetEnumerator();
  713. var current = default(TSource);
  714. var cts = new CancellationTokenDisposable();
  715. var d = new CompositeDisposable(cts, e);
  716. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  717. f = (tcs, ct) =>
  718. {
  719. if (done)
  720. tcs.TrySetResult(false);
  721. else
  722. e.MoveNext(ct).ContinueWith(t =>
  723. {
  724. t.Handle(tcs, res =>
  725. {
  726. if (res)
  727. {
  728. hasElements = true;
  729. current = e.Current;
  730. tcs.TrySetResult(true);
  731. }
  732. else
  733. {
  734. done = true;
  735. if (!hasElements)
  736. {
  737. current = defaultValue;
  738. tcs.TrySetResult(true);
  739. }
  740. else
  741. tcs.TrySetResult(false);
  742. }
  743. });
  744. });
  745. };
  746. return Create(
  747. (ct, tcs) =>
  748. {
  749. f(tcs, cts.Token);
  750. return tcs.Task.UsingEnumerator(e);
  751. },
  752. () => current,
  753. d.Dispose
  754. );
  755. });
  756. }
  757. public static IAsyncEnumerable<TSource> DefaultIfEmpty<TSource>(this IAsyncEnumerable<TSource> source)
  758. {
  759. if (source == null)
  760. throw new ArgumentNullException("source");
  761. return source.DefaultIfEmpty(default(TSource));
  762. }
  763. public static IAsyncEnumerable<TSource> Distinct<TSource>(this IAsyncEnumerable<TSource> source, IEqualityComparer<TSource> comparer)
  764. {
  765. if (source == null)
  766. throw new ArgumentNullException("source");
  767. if (comparer == null)
  768. throw new ArgumentNullException("comparer");
  769. return Defer(() =>
  770. {
  771. var set = new HashSet<TSource>(comparer);
  772. return source.Where(set.Add);
  773. });
  774. }
  775. public static IAsyncEnumerable<TSource> Distinct<TSource>(this IAsyncEnumerable<TSource> source)
  776. {
  777. if (source == null)
  778. throw new ArgumentNullException("source");
  779. return source.Distinct(EqualityComparer<TSource>.Default);
  780. }
  781. public static IAsyncEnumerable<TSource> Reverse<TSource>(this IAsyncEnumerable<TSource> source)
  782. {
  783. if (source == null)
  784. throw new ArgumentNullException("source");
  785. return Create(() =>
  786. {
  787. var e = source.GetEnumerator();
  788. var stack = default(Stack<TSource>);
  789. var cts = new CancellationTokenDisposable();
  790. var d = new CompositeDisposable(cts, e);
  791. return Create(
  792. (ct, tcs) =>
  793. {
  794. if (stack == null)
  795. {
  796. Create(() => e).Aggregate(new Stack<TSource>(), (s, x) => { s.Push(x); return s; }, cts.Token).ContinueWith(t =>
  797. {
  798. t.Handle(tcs, res =>
  799. {
  800. stack = res;
  801. tcs.TrySetResult(stack.Count > 0);
  802. });
  803. }, cts.Token);
  804. }
  805. else
  806. {
  807. stack.Pop();
  808. tcs.TrySetResult(stack.Count > 0);
  809. }
  810. return tcs.Task.UsingEnumerator(e);
  811. },
  812. () => stack.Peek(),
  813. d.Dispose
  814. );
  815. });
  816. }
  817. public static IOrderedAsyncEnumerable<TSource> OrderBy<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
  818. {
  819. if (source == null)
  820. throw new ArgumentNullException("source");
  821. if (keySelector == null)
  822. throw new ArgumentNullException("keySelector");
  823. if (comparer == null)
  824. throw new ArgumentNullException("comparer");
  825. return new OrderedAsyncEnumerable<TSource, TKey>(
  826. Create(() =>
  827. {
  828. var current = default(IEnumerable<TSource>);
  829. return Create(
  830. ct =>
  831. {
  832. var tcs = new TaskCompletionSource<bool>();
  833. if (current == null)
  834. {
  835. source.ToList(ct).ContinueWith(t =>
  836. {
  837. t.Handle(tcs, res =>
  838. {
  839. current = res;
  840. tcs.TrySetResult(true);
  841. });
  842. });
  843. }
  844. else
  845. tcs.TrySetResult(false);
  846. return tcs.Task;
  847. },
  848. () => current,
  849. () => { }
  850. );
  851. }),
  852. keySelector,
  853. comparer
  854. );
  855. }
  856. public static IOrderedAsyncEnumerable<TSource> OrderBy<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
  857. {
  858. if (source == null)
  859. throw new ArgumentNullException("source");
  860. if (keySelector == null)
  861. throw new ArgumentNullException("keySelector");
  862. return source.OrderBy(keySelector, Comparer<TKey>.Default);
  863. }
  864. public static IOrderedAsyncEnumerable<TSource> OrderByDescending<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
  865. {
  866. if (source == null)
  867. throw new ArgumentNullException("source");
  868. if (keySelector == null)
  869. throw new ArgumentNullException("keySelector");
  870. if (comparer == null)
  871. throw new ArgumentNullException("comparer");
  872. return source.OrderBy(keySelector, new ReverseComparer<TKey>(comparer));
  873. }
  874. public static IOrderedAsyncEnumerable<TSource> OrderByDescending<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
  875. {
  876. if (source == null)
  877. throw new ArgumentNullException("source");
  878. if (keySelector == null)
  879. throw new ArgumentNullException("keySelector");
  880. return source.OrderByDescending(keySelector, Comparer<TKey>.Default);
  881. }
  882. public static IOrderedAsyncEnumerable<TSource> ThenBy<TSource, TKey>(this IOrderedAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
  883. {
  884. if (source == null)
  885. throw new ArgumentNullException("source");
  886. if (keySelector == null)
  887. throw new ArgumentNullException("keySelector");
  888. return source.ThenBy(keySelector, Comparer<TKey>.Default);
  889. }
  890. public static IOrderedAsyncEnumerable<TSource> ThenBy<TSource, TKey>(this IOrderedAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
  891. {
  892. if (source == null)
  893. throw new ArgumentNullException("source");
  894. if (keySelector == null)
  895. throw new ArgumentNullException("keySelector");
  896. if (comparer == null)
  897. throw new ArgumentNullException("comparer");
  898. return source.CreateOrderedEnumerable(keySelector, comparer, false);
  899. }
  900. public static IOrderedAsyncEnumerable<TSource> ThenByDescending<TSource, TKey>(this IOrderedAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
  901. {
  902. if (source == null)
  903. throw new ArgumentNullException("source");
  904. if (keySelector == null)
  905. throw new ArgumentNullException("keySelector");
  906. return source.ThenByDescending(keySelector, Comparer<TKey>.Default);
  907. }
  908. public static IOrderedAsyncEnumerable<TSource> ThenByDescending<TSource, TKey>(this IOrderedAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
  909. {
  910. if (source == null)
  911. throw new ArgumentNullException("source");
  912. if (keySelector == null)
  913. throw new ArgumentNullException("keySelector");
  914. if (comparer == null)
  915. throw new ArgumentNullException("comparer");
  916. return source.CreateOrderedEnumerable(keySelector, comparer, true);
  917. }
  918. static IEnumerable<IGrouping<TKey, TElement>> GroupUntil<TSource, TKey, TElement>(this IEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, IComparer<TKey> comparer)
  919. {
  920. var group = default(EnumerableGrouping<TKey, TElement>);
  921. foreach (var x in source)
  922. {
  923. var key = keySelector(x);
  924. if (group == null || comparer.Compare(group.Key, key) != 0)
  925. {
  926. group = new EnumerableGrouping<TKey, TElement>(key);
  927. yield return group;
  928. }
  929. group.Add(elementSelector(x));
  930. }
  931. }
  932. class OrderedAsyncEnumerable<T, K> : IOrderedAsyncEnumerable<T>
  933. {
  934. private readonly IAsyncEnumerable<IEnumerable<T>> equivalenceClasses;
  935. private readonly Func<T, K> keySelector;
  936. private readonly IComparer<K> comparer;
  937. public OrderedAsyncEnumerable(IAsyncEnumerable<IEnumerable<T>> equivalenceClasses, Func<T, K> keySelector, IComparer<K> comparer)
  938. {
  939. this.equivalenceClasses = equivalenceClasses;
  940. this.keySelector = keySelector;
  941. this.comparer = comparer;
  942. }
  943. public IOrderedAsyncEnumerable<T> CreateOrderedEnumerable<TKey>(Func<T, TKey> keySelector, IComparer<TKey> comparer, bool descending)
  944. {
  945. if (descending)
  946. comparer = new ReverseComparer<TKey>(comparer);
  947. return new OrderedAsyncEnumerable<T, TKey>(Classes(), keySelector, comparer);
  948. }
  949. IAsyncEnumerable<IEnumerable<T>> Classes()
  950. {
  951. return Create(() =>
  952. {
  953. var e = equivalenceClasses.GetEnumerator();
  954. var list = new List<IEnumerable<T>>();
  955. var e1 = default(IEnumerator<IEnumerable<T>>);
  956. var cts = new CancellationTokenDisposable();
  957. var d1 = new AssignableDisposable();
  958. var d = new CompositeDisposable(cts, e, d1);
  959. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  960. var g = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  961. f = (tcs, ct) =>
  962. {
  963. e.MoveNext(ct).ContinueWith(t =>
  964. {
  965. t.Handle(tcs, res =>
  966. {
  967. if (res)
  968. {
  969. try
  970. {
  971. foreach (var group in e.Current.OrderBy(keySelector, comparer).GroupUntil(keySelector, x => x, comparer))
  972. list.Add(group);
  973. f(tcs, ct);
  974. }
  975. catch (Exception exception)
  976. {
  977. tcs.TrySetException(exception);
  978. return;
  979. }
  980. }
  981. else
  982. {
  983. e.Dispose();
  984. e1 = list.GetEnumerator();
  985. d1.Disposable = e1;
  986. g(tcs, ct);
  987. }
  988. });
  989. });
  990. };
  991. g = (tcs, ct) =>
  992. {
  993. var res = false;
  994. try
  995. {
  996. res = e1.MoveNext();
  997. }
  998. catch (Exception ex)
  999. {
  1000. tcs.TrySetException(ex);
  1001. return;
  1002. }
  1003. tcs.TrySetResult(res);
  1004. };
  1005. return Create(
  1006. (ct, tcs) =>
  1007. {
  1008. if (e1 != null)
  1009. {
  1010. g(tcs, cts.Token);
  1011. return tcs.Task.UsingEnumerator(e1);
  1012. }
  1013. else
  1014. {
  1015. f(tcs, cts.Token);
  1016. return tcs.Task.UsingEnumerator(e);
  1017. }
  1018. },
  1019. () => e1.Current,
  1020. d.Dispose
  1021. );
  1022. });
  1023. }
  1024. public IAsyncEnumerator<T> GetEnumerator()
  1025. {
  1026. return Classes().SelectMany(x => x.ToAsyncEnumerable()).GetEnumerator();
  1027. }
  1028. }
  1029. class ReverseComparer<T> : IComparer<T>
  1030. {
  1031. IComparer<T> comparer;
  1032. public ReverseComparer(IComparer<T> comparer)
  1033. {
  1034. this.comparer = comparer;
  1035. }
  1036. public int Compare(T x, T y)
  1037. {
  1038. return -comparer.Compare(x, y);
  1039. }
  1040. }
  1041. public static IAsyncEnumerable<IAsyncGrouping<TKey, TElement>> GroupBy<TSource, TKey, TElement>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, IEqualityComparer<TKey> comparer)
  1042. {
  1043. if (source == null)
  1044. throw new ArgumentNullException("source");
  1045. if (keySelector == null)
  1046. throw new ArgumentNullException("keySelector");
  1047. if (elementSelector == null)
  1048. throw new ArgumentNullException("elementSelector");
  1049. if (comparer == null)
  1050. throw new ArgumentNullException("comparer");
  1051. return Create(() =>
  1052. {
  1053. var gate = new object();
  1054. var e = source.GetEnumerator();
  1055. var count = 1;
  1056. var map = new Dictionary<TKey, Grouping<TKey, TElement>>(comparer);
  1057. var list = new List<IAsyncGrouping<TKey, TElement>>();
  1058. var index = 0;
  1059. var current = default(IAsyncGrouping<TKey, TElement>);
  1060. var faulted = default(Exception);
  1061. var task = default(Task<bool>);
  1062. var cts = new CancellationTokenDisposable();
  1063. var refCount = new Disposable(
  1064. () =>
  1065. {
  1066. if (Interlocked.Decrement(ref count) == 0)
  1067. e.Dispose();
  1068. }
  1069. );
  1070. var d = new CompositeDisposable(cts, refCount);
  1071. var iterateSource = default(Func<CancellationToken, Task<bool>>);
  1072. iterateSource = ct =>
  1073. {
  1074. var tcs = default(TaskCompletionSource<bool>);
  1075. lock (gate)
  1076. {
  1077. if (task != null)
  1078. {
  1079. return task;
  1080. }
  1081. else
  1082. {
  1083. tcs = new TaskCompletionSource<bool>();
  1084. task = tcs.Task.UsingEnumerator(e);
  1085. }
  1086. }
  1087. if (faulted != null)
  1088. {
  1089. tcs.TrySetException(faulted);
  1090. return task;
  1091. }
  1092. e.MoveNext(ct).ContinueWith(t =>
  1093. {
  1094. t.Handle(tcs,
  1095. res =>
  1096. {
  1097. if (res)
  1098. {
  1099. var key = default(TKey);
  1100. var element = default(TElement);
  1101. var cur = e.Current;
  1102. try
  1103. {
  1104. key = keySelector(cur);
  1105. element = elementSelector(cur);
  1106. }
  1107. catch (Exception exception)
  1108. {
  1109. foreach (var v in map.Values)
  1110. v.Error(exception);
  1111. tcs.TrySetException(exception);
  1112. return;
  1113. }
  1114. var group = default(Grouping<TKey, TElement>);
  1115. if (!map.TryGetValue(key, out group))
  1116. {
  1117. group = new Grouping<TKey, TElement>(key, iterateSource, refCount);
  1118. map.Add(key, group);
  1119. lock (list)
  1120. list.Add(group);
  1121. Interlocked.Increment(ref count);
  1122. }
  1123. group.Add(element);
  1124. }
  1125. tcs.TrySetResult(res);
  1126. },
  1127. ex =>
  1128. {
  1129. foreach (var v in map.Values)
  1130. v.Error(ex);
  1131. faulted = ex;
  1132. tcs.TrySetException(ex);
  1133. }
  1134. );
  1135. lock (gate)
  1136. {
  1137. task = null;
  1138. }
  1139. });
  1140. return tcs.Task.UsingEnumerator(e);
  1141. };
  1142. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1143. f = (tcs, ct) =>
  1144. {
  1145. iterateSource(ct).ContinueWith(t =>
  1146. {
  1147. t.Handle(tcs,
  1148. res =>
  1149. {
  1150. current = null;
  1151. lock (list)
  1152. {
  1153. if (index < list.Count)
  1154. current = list[index++];
  1155. }
  1156. if (current != null)
  1157. {
  1158. tcs.TrySetResult(true);
  1159. }
  1160. else
  1161. {
  1162. if (res)
  1163. f(tcs, ct);
  1164. else
  1165. tcs.TrySetResult(false);
  1166. }
  1167. }
  1168. );
  1169. });
  1170. };
  1171. return Create(
  1172. (ct, tcs) =>
  1173. {
  1174. f(tcs, cts.Token);
  1175. return tcs.Task;
  1176. },
  1177. () => current,
  1178. d.Dispose
  1179. );
  1180. });
  1181. }
  1182. public static IAsyncEnumerable<IAsyncGrouping<TKey, TElement>> GroupBy<TSource, TKey, TElement>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector)
  1183. {
  1184. if (source == null)
  1185. throw new ArgumentNullException("source");
  1186. if (keySelector == null)
  1187. throw new ArgumentNullException("keySelector");
  1188. if (elementSelector == null)
  1189. throw new ArgumentNullException("elementSelector");
  1190. return source.GroupBy(keySelector, elementSelector, EqualityComparer<TKey>.Default);
  1191. }
  1192. public static IAsyncEnumerable<IAsyncGrouping<TKey, TSource>> GroupBy<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
  1193. {
  1194. if (source == null)
  1195. throw new ArgumentNullException("source");
  1196. if (keySelector == null)
  1197. throw new ArgumentNullException("keySelector");
  1198. if (comparer == null)
  1199. throw new ArgumentNullException("comparer");
  1200. return source.GroupBy(keySelector, x => x, comparer);
  1201. }
  1202. public static IAsyncEnumerable<IAsyncGrouping<TKey, TSource>> GroupBy<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
  1203. {
  1204. if (source == null)
  1205. throw new ArgumentNullException("source");
  1206. if (keySelector == null)
  1207. throw new ArgumentNullException("keySelector");
  1208. return source.GroupBy(keySelector, x => x, EqualityComparer<TKey>.Default);
  1209. }
  1210. public static IAsyncEnumerable<TResult> GroupBy<TSource, TKey, TElement, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, Func<TKey, IAsyncEnumerable<TElement>, TResult> resultSelector, IEqualityComparer<TKey> comparer)
  1211. {
  1212. if (source == null)
  1213. throw new ArgumentNullException("source");
  1214. if (keySelector == null)
  1215. throw new ArgumentNullException("keySelector");
  1216. if (elementSelector == null)
  1217. throw new ArgumentNullException("elementSelector");
  1218. if (resultSelector == null)
  1219. throw new ArgumentNullException("resultSelector");
  1220. if (comparer == null)
  1221. throw new ArgumentNullException("comparer");
  1222. return source.GroupBy(keySelector, elementSelector, comparer).Select(g => resultSelector(g.Key, g));
  1223. }
  1224. public static IAsyncEnumerable<TResult> GroupBy<TSource, TKey, TElement, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, Func<TKey, IAsyncEnumerable<TElement>, TResult> resultSelector)
  1225. {
  1226. if (source == null)
  1227. throw new ArgumentNullException("source");
  1228. if (keySelector == null)
  1229. throw new ArgumentNullException("keySelector");
  1230. if (elementSelector == null)
  1231. throw new ArgumentNullException("elementSelector");
  1232. if (resultSelector == null)
  1233. throw new ArgumentNullException("resultSelector");
  1234. return source.GroupBy(keySelector, elementSelector, EqualityComparer<TKey>.Default).Select(g => resultSelector(g.Key, g));
  1235. }
  1236. public static IAsyncEnumerable<TResult> GroupBy<TSource, TKey, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TKey, IAsyncEnumerable<TSource>, TResult> resultSelector, IEqualityComparer<TKey> comparer)
  1237. {
  1238. if (source == null)
  1239. throw new ArgumentNullException("source");
  1240. if (keySelector == null)
  1241. throw new ArgumentNullException("keySelector");
  1242. if (resultSelector == null)
  1243. throw new ArgumentNullException("resultSelector");
  1244. if (comparer == null)
  1245. throw new ArgumentNullException("comparer");
  1246. return source.GroupBy(keySelector, x => x, comparer).Select(g => resultSelector(g.Key, g));
  1247. }
  1248. public static IAsyncEnumerable<TResult> GroupBy<TSource, TKey, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TKey, IAsyncEnumerable<TSource>, TResult> resultSelector)
  1249. {
  1250. if (source == null)
  1251. throw new ArgumentNullException("source");
  1252. if (keySelector == null)
  1253. throw new ArgumentNullException("keySelector");
  1254. if (resultSelector == null)
  1255. throw new ArgumentNullException("resultSelector");
  1256. return source.GroupBy(keySelector, x => x, EqualityComparer<TKey>.Default).Select(g => resultSelector(g.Key, g));
  1257. }
  1258. class Grouping<TKey, TElement> : IAsyncGrouping<TKey, TElement>
  1259. {
  1260. private readonly Func<CancellationToken, Task<bool>> iterateSource;
  1261. private readonly IDisposable sourceDisposable;
  1262. private readonly List<TElement> elements = new List<TElement>();
  1263. private bool done = false;
  1264. private Exception exception = null;
  1265. public Grouping(TKey key, Func<CancellationToken, Task<bool>> iterateSource, IDisposable sourceDisposable)
  1266. {
  1267. this.iterateSource = iterateSource;
  1268. this.sourceDisposable = sourceDisposable;
  1269. Key = key;
  1270. }
  1271. public TKey Key
  1272. {
  1273. get;
  1274. private set;
  1275. }
  1276. public void Add(TElement element)
  1277. {
  1278. lock (elements)
  1279. elements.Add(element);
  1280. }
  1281. public void Error(Exception exception)
  1282. {
  1283. done = true;
  1284. this.exception = exception;
  1285. }
  1286. public IAsyncEnumerator<TElement> GetEnumerator()
  1287. {
  1288. var index = -1;
  1289. var cts = new CancellationTokenDisposable();
  1290. var d = new CompositeDisposable(cts, sourceDisposable);
  1291. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1292. f = (tcs, ct) =>
  1293. {
  1294. var size = 0;
  1295. lock (elements)
  1296. size = elements.Count;
  1297. if (index < size)
  1298. {
  1299. tcs.TrySetResult(true);
  1300. }
  1301. else if (done)
  1302. {
  1303. if (exception != null)
  1304. tcs.TrySetException(exception);
  1305. else
  1306. tcs.TrySetResult(false);
  1307. }
  1308. else
  1309. {
  1310. iterateSource(ct).ContinueWith(t =>
  1311. {
  1312. t.Handle(tcs, res =>
  1313. {
  1314. if (res)
  1315. f(tcs, ct);
  1316. else
  1317. tcs.TrySetResult(false);
  1318. });
  1319. });
  1320. }
  1321. };
  1322. return Create(
  1323. (ct, tcs) =>
  1324. {
  1325. ++index;
  1326. f(tcs, cts.Token);
  1327. return tcs.Task;
  1328. },
  1329. () => elements[index],
  1330. d.Dispose
  1331. );
  1332. }
  1333. }
  1334. #region Ix
  1335. public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext)
  1336. {
  1337. if (source == null)
  1338. throw new ArgumentNullException("source");
  1339. if (onNext == null)
  1340. throw new ArgumentNullException("onNext");
  1341. return DoHelper(source, onNext, _ => { }, () => { });
  1342. }
  1343. public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext, Action onCompleted)
  1344. {
  1345. if (source == null)
  1346. throw new ArgumentNullException("source");
  1347. if (onNext == null)
  1348. throw new ArgumentNullException("onNext");
  1349. if (onCompleted == null)
  1350. throw new ArgumentNullException("onCompleted");
  1351. return DoHelper(source, onNext, _ => { }, onCompleted);
  1352. }
  1353. public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError)
  1354. {
  1355. if (source == null)
  1356. throw new ArgumentNullException("source");
  1357. if (onNext == null)
  1358. throw new ArgumentNullException("onNext");
  1359. if (onError == null)
  1360. throw new ArgumentNullException("onError");
  1361. return DoHelper(source, onNext, onError, () => { });
  1362. }
  1363. public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError, Action onCompleted)
  1364. {
  1365. if (source == null)
  1366. throw new ArgumentNullException("source");
  1367. if (onNext == null)
  1368. throw new ArgumentNullException("onNext");
  1369. if (onError == null)
  1370. throw new ArgumentNullException("onError");
  1371. if (onCompleted == null)
  1372. throw new ArgumentNullException("onCompleted");
  1373. return DoHelper(source, onNext, onError, onCompleted);
  1374. }
  1375. #if !NO_RXINTERFACES
  1376. public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, IObserver<TSource> observer)
  1377. {
  1378. if (source == null)
  1379. throw new ArgumentNullException("source");
  1380. if (observer == null)
  1381. throw new ArgumentNullException("observer");
  1382. return DoHelper(source, observer.OnNext, observer.OnError, observer.OnCompleted);
  1383. }
  1384. #endif
  1385. private static IAsyncEnumerable<TSource> DoHelper<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError, Action onCompleted)
  1386. {
  1387. return Create(() =>
  1388. {
  1389. var e = source.GetEnumerator();
  1390. var cts = new CancellationTokenDisposable();
  1391. var d = new CompositeDisposable(cts, e);
  1392. var current = default(TSource);
  1393. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1394. f = (tcs, ct) =>
  1395. {
  1396. e.MoveNext(ct).ContinueWith(t =>
  1397. {
  1398. if (!t.IsCanceled)
  1399. {
  1400. try
  1401. {
  1402. if (t.IsFaulted)
  1403. {
  1404. onError(t.Exception);
  1405. }
  1406. else if (!t.Result)
  1407. {
  1408. onCompleted();
  1409. }
  1410. else
  1411. {
  1412. current = e.Current;
  1413. onNext(current);
  1414. }
  1415. }
  1416. catch (Exception ex)
  1417. {
  1418. tcs.TrySetException(ex);
  1419. return;
  1420. }
  1421. }
  1422. t.Handle(tcs, res =>
  1423. {
  1424. tcs.TrySetResult(res);
  1425. });
  1426. });
  1427. };
  1428. return Create(
  1429. (ct, tcs) =>
  1430. {
  1431. f(tcs, cts.Token);
  1432. return tcs.Task.UsingEnumerator(e);
  1433. },
  1434. () => current,
  1435. d.Dispose
  1436. );
  1437. });
  1438. }
  1439. public static void ForEach<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> action, CancellationToken cancellationToken)
  1440. {
  1441. if (source == null)
  1442. throw new ArgumentNullException("source");
  1443. if (action == null)
  1444. throw new ArgumentNullException("action");
  1445. source.ForEachAsync(action, cancellationToken).Wait(cancellationToken);
  1446. }
  1447. public static Task ForEachAsync<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> action, CancellationToken cancellationToken)
  1448. {
  1449. if (source == null)
  1450. throw new ArgumentNullException("source");
  1451. if (action == null)
  1452. throw new ArgumentNullException("action");
  1453. return source.ForEachAsync((x, i) => action(x), cancellationToken);
  1454. }
  1455. public static void ForEach<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource, int> action, CancellationToken cancellationToken)
  1456. {
  1457. if (source == null)
  1458. throw new ArgumentNullException("source");
  1459. if (action == null)
  1460. throw new ArgumentNullException("action");
  1461. source.ForEachAsync(action, cancellationToken).Wait(cancellationToken);
  1462. }
  1463. public static Task ForEachAsync<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource, int> action, CancellationToken cancellationToken)
  1464. {
  1465. if (source == null)
  1466. throw new ArgumentNullException("source");
  1467. if (action == null)
  1468. throw new ArgumentNullException("action");
  1469. var tcs = new TaskCompletionSource<bool>();
  1470. var e = source.GetEnumerator();
  1471. var i = 0;
  1472. var f = default(Action<CancellationToken>);
  1473. f = ct =>
  1474. {
  1475. e.MoveNext(ct).ContinueWith(t =>
  1476. {
  1477. t.Handle(tcs, res =>
  1478. {
  1479. if (res)
  1480. {
  1481. try
  1482. {
  1483. action(e.Current, i++);
  1484. }
  1485. catch (Exception ex)
  1486. {
  1487. tcs.TrySetException(ex);
  1488. return;
  1489. }
  1490. f(ct);
  1491. }
  1492. else
  1493. tcs.TrySetResult(true);
  1494. });
  1495. });
  1496. };
  1497. f(cancellationToken);
  1498. return tcs.Task.UsingEnumerator(e);
  1499. }
  1500. public static IAsyncEnumerable<TSource> Repeat<TSource>(this IAsyncEnumerable<TSource> source, int count)
  1501. {
  1502. if (source == null)
  1503. throw new ArgumentNullException("source");
  1504. if (count < 0)
  1505. throw new ArgumentOutOfRangeException("count");
  1506. return Create(() =>
  1507. {
  1508. var e = default(IAsyncEnumerator<TSource>);
  1509. var a = new AssignableDisposable();
  1510. var n = count;
  1511. var current = default(TSource);
  1512. var cts = new CancellationTokenDisposable();
  1513. var d = new CompositeDisposable(cts, a);
  1514. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1515. f = (tcs, ct) =>
  1516. {
  1517. if (e == null)
  1518. {
  1519. if (n-- == 0)
  1520. {
  1521. tcs.TrySetResult(false);
  1522. return;
  1523. }
  1524. try
  1525. {
  1526. e = source.GetEnumerator();
  1527. }
  1528. catch (Exception ex)
  1529. {
  1530. tcs.TrySetException(ex);
  1531. return;
  1532. }
  1533. a.Disposable = e;
  1534. }
  1535. e.MoveNext(ct).ContinueWith(t =>
  1536. {
  1537. t.Handle(tcs, res =>
  1538. {
  1539. if (res)
  1540. {
  1541. current = e.Current;
  1542. tcs.TrySetResult(true);
  1543. }
  1544. else
  1545. {
  1546. e = null;
  1547. f(tcs, ct);
  1548. }
  1549. });
  1550. });
  1551. };
  1552. return Create(
  1553. (ct, tcs) =>
  1554. {
  1555. f(tcs, cts.Token);
  1556. return tcs.Task.UsingEnumerator(d);
  1557. },
  1558. () => current,
  1559. d.Dispose
  1560. );
  1561. });
  1562. }
  1563. public static IAsyncEnumerable<TSource> Repeat<TSource>(this IAsyncEnumerable<TSource> source)
  1564. {
  1565. if (source == null)
  1566. throw new ArgumentNullException("source");
  1567. return Create(() =>
  1568. {
  1569. var e = default(IAsyncEnumerator<TSource>);
  1570. var a = new AssignableDisposable();
  1571. var current = default(TSource);
  1572. var cts = new CancellationTokenDisposable();
  1573. var d = new CompositeDisposable(cts, a);
  1574. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1575. f = (tcs, ct) =>
  1576. {
  1577. if (e == null)
  1578. {
  1579. try
  1580. {
  1581. e = source.GetEnumerator();
  1582. }
  1583. catch (Exception ex)
  1584. {
  1585. tcs.TrySetException(ex);
  1586. return;
  1587. }
  1588. a.Disposable = e;
  1589. }
  1590. e.MoveNext(ct).ContinueWith(t =>
  1591. {
  1592. t.Handle(tcs, res =>
  1593. {
  1594. if (res)
  1595. {
  1596. current = e.Current;
  1597. tcs.TrySetResult(true);
  1598. }
  1599. else
  1600. {
  1601. e = null;
  1602. f(tcs, ct);
  1603. }
  1604. });
  1605. });
  1606. };
  1607. return Create(
  1608. (ct, tcs) =>
  1609. {
  1610. f(tcs, cts.Token);
  1611. return tcs.Task.UsingEnumerator(d);
  1612. },
  1613. () => current,
  1614. d.Dispose
  1615. );
  1616. });
  1617. }
  1618. public static IAsyncEnumerable<TSource> IgnoreElements<TSource>(this IAsyncEnumerable<TSource> source)
  1619. {
  1620. if (source == null)
  1621. throw new ArgumentNullException("source");
  1622. return Create(() =>
  1623. {
  1624. var e = source.GetEnumerator();
  1625. var cts = new CancellationTokenDisposable();
  1626. var d = new CompositeDisposable(cts, e);
  1627. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1628. f = (tcs, ct) =>
  1629. {
  1630. e.MoveNext(ct).ContinueWith(t =>
  1631. {
  1632. t.Handle(tcs, res =>
  1633. {
  1634. if (!res)
  1635. {
  1636. tcs.TrySetResult(false);
  1637. return;
  1638. }
  1639. f(tcs, ct);
  1640. });
  1641. });
  1642. };
  1643. return Create<TSource>(
  1644. (ct, tcs) =>
  1645. {
  1646. f(tcs, cts.Token);
  1647. return tcs.Task.UsingEnumerator(e);
  1648. },
  1649. () => { throw new InvalidOperationException(); },
  1650. d.Dispose
  1651. );
  1652. });
  1653. }
  1654. public static IAsyncEnumerable<TSource> StartWith<TSource>(this IAsyncEnumerable<TSource> source, params TSource[] values)
  1655. {
  1656. if (source == null)
  1657. throw new ArgumentNullException("source");
  1658. return values.ToAsyncEnumerable().Concat(source);
  1659. }
  1660. public static IAsyncEnumerable<IList<TSource>> Buffer<TSource>(this IAsyncEnumerable<TSource> source, int count)
  1661. {
  1662. if (source == null)
  1663. throw new ArgumentNullException("source");
  1664. if (count <= 0)
  1665. throw new ArgumentOutOfRangeException("count");
  1666. return source.Buffer_(count, count);
  1667. }
  1668. public static IAsyncEnumerable<IList<TSource>> Buffer<TSource>(this IAsyncEnumerable<TSource> source, int count, int skip)
  1669. {
  1670. if (source == null)
  1671. throw new ArgumentNullException("source");
  1672. if (count <= 0)
  1673. throw new ArgumentOutOfRangeException("count");
  1674. if (skip <= 0)
  1675. throw new ArgumentOutOfRangeException("skip");
  1676. return source.Buffer_(count, skip);
  1677. }
  1678. private static IAsyncEnumerable<IList<TSource>> Buffer_<TSource>(this IAsyncEnumerable<TSource> source, int count, int skip)
  1679. {
  1680. return Create(() =>
  1681. {
  1682. var e = source.GetEnumerator();
  1683. var cts = new CancellationTokenDisposable();
  1684. var d = new CompositeDisposable(cts, e);
  1685. var buffers = new Queue<IList<TSource>>();
  1686. var i = 0;
  1687. var current = default(IList<TSource>);
  1688. var stopped = false;
  1689. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1690. f = (tcs, ct) =>
  1691. {
  1692. if (!stopped)
  1693. {
  1694. e.MoveNext(ct).ContinueWith(t =>
  1695. {
  1696. t.Handle(tcs, res =>
  1697. {
  1698. if (res)
  1699. {
  1700. var item = e.Current;
  1701. if (i++ % skip == 0)
  1702. buffers.Enqueue(new List<TSource>(count));
  1703. foreach (var buffer in buffers)
  1704. buffer.Add(item);
  1705. if (buffers.Count > 0 && buffers.Peek().Count == count)
  1706. {
  1707. current = buffers.Dequeue();
  1708. tcs.TrySetResult(true);
  1709. return;
  1710. }
  1711. f(tcs, ct);
  1712. }
  1713. else
  1714. {
  1715. stopped = true;
  1716. e.Dispose();
  1717. f(tcs, ct);
  1718. }
  1719. });
  1720. });
  1721. }
  1722. else
  1723. {
  1724. if (buffers.Count > 0)
  1725. {
  1726. current = buffers.Dequeue();
  1727. tcs.TrySetResult(true);
  1728. }
  1729. else
  1730. {
  1731. tcs.TrySetResult(false);
  1732. }
  1733. }
  1734. };
  1735. return Create(
  1736. (ct, tcs) =>
  1737. {
  1738. f(tcs, cts.Token);
  1739. return tcs.Task.UsingEnumerator(e);
  1740. },
  1741. () => current,
  1742. d.Dispose
  1743. );
  1744. });
  1745. }
  1746. public static IAsyncEnumerable<TSource> Distinct<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
  1747. {
  1748. if (source == null)
  1749. throw new ArgumentNullException("source");
  1750. if (keySelector == null)
  1751. throw new ArgumentNullException("keySelector");
  1752. if (comparer == null)
  1753. throw new ArgumentNullException("comparer");
  1754. return Defer(() =>
  1755. {
  1756. var set = new HashSet<TKey>(comparer);
  1757. return source.Where(item => set.Add(keySelector(item)));
  1758. });
  1759. }
  1760. public static IAsyncEnumerable<TSource> Distinct<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
  1761. {
  1762. if (source == null)
  1763. throw new ArgumentNullException("source");
  1764. if (keySelector == null)
  1765. throw new ArgumentNullException("keySelector");
  1766. return source.Distinct(keySelector, EqualityComparer<TKey>.Default);
  1767. }
  1768. public static IAsyncEnumerable<TSource> DistinctUntilChanged<TSource>(this IAsyncEnumerable<TSource> source)
  1769. {
  1770. if (source == null)
  1771. throw new ArgumentNullException("source");
  1772. return source.DistinctUntilChanged_(x => x, EqualityComparer<TSource>.Default);
  1773. }
  1774. public static IAsyncEnumerable<TSource> DistinctUntilChanged<TSource>(this IAsyncEnumerable<TSource> source, IEqualityComparer<TSource> comparer)
  1775. {
  1776. if (source == null)
  1777. throw new ArgumentNullException("source");
  1778. if (comparer == null)
  1779. throw new ArgumentNullException("comparer");
  1780. return source.DistinctUntilChanged_(x => x, comparer);
  1781. }
  1782. public static IAsyncEnumerable<TSource> DistinctUntilChanged<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
  1783. {
  1784. if (source == null)
  1785. throw new ArgumentNullException("source");
  1786. if (keySelector == null)
  1787. throw new ArgumentNullException("keySelector");
  1788. return source.DistinctUntilChanged_(keySelector, EqualityComparer<TKey>.Default);
  1789. }
  1790. public static IAsyncEnumerable<TSource> DistinctUntilChanged<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
  1791. {
  1792. if (source == null)
  1793. throw new ArgumentNullException("source");
  1794. if (keySelector == null)
  1795. throw new ArgumentNullException("keySelector");
  1796. if (comparer == null)
  1797. throw new ArgumentNullException("comparer");
  1798. return source.DistinctUntilChanged_(keySelector, comparer);
  1799. }
  1800. private static IAsyncEnumerable<TSource> DistinctUntilChanged_<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
  1801. {
  1802. return Create(() =>
  1803. {
  1804. var e = source.GetEnumerator();
  1805. var cts = new CancellationTokenDisposable();
  1806. var d = new CompositeDisposable(cts, e);
  1807. var currentKey = default(TKey);
  1808. var hasCurrentKey = false;
  1809. var current = default(TSource);
  1810. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1811. f = (tcs, ct) =>
  1812. {
  1813. e.MoveNext(ct).ContinueWith(t =>
  1814. {
  1815. t.Handle(tcs, res =>
  1816. {
  1817. if (res)
  1818. {
  1819. var item = e.Current;
  1820. var key = default(TKey);
  1821. var comparerEquals = false;
  1822. try
  1823. {
  1824. key = keySelector(item);
  1825. if (hasCurrentKey)
  1826. {
  1827. comparerEquals = comparer.Equals(currentKey, key);
  1828. }
  1829. }
  1830. catch (Exception ex)
  1831. {
  1832. tcs.TrySetException(ex);
  1833. return;
  1834. }
  1835. if (!hasCurrentKey || !comparerEquals)
  1836. {
  1837. hasCurrentKey = true;
  1838. currentKey = key;
  1839. current = item;
  1840. tcs.TrySetResult(true);
  1841. }
  1842. else
  1843. {
  1844. f(tcs, ct);
  1845. }
  1846. }
  1847. else
  1848. {
  1849. tcs.TrySetResult(false);
  1850. }
  1851. });
  1852. });
  1853. };
  1854. return Create(
  1855. (ct, tcs) =>
  1856. {
  1857. f(tcs, cts.Token);
  1858. return tcs.Task.UsingEnumerator(e);
  1859. },
  1860. () => current,
  1861. d.Dispose
  1862. );
  1863. });
  1864. }
  1865. public static IAsyncEnumerable<TSource> Expand<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, IAsyncEnumerable<TSource>> selector)
  1866. {
  1867. if (source == null)
  1868. throw new ArgumentNullException("source");
  1869. if (selector == null)
  1870. throw new ArgumentNullException("selector");
  1871. return Create(() =>
  1872. {
  1873. var e = default(IAsyncEnumerator<TSource>);
  1874. var cts = new CancellationTokenDisposable();
  1875. var a = new AssignableDisposable();
  1876. var d = new CompositeDisposable(cts, a);
  1877. var queue = new Queue<IAsyncEnumerable<TSource>>();
  1878. queue.Enqueue(source);
  1879. var current = default(TSource);
  1880. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1881. f = (tcs, ct) =>
  1882. {
  1883. if (e == null)
  1884. {
  1885. if (queue.Count > 0)
  1886. {
  1887. var src = queue.Dequeue();
  1888. try
  1889. {
  1890. e = src.GetEnumerator();
  1891. }
  1892. catch (Exception ex)
  1893. {
  1894. tcs.TrySetException(ex);
  1895. return;
  1896. }
  1897. a.Disposable = e;
  1898. f(tcs, ct);
  1899. }
  1900. else
  1901. {
  1902. tcs.TrySetResult(false);
  1903. }
  1904. }
  1905. else
  1906. {
  1907. e.MoveNext(ct).ContinueWith(t =>
  1908. {
  1909. t.Handle(tcs, res =>
  1910. {
  1911. if (res)
  1912. {
  1913. var item = e.Current;
  1914. var next = default(IAsyncEnumerable<TSource>);
  1915. try
  1916. {
  1917. next = selector(item);
  1918. }
  1919. catch (Exception ex)
  1920. {
  1921. tcs.TrySetException(ex);
  1922. return;
  1923. }
  1924. queue.Enqueue(next);
  1925. current = item;
  1926. tcs.TrySetResult(true);
  1927. }
  1928. else
  1929. {
  1930. e = null;
  1931. f(tcs, ct);
  1932. }
  1933. });
  1934. });
  1935. }
  1936. };
  1937. return Create(
  1938. (ct, tcs) =>
  1939. {
  1940. f(tcs, cts.Token);
  1941. return tcs.Task.UsingEnumerator(a);
  1942. },
  1943. () => current,
  1944. d.Dispose
  1945. );
  1946. });
  1947. }
  1948. public static IAsyncEnumerable<TAccumulate> Scan<TSource, TAccumulate>(this IAsyncEnumerable<TSource> source, TAccumulate seed, Func<TAccumulate, TSource, TAccumulate> accumulator)
  1949. {
  1950. if (source == null)
  1951. throw new ArgumentNullException("source");
  1952. if (accumulator == null)
  1953. throw new ArgumentNullException("accumulator");
  1954. return Create(() =>
  1955. {
  1956. var e = source.GetEnumerator();
  1957. var cts = new CancellationTokenDisposable();
  1958. var d = new CompositeDisposable(cts, e);
  1959. var acc = seed;
  1960. var current = default(TAccumulate);
  1961. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  1962. f = (tcs, ct) =>
  1963. {
  1964. e.MoveNext(ct).ContinueWith(t =>
  1965. {
  1966. t.Handle(tcs, res =>
  1967. {
  1968. if (!res)
  1969. {
  1970. tcs.TrySetResult(false);
  1971. return;
  1972. }
  1973. var item = e.Current;
  1974. try
  1975. {
  1976. acc = accumulator(acc, item);
  1977. }
  1978. catch (Exception ex)
  1979. {
  1980. tcs.TrySetException(ex);
  1981. return;
  1982. }
  1983. current = acc;
  1984. tcs.TrySetResult(true);
  1985. });
  1986. });
  1987. };
  1988. return Create(
  1989. (ct, tcs) =>
  1990. {
  1991. f(tcs, cts.Token);
  1992. return tcs.Task.UsingEnumerator(e);
  1993. },
  1994. () => current,
  1995. d.Dispose
  1996. );
  1997. });
  1998. }
  1999. public static IAsyncEnumerable<TSource> Scan<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, TSource, TSource> accumulator)
  2000. {
  2001. if (source == null)
  2002. throw new ArgumentNullException("source");
  2003. if (accumulator == null)
  2004. throw new ArgumentNullException("accumulator");
  2005. return Create(() =>
  2006. {
  2007. var e = source.GetEnumerator();
  2008. var cts = new CancellationTokenDisposable();
  2009. var d = new CompositeDisposable(cts, e);
  2010. var hasSeed = false;
  2011. var acc = default(TSource);
  2012. var current = default(TSource);
  2013. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  2014. f = (tcs, ct) =>
  2015. {
  2016. e.MoveNext(ct).ContinueWith(t =>
  2017. {
  2018. t.Handle(tcs, res =>
  2019. {
  2020. if (!res)
  2021. {
  2022. tcs.TrySetResult(false);
  2023. return;
  2024. }
  2025. var item = e.Current;
  2026. if (!hasSeed)
  2027. {
  2028. hasSeed = true;
  2029. acc = item;
  2030. f(tcs, ct);
  2031. return;
  2032. }
  2033. try
  2034. {
  2035. acc = accumulator(acc, item);
  2036. }
  2037. catch (Exception ex)
  2038. {
  2039. tcs.TrySetException(ex);
  2040. return;
  2041. }
  2042. current = acc;
  2043. tcs.TrySetResult(true);
  2044. });
  2045. });
  2046. };
  2047. return Create(
  2048. (ct, tcs) =>
  2049. {
  2050. f(tcs, cts.Token);
  2051. return tcs.Task.UsingEnumerator(e);
  2052. },
  2053. () => current,
  2054. d.Dispose
  2055. );
  2056. });
  2057. }
  2058. public static IAsyncEnumerable<TSource> TakeLast<TSource>(this IAsyncEnumerable<TSource> source, int count)
  2059. {
  2060. if (source == null)
  2061. throw new ArgumentNullException("source");
  2062. if (count < 0)
  2063. throw new ArgumentOutOfRangeException("count");
  2064. return Create(() =>
  2065. {
  2066. var e = source.GetEnumerator();
  2067. var cts = new CancellationTokenDisposable();
  2068. var d = new CompositeDisposable(cts, e);
  2069. var q = new Queue<TSource>(count);
  2070. var done = false;
  2071. var current = default(TSource);
  2072. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  2073. f = (tcs, ct) =>
  2074. {
  2075. if (!done)
  2076. {
  2077. e.MoveNext(ct).ContinueWith(t =>
  2078. {
  2079. t.Handle(tcs, res =>
  2080. {
  2081. if (res)
  2082. {
  2083. var item = e.Current;
  2084. if (q.Count >= count)
  2085. q.Dequeue();
  2086. q.Enqueue(item);
  2087. }
  2088. else
  2089. {
  2090. done = true;
  2091. e.Dispose();
  2092. }
  2093. f(tcs, ct);
  2094. });
  2095. });
  2096. }
  2097. else
  2098. {
  2099. if (q.Count > 0)
  2100. {
  2101. current = q.Dequeue();
  2102. tcs.TrySetResult(true);
  2103. }
  2104. else
  2105. {
  2106. tcs.TrySetResult(false);
  2107. }
  2108. }
  2109. };
  2110. return Create(
  2111. (ct, tcs) =>
  2112. {
  2113. f(tcs, cts.Token);
  2114. return tcs.Task.UsingEnumerator(e);
  2115. },
  2116. () => current,
  2117. d.Dispose
  2118. );
  2119. });
  2120. }
  2121. public static IAsyncEnumerable<TSource> SkipLast<TSource>(this IAsyncEnumerable<TSource> source, int count)
  2122. {
  2123. if (source == null)
  2124. throw new ArgumentNullException("source");
  2125. if (count < 0)
  2126. throw new ArgumentOutOfRangeException("count");
  2127. return Create(() =>
  2128. {
  2129. var e = source.GetEnumerator();
  2130. var cts = new CancellationTokenDisposable();
  2131. var d = new CompositeDisposable(cts, e);
  2132. var q = new Queue<TSource>();
  2133. var current = default(TSource);
  2134. var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  2135. f = (tcs, ct) =>
  2136. {
  2137. e.MoveNext(ct).ContinueWith(t =>
  2138. {
  2139. t.Handle(tcs, res =>
  2140. {
  2141. if (res)
  2142. {
  2143. var item = e.Current;
  2144. q.Enqueue(item);
  2145. if (q.Count > count)
  2146. {
  2147. current = q.Dequeue();
  2148. tcs.TrySetResult(true);
  2149. }
  2150. else
  2151. {
  2152. f(tcs, ct);
  2153. }
  2154. }
  2155. else
  2156. {
  2157. tcs.TrySetResult(false);
  2158. }
  2159. });
  2160. });
  2161. };
  2162. return Create(
  2163. (ct, tcs) =>
  2164. {
  2165. f(tcs, cts.Token);
  2166. return tcs.Task.UsingEnumerator(e);
  2167. },
  2168. () => current,
  2169. d.Dispose
  2170. );
  2171. });
  2172. }
  2173. #endregion
  2174. }
  2175. }