123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355235623572358235923602361236223632364236523662367236823692370237123722373237423752376237723782379238023812382238323842385238623872388238923902391239223932394239523962397239823992400240124022403240424052406240724082409241024112412241324142415241624172418241924202421242224232424242524262427242824292430243124322433243424352436243724382439244024412442244324442445244624472448244924502451245224532454245524562457245824592460246124622463246424652466 |
- // Copyright (c) Microsoft Open Technologies, Inc. All rights reserved. See License.txt in the project root for license information.
- using System;
- using System.Collections.Generic;
- using System.Linq;
- using System.Threading;
- using System.Threading.Tasks;
- namespace System.Linq
- {
- public static partial class AsyncEnumerable
- {
- public static IAsyncEnumerable<TResult> Select<TSource, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, TResult> selector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (selector == null)
- throw new ArgumentNullException("selector");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var current = default(TResult);
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- return Create(
- (ct, tcs) =>
- {
- e.MoveNext(cts.Token).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- try
- {
- current = selector(e.Current);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- tcs.TrySetResult(true);
- }
- else
- {
- tcs.TrySetResult(false);
- }
- });
- });
-
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TResult> Select<TSource, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, int, TResult> selector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (selector == null)
- throw new ArgumentNullException("selector");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var current = default(TResult);
- var index = 0;
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- return Create(
- (ct, tcs) =>
- {
- e.MoveNext(cts.Token).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- try
- {
- current = selector(e.Current, index++);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- tcs.TrySetResult(true);
- }
- else
- {
- tcs.TrySetResult(false);
- }
- });
- });
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> AsAsyncEnumerable<TSource>(this IAsyncEnumerable<TSource> source)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return source.Select(x => x);
- }
- public static IAsyncEnumerable<TSource> Where<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, bool> predicate)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (predicate == null)
- throw new ArgumentNullException("predicate");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var b = false;
- try
- {
- b = predicate(e.Current);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- if (b)
- tcs.TrySetResult(true);
- else
- f(tcs, ct);
- }
- else
- tcs.TrySetResult(false);
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => e.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> Where<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, int, bool> predicate)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (predicate == null)
- throw new ArgumentNullException("predicate");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var index = 0;
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var b = false;
- try
- {
- b = predicate(e.Current, index++);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- if (b)
- tcs.TrySetResult(true);
- else
- f(tcs, ct);
- }
- else
- tcs.TrySetResult(false);
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => e.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TResult> SelectMany<TSource, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, IAsyncEnumerable<TResult>> selector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (selector == null)
- throw new ArgumentNullException("selector");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var ie = default(IAsyncEnumerator<TResult>);
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var outer = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- var inner = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- inner = (tcs, ct) =>
- {
- ie.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- tcs.TrySetResult(true);
- }
- else
- {
- ie = null;
- outer(tcs, ct);
- }
- });
- });
- };
- outer = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- try
- {
- ie = selector(e.Current).GetEnumerator();
- inner(tcs, ct);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- }
- }
- else
- tcs.TrySetResult(false);
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- if (ie == null)
- outer(tcs, cts.Token);
- else
- inner(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => ie.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TResult> SelectMany<TSource, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, int, IAsyncEnumerable<TResult>> selector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (selector == null)
- throw new ArgumentNullException("selector");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var ie = default(IAsyncEnumerator<TResult>);
- var index = 0;
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var outer = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- var inner = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- inner = (tcs, ct) =>
- {
- ie.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- tcs.TrySetResult(true);
- }
- else
- {
- ie = null;
- outer(tcs, ct);
- }
- });
- });
- };
- outer = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- try
- {
- ie = selector(e.Current, index++).GetEnumerator();
- inner(tcs, ct);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- }
- }
- else
- tcs.TrySetResult(false);
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- if (ie == null)
- outer(tcs, cts.Token);
- else
- inner(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => ie.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TResult> SelectMany<TSource, TCollection, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, IAsyncEnumerable<TCollection>> selector, Func<TSource, TCollection, TResult> resultSelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (selector == null)
- throw new ArgumentNullException("selector");
- if (resultSelector == null)
- throw new ArgumentNullException("resultSelector");
- return source.SelectMany(x => selector(x).Select(y => resultSelector(x, y)));
- }
- public static IAsyncEnumerable<TResult> SelectMany<TSource, TCollection, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, int, IAsyncEnumerable<TCollection>> selector, Func<TSource, TCollection, TResult> resultSelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (selector == null)
- throw new ArgumentNullException("selector");
- if (resultSelector == null)
- throw new ArgumentNullException("resultSelector");
- return source.SelectMany((x, i) => selector(x, i).Select(y => resultSelector(x, y)));
- }
- public static IAsyncEnumerable<TType> OfType<TType>(this IAsyncEnumerable<object> source)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return source.Where(x => x is TType).Cast<TType>();
- }
- public static IAsyncEnumerable<TResult> Cast<TResult>(this IAsyncEnumerable<object> source)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return source.Select(x => (TResult)x);
- }
- public static IAsyncEnumerable<TSource> Take<TSource>(this IAsyncEnumerable<TSource> source, int count)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (count < 0)
- throw new ArgumentOutOfRangeException("count");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var n = count;
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- return Create(
- (ct, tcs) =>
- {
- if (n == 0)
- return TaskExt.Return(false, cts.Token);
- e.MoveNext(cts.Token).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- --n;
- tcs.TrySetResult(res);
- });
- });
- return tcs.Task.UsingEnumerator(e);
- },
- () => e.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> TakeWhile<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, bool> predicate)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (predicate == null)
- throw new ArgumentNullException("predicate");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- return Create(
- (ct, tcs) =>
- {
- e.MoveNext(cts.Token).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var b = false;
- try
- {
- b = predicate(e.Current);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- if (b)
- {
- tcs.TrySetResult(true);
- return;
- }
- }
- tcs.TrySetResult(false);
- });
- });
- return tcs.Task.UsingEnumerator(e);
- },
- () => e.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> TakeWhile<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, int, bool> predicate)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (predicate == null)
- throw new ArgumentNullException("predicate");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var index = 0;
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- return Create(
- (ct, tcs) =>
- {
- e.MoveNext(cts.Token).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var b = false;
- try
- {
- b = predicate(e.Current, index++);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- if (b)
- {
- tcs.TrySetResult(true);
- return;
- }
- }
- tcs.TrySetResult(false);
- });
- });
- return tcs.Task.UsingEnumerator(e);
- },
- () => e.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> Skip<TSource>(this IAsyncEnumerable<TSource> source, int count)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (count < 0)
- throw new ArgumentOutOfRangeException("count");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var n = count;
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- if (n == 0)
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, x => tcs.TrySetResult(x));
- });
- else
- {
- --n;
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (!res)
- tcs.TrySetResult(false);
- else
- f(tcs, ct);
- });
- });
- }
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => e.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> SkipWhile<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, bool> predicate)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (predicate == null)
- throw new ArgumentNullException("predicate");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var skipping = true;
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- if (skipping)
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var result = false;
- try
- {
- result = predicate(e.Current);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- if (result)
- f(tcs, ct);
- else
- {
- skipping = false;
- tcs.TrySetResult(true);
- }
- }
- else
- tcs.TrySetResult(false);
- });
- });
- else
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, x => tcs.TrySetResult(x));
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => e.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> SkipWhile<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, int, bool> predicate)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (predicate == null)
- throw new ArgumentNullException("predicate");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var skipping = true;
- var index = 0;
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- if (skipping)
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var result = false;
- try
- {
- result = predicate(e.Current, index++);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- if (result)
- f(tcs, ct);
- else
- {
- skipping = false;
- tcs.TrySetResult(true);
- }
- }
- else
- tcs.TrySetResult(false);
- });
- });
- else
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, x => tcs.TrySetResult(x));
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => e.Current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> DefaultIfEmpty<TSource>(this IAsyncEnumerable<TSource> source, TSource defaultValue)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return Create(() =>
- {
- var done = false;
- var hasElements = false;
- var e = source.GetEnumerator();
- var current = default(TSource);
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- if (done)
- tcs.TrySetResult(false);
- else
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- hasElements = true;
- current = e.Current;
- tcs.TrySetResult(true);
- }
- else
- {
- done = true;
- if (!hasElements)
- {
- current = defaultValue;
- tcs.TrySetResult(true);
- }
- else
- tcs.TrySetResult(false);
- }
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> DefaultIfEmpty<TSource>(this IAsyncEnumerable<TSource> source)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return source.DefaultIfEmpty(default(TSource));
- }
- public static IAsyncEnumerable<TSource> Distinct<TSource>(this IAsyncEnumerable<TSource> source, IEqualityComparer<TSource> comparer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return Defer(() =>
- {
- var set = new HashSet<TSource>(comparer);
- return source.Where(set.Add);
- });
- }
- public static IAsyncEnumerable<TSource> Distinct<TSource>(this IAsyncEnumerable<TSource> source)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return source.Distinct(EqualityComparer<TSource>.Default);
- }
- public static IAsyncEnumerable<TSource> Reverse<TSource>(this IAsyncEnumerable<TSource> source)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var stack = default(Stack<TSource>);
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- return Create(
- (ct, tcs) =>
- {
- if (stack == null)
- {
- Create(() => e).Aggregate(new Stack<TSource>(), (s, x) => { s.Push(x); return s; }, cts.Token).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- stack = res;
- tcs.TrySetResult(stack.Count > 0);
- });
- }, cts.Token);
- }
- else
- {
- stack.Pop();
- tcs.TrySetResult(stack.Count > 0);
- }
- return tcs.Task.UsingEnumerator(e);
- },
- () => stack.Peek(),
- d.Dispose
- );
- });
- }
- public static IOrderedAsyncEnumerable<TSource> OrderBy<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return new OrderedAsyncEnumerable<TSource, TKey>(
- Create(() =>
- {
- var current = default(IEnumerable<TSource>);
- return Create(
- ct =>
- {
- var tcs = new TaskCompletionSource<bool>();
- if (current == null)
- {
- source.ToList(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- current = res;
- tcs.TrySetResult(true);
- });
- });
- }
- else
- tcs.TrySetResult(false);
- return tcs.Task;
- },
- () => current,
- () => { }
- );
- }),
- keySelector,
- comparer
- );
- }
- public static IOrderedAsyncEnumerable<TSource> OrderBy<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- return source.OrderBy(keySelector, Comparer<TKey>.Default);
- }
- public static IOrderedAsyncEnumerable<TSource> OrderByDescending<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return source.OrderBy(keySelector, new ReverseComparer<TKey>(comparer));
- }
- public static IOrderedAsyncEnumerable<TSource> OrderByDescending<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- return source.OrderByDescending(keySelector, Comparer<TKey>.Default);
- }
- public static IOrderedAsyncEnumerable<TSource> ThenBy<TSource, TKey>(this IOrderedAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- return source.ThenBy(keySelector, Comparer<TKey>.Default);
- }
- public static IOrderedAsyncEnumerable<TSource> ThenBy<TSource, TKey>(this IOrderedAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return source.CreateOrderedEnumerable(keySelector, comparer, false);
- }
- public static IOrderedAsyncEnumerable<TSource> ThenByDescending<TSource, TKey>(this IOrderedAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- return source.ThenByDescending(keySelector, Comparer<TKey>.Default);
- }
- public static IOrderedAsyncEnumerable<TSource> ThenByDescending<TSource, TKey>(this IOrderedAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IComparer<TKey> comparer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return source.CreateOrderedEnumerable(keySelector, comparer, true);
- }
- static IEnumerable<IGrouping<TKey, TElement>> GroupUntil<TSource, TKey, TElement>(this IEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector, IComparer<TKey> comparer)
- {
- var group = default(EnumerableGrouping<TKey, TElement>);
- foreach (var x in source)
- {
- var key = keySelector(x);
- if (group == null || comparer.Compare(group.Key, key) != 0)
- {
- group = new EnumerableGrouping<TKey, TElement>(key);
- yield return group;
- }
- group.Add(elementSelector(x));
- }
- }
- class OrderedAsyncEnumerable<T, K> : IOrderedAsyncEnumerable<T>
- {
- private readonly IAsyncEnumerable<IEnumerable<T>> equivalenceClasses;
- private readonly Func<T, K> keySelector;
- private readonly IComparer<K> comparer;
- public OrderedAsyncEnumerable(IAsyncEnumerable<IEnumerable<T>> equivalenceClasses, Func<T, K> keySelector, IComparer<K> comparer)
- {
- this.equivalenceClasses = equivalenceClasses;
- this.keySelector = keySelector;
- this.comparer = comparer;
- }
- public IOrderedAsyncEnumerable<T> CreateOrderedEnumerable<TKey>(Func<T, TKey> keySelector, IComparer<TKey> comparer, bool descending)
- {
- if (descending)
- comparer = new ReverseComparer<TKey>(comparer);
- return new OrderedAsyncEnumerable<T, TKey>(Classes(), keySelector, comparer);
- }
- IAsyncEnumerable<IEnumerable<T>> Classes()
- {
- return Create(() =>
- {
- var e = equivalenceClasses.GetEnumerator();
- var list = new List<IEnumerable<T>>();
- var e1 = default(IEnumerator<IEnumerable<T>>);
- var cts = new CancellationTokenDisposable();
- var d1 = new AssignableDisposable();
- var d = new CompositeDisposable(cts, e, d1);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- var g = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- try
- {
- foreach (var group in e.Current.OrderBy(keySelector, comparer).GroupUntil(keySelector, x => x, comparer))
- list.Add(group);
- f(tcs, ct);
- }
- catch (Exception exception)
- {
- tcs.TrySetException(exception);
- return;
- }
- }
- else
- {
- e.Dispose();
- e1 = list.GetEnumerator();
- d1.Disposable = e1;
- g(tcs, ct);
- }
- });
- });
- };
- g = (tcs, ct) =>
- {
- var res = false;
- try
- {
- res = e1.MoveNext();
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- tcs.TrySetResult(res);
- };
- return Create(
- (ct, tcs) =>
- {
- if (e1 != null)
- {
- g(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e1);
- }
- else
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- }
- },
- () => e1.Current,
- d.Dispose
- );
- });
- }
- public IAsyncEnumerator<T> GetEnumerator()
- {
- return Classes().SelectMany(x => x.ToAsyncEnumerable()).GetEnumerator();
- }
- }
- class ReverseComparer<T> : IComparer<T>
- {
- IComparer<T> comparer;
- public ReverseComparer(IComparer<T> comparer)
- {
- this.comparer = comparer;
- }
- public int Compare(T x, T y)
- {
- return -comparer.Compare(x, y);
- }
- }
- 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)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (elementSelector == null)
- throw new ArgumentNullException("elementSelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return Create(() =>
- {
- var gate = new object();
-
- var e = source.GetEnumerator();
- var count = 1;
- var map = new Dictionary<TKey, Grouping<TKey, TElement>>(comparer);
- var list = new List<IAsyncGrouping<TKey, TElement>>();
- var index = 0;
- var current = default(IAsyncGrouping<TKey, TElement>);
- var faulted = default(Exception);
- var task = default(Task<bool>);
- var cts = new CancellationTokenDisposable();
- var refCount = new Disposable(
- () =>
- {
- if (Interlocked.Decrement(ref count) == 0)
- e.Dispose();
- }
- );
- var d = new CompositeDisposable(cts, refCount);
- var iterateSource = default(Func<CancellationToken, Task<bool>>);
- iterateSource = ct =>
- {
- var tcs = default(TaskCompletionSource<bool>);
- lock (gate)
- {
- if (task != null)
- {
- return task;
- }
- else
- {
- tcs = new TaskCompletionSource<bool>();
- task = tcs.Task.UsingEnumerator(e);
- }
- }
- if (faulted != null)
- {
- tcs.TrySetException(faulted);
- return task;
- }
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs,
- res =>
- {
- if (res)
- {
- var key = default(TKey);
- var element = default(TElement);
- var cur = e.Current;
- try
- {
- key = keySelector(cur);
- element = elementSelector(cur);
- }
- catch (Exception exception)
- {
- foreach (var v in map.Values)
- v.Error(exception);
- tcs.TrySetException(exception);
- return;
- }
- var group = default(Grouping<TKey, TElement>);
- if (!map.TryGetValue(key, out group))
- {
- group = new Grouping<TKey, TElement>(key, iterateSource, refCount);
- map.Add(key, group);
- lock (list)
- list.Add(group);
- Interlocked.Increment(ref count);
- }
- group.Add(element);
- }
- tcs.TrySetResult(res);
- },
- ex =>
- {
- foreach (var v in map.Values)
- v.Error(ex);
- faulted = ex;
- tcs.TrySetException(ex);
- }
- );
- lock (gate)
- {
- task = null;
- }
- });
- return tcs.Task.UsingEnumerator(e);
- };
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- iterateSource(ct).ContinueWith(t =>
- {
- t.Handle(tcs,
- res =>
- {
- current = null;
- lock (list)
- {
- if (index < list.Count)
- current = list[index++];
- }
- if (current != null)
- {
- tcs.TrySetResult(true);
- }
- else
- {
- if (res)
- f(tcs, ct);
- else
- tcs.TrySetResult(false);
- }
- }
- );
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task;
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<IAsyncGrouping<TKey, TElement>> GroupBy<TSource, TKey, TElement>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TSource, TElement> elementSelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (elementSelector == null)
- throw new ArgumentNullException("elementSelector");
-
- return source.GroupBy(keySelector, elementSelector, EqualityComparer<TKey>.Default);
- }
- public static IAsyncEnumerable<IAsyncGrouping<TKey, TSource>> GroupBy<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return source.GroupBy(keySelector, x => x, comparer);
- }
- public static IAsyncEnumerable<IAsyncGrouping<TKey, TSource>> GroupBy<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
-
- return source.GroupBy(keySelector, x => x, EqualityComparer<TKey>.Default);
- }
- 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)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (elementSelector == null)
- throw new ArgumentNullException("elementSelector");
- if (resultSelector == null)
- throw new ArgumentNullException("resultSelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return source.GroupBy(keySelector, elementSelector, comparer).Select(g => resultSelector(g.Key, g));
- }
- 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)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (elementSelector == null)
- throw new ArgumentNullException("elementSelector");
- if (resultSelector == null)
- throw new ArgumentNullException("resultSelector");
- return source.GroupBy(keySelector, elementSelector, EqualityComparer<TKey>.Default).Select(g => resultSelector(g.Key, g));
- }
- 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)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (resultSelector == null)
- throw new ArgumentNullException("resultSelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return source.GroupBy(keySelector, x => x, comparer).Select(g => resultSelector(g.Key, g));
- }
- public static IAsyncEnumerable<TResult> GroupBy<TSource, TKey, TResult>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, Func<TKey, IAsyncEnumerable<TSource>, TResult> resultSelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (resultSelector == null)
- throw new ArgumentNullException("resultSelector");
-
- return source.GroupBy(keySelector, x => x, EqualityComparer<TKey>.Default).Select(g => resultSelector(g.Key, g));
- }
- class Grouping<TKey, TElement> : IAsyncGrouping<TKey, TElement>
- {
- private readonly Func<CancellationToken, Task<bool>> iterateSource;
- private readonly IDisposable sourceDisposable;
- private readonly List<TElement> elements = new List<TElement>();
- private bool done = false;
- private Exception exception = null;
- public Grouping(TKey key, Func<CancellationToken, Task<bool>> iterateSource, IDisposable sourceDisposable)
- {
- this.iterateSource = iterateSource;
- this.sourceDisposable = sourceDisposable;
- Key = key;
- }
- public TKey Key
- {
- get;
- private set;
- }
- public void Add(TElement element)
- {
- lock (elements)
- elements.Add(element);
- }
- public void Error(Exception exception)
- {
- done = true;
- this.exception = exception;
- }
- public IAsyncEnumerator<TElement> GetEnumerator()
- {
- var index = -1;
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, sourceDisposable);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- var size = 0;
- lock (elements)
- size = elements.Count;
- if (index < size)
- {
- tcs.TrySetResult(true);
- }
- else if (done)
- {
- if (exception != null)
- tcs.TrySetException(exception);
- else
- tcs.TrySetResult(false);
- }
- else
- {
- iterateSource(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- f(tcs, ct);
- else
- tcs.TrySetResult(false);
- });
- });
- }
- };
- return Create(
- (ct, tcs) =>
- {
- ++index;
- f(tcs, cts.Token);
- return tcs.Task;
- },
- () => elements[index],
- d.Dispose
- );
- }
- }
- #region Ix
- public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (onNext == null)
- throw new ArgumentNullException("onNext");
- return DoHelper(source, onNext, _ => { }, () => { });
- }
- public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext, Action onCompleted)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (onNext == null)
- throw new ArgumentNullException("onNext");
- if (onCompleted == null)
- throw new ArgumentNullException("onCompleted");
- return DoHelper(source, onNext, _ => { }, onCompleted);
- }
- public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (onNext == null)
- throw new ArgumentNullException("onNext");
- if (onError == null)
- throw new ArgumentNullException("onError");
- return DoHelper(source, onNext, onError, () => { });
- }
- public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError, Action onCompleted)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (onNext == null)
- throw new ArgumentNullException("onNext");
- if (onError == null)
- throw new ArgumentNullException("onError");
- if (onCompleted == null)
- throw new ArgumentNullException("onCompleted");
- return DoHelper(source, onNext, onError, onCompleted);
- }
- #if !NO_RXINTERFACES
- public static IAsyncEnumerable<TSource> Do<TSource>(this IAsyncEnumerable<TSource> source, IObserver<TSource> observer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (observer == null)
- throw new ArgumentNullException("observer");
- return DoHelper(source, observer.OnNext, observer.OnError, observer.OnCompleted);
- }
- #endif
- private static IAsyncEnumerable<TSource> DoHelper<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> onNext, Action<Exception> onError, Action onCompleted)
- {
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var current = default(TSource);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- if (!t.IsCanceled)
- {
- try
- {
- if (t.IsFaulted)
- {
- onError(t.Exception);
- }
- else if (!t.Result)
- {
- onCompleted();
- }
- else
- {
- current = e.Current;
- onNext(current);
- }
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- }
- t.Handle(tcs, res =>
- {
- tcs.TrySetResult(res);
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static void ForEach<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> action, CancellationToken cancellationToken)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (action == null)
- throw new ArgumentNullException("action");
- source.ForEachAsync(action, cancellationToken).Wait(cancellationToken);
- }
- public static Task ForEachAsync<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource> action, CancellationToken cancellationToken)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (action == null)
- throw new ArgumentNullException("action");
- return source.ForEachAsync((x, i) => action(x), cancellationToken);
- }
- public static void ForEach<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource, int> action, CancellationToken cancellationToken)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (action == null)
- throw new ArgumentNullException("action");
- source.ForEachAsync(action, cancellationToken).Wait(cancellationToken);
- }
- public static Task ForEachAsync<TSource>(this IAsyncEnumerable<TSource> source, Action<TSource, int> action, CancellationToken cancellationToken)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (action == null)
- throw new ArgumentNullException("action");
- var tcs = new TaskCompletionSource<bool>();
- var e = source.GetEnumerator();
- var i = 0;
- var f = default(Action<CancellationToken>);
- f = ct =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- try
- {
- action(e.Current, i++);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- f(ct);
- }
- else
- tcs.TrySetResult(true);
- });
- });
- };
- f(cancellationToken);
- return tcs.Task.UsingEnumerator(e);
- }
- public static IAsyncEnumerable<TSource> Repeat<TSource>(this IAsyncEnumerable<TSource> source, int count)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (count < 0)
- throw new ArgumentOutOfRangeException("count");
- return Create(() =>
- {
- var e = default(IAsyncEnumerator<TSource>);
- var a = new AssignableDisposable();
- var n = count;
- var current = default(TSource);
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, a);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- if (e == null)
- {
- if (n-- == 0)
- {
- tcs.TrySetResult(false);
- return;
- }
- try
- {
- e = source.GetEnumerator();
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- a.Disposable = e;
- }
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- current = e.Current;
- tcs.TrySetResult(true);
- }
- else
- {
- e = null;
- f(tcs, ct);
- }
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(d);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> Repeat<TSource>(this IAsyncEnumerable<TSource> source)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return Create(() =>
- {
- var e = default(IAsyncEnumerator<TSource>);
- var a = new AssignableDisposable();
- var current = default(TSource);
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, a);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- if (e == null)
- {
- try
- {
- e = source.GetEnumerator();
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- a.Disposable = e;
- }
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- current = e.Current;
- tcs.TrySetResult(true);
- }
- else
- {
- e = null;
- f(tcs, ct);
- }
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(d);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> IgnoreElements<TSource>(this IAsyncEnumerable<TSource> source)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (!res)
- {
- tcs.TrySetResult(false);
- return;
- }
- f(tcs, ct);
- });
- });
- };
- return Create<TSource>(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => { throw new InvalidOperationException(); },
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> StartWith<TSource>(this IAsyncEnumerable<TSource> source, params TSource[] values)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return values.ToAsyncEnumerable().Concat(source);
- }
- public static IAsyncEnumerable<IList<TSource>> Buffer<TSource>(this IAsyncEnumerable<TSource> source, int count)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (count <= 0)
- throw new ArgumentOutOfRangeException("count");
- return source.Buffer_(count, count);
- }
- public static IAsyncEnumerable<IList<TSource>> Buffer<TSource>(this IAsyncEnumerable<TSource> source, int count, int skip)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (count <= 0)
- throw new ArgumentOutOfRangeException("count");
- if (skip <= 0)
- throw new ArgumentOutOfRangeException("skip");
- return source.Buffer_(count, skip);
- }
- private static IAsyncEnumerable<IList<TSource>> Buffer_<TSource>(this IAsyncEnumerable<TSource> source, int count, int skip)
- {
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var buffers = new Queue<IList<TSource>>();
- var i = 0;
- var current = default(IList<TSource>);
- var stopped = false;
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- if (!stopped)
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var item = e.Current;
- if (i++ % skip == 0)
- buffers.Enqueue(new List<TSource>(count));
- foreach (var buffer in buffers)
- buffer.Add(item);
- if (buffers.Count > 0 && buffers.Peek().Count == count)
- {
- current = buffers.Dequeue();
- tcs.TrySetResult(true);
- return;
- }
- f(tcs, ct);
- }
- else
- {
- stopped = true;
- e.Dispose();
- f(tcs, ct);
- }
- });
- });
- }
- else
- {
- if (buffers.Count > 0)
- {
- current = buffers.Dequeue();
- tcs.TrySetResult(true);
- }
- else
- {
- tcs.TrySetResult(false);
- }
- }
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> Distinct<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return Defer(() =>
- {
- var set = new HashSet<TKey>(comparer);
- return source.Where(item => set.Add(keySelector(item)));
- });
- }
- public static IAsyncEnumerable<TSource> Distinct<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- return source.Distinct(keySelector, EqualityComparer<TKey>.Default);
- }
- public static IAsyncEnumerable<TSource> DistinctUntilChanged<TSource>(this IAsyncEnumerable<TSource> source)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- return source.DistinctUntilChanged_(x => x, EqualityComparer<TSource>.Default);
- }
- public static IAsyncEnumerable<TSource> DistinctUntilChanged<TSource>(this IAsyncEnumerable<TSource> source, IEqualityComparer<TSource> comparer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return source.DistinctUntilChanged_(x => x, comparer);
- }
- public static IAsyncEnumerable<TSource> DistinctUntilChanged<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- return source.DistinctUntilChanged_(keySelector, EqualityComparer<TKey>.Default);
- }
- public static IAsyncEnumerable<TSource> DistinctUntilChanged<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (keySelector == null)
- throw new ArgumentNullException("keySelector");
- if (comparer == null)
- throw new ArgumentNullException("comparer");
- return source.DistinctUntilChanged_(keySelector, comparer);
- }
- private static IAsyncEnumerable<TSource> DistinctUntilChanged_<TSource, TKey>(this IAsyncEnumerable<TSource> source, Func<TSource, TKey> keySelector, IEqualityComparer<TKey> comparer)
- {
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var currentKey = default(TKey);
- var hasCurrentKey = false;
- var current = default(TSource);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var item = e.Current;
- var key = default(TKey);
- var comparerEquals = false;
- try
- {
- key = keySelector(item);
- if (hasCurrentKey)
- {
- comparerEquals = comparer.Equals(currentKey, key);
- }
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- if (!hasCurrentKey || !comparerEquals)
- {
- hasCurrentKey = true;
- currentKey = key;
- current = item;
- tcs.TrySetResult(true);
- }
- else
- {
- f(tcs, ct);
- }
- }
- else
- {
- tcs.TrySetResult(false);
- }
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> Expand<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, IAsyncEnumerable<TSource>> selector)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (selector == null)
- throw new ArgumentNullException("selector");
- return Create(() =>
- {
- var e = default(IAsyncEnumerator<TSource>);
- var cts = new CancellationTokenDisposable();
- var a = new AssignableDisposable();
- var d = new CompositeDisposable(cts, a);
- var queue = new Queue<IAsyncEnumerable<TSource>>();
- queue.Enqueue(source);
- var current = default(TSource);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- if (e == null)
- {
- if (queue.Count > 0)
- {
- var src = queue.Dequeue();
- try
- {
- e = src.GetEnumerator();
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- a.Disposable = e;
- f(tcs, ct);
- }
- else
- {
- tcs.TrySetResult(false);
- }
- }
- else
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var item = e.Current;
- var next = default(IAsyncEnumerable<TSource>);
- try
- {
- next = selector(item);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- queue.Enqueue(next);
- current = item;
- tcs.TrySetResult(true);
- }
- else
- {
- e = null;
- f(tcs, ct);
- }
- });
- });
- }
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(a);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TAccumulate> Scan<TSource, TAccumulate>(this IAsyncEnumerable<TSource> source, TAccumulate seed, Func<TAccumulate, TSource, TAccumulate> accumulator)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (accumulator == null)
- throw new ArgumentNullException("accumulator");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var acc = seed;
- var current = default(TAccumulate);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (!res)
- {
- tcs.TrySetResult(false);
- return;
- }
- var item = e.Current;
- try
- {
- acc = accumulator(acc, item);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- current = acc;
- tcs.TrySetResult(true);
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> Scan<TSource>(this IAsyncEnumerable<TSource> source, Func<TSource, TSource, TSource> accumulator)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (accumulator == null)
- throw new ArgumentNullException("accumulator");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var hasSeed = false;
- var acc = default(TSource);
- var current = default(TSource);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (!res)
- {
- tcs.TrySetResult(false);
- return;
- }
- var item = e.Current;
- if (!hasSeed)
- {
- hasSeed = true;
- acc = item;
- f(tcs, ct);
- return;
- }
- try
- {
- acc = accumulator(acc, item);
- }
- catch (Exception ex)
- {
- tcs.TrySetException(ex);
- return;
- }
- current = acc;
- tcs.TrySetResult(true);
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> TakeLast<TSource>(this IAsyncEnumerable<TSource> source, int count)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (count < 0)
- throw new ArgumentOutOfRangeException("count");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var q = new Queue<TSource>(count);
- var done = false;
- var current = default(TSource);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- if (!done)
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var item = e.Current;
- if (q.Count >= count)
- q.Dequeue();
- q.Enqueue(item);
- }
- else
- {
- done = true;
- e.Dispose();
- }
- f(tcs, ct);
- });
- });
- }
- else
- {
- if (q.Count > 0)
- {
- current = q.Dequeue();
- tcs.TrySetResult(true);
- }
- else
- {
- tcs.TrySetResult(false);
- }
- }
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- public static IAsyncEnumerable<TSource> SkipLast<TSource>(this IAsyncEnumerable<TSource> source, int count)
- {
- if (source == null)
- throw new ArgumentNullException("source");
- if (count < 0)
- throw new ArgumentOutOfRangeException("count");
- return Create(() =>
- {
- var e = source.GetEnumerator();
- var cts = new CancellationTokenDisposable();
- var d = new CompositeDisposable(cts, e);
- var q = new Queue<TSource>();
- var current = default(TSource);
- var f = default(Action<TaskCompletionSource<bool>, CancellationToken>);
- f = (tcs, ct) =>
- {
- e.MoveNext(ct).ContinueWith(t =>
- {
- t.Handle(tcs, res =>
- {
- if (res)
- {
- var item = e.Current;
- q.Enqueue(item);
- if (q.Count > count)
- {
- current = q.Dequeue();
- tcs.TrySetResult(true);
- }
- else
- {
- f(tcs, ct);
- }
- }
- else
- {
- tcs.TrySetResult(false);
- }
- });
- });
- };
- return Create(
- (ct, tcs) =>
- {
- f(tcs, cts.Token);
- return tcs.Task.UsingEnumerator(e);
- },
- () => current,
- d.Dispose
- );
- });
- }
- #endregion
- }
- }
|