GroupByUntilTest.cs 134 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355235623572358235923602361236223632364236523662367236823692370237123722373237423752376237723782379238023812382238323842385238623872388238923902391239223932394239523962397239823992400240124022403240424052406240724082409241024112412241324142415241624172418241924202421242224232424242524262427242824292430243124322433243424352436243724382439244024412442244324442445244624472448244924502451245224532454245524562457245824592460246124622463246424652466246724682469247024712472247324742475247624772478247924802481248224832484248524862487248824892490249124922493249424952496249724982499250025012502250325042505250625072508250925102511251225132514251525162517251825192520252125222523252425252526252725282529253025312532253325342535253625372538253925402541254225432544254525462547254825492550255125522553255425552556255725582559256025612562256325642565256625672568256925702571257225732574257525762577257825792580258125822583258425852586258725882589259025912592259325942595259625972598259926002601260226032604260526062607260826092610261126122613261426152616261726182619262026212622262326242625262626272628262926302631263226332634263526362637263826392640264126422643264426452646264726482649265026512652265326542655265626572658265926602661266226632664266526662667266826692670267126722673267426752676267726782679268026812682268326842685268626872688268926902691269226932694269526962697269826992700270127022703270427052706270727082709271027112712271327142715271627172718271927202721272227232724272527262727272827292730273127322733273427352736273727382739274027412742274327442745274627472748274927502751275227532754275527562757275827592760276127622763276427652766276727682769277027712772277327742775277627772778277927802781278227832784278527862787278827892790279127922793279427952796279727982799280028012802280328042805280628072808280928102811281228132814281528162817281828192820282128222823282428252826282728282829283028312832283328342835283628372838283928402841284228432844284528462847284828492850285128522853285428552856285728582859286028612862286328642865286628672868286928702871287228732874287528762877287828792880288128822883288428852886288728882889289028912892289328942895289628972898289929002901290229032904290529062907290829092910291129122913291429152916291729182919292029212922292329242925292629272928292929302931293229332934293529362937293829392940294129422943294429452946294729482949295029512952295329542955295629572958295929602961296229632964296529662967296829692970297129722973297429752976297729782979298029812982298329842985298629872988298929902991299229932994299529962997299829993000300130023003300430053006300730083009301030113012301330143015301630173018301930203021302230233024302530263027302830293030303130323033303430353036303730383039304030413042304330443045304630473048304930503051305230533054305530563057305830593060306130623063306430653066306730683069307030713072307330743075307630773078307930803081308230833084308530863087308830893090309130923093309430953096309730983099310031013102310331043105310631073108310931103111311231133114311531163117311831193120312131223123312431253126312731283129313031313132313331343135313631373138313931403141314231433144314531463147314831493150315131523153315431553156315731583159316031613162316331643165316631673168316931703171317231733174317531763177317831793180318131823183318431853186318731883189319031913192319331943195319631973198319932003201320232033204320532063207320832093210321132123213321432153216321732183219322032213222322332243225322632273228322932303231323232333234323532363237323832393240324132423243324432453246324732483249325032513252325332543255325632573258325932603261326232633264326532663267326832693270327132723273327432753276327732783279328032813282328332843285328632873288328932903291329232933294329532963297329832993300330133023303330433053306330733083309331033113312331333143315331633173318331933203321332233233324332533263327332833293330333133323333333433353336333733383339334033413342334333443345334633473348334933503351335233533354335533563357335833593360336133623363336433653366336733683369337033713372337333743375337633773378337933803381338233833384338533863387338833893390339133923393339433953396339733983399340034013402340334043405340634073408340934103411341234133414341534163417341834193420342134223423342434253426342734283429343034313432343334343435343634373438343934403441344234433444344534463447344834493450345134523453345434553456345734583459346034613462346334643465346634673468346934703471347234733474347534763477347834793480348134823483348434853486348734883489349034913492349334943495349634973498349935003501350235033504350535063507350835093510351135123513351435153516351735183519352035213522352335243525352635273528352935303531353235333534353535363537353835393540354135423543354435453546354735483549355035513552355335543555355635573558355935603561356235633564356535663567356835693570357135723573357435753576357735783579358035813582358335843585358635873588358935903591359235933594359535963597359835993600360136023603360436053606360736083609361036113612361336143615361636173618361936203621362236233624362536263627362836293630363136323633363436353636363736383639364036413642364336443645364636473648364936503651365236533654365536563657365836593660366136623663366436653666366736683669367036713672367336743675367636773678367936803681368236833684368536863687368836893690369136923693369436953696369736983699370037013702370337043705370637073708370937103711371237133714371537163717371837193720372137223723372437253726372737283729373037313732373337343735373637373738373937403741374237433744374537463747374837493750375137523753375437553756375737583759376037613762376337643765376637673768376937703771377237733774377537763777377837793780378137823783378437853786378737883789379037913792379337943795379637973798379938003801380238033804380538063807380838093810381138123813381438153816381738183819382038213822382338243825382638273828
  1. // Licensed to the .NET Foundation under one or more agreements.
  2. // The .NET Foundation licenses this file to you under the Apache 2.0 License.
  3. // See the LICENSE file in the project root for more information.
  4. using System;
  5. using System.Collections.Generic;
  6. using System.Linq;
  7. using System.Reactive;
  8. using System.Reactive.Concurrency;
  9. using System.Reactive.Linq;
  10. using System.Text;
  11. using Microsoft.Reactive.Testing;
  12. using ReactiveTests.Dummies;
  13. using Xunit;
  14. namespace ReactiveTests.Tests
  15. {
  16. public class GroupByUntilTest : ReactiveTest
  17. {
  18. #region + GroupByUntil +
  19. [Fact]
  20. public void GroupByUntil_ArgumentChecking()
  21. {
  22. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(default, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, EqualityComparer<int>.Default));
  23. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, default, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, EqualityComparer<int>.Default));
  24. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, default, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, EqualityComparer<int>.Default));
  25. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, default(Func<IGroupedObservable<int, int>, IObservable<int>>), EqualityComparer<int>.Default));
  26. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, default(IEqualityComparer<int>)));
  27. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(default, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance));
  28. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, default, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance));
  29. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, default, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance));
  30. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, default(Func<IGroupedObservable<int, int>, IObservable<int>>)));
  31. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(default, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, EqualityComparer<int>.Default));
  32. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, default, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, EqualityComparer<int>.Default));
  33. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, default(Func<IGroupedObservable<int, int>, IObservable<int>>), EqualityComparer<int>.Default));
  34. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, default(IEqualityComparer<int>)));
  35. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(default, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance));
  36. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, default, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance));
  37. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, default(Func<IGroupedObservable<int, int>, IObservable<int>>)));
  38. }
  39. [Fact]
  40. public void GroupByUntil_WithKeyComparer()
  41. {
  42. var scheduler = new TestScheduler();
  43. var keyInvoked = 0;
  44. var xs = scheduler.CreateHotObservable(
  45. OnNext(90, "error"),
  46. OnNext(110, "error"),
  47. OnNext(130, "error"),
  48. OnNext(220, " foo"),
  49. OnNext(240, " FoO "),
  50. OnNext(270, "baR "),
  51. OnNext(310, "foO "),
  52. OnNext(350, " Baz "),
  53. OnNext(360, " qux "),
  54. OnNext(390, " bar"),
  55. OnNext(420, " BAR "),
  56. OnNext(470, "FOO "),
  57. OnNext(480, "baz "),
  58. OnNext(510, " bAZ "),
  59. OnNext(530, " fOo "),
  60. OnCompleted<string>(570),
  61. OnNext(580, "error"),
  62. OnCompleted<string>(600),
  63. OnError<string>(650, new Exception())
  64. );
  65. var comparer = new GroupByComparer(scheduler);
  66. var res = scheduler.Start(() =>
  67. xs.GroupByUntil(
  68. x =>
  69. {
  70. keyInvoked++;
  71. return x.Trim();
  72. },
  73. g => g.Skip(2),
  74. comparer
  75. ).Select(g => g.Key)
  76. );
  77. res.Messages.AssertEqual(
  78. OnNext(220, "foo"),
  79. OnNext(270, "baR"),
  80. OnNext(350, "Baz"),
  81. OnNext(360, "qux"),
  82. OnNext(470, "FOO"),
  83. OnCompleted<string>(570)
  84. );
  85. xs.Subscriptions.AssertEqual(
  86. Subscribe(200, 570)
  87. );
  88. Assert.Equal(12, keyInvoked);
  89. }
  90. [Fact]
  91. public void GroupByUntil_Outer_Complete()
  92. {
  93. var scheduler = new TestScheduler();
  94. var keyInvoked = 0;
  95. var eleInvoked = 0;
  96. var xs = scheduler.CreateHotObservable(
  97. OnNext(90, "error"),
  98. OnNext(110, "error"),
  99. OnNext(130, "error"),
  100. OnNext(220, " foo"),
  101. OnNext(240, " FoO "),
  102. OnNext(270, "baR "),
  103. OnNext(310, "foO "),
  104. OnNext(350, " Baz "),
  105. OnNext(360, " qux "),
  106. OnNext(390, " bar"),
  107. OnNext(420, " BAR "),
  108. OnNext(470, "FOO "),
  109. OnNext(480, "baz "),
  110. OnNext(510, " bAZ "),
  111. OnNext(530, " fOo "),
  112. OnCompleted<string>(570),
  113. OnNext(580, "error"),
  114. OnCompleted<string>(600),
  115. OnError<string>(650, new Exception())
  116. );
  117. var comparer = new GroupByComparer(scheduler);
  118. var res = scheduler.Start(() =>
  119. xs.GroupByUntil(
  120. x =>
  121. {
  122. keyInvoked++;
  123. return x.Trim();
  124. },
  125. x =>
  126. {
  127. eleInvoked++;
  128. return Reverse(x);
  129. },
  130. g => g.Skip(2),
  131. comparer
  132. ).Select(g => g.Key)
  133. );
  134. res.Messages.AssertEqual(
  135. OnNext(220, "foo"),
  136. OnNext(270, "baR"),
  137. OnNext(350, "Baz"),
  138. OnNext(360, "qux"),
  139. OnNext(470, "FOO"),
  140. OnCompleted<string>(570)
  141. );
  142. xs.Subscriptions.AssertEqual(
  143. Subscribe(200, 570)
  144. );
  145. Assert.Equal(12, keyInvoked);
  146. Assert.Equal(12, eleInvoked);
  147. }
  148. [Fact]
  149. public void GroupByUntil_Outer_Error()
  150. {
  151. var scheduler = new TestScheduler();
  152. var keyInvoked = 0;
  153. var eleInvoked = 0;
  154. var ex = new Exception();
  155. var xs = scheduler.CreateHotObservable(
  156. OnNext(90, "error"),
  157. OnNext(110, "error"),
  158. OnNext(130, "error"),
  159. OnNext(220, " foo"),
  160. OnNext(240, " FoO "),
  161. OnNext(270, "baR "),
  162. OnNext(310, "foO "),
  163. OnNext(350, " Baz "),
  164. OnNext(360, " qux "),
  165. OnNext(390, " bar"),
  166. OnNext(420, " BAR "),
  167. OnNext(470, "FOO "),
  168. OnNext(480, "baz "),
  169. OnNext(510, " bAZ "),
  170. OnNext(530, " fOo "),
  171. OnError<string>(570, ex),
  172. OnNext(580, "error"),
  173. OnCompleted<string>(600),
  174. OnError<string>(650, new Exception())
  175. );
  176. var comparer = new GroupByComparer(scheduler);
  177. var res = scheduler.Start(() =>
  178. xs.GroupByUntil(
  179. x =>
  180. {
  181. keyInvoked++;
  182. return x.Trim();
  183. },
  184. x =>
  185. {
  186. eleInvoked++;
  187. return Reverse(x);
  188. },
  189. g => g.Skip(2),
  190. comparer
  191. ).Select(g => g.Key)
  192. );
  193. res.Messages.AssertEqual(
  194. OnNext(220, "foo"),
  195. OnNext(270, "baR"),
  196. OnNext(350, "Baz"),
  197. OnNext(360, "qux"),
  198. OnNext(470, "FOO"),
  199. OnError<string>(570, ex)
  200. );
  201. xs.Subscriptions.AssertEqual(
  202. Subscribe(200, 570)
  203. );
  204. Assert.Equal(12, keyInvoked);
  205. Assert.Equal(12, eleInvoked);
  206. }
  207. [Fact]
  208. public void GroupByUntil_Outer_Dispose()
  209. {
  210. var scheduler = new TestScheduler();
  211. var keyInvoked = 0;
  212. var eleInvoked = 0;
  213. var xs = scheduler.CreateHotObservable(
  214. OnNext(90, "error"),
  215. OnNext(110, "error"),
  216. OnNext(130, "error"),
  217. OnNext(220, " foo"),
  218. OnNext(240, " FoO "),
  219. OnNext(270, "baR "),
  220. OnNext(310, "foO "),
  221. OnNext(350, " Baz "),
  222. OnNext(360, " qux "),
  223. OnNext(390, " bar"),
  224. OnNext(420, " BAR "),
  225. OnNext(470, "FOO "),
  226. OnNext(480, "baz "),
  227. OnNext(510, " bAZ "),
  228. OnNext(530, " fOo "),
  229. OnCompleted<string>(570),
  230. OnNext(580, "error"),
  231. OnCompleted<string>(600),
  232. OnError<string>(650, new Exception())
  233. );
  234. var comparer = new GroupByComparer(scheduler);
  235. var res = scheduler.Start(() =>
  236. xs.GroupByUntil(
  237. x =>
  238. {
  239. keyInvoked++;
  240. return x.Trim();
  241. },
  242. x =>
  243. {
  244. eleInvoked++;
  245. return Reverse(x);
  246. },
  247. g => g.Skip(2),
  248. comparer
  249. ).Select(g => g.Key),
  250. 355
  251. );
  252. res.Messages.AssertEqual(
  253. OnNext(220, "foo"),
  254. OnNext(270, "baR"),
  255. OnNext(350, "Baz")
  256. );
  257. xs.Subscriptions.AssertEqual(
  258. Subscribe(200, 355)
  259. );
  260. Assert.Equal(5, keyInvoked);
  261. Assert.Equal(5, eleInvoked);
  262. }
  263. [Fact]
  264. public void GroupByUntil_Outer_KeyThrow()
  265. {
  266. var scheduler = new TestScheduler();
  267. var keyInvoked = 0;
  268. var eleInvoked = 0;
  269. var ex = new Exception();
  270. var xs = scheduler.CreateHotObservable(
  271. OnNext(90, "error"),
  272. OnNext(110, "error"),
  273. OnNext(130, "error"),
  274. OnNext(220, " foo"),
  275. OnNext(240, " FoO "),
  276. OnNext(270, "baR "),
  277. OnNext(310, "foO "),
  278. OnNext(350, " Baz "),
  279. OnNext(360, " qux "),
  280. OnNext(390, " bar"),
  281. OnNext(420, " BAR "),
  282. OnNext(470, "FOO "),
  283. OnNext(480, "baz "),
  284. OnNext(510, " bAZ "),
  285. OnNext(530, " fOo "),
  286. OnCompleted<string>(570),
  287. OnNext(580, "error"),
  288. OnCompleted<string>(600),
  289. OnError<string>(650, new Exception())
  290. );
  291. var comparer = new GroupByComparer(scheduler);
  292. var res = scheduler.Start(() =>
  293. xs.GroupByUntil(
  294. x =>
  295. {
  296. keyInvoked++;
  297. if (keyInvoked == 10)
  298. {
  299. throw ex;
  300. }
  301. return x.Trim();
  302. },
  303. x =>
  304. {
  305. eleInvoked++;
  306. return Reverse(x);
  307. },
  308. g => g.Skip(2),
  309. comparer
  310. ).Select(g => g.Key)
  311. );
  312. res.Messages.AssertEqual(
  313. OnNext(220, "foo"),
  314. OnNext(270, "baR"),
  315. OnNext(350, "Baz"),
  316. OnNext(360, "qux"),
  317. OnNext(470, "FOO"),
  318. OnError<string>(480, ex)
  319. );
  320. xs.Subscriptions.AssertEqual(
  321. Subscribe(200, 480)
  322. );
  323. Assert.Equal(10, keyInvoked);
  324. Assert.Equal(9, eleInvoked);
  325. }
  326. [Fact]
  327. public void GroupByUntil_Outer_EleThrow()
  328. {
  329. var scheduler = new TestScheduler();
  330. var keyInvoked = 0;
  331. var eleInvoked = 0;
  332. var ex = new Exception();
  333. var xs = scheduler.CreateHotObservable(
  334. OnNext(90, "error"),
  335. OnNext(110, "error"),
  336. OnNext(130, "error"),
  337. OnNext(220, " foo"),
  338. OnNext(240, " FoO "),
  339. OnNext(270, "baR "),
  340. OnNext(310, "foO "),
  341. OnNext(350, " Baz "),
  342. OnNext(360, " qux "),
  343. OnNext(390, " bar"),
  344. OnNext(420, " BAR "),
  345. OnNext(470, "FOO "),
  346. OnNext(480, "baz "),
  347. OnNext(510, " bAZ "),
  348. OnNext(530, " fOo "),
  349. OnCompleted<string>(570),
  350. OnNext(580, "error"),
  351. OnCompleted<string>(600),
  352. OnError<string>(650, new Exception())
  353. );
  354. var comparer = new GroupByComparer(scheduler);
  355. var res = scheduler.Start(() =>
  356. xs.GroupByUntil(
  357. x =>
  358. {
  359. keyInvoked++;
  360. return x.Trim();
  361. },
  362. x =>
  363. {
  364. eleInvoked++;
  365. if (eleInvoked == 10)
  366. {
  367. throw ex;
  368. }
  369. return Reverse(x);
  370. },
  371. g => g.Skip(2),
  372. comparer
  373. ).Select(g => g.Key)
  374. );
  375. res.Messages.AssertEqual(
  376. OnNext(220, "foo"),
  377. OnNext(270, "baR"),
  378. OnNext(350, "Baz"),
  379. OnNext(360, "qux"),
  380. OnNext(470, "FOO"),
  381. OnError<string>(480, ex)
  382. );
  383. xs.Subscriptions.AssertEqual(
  384. Subscribe(200, 480)
  385. );
  386. Assert.Equal(10, keyInvoked);
  387. Assert.Equal(10, eleInvoked);
  388. }
  389. [Fact]
  390. public void GroupByUntil_Outer_ComparerEqualsThrow()
  391. {
  392. var scheduler = new TestScheduler();
  393. var keyInvoked = 0;
  394. var eleInvoked = 0;
  395. var xs = scheduler.CreateHotObservable(
  396. OnNext(90, "error"),
  397. OnNext(110, "error"),
  398. OnNext(130, "error"),
  399. OnNext(220, " foo"),
  400. OnNext(240, " FoO "),
  401. OnNext(270, "baR "),
  402. OnNext(310, "foO "),
  403. OnNext(350, " Baz "),
  404. OnNext(360, " qux "),
  405. OnNext(390, " bar"),
  406. OnNext(420, " BAR "),
  407. OnNext(470, "FOO "),
  408. OnNext(480, "baz "),
  409. OnNext(510, " bAZ "),
  410. OnNext(530, " fOo "),
  411. OnCompleted<string>(570),
  412. OnNext(580, "error"),
  413. OnCompleted<string>(600),
  414. OnError<string>(650, new Exception())
  415. );
  416. var comparer = new GroupByComparer(scheduler, 250, ushort.MaxValue);
  417. var res = scheduler.Start(() =>
  418. xs.GroupByUntil(
  419. x =>
  420. {
  421. keyInvoked++;
  422. return x.Trim();
  423. },
  424. x =>
  425. {
  426. eleInvoked++;
  427. return Reverse(x);
  428. },
  429. g => g.Skip(2),
  430. comparer
  431. ).Select(g => g.Key)
  432. );
  433. res.Messages.AssertEqual(
  434. OnNext(220, "foo"),
  435. OnNext(270, "baR"),
  436. OnError<string>(310, comparer.EqualsException)
  437. );
  438. xs.Subscriptions.AssertEqual(
  439. Subscribe(200, 310)
  440. );
  441. Assert.Equal(4, keyInvoked);
  442. Assert.Equal(3, eleInvoked);
  443. }
  444. [Fact]
  445. public void GroupByUntil_Outer_ComparerGetHashCodeThrow()
  446. {
  447. var scheduler = new TestScheduler();
  448. var keyInvoked = 0;
  449. var eleInvoked = 0;
  450. var xs = scheduler.CreateHotObservable(
  451. OnNext(90, "error"),
  452. OnNext(110, "error"),
  453. OnNext(130, "error"),
  454. OnNext(220, " foo"),
  455. OnNext(240, " FoO "),
  456. OnNext(270, "baR "),
  457. OnNext(310, "foO "),
  458. OnNext(350, " Baz "),
  459. OnNext(360, " qux "),
  460. OnNext(390, " bar"),
  461. OnNext(420, " BAR "),
  462. OnNext(470, "FOO "),
  463. OnNext(480, "baz "),
  464. OnNext(510, " bAZ "),
  465. OnNext(530, " fOo "),
  466. OnCompleted<string>(570),
  467. OnNext(580, "error"),
  468. OnCompleted<string>(600),
  469. OnError<string>(650, new Exception())
  470. );
  471. var comparer = new GroupByComparer(scheduler, ushort.MaxValue, 410);
  472. var res = scheduler.Start(() =>
  473. xs.GroupByUntil(
  474. x =>
  475. {
  476. keyInvoked++;
  477. return x.Trim();
  478. },
  479. x =>
  480. {
  481. eleInvoked++;
  482. return Reverse(x);
  483. },
  484. g => g.Skip(2),
  485. comparer
  486. ).Select(g => g.Key)
  487. );
  488. res.Messages.AssertEqual(
  489. OnNext(220, "foo"),
  490. OnNext(270, "baR"),
  491. OnNext(350, "Baz"),
  492. OnNext(360, "qux"),
  493. OnError<string>(420, comparer.HashCodeException)
  494. );
  495. xs.Subscriptions.AssertEqual(
  496. Subscribe(200, 420)
  497. );
  498. Assert.Equal(8, keyInvoked);
  499. Assert.Equal(7, eleInvoked);
  500. }
  501. [Fact]
  502. public void GroupByUntil_Inner_Complete()
  503. {
  504. var scheduler = new TestScheduler();
  505. var xs = scheduler.CreateHotObservable(
  506. OnNext(90, "error"),
  507. OnNext(110, "error"),
  508. OnNext(130, "error"),
  509. OnNext(220, " foo"),
  510. OnNext(240, " FoO "),
  511. OnNext(270, "baR "),
  512. OnNext(310, "foO "),
  513. OnNext(350, " Baz "),
  514. OnNext(360, " qux "),
  515. OnNext(390, " bar"),
  516. OnNext(420, " BAR "),
  517. OnNext(470, "FOO "),
  518. OnNext(480, "baz "),
  519. OnNext(510, " bAZ "),
  520. OnNext(530, " fOo "),
  521. OnCompleted<string>(570),
  522. OnNext(580, "error"),
  523. OnCompleted<string>(600),
  524. OnError<string>(650, new Exception())
  525. );
  526. var comparer = new GroupByComparer(scheduler);
  527. var outer = default(IObservable<IGroupedObservable<string, string>>);
  528. var outerSubscription = default(IDisposable);
  529. var inners = new Dictionary<string, IObservable<string>>();
  530. var innerSubscriptions = new Dictionary<string, IDisposable>();
  531. var res = new Dictionary<string, ITestableObserver<string>>();
  532. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), comparer));
  533. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  534. {
  535. var result = scheduler.CreateObserver<string>();
  536. inners[group.Key] = group;
  537. res[group.Key] = result;
  538. scheduler.ScheduleRelative(100, () => innerSubscriptions[group.Key] = group.Subscribe(result));
  539. }));
  540. scheduler.ScheduleAbsolute(Disposed, () =>
  541. {
  542. outerSubscription.Dispose();
  543. foreach (var d in innerSubscriptions.Values)
  544. {
  545. d.Dispose();
  546. }
  547. });
  548. scheduler.Start();
  549. Assert.Equal(5, inners.Count);
  550. res["foo"].Messages.AssertEqual(
  551. OnCompleted<string>(320)
  552. );
  553. res["baR"].Messages.AssertEqual(
  554. OnNext(390, "rab "),
  555. OnNext(420, " RAB "), // Breaking change > v2.2 - prior to resolving a deadlock, the group would get closed prior to letting this message through
  556. OnCompleted<string>(420)
  557. );
  558. res["Baz"].Messages.AssertEqual(
  559. OnNext(480, " zab"),
  560. OnNext(510, " ZAb "), // Breaking change > v2.2 - prior to resolving a deadlock, the group would get closed prior to letting this message through
  561. OnCompleted<string>(510)
  562. );
  563. res["qux"].Messages.AssertEqual(
  564. OnCompleted<string>(570)
  565. );
  566. res["FOO"].Messages.AssertEqual(
  567. OnCompleted<string>(570)
  568. );
  569. xs.Subscriptions.AssertEqual(
  570. Subscribe(200, 570)
  571. );
  572. }
  573. [Fact]
  574. public void GroupByUntil_Inner_Complete_All()
  575. {
  576. var scheduler = new TestScheduler();
  577. var xs = scheduler.CreateHotObservable(
  578. OnNext(90, "error"),
  579. OnNext(110, "error"),
  580. OnNext(130, "error"),
  581. OnNext(220, " foo"),
  582. OnNext(240, " FoO "),
  583. OnNext(270, "baR "),
  584. OnNext(310, "foO "),
  585. OnNext(350, " Baz "),
  586. OnNext(360, " qux "),
  587. OnNext(390, " bar"),
  588. OnNext(420, " BAR "),
  589. OnNext(470, "FOO "),
  590. OnNext(480, "baz "),
  591. OnNext(510, " bAZ "),
  592. OnNext(530, " fOo "),
  593. OnCompleted<string>(570),
  594. OnNext(580, "error"),
  595. OnCompleted<string>(600),
  596. OnError<string>(650, new Exception())
  597. );
  598. var comparer = new GroupByComparer(scheduler);
  599. var outer = default(IObservable<IGroupedObservable<string, string>>);
  600. var outerSubscription = default(IDisposable);
  601. var inners = new Dictionary<string, IObservable<string>>();
  602. var innerSubscriptions = new Dictionary<string, IDisposable>();
  603. var res = new Dictionary<string, ITestableObserver<string>>();
  604. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), comparer));
  605. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  606. {
  607. var result = scheduler.CreateObserver<string>();
  608. inners[group.Key] = group;
  609. res[group.Key] = result;
  610. innerSubscriptions[group.Key] = group.Subscribe(result);
  611. }));
  612. scheduler.ScheduleAbsolute(Disposed, () =>
  613. {
  614. outerSubscription.Dispose();
  615. foreach (var d in innerSubscriptions.Values)
  616. {
  617. d.Dispose();
  618. }
  619. });
  620. scheduler.Start();
  621. Assert.Equal(5, inners.Count);
  622. res["foo"].Messages.AssertEqual(
  623. OnNext(220, "oof "),
  624. OnNext(240, " OoF "),
  625. OnNext(310, " Oof"),
  626. OnCompleted<string>(310)
  627. );
  628. res["baR"].Messages.AssertEqual(
  629. OnNext(270, " Rab"),
  630. OnNext(390, "rab "),
  631. OnNext(420, " RAB "),
  632. OnCompleted<string>(420)
  633. );
  634. res["Baz"].Messages.AssertEqual(
  635. OnNext(350, " zaB "),
  636. OnNext(480, " zab"),
  637. OnNext(510, " ZAb "),
  638. OnCompleted<string>(510)
  639. );
  640. res["qux"].Messages.AssertEqual(
  641. OnNext(360, " xuq "),
  642. OnCompleted<string>(570)
  643. );
  644. res["FOO"].Messages.AssertEqual(
  645. OnNext(470, " OOF"),
  646. OnNext(530, " oOf "),
  647. OnCompleted<string>(570)
  648. );
  649. xs.Subscriptions.AssertEqual(
  650. Subscribe(200, 570)
  651. );
  652. }
  653. [Fact]
  654. public void GroupByUntil_Inner_Error()
  655. {
  656. var scheduler = new TestScheduler();
  657. var ex1 = new Exception();
  658. var xs = scheduler.CreateHotObservable(
  659. OnNext(90, "error"),
  660. OnNext(110, "error"),
  661. OnNext(130, "error"),
  662. OnNext(220, " foo"),
  663. OnNext(240, " FoO "),
  664. OnNext(270, "baR "),
  665. OnNext(310, "foO "),
  666. OnNext(350, " Baz "),
  667. OnNext(360, " qux "),
  668. OnNext(390, " bar"),
  669. OnNext(420, " BAR "),
  670. OnNext(470, "FOO "),
  671. OnNext(480, "baz "),
  672. OnNext(510, " bAZ "),
  673. OnNext(530, " fOo "),
  674. OnError<string>(570, ex1),
  675. OnNext(580, "error"),
  676. OnCompleted<string>(600),
  677. OnError<string>(650, new Exception())
  678. );
  679. var comparer = new GroupByComparer(scheduler);
  680. var outer = default(IObservable<IGroupedObservable<string, string>>);
  681. var outerSubscription = default(IDisposable);
  682. var inners = new Dictionary<string, IObservable<string>>();
  683. var innerSubscriptions = new Dictionary<string, IDisposable>();
  684. var res = new Dictionary<string, ITestableObserver<string>>();
  685. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), comparer));
  686. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  687. {
  688. var result = scheduler.CreateObserver<string>();
  689. inners[group.Key] = group;
  690. res[group.Key] = result;
  691. scheduler.ScheduleRelative(100, () => innerSubscriptions[group.Key] = group.Subscribe(result));
  692. }, ex => { }));
  693. scheduler.ScheduleAbsolute(Disposed, () =>
  694. {
  695. outerSubscription.Dispose();
  696. foreach (var d in innerSubscriptions.Values)
  697. {
  698. d.Dispose();
  699. }
  700. });
  701. scheduler.Start();
  702. Assert.Equal(5, inners.Count);
  703. res["foo"].Messages.AssertEqual(
  704. OnCompleted<string>(320)
  705. );
  706. res["baR"].Messages.AssertEqual(
  707. OnNext(390, "rab "),
  708. OnNext(420, " RAB "), // Breaking change > v2.2 - prior to resolving a deadlock, the group would get closed prior to letting this message through
  709. OnCompleted<string>(420)
  710. );
  711. res["Baz"].Messages.AssertEqual(
  712. OnNext(480, " zab"),
  713. OnNext(510, " ZAb "), // Breaking change > v2.2 - prior to resolving a deadlock, the group would get closed prior to letting this message through
  714. OnCompleted<string>(510)
  715. );
  716. res["qux"].Messages.AssertEqual(
  717. OnError<string>(570, ex1)
  718. );
  719. res["FOO"].Messages.AssertEqual(
  720. OnError<string>(570, ex1)
  721. );
  722. xs.Subscriptions.AssertEqual(
  723. Subscribe(200, 570)
  724. );
  725. }
  726. [Fact]
  727. public void GroupByUntil_Inner_Dispose()
  728. {
  729. var scheduler = new TestScheduler();
  730. var xs = scheduler.CreateHotObservable(
  731. OnNext(90, "error"),
  732. OnNext(110, "error"),
  733. OnNext(130, "error"),
  734. OnNext(220, " foo"),
  735. OnNext(240, " FoO "),
  736. OnNext(270, "baR "),
  737. OnNext(310, "foO "),
  738. OnNext(350, " Baz "),
  739. OnNext(360, " qux "),
  740. OnNext(390, " bar"),
  741. OnNext(420, " BAR "),
  742. OnNext(470, "FOO "),
  743. OnNext(480, "baz "),
  744. OnNext(510, " bAZ "),
  745. OnNext(530, " fOo "),
  746. OnCompleted<string>(570),
  747. OnNext(580, "error"),
  748. OnCompleted<string>(600),
  749. OnError<string>(650, new Exception())
  750. );
  751. var comparer = new GroupByComparer(scheduler);
  752. var outer = default(IObservable<IGroupedObservable<string, string>>);
  753. var outerSubscription = default(IDisposable);
  754. var inners = new Dictionary<string, IObservable<string>>();
  755. var innerSubscriptions = new Dictionary<string, IDisposable>();
  756. var res = new Dictionary<string, ITestableObserver<string>>();
  757. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), comparer));
  758. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  759. {
  760. var result = scheduler.CreateObserver<string>();
  761. inners[group.Key] = group;
  762. res[group.Key] = result;
  763. innerSubscriptions[group.Key] = group.Subscribe(result);
  764. }));
  765. scheduler.ScheduleAbsolute(400, () =>
  766. {
  767. outerSubscription.Dispose();
  768. foreach (var d in innerSubscriptions.Values)
  769. {
  770. d.Dispose();
  771. }
  772. });
  773. scheduler.Start();
  774. Assert.Equal(4, inners.Count);
  775. res["foo"].Messages.AssertEqual(
  776. OnNext(220, "oof "),
  777. OnNext(240, " OoF "),
  778. OnNext(310, " Oof"),
  779. OnCompleted<string>(310)
  780. );
  781. res["baR"].Messages.AssertEqual(
  782. OnNext(270, " Rab"),
  783. OnNext(390, "rab ")
  784. );
  785. res["Baz"].Messages.AssertEqual(
  786. OnNext(350, " zaB ")
  787. );
  788. res["qux"].Messages.AssertEqual(
  789. OnNext(360, " xuq ")
  790. );
  791. xs.Subscriptions.AssertEqual(
  792. Subscribe(200, 400)
  793. );
  794. }
  795. [Fact]
  796. public void GroupByUntil_Inner_KeyThrow()
  797. {
  798. var scheduler = new TestScheduler();
  799. var xs = scheduler.CreateHotObservable(
  800. OnNext(90, "error"),
  801. OnNext(110, "error"),
  802. OnNext(130, "error"),
  803. OnNext(220, " foo"),
  804. OnNext(240, " FoO "),
  805. OnNext(270, "baR "),
  806. OnNext(310, "foO "),
  807. OnNext(350, " Baz "),
  808. OnNext(360, " qux "),
  809. OnNext(390, " bar"),
  810. OnNext(420, " BAR "),
  811. OnNext(470, "FOO "),
  812. OnNext(480, "baz "),
  813. OnNext(510, " bAZ "),
  814. OnNext(530, " fOo "),
  815. OnCompleted<string>(570),
  816. OnNext(580, "error"),
  817. OnCompleted<string>(600),
  818. OnError<string>(650, new Exception())
  819. );
  820. var comparer = new GroupByComparer(scheduler);
  821. var outer = default(IObservable<IGroupedObservable<string, string>>);
  822. var outerSubscription = default(IDisposable);
  823. var inners = new Dictionary<string, IObservable<string>>();
  824. var innerSubscriptions = new Dictionary<string, IDisposable>();
  825. var res = new Dictionary<string, ITestableObserver<string>>();
  826. var keyInvoked = 0;
  827. var ex = new Exception();
  828. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x =>
  829. {
  830. keyInvoked++;
  831. if (keyInvoked == 6)
  832. {
  833. throw ex;
  834. }
  835. return x.Trim();
  836. }, x => Reverse(x), g => g.Skip(2), comparer));
  837. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  838. {
  839. var result = scheduler.CreateObserver<string>();
  840. inners[group.Key] = group;
  841. res[group.Key] = result;
  842. innerSubscriptions[group.Key] = group.Subscribe(result);
  843. }, _ => { }));
  844. scheduler.ScheduleAbsolute(Disposed, () =>
  845. {
  846. outerSubscription.Dispose();
  847. foreach (var d in innerSubscriptions.Values)
  848. {
  849. d.Dispose();
  850. }
  851. });
  852. scheduler.Start();
  853. Assert.Equal(3, inners.Count);
  854. res["foo"].Messages.AssertEqual(
  855. OnNext(220, "oof "),
  856. OnNext(240, " OoF "),
  857. OnNext(310, " Oof"),
  858. OnCompleted<string>(310)
  859. );
  860. res["baR"].Messages.AssertEqual(
  861. OnNext(270, " Rab"),
  862. OnError<string>(360, ex)
  863. );
  864. res["Baz"].Messages.AssertEqual(
  865. OnNext(350, " zaB "),
  866. OnError<string>(360, ex)
  867. );
  868. xs.Subscriptions.AssertEqual(
  869. Subscribe(200, 360)
  870. );
  871. }
  872. [Fact]
  873. public void GroupByUntil_Inner_EleThrow()
  874. {
  875. var scheduler = new TestScheduler();
  876. var xs = scheduler.CreateHotObservable(
  877. OnNext(90, "error"),
  878. OnNext(110, "error"),
  879. OnNext(130, "error"),
  880. OnNext(220, " foo"),
  881. OnNext(240, " FoO "),
  882. OnNext(270, "baR "),
  883. OnNext(310, "foO "),
  884. OnNext(350, " Baz "),
  885. OnNext(360, " qux "),
  886. OnNext(390, " bar"),
  887. OnNext(420, " BAR "),
  888. OnNext(470, "FOO "),
  889. OnNext(480, "baz "),
  890. OnNext(510, " bAZ "),
  891. OnNext(530, " fOo "),
  892. OnCompleted<string>(570),
  893. OnNext(580, "error"),
  894. OnCompleted<string>(600),
  895. OnError<string>(650, new Exception())
  896. );
  897. var comparer = new GroupByComparer(scheduler);
  898. var outer = default(IObservable<IGroupedObservable<string, string>>);
  899. var outerSubscription = default(IDisposable);
  900. var inners = new Dictionary<string, IObservable<string>>();
  901. var innerSubscriptions = new Dictionary<string, IDisposable>();
  902. var res = new Dictionary<string, ITestableObserver<string>>();
  903. var eleInvoked = 0;
  904. var ex = new Exception();
  905. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x =>
  906. {
  907. eleInvoked++;
  908. if (eleInvoked == 6)
  909. {
  910. throw ex;
  911. }
  912. return Reverse(x);
  913. }, g => g.Skip(2), comparer));
  914. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  915. {
  916. var result = scheduler.CreateObserver<string>();
  917. inners[group.Key] = group;
  918. res[group.Key] = result;
  919. innerSubscriptions[group.Key] = group.Subscribe(result);
  920. }, _ => { }));
  921. scheduler.ScheduleAbsolute(Disposed, () =>
  922. {
  923. outerSubscription.Dispose();
  924. foreach (var d in innerSubscriptions.Values)
  925. {
  926. d.Dispose();
  927. }
  928. });
  929. scheduler.Start();
  930. Assert.Equal(4, inners.Count);
  931. res["foo"].Messages.AssertEqual(
  932. OnNext(220, "oof "),
  933. OnNext(240, " OoF "),
  934. OnNext(310, " Oof"),
  935. OnCompleted<string>(310)
  936. );
  937. res["baR"].Messages.AssertEqual(
  938. OnNext(270, " Rab"),
  939. OnError<string>(360, ex)
  940. );
  941. res["Baz"].Messages.AssertEqual(
  942. OnNext(350, " zaB "),
  943. OnError<string>(360, ex)
  944. );
  945. res["qux"].Messages.AssertEqual(
  946. OnError<string>(360, ex)
  947. );
  948. xs.Subscriptions.AssertEqual(
  949. Subscribe(200, 360)
  950. );
  951. }
  952. [Fact]
  953. public void GroupByUntil_Inner_Comparer_EqualsThrow()
  954. {
  955. var scheduler = new TestScheduler();
  956. var xs = scheduler.CreateHotObservable(
  957. OnNext(90, "error"),
  958. OnNext(110, "error"),
  959. OnNext(130, "error"),
  960. OnNext(220, " foo"),
  961. OnNext(240, " FoO "),
  962. OnNext(270, "baR "),
  963. OnNext(310, "foO "),
  964. OnNext(350, " Baz "),
  965. OnNext(360, " qux "),
  966. OnNext(390, " bar"),
  967. OnNext(420, " BAR "),
  968. OnNext(470, "FOO "),
  969. OnNext(480, "baz "),
  970. OnNext(510, " bAZ "),
  971. OnNext(530, " fOo "),
  972. OnCompleted<string>(570),
  973. OnNext(580, "error"),
  974. OnCompleted<string>(600),
  975. OnError<string>(650, new Exception())
  976. );
  977. var comparer = new GroupByComparer(scheduler, 400, ushort.MaxValue);
  978. var outer = default(IObservable<IGroupedObservable<string, string>>);
  979. var outerSubscription = default(IDisposable);
  980. var inners = new Dictionary<string, IObservable<string>>();
  981. var innerSubscriptions = new Dictionary<string, IDisposable>();
  982. var res = new Dictionary<string, ITestableObserver<string>>();
  983. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), comparer));
  984. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  985. {
  986. var result = scheduler.CreateObserver<string>();
  987. inners[group.Key] = group;
  988. res[group.Key] = result;
  989. innerSubscriptions[group.Key] = group.Subscribe(result);
  990. }, _ => { }));
  991. scheduler.ScheduleAbsolute(Disposed, () =>
  992. {
  993. outerSubscription.Dispose();
  994. foreach (var d in innerSubscriptions.Values)
  995. {
  996. d.Dispose();
  997. }
  998. });
  999. scheduler.Start();
  1000. Assert.Equal(4, inners.Count);
  1001. res["foo"].Messages.AssertEqual(
  1002. OnNext(220, "oof "),
  1003. OnNext(240, " OoF "),
  1004. OnNext(310, " Oof"),
  1005. OnCompleted<string>(310)
  1006. );
  1007. res["baR"].Messages.AssertEqual(
  1008. OnNext(270, " Rab"),
  1009. OnNext(390, "rab "),
  1010. OnError<string>(420, comparer.EqualsException)
  1011. );
  1012. res["Baz"].Messages.AssertEqual(
  1013. OnNext(350, " zaB "),
  1014. OnError<string>(420, comparer.EqualsException)
  1015. );
  1016. res["qux"].Messages.AssertEqual(
  1017. OnNext(360, " xuq "),
  1018. OnError<string>(420, comparer.EqualsException)
  1019. );
  1020. xs.Subscriptions.AssertEqual(
  1021. Subscribe(200, 420)
  1022. );
  1023. }
  1024. [Fact]
  1025. public void GroupByUntil_Inner_Comparer_GetHashCodeThrow()
  1026. {
  1027. var scheduler = new TestScheduler();
  1028. var xs = scheduler.CreateHotObservable(
  1029. OnNext(90, "error"),
  1030. OnNext(110, "error"),
  1031. OnNext(130, "error"),
  1032. OnNext(220, " foo"),
  1033. OnNext(240, " FoO "),
  1034. OnNext(270, "baR "),
  1035. OnNext(310, "foO "),
  1036. OnNext(350, " Baz "),
  1037. OnNext(360, " qux "),
  1038. OnNext(390, " bar"),
  1039. OnNext(420, " BAR "),
  1040. OnNext(470, "FOO "),
  1041. OnNext(480, "baz "),
  1042. OnNext(510, " bAZ "),
  1043. OnNext(530, " fOo "),
  1044. OnCompleted<string>(570),
  1045. OnNext(580, "error"),
  1046. OnCompleted<string>(600),
  1047. OnError<string>(650, new Exception())
  1048. );
  1049. var comparer = new GroupByComparer(scheduler, ushort.MaxValue, 400);
  1050. var outer = default(IObservable<IGroupedObservable<string, string>>);
  1051. var outerSubscription = default(IDisposable);
  1052. var inners = new Dictionary<string, IObservable<string>>();
  1053. var innerSubscriptions = new Dictionary<string, IDisposable>();
  1054. var res = new Dictionary<string, ITestableObserver<string>>();
  1055. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), comparer));
  1056. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  1057. {
  1058. var result = scheduler.CreateObserver<string>();
  1059. inners[group.Key] = group;
  1060. res[group.Key] = result;
  1061. innerSubscriptions[group.Key] = group.Subscribe(result);
  1062. }, _ => { }));
  1063. scheduler.ScheduleAbsolute(Disposed, () =>
  1064. {
  1065. outerSubscription.Dispose();
  1066. foreach (var d in innerSubscriptions.Values)
  1067. {
  1068. d.Dispose();
  1069. }
  1070. });
  1071. scheduler.Start();
  1072. Assert.Equal(4, inners.Count);
  1073. res["foo"].Messages.AssertEqual(
  1074. OnNext(220, "oof "),
  1075. OnNext(240, " OoF "),
  1076. OnNext(310, " Oof"),
  1077. OnCompleted<string>(310)
  1078. );
  1079. res["baR"].Messages.AssertEqual(
  1080. OnNext(270, " Rab"),
  1081. OnNext(390, "rab "),
  1082. OnError<string>(420, comparer.HashCodeException)
  1083. );
  1084. res["Baz"].Messages.AssertEqual(
  1085. OnNext(350, " zaB "),
  1086. OnError<string>(420, comparer.HashCodeException)
  1087. );
  1088. res["qux"].Messages.AssertEqual(
  1089. OnNext(360, " xuq "),
  1090. OnError<string>(420, comparer.HashCodeException)
  1091. );
  1092. xs.Subscriptions.AssertEqual(
  1093. Subscribe(200, 420)
  1094. );
  1095. }
  1096. [Fact]
  1097. public void GroupByUntil_Outer_Independence()
  1098. {
  1099. var scheduler = new TestScheduler();
  1100. var xs = scheduler.CreateHotObservable(
  1101. OnNext(90, "error"),
  1102. OnNext(110, "error"),
  1103. OnNext(130, "error"),
  1104. OnNext(220, " foo"),
  1105. OnNext(240, " FoO "),
  1106. OnNext(270, "baR "),
  1107. OnNext(310, "foO "),
  1108. OnNext(350, " Baz "),
  1109. OnNext(360, " qux "),
  1110. OnNext(390, " bar"),
  1111. OnNext(420, " BAR "),
  1112. OnNext(470, "FOO "),
  1113. OnNext(480, "baz "),
  1114. OnNext(510, " bAZ "),
  1115. OnNext(530, " fOo "),
  1116. OnCompleted<string>(570),
  1117. OnNext(580, "error"),
  1118. OnCompleted<string>(600),
  1119. OnError<string>(650, new Exception())
  1120. );
  1121. var comparer = new GroupByComparer(scheduler);
  1122. var outer = default(IObservable<IGroupedObservable<string, string>>);
  1123. var outerSubscription = default(IDisposable);
  1124. var inners = new Dictionary<string, IObservable<string>>();
  1125. var innerSubscriptions = new Dictionary<string, IDisposable>();
  1126. var res = new Dictionary<string, ITestableObserver<string>>();
  1127. var outerResults = scheduler.CreateObserver<string>();
  1128. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), comparer));
  1129. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  1130. {
  1131. outerResults.OnNext(group.Key);
  1132. var result = scheduler.CreateObserver<string>();
  1133. inners[group.Key] = group;
  1134. res[group.Key] = result;
  1135. innerSubscriptions[group.Key] = group.Subscribe(result);
  1136. }, outerResults.OnError, outerResults.OnCompleted));
  1137. scheduler.ScheduleAbsolute(Disposed, () =>
  1138. {
  1139. outerSubscription.Dispose();
  1140. foreach (var d in innerSubscriptions.Values)
  1141. {
  1142. d.Dispose();
  1143. }
  1144. });
  1145. scheduler.ScheduleAbsolute(320, () => outerSubscription.Dispose());
  1146. scheduler.Start();
  1147. Assert.Equal(2, inners.Count);
  1148. outerResults.Messages.AssertEqual(
  1149. OnNext(220, "foo"),
  1150. OnNext(270, "baR")
  1151. );
  1152. res["foo"].Messages.AssertEqual(
  1153. OnNext(220, "oof "),
  1154. OnNext(240, " OoF "),
  1155. OnNext(310, " Oof"),
  1156. OnCompleted<string>(310)
  1157. );
  1158. res["baR"].Messages.AssertEqual(
  1159. OnNext(270, " Rab"),
  1160. OnNext(390, "rab "),
  1161. OnNext(420, " RAB "),
  1162. OnCompleted<string>(420)
  1163. );
  1164. xs.Subscriptions.AssertEqual(
  1165. Subscribe(200, 420)
  1166. );
  1167. }
  1168. [Fact]
  1169. public void GroupByUntil_Inner_Independence()
  1170. {
  1171. var scheduler = new TestScheduler();
  1172. var xs = scheduler.CreateHotObservable(
  1173. OnNext(90, "error"),
  1174. OnNext(110, "error"),
  1175. OnNext(130, "error"),
  1176. OnNext(220, " foo"),
  1177. OnNext(240, " FoO "),
  1178. OnNext(270, "baR "),
  1179. OnNext(310, "foO "),
  1180. OnNext(350, " Baz "),
  1181. OnNext(360, " qux "),
  1182. OnNext(390, " bar"),
  1183. OnNext(420, " BAR "),
  1184. OnNext(470, "FOO "),
  1185. OnNext(480, "baz "),
  1186. OnNext(510, " bAZ "),
  1187. OnNext(530, " fOo "),
  1188. OnCompleted<string>(570),
  1189. OnNext(580, "error"),
  1190. OnCompleted<string>(600),
  1191. OnError<string>(650, new Exception())
  1192. );
  1193. var comparer = new GroupByComparer(scheduler);
  1194. var outer = default(IObservable<IGroupedObservable<string, string>>);
  1195. var outerSubscription = default(IDisposable);
  1196. var inners = new Dictionary<string, IObservable<string>>();
  1197. var innerSubscriptions = new Dictionary<string, IDisposable>();
  1198. var res = new Dictionary<string, ITestableObserver<string>>();
  1199. var outerResults = scheduler.CreateObserver<string>();
  1200. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), comparer));
  1201. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  1202. {
  1203. outerResults.OnNext(group.Key);
  1204. var result = scheduler.CreateObserver<string>();
  1205. inners[group.Key] = group;
  1206. res[group.Key] = result;
  1207. innerSubscriptions[group.Key] = group.Subscribe(result);
  1208. }, outerResults.OnError, outerResults.OnCompleted));
  1209. scheduler.ScheduleAbsolute(Disposed, () =>
  1210. {
  1211. outerSubscription.Dispose();
  1212. foreach (var d in innerSubscriptions.Values)
  1213. {
  1214. d.Dispose();
  1215. }
  1216. });
  1217. scheduler.ScheduleAbsolute(320, () => innerSubscriptions["foo"].Dispose());
  1218. scheduler.Start();
  1219. Assert.Equal(5, inners.Count);
  1220. res["foo"].Messages.AssertEqual(
  1221. OnNext(220, "oof "),
  1222. OnNext(240, " OoF "),
  1223. OnNext(310, " Oof"),
  1224. OnCompleted<string>(310)
  1225. );
  1226. res["baR"].Messages.AssertEqual(
  1227. OnNext(270, " Rab"),
  1228. OnNext(390, "rab "),
  1229. OnNext(420, " RAB "),
  1230. OnCompleted<string>(420)
  1231. );
  1232. res["Baz"].Messages.AssertEqual(
  1233. OnNext(350, " zaB "),
  1234. OnNext(480, " zab"),
  1235. OnNext(510, " ZAb "),
  1236. OnCompleted<string>(510)
  1237. );
  1238. res["qux"].Messages.AssertEqual(
  1239. OnNext(360, " xuq "),
  1240. OnCompleted<string>(570)
  1241. );
  1242. res["FOO"].Messages.AssertEqual(
  1243. OnNext(470, " OOF"),
  1244. OnNext(530, " oOf "),
  1245. OnCompleted<string>(570)
  1246. );
  1247. xs.Subscriptions.AssertEqual(
  1248. Subscribe(200, 570)
  1249. );
  1250. }
  1251. [Fact]
  1252. public void GroupByUntil_Inner_Multiple_Independence()
  1253. {
  1254. var scheduler = new TestScheduler();
  1255. var xs = scheduler.CreateHotObservable(
  1256. OnNext(90, "error"),
  1257. OnNext(110, "error"),
  1258. OnNext(130, "error"),
  1259. OnNext(220, " foo"),
  1260. OnNext(240, " FoO "),
  1261. OnNext(270, "baR "),
  1262. OnNext(310, "foO "),
  1263. OnNext(350, " Baz "),
  1264. OnNext(360, " qux "),
  1265. OnNext(390, " bar"),
  1266. OnNext(420, " BAR "),
  1267. OnNext(470, "FOO "),
  1268. OnNext(480, "baz "),
  1269. OnNext(510, " bAZ "),
  1270. OnNext(530, " fOo "),
  1271. OnCompleted<string>(570),
  1272. OnNext(580, "error"),
  1273. OnCompleted<string>(600),
  1274. OnError<string>(650, new Exception())
  1275. );
  1276. var comparer = new GroupByComparer(scheduler);
  1277. var outer = default(IObservable<IGroupedObservable<string, string>>);
  1278. var outerSubscription = default(IDisposable);
  1279. var inners = new Dictionary<string, IObservable<string>>();
  1280. var innerSubscriptions = new Dictionary<string, IDisposable>();
  1281. var res = new Dictionary<string, ITestableObserver<string>>();
  1282. var outerResults = scheduler.CreateObserver<string>();
  1283. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), comparer));
  1284. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  1285. {
  1286. outerResults.OnNext(group.Key);
  1287. var result = scheduler.CreateObserver<string>();
  1288. inners[group.Key] = group;
  1289. res[group.Key] = result;
  1290. innerSubscriptions[group.Key] = group.Subscribe(result);
  1291. }, outerResults.OnError, outerResults.OnCompleted));
  1292. scheduler.ScheduleAbsolute(Disposed, () =>
  1293. {
  1294. outerSubscription.Dispose();
  1295. foreach (var d in innerSubscriptions.Values)
  1296. {
  1297. d.Dispose();
  1298. }
  1299. });
  1300. scheduler.ScheduleAbsolute(320, () => innerSubscriptions["foo"].Dispose());
  1301. scheduler.ScheduleAbsolute(280, () => innerSubscriptions["baR"].Dispose());
  1302. scheduler.ScheduleAbsolute(355, () => innerSubscriptions["Baz"].Dispose());
  1303. scheduler.ScheduleAbsolute(400, () => innerSubscriptions["qux"].Dispose());
  1304. scheduler.Start();
  1305. Assert.Equal(5, inners.Count);
  1306. res["foo"].Messages.AssertEqual(
  1307. OnNext(220, "oof "),
  1308. OnNext(240, " OoF "),
  1309. OnNext(310, " Oof"),
  1310. OnCompleted<string>(310)
  1311. );
  1312. res["baR"].Messages.AssertEqual(
  1313. OnNext(270, " Rab")
  1314. );
  1315. res["Baz"].Messages.AssertEqual(
  1316. OnNext(350, " zaB ")
  1317. );
  1318. res["qux"].Messages.AssertEqual(
  1319. OnNext(360, " xuq ")
  1320. );
  1321. res["FOO"].Messages.AssertEqual(
  1322. OnNext(470, " OOF"),
  1323. OnNext(530, " oOf "),
  1324. OnCompleted<string>(570)
  1325. );
  1326. xs.Subscriptions.AssertEqual(
  1327. Subscribe(200, 570)
  1328. );
  1329. }
  1330. [Fact]
  1331. public void GroupByUntil_Inner_Escape_Complete()
  1332. {
  1333. var scheduler = new TestScheduler();
  1334. var xs = scheduler.CreateHotObservable(
  1335. OnNext(220, " foo"),
  1336. OnNext(240, " FoO "),
  1337. OnNext(310, "foO "),
  1338. OnNext(470, "FOO "),
  1339. OnNext(530, " fOo "),
  1340. OnCompleted<string>(570)
  1341. );
  1342. var outer = default(IObservable<IGroupedObservable<string, string>>);
  1343. var outerSubscription = default(IDisposable);
  1344. var inner = default(IObservable<string>);
  1345. var innerSubscription = default(IDisposable);
  1346. var res = scheduler.CreateObserver<string>();
  1347. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), g => g.Skip(2)));
  1348. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  1349. {
  1350. inner = group;
  1351. }));
  1352. scheduler.ScheduleAbsolute(600, () => innerSubscription = inner.Subscribe(res));
  1353. scheduler.ScheduleAbsolute(Disposed, () =>
  1354. {
  1355. outerSubscription.Dispose();
  1356. innerSubscription.Dispose();
  1357. });
  1358. scheduler.Start();
  1359. xs.Subscriptions.AssertEqual(
  1360. Subscribe(200, 570)
  1361. );
  1362. res.Messages.AssertEqual(
  1363. OnCompleted<string>(600)
  1364. );
  1365. }
  1366. [Fact]
  1367. public void GroupByUntil_Inner_Escape_Error()
  1368. {
  1369. var scheduler = new TestScheduler();
  1370. var ex = new Exception();
  1371. var xs = scheduler.CreateHotObservable(
  1372. OnNext(220, " foo"),
  1373. OnNext(240, " FoO "),
  1374. OnNext(310, "foO "),
  1375. OnNext(470, "FOO "),
  1376. OnNext(530, " fOo "),
  1377. OnError<string>(570, ex)
  1378. );
  1379. var outer = default(IObservable<IGroupedObservable<string, string>>);
  1380. var outerSubscription = default(IDisposable);
  1381. var inner = default(IObservable<string>);
  1382. var innerSubscription = default(IDisposable);
  1383. var res = scheduler.CreateObserver<string>();
  1384. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), g => g.Skip(2)));
  1385. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  1386. {
  1387. inner = group;
  1388. }, _ => { }));
  1389. scheduler.ScheduleAbsolute(600, () => innerSubscription = inner.Subscribe(res));
  1390. scheduler.ScheduleAbsolute(Disposed, () =>
  1391. {
  1392. outerSubscription.Dispose();
  1393. innerSubscription.Dispose();
  1394. });
  1395. scheduler.Start();
  1396. xs.Subscriptions.AssertEqual(
  1397. Subscribe(200, 570)
  1398. );
  1399. res.Messages.AssertEqual(
  1400. OnError<string>(600, ex)
  1401. );
  1402. }
  1403. [Fact]
  1404. public void GroupByUntil_Inner_Escape_Dispose()
  1405. {
  1406. var scheduler = new TestScheduler();
  1407. var xs = scheduler.CreateHotObservable(
  1408. OnNext(220, " foo"),
  1409. OnNext(240, " FoO "),
  1410. OnNext(310, "foO "),
  1411. OnNext(470, "FOO "),
  1412. OnNext(530, " fOo "),
  1413. OnError<string>(570, new Exception())
  1414. );
  1415. var outer = default(IObservable<IGroupedObservable<string, string>>);
  1416. var outerSubscription = default(IDisposable);
  1417. var inner = default(IObservable<string>);
  1418. var innerSubscription = default(IDisposable);
  1419. var res = scheduler.CreateObserver<string>();
  1420. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), g => g.Skip(2)));
  1421. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  1422. {
  1423. inner = group;
  1424. }));
  1425. scheduler.ScheduleAbsolute(290, () => outerSubscription.Dispose());
  1426. scheduler.ScheduleAbsolute(600, () => innerSubscription = inner.Subscribe(res));
  1427. scheduler.ScheduleAbsolute(Disposed, () =>
  1428. {
  1429. innerSubscription.Dispose();
  1430. });
  1431. scheduler.Start();
  1432. xs.Subscriptions.AssertEqual(
  1433. Subscribe(200, 290)
  1434. );
  1435. res.Messages.AssertEqual(
  1436. );
  1437. }
  1438. [Fact]
  1439. public void GroupByUntil_Default()
  1440. {
  1441. var scheduler = new TestScheduler();
  1442. var keyInvoked = 0;
  1443. var eleInvoked = 0;
  1444. var xs = scheduler.CreateHotObservable(
  1445. OnNext(90, "error"),
  1446. OnNext(110, "error"),
  1447. OnNext(130, "error"),
  1448. OnNext(220, " foo"),
  1449. OnNext(240, " FoO "),
  1450. OnNext(270, "baR "),
  1451. OnNext(310, "foO "),
  1452. OnNext(350, " Baz "),
  1453. OnNext(360, " qux "),
  1454. OnNext(390, " bar"),
  1455. OnNext(420, " BAR "),
  1456. OnNext(470, "FOO "),
  1457. OnNext(480, "baz "),
  1458. OnNext(510, " bAZ "),
  1459. OnNext(530, " fOo "),
  1460. OnCompleted<string>(570),
  1461. OnNext(580, "error"),
  1462. OnCompleted<string>(600),
  1463. OnError<string>(650, new Exception())
  1464. );
  1465. var res = scheduler.Start(() =>
  1466. xs.GroupByUntil(
  1467. x =>
  1468. {
  1469. keyInvoked++;
  1470. return x.Trim().ToLower();
  1471. },
  1472. x =>
  1473. {
  1474. eleInvoked++;
  1475. return Reverse(x);
  1476. },
  1477. g => g.Skip(2)
  1478. ).Select(g => g.Key)
  1479. );
  1480. res.Messages.AssertEqual(
  1481. OnNext(220, "foo"),
  1482. OnNext(270, "bar"),
  1483. OnNext(350, "baz"),
  1484. OnNext(360, "qux"),
  1485. OnNext(470, "foo"),
  1486. OnCompleted<string>(570)
  1487. );
  1488. xs.Subscriptions.AssertEqual(
  1489. Subscribe(200, 570)
  1490. );
  1491. Assert.Equal(12, keyInvoked);
  1492. Assert.Equal(12, eleInvoked);
  1493. }
  1494. [Fact]
  1495. public void GroupByUntil_DurationSelector_Throws()
  1496. {
  1497. var scheduler = new TestScheduler();
  1498. var xs = scheduler.CreateHotObservable(
  1499. OnNext(210, "foo")
  1500. );
  1501. var ex = new Exception();
  1502. var res = scheduler.Start(() =>
  1503. xs.GroupByUntil<string, string, string>(x => x, g => { throw ex; })
  1504. );
  1505. res.Messages.AssertEqual(
  1506. OnError<IGroupedObservable<string, string>>(210, ex)
  1507. );
  1508. xs.Subscriptions.AssertEqual(
  1509. Subscribe(200, 210)
  1510. );
  1511. }
  1512. [Fact]
  1513. public void GroupByUntil_NullKeys_Simple_Never()
  1514. {
  1515. var scheduler = new TestScheduler();
  1516. var xs = scheduler.CreateHotObservable(
  1517. OnNext(220, "bar"),
  1518. OnNext(240, "foo"),
  1519. OnNext(310, "qux"),
  1520. OnNext(470, "baz"),
  1521. OnCompleted<string>(500)
  1522. );
  1523. var res = scheduler.Start(() => xs.GroupByUntil(x => x[0] == 'b' ? null : x.ToUpper(), g => Observable.Never<Unit>()).SelectMany(g => g, (g, x) => (g.Key ?? "(null)") + x));
  1524. res.Messages.AssertEqual(
  1525. OnNext(220, "(null)bar"),
  1526. OnNext(240, "FOOfoo"),
  1527. OnNext(310, "QUXqux"),
  1528. OnNext(470, "(null)baz"),
  1529. OnCompleted<string>(500)
  1530. );
  1531. xs.Subscriptions.AssertEqual(
  1532. Subscribe(200, 500)
  1533. );
  1534. }
  1535. [Fact]
  1536. public void GroupByUntil_NullKeys_Simple_Expire1()
  1537. {
  1538. var scheduler = new TestScheduler();
  1539. var xs = scheduler.CreateHotObservable(
  1540. OnNext(220, "bar"),
  1541. OnNext(240, "foo"),
  1542. OnNext(310, "qux"),
  1543. OnNext(470, "baz"),
  1544. OnCompleted<string>(500)
  1545. );
  1546. var n = 0;
  1547. var res = scheduler.Start(() => xs.GroupByUntil(x => x[0] == 'b' ? null : x.ToUpper(), g => { if (g.Key == null) { n++; } return Observable.Timer(TimeSpan.FromTicks(50), scheduler); }).SelectMany(g => g, (g, x) => (g.Key ?? "(null)") + x));
  1548. Assert.Equal(2, n);
  1549. res.Messages.AssertEqual(
  1550. OnNext(220, "(null)bar"),
  1551. OnNext(240, "FOOfoo"),
  1552. OnNext(310, "QUXqux"),
  1553. OnNext(470, "(null)baz"),
  1554. OnCompleted<string>(500)
  1555. );
  1556. xs.Subscriptions.AssertEqual(
  1557. Subscribe(200, 500)
  1558. );
  1559. }
  1560. [Fact]
  1561. public void GroupByUntil_NullKeys_Simple_Expire2()
  1562. {
  1563. var scheduler = new TestScheduler();
  1564. var xs = scheduler.CreateHotObservable(
  1565. OnNext(220, "bar"),
  1566. OnNext(240, "foo"),
  1567. OnNext(310, "qux"),
  1568. OnNext(470, "baz"),
  1569. OnCompleted<string>(500)
  1570. );
  1571. var n = 0;
  1572. var res = scheduler.Start(() => xs.GroupByUntil(x => x[0] == 'b' ? null : x.ToUpper(), g => { if (g.Key == null) { n++; } return Observable.Timer(TimeSpan.FromTicks(50), scheduler).IgnoreElements(); }).SelectMany(g => g, (g, x) => (g.Key ?? "(null)") + x));
  1573. Assert.Equal(2, n);
  1574. res.Messages.AssertEqual(
  1575. OnNext(220, "(null)bar"),
  1576. OnNext(240, "FOOfoo"),
  1577. OnNext(310, "QUXqux"),
  1578. OnNext(470, "(null)baz"),
  1579. OnCompleted<string>(500)
  1580. );
  1581. xs.Subscriptions.AssertEqual(
  1582. Subscribe(200, 500)
  1583. );
  1584. }
  1585. [Fact]
  1586. public void GroupByUntil_NullKeys_Error()
  1587. {
  1588. var scheduler = new TestScheduler();
  1589. var ex = new Exception();
  1590. var xs = scheduler.CreateHotObservable(
  1591. OnNext(220, "bar"),
  1592. OnNext(240, "foo"),
  1593. OnNext(310, "qux"),
  1594. OnNext(470, "baz"),
  1595. OnError<string>(500, ex)
  1596. );
  1597. var nullGroup = scheduler.CreateObserver<string>();
  1598. var err = default(Exception);
  1599. scheduler.ScheduleAbsolute(200, () => xs.GroupByUntil(x => x[0] == 'b' ? null : x.ToUpper(), g => Observable.Never<Unit>()).Where(g => g.Key == null).Subscribe(g => g.Subscribe(nullGroup), ex_ => err = ex_));
  1600. scheduler.Start();
  1601. Assert.Same(ex, err);
  1602. nullGroup.Messages.AssertEqual(
  1603. OnNext(220, "bar"),
  1604. OnNext(470, "baz"),
  1605. OnError<string>(500, ex)
  1606. );
  1607. xs.Subscriptions.AssertEqual(
  1608. Subscribe(200, 500)
  1609. );
  1610. }
  1611. #endregion
  1612. #region + GroupByUntil w/capacity +
  1613. private const int _groupByUntilCapacity = 1024;
  1614. [Fact]
  1615. public void GroupByUntil_Capacity_ArgumentChecking()
  1616. {
  1617. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(default, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity, EqualityComparer<int>.Default));
  1618. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, default, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity, EqualityComparer<int>.Default));
  1619. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, default, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity, EqualityComparer<int>.Default));
  1620. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, default(Func<IGroupedObservable<int, int>, IObservable<int>>), _groupByUntilCapacity, EqualityComparer<int>.Default));
  1621. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity, default));
  1622. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(default, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity));
  1623. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, default, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity));
  1624. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, default, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity));
  1625. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, default(Func<IGroupedObservable<int, int>, IObservable<int>>), _groupByUntilCapacity));
  1626. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(default, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity, EqualityComparer<int>.Default));
  1627. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, default, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity, EqualityComparer<int>.Default));
  1628. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, default(Func<IGroupedObservable<int, int>, IObservable<int>>), _groupByUntilCapacity, EqualityComparer<int>.Default));
  1629. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity, default));
  1630. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(default, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity));
  1631. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, default, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, _groupByUntilCapacity));
  1632. ReactiveAssert.Throws<ArgumentNullException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, default(Func<IGroupedObservable<int, int>, IObservable<int>>), _groupByUntilCapacity));
  1633. ReactiveAssert.Throws<ArgumentOutOfRangeException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, -1, EqualityComparer<int>.Default));
  1634. ReactiveAssert.Throws<ArgumentOutOfRangeException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, -1));
  1635. ReactiveAssert.Throws<ArgumentOutOfRangeException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, -1, EqualityComparer<int>.Default));
  1636. ReactiveAssert.Throws<ArgumentOutOfRangeException>(() => Observable.GroupByUntil(DummyObservable<int>.Instance, DummyFunc<int, int>.Instance, DummyFunc<IGroupedObservable<int, int>, IObservable<int>>.Instance, -1));
  1637. }
  1638. [Fact]
  1639. public void GroupByUntil_Capacity_WithKeyComparer()
  1640. {
  1641. var scheduler = new TestScheduler();
  1642. var keyInvoked = 0;
  1643. var xs = scheduler.CreateHotObservable(
  1644. OnNext(90, "error"),
  1645. OnNext(110, "error"),
  1646. OnNext(130, "error"),
  1647. OnNext(220, " foo"),
  1648. OnNext(240, " FoO "),
  1649. OnNext(270, "baR "),
  1650. OnNext(310, "foO "),
  1651. OnNext(350, " Baz "),
  1652. OnNext(360, " qux "),
  1653. OnNext(390, " bar"),
  1654. OnNext(420, " BAR "),
  1655. OnNext(470, "FOO "),
  1656. OnNext(480, "baz "),
  1657. OnNext(510, " bAZ "),
  1658. OnNext(530, " fOo "),
  1659. OnCompleted<string>(570),
  1660. OnNext(580, "error"),
  1661. OnCompleted<string>(600),
  1662. OnError<string>(650, new Exception())
  1663. );
  1664. var comparer = new GroupByComparer(scheduler);
  1665. var res = scheduler.Start(() =>
  1666. xs.GroupByUntil(
  1667. x =>
  1668. {
  1669. keyInvoked++;
  1670. return x.Trim();
  1671. },
  1672. g => g.Skip(2),
  1673. _groupByUntilCapacity,
  1674. comparer
  1675. ).Select(g => g.Key)
  1676. );
  1677. res.Messages.AssertEqual(
  1678. OnNext(220, "foo"),
  1679. OnNext(270, "baR"),
  1680. OnNext(350, "Baz"),
  1681. OnNext(360, "qux"),
  1682. OnNext(470, "FOO"),
  1683. OnCompleted<string>(570)
  1684. );
  1685. xs.Subscriptions.AssertEqual(
  1686. Subscribe(200, 570)
  1687. );
  1688. Assert.Equal(12, keyInvoked);
  1689. }
  1690. [Fact]
  1691. public void GroupByUntil_Capacity_Outer_Complete()
  1692. {
  1693. var scheduler = new TestScheduler();
  1694. var keyInvoked = 0;
  1695. var eleInvoked = 0;
  1696. var xs = scheduler.CreateHotObservable(
  1697. OnNext(90, "error"),
  1698. OnNext(110, "error"),
  1699. OnNext(130, "error"),
  1700. OnNext(220, " foo"),
  1701. OnNext(240, " FoO "),
  1702. OnNext(270, "baR "),
  1703. OnNext(310, "foO "),
  1704. OnNext(350, " Baz "),
  1705. OnNext(360, " qux "),
  1706. OnNext(390, " bar"),
  1707. OnNext(420, " BAR "),
  1708. OnNext(470, "FOO "),
  1709. OnNext(480, "baz "),
  1710. OnNext(510, " bAZ "),
  1711. OnNext(530, " fOo "),
  1712. OnCompleted<string>(570),
  1713. OnNext(580, "error"),
  1714. OnCompleted<string>(600),
  1715. OnError<string>(650, new Exception())
  1716. );
  1717. var comparer = new GroupByComparer(scheduler);
  1718. var res = scheduler.Start(() =>
  1719. xs.GroupByUntil(
  1720. x =>
  1721. {
  1722. keyInvoked++;
  1723. return x.Trim();
  1724. },
  1725. x =>
  1726. {
  1727. eleInvoked++;
  1728. return Reverse(x);
  1729. },
  1730. g => g.Skip(2),
  1731. _groupByUntilCapacity,
  1732. comparer
  1733. ).Select(g => g.Key)
  1734. );
  1735. res.Messages.AssertEqual(
  1736. OnNext(220, "foo"),
  1737. OnNext(270, "baR"),
  1738. OnNext(350, "Baz"),
  1739. OnNext(360, "qux"),
  1740. OnNext(470, "FOO"),
  1741. OnCompleted<string>(570)
  1742. );
  1743. xs.Subscriptions.AssertEqual(
  1744. Subscribe(200, 570)
  1745. );
  1746. Assert.Equal(12, keyInvoked);
  1747. Assert.Equal(12, eleInvoked);
  1748. }
  1749. [Fact]
  1750. public void GroupByUntil_Capacity_Outer_Error()
  1751. {
  1752. var scheduler = new TestScheduler();
  1753. var keyInvoked = 0;
  1754. var eleInvoked = 0;
  1755. var ex = new Exception();
  1756. var xs = scheduler.CreateHotObservable(
  1757. OnNext(90, "error"),
  1758. OnNext(110, "error"),
  1759. OnNext(130, "error"),
  1760. OnNext(220, " foo"),
  1761. OnNext(240, " FoO "),
  1762. OnNext(270, "baR "),
  1763. OnNext(310, "foO "),
  1764. OnNext(350, " Baz "),
  1765. OnNext(360, " qux "),
  1766. OnNext(390, " bar"),
  1767. OnNext(420, " BAR "),
  1768. OnNext(470, "FOO "),
  1769. OnNext(480, "baz "),
  1770. OnNext(510, " bAZ "),
  1771. OnNext(530, " fOo "),
  1772. OnError<string>(570, ex),
  1773. OnNext(580, "error"),
  1774. OnCompleted<string>(600),
  1775. OnError<string>(650, new Exception())
  1776. );
  1777. var comparer = new GroupByComparer(scheduler);
  1778. var res = scheduler.Start(() =>
  1779. xs.GroupByUntil(
  1780. x =>
  1781. {
  1782. keyInvoked++;
  1783. return x.Trim();
  1784. },
  1785. x =>
  1786. {
  1787. eleInvoked++;
  1788. return Reverse(x);
  1789. },
  1790. g => g.Skip(2),
  1791. _groupByUntilCapacity,
  1792. comparer
  1793. ).Select(g => g.Key)
  1794. );
  1795. res.Messages.AssertEqual(
  1796. OnNext(220, "foo"),
  1797. OnNext(270, "baR"),
  1798. OnNext(350, "Baz"),
  1799. OnNext(360, "qux"),
  1800. OnNext(470, "FOO"),
  1801. OnError<string>(570, ex)
  1802. );
  1803. xs.Subscriptions.AssertEqual(
  1804. Subscribe(200, 570)
  1805. );
  1806. Assert.Equal(12, keyInvoked);
  1807. Assert.Equal(12, eleInvoked);
  1808. }
  1809. [Fact]
  1810. public void GroupByUntil_Capacity_Outer_Dispose()
  1811. {
  1812. var scheduler = new TestScheduler();
  1813. var keyInvoked = 0;
  1814. var eleInvoked = 0;
  1815. var xs = scheduler.CreateHotObservable(
  1816. OnNext(90, "error"),
  1817. OnNext(110, "error"),
  1818. OnNext(130, "error"),
  1819. OnNext(220, " foo"),
  1820. OnNext(240, " FoO "),
  1821. OnNext(270, "baR "),
  1822. OnNext(310, "foO "),
  1823. OnNext(350, " Baz "),
  1824. OnNext(360, " qux "),
  1825. OnNext(390, " bar"),
  1826. OnNext(420, " BAR "),
  1827. OnNext(470, "FOO "),
  1828. OnNext(480, "baz "),
  1829. OnNext(510, " bAZ "),
  1830. OnNext(530, " fOo "),
  1831. OnCompleted<string>(570),
  1832. OnNext(580, "error"),
  1833. OnCompleted<string>(600),
  1834. OnError<string>(650, new Exception())
  1835. );
  1836. var comparer = new GroupByComparer(scheduler);
  1837. var res = scheduler.Start(() =>
  1838. xs.GroupByUntil(
  1839. x =>
  1840. {
  1841. keyInvoked++;
  1842. return x.Trim();
  1843. },
  1844. x =>
  1845. {
  1846. eleInvoked++;
  1847. return Reverse(x);
  1848. },
  1849. g => g.Skip(2),
  1850. _groupByUntilCapacity,
  1851. comparer
  1852. ).Select(g => g.Key),
  1853. 355
  1854. );
  1855. res.Messages.AssertEqual(
  1856. OnNext(220, "foo"),
  1857. OnNext(270, "baR"),
  1858. OnNext(350, "Baz")
  1859. );
  1860. xs.Subscriptions.AssertEqual(
  1861. Subscribe(200, 355)
  1862. );
  1863. Assert.Equal(5, keyInvoked);
  1864. Assert.Equal(5, eleInvoked);
  1865. }
  1866. [Fact]
  1867. public void GroupByUntil_Capacity_Outer_KeyThrow()
  1868. {
  1869. var scheduler = new TestScheduler();
  1870. var keyInvoked = 0;
  1871. var eleInvoked = 0;
  1872. var ex = new Exception();
  1873. var xs = scheduler.CreateHotObservable(
  1874. OnNext(90, "error"),
  1875. OnNext(110, "error"),
  1876. OnNext(130, "error"),
  1877. OnNext(220, " foo"),
  1878. OnNext(240, " FoO "),
  1879. OnNext(270, "baR "),
  1880. OnNext(310, "foO "),
  1881. OnNext(350, " Baz "),
  1882. OnNext(360, " qux "),
  1883. OnNext(390, " bar"),
  1884. OnNext(420, " BAR "),
  1885. OnNext(470, "FOO "),
  1886. OnNext(480, "baz "),
  1887. OnNext(510, " bAZ "),
  1888. OnNext(530, " fOo "),
  1889. OnCompleted<string>(570),
  1890. OnNext(580, "error"),
  1891. OnCompleted<string>(600),
  1892. OnError<string>(650, new Exception())
  1893. );
  1894. var comparer = new GroupByComparer(scheduler);
  1895. var res = scheduler.Start(() =>
  1896. xs.GroupByUntil(
  1897. x =>
  1898. {
  1899. keyInvoked++;
  1900. if (keyInvoked == 10)
  1901. {
  1902. throw ex;
  1903. }
  1904. return x.Trim();
  1905. },
  1906. x =>
  1907. {
  1908. eleInvoked++;
  1909. return Reverse(x);
  1910. },
  1911. g => g.Skip(2),
  1912. _groupByUntilCapacity,
  1913. comparer
  1914. ).Select(g => g.Key)
  1915. );
  1916. res.Messages.AssertEqual(
  1917. OnNext(220, "foo"),
  1918. OnNext(270, "baR"),
  1919. OnNext(350, "Baz"),
  1920. OnNext(360, "qux"),
  1921. OnNext(470, "FOO"),
  1922. OnError<string>(480, ex)
  1923. );
  1924. xs.Subscriptions.AssertEqual(
  1925. Subscribe(200, 480)
  1926. );
  1927. Assert.Equal(10, keyInvoked);
  1928. Assert.Equal(9, eleInvoked);
  1929. }
  1930. [Fact]
  1931. public void GroupByUntil_Capacity_Outer_EleThrow()
  1932. {
  1933. var scheduler = new TestScheduler();
  1934. var keyInvoked = 0;
  1935. var eleInvoked = 0;
  1936. var ex = new Exception();
  1937. var xs = scheduler.CreateHotObservable(
  1938. OnNext(90, "error"),
  1939. OnNext(110, "error"),
  1940. OnNext(130, "error"),
  1941. OnNext(220, " foo"),
  1942. OnNext(240, " FoO "),
  1943. OnNext(270, "baR "),
  1944. OnNext(310, "foO "),
  1945. OnNext(350, " Baz "),
  1946. OnNext(360, " qux "),
  1947. OnNext(390, " bar"),
  1948. OnNext(420, " BAR "),
  1949. OnNext(470, "FOO "),
  1950. OnNext(480, "baz "),
  1951. OnNext(510, " bAZ "),
  1952. OnNext(530, " fOo "),
  1953. OnCompleted<string>(570),
  1954. OnNext(580, "error"),
  1955. OnCompleted<string>(600),
  1956. OnError<string>(650, new Exception())
  1957. );
  1958. var comparer = new GroupByComparer(scheduler);
  1959. var res = scheduler.Start(() =>
  1960. xs.GroupByUntil(
  1961. x =>
  1962. {
  1963. keyInvoked++;
  1964. return x.Trim();
  1965. },
  1966. x =>
  1967. {
  1968. eleInvoked++;
  1969. if (eleInvoked == 10)
  1970. {
  1971. throw ex;
  1972. }
  1973. return Reverse(x);
  1974. },
  1975. g => g.Skip(2),
  1976. _groupByUntilCapacity,
  1977. comparer
  1978. ).Select(g => g.Key)
  1979. );
  1980. res.Messages.AssertEqual(
  1981. OnNext(220, "foo"),
  1982. OnNext(270, "baR"),
  1983. OnNext(350, "Baz"),
  1984. OnNext(360, "qux"),
  1985. OnNext(470, "FOO"),
  1986. OnError<string>(480, ex)
  1987. );
  1988. xs.Subscriptions.AssertEqual(
  1989. Subscribe(200, 480)
  1990. );
  1991. Assert.Equal(10, keyInvoked);
  1992. Assert.Equal(10, eleInvoked);
  1993. }
  1994. [Fact]
  1995. public void GroupByUntil_Capacity_Outer_ComparerEqualsThrow()
  1996. {
  1997. var scheduler = new TestScheduler();
  1998. var keyInvoked = 0;
  1999. var eleInvoked = 0;
  2000. var xs = scheduler.CreateHotObservable(
  2001. OnNext(90, "error"),
  2002. OnNext(110, "error"),
  2003. OnNext(130, "error"),
  2004. OnNext(220, " foo"),
  2005. OnNext(240, " FoO "),
  2006. OnNext(270, "baR "),
  2007. OnNext(310, "foO "),
  2008. OnNext(350, " Baz "),
  2009. OnNext(360, " qux "),
  2010. OnNext(390, " bar"),
  2011. OnNext(420, " BAR "),
  2012. OnNext(470, "FOO "),
  2013. OnNext(480, "baz "),
  2014. OnNext(510, " bAZ "),
  2015. OnNext(530, " fOo "),
  2016. OnCompleted<string>(570),
  2017. OnNext(580, "error"),
  2018. OnCompleted<string>(600),
  2019. OnError<string>(650, new Exception())
  2020. );
  2021. var comparer = new GroupByComparer(scheduler, 250, ushort.MaxValue);
  2022. var res = scheduler.Start(() =>
  2023. xs.GroupByUntil(
  2024. x =>
  2025. {
  2026. keyInvoked++;
  2027. return x.Trim();
  2028. },
  2029. x =>
  2030. {
  2031. eleInvoked++;
  2032. return Reverse(x);
  2033. },
  2034. g => g.Skip(2),
  2035. _groupByUntilCapacity,
  2036. comparer
  2037. ).Select(g => g.Key)
  2038. );
  2039. res.Messages.AssertEqual(
  2040. OnNext(220, "foo"),
  2041. OnNext(270, "baR"),
  2042. OnError<string>(310, comparer.EqualsException)
  2043. );
  2044. xs.Subscriptions.AssertEqual(
  2045. Subscribe(200, 310)
  2046. );
  2047. Assert.Equal(4, keyInvoked);
  2048. Assert.Equal(3, eleInvoked);
  2049. }
  2050. [Fact]
  2051. public void GroupByUntil_Capacity_Outer_ComparerGetHashCodeThrow()
  2052. {
  2053. var scheduler = new TestScheduler();
  2054. var keyInvoked = 0;
  2055. var eleInvoked = 0;
  2056. var xs = scheduler.CreateHotObservable(
  2057. OnNext(90, "error"),
  2058. OnNext(110, "error"),
  2059. OnNext(130, "error"),
  2060. OnNext(220, " foo"),
  2061. OnNext(240, " FoO "),
  2062. OnNext(270, "baR "),
  2063. OnNext(310, "foO "),
  2064. OnNext(350, " Baz "),
  2065. OnNext(360, " qux "),
  2066. OnNext(390, " bar"),
  2067. OnNext(420, " BAR "),
  2068. OnNext(470, "FOO "),
  2069. OnNext(480, "baz "),
  2070. OnNext(510, " bAZ "),
  2071. OnNext(530, " fOo "),
  2072. OnCompleted<string>(570),
  2073. OnNext(580, "error"),
  2074. OnCompleted<string>(600),
  2075. OnError<string>(650, new Exception())
  2076. );
  2077. var comparer = new GroupByComparer(scheduler, ushort.MaxValue, 410);
  2078. var res = scheduler.Start(() =>
  2079. xs.GroupByUntil(
  2080. x =>
  2081. {
  2082. keyInvoked++;
  2083. return x.Trim();
  2084. },
  2085. x =>
  2086. {
  2087. eleInvoked++;
  2088. return Reverse(x);
  2089. },
  2090. g => g.Skip(2),
  2091. _groupByUntilCapacity,
  2092. comparer
  2093. ).Select(g => g.Key)
  2094. );
  2095. res.Messages.AssertEqual(
  2096. OnNext(220, "foo"),
  2097. OnNext(270, "baR"),
  2098. OnNext(350, "Baz"),
  2099. OnNext(360, "qux"),
  2100. OnError<string>(420, comparer.HashCodeException)
  2101. );
  2102. xs.Subscriptions.AssertEqual(
  2103. Subscribe(200, 420)
  2104. );
  2105. Assert.Equal(8, keyInvoked);
  2106. Assert.Equal(7, eleInvoked);
  2107. }
  2108. [Fact]
  2109. public void GroupByUntil_Capacity_Inner_Complete()
  2110. {
  2111. var scheduler = new TestScheduler();
  2112. var xs = scheduler.CreateHotObservable(
  2113. OnNext(90, "error"),
  2114. OnNext(110, "error"),
  2115. OnNext(130, "error"),
  2116. OnNext(220, " foo"),
  2117. OnNext(240, " FoO "),
  2118. OnNext(270, "baR "),
  2119. OnNext(310, "foO "),
  2120. OnNext(350, " Baz "),
  2121. OnNext(360, " qux "),
  2122. OnNext(390, " bar"),
  2123. OnNext(420, " BAR "),
  2124. OnNext(470, "FOO "),
  2125. OnNext(480, "baz "),
  2126. OnNext(510, " bAZ "),
  2127. OnNext(530, " fOo "),
  2128. OnCompleted<string>(570),
  2129. OnNext(580, "error"),
  2130. OnCompleted<string>(600),
  2131. OnError<string>(650, new Exception())
  2132. );
  2133. var comparer = new GroupByComparer(scheduler);
  2134. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2135. var outerSubscription = default(IDisposable);
  2136. var inners = new Dictionary<string, IObservable<string>>();
  2137. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2138. var res = new Dictionary<string, ITestableObserver<string>>();
  2139. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2140. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2141. {
  2142. var result = scheduler.CreateObserver<string>();
  2143. inners[group.Key] = group;
  2144. res[group.Key] = result;
  2145. scheduler.ScheduleRelative(100, () => innerSubscriptions[group.Key] = group.Subscribe(result));
  2146. }));
  2147. scheduler.ScheduleAbsolute(Disposed, () =>
  2148. {
  2149. outerSubscription.Dispose();
  2150. foreach (var d in innerSubscriptions.Values)
  2151. {
  2152. d.Dispose();
  2153. }
  2154. });
  2155. scheduler.Start();
  2156. Assert.Equal(5, inners.Count);
  2157. res["foo"].Messages.AssertEqual(
  2158. OnCompleted<string>(320)
  2159. );
  2160. res["baR"].Messages.AssertEqual(
  2161. OnNext(390, "rab "),
  2162. OnNext(420, " RAB "), // Breaking change > v2.2 - prior to resolving a deadlock, the group would get closed prior to letting this message through
  2163. OnCompleted<string>(420)
  2164. );
  2165. res["Baz"].Messages.AssertEqual(
  2166. OnNext(480, " zab"),
  2167. OnNext(510, " ZAb "), // Breaking change > v2.2 - prior to resolving a deadlock, the group would get closed prior to letting this message through
  2168. OnCompleted<string>(510)
  2169. );
  2170. res["qux"].Messages.AssertEqual(
  2171. OnCompleted<string>(570)
  2172. );
  2173. res["FOO"].Messages.AssertEqual(
  2174. OnCompleted<string>(570)
  2175. );
  2176. xs.Subscriptions.AssertEqual(
  2177. Subscribe(200, 570)
  2178. );
  2179. }
  2180. [Fact]
  2181. public void GroupByUntil_Capacity_Inner_Complete_All()
  2182. {
  2183. var scheduler = new TestScheduler();
  2184. var xs = scheduler.CreateHotObservable(
  2185. OnNext(90, "error"),
  2186. OnNext(110, "error"),
  2187. OnNext(130, "error"),
  2188. OnNext(220, " foo"),
  2189. OnNext(240, " FoO "),
  2190. OnNext(270, "baR "),
  2191. OnNext(310, "foO "),
  2192. OnNext(350, " Baz "),
  2193. OnNext(360, " qux "),
  2194. OnNext(390, " bar"),
  2195. OnNext(420, " BAR "),
  2196. OnNext(470, "FOO "),
  2197. OnNext(480, "baz "),
  2198. OnNext(510, " bAZ "),
  2199. OnNext(530, " fOo "),
  2200. OnCompleted<string>(570),
  2201. OnNext(580, "error"),
  2202. OnCompleted<string>(600),
  2203. OnError<string>(650, new Exception())
  2204. );
  2205. var comparer = new GroupByComparer(scheduler);
  2206. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2207. var outerSubscription = default(IDisposable);
  2208. var inners = new Dictionary<string, IObservable<string>>();
  2209. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2210. var res = new Dictionary<string, ITestableObserver<string>>();
  2211. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2212. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2213. {
  2214. var result = scheduler.CreateObserver<string>();
  2215. inners[group.Key] = group;
  2216. res[group.Key] = result;
  2217. innerSubscriptions[group.Key] = group.Subscribe(result);
  2218. }));
  2219. scheduler.ScheduleAbsolute(Disposed, () =>
  2220. {
  2221. outerSubscription.Dispose();
  2222. foreach (var d in innerSubscriptions.Values)
  2223. {
  2224. d.Dispose();
  2225. }
  2226. });
  2227. scheduler.Start();
  2228. Assert.Equal(5, inners.Count);
  2229. res["foo"].Messages.AssertEqual(
  2230. OnNext(220, "oof "),
  2231. OnNext(240, " OoF "),
  2232. OnNext(310, " Oof"),
  2233. OnCompleted<string>(310)
  2234. );
  2235. res["baR"].Messages.AssertEqual(
  2236. OnNext(270, " Rab"),
  2237. OnNext(390, "rab "),
  2238. OnNext(420, " RAB "),
  2239. OnCompleted<string>(420)
  2240. );
  2241. res["Baz"].Messages.AssertEqual(
  2242. OnNext(350, " zaB "),
  2243. OnNext(480, " zab"),
  2244. OnNext(510, " ZAb "),
  2245. OnCompleted<string>(510)
  2246. );
  2247. res["qux"].Messages.AssertEqual(
  2248. OnNext(360, " xuq "),
  2249. OnCompleted<string>(570)
  2250. );
  2251. res["FOO"].Messages.AssertEqual(
  2252. OnNext(470, " OOF"),
  2253. OnNext(530, " oOf "),
  2254. OnCompleted<string>(570)
  2255. );
  2256. xs.Subscriptions.AssertEqual(
  2257. Subscribe(200, 570)
  2258. );
  2259. }
  2260. [Fact]
  2261. public void GroupByUntil_Capacity_Inner_Error()
  2262. {
  2263. var scheduler = new TestScheduler();
  2264. var ex1 = new Exception();
  2265. var xs = scheduler.CreateHotObservable(
  2266. OnNext(90, "error"),
  2267. OnNext(110, "error"),
  2268. OnNext(130, "error"),
  2269. OnNext(220, " foo"),
  2270. OnNext(240, " FoO "),
  2271. OnNext(270, "baR "),
  2272. OnNext(310, "foO "),
  2273. OnNext(350, " Baz "),
  2274. OnNext(360, " qux "),
  2275. OnNext(390, " bar"),
  2276. OnNext(420, " BAR "),
  2277. OnNext(470, "FOO "),
  2278. OnNext(480, "baz "),
  2279. OnNext(510, " bAZ "),
  2280. OnNext(530, " fOo "),
  2281. OnError<string>(570, ex1),
  2282. OnNext(580, "error"),
  2283. OnCompleted<string>(600),
  2284. OnError<string>(650, new Exception())
  2285. );
  2286. var comparer = new GroupByComparer(scheduler);
  2287. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2288. var outerSubscription = default(IDisposable);
  2289. var inners = new Dictionary<string, IObservable<string>>();
  2290. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2291. var res = new Dictionary<string, ITestableObserver<string>>();
  2292. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2293. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2294. {
  2295. var result = scheduler.CreateObserver<string>();
  2296. inners[group.Key] = group;
  2297. res[group.Key] = result;
  2298. scheduler.ScheduleRelative(100, () => innerSubscriptions[group.Key] = group.Subscribe(result));
  2299. }, ex => { }));
  2300. scheduler.ScheduleAbsolute(Disposed, () =>
  2301. {
  2302. outerSubscription.Dispose();
  2303. foreach (var d in innerSubscriptions.Values)
  2304. {
  2305. d.Dispose();
  2306. }
  2307. });
  2308. scheduler.Start();
  2309. Assert.Equal(5, inners.Count);
  2310. res["foo"].Messages.AssertEqual(
  2311. OnCompleted<string>(320)
  2312. );
  2313. res["baR"].Messages.AssertEqual(
  2314. OnNext(390, "rab "),
  2315. OnNext(420, " RAB "), // Breaking change > v2.2 - prior to resolving a deadlock, the group would get closed prior to letting this message through
  2316. OnCompleted<string>(420)
  2317. );
  2318. res["Baz"].Messages.AssertEqual(
  2319. OnNext(480, " zab"),
  2320. OnNext(510, " ZAb "), // Breaking change > v2.2 - prior to resolving a deadlock, the group would get closed prior to letting this message through
  2321. OnCompleted<string>(510)
  2322. );
  2323. res["qux"].Messages.AssertEqual(
  2324. OnError<string>(570, ex1)
  2325. );
  2326. res["FOO"].Messages.AssertEqual(
  2327. OnError<string>(570, ex1)
  2328. );
  2329. xs.Subscriptions.AssertEqual(
  2330. Subscribe(200, 570)
  2331. );
  2332. }
  2333. [Fact]
  2334. public void GroupByUntil_Capacity_Inner_Dispose()
  2335. {
  2336. var scheduler = new TestScheduler();
  2337. var xs = scheduler.CreateHotObservable(
  2338. OnNext(90, "error"),
  2339. OnNext(110, "error"),
  2340. OnNext(130, "error"),
  2341. OnNext(220, " foo"),
  2342. OnNext(240, " FoO "),
  2343. OnNext(270, "baR "),
  2344. OnNext(310, "foO "),
  2345. OnNext(350, " Baz "),
  2346. OnNext(360, " qux "),
  2347. OnNext(390, " bar"),
  2348. OnNext(420, " BAR "),
  2349. OnNext(470, "FOO "),
  2350. OnNext(480, "baz "),
  2351. OnNext(510, " bAZ "),
  2352. OnNext(530, " fOo "),
  2353. OnCompleted<string>(570),
  2354. OnNext(580, "error"),
  2355. OnCompleted<string>(600),
  2356. OnError<string>(650, new Exception())
  2357. );
  2358. var comparer = new GroupByComparer(scheduler);
  2359. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2360. var outerSubscription = default(IDisposable);
  2361. var inners = new Dictionary<string, IObservable<string>>();
  2362. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2363. var res = new Dictionary<string, ITestableObserver<string>>();
  2364. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2365. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2366. {
  2367. var result = scheduler.CreateObserver<string>();
  2368. inners[group.Key] = group;
  2369. res[group.Key] = result;
  2370. innerSubscriptions[group.Key] = group.Subscribe(result);
  2371. }));
  2372. scheduler.ScheduleAbsolute(400, () =>
  2373. {
  2374. outerSubscription.Dispose();
  2375. foreach (var d in innerSubscriptions.Values)
  2376. {
  2377. d.Dispose();
  2378. }
  2379. });
  2380. scheduler.Start();
  2381. Assert.Equal(4, inners.Count);
  2382. res["foo"].Messages.AssertEqual(
  2383. OnNext(220, "oof "),
  2384. OnNext(240, " OoF "),
  2385. OnNext(310, " Oof"),
  2386. OnCompleted<string>(310)
  2387. );
  2388. res["baR"].Messages.AssertEqual(
  2389. OnNext(270, " Rab"),
  2390. OnNext(390, "rab ")
  2391. );
  2392. res["Baz"].Messages.AssertEqual(
  2393. OnNext(350, " zaB ")
  2394. );
  2395. res["qux"].Messages.AssertEqual(
  2396. OnNext(360, " xuq ")
  2397. );
  2398. xs.Subscriptions.AssertEqual(
  2399. Subscribe(200, 400)
  2400. );
  2401. }
  2402. [Fact]
  2403. public void GroupByUntil_Capacity_Inner_KeyThrow()
  2404. {
  2405. var scheduler = new TestScheduler();
  2406. var xs = scheduler.CreateHotObservable(
  2407. OnNext(90, "error"),
  2408. OnNext(110, "error"),
  2409. OnNext(130, "error"),
  2410. OnNext(220, " foo"),
  2411. OnNext(240, " FoO "),
  2412. OnNext(270, "baR "),
  2413. OnNext(310, "foO "),
  2414. OnNext(350, " Baz "),
  2415. OnNext(360, " qux "),
  2416. OnNext(390, " bar"),
  2417. OnNext(420, " BAR "),
  2418. OnNext(470, "FOO "),
  2419. OnNext(480, "baz "),
  2420. OnNext(510, " bAZ "),
  2421. OnNext(530, " fOo "),
  2422. OnCompleted<string>(570),
  2423. OnNext(580, "error"),
  2424. OnCompleted<string>(600),
  2425. OnError<string>(650, new Exception())
  2426. );
  2427. var comparer = new GroupByComparer(scheduler);
  2428. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2429. var outerSubscription = default(IDisposable);
  2430. var inners = new Dictionary<string, IObservable<string>>();
  2431. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2432. var res = new Dictionary<string, ITestableObserver<string>>();
  2433. var keyInvoked = 0;
  2434. var ex = new Exception();
  2435. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x =>
  2436. {
  2437. keyInvoked++;
  2438. if (keyInvoked == 6)
  2439. {
  2440. throw ex;
  2441. }
  2442. return x.Trim();
  2443. }, x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2444. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2445. {
  2446. var result = scheduler.CreateObserver<string>();
  2447. inners[group.Key] = group;
  2448. res[group.Key] = result;
  2449. innerSubscriptions[group.Key] = group.Subscribe(result);
  2450. }, _ => { }));
  2451. scheduler.ScheduleAbsolute(Disposed, () =>
  2452. {
  2453. outerSubscription.Dispose();
  2454. foreach (var d in innerSubscriptions.Values)
  2455. {
  2456. d.Dispose();
  2457. }
  2458. });
  2459. scheduler.Start();
  2460. Assert.Equal(3, inners.Count);
  2461. res["foo"].Messages.AssertEqual(
  2462. OnNext(220, "oof "),
  2463. OnNext(240, " OoF "),
  2464. OnNext(310, " Oof"),
  2465. OnCompleted<string>(310)
  2466. );
  2467. res["baR"].Messages.AssertEqual(
  2468. OnNext(270, " Rab"),
  2469. OnError<string>(360, ex)
  2470. );
  2471. res["Baz"].Messages.AssertEqual(
  2472. OnNext(350, " zaB "),
  2473. OnError<string>(360, ex)
  2474. );
  2475. xs.Subscriptions.AssertEqual(
  2476. Subscribe(200, 360)
  2477. );
  2478. }
  2479. [Fact]
  2480. public void GroupByUntil_Capacity_Inner_EleThrow()
  2481. {
  2482. var scheduler = new TestScheduler();
  2483. var xs = scheduler.CreateHotObservable(
  2484. OnNext(90, "error"),
  2485. OnNext(110, "error"),
  2486. OnNext(130, "error"),
  2487. OnNext(220, " foo"),
  2488. OnNext(240, " FoO "),
  2489. OnNext(270, "baR "),
  2490. OnNext(310, "foO "),
  2491. OnNext(350, " Baz "),
  2492. OnNext(360, " qux "),
  2493. OnNext(390, " bar"),
  2494. OnNext(420, " BAR "),
  2495. OnNext(470, "FOO "),
  2496. OnNext(480, "baz "),
  2497. OnNext(510, " bAZ "),
  2498. OnNext(530, " fOo "),
  2499. OnCompleted<string>(570),
  2500. OnNext(580, "error"),
  2501. OnCompleted<string>(600),
  2502. OnError<string>(650, new Exception())
  2503. );
  2504. var comparer = new GroupByComparer(scheduler);
  2505. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2506. var outerSubscription = default(IDisposable);
  2507. var inners = new Dictionary<string, IObservable<string>>();
  2508. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2509. var res = new Dictionary<string, ITestableObserver<string>>();
  2510. var eleInvoked = 0;
  2511. var ex = new Exception();
  2512. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x =>
  2513. {
  2514. eleInvoked++;
  2515. if (eleInvoked == 6)
  2516. {
  2517. throw ex;
  2518. }
  2519. return Reverse(x);
  2520. }, g => g.Skip(2), _groupByUntilCapacity, comparer));
  2521. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2522. {
  2523. var result = scheduler.CreateObserver<string>();
  2524. inners[group.Key] = group;
  2525. res[group.Key] = result;
  2526. innerSubscriptions[group.Key] = group.Subscribe(result);
  2527. }, _ => { }));
  2528. scheduler.ScheduleAbsolute(Disposed, () =>
  2529. {
  2530. outerSubscription.Dispose();
  2531. foreach (var d in innerSubscriptions.Values)
  2532. {
  2533. d.Dispose();
  2534. }
  2535. });
  2536. scheduler.Start();
  2537. Assert.Equal(4, inners.Count);
  2538. res["foo"].Messages.AssertEqual(
  2539. OnNext(220, "oof "),
  2540. OnNext(240, " OoF "),
  2541. OnNext(310, " Oof"),
  2542. OnCompleted<string>(310)
  2543. );
  2544. res["baR"].Messages.AssertEqual(
  2545. OnNext(270, " Rab"),
  2546. OnError<string>(360, ex)
  2547. );
  2548. res["Baz"].Messages.AssertEqual(
  2549. OnNext(350, " zaB "),
  2550. OnError<string>(360, ex)
  2551. );
  2552. res["qux"].Messages.AssertEqual(
  2553. OnError<string>(360, ex)
  2554. );
  2555. xs.Subscriptions.AssertEqual(
  2556. Subscribe(200, 360)
  2557. );
  2558. }
  2559. [Fact]
  2560. public void GroupByUntil_Capacity_Inner_Comparer_EqualsThrow()
  2561. {
  2562. var scheduler = new TestScheduler();
  2563. var xs = scheduler.CreateHotObservable(
  2564. OnNext(90, "error"),
  2565. OnNext(110, "error"),
  2566. OnNext(130, "error"),
  2567. OnNext(220, " foo"),
  2568. OnNext(240, " FoO "),
  2569. OnNext(270, "baR "),
  2570. OnNext(310, "foO "),
  2571. OnNext(350, " Baz "),
  2572. OnNext(360, " qux "),
  2573. OnNext(390, " bar"),
  2574. OnNext(420, " BAR "),
  2575. OnNext(470, "FOO "),
  2576. OnNext(480, "baz "),
  2577. OnNext(510, " bAZ "),
  2578. OnNext(530, " fOo "),
  2579. OnCompleted<string>(570),
  2580. OnNext(580, "error"),
  2581. OnCompleted<string>(600),
  2582. OnError<string>(650, new Exception())
  2583. );
  2584. var comparer = new GroupByComparer(scheduler, 400, ushort.MaxValue);
  2585. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2586. var outerSubscription = default(IDisposable);
  2587. var inners = new Dictionary<string, IObservable<string>>();
  2588. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2589. var res = new Dictionary<string, ITestableObserver<string>>();
  2590. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2591. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2592. {
  2593. var result = scheduler.CreateObserver<string>();
  2594. inners[group.Key] = group;
  2595. res[group.Key] = result;
  2596. innerSubscriptions[group.Key] = group.Subscribe(result);
  2597. }, _ => { }));
  2598. scheduler.ScheduleAbsolute(Disposed, () =>
  2599. {
  2600. outerSubscription.Dispose();
  2601. foreach (var d in innerSubscriptions.Values)
  2602. {
  2603. d.Dispose();
  2604. }
  2605. });
  2606. scheduler.Start();
  2607. Assert.Equal(4, inners.Count);
  2608. res["foo"].Messages.AssertEqual(
  2609. OnNext(220, "oof "),
  2610. OnNext(240, " OoF "),
  2611. OnNext(310, " Oof"),
  2612. OnCompleted<string>(310)
  2613. );
  2614. res["baR"].Messages.AssertEqual(
  2615. OnNext(270, " Rab"),
  2616. OnNext(390, "rab "),
  2617. OnError<string>(420, comparer.EqualsException)
  2618. );
  2619. res["Baz"].Messages.AssertEqual(
  2620. OnNext(350, " zaB "),
  2621. OnError<string>(420, comparer.EqualsException)
  2622. );
  2623. res["qux"].Messages.AssertEqual(
  2624. OnNext(360, " xuq "),
  2625. OnError<string>(420, comparer.EqualsException)
  2626. );
  2627. xs.Subscriptions.AssertEqual(
  2628. Subscribe(200, 420)
  2629. );
  2630. }
  2631. [Fact]
  2632. public void GroupByUntil_Capacity_Inner_Comparer_GetHashCodeThrow()
  2633. {
  2634. var scheduler = new TestScheduler();
  2635. var xs = scheduler.CreateHotObservable(
  2636. OnNext(90, "error"),
  2637. OnNext(110, "error"),
  2638. OnNext(130, "error"),
  2639. OnNext(220, " foo"),
  2640. OnNext(240, " FoO "),
  2641. OnNext(270, "baR "),
  2642. OnNext(310, "foO "),
  2643. OnNext(350, " Baz "),
  2644. OnNext(360, " qux "),
  2645. OnNext(390, " bar"),
  2646. OnNext(420, " BAR "),
  2647. OnNext(470, "FOO "),
  2648. OnNext(480, "baz "),
  2649. OnNext(510, " bAZ "),
  2650. OnNext(530, " fOo "),
  2651. OnCompleted<string>(570),
  2652. OnNext(580, "error"),
  2653. OnCompleted<string>(600),
  2654. OnError<string>(650, new Exception())
  2655. );
  2656. var comparer = new GroupByComparer(scheduler, ushort.MaxValue, 400);
  2657. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2658. var outerSubscription = default(IDisposable);
  2659. var inners = new Dictionary<string, IObservable<string>>();
  2660. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2661. var res = new Dictionary<string, ITestableObserver<string>>();
  2662. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2663. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2664. {
  2665. var result = scheduler.CreateObserver<string>();
  2666. inners[group.Key] = group;
  2667. res[group.Key] = result;
  2668. innerSubscriptions[group.Key] = group.Subscribe(result);
  2669. }, _ => { }));
  2670. scheduler.ScheduleAbsolute(Disposed, () =>
  2671. {
  2672. outerSubscription.Dispose();
  2673. foreach (var d in innerSubscriptions.Values)
  2674. {
  2675. d.Dispose();
  2676. }
  2677. });
  2678. scheduler.Start();
  2679. Assert.Equal(4, inners.Count);
  2680. res["foo"].Messages.AssertEqual(
  2681. OnNext(220, "oof "),
  2682. OnNext(240, " OoF "),
  2683. OnNext(310, " Oof"),
  2684. OnCompleted<string>(310)
  2685. );
  2686. res["baR"].Messages.AssertEqual(
  2687. OnNext(270, " Rab"),
  2688. OnNext(390, "rab "),
  2689. OnError<string>(420, comparer.HashCodeException)
  2690. );
  2691. res["Baz"].Messages.AssertEqual(
  2692. OnNext(350, " zaB "),
  2693. OnError<string>(420, comparer.HashCodeException)
  2694. );
  2695. res["qux"].Messages.AssertEqual(
  2696. OnNext(360, " xuq "),
  2697. OnError<string>(420, comparer.HashCodeException)
  2698. );
  2699. xs.Subscriptions.AssertEqual(
  2700. Subscribe(200, 420)
  2701. );
  2702. }
  2703. [Fact]
  2704. public void GroupByUntil_Capacity_Outer_Independence()
  2705. {
  2706. var scheduler = new TestScheduler();
  2707. var xs = scheduler.CreateHotObservable(
  2708. OnNext(90, "error"),
  2709. OnNext(110, "error"),
  2710. OnNext(130, "error"),
  2711. OnNext(220, " foo"),
  2712. OnNext(240, " FoO "),
  2713. OnNext(270, "baR "),
  2714. OnNext(310, "foO "),
  2715. OnNext(350, " Baz "),
  2716. OnNext(360, " qux "),
  2717. OnNext(390, " bar"),
  2718. OnNext(420, " BAR "),
  2719. OnNext(470, "FOO "),
  2720. OnNext(480, "baz "),
  2721. OnNext(510, " bAZ "),
  2722. OnNext(530, " fOo "),
  2723. OnCompleted<string>(570),
  2724. OnNext(580, "error"),
  2725. OnCompleted<string>(600),
  2726. OnError<string>(650, new Exception())
  2727. );
  2728. var comparer = new GroupByComparer(scheduler);
  2729. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2730. var outerSubscription = default(IDisposable);
  2731. var inners = new Dictionary<string, IObservable<string>>();
  2732. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2733. var res = new Dictionary<string, ITestableObserver<string>>();
  2734. var outerResults = scheduler.CreateObserver<string>();
  2735. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2736. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2737. {
  2738. outerResults.OnNext(group.Key);
  2739. var result = scheduler.CreateObserver<string>();
  2740. inners[group.Key] = group;
  2741. res[group.Key] = result;
  2742. innerSubscriptions[group.Key] = group.Subscribe(result);
  2743. }, outerResults.OnError, outerResults.OnCompleted));
  2744. scheduler.ScheduleAbsolute(Disposed, () =>
  2745. {
  2746. outerSubscription.Dispose();
  2747. foreach (var d in innerSubscriptions.Values)
  2748. {
  2749. d.Dispose();
  2750. }
  2751. });
  2752. scheduler.ScheduleAbsolute(320, () => outerSubscription.Dispose());
  2753. scheduler.Start();
  2754. Assert.Equal(2, inners.Count);
  2755. outerResults.Messages.AssertEqual(
  2756. OnNext(220, "foo"),
  2757. OnNext(270, "baR")
  2758. );
  2759. res["foo"].Messages.AssertEqual(
  2760. OnNext(220, "oof "),
  2761. OnNext(240, " OoF "),
  2762. OnNext(310, " Oof"),
  2763. OnCompleted<string>(310)
  2764. );
  2765. res["baR"].Messages.AssertEqual(
  2766. OnNext(270, " Rab"),
  2767. OnNext(390, "rab "),
  2768. OnNext(420, " RAB "),
  2769. OnCompleted<string>(420)
  2770. );
  2771. xs.Subscriptions.AssertEqual(
  2772. Subscribe(200, 420)
  2773. );
  2774. }
  2775. [Fact]
  2776. public void GroupByUntil_Capacity_Inner_Independence()
  2777. {
  2778. var scheduler = new TestScheduler();
  2779. var xs = scheduler.CreateHotObservable(
  2780. OnNext(90, "error"),
  2781. OnNext(110, "error"),
  2782. OnNext(130, "error"),
  2783. OnNext(220, " foo"),
  2784. OnNext(240, " FoO "),
  2785. OnNext(270, "baR "),
  2786. OnNext(310, "foO "),
  2787. OnNext(350, " Baz "),
  2788. OnNext(360, " qux "),
  2789. OnNext(390, " bar"),
  2790. OnNext(420, " BAR "),
  2791. OnNext(470, "FOO "),
  2792. OnNext(480, "baz "),
  2793. OnNext(510, " bAZ "),
  2794. OnNext(530, " fOo "),
  2795. OnCompleted<string>(570),
  2796. OnNext(580, "error"),
  2797. OnCompleted<string>(600),
  2798. OnError<string>(650, new Exception())
  2799. );
  2800. var comparer = new GroupByComparer(scheduler);
  2801. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2802. var outerSubscription = default(IDisposable);
  2803. var inners = new Dictionary<string, IObservable<string>>();
  2804. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2805. var res = new Dictionary<string, ITestableObserver<string>>();
  2806. var outerResults = scheduler.CreateObserver<string>();
  2807. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2808. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2809. {
  2810. outerResults.OnNext(group.Key);
  2811. var result = scheduler.CreateObserver<string>();
  2812. inners[group.Key] = group;
  2813. res[group.Key] = result;
  2814. innerSubscriptions[group.Key] = group.Subscribe(result);
  2815. }, outerResults.OnError, outerResults.OnCompleted));
  2816. scheduler.ScheduleAbsolute(Disposed, () =>
  2817. {
  2818. outerSubscription.Dispose();
  2819. foreach (var d in innerSubscriptions.Values)
  2820. {
  2821. d.Dispose();
  2822. }
  2823. });
  2824. scheduler.ScheduleAbsolute(320, () => innerSubscriptions["foo"].Dispose());
  2825. scheduler.Start();
  2826. Assert.Equal(5, inners.Count);
  2827. res["foo"].Messages.AssertEqual(
  2828. OnNext(220, "oof "),
  2829. OnNext(240, " OoF "),
  2830. OnNext(310, " Oof"),
  2831. OnCompleted<string>(310)
  2832. );
  2833. res["baR"].Messages.AssertEqual(
  2834. OnNext(270, " Rab"),
  2835. OnNext(390, "rab "),
  2836. OnNext(420, " RAB "),
  2837. OnCompleted<string>(420)
  2838. );
  2839. res["Baz"].Messages.AssertEqual(
  2840. OnNext(350, " zaB "),
  2841. OnNext(480, " zab"),
  2842. OnNext(510, " ZAb "),
  2843. OnCompleted<string>(510)
  2844. );
  2845. res["qux"].Messages.AssertEqual(
  2846. OnNext(360, " xuq "),
  2847. OnCompleted<string>(570)
  2848. );
  2849. res["FOO"].Messages.AssertEqual(
  2850. OnNext(470, " OOF"),
  2851. OnNext(530, " oOf "),
  2852. OnCompleted<string>(570)
  2853. );
  2854. xs.Subscriptions.AssertEqual(
  2855. Subscribe(200, 570)
  2856. );
  2857. }
  2858. [Fact]
  2859. public void GroupByUntil_Capacity_Inner_Multiple_Independence()
  2860. {
  2861. var scheduler = new TestScheduler();
  2862. var xs = scheduler.CreateHotObservable(
  2863. OnNext(90, "error"),
  2864. OnNext(110, "error"),
  2865. OnNext(130, "error"),
  2866. OnNext(220, " foo"),
  2867. OnNext(240, " FoO "),
  2868. OnNext(270, "baR "),
  2869. OnNext(310, "foO "),
  2870. OnNext(350, " Baz "),
  2871. OnNext(360, " qux "),
  2872. OnNext(390, " bar"),
  2873. OnNext(420, " BAR "),
  2874. OnNext(470, "FOO "),
  2875. OnNext(480, "baz "),
  2876. OnNext(510, " bAZ "),
  2877. OnNext(530, " fOo "),
  2878. OnCompleted<string>(570),
  2879. OnNext(580, "error"),
  2880. OnCompleted<string>(600),
  2881. OnError<string>(650, new Exception())
  2882. );
  2883. var comparer = new GroupByComparer(scheduler);
  2884. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2885. var outerSubscription = default(IDisposable);
  2886. var inners = new Dictionary<string, IObservable<string>>();
  2887. var innerSubscriptions = new Dictionary<string, IDisposable>();
  2888. var res = new Dictionary<string, ITestableObserver<string>>();
  2889. var outerResults = scheduler.CreateObserver<string>();
  2890. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), x => Reverse(x), g => g.Skip(2), _groupByUntilCapacity, comparer));
  2891. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2892. {
  2893. outerResults.OnNext(group.Key);
  2894. var result = scheduler.CreateObserver<string>();
  2895. inners[group.Key] = group;
  2896. res[group.Key] = result;
  2897. innerSubscriptions[group.Key] = group.Subscribe(result);
  2898. }, outerResults.OnError, outerResults.OnCompleted));
  2899. scheduler.ScheduleAbsolute(Disposed, () =>
  2900. {
  2901. outerSubscription.Dispose();
  2902. foreach (var d in innerSubscriptions.Values)
  2903. {
  2904. d.Dispose();
  2905. }
  2906. });
  2907. scheduler.ScheduleAbsolute(320, () => innerSubscriptions["foo"].Dispose());
  2908. scheduler.ScheduleAbsolute(280, () => innerSubscriptions["baR"].Dispose());
  2909. scheduler.ScheduleAbsolute(355, () => innerSubscriptions["Baz"].Dispose());
  2910. scheduler.ScheduleAbsolute(400, () => innerSubscriptions["qux"].Dispose());
  2911. scheduler.Start();
  2912. Assert.Equal(5, inners.Count);
  2913. res["foo"].Messages.AssertEqual(
  2914. OnNext(220, "oof "),
  2915. OnNext(240, " OoF "),
  2916. OnNext(310, " Oof"),
  2917. OnCompleted<string>(310)
  2918. );
  2919. res["baR"].Messages.AssertEqual(
  2920. OnNext(270, " Rab")
  2921. );
  2922. res["Baz"].Messages.AssertEqual(
  2923. OnNext(350, " zaB ")
  2924. );
  2925. res["qux"].Messages.AssertEqual(
  2926. OnNext(360, " xuq ")
  2927. );
  2928. res["FOO"].Messages.AssertEqual(
  2929. OnNext(470, " OOF"),
  2930. OnNext(530, " oOf "),
  2931. OnCompleted<string>(570)
  2932. );
  2933. xs.Subscriptions.AssertEqual(
  2934. Subscribe(200, 570)
  2935. );
  2936. }
  2937. [Fact]
  2938. public void GroupByUntil_Capacity_Inner_Escape_Complete()
  2939. {
  2940. var scheduler = new TestScheduler();
  2941. var xs = scheduler.CreateHotObservable(
  2942. OnNext(220, " foo"),
  2943. OnNext(240, " FoO "),
  2944. OnNext(310, "foO "),
  2945. OnNext(470, "FOO "),
  2946. OnNext(530, " fOo "),
  2947. OnCompleted<string>(570)
  2948. );
  2949. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2950. var outerSubscription = default(IDisposable);
  2951. var inner = default(IObservable<string>);
  2952. var innerSubscription = default(IDisposable);
  2953. var res = scheduler.CreateObserver<string>();
  2954. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), g => g.Skip(2), _groupByUntilCapacity));
  2955. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2956. {
  2957. inner = group;
  2958. }));
  2959. scheduler.ScheduleAbsolute(600, () => innerSubscription = inner.Subscribe(res));
  2960. scheduler.ScheduleAbsolute(Disposed, () =>
  2961. {
  2962. outerSubscription.Dispose();
  2963. innerSubscription.Dispose();
  2964. });
  2965. scheduler.Start();
  2966. xs.Subscriptions.AssertEqual(
  2967. Subscribe(200, 570)
  2968. );
  2969. res.Messages.AssertEqual(
  2970. OnCompleted<string>(600)
  2971. );
  2972. }
  2973. [Fact]
  2974. public void GroupByUntil_Capacity_Inner_Escape_Error()
  2975. {
  2976. var scheduler = new TestScheduler();
  2977. var ex = new Exception();
  2978. var xs = scheduler.CreateHotObservable(
  2979. OnNext(220, " foo"),
  2980. OnNext(240, " FoO "),
  2981. OnNext(310, "foO "),
  2982. OnNext(470, "FOO "),
  2983. OnNext(530, " fOo "),
  2984. OnError<string>(570, ex)
  2985. );
  2986. var outer = default(IObservable<IGroupedObservable<string, string>>);
  2987. var outerSubscription = default(IDisposable);
  2988. var inner = default(IObservable<string>);
  2989. var innerSubscription = default(IDisposable);
  2990. var res = scheduler.CreateObserver<string>();
  2991. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), g => g.Skip(2), _groupByUntilCapacity));
  2992. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  2993. {
  2994. inner = group;
  2995. }, _ => { }));
  2996. scheduler.ScheduleAbsolute(600, () => innerSubscription = inner.Subscribe(res));
  2997. scheduler.ScheduleAbsolute(Disposed, () =>
  2998. {
  2999. outerSubscription.Dispose();
  3000. innerSubscription.Dispose();
  3001. });
  3002. scheduler.Start();
  3003. xs.Subscriptions.AssertEqual(
  3004. Subscribe(200, 570)
  3005. );
  3006. res.Messages.AssertEqual(
  3007. OnError<string>(600, ex)
  3008. );
  3009. }
  3010. [Fact]
  3011. public void GroupByUntil_Capacity_Inner_Escape_Dispose()
  3012. {
  3013. var scheduler = new TestScheduler();
  3014. var xs = scheduler.CreateHotObservable(
  3015. OnNext(220, " foo"),
  3016. OnNext(240, " FoO "),
  3017. OnNext(310, "foO "),
  3018. OnNext(470, "FOO "),
  3019. OnNext(530, " fOo "),
  3020. OnError<string>(570, new Exception())
  3021. );
  3022. var outer = default(IObservable<IGroupedObservable<string, string>>);
  3023. var outerSubscription = default(IDisposable);
  3024. var inner = default(IObservable<string>);
  3025. var innerSubscription = default(IDisposable);
  3026. var res = scheduler.CreateObserver<string>();
  3027. scheduler.ScheduleAbsolute(Created, () => outer = xs.GroupByUntil(x => x.Trim(), g => g.Skip(2), _groupByUntilCapacity));
  3028. scheduler.ScheduleAbsolute(Subscribed, () => outerSubscription = outer.Subscribe(group =>
  3029. {
  3030. inner = group;
  3031. }));
  3032. scheduler.ScheduleAbsolute(290, () => outerSubscription.Dispose());
  3033. scheduler.ScheduleAbsolute(600, () => innerSubscription = inner.Subscribe(res));
  3034. scheduler.ScheduleAbsolute(Disposed, () =>
  3035. {
  3036. innerSubscription.Dispose();
  3037. });
  3038. scheduler.Start();
  3039. xs.Subscriptions.AssertEqual(
  3040. Subscribe(200, 290)
  3041. );
  3042. res.Messages.AssertEqual(
  3043. );
  3044. }
  3045. [Fact]
  3046. public void GroupByUntil_Capacity_Default()
  3047. {
  3048. var scheduler = new TestScheduler();
  3049. var keyInvoked = 0;
  3050. var eleInvoked = 0;
  3051. var xs = scheduler.CreateHotObservable(
  3052. OnNext(90, "error"),
  3053. OnNext(110, "error"),
  3054. OnNext(130, "error"),
  3055. OnNext(220, " foo"),
  3056. OnNext(240, " FoO "),
  3057. OnNext(270, "baR "),
  3058. OnNext(310, "foO "),
  3059. OnNext(350, " Baz "),
  3060. OnNext(360, " qux "),
  3061. OnNext(390, " bar"),
  3062. OnNext(420, " BAR "),
  3063. OnNext(470, "FOO "),
  3064. OnNext(480, "baz "),
  3065. OnNext(510, " bAZ "),
  3066. OnNext(530, " fOo "),
  3067. OnCompleted<string>(570),
  3068. OnNext(580, "error"),
  3069. OnCompleted<string>(600),
  3070. OnError<string>(650, new Exception())
  3071. );
  3072. var res = scheduler.Start(() =>
  3073. xs.GroupByUntil(
  3074. x =>
  3075. {
  3076. keyInvoked++;
  3077. return x.Trim().ToLower();
  3078. },
  3079. x =>
  3080. {
  3081. eleInvoked++;
  3082. return Reverse(x);
  3083. },
  3084. g => g.Skip(2),
  3085. _groupByUntilCapacity
  3086. ).Select(g => g.Key)
  3087. );
  3088. res.Messages.AssertEqual(
  3089. OnNext(220, "foo"),
  3090. OnNext(270, "bar"),
  3091. OnNext(350, "baz"),
  3092. OnNext(360, "qux"),
  3093. OnNext(470, "foo"),
  3094. OnCompleted<string>(570)
  3095. );
  3096. xs.Subscriptions.AssertEqual(
  3097. Subscribe(200, 570)
  3098. );
  3099. Assert.Equal(12, keyInvoked);
  3100. Assert.Equal(12, eleInvoked);
  3101. }
  3102. [Fact]
  3103. public void GroupByUntil_Capacity_DurationSelector_Throws()
  3104. {
  3105. var scheduler = new TestScheduler();
  3106. var xs = scheduler.CreateHotObservable(
  3107. OnNext(210, "foo")
  3108. );
  3109. var ex = new Exception();
  3110. var res = scheduler.Start(() =>
  3111. xs.GroupByUntil<string, string, string>(x => x, g => { throw ex; }, _groupByUntilCapacity)
  3112. );
  3113. res.Messages.AssertEqual(
  3114. OnError<IGroupedObservable<string, string>>(210, ex)
  3115. );
  3116. xs.Subscriptions.AssertEqual(
  3117. Subscribe(200, 210)
  3118. );
  3119. }
  3120. [Fact]
  3121. public void GroupByUntil_Capacity_NullKeys_Simple_Never()
  3122. {
  3123. var scheduler = new TestScheduler();
  3124. var xs = scheduler.CreateHotObservable(
  3125. OnNext(220, "bar"),
  3126. OnNext(240, "foo"),
  3127. OnNext(310, "qux"),
  3128. OnNext(470, "baz"),
  3129. OnCompleted<string>(500)
  3130. );
  3131. var res = scheduler.Start(() => xs.GroupByUntil(x => x[0] == 'b' ? null : x.ToUpper(), g => Observable.Never<Unit>(), _groupByUntilCapacity).SelectMany(g => g, (g, x) => (g.Key ?? "(null)") + x));
  3132. res.Messages.AssertEqual(
  3133. OnNext(220, "(null)bar"),
  3134. OnNext(240, "FOOfoo"),
  3135. OnNext(310, "QUXqux"),
  3136. OnNext(470, "(null)baz"),
  3137. OnCompleted<string>(500)
  3138. );
  3139. xs.Subscriptions.AssertEqual(
  3140. Subscribe(200, 500)
  3141. );
  3142. }
  3143. [Fact]
  3144. public void GroupByUntil_Capacity_NullKeys_Simple_Expire1()
  3145. {
  3146. var scheduler = new TestScheduler();
  3147. var xs = scheduler.CreateHotObservable(
  3148. OnNext(220, "bar"),
  3149. OnNext(240, "foo"),
  3150. OnNext(310, "qux"),
  3151. OnNext(470, "baz"),
  3152. OnCompleted<string>(500)
  3153. );
  3154. var n = 0;
  3155. var res = scheduler.Start(() => xs.GroupByUntil(x => x[0] == 'b' ? null : x.ToUpper(), g => { if (g.Key == null) { n++; } return Observable.Timer(TimeSpan.FromTicks(50), scheduler); }, _groupByUntilCapacity).SelectMany(g => g, (g, x) => (g.Key ?? "(null)") + x));
  3156. Assert.Equal(2, n);
  3157. res.Messages.AssertEqual(
  3158. OnNext(220, "(null)bar"),
  3159. OnNext(240, "FOOfoo"),
  3160. OnNext(310, "QUXqux"),
  3161. OnNext(470, "(null)baz"),
  3162. OnCompleted<string>(500)
  3163. );
  3164. xs.Subscriptions.AssertEqual(
  3165. Subscribe(200, 500)
  3166. );
  3167. }
  3168. [Fact]
  3169. public void GroupByUntil_Capacity_NullKeys_Simple_Expire2()
  3170. {
  3171. var scheduler = new TestScheduler();
  3172. var xs = scheduler.CreateHotObservable(
  3173. OnNext(220, "bar"),
  3174. OnNext(240, "foo"),
  3175. OnNext(310, "qux"),
  3176. OnNext(470, "baz"),
  3177. OnCompleted<string>(500)
  3178. );
  3179. var n = 0;
  3180. var res = scheduler.Start(() => xs.GroupByUntil(x => x[0] == 'b' ? null : x.ToUpper(), g => { if (g.Key == null) { n++; } return Observable.Timer(TimeSpan.FromTicks(50), scheduler).IgnoreElements(); }, _groupByUntilCapacity).SelectMany(g => g, (g, x) => (g.Key ?? "(null)") + x));
  3181. Assert.Equal(2, n);
  3182. res.Messages.AssertEqual(
  3183. OnNext(220, "(null)bar"),
  3184. OnNext(240, "FOOfoo"),
  3185. OnNext(310, "QUXqux"),
  3186. OnNext(470, "(null)baz"),
  3187. OnCompleted<string>(500)
  3188. );
  3189. xs.Subscriptions.AssertEqual(
  3190. Subscribe(200, 500)
  3191. );
  3192. }
  3193. [Fact]
  3194. public void GroupByUntil_Capacity_NullKeys_Error()
  3195. {
  3196. var scheduler = new TestScheduler();
  3197. var ex = new Exception();
  3198. var xs = scheduler.CreateHotObservable(
  3199. OnNext(220, "bar"),
  3200. OnNext(240, "foo"),
  3201. OnNext(310, "qux"),
  3202. OnNext(470, "baz"),
  3203. OnError<string>(500, ex)
  3204. );
  3205. var nullGroup = scheduler.CreateObserver<string>();
  3206. var err = default(Exception);
  3207. scheduler.ScheduleAbsolute(200, () => xs.GroupByUntil(x => x[0] == 'b' ? null : x.ToUpper(), g => Observable.Never<Unit>(), _groupByUntilCapacity).Where(g => g.Key == null).Subscribe(g => g.Subscribe(nullGroup), ex_ => err = ex_));
  3208. scheduler.Start();
  3209. Assert.Same(ex, err);
  3210. nullGroup.Messages.AssertEqual(
  3211. OnNext(220, "bar"),
  3212. OnNext(470, "baz"),
  3213. OnError<string>(500, ex)
  3214. );
  3215. xs.Subscriptions.AssertEqual(
  3216. Subscribe(200, 500)
  3217. );
  3218. }
  3219. #endregion
  3220. private static string Reverse(string s)
  3221. {
  3222. var sb = new StringBuilder();
  3223. for (var i = s.Length - 1; i >= 0; i--)
  3224. {
  3225. sb.Append(s[i]);
  3226. }
  3227. return sb.ToString();
  3228. }
  3229. }
  3230. }