AsyncEnumerable.Single.cs 92 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503
  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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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. var localIe = default(IAsyncEnumerator<TResult>);
  227. lock (syncRoot)
  228. {
  229. localIe = ie;
  230. }
  231. if (localIe != null)
  232. localIe.Dispose();
  233. });
  234. var cts = new CancellationTokenDisposable();
  235. var d = Disposable.Create(cts, new Disposable(disposeIe), e);
  236. var outer = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  237. var inner = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  238. inner = (tcs, ct) =>
  239. {
  240. ie.MoveNext(ct).ContinueWith(t =>
  241. {
  242. t.Handle(tcs, res =>
  243. {
  244. if (res)
  245. {
  246. tcs.TrySetResult(true);
  247. }
  248. else
  249. {
  250. disposeIe();
  251. outer(tcs, ct);
  252. }
  253. });
  254. });
  255. };
  256. outer = (tcs, ct) =>
  257. {
  258. e.MoveNext(ct).ContinueWith(t =>
  259. {
  260. t.Handle(tcs, res =>
  261. {
  262. if (res)
  263. {
  264. try
  265. {
  266. ie = selector(e.Current).GetEnumerator();
  267. inner(tcs, ct);
  268. }
  269. catch (Exception ex)
  270. {
  271. tcs.TrySetException(ex);
  272. }
  273. }
  274. else
  275. tcs.TrySetResult(false);
  276. });
  277. });
  278. };
  279. return Create(
  280. (ct, tcs) =>
  281. {
  282. if (ie == null)
  283. outer(tcs, cts.Token);
  284. else
  285. inner(tcs, cts.Token);
  286. return tcs.Task.UsingEnumerator(e);
  287. },
  288. () => ie.Current,
  289. d.Dispose
  290. );
  291. });
  292. }
  293. public static IAsyncEnumerable<TResult> SelectMany<TSource, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, int, IAsyncEnumerable<TResult>> selector)
  294. {
  295. if (source == null)
  296. throw new ArgumentNullException("source");
  297. if (selector == null)
  298. throw new ArgumentNullException("selector");
  299. return Create(() =>
  300. {
  301. // A lock seems inevitable. Disposal of the outer enumerator and completion of
  302. // MoveNext of the inner enumerator can happen concurrently.
  303. var syncRoot = new object();
  304. var e = source.GetEnumerator();
  305. var ie = default(IAsyncEnumerator<TResult>);
  306. var disposeIe = new Action(() =>
  307. {
  308. var localIe = default(IAsyncEnumerator<TResult>);
  309. lock (syncRoot)
  310. {
  311. localIe = ie;
  312. }
  313. if (localIe != null)
  314. localIe.Dispose();
  315. });
  316. var index = 0;
  317. var cts = new CancellationTokenDisposable();
  318. var d = Disposable.Create(cts, new Disposable(disposeIe), e);
  319. var outer = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  320. var inner = default(Action<TaskCompletionSource<bool>, CancellationToken>);
  321. inner = (tcs, ct) =>
  322. {
  323. ie.MoveNext(ct).ContinueWith(t =>
  324. {
  325. t.Handle(tcs, res =>
  326. {
  327. if (res)
  328. {
  329. tcs.TrySetResult(true);
  330. }
  331. else
  332. {
  333. disposeIe();
  334. outer(tcs, ct);
  335. }
  336. });
  337. });
  338. };
  339. outer = (tcs, ct) =>
  340. {
  341. e.MoveNext(ct).ContinueWith(t =>
  342. {
  343. t.Handle(tcs, res =>
  344. {
  345. if (res)
  346. {
  347. try
  348. {
  349. ie = selector(e.Current, index++).GetEnumerator();
  350. inner(tcs, ct);
  351. }
  352. catch (Exception ex)
  353. {
  354. tcs.TrySetException(ex);
  355. }
  356. }
  357. else
  358. tcs.TrySetResult(false);
  359. });
  360. });
  361. };
  362. return Create(
  363. (ct, tcs) =>
  364. {
  365. if (ie == null)
  366. outer(tcs, cts.Token);
  367. else
  368. inner(tcs, cts.Token);
  369. return tcs.Task.UsingEnumerator(e);
  370. },
  371. () => ie.Current,
  372. d.Dispose
  373. );
  374. });
  375. }
  376. public static IAsyncEnumerable<TResult> SelectMany<TSource, TCollection, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, IAsyncEnumerable<TCollection>> selector, Func<TSource, TCollection, TResult> resultSelector)
  377. {
  378. if (source == null)
  379. throw new ArgumentNullException("source");
  380. if (selector == null)
  381. throw new ArgumentNullException("selector");
  382. if (resultSelector == null)
  383. throw new ArgumentNullException("resultSelector");
  384. return source.SelectMany(x => selector(x).Select(y => resultSelector(x, y)));
  385. }
  386. public static IAsyncEnumerable<TResult> SelectMany<TSource, TCollection, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, int, IAsyncEnumerable<TCollection>> selector, Func<TSource, TCollection, TResult> resultSelector)
  387. {
  388. if (source == null)
  389. throw new ArgumentNullException("source");
  390. if (selector == null)
  391. throw new ArgumentNullException("selector");
  392. if (resultSelector == null)
  393. throw new ArgumentNullException("resultSelector");
  394. return source.SelectMany((x, i) => selector(x, i).Select(y => resultSelector(x, y)));
  395. }
  396. public static IAsyncEnumerable<TType> OfType<TType>(this IAsyncEnumerable<object> source)
  397. {
  398. if (source == null)
  399. throw new ArgumentNullException("source");
  400. return source.Where(x => x is TType).Cast<TType>();
  401. }
  402. public static IAsyncEnumerable<TResult> Cast<TResult>(this IAsyncEnumerable<object> source)
  403. {
  404. if (source == null)
  405. throw new ArgumentNullException("source");
  406. return source.Select(x => (TResult)x);
  407. }
  408. public static IAsyncEnumerable<TSource> Take<TSource>(this IAsyncEnumerable<TSource> source, int count)
  409. {
  410. if (source == null)
  411. throw new ArgumentNullException("source");
  412. if (count < 0)
  413. throw new ArgumentOutOfRangeException("count");
  414. return Create(() =>
  415. {
  416. var e = source.GetEnumerator();
  417. var n = count;
  418. var cts = new CancellationTokenDisposable();
  419. var d = Disposable.Create(cts, e);
  420. return Create(
  421. (ct, tcs) =>
  422. {
  423. if (n == 0)
  424. return TaskExt.Return(false, cts.Token);
  425. e.MoveNext(cts.Token).ContinueWith(t =>
  426. {
  427. t.Handle(tcs, res =>
  428. {
  429. --n;
  430. if (n == 0)
  431. e.Dispose();
  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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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 = Disposable.Create(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. }