session-runner.test.ts 188 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355235623572358235923602361236223632364236523662367236823692370237123722373237423752376237723782379238023812382238323842385238623872388238923902391239223932394239523962397239823992400240124022403240424052406240724082409241024112412241324142415241624172418241924202421242224232424242524262427242824292430243124322433243424352436243724382439244024412442244324442445244624472448244924502451245224532454245524562457245824592460246124622463246424652466246724682469247024712472247324742475247624772478247924802481248224832484248524862487248824892490249124922493249424952496249724982499250025012502250325042505250625072508250925102511251225132514251525162517251825192520252125222523252425252526252725282529253025312532253325342535253625372538253925402541254225432544254525462547254825492550255125522553255425552556255725582559256025612562256325642565256625672568256925702571257225732574257525762577257825792580258125822583258425852586258725882589259025912592259325942595259625972598259926002601260226032604260526062607260826092610261126122613261426152616261726182619262026212622262326242625262626272628262926302631263226332634263526362637263826392640264126422643264426452646264726482649265026512652265326542655265626572658265926602661266226632664266526662667266826692670267126722673267426752676267726782679268026812682268326842685268626872688268926902691269226932694269526962697269826992700270127022703270427052706270727082709271027112712271327142715271627172718271927202721272227232724272527262727272827292730273127322733273427352736273727382739274027412742274327442745274627472748274927502751275227532754275527562757275827592760276127622763276427652766276727682769277027712772277327742775277627772778277927802781278227832784278527862787278827892790279127922793279427952796279727982799280028012802280328042805280628072808280928102811281228132814281528162817281828192820282128222823282428252826282728282829283028312832283328342835283628372838283928402841284228432844284528462847284828492850285128522853285428552856285728582859286028612862286328642865286628672868286928702871287228732874287528762877287828792880288128822883288428852886288728882889289028912892289328942895289628972898289929002901290229032904290529062907290829092910291129122913291429152916291729182919292029212922292329242925292629272928292929302931293229332934293529362937293829392940294129422943294429452946294729482949295029512952295329542955295629572958295929602961296229632964296529662967296829692970297129722973297429752976297729782979298029812982298329842985298629872988298929902991299229932994299529962997299829993000300130023003300430053006300730083009301030113012301330143015301630173018301930203021302230233024302530263027302830293030303130323033303430353036303730383039304030413042304330443045304630473048304930503051305230533054305530563057305830593060306130623063306430653066306730683069307030713072307330743075307630773078307930803081308230833084308530863087308830893090309130923093309430953096309730983099310031013102310331043105310631073108310931103111311231133114311531163117311831193120312131223123312431253126312731283129313031313132313331343135313631373138313931403141314231433144314531463147314831493150315131523153315431553156315731583159316031613162316331643165316631673168316931703171317231733174317531763177317831793180318131823183318431853186318731883189319031913192319331943195319631973198319932003201320232033204320532063207320832093210321132123213321432153216321732183219322032213222322332243225322632273228322932303231323232333234323532363237323832393240324132423243324432453246324732483249325032513252325332543255325632573258325932603261326232633264326532663267326832693270327132723273327432753276327732783279328032813282328332843285328632873288328932903291329232933294329532963297329832993300330133023303330433053306330733083309331033113312331333143315331633173318331933203321332233233324332533263327332833293330333133323333333433353336333733383339334033413342334333443345334633473348334933503351335233533354335533563357335833593360336133623363336433653366336733683369337033713372337333743375337633773378337933803381338233833384338533863387338833893390339133923393339433953396339733983399340034013402340334043405340634073408340934103411341234133414341534163417341834193420342134223423342434253426342734283429343034313432343334343435343634373438343934403441344234433444344534463447344834493450345134523453345434553456345734583459346034613462346334643465346634673468346934703471347234733474347534763477347834793480348134823483348434853486348734883489349034913492349334943495349634973498349935003501350235033504350535063507350835093510351135123513351435153516351735183519352035213522352335243525352635273528352935303531353235333534353535363537353835393540354135423543354435453546354735483549355035513552355335543555355635573558355935603561356235633564356535663567356835693570357135723573357435753576357735783579358035813582358335843585358635873588358935903591359235933594359535963597359835993600360136023603360436053606360736083609361036113612361336143615361636173618361936203621362236233624362536263627362836293630363136323633363436353636363736383639364036413642364336443645364636473648364936503651365236533654365536563657365836593660366136623663366436653666366736683669367036713672367336743675367636773678367936803681368236833684368536863687368836893690369136923693369436953696369736983699370037013702370337043705370637073708370937103711371237133714371537163717371837193720372137223723372437253726372737283729373037313732373337343735373637373738373937403741374237433744374537463747374837493750375137523753375437553756375737583759376037613762376337643765376637673768376937703771377237733774377537763777377837793780378137823783378437853786378737883789379037913792379337943795379637973798379938003801380238033804380538063807380838093810381138123813381438153816381738183819382038213822382338243825382638273828382938303831383238333834383538363837383838393840384138423843384438453846384738483849385038513852385338543855385638573858385938603861386238633864386538663867386838693870387138723873387438753876387738783879388038813882388338843885388638873888388938903891389238933894389538963897389838993900390139023903390439053906390739083909391039113912391339143915391639173918391939203921392239233924392539263927392839293930393139323933393439353936393739383939394039413942394339443945394639473948394939503951395239533954395539563957395839593960396139623963396439653966396739683969397039713972397339743975397639773978397939803981398239833984398539863987398839893990399139923993399439953996399739983999400040014002400340044005400640074008400940104011401240134014401540164017401840194020402140224023402440254026402740284029403040314032403340344035403640374038403940404041404240434044404540464047404840494050405140524053405440554056405740584059406040614062406340644065406640674068406940704071407240734074407540764077407840794080408140824083408440854086408740884089409040914092409340944095409640974098409941004101410241034104410541064107410841094110411141124113411441154116411741184119412041214122412341244125412641274128412941304131413241334134413541364137413841394140414141424143414441454146414741484149415041514152415341544155415641574158415941604161416241634164416541664167416841694170417141724173417441754176417741784179418041814182418341844185418641874188418941904191419241934194419541964197419841994200420142024203420442054206420742084209421042114212421342144215421642174218421942204221422242234224422542264227422842294230423142324233423442354236423742384239424042414242424342444245424642474248424942504251425242534254425542564257425842594260426142624263426442654266426742684269427042714272427342744275427642774278427942804281428242834284428542864287428842894290429142924293429442954296429742984299430043014302430343044305430643074308430943104311431243134314431543164317431843194320432143224323432443254326432743284329433043314332433343344335433643374338433943404341434243434344434543464347434843494350435143524353435443554356435743584359436043614362436343644365436643674368436943704371437243734374437543764377437843794380438143824383438443854386438743884389439043914392439343944395439643974398439944004401440244034404440544064407440844094410441144124413441444154416441744184419442044214422442344244425442644274428442944304431443244334434443544364437443844394440444144424443444444454446444744484449445044514452445344544455445644574458445944604461446244634464446544664467446844694470447144724473447444754476447744784479448044814482448344844485448644874488448944904491449244934494449544964497449844994500450145024503450445054506450745084509451045114512451345144515451645174518451945204521452245234524452545264527452845294530453145324533453445354536453745384539454045414542454345444545454645474548454945504551455245534554455545564557455845594560456145624563456445654566456745684569457045714572457345744575457645774578457945804581458245834584458545864587458845894590459145924593459445954596459745984599460046014602460346044605460646074608460946104611461246134614461546164617461846194620462146224623462446254626462746284629463046314632463346344635463646374638463946404641464246434644464546464647464846494650465146524653465446554656465746584659466046614662466346644665466646674668466946704671467246734674467546764677467846794680468146824683468446854686468746884689469046914692469346944695469646974698469947004701470247034704470547064707470847094710471147124713471447154716471747184719472047214722472347244725472647274728472947304731473247334734473547364737473847394740474147424743474447454746474747484749475047514752475347544755475647574758475947604761476247634764476547664767476847694770477147724773477447754776477747784779478047814782478347844785478647874788478947904791479247934794479547964797479847994800480148024803480448054806480748084809481048114812481348144815481648174818481948204821482248234824482548264827482848294830483148324833483448354836483748384839484048414842484348444845484648474848484948504851485248534854485548564857485848594860486148624863486448654866486748684869487048714872487348744875487648774878487948804881488248834884488548864887488848894890489148924893489448954896489748984899490049014902490349044905490649074908490949104911491249134914491549164917491849194920492149224923492449254926492749284929493049314932493349344935493649374938493949404941494249434944494549464947494849494950495149524953495449554956495749584959496049614962496349644965496649674968496949704971497249734974497549764977497849794980498149824983498449854986498749884989499049914992499349944995499649974998499950005001500250035004500550065007500850095010501150125013501450155016501750185019502050215022502350245025502650275028502950305031503250335034503550365037503850395040504150425043504450455046504750485049505050515052505350545055505650575058505950605061506250635064506550665067506850695070
  1. import { describe, expect, test } from "bun:test"
  2. import {
  3. LLMClient,
  4. LLMError,
  5. LLMEvent,
  6. Message,
  7. Model,
  8. SystemPart,
  9. ToolFailure,
  10. TransportReason,
  11. InvalidProviderOutputReason,
  12. InvalidRequestReason,
  13. RateLimitReason,
  14. type LLMClientShape,
  15. type LLMRequest,
  16. } from "@opencode-ai/ai"
  17. import * as OpenAIChat from "@opencode-ai/ai/protocols/openai-chat"
  18. import { Catalog } from "@opencode-ai/core/catalog"
  19. import { Database } from "@opencode-ai/core/database/database"
  20. import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
  21. import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
  22. import { LayerNodePlatform } from "@opencode-ai/core/effect/app-node-platform"
  23. import { LayerNode } from "@opencode-ai/util/effect/layer-node"
  24. import { EventV2 } from "@opencode-ai/core/event"
  25. import { App } from "@opencode-ai/core/app"
  26. import { PermissionV2 } from "@opencode-ai/core/permission"
  27. import { EventTable } from "@opencode-ai/core/event/sql"
  28. import { Project } from "@opencode-ai/core/project"
  29. import { ProjectTable } from "@opencode-ai/core/project/sql"
  30. import { Form } from "@opencode-ai/core/form"
  31. import { AbsolutePath } from "@opencode-ai/core/schema"
  32. import { SessionV2 } from "@opencode-ai/core/session"
  33. import { Snapshot } from "@opencode-ai/core/snapshot"
  34. import { SessionEvent } from "@opencode-ai/core/session/event"
  35. import { SessionPending } from "@opencode-ai/core/session/pending"
  36. import { SessionMessage } from "@opencode-ai/core/session/message"
  37. import { Money } from "@opencode-ai/schema/money"
  38. import { SessionProjector } from "@opencode-ai/core/session/projector"
  39. import { SessionExecution } from "@opencode-ai/core/session/execution"
  40. import { SessionRunCoordinator } from "@opencode-ai/core/session/run-coordinator"
  41. import { SessionRunner } from "@opencode-ai/core/session/runner"
  42. import * as SessionRunnerLLM from "@opencode-ai/core/session/runner/llm"
  43. import { SessionRunnerModel } from "@opencode-ai/core/session/runner/model"
  44. import { SessionUsage } from "@opencode-ai/core/session/usage"
  45. import { ToolRegistry } from "@opencode-ai/core/tool/registry"
  46. import { CodeMode } from "@opencode-ai/core/codemode"
  47. import { PluginSupervisor } from "@opencode-ai/core/plugin/supervisor"
  48. import { PluginHooks } from "@opencode-ai/core/plugin/hooks"
  49. import { SystemPromptPlugin } from "@opencode-ai/core/plugin/system-prompt"
  50. import { QuestionTool } from "@opencode-ai/core/tool/question"
  51. import { ToolOutputStore } from "@opencode-ai/core/tool-output-store"
  52. import { AgentV2 } from "@opencode-ai/core/agent"
  53. import { Config } from "@opencode-ai/core/config"
  54. import { ConfigCompaction } from "@opencode-ai/core/config/compaction"
  55. import { Tool } from "@opencode-ai/core/tool/tool"
  56. import { ToolHooks } from "@opencode-ai/core/tool/hooks"
  57. import {
  58. InstructionStateTable,
  59. SessionPendingTable,
  60. SessionMessageTable,
  61. SessionTable,
  62. } from "@opencode-ai/core/session/sql"
  63. import { InstructionEntry } from "@opencode-ai/core/session/instruction-entry"
  64. import { SessionStore } from "@opencode-ai/core/session/store"
  65. import { Instructions } from "@opencode-ai/core/instructions"
  66. import { InstructionBuiltIns } from "@opencode-ai/core/instructions/builtins"
  67. import { InstructionDiscovery } from "@opencode-ai/core/instruction-discovery"
  68. import { SkillInstructions } from "@opencode-ai/core/skill/instructions"
  69. import { ReferenceInstructions } from "@opencode-ai/core/reference/instructions"
  70. import { McpInstructions } from "@opencode-ai/core/mcp/instructions"
  71. import { ModelV2 } from "@opencode-ai/core/model"
  72. import { Location } from "@opencode-ai/core/location"
  73. import { ProviderV2 } from "@opencode-ai/core/provider"
  74. import { Cause, DateTime, Deferred, Effect, Exit, Fiber, Layer, Schema, Scope, Stream } from "effect"
  75. import { TestClock } from "effect/testing"
  76. import { asc, eq } from "drizzle-orm"
  77. import { testEffect } from "./lib/effect"
  78. import { agentHost, catalogHost, host } from "./plugin/host"
  79. import PROMPT_DEFAULT from "../src/session/runner/prompt/base.txt"
  80. const requests: LLMRequest[] = []
  81. let response: LLMEvent[] = []
  82. let responses: LLMEvent[][] | undefined
  83. let responseStream: Stream.Stream<LLMEvent, LLMError> | undefined
  84. let responseStreams: Stream.Stream<LLMEvent, LLMError>[] | undefined
  85. let streamGate: Deferred.Deferred<void> | undefined
  86. let streamStarted: Deferred.Deferred<void> | undefined
  87. let streamFailure: LLMError | undefined
  88. let toolExecutionGate: Deferred.Deferred<void> | undefined
  89. let toolExecutionsStarted: Deferred.Deferred<void> | undefined
  90. let toolExecutionsReady = 5
  91. let activeToolExecutions = 0
  92. let maxActiveToolExecutions = 0
  93. const client = Layer.succeed(
  94. LLMClient.Service,
  95. LLMClient.Service.of({
  96. prepare: () => Effect.die("unused"),
  97. stream: ((request: LLMRequest) => {
  98. requests.push(request)
  99. if (responseStreams) return responseStreams.shift() ?? Stream.empty
  100. if (responseStream) {
  101. const stream = responseStream
  102. responseStream = undefined
  103. return stream
  104. }
  105. const events = streamFailure
  106. ? Stream.fail(streamFailure)
  107. : Stream.fromIterable(responses === undefined ? response : (responses.shift() ?? []))
  108. if (!streamGate) return events
  109. return Stream.unwrap(
  110. (streamStarted ? Deferred.succeed(streamStarted, undefined) : Effect.void).pipe(
  111. Effect.andThen(Deferred.await(streamGate)),
  112. Effect.as(events),
  113. ),
  114. )
  115. }) as unknown as LLMClientShape["stream"],
  116. generate: () => Effect.die("unused"),
  117. }),
  118. )
  119. const reply = {
  120. stop: () => [
  121. LLMEvent.stepStart({ index: 0 }),
  122. LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
  123. LLMEvent.finish({ reason: { normalized: "stop" } }),
  124. ],
  125. text: (text: string, id: string) => fragmentFixture("text", id, [text]).completeEvents,
  126. textWithUsage: (text: string, id: string, inputTokens: number) =>
  127. fragmentFixture("text", id, [text]).completeEvents.map((event) =>
  128. LLMEvent.is.stepFinish(event)
  129. ? LLMEvent.stepFinish({
  130. index: event.index,
  131. reason: event.reason,
  132. usage: { inputTokens, nonCachedInputTokens: inputTokens },
  133. })
  134. : event,
  135. ),
  136. tool: (id: string, name: string, input: unknown) => [
  137. LLMEvent.stepStart({ index: 0 }),
  138. LLMEvent.toolCall({ id, name, input }),
  139. LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }),
  140. LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
  141. ],
  142. }
  143. const model = Model.make({ id: "fake-model", provider: "fake", route: OpenAIChat.route })
  144. const defaultSystem = PROMPT_DEFAULT
  145. const replacementModel = Model.make({ id: "replacement", provider: "fake", route: OpenAIChat.route })
  146. const compactModel = Model.make({
  147. id: "compact",
  148. provider: "fake",
  149. route: OpenAIChat.route.with({ limits: { context: 4_000, output: 50 } }),
  150. })
  151. const fullOutputModel = Model.make({
  152. id: "full-output",
  153. provider: "fake",
  154. route: OpenAIChat.route.with({ limits: { context: 262_144, output: 262_144 } }),
  155. })
  156. const undersizedContextModel = Model.make({
  157. id: "undersized-context",
  158. provider: "fake",
  159. route: OpenAIChat.route.with({ limits: { context: 1, output: 1_000 } }),
  160. })
  161. const recoveryModel = Model.make({
  162. id: "recovery",
  163. provider: "fake",
  164. route: OpenAIChat.route.with({ limits: { context: 20_000, output: 1_000 } }),
  165. })
  166. test("calculates step cost using the matching context tier", () => {
  167. expect(
  168. SessionUsage.calculateCost(
  169. [
  170. {
  171. input: Money.USDPerMillionTokens.make(1),
  172. output: Money.USDPerMillionTokens.make(2),
  173. cache: {
  174. read: Money.USDPerMillionTokens.make(0.1),
  175. write: Money.USDPerMillionTokens.make(0.5),
  176. },
  177. },
  178. {
  179. tier: { type: "context", size: 100 },
  180. input: Money.USDPerMillionTokens.make(3),
  181. output: Money.USDPerMillionTokens.make(4),
  182. cache: {
  183. read: Money.USDPerMillionTokens.make(0.2),
  184. write: Money.USDPerMillionTokens.make(0.6),
  185. },
  186. },
  187. ],
  188. { input: 80, output: 10, reasoning: 2, cache: { read: 20, write: 1 } },
  189. ),
  190. ).toBeCloseTo(0.0002926)
  191. })
  192. test("does not apply an ineligible tier without base pricing", () => {
  193. expect(
  194. SessionUsage.calculateCost(
  195. [
  196. {
  197. tier: { type: "context", size: 100 },
  198. input: Money.USDPerMillionTokens.make(3),
  199. output: Money.USDPerMillionTokens.make(4),
  200. cache: {
  201. read: Money.USDPerMillionTokens.make(0.2),
  202. write: Money.USDPerMillionTokens.make(0.6),
  203. },
  204. },
  205. ],
  206. { input: 80, output: 10, reasoning: 2, cache: { read: 20, write: 0 } },
  207. ),
  208. ).toBe(Money.USD.zero)
  209. })
  210. const authorizations: Tool.Context[] = []
  211. const executions: string[] = []
  212. const permissionFail = Tool.make({
  213. description: "Reject a permission",
  214. input: Schema.Struct({}),
  215. output: Schema.Struct({}),
  216. execute: () =>
  217. new ToolFailure({
  218. message: "Permission denied: edit",
  219. error: new PermissionV2.BlockedError({
  220. rules: [],
  221. permission: "edit",
  222. resources: ["src/index.ts"],
  223. }),
  224. }),
  225. })
  226. const permission = Layer.succeed(
  227. PermissionV2.Service,
  228. PermissionV2.Service.of({
  229. assert: () => Effect.die("unused"),
  230. ask: () => Effect.die("unused"),
  231. reply: () => Effect.die("unused"),
  232. get: () => Effect.die("unused"),
  233. forSession: () => Effect.die("unused"),
  234. list: () => Effect.die("unused"),
  235. }),
  236. )
  237. const echo = Layer.effectDiscard(
  238. ToolRegistry.Service.use((registry) =>
  239. registry.register(
  240. {
  241. echo: Tool.make({
  242. description: "Echo text",
  243. input: Schema.Struct({ text: Schema.String }),
  244. output: Schema.Struct({ text: Schema.String }),
  245. execute: ({ text }, context) =>
  246. Effect.gen(function* () {
  247. authorizations.push(context)
  248. executions.push(text)
  249. activeToolExecutions++
  250. maxActiveToolExecutions = Math.max(maxActiveToolExecutions, activeToolExecutions)
  251. if (activeToolExecutions === toolExecutionsReady && toolExecutionsStarted) {
  252. yield* Deferred.succeed(toolExecutionsStarted, undefined)
  253. }
  254. if (toolExecutionGate) yield* Deferred.await(toolExecutionGate)
  255. return { output: { text }, content: text }
  256. }).pipe(Effect.ensuring(Effect.sync(() => activeToolExecutions--))),
  257. }),
  258. defect: Tool.make({
  259. description: "Fail unexpectedly",
  260. input: Schema.Struct({}),
  261. output: Schema.Struct({}),
  262. execute: () =>
  263. (toolExecutionGate ? Deferred.await(toolExecutionGate) : Effect.void).pipe(
  264. Effect.andThen(Effect.die("unexpected tool defect")),
  265. ),
  266. }),
  267. // The wrapped ToolOutputStore below fails bound for this call ID with a
  268. // typed StorageError, exercising the infrastructure failure channel.
  269. storefail: Tool.make({
  270. description: "Produce output that cannot be persisted",
  271. input: Schema.Struct({}),
  272. output: Schema.Struct({}),
  273. execute: () => Effect.succeed({ output: {} }),
  274. }),
  275. },
  276. { codemode: false },
  277. ),
  278. ),
  279. )
  280. const echoNode = makeLocationNode({ name: "test/session-runner-tools", layer: echo, deps: [ToolRegistry.node] })
  281. let modelResolveHook = Effect.void
  282. let currentModel = model
  283. const models = Layer.mock(SessionRunnerModel.Service)({
  284. resolve: (session) =>
  285. modelResolveHook.pipe(
  286. Effect.as(
  287. SessionRunnerModel.resolved(session.model?.id === "replacement" ? replacementModel : currentModel, {
  288. capabilities: { tools: true, input: ["text", "image"], output: ["text"] },
  289. cost: [],
  290. variant: session.model?.variant,
  291. }),
  292. ),
  293. ),
  294. })
  295. const systemContextKey = Instructions.Key.make("test/context")
  296. let systemBaseline = "Initial context"
  297. let systemRemoved = false
  298. let systemUnavailable = false
  299. let systemLoadHook = Effect.void
  300. const skillBaselines = new Map<AgentV2.ID, string>()
  301. const systemContext = Layer.mock(InstructionBuiltIns.Service, {
  302. load: () =>
  303. Effect.sync(() =>
  304. Instructions.make({
  305. key: systemContextKey,
  306. codec: Schema.toCodecJson(Schema.String),
  307. read: systemLoadHook.pipe(
  308. Effect.andThen(
  309. Effect.sync(() =>
  310. systemUnavailable ? Instructions.unavailable : systemRemoved ? Instructions.removed : systemBaseline,
  311. ),
  312. ),
  313. ),
  314. render: {
  315. initial: String,
  316. changed: (_previous, current) => current,
  317. removed: () => "System context source removed: test/context",
  318. },
  319. }),
  320. ),
  321. })
  322. const instructionContext = Layer.mock(InstructionDiscovery.Service, { load: () => Effect.succeed(Instructions.empty) })
  323. const skillInstructions = Layer.mock(SkillInstructions.Service, {
  324. load: (agent) =>
  325. Effect.succeed(
  326. skillBaselines.has(agent.id)
  327. ? Instructions.make({
  328. key: Instructions.Key.make("test/skill-guidance"),
  329. codec: Schema.toCodecJson(Schema.String),
  330. read: Effect.succeed(skillBaselines.get(agent.id)!),
  331. render: {
  332. initial: String,
  333. changed: (_previous, current) => current,
  334. removed: () => "Skill guidance removed",
  335. },
  336. })
  337. : Instructions.empty,
  338. ),
  339. })
  340. const referenceInstructions = Layer.mock(ReferenceInstructions.Service, {
  341. load: () => Effect.succeed(Instructions.empty),
  342. })
  343. const mcpInstructions = Layer.mock(McpInstructions.Service, { load: () => Effect.succeed(Instructions.empty) })
  344. const config = Layer.succeed(
  345. Config.Service,
  346. Config.Service.of({
  347. entries: () =>
  348. Effect.succeed([
  349. new Config.Document({
  350. type: "document",
  351. info: new Config.Info({
  352. compaction: new ConfigCompaction.Info({
  353. buffer: 3_000,
  354. keep: new ConfigCompaction.Keep({ tokens: 1_000 }),
  355. }),
  356. }),
  357. }),
  358. ]),
  359. }),
  360. )
  361. let pluginFlushHook = Effect.void
  362. const pluginSupervisor = Layer.succeed(
  363. PluginSupervisor.Service,
  364. PluginSupervisor.Service.of({
  365. flush: Effect.suspend(() => pluginFlushHook),
  366. }),
  367. )
  368. let codeModeMaterializations: ReadonlyArray<CodeMode.Materialization> = []
  369. let codeModeMaterializationCount = 0
  370. const codeMode = Layer.mock(CodeMode.Service, {
  371. register: () => Effect.void,
  372. materialize: () => Effect.sync(() => codeModeMaterializations[codeModeMaterializationCount++] ?? {}),
  373. })
  374. const promptCatalog = Layer.mock(Catalog.Service, {
  375. provider: {
  376. get: () => Effect.succeed(undefined),
  377. all: () => Effect.succeed([]),
  378. available: () => Effect.succeed([]),
  379. },
  380. model: {
  381. get: () => Effect.succeed(undefined),
  382. all: () => Effect.succeed([]),
  383. available: () => Effect.succeed([]),
  384. default: () => Effect.succeed(undefined),
  385. small: () => Effect.succeed(undefined),
  386. },
  387. })
  388. // Pass-through bounding that fails "call-storefail" with a typed StorageError so
  389. // runner tests can exercise the infrastructure failure channel deterministically.
  390. const toolOutputStore = Layer.mock(ToolOutputStore.Service, {
  391. limits: () => Effect.succeed({ maxLines: ToolOutputStore.MAX_LINES, maxBytes: ToolOutputStore.MAX_BYTES }),
  392. bound: (input) =>
  393. input.callID === "call-storefail"
  394. ? Effect.fail(new ToolOutputStore.StorageError({ operation: "write", cause: new Error("disk full") }))
  395. : Effect.succeed({ content: input.content, outputPaths: [] }),
  396. })
  397. const runnerLayer = AppNodeBuilder.build(SessionRunnerLLM.node, [
  398. [Snapshot.node, Snapshot.noopLayer],
  399. [LayerNodePlatform.llmClient, client],
  400. [SessionRunnerModel.node, models],
  401. [InstructionBuiltIns.node, systemContext],
  402. [InstructionDiscovery.node, instructionContext],
  403. [Location.node, Location.boundNode({ directory: AbsolutePath.make("/project") })],
  404. [SkillInstructions.node, skillInstructions],
  405. [ReferenceInstructions.node, referenceInstructions],
  406. [PermissionV2.node, permission],
  407. [Config.node, config],
  408. [McpInstructions.node, mcpInstructions],
  409. [ToolOutputStore.node, toolOutputStore],
  410. [PluginSupervisor.node, pluginSupervisor],
  411. [CodeMode.node, codeMode],
  412. ])
  413. const execution = Layer.effect(
  414. SessionExecution.Service,
  415. Effect.gen(function* () {
  416. const sessionRunner = yield* SessionRunner.Service
  417. const coordinator = yield* SessionRunCoordinator.make<SessionV2.ID, SessionRunner.RunError>({
  418. drain: (sessionID, force) => sessionRunner.drain({ sessionID, force }),
  419. })
  420. return SessionExecution.Service.of({
  421. active: coordinator.active,
  422. resume: coordinator.run,
  423. wake: coordinator.wake,
  424. interrupt: coordinator.interrupt,
  425. awaitIdle: coordinator.awaitIdle,
  426. })
  427. }),
  428. ).pipe(Layer.provide(runnerLayer))
  429. const it = testEffect(
  430. AppNodeBuilder.build(
  431. LayerNode.group([
  432. Database.node,
  433. EventV2.node,
  434. Form.node,
  435. SessionProjector.node,
  436. SessionStore.node,
  437. AgentV2.node,
  438. Catalog.node,
  439. ToolRegistry.node,
  440. ToolRegistry.toolsNode,
  441. ToolHooks.node,
  442. PluginHooks.node,
  443. echoNode,
  444. SessionRunnerModel.node,
  445. InstructionBuiltIns.node,
  446. InstructionDiscovery.node,
  447. InstructionEntry.node,
  448. SkillInstructions.node,
  449. ReferenceInstructions.node,
  450. Config.node,
  451. Snapshot.node,
  452. SessionRunnerLLM.node,
  453. SessionExecution.node,
  454. SessionV2.node,
  455. ]),
  456. [
  457. [LayerNodePlatform.llmClient, client],
  458. [PermissionV2.node, permission],
  459. [Catalog.node, promptCatalog],
  460. [SessionRunnerModel.node, models],
  461. [InstructionBuiltIns.node, systemContext],
  462. [InstructionDiscovery.node, instructionContext],
  463. [Location.node, Location.boundNode({ directory: AbsolutePath.make("/project") })],
  464. [SkillInstructions.node, skillInstructions],
  465. [ReferenceInstructions.node, referenceInstructions],
  466. [Snapshot.node, Snapshot.noopLayer],
  467. [SessionExecution.node, execution],
  468. [Config.node, config],
  469. [ToolOutputStore.node, toolOutputStore],
  470. [PluginSupervisor.node, pluginSupervisor],
  471. [CodeMode.node, codeMode],
  472. ],
  473. ),
  474. )
  475. const sessionID = SessionV2.ID.make("ses_runner_test")
  476. const otherSessionID = SessionV2.ID.make("ses_runner_other")
  477. const admit = (session: SessionV2.Interface, text: string) => session.prompt({ sessionID, text, resume: false })
  478. const insertSession = (id: SessionV2.ID) =>
  479. Effect.gen(function* () {
  480. const { db } = yield* Database.Service
  481. yield* db
  482. .insert(SessionTable)
  483. .values({
  484. id,
  485. project_id: Project.ID.global,
  486. slug: id,
  487. directory: "/project",
  488. title: "test",
  489. version: "test",
  490. })
  491. .onConflictDoNothing()
  492. .run()
  493. .pipe(Effect.orDie)
  494. })
  495. const setup = Effect.gen(function* () {
  496. const { db } = yield* Database.Service
  497. const agents = yield* AgentV2.Service
  498. const catalog = yield* Catalog.Service
  499. const hooks = yield* PluginHooks.Service
  500. const pluginHost = host({
  501. agent: agentHost(agents),
  502. catalog: catalogHost(catalog),
  503. session: { hook: (name, callback) => hooks.register("session", name, callback) },
  504. })
  505. yield* Effect.forEach(SystemPromptPlugin.Plugins, (plugin) => plugin.effect(pluginHost), {
  506. discard: true,
  507. })
  508. requests.length = 0
  509. authorizations.length = 0
  510. executions.length = 0
  511. response = []
  512. systemBaseline = "Initial context"
  513. systemRemoved = false
  514. systemUnavailable = false
  515. systemLoadHook = Effect.void
  516. modelResolveHook = Effect.void
  517. pluginFlushHook = Effect.void
  518. codeModeMaterializations = []
  519. codeModeMaterializationCount = 0
  520. currentModel = model
  521. skillBaselines.clear()
  522. responses = undefined
  523. streamFailure = undefined
  524. responseStream = undefined
  525. responseStreams = undefined
  526. streamGate = undefined
  527. streamStarted = undefined
  528. toolExecutionGate = undefined
  529. toolExecutionsStarted = undefined
  530. toolExecutionsReady = 5
  531. activeToolExecutions = 0
  532. maxActiveToolExecutions = 0
  533. yield* agents.transform((draft) =>
  534. draft.update(AgentV2.ID.make("build"), (agent) => {
  535. agent.mode = "primary"
  536. }),
  537. )
  538. yield* db
  539. .insert(ProjectTable)
  540. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  541. .onConflictDoNothing()
  542. .run()
  543. .pipe(Effect.orDie)
  544. yield* insertSession(sessionID)
  545. return yield* SessionV2.Service
  546. })
  547. const providerUnavailable = () =>
  548. new LLMError({
  549. module: "test",
  550. method: "stream",
  551. reason: new TransportReason({ message: "Provider unavailable" }),
  552. })
  553. const invalidRequest = () =>
  554. new LLMError({
  555. module: "test",
  556. method: "stream",
  557. reason: new InvalidRequestReason({ message: "Invalid request" }),
  558. })
  559. const rateLimited = (retryAfterMs?: number) =>
  560. new LLMError({
  561. module: "test",
  562. method: "stream",
  563. reason: new RateLimitReason({ message: "Rate limited", retryAfterMs }),
  564. })
  565. const setupOverflowRecovery = Effect.gen(function* () {
  566. const session = yield* setup
  567. response = reply.text("Earlier answer", "text-earlier")
  568. yield* admit(session, "Earlier question ".repeat(700))
  569. yield* session.resume(sessionID)
  570. currentModel = recoveryModel
  571. requests.length = 0
  572. return session
  573. })
  574. const messageTexts = (request: LLMRequest, role: "user" | "system") =>
  575. request.messages.flatMap((message) =>
  576. message.role === role ? message.content.flatMap((content) => (content.type === "text" ? [content.text] : [])) : [],
  577. )
  578. const userTexts = (request: LLMRequest) => messageTexts(request, "user")
  579. const systemTexts = (request: LLMRequest) => messageTexts(request, "system")
  580. const recordedEventTypes = (id: SessionV2.ID) =>
  581. Effect.gen(function* () {
  582. const { db } = yield* Database.Service
  583. return yield* db
  584. .select({ type: EventTable.type })
  585. .from(EventTable)
  586. .where(eq(EventTable.aggregate_id, id))
  587. .orderBy(asc(EventTable.seq))
  588. .all()
  589. .pipe(
  590. Effect.orDie,
  591. Effect.map((rows) => rows.map((row) => row.type)),
  592. )
  593. })
  594. const recordedStepSettlementEvents = (id: SessionV2.ID, assistantMessageID: SessionMessage.ID) =>
  595. Effect.gen(function* () {
  596. const { db } = yield* Database.Service
  597. const settlementTypes = new Set([
  598. "session.step.started.1",
  599. "session.tool.called.1",
  600. "session.tool.success.2",
  601. "session.tool.failed.2",
  602. "session.step.ended.1",
  603. "session.step.failed.1",
  604. ])
  605. return (yield* db
  606. .select({ type: EventTable.type, data: EventTable.data })
  607. .from(EventTable)
  608. .where(eq(EventTable.aggregate_id, id))
  609. .orderBy(asc(EventTable.seq))
  610. .all()
  611. .pipe(Effect.orDie)).filter(
  612. (event) => settlementTypes.has(event.type) && event.data.assistantMessageID === assistantMessageID,
  613. )
  614. })
  615. const hostedCall = (id: string, query: string) =>
  616. LLMEvent.toolCall({ id, name: "web_search", input: { query }, providerExecuted: true })
  617. const requireAssistant = (messages: readonly SessionMessage.Info[]) => {
  618. const assistant = messages.find((message) => message.type === "assistant")
  619. if (!assistant) throw new Error("Assistant message missing")
  620. return assistant
  621. }
  622. const replaySessionProjection = (id: SessionV2.ID) =>
  623. Effect.gen(function* () {
  624. const { db } = yield* Database.Service
  625. const events = yield* EventV2.Service
  626. const recorded = yield* db
  627. .select()
  628. .from(EventTable)
  629. .where(eq(EventTable.aggregate_id, id))
  630. .orderBy(asc(EventTable.seq))
  631. .all()
  632. .pipe(Effect.orDie)
  633. yield* events.remove(id)
  634. yield* db.delete(InstructionStateTable).where(eq(InstructionStateTable.session_id, id)).run().pipe(Effect.orDie)
  635. yield* db.delete(SessionPendingTable).where(eq(SessionPendingTable.session_id, id)).run().pipe(Effect.orDie)
  636. yield* db.delete(SessionMessageTable).where(eq(SessionMessageTable.session_id, id)).run().pipe(Effect.orDie)
  637. yield* events.replayAll(
  638. recorded.map((event) => ({
  639. id: event.id,
  640. created: DateTime.makeUnsafe(event.created),
  641. aggregateID: event.aggregate_id,
  642. seq: event.seq,
  643. type: event.type,
  644. data: event.data,
  645. })),
  646. )
  647. })
  648. type FragmentKind = "text" | "reasoning" | "tool input"
  649. type FragmentFixture = {
  650. readonly delta: EventV2.Definition
  651. readonly completeEvents: LLMEvent[]
  652. readonly partialEvents: LLMEvent[]
  653. readonly expectedAssistant: unknown
  654. readonly expectedContent: unknown
  655. }
  656. const fragmentKinds: readonly FragmentKind[] = ["text", "reasoning", "tool input"]
  657. const fragmentID = (kind: FragmentKind, suffix: string) => `${kind === "tool input" ? "call" : kind}-${suffix}`
  658. const fragmentFixture = (kind: FragmentKind, id: string, chunks: readonly string[]): FragmentFixture => {
  659. const text = chunks.join("")
  660. switch (kind) {
  661. case "text": {
  662. const partialEvents = [
  663. LLMEvent.stepStart({ index: 0 }),
  664. LLMEvent.textStart({ id }),
  665. ...chunks.map((text) => LLMEvent.textDelta({ id, text })),
  666. ]
  667. const expectedContent = { type: "text", text }
  668. return {
  669. delta: SessionEvent.Text.Delta,
  670. partialEvents,
  671. completeEvents: [
  672. ...partialEvents,
  673. LLMEvent.textEnd({ id }),
  674. LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
  675. LLMEvent.finish({ reason: { normalized: "stop" } }),
  676. ],
  677. expectedAssistant: { type: "assistant", finish: "stop", content: [expectedContent] },
  678. expectedContent,
  679. }
  680. }
  681. case "reasoning": {
  682. const partialEvents = [
  683. LLMEvent.stepStart({ index: 0 }),
  684. LLMEvent.reasoningStart({ id }),
  685. ...chunks.map((text) => LLMEvent.reasoningDelta({ id, text })),
  686. ]
  687. const expectedContent = { type: "reasoning", text }
  688. return {
  689. delta: SessionEvent.Reasoning.Delta,
  690. partialEvents,
  691. completeEvents: [
  692. ...partialEvents,
  693. LLMEvent.reasoningEnd({ id }),
  694. LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
  695. LLMEvent.finish({ reason: { normalized: "stop" } }),
  696. ],
  697. expectedAssistant: { type: "assistant", finish: "stop", content: [expectedContent] },
  698. expectedContent,
  699. }
  700. }
  701. case "tool input": {
  702. const partialEvents = [
  703. LLMEvent.stepStart({ index: 0 }),
  704. LLMEvent.toolInputStart({ id, name: "echo" }),
  705. ...chunks.map((text) => LLMEvent.toolInputDelta({ id, name: "echo", text })),
  706. ]
  707. const expectedContent = { type: "tool", id, state: { status: "streaming", input: text } }
  708. return {
  709. delta: SessionEvent.Tool.Input.Delta,
  710. partialEvents,
  711. completeEvents: [...partialEvents, LLMEvent.toolInputEnd({ id, name: "echo" })],
  712. expectedAssistant: { type: "assistant", content: [expectedContent] },
  713. expectedContent,
  714. }
  715. }
  716. }
  717. }
  718. const verifyEphemeralDeltas = (kind: FragmentKind) =>
  719. Effect.gen(function* () {
  720. const session = yield* setup
  721. const prompt = `Stream ${kind}`
  722. const chunks = Array.from({ length: 32 }, (_, index) => `${index},`)
  723. const fixture = fragmentFixture(kind, fragmentID(kind, "many"), chunks)
  724. const expectedContext = [{ type: "user", text: prompt }, fixture.expectedAssistant]
  725. yield* admit(session, prompt)
  726. const events = yield* EventV2.Service
  727. const live = yield* events.subscribe(fixture.delta).pipe(Stream.take(32), Stream.runCollect, Effect.forkScoped)
  728. yield* Effect.yieldNow
  729. response = fixture.completeEvents
  730. yield* session.resume(sessionID)
  731. const { db } = yield* Database.Service
  732. const deltas = yield* db
  733. .select({ type: EventTable.type })
  734. .from(EventTable)
  735. .where(eq(EventTable.type, EventV2.versionedType(fixture.delta.type, 1)))
  736. .all()
  737. .pipe(Effect.orDie)
  738. expect(Array.from(yield* Fiber.join(live))).toHaveLength(32)
  739. expect(deltas).toHaveLength(0)
  740. expect(yield* session.context(sessionID)).toMatchObject(expectedContext)
  741. yield* replaySessionProjection(sessionID)
  742. expect(yield* session.context(sessionID)).toMatchObject(expectedContext)
  743. })
  744. const verifyPartialFlushOnFailure = (kind: FragmentKind) =>
  745. Effect.gen(function* () {
  746. const session = yield* setup
  747. const prompt = `Fail after ${kind}`
  748. const fixture = fragmentFixture(kind, fragmentID(kind, "partial"), ["Partial"])
  749. const failure = providerUnavailable()
  750. yield* admit(session, prompt)
  751. responseStream = Stream.concat(Stream.fromIterable(fixture.partialEvents), Stream.fail(failure))
  752. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  753. expect(yield* session.context(sessionID)).toMatchObject([
  754. { type: "user", text: prompt },
  755. {
  756. type: "assistant",
  757. finish: "error",
  758. error: { type: "provider.transport", message: "Provider unavailable" },
  759. content: [
  760. kind === "tool input"
  761. ? {
  762. type: "tool",
  763. id: fragmentID(kind, "partial"),
  764. state: {
  765. status: "error",
  766. error: { type: "provider.transport", message: "Provider unavailable" },
  767. },
  768. }
  769. : fixture.expectedContent,
  770. ],
  771. },
  772. ])
  773. expect(requests).toHaveLength(1)
  774. })
  775. const verifyPartialFlushOnInterruption = (kind: FragmentKind) =>
  776. Effect.gen(function* () {
  777. const session = yield* setup
  778. const prompt = `Interrupt after ${kind}`
  779. const fixture = fragmentFixture(kind, fragmentID(kind, "interrupted"), ["Partial"])
  780. const streamed = yield* Deferred.make<void>()
  781. yield* admit(session, prompt)
  782. responseStream = Stream.concat(
  783. Stream.fromIterable(fixture.partialEvents),
  784. Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)),
  785. )
  786. const runner = yield* SessionRunner.Service
  787. const fiber = yield* runner.drain({ sessionID, force: true }).pipe(Effect.forkChild)
  788. yield* Deferred.await(streamed)
  789. yield* Fiber.interrupt(fiber)
  790. expect(yield* session.context(sessionID)).toMatchObject([
  791. { type: "user", text: prompt },
  792. {
  793. type: "assistant",
  794. finish: "error",
  795. error: { type: "aborted", message: "Step interrupted" },
  796. content: [
  797. kind === "tool input"
  798. ? { type: "tool", id: fragmentID(kind, "interrupted"), state: { status: "error" } }
  799. : fixture.expectedContent,
  800. ],
  801. },
  802. ])
  803. })
  804. describe("SessionRunnerLLM", () => {
  805. it.effect("uses one Code Mode materialization per request for instructions and execution", () =>
  806. Effect.gen(function* () {
  807. const executed: string[] = []
  808. const execute = (name: string) =>
  809. Tool.make({
  810. description: `Execute ${name}`,
  811. input: Schema.Struct({}),
  812. output: Schema.String,
  813. execute: () => Effect.sync(() => executed.push(name)).pipe(Effect.as({ output: name })),
  814. })
  815. const catalog = (name: string) => [
  816. {
  817. path: `catalog.${name.toLowerCase()}`,
  818. description: `Code Mode catalog ${name}`,
  819. signature: `tools.catalog.${name.toLowerCase()}(input: {}): Promise<string>`,
  820. },
  821. ]
  822. const session = yield* setup
  823. codeModeMaterializations = [
  824. { catalog: catalog("A"), tool: execute("A") },
  825. { catalog: catalog("B"), tool: execute("B") },
  826. { catalog: catalog("C"), tool: execute("C") },
  827. { catalog: catalog("D"), tool: execute("D") },
  828. ]
  829. yield* admit(session, "Use Code Mode")
  830. responses = [reply.tool("call-execute", "execute", {}), reply.stop()]
  831. yield* session.resume(sessionID)
  832. expect(requests).toHaveLength(2)
  833. expect(codeModeMaterializationCount).toBe(2)
  834. expect(requests[0]?.system.some((part) => part.text.includes("Code Mode catalog A"))).toBe(true)
  835. expect(requests[0]?.system.some((part) => part.text.includes("Code Mode catalog B"))).toBe(false)
  836. expect(requests[0]?.tools.find((tool) => tool.name === "execute")?.description).toBe("Execute A")
  837. expect(executed).toEqual(["A"])
  838. expect(requests[1]?.tools.find((tool) => tool.name === "execute")?.description).toBe("Execute B")
  839. expect(
  840. requests[1]?.messages.some(
  841. (message) =>
  842. message.role === "system" &&
  843. message.content.some((part) => part.type === "text" && part.text.includes("Code Mode catalog B")),
  844. ),
  845. ).toBe(true)
  846. }),
  847. )
  848. it.effect("advertises execute and durable guidance for an empty Code Mode catalog", () =>
  849. Effect.gen(function* () {
  850. const session = yield* setup
  851. const empty = {
  852. catalog: [],
  853. tool: Tool.make({
  854. description: "Execute Code Mode",
  855. input: Schema.Struct({ code: Schema.String }),
  856. output: Schema.String,
  857. execute: () => Effect.succeed({ output: "unused" }),
  858. }),
  859. }
  860. codeModeMaterializations = [empty, empty, {}]
  861. yield* admit(session, "Continue without Code Mode tools")
  862. response = reply.stop()
  863. yield* session.resume(sessionID)
  864. yield* admit(session, "Still no Code Mode tools")
  865. yield* session.resume(sessionID)
  866. yield* admit(session, "Code Mode denied")
  867. yield* session.resume(sessionID)
  868. expect(requests).toHaveLength(3)
  869. expect(requests[0]?.tools.map((tool) => tool.name)).toEqual(["defect", "echo", "storefail", "execute"])
  870. expect(requests[0]?.system.some((part) => part.text.includes("Do not call `execute`"))).toBe(true)
  871. expect(requests[1]?.tools.map((tool) => tool.name)).toEqual(["defect", "echo", "storefail", "execute"])
  872. expect(requests[1]?.messages.filter((message) => message.role === "system")).toEqual([])
  873. expect(requests[2]?.tools.map((tool) => tool.name)).toEqual(["defect", "echo", "storefail"])
  874. expect(
  875. requests[2]?.messages.some(
  876. (message) =>
  877. message.role === "system" &&
  878. message.content.some(
  879. (part) => part.type === "text" && part.text.includes("Code Mode tools are no longer available"),
  880. ),
  881. ),
  882. ).toBe(true)
  883. }),
  884. )
  885. it.effect("applies session context hooks without exposing unavailable tools", () =>
  886. Effect.gen(function* () {
  887. const session = yield* setup
  888. const hooks = yield* PluginHooks.Service
  889. yield* hooks.register("session", "context", (event) =>
  890. Effect.sync(() => {
  891. event.system = [SystemPart.make("Hooked system")]
  892. event.messages = [Message.user("Hooked message")]
  893. delete event.tools.echo
  894. event.tools.unregistered = { description: "Unavailable", input: { type: "object" } }
  895. }),
  896. )
  897. yield* admit(session, "Original message")
  898. responses = [reply.tool("call-removed", "echo", { text: "blocked" })]
  899. yield* session.resume(sessionID)
  900. // A hook-removed call fails independently and continues while step allowance remains.
  901. expect(requests).toHaveLength(2)
  902. expect(requests[0]?.system.map((part) => part.text)).toEqual(["Hooked system"])
  903. expect(requests[0]?.messages).toEqual([Message.user("Hooked message")])
  904. expect(requests[0]?.tools.map((tool) => tool.name)).not.toContain("echo")
  905. expect(requests[0]?.tools.map((tool) => tool.name)).not.toContain("unregistered")
  906. expect(executions).toEqual([])
  907. expect(yield* session.context(sessionID)).toMatchObject([
  908. { type: "user", text: "Original message" },
  909. {
  910. type: "assistant",
  911. content: [
  912. {
  913. type: "tool",
  914. id: "call-removed",
  915. state: { status: "error", error: { type: "tool.unknown" } },
  916. },
  917. ],
  918. },
  919. ])
  920. }),
  921. )
  922. it.effect("advertises and executes a location registered tool", () =>
  923. Effect.gen(function* () {
  924. const session = yield* setup
  925. const registry = yield* ToolRegistry.Service
  926. const contexts: Tool.Context[] = []
  927. yield* registry.register(
  928. {
  929. location_context: Tool.make({
  930. description: "Read application context",
  931. input: Schema.Struct({ query: Schema.String }),
  932. output: Schema.Struct({ answer: Schema.String }),
  933. execute: ({ query }, context) =>
  934. Effect.gen(function* () {
  935. contexts.push(context)
  936. yield* context.progress({ phase: "reading" })
  937. return { output: { answer: query.toUpperCase() } }
  938. }),
  939. }),
  940. },
  941. { codemode: false },
  942. )
  943. yield* admit(session, "Use application context")
  944. responses = [reply.tool("call-location", "location_context", { query: "hello" }), []]
  945. const events = yield* EventV2.Service
  946. const progressFiber = yield* events.subscribe(SessionEvent.Tool.Progress).pipe(
  947. Stream.filter((event) => event.data.sessionID === sessionID && event.data.callID === "call-location"),
  948. Stream.take(1),
  949. Stream.runCollect,
  950. Effect.forkScoped({ startImmediately: true }),
  951. )
  952. yield* session.resume(sessionID)
  953. expect(requests[0]?.tools.map((tool) => tool.name)).toContain("location_context")
  954. expect(contexts).toEqual([
  955. {
  956. sessionID,
  957. agent: AgentV2.ID.make("build"),
  958. messageID: expect.stringMatching(/^msg_/),
  959. callID: "call-location",
  960. progress: expect.any(Function),
  961. },
  962. ])
  963. expect(Array.from(yield* Fiber.join(progressFiber))[0]?.data.metadata).toEqual({ phase: "reading" })
  964. expect(yield* session.context(sessionID)).toMatchObject([
  965. { type: "user", text: "Use application context" },
  966. {
  967. type: "assistant",
  968. content: [
  969. {
  970. type: "tool",
  971. id: "call-location",
  972. state: { status: "completed", content: [{ type: "text", text: '{"answer":"HELLO"}' }] },
  973. },
  974. ],
  975. },
  976. ])
  977. }),
  978. )
  979. it.effect("prefers failure outcome metadata over retained progress", () =>
  980. Effect.gen(function* () {
  981. const session = yield* setup
  982. const registry = yield* ToolRegistry.Service
  983. const hooks = yield* ToolHooks.Service
  984. yield* hooks.hook.after((event) => {
  985. if (event.status === "error") event.metadata = { phase: "failed" }
  986. })
  987. yield* registry.register(
  988. {
  989. failing_progress: Tool.make({
  990. description: "Report progress and fail",
  991. input: Schema.Struct({}),
  992. output: Schema.Struct({}),
  993. execute: (_, context) =>
  994. Effect.gen(function* () {
  995. yield* context.progress({ phase: "running" })
  996. return yield* new ToolFailure({ message: "failed after progress" })
  997. }),
  998. }),
  999. },
  1000. { codemode: false },
  1001. )
  1002. yield* admit(session, "Run failing progress")
  1003. responses = [reply.tool("call-failing-progress", "failing_progress", {}), reply.stop()]
  1004. yield* session.resume(sessionID)
  1005. expect(yield* session.context(sessionID)).toMatchObject([
  1006. { type: "user", text: "Run failing progress" },
  1007. {
  1008. type: "assistant",
  1009. content: [
  1010. {
  1011. type: "tool",
  1012. id: "call-failing-progress",
  1013. state: {
  1014. status: "error",
  1015. metadata: { phase: "failed" },
  1016. error: { message: "failed after progress" },
  1017. },
  1018. },
  1019. ],
  1020. },
  1021. { type: "assistant", finish: "stop" },
  1022. ])
  1023. }),
  1024. )
  1025. it.effect("executes the tool advertised before a registry reload", () =>
  1026. Effect.gen(function* () {
  1027. const session = yield* setup
  1028. const registry = yield* ToolRegistry.Service
  1029. const scope = yield* Scope.make()
  1030. const executions: string[] = []
  1031. yield* registry
  1032. .register(
  1033. {
  1034. reloaded: Tool.make({
  1035. description: "Record the advertised tool",
  1036. input: Schema.Struct({}),
  1037. output: Schema.Struct({ value: Schema.String }),
  1038. execute: () =>
  1039. Effect.sync(() => executions.push("advertised")).pipe(Effect.as({ output: { value: "advertised" } })),
  1040. }),
  1041. },
  1042. { codemode: false },
  1043. )
  1044. .pipe(Scope.provide(scope))
  1045. yield* admit(session, "Use the reloaded tool")
  1046. responses = [
  1047. [
  1048. LLMEvent.stepStart({ index: 0 }),
  1049. LLMEvent.toolCall({ id: "call-reloaded", name: "reloaded", input: {} }),
  1050. LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }),
  1051. LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
  1052. ],
  1053. [],
  1054. ]
  1055. streamGate = yield* Deferred.make<void>()
  1056. streamStarted = yield* Deferred.make<void>()
  1057. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1058. yield* Deferred.await(streamStarted)
  1059. yield* Scope.close(scope, Exit.void)
  1060. yield* registry.register(
  1061. {
  1062. reloaded: Tool.make({
  1063. description: "Record the replacement tool",
  1064. input: Schema.Struct({}),
  1065. output: Schema.Struct({ value: Schema.String }),
  1066. execute: () =>
  1067. Effect.sync(() => executions.push("replacement")).pipe(Effect.as({ output: { value: "replacement" } })),
  1068. }),
  1069. },
  1070. { codemode: false },
  1071. )
  1072. yield* Deferred.succeed(streamGate, undefined)
  1073. yield* Fiber.join(run)
  1074. expect(executions).toEqual(["advertised"])
  1075. expect(yield* session.context(sessionID)).toMatchObject([
  1076. { type: "user", text: "Use the reloaded tool" },
  1077. {
  1078. type: "assistant",
  1079. content: [
  1080. {
  1081. type: "tool",
  1082. id: "call-reloaded",
  1083. state: { status: "completed", content: [{ type: "text", text: '{"value":"advertised"}' }] },
  1084. },
  1085. ],
  1086. },
  1087. ])
  1088. }),
  1089. )
  1090. it.effect("starts a real runner step after default prompt recording", () =>
  1091. Effect.gen(function* () {
  1092. const session = yield* setup
  1093. const message = yield* session.prompt({
  1094. sessionID,
  1095. text: "Run automatically",
  1096. })
  1097. yield* session.wait(sessionID)
  1098. expect(requests).toHaveLength(1)
  1099. expect(yield* session.messages({ sessionID })).toMatchObject([
  1100. { id: message.id, type: "user", text: "Run automatically" },
  1101. ])
  1102. }),
  1103. )
  1104. it.effect("runs a follow-up when synthetic input arrives during an active continuation", () =>
  1105. Effect.gen(function* () {
  1106. const session = yield* setup
  1107. const secondStarted = yield* Deferred.make<void>()
  1108. const releaseSecond = yield* Deferred.make<void>()
  1109. responseStreams = [
  1110. Stream.fromIterable(reply.tool("call-echo", "echo", { text: "background started" })),
  1111. Stream.unwrap(
  1112. Deferred.succeed(secondStarted, undefined).pipe(
  1113. Effect.andThen(Deferred.await(releaseSecond)),
  1114. Effect.as(Stream.fromIterable(reply.stop())),
  1115. ),
  1116. ),
  1117. Stream.fromIterable(reply.text("Handled completion", "text-completion")),
  1118. ]
  1119. yield* admit(session, "Start background work")
  1120. const running = yield* session.resume(sessionID).pipe(Effect.forkChild({ startImmediately: true }))
  1121. yield* Deferred.await(secondStarted)
  1122. yield* session.synthetic({ sessionID, text: "Background work completed" })
  1123. yield* Deferred.succeed(releaseSecond, undefined)
  1124. yield* Fiber.join(running)
  1125. expect(requests).toHaveLength(3)
  1126. expect(userTexts(requests[2]!)).toContain("Background work completed")
  1127. }),
  1128. )
  1129. it.effect("streams one request with registry definitions from chronological V2 user history", () =>
  1130. Effect.gen(function* () {
  1131. const session = yield* setup
  1132. yield* admit(session, "First")
  1133. yield* admit(session, "Second")
  1134. yield* session.resume(sessionID)
  1135. expect(requests).toHaveLength(1)
  1136. expect(requests[0]?.model).toBe(model)
  1137. expect(requests[0]?.tools.map((tool) => tool.name)).toEqual(["defect", "echo", "storefail"])
  1138. expect(requests[0]?.messages.map((message) => ({ role: message.role, content: message.content }))).toEqual([
  1139. { role: "user", content: [{ type: "text", text: "First" }] },
  1140. { role: "user", content: [{ type: "text", text: "Second" }] },
  1141. ])
  1142. expect(yield* session.messages({ sessionID })).toHaveLength(2)
  1143. }),
  1144. )
  1145. it.effect("marks the initial instruction sync as baseline metadata", () =>
  1146. Effect.gen(function* () {
  1147. const session = yield* setup
  1148. const events = yield* EventV2.Service
  1149. const instructionEvents: EventV2.Payload[] = []
  1150. const unsubscribe = yield* events.listen((event) =>
  1151. Effect.sync(() => {
  1152. if (event.type === "session.instructions.updated") instructionEvents.push(event)
  1153. }),
  1154. )
  1155. yield* admit(session, "First")
  1156. yield* session.resume(sessionID)
  1157. systemBaseline = "Changed context"
  1158. yield* admit(session, "Second")
  1159. yield* session.resume(sessionID)
  1160. yield* unsubscribe
  1161. expect(instructionEvents).toHaveLength(2)
  1162. expect(instructionEvents[0]?.metadata).toEqual({ instructions: { initial: true } })
  1163. expect(instructionEvents[1]?.metadata).toBeUndefined()
  1164. }),
  1165. )
  1166. it.effect("retries the first request after system context becomes available", () =>
  1167. Effect.gen(function* () {
  1168. const session = yield* setup
  1169. const { db } = yield* Database.Service
  1170. const messageID = SessionMessage.ID.create()
  1171. systemUnavailable = true
  1172. yield* session.prompt({
  1173. id: messageID,
  1174. sessionID,
  1175. text: "First",
  1176. resume: false,
  1177. })
  1178. const exit = yield* session.resume(sessionID).pipe(Effect.exit)
  1179. expect(Exit.isFailure(exit)).toBe(true)
  1180. if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(Instructions.InitializationBlocked)
  1181. expect(requests).toHaveLength(0)
  1182. expect(yield* SessionPending.has(db, sessionID, "steer")).toBe(true)
  1183. expect(
  1184. yield* db.select().from(InstructionStateTable).where(eq(InstructionStateTable.session_id, sessionID)).get(),
  1185. ).toBeUndefined()
  1186. systemUnavailable = false
  1187. yield* session.prompt({ id: messageID, sessionID, text: "First" })
  1188. yield* session.wait(sessionID)
  1189. expect(requests).toHaveLength(1)
  1190. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user"])
  1191. }),
  1192. )
  1193. it.effect("interrupts a source Location runner after a Session moves", () =>
  1194. Effect.gen(function* () {
  1195. const session = yield* setup
  1196. const events = yield* EventV2.Service
  1197. const { db } = yield* Database.Service
  1198. yield* admit(session, "First")
  1199. yield* session.resume(sessionID)
  1200. yield* events.publish(SessionEvent.Moved, {
  1201. sessionID,
  1202. location: Location.Ref.make({ directory: AbsolutePath.make("/moved") }),
  1203. })
  1204. expect(
  1205. yield* db.select().from(InstructionStateTable).where(eq(InstructionStateTable.session_id, sessionID)).get(),
  1206. ).toBeUndefined()
  1207. yield* admit(session, "Second")
  1208. const exit = yield* session.resume(sessionID).pipe(Effect.exit)
  1209. expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true)
  1210. expect(requests).toHaveLength(1)
  1211. expect(yield* SessionPending.has(db, sessionID, "steer")).toBe(true)
  1212. }),
  1213. )
  1214. it.effect("forks instruction values at the selected message instead of the parent's latest state", () =>
  1215. Effect.gen(function* () {
  1216. const session = yield* setup
  1217. const first = yield* admit(session, "First")
  1218. yield* session.resume(sessionID)
  1219. systemBaseline = "Changed context"
  1220. const second = yield* admit(session, "Second")
  1221. yield* session.resume(sessionID)
  1222. systemBaseline = "Latest context"
  1223. yield* admit(session, "Third")
  1224. yield* session.resume(sessionID)
  1225. const forked = yield* session.fork({ sessionID, messageID: second.id })
  1226. expect(
  1227. yield* (yield* Database.Service).db
  1228. .select()
  1229. .from(InstructionStateTable)
  1230. .where(eq(InstructionStateTable.session_id, forked.id))
  1231. .get(),
  1232. ).toMatchObject({
  1233. initial_values: { "test/context": Instructions.hash("Initial context") },
  1234. current_values: { "test/context": Instructions.hash("Changed context") },
  1235. })
  1236. yield* session.prompt({ sessionID: forked.id, text: "Forked", resume: false })
  1237. yield* session.resume(forked.id)
  1238. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
  1239. expect(systemTexts(requests.at(-1)!)).toContain("Changed context")
  1240. expect(systemTexts(requests.at(-1)!)).toContain("Latest context")
  1241. const { db } = yield* Database.Service
  1242. const events = yield* EventV2.Service
  1243. const recorded = yield* db
  1244. .select()
  1245. .from(EventTable)
  1246. .where(eq(EventTable.aggregate_id, forked.id))
  1247. .orderBy(asc(EventTable.seq))
  1248. .all()
  1249. yield* events.remove(forked.id)
  1250. yield* db.delete(SessionTable).where(eq(SessionTable.id, forked.id)).run()
  1251. yield* events.replayAll(
  1252. recorded.map((event) => ({
  1253. id: event.id,
  1254. created: DateTime.makeUnsafe(event.created),
  1255. aggregateID: event.aggregate_id,
  1256. seq: event.seq,
  1257. type: event.type,
  1258. data: event.data,
  1259. })),
  1260. )
  1261. expect(
  1262. yield* db.select().from(InstructionStateTable).where(eq(InstructionStateTable.session_id, forked.id)).get(),
  1263. ).toMatchObject({ current_values: { "test/context": Instructions.hash("Latest context") } })
  1264. }),
  1265. )
  1266. it.effect("caps nested fork instruction ancestry at the selected message", () =>
  1267. Effect.gen(function* () {
  1268. const session = yield* setup
  1269. yield* admit(session, "First")
  1270. yield* session.resume(sessionID)
  1271. systemBaseline = "Changed context"
  1272. const second = yield* admit(session, "Second")
  1273. yield* session.resume(sessionID)
  1274. const child = yield* session.fork({ sessionID, messageID: second.id })
  1275. const inheritedFirst = (yield* session.messages({ sessionID: child.id })).find(
  1276. (message) => message.type === "user" && message.text === "First",
  1277. )
  1278. if (!inheritedFirst) return yield* Effect.die(new Error("Nested fork boundary message not found"))
  1279. const grandchild = yield* session.fork({ sessionID: child.id, messageID: inheritedFirst.id })
  1280. expect(
  1281. yield* (yield* Database.Service).db
  1282. .select()
  1283. .from(InstructionStateTable)
  1284. .where(eq(InstructionStateTable.session_id, grandchild.id))
  1285. .get(),
  1286. ).toMatchObject({
  1287. initial_values: { "test/context": Instructions.hash("Initial context") },
  1288. current_values: { "test/context": Instructions.hash("Initial context") },
  1289. })
  1290. }),
  1291. )
  1292. it.effect("rebuilds a missing instruction cache without admitting another delta", () =>
  1293. Effect.gen(function* () {
  1294. const session = yield* setup
  1295. const { db } = yield* Database.Service
  1296. yield* admit(session, "First")
  1297. yield* session.resume(sessionID)
  1298. yield* db.delete(InstructionStateTable).where(eq(InstructionStateTable.session_id, sessionID)).run()
  1299. yield* admit(session, "Second")
  1300. requests.length = 0
  1301. yield* session.resume(sessionID)
  1302. expect(requests).toHaveLength(1)
  1303. expect(requests[0]?.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
  1304. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "user"])
  1305. expect(
  1306. yield* db
  1307. .select({ id: EventTable.id })
  1308. .from(EventTable)
  1309. .where(eq(EventTable.type, "session.instructions.updated.2"))
  1310. .all(),
  1311. ).toHaveLength(1)
  1312. expect(yield* db.select().from(InstructionStateTable).get()).toMatchObject({
  1313. initial_values: { "test/context": Instructions.hash("Initial context") },
  1314. current_values: { "test/context": Instructions.hash("Initial context") },
  1315. })
  1316. }),
  1317. )
  1318. it.effect("keeps the initial instructions stable and derives a chronological update from values", () =>
  1319. Effect.gen(function* () {
  1320. const session = yield* setup
  1321. yield* admit(session, "First")
  1322. yield* session.resume(sessionID)
  1323. systemBaseline = "Changed context"
  1324. yield* admit(session, "Second")
  1325. yield* session.resume(sessionID)
  1326. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1327. [defaultSystem, "Initial context"],
  1328. [defaultSystem, "Initial context"],
  1329. ])
  1330. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
  1331. expect(requests[1]?.messages.at(1)?.content).toEqual([{ type: "text", text: "Changed context" }])
  1332. expect(yield* session.messages({ sessionID })).toHaveLength(2)
  1333. const { db } = yield* Database.Service
  1334. const updates = yield* db
  1335. .select({ data: EventTable.data })
  1336. .from(EventTable)
  1337. .where(eq(EventTable.type, "session.instructions.updated.2"))
  1338. .orderBy(asc(EventTable.seq))
  1339. .all()
  1340. .pipe(Effect.orDie)
  1341. expect(updates).toHaveLength(2)
  1342. expect(updates[0]?.data).toEqual({
  1343. sessionID,
  1344. delta: { "test/context": Instructions.hash("Initial context") },
  1345. })
  1346. expect(updates[1]?.data).toEqual({
  1347. sessionID,
  1348. delta: { "test/context": Instructions.hash("Changed context") },
  1349. })
  1350. yield* replaySessionProjection(sessionID)
  1351. expect(yield* session.messages({ sessionID })).toHaveLength(2)
  1352. }),
  1353. )
  1354. it.effect("uses the selected model family prompt when the agent does not override it", () =>
  1355. Effect.gen(function* () {
  1356. const session = yield* setup
  1357. currentModel = Model.make({ id: "gpt-5", provider: "openai", route: OpenAIChat.route })
  1358. yield* admit(session, "First")
  1359. response = reply.text("Done", "text-provider-prompt")
  1360. yield* session.resume(sessionID)
  1361. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([
  1362. expect.stringContaining("You are OpenCode, You and the user share the same workspace"),
  1363. "Initial context",
  1364. ])
  1365. }),
  1366. )
  1367. it.effect("uses the selected model family prompt when the agent system override is empty", () =>
  1368. Effect.gen(function* () {
  1369. const session = yield* setup
  1370. currentModel = Model.make({ id: "gpt-5", provider: "openai", route: OpenAIChat.route })
  1371. const agent = yield* AgentV2.Service
  1372. yield* agent.transform((editor) =>
  1373. editor.update(AgentV2.ID.make("build"), (agent) => {
  1374. agent.system = ""
  1375. agent.mode = "primary"
  1376. }),
  1377. )
  1378. yield* admit(session, "First")
  1379. response = reply.text("Done", "text-empty-agent-system")
  1380. yield* session.resume(sessionID)
  1381. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([
  1382. expect.stringContaining("You are OpenCode, You and the user share the same workspace"),
  1383. "Initial context",
  1384. ])
  1385. }),
  1386. )
  1387. it.effect("includes the effective default agent system before durable context", () =>
  1388. Effect.gen(function* () {
  1389. const session = yield* setup
  1390. const agent = yield* AgentV2.Service
  1391. yield* agent.transform((editor) =>
  1392. editor.update(AgentV2.ID.make("build"), (agent) => {
  1393. agent.system = "Build agent instructions"
  1394. agent.mode = "primary"
  1395. }),
  1396. )
  1397. yield* admit(session, "First")
  1398. response = reply.text("Done", "text-build")
  1399. yield* session.resume(sessionID)
  1400. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Build agent instructions", "Initial context"])
  1401. }),
  1402. )
  1403. it.effect("uses the configured default agent system for omitted-agent sessions", () =>
  1404. Effect.gen(function* () {
  1405. const session = yield* setup
  1406. const agent = yield* AgentV2.Service
  1407. yield* agent.transform((editor) => {
  1408. editor.update(AgentV2.ID.make("build"), (agent) => {
  1409. agent.system = "Build agent instructions"
  1410. agent.mode = "primary"
  1411. })
  1412. editor.update(AgentV2.ID.make("reviewer"), (agent) => {
  1413. agent.system = "Reviewer instructions"
  1414. agent.mode = "primary"
  1415. })
  1416. editor.default(AgentV2.ID.make("reviewer"))
  1417. })
  1418. yield* admit(session, "First")
  1419. response = reply.text("Done", "text-reviewer")
  1420. yield* session.resume(sessionID)
  1421. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"])
  1422. expect((yield* session.messages({ sessionID }))[0]).toMatchObject({ type: "assistant", agent: "reviewer" })
  1423. }),
  1424. )
  1425. it.effect("uses only the agent prompt and initial instructions as system parts", () =>
  1426. Effect.gen(function* () {
  1427. const session = yield* setup
  1428. const agent = yield* AgentV2.Service
  1429. yield* agent.transform((editor) =>
  1430. editor.update(AgentV2.ID.make("build"), (agent) => {
  1431. agent.system = "Build agent instructions"
  1432. agent.mode = "primary"
  1433. }),
  1434. )
  1435. yield* admit(session, "First")
  1436. response = reply.text("Done", "text-no-system")
  1437. yield* session.resume(sessionID)
  1438. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Build agent instructions", "Initial context"])
  1439. }),
  1440. )
  1441. it.effect("uses an explicitly selected non-build agent system", () =>
  1442. Effect.gen(function* () {
  1443. const session = yield* setup
  1444. const { db } = yield* Database.Service
  1445. const agent = yield* AgentV2.Service
  1446. yield* agent.transform((editor) =>
  1447. editor.update(AgentV2.ID.make("reviewer"), (agent) => {
  1448. agent.system = "Reviewer instructions"
  1449. agent.mode = "primary"
  1450. }),
  1451. )
  1452. yield* db
  1453. .update(SessionTable)
  1454. .set({ agent: "reviewer" })
  1455. .where(eq(SessionTable.id, sessionID))
  1456. .run()
  1457. .pipe(Effect.orDie)
  1458. yield* admit(session, "First")
  1459. response = reply.text("Done", "text-selected")
  1460. yield* session.resume(sessionID)
  1461. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"])
  1462. expect((yield* session.messages({ sessionID }))[0]).toMatchObject({ type: "assistant", agent: "reviewer" })
  1463. }),
  1464. )
  1465. it.effect("fails before the model request when the selected agent is unavailable", () =>
  1466. Effect.gen(function* () {
  1467. yield* setup
  1468. const { db } = yield* Database.Service
  1469. yield* db
  1470. .update(SessionTable)
  1471. .set({ agent: "explore" })
  1472. .where(eq(SessionTable.id, sessionID))
  1473. .run()
  1474. .pipe(Effect.orDie)
  1475. const session = yield* SessionV2.Service
  1476. yield* session.prompt({ sessionID, text: "Inspect files", resume: false })
  1477. requests.length = 0
  1478. response = []
  1479. const failure = yield* session.resume(sessionID).pipe(Effect.flip)
  1480. expect(failure).toMatchObject({
  1481. _tag: "Session.AgentNotFoundError",
  1482. sessionID,
  1483. agent: "explore",
  1484. })
  1485. expect(requests).toHaveLength(0)
  1486. }),
  1487. )
  1488. it.effect("waits for initial plugin readiness before constructing the model request", () =>
  1489. Effect.gen(function* () {
  1490. yield* setup
  1491. const release = yield* Deferred.make<void>()
  1492. pluginFlushHook = Deferred.await(release)
  1493. const session = yield* SessionV2.Service
  1494. yield* session.prompt({ sessionID, text: "Wait for plugins", resume: false })
  1495. requests.length = 0
  1496. response = []
  1497. const running = yield* session.resume(sessionID).pipe(Effect.forkChild({ startImmediately: true }))
  1498. yield* Effect.yieldNow
  1499. expect(requests).toHaveLength(0)
  1500. expect(running.pollUnsafe()).toBeUndefined()
  1501. yield* Deferred.succeed(release, undefined)
  1502. yield* Fiber.join(running)
  1503. expect(requests).toHaveLength(1)
  1504. }),
  1505. )
  1506. it.effect("updates selected-agent skill instructions after an agent switch", () =>
  1507. Effect.gen(function* () {
  1508. const session = yield* setup
  1509. const events = yield* EventV2.Service
  1510. const agents = yield* AgentV2.Service
  1511. yield* agents.transform((draft) =>
  1512. draft.update(AgentV2.ID.make("reviewer"), (agent) => {
  1513. agent.mode = "primary"
  1514. }),
  1515. )
  1516. skillBaselines.set(AgentV2.ID.make("build"), "Build skills")
  1517. yield* admit(session, "First")
  1518. yield* session.resume(sessionID)
  1519. skillBaselines.set(AgentV2.ID.make("reviewer"), "Reviewer skills")
  1520. yield* events.publish(SessionEvent.AgentSelected, {
  1521. sessionID,
  1522. agent: AgentV2.ID.make("reviewer"),
  1523. })
  1524. yield* admit(session, "Second")
  1525. yield* session.resume(sessionID)
  1526. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1527. [defaultSystem, "Initial context\n\nBuild skills"],
  1528. [defaultSystem, "Initial context\n\nBuild skills"],
  1529. ])
  1530. expect(systemTexts(requests[1]!)).toContainEqual(expect.stringContaining("Reviewer skills"))
  1531. }),
  1532. )
  1533. it.effect("keeps the sampled agent when selection changes during observation", () =>
  1534. Effect.gen(function* () {
  1535. const session = yield* setup
  1536. const events = yield* EventV2.Service
  1537. skillBaselines.set(AgentV2.ID.make("build"), "Build skills")
  1538. skillBaselines.set(AgentV2.ID.make("reviewer"), "Reviewer skills")
  1539. let switched = false
  1540. systemLoadHook = Effect.suspend(() => {
  1541. if (switched) return Effect.void
  1542. switched = true
  1543. return events
  1544. .publish(SessionEvent.AgentSelected, {
  1545. sessionID,
  1546. agent: AgentV2.ID.make("reviewer"),
  1547. })
  1548. .pipe(Effect.asVoid)
  1549. })
  1550. yield* admit(session, "First")
  1551. yield* session.resume(sessionID)
  1552. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1553. [defaultSystem, "Initial context\n\nBuild skills"],
  1554. ])
  1555. }),
  1556. )
  1557. it.effect("keeps the sampled model when selection changes during model resolution", () =>
  1558. Effect.gen(function* () {
  1559. const session = yield* setup
  1560. const events = yield* EventV2.Service
  1561. let switched = false
  1562. modelResolveHook = Effect.suspend(() => {
  1563. if (switched) return Effect.void
  1564. switched = true
  1565. return events
  1566. .publish(SessionEvent.ModelSelected, {
  1567. sessionID,
  1568. model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") },
  1569. })
  1570. .pipe(Effect.asVoid)
  1571. })
  1572. yield* admit(session, "First")
  1573. yield* session.resume(sessionID)
  1574. expect(requests.map((request) => request.model)).toEqual([model])
  1575. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1576. [defaultSystem, "Initial context"],
  1577. ])
  1578. }),
  1579. )
  1580. it.effect("admits removed context as a chronological System message", () =>
  1581. Effect.gen(function* () {
  1582. const session = yield* setup
  1583. yield* admit(session, "First")
  1584. yield* session.resume(sessionID)
  1585. systemRemoved = true
  1586. yield* admit(session, "Second")
  1587. yield* session.resume(sessionID)
  1588. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
  1589. expect(requests[1]?.messages.at(1)?.content).toEqual([
  1590. { type: "text", text: "System context source removed: test/context" },
  1591. ])
  1592. expect(yield* session.messages({ sessionID })).toHaveLength(2)
  1593. }),
  1594. )
  1595. it.effect("renders API context entries through add, change, and removal", () =>
  1596. Effect.gen(function* () {
  1597. const session = yield* setup
  1598. const contextEntries = yield* InstructionEntry.Service
  1599. yield* contextEntries.put({ sessionID, key: "deploy-target", value: "production" })
  1600. yield* admit(session, "First")
  1601. yield* session.resume(sessionID)
  1602. // String values render verbatim inside the initial tagged block.
  1603. expect(requests[0]?.system.map((part) => part.text)).toEqual([
  1604. defaultSystem,
  1605. ["Initial context", "", '<context key="deploy-target">', "production", "</context>"].join("\n"),
  1606. ])
  1607. // Non-string JSON pretty-prints; the change narrates as a System update.
  1608. yield* contextEntries.put({ sessionID, key: "deploy-target", value: { region: "us-east-1" } })
  1609. yield* admit(session, "Second")
  1610. yield* session.resume(sessionID)
  1611. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
  1612. expect(requests[1]?.messages.at(1)?.content).toEqual([
  1613. {
  1614. type: "text",
  1615. text: [
  1616. 'The context under "deploy-target" changed and supersedes the previous value:',
  1617. '<context key="deploy-target">',
  1618. "{",
  1619. ' "region": "us-east-1"',
  1620. "}",
  1621. "</context>",
  1622. ].join("\n"),
  1623. },
  1624. ])
  1625. expect(yield* contextEntries.list(sessionID)).toEqual([{ key: "deploy-target", value: { region: "us-east-1" } }])
  1626. // Deleting the row announces removal through the stored removal text.
  1627. yield* contextEntries.remove({ sessionID, key: "deploy-target" })
  1628. yield* admit(session, "Third")
  1629. yield* session.resume(sessionID)
  1630. expect(requests[2]?.messages.map((message) => message.role)).toEqual(["user", "system", "user", "system", "user"])
  1631. expect(requests[2]?.messages.at(-2)?.content).toEqual([
  1632. { type: "text", text: 'The context under "deploy-target" no longer applies. Disregard it.' },
  1633. ])
  1634. expect(yield* contextEntries.list(sessionID)).toEqual([])
  1635. }),
  1636. )
  1637. it.effect("retains JSON null API entries as values", () =>
  1638. Effect.gen(function* () {
  1639. const session = yield* setup
  1640. const entries = yield* InstructionEntry.Service
  1641. yield* entries.put({ sessionID, key: "nullable", value: "present" })
  1642. yield* admit(session, "First")
  1643. yield* session.resume(sessionID)
  1644. yield* entries.put({ sessionID, key: "nullable", value: null })
  1645. yield* admit(session, "Second")
  1646. yield* session.resume(sessionID)
  1647. expect(requests[1]?.messages.at(1)?.content).toEqual([
  1648. {
  1649. type: "text",
  1650. text: [
  1651. 'The context under "nullable" changed and supersedes the previous value:',
  1652. '<context key="nullable">',
  1653. "null",
  1654. "</context>",
  1655. ].join("\n"),
  1656. },
  1657. ])
  1658. expect(yield* entries.list(sessionID)).toEqual([{ key: "nullable", value: null }])
  1659. }),
  1660. )
  1661. it.effect("rejects API instruction entries larger than 8KB", () =>
  1662. Effect.gen(function* () {
  1663. yield* setup
  1664. const entries = yield* InstructionEntry.Service
  1665. const exit = yield* entries
  1666. .put({ sessionID, key: "oversized", value: "x".repeat(InstructionEntry.MaxValueBytes) })
  1667. .pipe(Effect.exit)
  1668. expect(Exit.isFailure(exit)).toBe(true)
  1669. if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(InstructionEntry.ValueTooLargeError)
  1670. expect(yield* entries.list(sessionID)).toEqual([])
  1671. }),
  1672. )
  1673. it.effect("keeps initial instructions and chronological updates after a model switch", () =>
  1674. Effect.gen(function* () {
  1675. const session = yield* setup
  1676. const events = yield* EventV2.Service
  1677. yield* admit(session, "First")
  1678. yield* session.resume(sessionID)
  1679. systemBaseline = "Changed context"
  1680. yield* admit(session, "Second")
  1681. yield* session.resume(sessionID)
  1682. yield* events.publish(SessionEvent.ModelSelected, {
  1683. sessionID,
  1684. model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") },
  1685. })
  1686. systemBaseline = "Replacement context"
  1687. yield* admit(session, "Third")
  1688. yield* session.resume(sessionID)
  1689. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1690. [defaultSystem, "Initial context"],
  1691. [defaultSystem, "Initial context"],
  1692. [defaultSystem, "Initial context"],
  1693. ])
  1694. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
  1695. expect(requests[2]?.messages.filter((message) => message.role === "system")).toHaveLength(2)
  1696. expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([
  1697. "user",
  1698. "user",
  1699. "model-switched",
  1700. "user",
  1701. ])
  1702. yield* replaySessionProjection(sessionID)
  1703. expect(yield* session.messages({ sessionID })).toHaveLength(4)
  1704. yield* admit(session, "Fourth")
  1705. yield* session.resume(sessionID)
  1706. }),
  1707. )
  1708. it.effect("preserves instruction values while a source is temporarily unavailable", () =>
  1709. Effect.gen(function* () {
  1710. const session = yield* setup
  1711. const events = yield* EventV2.Service
  1712. yield* admit(session, "First")
  1713. yield* session.resume(sessionID)
  1714. yield* events.publish(SessionEvent.ModelSelected, {
  1715. sessionID,
  1716. model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") },
  1717. })
  1718. systemUnavailable = true
  1719. yield* admit(session, "Second")
  1720. yield* session.resume(sessionID)
  1721. systemUnavailable = false
  1722. systemBaseline = "Replacement context"
  1723. yield* admit(session, "Third")
  1724. yield* session.resume(sessionID)
  1725. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1726. [defaultSystem, "Initial context"],
  1727. [defaultSystem, "Initial context"],
  1728. [defaultSystem, "Initial context"],
  1729. ])
  1730. }),
  1731. )
  1732. it.effect("moves the epoch at compaction and narrates later changes", () =>
  1733. Effect.gen(function* () {
  1734. const session = yield* setup
  1735. const events = yield* EventV2.Service
  1736. yield* admit(session, "First")
  1737. yield* session.resume(sessionID)
  1738. yield* events.publish(SessionEvent.Compaction.Started, {
  1739. sessionID,
  1740. reason: "manual",
  1741. recent: "",
  1742. })
  1743. yield* events.publish(SessionEvent.Compaction.Ended, {
  1744. sessionID,
  1745. reason: "manual",
  1746. text: "summary",
  1747. recent: "",
  1748. })
  1749. systemBaseline = "Replacement context"
  1750. yield* admit(session, "Second")
  1751. yield* session.resume(sessionID)
  1752. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1753. [defaultSystem, "Initial context"],
  1754. [defaultSystem, "Initial context"],
  1755. ])
  1756. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
  1757. expect(requests[1]?.messages.at(1)?.content).toEqual([{ type: "text", text: "Replacement context" }])
  1758. yield* replaySessionProjection(sessionID)
  1759. yield* admit(session, "Third")
  1760. yield* session.resume(sessionID)
  1761. }),
  1762. )
  1763. it.effect("runs one durable compaction barrier after tool settlement and before later inputs", () =>
  1764. Effect.gen(function* () {
  1765. const session = yield* setup
  1766. currentModel = recoveryModel
  1767. streamGate = yield* Deferred.make<void>()
  1768. streamStarted = yield* Deferred.make<void>()
  1769. responses = [
  1770. reply.tool("call-active", "echo", { text: "active" }),
  1771. [LLMEvent.textDelta({ id: "summary", text: "durable summary" })],
  1772. reply.text("Steer complete", "text-steer"),
  1773. reply.text("Queue complete", "text-queue"),
  1774. ]
  1775. yield* admit(session, "Active work")
  1776. const active = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1777. yield* Deferred.await(streamStarted)
  1778. const first = yield* session.compact({ sessionID })
  1779. const second = yield* session.compact({ sessionID })
  1780. expect(second.id).toBe(first.id)
  1781. expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toMatchObject({
  1782. id: first.id,
  1783. })
  1784. expect((yield* session.messages({ sessionID })).find((message) => message.id === first.id)).toBeUndefined()
  1785. yield* admit(session, "Steer after compaction")
  1786. yield* session.synthetic({ sessionID, text: "Completion after compaction", resume: false })
  1787. yield* session.prompt({
  1788. sessionID,
  1789. text: "Queue after compaction",
  1790. delivery: "queue",
  1791. resume: false,
  1792. })
  1793. expect(yield* SessionPending.has((yield* Database.Service).db, sessionID, "steer")).toBe(false)
  1794. yield* Deferred.succeed(streamGate, undefined)
  1795. yield* Fiber.join(active)
  1796. expect(requests).toHaveLength(4)
  1797. expect(userTexts(requests[1])[0]).toContain("Create a new anchored summary")
  1798. expect(userTexts(requests[2])).toContain("Steer after compaction")
  1799. expect(userTexts(requests[2])).toContain("Completion after compaction")
  1800. expect(userTexts(requests[3])).toContain("Queue after compaction")
  1801. expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
  1802. expect((yield* session.messages({ sessionID })).find((message) => message.id === first.id)).toMatchObject({
  1803. type: "compaction",
  1804. status: "completed",
  1805. summary: "durable summary",
  1806. })
  1807. }),
  1808. )
  1809. it.effect("releases queued prompts when durable compaction fails", () =>
  1810. Effect.gen(function* () {
  1811. const session = yield* setup
  1812. currentModel = recoveryModel
  1813. streamGate = yield* Deferred.make<void>()
  1814. streamStarted = yield* Deferred.make<void>()
  1815. responses = [
  1816. reply.text("Active complete", "text-active-failure"),
  1817. [],
  1818. reply.text("Continued", "text-after-failure"),
  1819. ]
  1820. yield* admit(session, "Active work")
  1821. const active = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1822. yield* Deferred.await(streamStarted)
  1823. const compaction = yield* session.compact({ sessionID })
  1824. yield* session.prompt({
  1825. sessionID,
  1826. text: "Continue after failure",
  1827. delivery: "queue",
  1828. resume: false,
  1829. })
  1830. yield* Deferred.succeed(streamGate, undefined)
  1831. yield* Fiber.join(active)
  1832. expect(requests).toHaveLength(3)
  1833. expect(userTexts(requests[2])).toContain("Continue after failure")
  1834. expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
  1835. expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
  1836. type: "compaction",
  1837. status: "failed",
  1838. })
  1839. expect(
  1840. (yield* recordedEventTypes(sessionID)).filter(
  1841. (type) => type === EventV2.versionedType(SessionEvent.Compaction.Failed.type, 1),
  1842. ),
  1843. ).toHaveLength(1)
  1844. }),
  1845. )
  1846. it.effect("explains when manual compaction has no history", () =>
  1847. Effect.gen(function* () {
  1848. yield* setup
  1849. const session = yield* SessionV2.Service
  1850. const compaction = yield* session.compact({ sessionID })
  1851. modelResolveHook = Effect.die("model resolution should not run")
  1852. yield* session.resume(sessionID)
  1853. expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
  1854. expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
  1855. type: "compaction",
  1856. status: "failed",
  1857. reason: "manual",
  1858. error: { type: "compaction.unavailable", message: "Nothing to compact yet" },
  1859. })
  1860. expect(
  1861. (yield* recordedEventTypes(sessionID)).filter(
  1862. (type) => type === EventV2.versionedType(SessionEvent.Compaction.Failed.type, 1),
  1863. ),
  1864. ).toHaveLength(1)
  1865. }),
  1866. )
  1867. it.effect("manually compacts when the model has no context limit", () =>
  1868. Effect.gen(function* () {
  1869. const session = yield* setup
  1870. response = reply.text("Earlier answer", "text-manual-unknown-history")
  1871. yield* admit(session, "Earlier question")
  1872. yield* session.resume(sessionID)
  1873. requests.length = 0
  1874. response = reply.text("Manual summary", "text-manual-unknown-summary")
  1875. const compaction = yield* session.compact({ sessionID })
  1876. yield* session.resume(sessionID)
  1877. expect(requests).toHaveLength(1)
  1878. expect(userTexts(requests[0])[0]).toContain("Earlier question")
  1879. expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
  1880. type: "compaction",
  1881. status: "completed",
  1882. summary: "Manual summary",
  1883. })
  1884. }),
  1885. )
  1886. it.effect("preserves provider errors from manual compaction", () =>
  1887. Effect.gen(function* () {
  1888. const session = yield* setup
  1889. response = reply.text("Earlier answer", "text-manual-provider-history")
  1890. yield* admit(session, "Earlier question")
  1891. yield* session.resume(sessionID)
  1892. response = [LLMEvent.providerError({ message: "summary unavailable" })]
  1893. const compaction = yield* session.compact({ sessionID })
  1894. yield* session.resume(sessionID)
  1895. expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
  1896. type: "compaction",
  1897. status: "failed",
  1898. error: { type: "provider.error", message: "summary unavailable" },
  1899. })
  1900. }),
  1901. )
  1902. it.effect("preserves typed provider failures from manual compaction", () =>
  1903. Effect.gen(function* () {
  1904. const session = yield* setup
  1905. response = reply.text("Earlier answer", "text-manual-failure-history")
  1906. yield* admit(session, "Earlier question")
  1907. yield* session.resume(sessionID)
  1908. responseStream = Stream.fail(providerUnavailable())
  1909. const compaction = yield* session.compact({ sessionID })
  1910. yield* session.resume(sessionID)
  1911. expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
  1912. type: "compaction",
  1913. status: "failed",
  1914. error: { type: "provider.transport", message: "Provider unavailable" },
  1915. })
  1916. }),
  1917. )
  1918. it.effect("records cancelled manual compaction without surfacing an internal failure", () =>
  1919. Effect.gen(function* () {
  1920. const session = yield* setup
  1921. response = reply.text("Earlier answer", "text-manual-interrupt-history")
  1922. yield* admit(session, "Earlier question")
  1923. yield* session.resume(sessionID)
  1924. const streamed = yield* Deferred.make<void>()
  1925. const partial = fragmentFixture("text", "text-manual-interrupt-summary", ["Partial summary"])
  1926. responseStream = Stream.concat(
  1927. Stream.fromIterable(partial.partialEvents),
  1928. Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)),
  1929. )
  1930. const compaction = yield* session.compact({ sessionID })
  1931. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1932. yield* Deferred.await(streamed)
  1933. yield* session.interrupt(sessionID)
  1934. yield* Fiber.await(run)
  1935. expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
  1936. expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
  1937. type: "compaction",
  1938. status: "failed",
  1939. reason: "manual",
  1940. error: { type: "aborted", message: "Compaction cancelled" },
  1941. })
  1942. }),
  1943. )
  1944. it.effect("settles an admitted manual compaction when pre-start resolution throws", () =>
  1945. Effect.gen(function* () {
  1946. const session = yield* setup
  1947. response = reply.text("Earlier answer", "text-manual-resolution-history")
  1948. yield* admit(session, "Earlier question")
  1949. yield* session.resume(sessionID)
  1950. const compaction = yield* session.compact({ sessionID })
  1951. modelResolveHook = Effect.die("model resolution failed")
  1952. expect(yield* Effect.exit(session.resume(sessionID))).toMatchObject({ _tag: "Failure" })
  1953. expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
  1954. expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
  1955. type: "compaction",
  1956. status: "failed",
  1957. reason: "manual",
  1958. })
  1959. expect(
  1960. (yield* recordedEventTypes(sessionID)).filter(
  1961. (type) => type === EventV2.versionedType(SessionEvent.Compaction.Failed.type, 1),
  1962. ),
  1963. ).toHaveLength(1)
  1964. }),
  1965. )
  1966. it.effect("automatically compacts into a completed summary and retained recent turn", () =>
  1967. Effect.gen(function* () {
  1968. const session = yield* setup
  1969. response = reply.textWithUsage("Earlier answer", "text-first", 3_950)
  1970. yield* admit(session, "Earlier question ".repeat(180))
  1971. yield* session.resume(sessionID)
  1972. currentModel = compactModel
  1973. requests.length = 0
  1974. responses = [
  1975. reply.text("## Objective\n- Preserve the task", "text-summary"),
  1976. reply.textWithUsage("Continued", "text-final", 3_950),
  1977. ]
  1978. yield* admit(session, "Recent exact request ".repeat(180))
  1979. yield* session.resume(sessionID)
  1980. expect(requests).toHaveLength(2)
  1981. expect(userTexts(requests[0])[0]).toContain("## Objective")
  1982. expect(userTexts(requests[1])).toHaveLength(1)
  1983. expect(userTexts(requests[1])[0]).toContain("<summary>\n## Objective\n- Preserve the task\n</summary>")
  1984. expect(userTexts(requests[1])[0]).toContain(`[User]: ${"Recent exact request ".repeat(180)}`)
  1985. const context = yield* (yield* SessionStore.Service).context(sessionID)
  1986. expect(context.map((message) => message.type)).toEqual(["compaction", "assistant"])
  1987. expect(context[0]).toMatchObject({
  1988. type: "compaction",
  1989. summary: "## Objective\n- Preserve the task",
  1990. })
  1991. requests.length = 0
  1992. executions.length = 0
  1993. responses = [
  1994. reply.text("## Objective\n- Preserve the updated task", "text-summary-2"),
  1995. reply.text("Continued again", "text-final-2"),
  1996. ]
  1997. yield* admit(session, "Newest exact request ".repeat(180))
  1998. yield* session.resume(sessionID)
  1999. expect(requests).toHaveLength(2)
  2000. expect(userTexts(requests[0])[0]).toContain(
  2001. "<previous-summary>\n## Objective\n- Preserve the task\n</previous-summary>",
  2002. )
  2003. expect(userTexts(requests[0])[0]).toContain("Recent exact request")
  2004. expect((yield* (yield* SessionStore.Service).context(sessionID))[0]).toMatchObject({
  2005. type: "compaction",
  2006. summary: "## Objective\n- Preserve the updated task",
  2007. })
  2008. }),
  2009. )
  2010. it.effect("does not compact immediately when the advertised output limit fills the context", () =>
  2011. Effect.gen(function* () {
  2012. const session = yield* setup
  2013. currentModel = fullOutputModel
  2014. response = reply.textWithUsage("Earlier answer", "text-full-output-first", 9_500)
  2015. yield* admit(session, "Earlier question")
  2016. yield* session.resume(sessionID)
  2017. requests.length = 0
  2018. response = reply.text("Continued", "text-full-output-final")
  2019. yield* admit(session, "Continue")
  2020. yield* session.resume(sessionID)
  2021. expect(requests).toHaveLength(1)
  2022. expect(userTexts(requests[0])).toContain("Continue")
  2023. expect(yield* session.context(sessionID)).not.toContainEqual(expect.objectContaining({ type: "compaction" }))
  2024. }),
  2025. )
  2026. it.effect("stops after required automatic compaction fails", () =>
  2027. Effect.gen(function* () {
  2028. const session = yield* setup
  2029. response = reply.textWithUsage("Earlier answer", "text-before-failed-compaction", 3_950)
  2030. yield* admit(session, "Earlier question ".repeat(180))
  2031. yield* session.resume(sessionID)
  2032. currentModel = compactModel
  2033. requests.length = 0
  2034. responses = [
  2035. [LLMEvent.providerError({ message: "Unsupported parameter: max_output_tokens" })],
  2036. reply.text("Must not run", "text-after-failed-compaction"),
  2037. ]
  2038. yield* admit(session, "Recent exact request ".repeat(180))
  2039. expect(yield* Effect.exit(session.resume(sessionID))).toMatchObject({ _tag: "Failure" })
  2040. expect(requests).toHaveLength(1)
  2041. expect(requests[0]?.generation).toBeUndefined()
  2042. expect(yield* session.context(sessionID)).toContainEqual(
  2043. expect.objectContaining({
  2044. type: "compaction",
  2045. status: "failed",
  2046. reason: "auto",
  2047. error: expect.objectContaining({ message: "Unsupported parameter: max_output_tokens" }),
  2048. }),
  2049. )
  2050. }),
  2051. )
  2052. it.effect("forces one compaction and retries after provider context overflow", () =>
  2053. Effect.gen(function* () {
  2054. const session = yield* setupOverflowRecovery
  2055. responses = [
  2056. [
  2057. LLMEvent.stepStart({ index: 0 }),
  2058. LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
  2059. ],
  2060. reply.text("## Objective\n- Recover overflow", "text-summary"),
  2061. reply.text("Recovered", "text-final"),
  2062. ]
  2063. yield* admit(session, "Continue")
  2064. yield* session.resume(sessionID)
  2065. expect(requests).toHaveLength(3)
  2066. expect(userTexts(requests[1])[0]).toContain("## Objective")
  2067. expect(userTexts(requests[2])[0]).toContain("<summary>\n## Objective\n- Recover overflow\n</summary>")
  2068. expect(yield* session.context(sessionID)).toMatchObject([
  2069. { type: "compaction", summary: "## Objective\n- Recover overflow" },
  2070. { type: "assistant", finish: "stop" },
  2071. ])
  2072. yield* replaySessionProjection(sessionID)
  2073. expect(yield* session.context(sessionID)).toMatchObject([
  2074. { type: "compaction" },
  2075. { type: "assistant", finish: "stop" },
  2076. ])
  2077. }),
  2078. )
  2079. it.effect("recovers from provider context overflow without a configured context limit", () =>
  2080. Effect.gen(function* () {
  2081. const session = yield* setupOverflowRecovery
  2082. currentModel = model
  2083. responses = [
  2084. [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })],
  2085. reply.text("## Objective\n- Recover unknown limit", "text-summary-unknown-limit"),
  2086. reply.text("Recovered", "text-final-unknown-limit"),
  2087. ]
  2088. yield* admit(session, "Continue")
  2089. yield* session.resume(sessionID)
  2090. expect(requests).toHaveLength(3)
  2091. expect(yield* session.context(sessionID)).toMatchObject([
  2092. { type: "compaction", summary: "## Objective\n- Recover unknown limit" },
  2093. { type: "assistant", finish: "stop" },
  2094. ])
  2095. }),
  2096. )
  2097. it.effect("recovers from provider context overflow despite an undersized configured context limit", () =>
  2098. Effect.gen(function* () {
  2099. const session = yield* setupOverflowRecovery
  2100. currentModel = undersizedContextModel
  2101. responses = [
  2102. [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })],
  2103. reply.text("## Objective\n- Recover undersized limit", "text-summary-undersized-limit"),
  2104. reply.text("Recovered", "text-final-undersized-limit"),
  2105. ]
  2106. yield* admit(session, "Continue")
  2107. yield* session.resume(sessionID)
  2108. expect(requests).toHaveLength(3)
  2109. expect(yield* session.context(sessionID)).toMatchObject([
  2110. { type: "compaction", summary: "## Objective\n- Recover undersized limit" },
  2111. { type: "assistant", finish: "stop" },
  2112. ])
  2113. }),
  2114. )
  2115. it.effect("persists a second context overflow after one recovery", () =>
  2116. Effect.gen(function* () {
  2117. const session = yield* setupOverflowRecovery
  2118. const overflow = () => [
  2119. LLMEvent.stepStart({ index: 0 }),
  2120. LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
  2121. ]
  2122. responses = [overflow(), reply.text("## Objective\n- Recover once", "text-summary"), overflow()]
  2123. yield* admit(session, "Continue")
  2124. expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long")
  2125. expect(requests).toHaveLength(3)
  2126. expect(yield* session.context(sessionID)).toMatchObject([
  2127. { type: "compaction" },
  2128. { type: "assistant", finish: "error", error: { message: "prompt too long" } },
  2129. ])
  2130. }),
  2131. )
  2132. it.effect("recovers once from a raw context overflow failure", () =>
  2133. Effect.gen(function* () {
  2134. const session = yield* setupOverflowRecovery
  2135. responseStream = Stream.fail(
  2136. new LLMError({
  2137. module: "test",
  2138. method: "stream",
  2139. reason: new InvalidRequestReason({
  2140. message: "prompt too long",
  2141. classification: "context-overflow",
  2142. }),
  2143. }),
  2144. )
  2145. responses = [
  2146. reply.text("## Objective\n- Recover raw overflow", "text-summary"),
  2147. reply.text("Recovered", "text-final"),
  2148. ]
  2149. yield* admit(session, "Continue")
  2150. yield* session.resume(sessionID)
  2151. expect(requests).toHaveLength(3)
  2152. expect(yield* session.context(sessionID)).toMatchObject([
  2153. { type: "compaction", summary: "## Objective\n- Recover raw overflow" },
  2154. { type: "assistant", finish: "stop" },
  2155. ])
  2156. }),
  2157. )
  2158. it.effect("publishes the original overflow when recovery summarization fails", () =>
  2159. Effect.gen(function* () {
  2160. const session = yield* setupOverflowRecovery
  2161. responses = [
  2162. [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })],
  2163. [LLMEvent.providerError({ message: "summary unavailable" })],
  2164. ]
  2165. yield* admit(session, "Continue")
  2166. expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long")
  2167. expect(requests).toHaveLength(2)
  2168. const context = yield* session.context(sessionID)
  2169. expect(context).toContainEqual(
  2170. expect.objectContaining({
  2171. type: "compaction",
  2172. status: "failed",
  2173. reason: "auto",
  2174. error: { type: "provider.error", message: "summary unavailable" },
  2175. }),
  2176. )
  2177. expect(context.slice(-3)).toMatchObject([
  2178. { type: "user", text: "Continue" },
  2179. { type: "compaction", status: "failed", reason: "auto" },
  2180. { type: "assistant", finish: "error", error: { message: "prompt too long" } },
  2181. ])
  2182. }),
  2183. )
  2184. it.effect("interrupts overflow recovery while the summary provider is running", () =>
  2185. Effect.gen(function* () {
  2186. const session = yield* setupOverflowRecovery
  2187. responses = [
  2188. [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })],
  2189. reply.text("## Objective\n- Interrupted", "text-summary"),
  2190. ]
  2191. const firstGate = yield* Deferred.make<void>()
  2192. const summaryGate = yield* Deferred.make<void>()
  2193. streamGate = firstGate
  2194. yield* admit(session, "Continue")
  2195. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2196. while (requests.length < 1) yield* Effect.yieldNow
  2197. streamGate = summaryGate
  2198. yield* Deferred.succeed(firstGate, undefined)
  2199. while (requests.length < 2) yield* Effect.yieldNow
  2200. yield* session.interrupt(sessionID)
  2201. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  2202. streamGate = undefined
  2203. expect(requests).toHaveLength(2)
  2204. expect(yield* session.context(sessionID)).toContainEqual(
  2205. expect.objectContaining({ type: "compaction", status: "failed", reason: "auto" }),
  2206. )
  2207. }),
  2208. )
  2209. it.effect("uses epoch values after compaction while a source is unavailable", () =>
  2210. Effect.gen(function* () {
  2211. const session = yield* setup
  2212. const events = yield* EventV2.Service
  2213. yield* admit(session, "First")
  2214. yield* session.resume(sessionID)
  2215. systemBaseline = "Changed context"
  2216. yield* admit(session, "Second")
  2217. yield* session.resume(sessionID)
  2218. yield* events.publish(SessionEvent.Compaction.Started, {
  2219. sessionID,
  2220. reason: "manual",
  2221. recent: "",
  2222. })
  2223. yield* events.publish(SessionEvent.Compaction.Ended, {
  2224. sessionID,
  2225. reason: "manual",
  2226. text: "summary",
  2227. recent: "",
  2228. })
  2229. systemUnavailable = true
  2230. yield* admit(session, "Third")
  2231. yield* session.resume(sessionID)
  2232. // Compaction already moved current values into the new epoch before the unavailable read.
  2233. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([defaultSystem, "Changed context"])
  2234. expect(systemTexts(requests.at(-1)!)).not.toContain("Changed context")
  2235. }),
  2236. )
  2237. it.effect("projects reasoning and tool events without executing or continuing tools", () =>
  2238. Effect.gen(function* () {
  2239. const session = yield* setup
  2240. yield* admit(session, "Use tools")
  2241. response = [
  2242. LLMEvent.stepStart({ index: 0 }),
  2243. LLMEvent.reasoningStart({ id: "reasoning-1" }),
  2244. LLMEvent.reasoningDelta({ id: "reasoning-1", text: "Think" }),
  2245. LLMEvent.reasoningEnd({ id: "reasoning-1" }),
  2246. LLMEvent.toolInputStart({ id: "call-error", name: "write" }),
  2247. LLMEvent.toolInputDelta({ id: "call-error", name: "write", text: '{"path":"README.md"}' }),
  2248. LLMEvent.toolInputEnd({ id: "call-error", name: "write" }),
  2249. LLMEvent.toolCall({ id: "call-error", name: "write", input: { path: "README.md" }, providerExecuted: true }),
  2250. LLMEvent.toolError({ id: "call-error", name: "write", message: "Denied" }),
  2251. LLMEvent.toolResult({ id: "call-error", name: "write", result: { type: "error", value: "Denied" } }),
  2252. LLMEvent.toolCall({
  2253. id: "call-provider",
  2254. name: "web_search",
  2255. input: { query: "hello" },
  2256. providerExecuted: true,
  2257. providerMetadata: { openai: { source: "provider" } },
  2258. }),
  2259. LLMEvent.toolResult({
  2260. id: "call-provider",
  2261. name: "web_search",
  2262. result: {
  2263. type: "content",
  2264. value: [
  2265. { type: "text", text: "Hello" },
  2266. { type: "file", uri: "data:image/png;base64,aGVsbG8=", mime: "image/png", name: "hello.png" },
  2267. ],
  2268. },
  2269. providerExecuted: true,
  2270. providerMetadata: { openai: { source: "provider" } },
  2271. }),
  2272. LLMEvent.stepFinish({
  2273. index: 0,
  2274. reason: { normalized: "tool-calls" },
  2275. usage: {
  2276. inputTokens: 10,
  2277. nonCachedInputTokens: 8,
  2278. outputTokens: 4,
  2279. reasoningTokens: 1,
  2280. cacheReadInputTokens: 2,
  2281. },
  2282. }),
  2283. LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
  2284. ]
  2285. yield* session.resume(sessionID)
  2286. expect(requests).toHaveLength(1)
  2287. expect(requests[0]?.tools.map((tool) => tool.name)).toEqual(["defect", "echo", "storefail"])
  2288. expect(yield* session.context(sessionID)).toMatchObject([
  2289. { type: "user", text: "Use tools" },
  2290. {
  2291. type: "assistant",
  2292. finish: "tool-calls",
  2293. cost: 0,
  2294. tokens: { input: 8, output: 3, reasoning: 1, cache: { read: 2, write: 0 } },
  2295. content: [
  2296. { type: "reasoning", text: "Think" },
  2297. {
  2298. type: "tool",
  2299. id: "call-error",
  2300. name: "write",
  2301. state: {
  2302. status: "error",
  2303. input: { path: "README.md" },
  2304. error: { type: "tool.execution", message: "Denied" },
  2305. },
  2306. },
  2307. {
  2308. type: "tool",
  2309. id: "call-provider",
  2310. name: "web_search",
  2311. executed: true,
  2312. providerState: { source: "provider" },
  2313. providerResultState: { source: "provider" },
  2314. state: {
  2315. status: "completed",
  2316. input: { query: "hello" },
  2317. content: [
  2318. { type: "text", text: "Hello" },
  2319. { type: "file", mime: "image/png", uri: "data:image/png;base64,aGVsbG8=", name: "hello.png" },
  2320. ],
  2321. },
  2322. },
  2323. ],
  2324. },
  2325. ])
  2326. }),
  2327. )
  2328. it.effect("continues with reloaded history after durably settling one local tool call", () =>
  2329. Effect.gen(function* () {
  2330. const session = yield* setup
  2331. yield* admit(session, "Echo this")
  2332. responses = [reply.tool("call-echo", "echo", { text: "hello" }), reply.text("Done", "text-final")]
  2333. yield* session.resume(sessionID)
  2334. expect(requests).toHaveLength(2)
  2335. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
  2336. expect(authorizations).toMatchObject([{ sessionID, callID: "call-echo" }])
  2337. expect(executions).toEqual(["hello"])
  2338. const context = yield* session.context(sessionID)
  2339. expect(context).toMatchObject([
  2340. { type: "user", text: "Echo this" },
  2341. {
  2342. type: "assistant",
  2343. finish: "tool-calls",
  2344. content: [
  2345. {
  2346. type: "tool",
  2347. id: "call-echo",
  2348. name: "echo",
  2349. state: {
  2350. status: "completed",
  2351. input: { text: "hello" },
  2352. content: [{ type: "text", text: "hello" }],
  2353. },
  2354. },
  2355. ],
  2356. },
  2357. { type: "assistant", finish: "stop", content: [{ type: "text", text: "Done" }] },
  2358. ])
  2359. const assistant = requireAssistant(context)
  2360. expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([
  2361. "session.step.started.1",
  2362. "session.tool.called.1",
  2363. "session.tool.success.2",
  2364. "session.step.ended.1",
  2365. ])
  2366. }),
  2367. )
  2368. it.effect("reloads a model switch before a tool-driven continuation step", () =>
  2369. Effect.gen(function* () {
  2370. const session = yield* setup
  2371. const events = yield* EventV2.Service
  2372. yield* admit(session, "Echo this")
  2373. responses = [reply.tool("call-echo", "echo", { text: "hello" }), reply.stop()]
  2374. toolExecutionGate = yield* Deferred.make<void>()
  2375. toolExecutionsStarted = yield* Deferred.make<void>()
  2376. toolExecutionsReady = 1
  2377. const run = yield* Effect.forkChild(session.resume(sessionID))
  2378. yield* Deferred.await(toolExecutionsStarted)
  2379. yield* events.publish(SessionEvent.ModelSelected, {
  2380. sessionID,
  2381. model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") },
  2382. })
  2383. systemBaseline = "Replacement context"
  2384. yield* Deferred.succeed(toolExecutionGate, undefined)
  2385. yield* Fiber.join(run)
  2386. expect(requests.map((request) => request.model)).toEqual([model, replacementModel])
  2387. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  2388. [defaultSystem, "Initial context"],
  2389. [defaultSystem, "Initial context"],
  2390. ])
  2391. expect(systemTexts(requests[1]!)).toContain("Replacement context")
  2392. }),
  2393. )
  2394. it.effect("restores durable reasoning provider metadata in the next request", () =>
  2395. Effect.gen(function* () {
  2396. const session = yield* setup
  2397. yield* admit(session, "Think first")
  2398. response = [
  2399. LLMEvent.stepStart({ index: 0 }),
  2400. LLMEvent.reasoningStart({ id: "reasoning-anthropic" }),
  2401. LLMEvent.reasoningDelta({ id: "reasoning-anthropic", text: "Signed thought" }),
  2402. LLMEvent.reasoningEnd({
  2403. id: "reasoning-anthropic",
  2404. providerMetadata: { openai: { signature: "sig_1" }, anthropic: { ignored: true } },
  2405. }),
  2406. LLMEvent.reasoningStart({
  2407. id: "reasoning-openai",
  2408. providerMetadata: {
  2409. openai: { itemId: "rs_1", reasoningEncryptedContent: null },
  2410. anthropic: { ignored: true },
  2411. },
  2412. }),
  2413. LLMEvent.reasoningDelta({ id: "reasoning-openai", text: "Encrypted thought" }),
  2414. LLMEvent.reasoningEnd({
  2415. id: "reasoning-openai",
  2416. providerMetadata: {
  2417. openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" },
  2418. anthropic: { ignored: true },
  2419. },
  2420. }),
  2421. LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
  2422. LLMEvent.finish({ reason: { normalized: "stop" } }),
  2423. ]
  2424. yield* session.resume(sessionID)
  2425. yield* replaySessionProjection(sessionID)
  2426. expect(yield* session.context(sessionID)).toMatchObject([
  2427. { type: "user", text: "Think first" },
  2428. {
  2429. type: "assistant",
  2430. content: [
  2431. {
  2432. type: "reasoning",
  2433. text: "Signed thought",
  2434. state: { signature: "sig_1" },
  2435. },
  2436. {
  2437. type: "reasoning",
  2438. text: "Encrypted thought",
  2439. state: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" },
  2440. },
  2441. ],
  2442. },
  2443. ])
  2444. yield* admit(session, "Continue")
  2445. response = []
  2446. yield* session.resume(sessionID)
  2447. expect(requests[1]?.messages[1]?.content).toEqual([
  2448. {
  2449. type: "reasoning",
  2450. text: "Signed thought",
  2451. providerMetadata: { openai: { signature: "sig_1" } },
  2452. },
  2453. {
  2454. type: "reasoning",
  2455. text: "Encrypted thought",
  2456. providerMetadata: { openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" } },
  2457. },
  2458. ])
  2459. }),
  2460. )
  2461. it.effect("restores durable text provider metadata in the next request", () =>
  2462. Effect.gen(function* () {
  2463. const session = yield* setup
  2464. yield* admit(session, "Check first")
  2465. response = [
  2466. LLMEvent.stepStart({ index: 0 }),
  2467. LLMEvent.textStart({ id: "commentary", providerMetadata: { openai: { phase: "commentary" } } }),
  2468. LLMEvent.textDelta({ id: "commentary", text: "Checking." }),
  2469. LLMEvent.textEnd({
  2470. id: "commentary",
  2471. providerMetadata: { openai: { phase: "commentary" }, anthropic: { ignored: true } },
  2472. }),
  2473. LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
  2474. LLMEvent.finish({ reason: { normalized: "stop" } }),
  2475. ]
  2476. yield* session.resume(sessionID)
  2477. yield* replaySessionProjection(sessionID)
  2478. expect(yield* session.context(sessionID)).toMatchObject([
  2479. { type: "user", text: "Check first" },
  2480. {
  2481. type: "assistant",
  2482. content: [{ type: "text", text: "Checking.", state: { phase: "commentary" } }],
  2483. },
  2484. ])
  2485. yield* admit(session, "Continue")
  2486. response = []
  2487. yield* session.resume(sessionID)
  2488. expect(requests[1]?.messages[1]?.content).toEqual([
  2489. {
  2490. type: "text",
  2491. text: "Checking.",
  2492. providerMetadata: { openai: { phase: "commentary" } },
  2493. },
  2494. ])
  2495. }),
  2496. )
  2497. it.effect("replays durable provider-executed tool results inline in the next request", () =>
  2498. Effect.gen(function* () {
  2499. const session = yield* setup
  2500. yield* admit(session, "Search first")
  2501. response = [
  2502. LLMEvent.stepStart({ index: 0 }),
  2503. LLMEvent.toolCall({
  2504. id: "hosted-search",
  2505. name: "web_search",
  2506. input: { query: "Effect" },
  2507. providerExecuted: true,
  2508. providerMetadata: { openai: { itemId: "hosted-search" }, fake: { ignored: true } },
  2509. }),
  2510. LLMEvent.toolResult({
  2511. id: "hosted-search",
  2512. name: "web_search",
  2513. result: { type: "json", value: [{ title: "Effect" }] },
  2514. providerExecuted: true,
  2515. providerMetadata: { openai: { blockType: "web_search_tool_result" }, anthropic: { ignored: true } },
  2516. }),
  2517. LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
  2518. LLMEvent.finish({ reason: { normalized: "stop" } }),
  2519. ]
  2520. yield* session.resume(sessionID)
  2521. yield* replaySessionProjection(sessionID)
  2522. yield* admit(session, "Continue")
  2523. response = []
  2524. yield* session.resume(sessionID)
  2525. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "user"])
  2526. expect(requests[1]?.messages[1]?.content).toMatchObject([
  2527. {
  2528. type: "tool-call",
  2529. id: "hosted-search",
  2530. name: "web_search",
  2531. input: { query: "Effect" },
  2532. providerExecuted: true,
  2533. providerMetadata: { openai: { itemId: "hosted-search" } },
  2534. },
  2535. {
  2536. type: "tool-result",
  2537. id: "hosted-search",
  2538. name: "web_search",
  2539. // The generic replay result derives from canonical stored content.
  2540. result: { type: "text", value: '[{"title":"Effect"}]' },
  2541. providerExecuted: true,
  2542. providerMetadata: { openai: { blockType: "web_search_tool_result" } },
  2543. },
  2544. ])
  2545. }),
  2546. )
  2547. it.effect("starts recorded local tools eagerly and awaits settlement before continuing", () =>
  2548. Effect.gen(function* () {
  2549. const session = yield* setup
  2550. yield* admit(session, "Echo five times")
  2551. toolExecutionGate = yield* Deferred.make<void>()
  2552. toolExecutionsStarted = yield* Deferred.make<void>()
  2553. const providerGate = yield* Deferred.make<void>()
  2554. const initial = Stream.fromIterable([
  2555. LLMEvent.stepStart({ index: 0 }),
  2556. ...Array.from({ length: 5 }, (_, index) =>
  2557. LLMEvent.toolCall({ id: `call-echo-${index}`, name: "echo", input: { text: `${index}` } }),
  2558. ),
  2559. ])
  2560. const final = Stream.fromIterable([
  2561. LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }),
  2562. LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
  2563. ])
  2564. responseStream = Stream.concat(
  2565. initial,
  2566. Stream.fromEffect(Deferred.await(providerGate)).pipe(Stream.flatMap(() => final)),
  2567. )
  2568. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2569. yield* Deferred.await(toolExecutionsStarted)
  2570. expect(executions).toHaveLength(5)
  2571. expect(maxActiveToolExecutions).toBe(5)
  2572. expect(yield* session.context(sessionID)).toMatchObject([
  2573. { type: "user", text: "Echo five times" },
  2574. {
  2575. type: "assistant",
  2576. content: Array.from({ length: 5 }, (_, index) => ({
  2577. type: "tool",
  2578. id: `call-echo-${index}`,
  2579. state: { status: "running", input: { text: `${index}` } },
  2580. })),
  2581. },
  2582. ])
  2583. yield* Deferred.succeed(providerGate, undefined)
  2584. yield* Effect.yieldNow
  2585. expect(requests).toHaveLength(1)
  2586. yield* Deferred.succeed(toolExecutionGate, undefined)
  2587. yield* Fiber.join(run)
  2588. toolExecutionGate = undefined
  2589. toolExecutionsStarted = undefined
  2590. expect(executions).toHaveLength(5)
  2591. expect(maxActiveToolExecutions).toBe(5)
  2592. expect(requests).toHaveLength(2)
  2593. }),
  2594. )
  2595. it.effect("settles repeated provider-local tool call IDs against their owning assistant messages", () =>
  2596. Effect.gen(function* () {
  2597. const session = yield* setup
  2598. yield* admit(session, "Echo twice")
  2599. responses = [
  2600. reply.tool("tool_0", "echo", { text: "first" }),
  2601. reply.tool("tool_0", "echo", { text: "second" }),
  2602. [],
  2603. ]
  2604. yield* session.resume(sessionID)
  2605. expect(executions).toEqual(["first", "second"])
  2606. expect(requests).toHaveLength(3)
  2607. expect(yield* session.context(sessionID)).toMatchObject([
  2608. { type: "user", text: "Echo twice" },
  2609. {
  2610. type: "assistant",
  2611. content: [
  2612. {
  2613. type: "tool",
  2614. id: "tool_0",
  2615. state: { status: "completed", content: [{ type: "text", text: "first" }] },
  2616. },
  2617. ],
  2618. },
  2619. {
  2620. type: "assistant",
  2621. content: [
  2622. {
  2623. type: "tool",
  2624. id: "tool_0",
  2625. state: { status: "completed", content: [{ type: "text", text: "second" }] },
  2626. },
  2627. ],
  2628. },
  2629. ])
  2630. yield* replaySessionProjection(sessionID)
  2631. expect(yield* session.context(sessionID)).toMatchObject([
  2632. { type: "user", text: "Echo twice" },
  2633. {
  2634. type: "assistant",
  2635. content: [
  2636. {
  2637. type: "tool",
  2638. id: "tool_0",
  2639. state: { status: "completed", content: [{ type: "text", text: "first" }] },
  2640. },
  2641. ],
  2642. },
  2643. {
  2644. type: "assistant",
  2645. content: [
  2646. {
  2647. type: "tool",
  2648. id: "tool_0",
  2649. state: { status: "completed", content: [{ type: "text", text: "second" }] },
  2650. },
  2651. ],
  2652. },
  2653. ])
  2654. }),
  2655. )
  2656. it.effect("joins concurrent resume calls into one active provider run", () =>
  2657. Effect.gen(function* () {
  2658. const session = yield* setup
  2659. yield* admit(session, "Run once")
  2660. response = reply.text("Once", "text-once")
  2661. streamGate = yield* Deferred.make<void>()
  2662. streamStarted = yield* Deferred.make<void>()
  2663. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2664. yield* Deferred.await(streamStarted)
  2665. const second = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2666. yield* Effect.yieldNow
  2667. expect(requests).toHaveLength(1)
  2668. yield* Deferred.succeed(streamGate, undefined)
  2669. yield* Fiber.join(first)
  2670. yield* Fiber.join(second)
  2671. streamGate = undefined
  2672. streamStarted = undefined
  2673. expect(requests).toHaveLength(1)
  2674. expect(yield* session.context(sessionID)).toMatchObject([
  2675. { type: "user", text: "Run once" },
  2676. { type: "assistant", finish: "stop", content: [{ type: "text", text: "Once" }] },
  2677. ])
  2678. }),
  2679. )
  2680. it.effect("steers an active step with newly recorded prompts", () =>
  2681. Effect.gen(function* () {
  2682. const session = yield* setup
  2683. yield* admit(session, "Start working")
  2684. responses = [reply.stop(), reply.stop()]
  2685. streamGate = yield* Deferred.make<void>()
  2686. streamStarted = yield* Deferred.make<void>()
  2687. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2688. yield* Deferred.await(streamStarted)
  2689. yield* session.prompt({ sessionID, text: "Change direction" })
  2690. yield* Deferred.succeed(streamGate, undefined)
  2691. yield* Fiber.join(first)
  2692. streamGate = undefined
  2693. streamStarted = undefined
  2694. yield* Effect.yieldNow
  2695. expect(requests).toHaveLength(2)
  2696. expect(userTexts(requests[0]!)).toEqual(["Start working"])
  2697. expect(userTexts(requests[1]!)).toEqual(["Start working", "Change direction"])
  2698. expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([
  2699. "user",
  2700. "assistant",
  2701. "user",
  2702. "assistant",
  2703. ])
  2704. }),
  2705. )
  2706. it.effect("promotes queued input after continuation ends", () =>
  2707. Effect.gen(function* () {
  2708. const session = yield* setup
  2709. yield* admit(session, "Start working")
  2710. responses = [reply.tool("call-echo", "echo", { text: "hello" }), reply.stop(), reply.stop()]
  2711. streamGate = yield* Deferred.make<void>()
  2712. streamStarted = yield* Deferred.make<void>()
  2713. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2714. yield* Deferred.await(streamStarted)
  2715. yield* session.prompt({
  2716. sessionID,
  2717. text: "Wait until continuation ends",
  2718. delivery: "queue",
  2719. })
  2720. yield* Deferred.succeed(streamGate, undefined)
  2721. yield* Fiber.join(first)
  2722. streamGate = undefined
  2723. streamStarted = undefined
  2724. expect(requests).toHaveLength(3)
  2725. expect(userTexts(requests[0]!)).toEqual(["Start working"])
  2726. expect(userTexts(requests[1]!)).toEqual(["Start working"])
  2727. expect(userTexts(requests[2]!)).toEqual(["Start working", "Wait until continuation ends"])
  2728. }),
  2729. )
  2730. it.effect("preserves durable queued input for a later wake after interruption", () =>
  2731. Effect.gen(function* () {
  2732. const session = yield* setup
  2733. const { db } = yield* Database.Service
  2734. yield* admit(session, "Interrupt current work")
  2735. responses = [[], reply.stop()]
  2736. streamGate = yield* Deferred.make<void>()
  2737. streamStarted = yield* Deferred.make<void>()
  2738. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2739. yield* Deferred.await(streamStarted)
  2740. yield* session.prompt({
  2741. sessionID,
  2742. text: "Run after interrupt",
  2743. delivery: "queue",
  2744. })
  2745. yield* session.interrupt(sessionID)
  2746. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  2747. expect(requests).toHaveLength(1)
  2748. expect(yield* SessionPending.has(db, sessionID, "queue")).toBe(true)
  2749. const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2750. while (requests.length < 2) yield* Effect.yieldNow
  2751. yield* Deferred.succeed(streamGate, undefined)
  2752. yield* Fiber.join(resumed)
  2753. streamGate = undefined
  2754. streamStarted = undefined
  2755. expect(requests).toHaveLength(2)
  2756. expect(userTexts(requests[0]!)).toEqual(["Interrupt current work"])
  2757. expect(userTexts(requests[1]!)).toEqual(["Interrupt current work", "Run after interrupt"])
  2758. }),
  2759. )
  2760. it.effect("preserves durable steering input for a later resume after interruption", () =>
  2761. Effect.gen(function* () {
  2762. const session = yield* setup
  2763. const { db } = yield* Database.Service
  2764. yield* admit(session, "Interrupt current work")
  2765. responses = [[], reply.stop()]
  2766. streamGate = yield* Deferred.make<void>()
  2767. streamStarted = yield* Deferred.make<void>()
  2768. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2769. yield* Deferred.await(streamStarted)
  2770. yield* session.prompt({
  2771. sessionID,
  2772. text: "Steer after interrupt",
  2773. })
  2774. yield* session.interrupt(sessionID)
  2775. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  2776. expect(requests).toHaveLength(1)
  2777. expect(yield* SessionPending.has(db, sessionID, "steer")).toBe(true)
  2778. const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2779. while (requests.length < 2) yield* Effect.yieldNow
  2780. yield* Deferred.succeed(streamGate, undefined)
  2781. yield* Fiber.join(resumed)
  2782. streamGate = undefined
  2783. streamStarted = undefined
  2784. expect(requests).toHaveLength(2)
  2785. expect(userTexts(requests[0]!)).toEqual(["Interrupt current work"])
  2786. expect(userTexts(requests[1]!)).toEqual(["Interrupt current work", "Steer after interrupt"])
  2787. }),
  2788. )
  2789. it.effect("promotes queued inputs one at a time in FIFO order", () =>
  2790. Effect.gen(function* () {
  2791. const session = yield* setup
  2792. yield* admit(session, "Start working")
  2793. responses = [reply.stop(), reply.stop(), reply.stop()]
  2794. streamGate = yield* Deferred.make<void>()
  2795. streamStarted = yield* Deferred.make<void>()
  2796. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2797. yield* Deferred.await(streamStarted)
  2798. yield* session.prompt({ sessionID, text: "Queue first", delivery: "queue" })
  2799. yield* session.prompt({ sessionID, text: "Queue second", delivery: "queue" })
  2800. yield* Deferred.succeed(streamGate, undefined)
  2801. yield* Fiber.join(first)
  2802. streamGate = undefined
  2803. streamStarted = undefined
  2804. expect(requests).toHaveLength(3)
  2805. expect(userTexts(requests[0]!)).toEqual(["Start working"])
  2806. expect(userTexts(requests[1]!)).toEqual(["Start working", "Queue first"])
  2807. expect(userTexts(requests[2]!)).toEqual(["Start working", "Queue first", "Queue second"])
  2808. }),
  2809. )
  2810. it.effect("promotes queued input after steering continuation ends", () =>
  2811. Effect.gen(function* () {
  2812. const session = yield* setup
  2813. yield* admit(session, "Start steering")
  2814. yield* session.prompt({
  2815. sessionID,
  2816. text: "Queue for later",
  2817. delivery: "queue",
  2818. resume: false,
  2819. })
  2820. responses = [reply.stop(), reply.stop()]
  2821. yield* session.resume(sessionID)
  2822. expect(requests).toHaveLength(2)
  2823. expect(userTexts(requests[0]!)).toEqual(["Start steering"])
  2824. expect(userTexts(requests[1]!)).toEqual(["Start steering", "Queue for later"])
  2825. }),
  2826. )
  2827. it.effect("promotes steers before the next queued input", () =>
  2828. Effect.gen(function* () {
  2829. const session = yield* setup
  2830. yield* admit(session, "Start working")
  2831. responses = [reply.stop(), reply.stop(), reply.stop(), reply.stop()]
  2832. const firstGate = yield* Deferred.make<void>()
  2833. const secondGate = yield* Deferred.make<void>()
  2834. streamGate = firstGate
  2835. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2836. while (requests.length < 1) yield* Effect.yieldNow
  2837. yield* session.prompt({ sessionID, text: "Queue first", delivery: "queue" })
  2838. yield* session.prompt({ sessionID, text: "Queue second", delivery: "queue" })
  2839. streamGate = secondGate
  2840. yield* Deferred.succeed(firstGate, undefined)
  2841. while (requests.length < 2) yield* Effect.yieldNow
  2842. yield* session.prompt({ sessionID, text: "Steer before next queued input" })
  2843. yield* session.prompt({
  2844. sessionID,
  2845. text: "Also steer before next queued input",
  2846. })
  2847. yield* session.synthetic({ sessionID, text: "Background completion before next queued input" })
  2848. yield* Deferred.succeed(secondGate, undefined)
  2849. yield* Fiber.join(first)
  2850. streamGate = undefined
  2851. expect(requests).toHaveLength(4)
  2852. expect(userTexts(requests[0]!)).toEqual(["Start working"])
  2853. expect(userTexts(requests[1]!)).toEqual(["Start working", "Queue first"])
  2854. expect(userTexts(requests[2]!)).toEqual([
  2855. "Start working",
  2856. "Queue first",
  2857. "Steer before next queued input",
  2858. "Also steer before next queued input",
  2859. "Background completion before next queued input",
  2860. ])
  2861. expect(userTexts(requests[3]!)).toEqual([
  2862. "Start working",
  2863. "Queue first",
  2864. "Steer before next queued input",
  2865. "Also steer before next queued input",
  2866. "Background completion before next queued input",
  2867. "Queue second",
  2868. ])
  2869. }),
  2870. )
  2871. it.effect("coalesces multiple active steering prompts into one continuation step", () =>
  2872. Effect.gen(function* () {
  2873. const session = yield* setup
  2874. yield* admit(session, "Start working")
  2875. responses = [reply.stop(), reply.stop()]
  2876. streamGate = yield* Deferred.make<void>()
  2877. streamStarted = yield* Deferred.make<void>()
  2878. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2879. yield* Deferred.await(streamStarted)
  2880. yield* session.prompt({ sessionID, text: "First steer" })
  2881. yield* session.prompt({ sessionID, text: "Second steer" })
  2882. yield* Deferred.succeed(streamGate, undefined)
  2883. yield* Fiber.join(first)
  2884. streamGate = undefined
  2885. streamStarted = undefined
  2886. yield* Effect.yieldNow
  2887. expect(requests).toHaveLength(2)
  2888. expect(userTexts(requests[1]!)).toEqual(["Start working", "First steer", "Second steer"])
  2889. yield* (yield* SessionExecution.Service).wake(sessionID)
  2890. yield* Effect.yieldNow
  2891. expect(requests).toHaveLength(2)
  2892. }),
  2893. )
  2894. it.effect("runs steering input accepted while the active step fails", () =>
  2895. Effect.gen(function* () {
  2896. const session = yield* setup
  2897. yield* admit(session, "Start working")
  2898. streamFailure = invalidRequest()
  2899. streamGate = yield* Deferred.make<void>()
  2900. streamStarted = yield* Deferred.make<void>()
  2901. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2902. yield* Deferred.await(streamStarted)
  2903. yield* session.prompt({ sessionID, text: "Recover with this" })
  2904. yield* Deferred.succeed(streamGate, undefined)
  2905. expect(yield* Fiber.join(first).pipe(Effect.flip)).toBe(streamFailure)
  2906. streamFailure = undefined
  2907. streamGate = undefined
  2908. streamStarted = undefined
  2909. yield* session.wait(sessionID)
  2910. expect(requests).toHaveLength(2)
  2911. expect(userTexts(requests[1]!)).toEqual(["Start working", "Recover with this"])
  2912. }),
  2913. )
  2914. it.effect("durably fails local tools left running by a prior process before continuing", () =>
  2915. Effect.gen(function* () {
  2916. const session = yield* setup
  2917. const events = yield* EventV2.Service
  2918. yield* admit(session, "Recover interrupted tool")
  2919. yield* SessionPending.promote((yield* Database.Service).db, events, sessionID, "steer")
  2920. const assistantMessageID = SessionMessage.ID.create()
  2921. yield* events.publish(SessionEvent.Step.Started, {
  2922. sessionID,
  2923. assistantMessageID,
  2924. agent: AgentV2.ID.make("build"),
  2925. model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
  2926. })
  2927. yield* events.publish(SessionEvent.Tool.Input.Started, {
  2928. sessionID,
  2929. assistantMessageID,
  2930. callID: "call-interrupted",
  2931. name: "echo",
  2932. })
  2933. yield* events.publish(SessionEvent.Tool.Input.Ended, {
  2934. sessionID,
  2935. assistantMessageID,
  2936. callID: "call-interrupted",
  2937. text: '{"text":"stale"}',
  2938. })
  2939. yield* events.publish(SessionEvent.Tool.Called, {
  2940. sessionID,
  2941. assistantMessageID,
  2942. callID: "call-interrupted",
  2943. input: { text: "stale" },
  2944. executed: false,
  2945. })
  2946. requests.length = 0
  2947. response = []
  2948. yield* session.resume(sessionID)
  2949. expect(requests).toHaveLength(1)
  2950. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
  2951. expect(yield* session.context(sessionID)).toMatchObject([
  2952. { type: "user", text: "Recover interrupted tool" },
  2953. {
  2954. type: "assistant",
  2955. content: [
  2956. {
  2957. type: "tool",
  2958. id: "call-interrupted",
  2959. state: {
  2960. status: "error",
  2961. error: { type: "aborted", message: "Tool execution interrupted: echo" },
  2962. },
  2963. },
  2964. ],
  2965. },
  2966. ])
  2967. }),
  2968. )
  2969. it.effect("durably fails hosted tools left running by a prior process before continuing inline", () =>
  2970. Effect.gen(function* () {
  2971. const session = yield* setup
  2972. const events = yield* EventV2.Service
  2973. yield* admit(session, "Recover interrupted hosted tool")
  2974. yield* SessionPending.promote((yield* Database.Service).db, events, sessionID, "steer")
  2975. const assistantMessageID = SessionMessage.ID.create()
  2976. yield* events.publish(SessionEvent.Step.Started, {
  2977. sessionID,
  2978. assistantMessageID,
  2979. agent: AgentV2.ID.make("build"),
  2980. model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
  2981. })
  2982. yield* events.publish(SessionEvent.Tool.Input.Started, {
  2983. sessionID,
  2984. assistantMessageID,
  2985. callID: "call-hosted-interrupted",
  2986. name: "web_search",
  2987. })
  2988. yield* events.publish(SessionEvent.Tool.Input.Ended, {
  2989. sessionID,
  2990. assistantMessageID,
  2991. callID: "call-hosted-interrupted",
  2992. text: '{"query":"stale"}',
  2993. })
  2994. yield* events.publish(SessionEvent.Tool.Called, {
  2995. sessionID,
  2996. assistantMessageID,
  2997. callID: "call-hosted-interrupted",
  2998. input: { query: "stale" },
  2999. executed: true,
  3000. state: { itemId: "call-hosted-interrupted" },
  3001. })
  3002. requests.length = 0
  3003. response = []
  3004. yield* session.resume(sessionID)
  3005. expect(requests).toHaveLength(1)
  3006. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant"])
  3007. expect(requests[0]?.messages[1]?.content).toMatchObject([
  3008. {
  3009. type: "tool-call",
  3010. id: "call-hosted-interrupted",
  3011. providerExecuted: true,
  3012. providerMetadata: { openai: { itemId: "call-hosted-interrupted" } },
  3013. },
  3014. { type: "tool-result", id: "call-hosted-interrupted", providerExecuted: true, result: { type: "error" } },
  3015. ])
  3016. }),
  3017. )
  3018. it.effect("durably fails pending tool input left by a prior process before continuing", () =>
  3019. Effect.gen(function* () {
  3020. const session = yield* setup
  3021. const events = yield* EventV2.Service
  3022. yield* admit(session, "Recover interrupted tool input")
  3023. yield* SessionPending.promote((yield* Database.Service).db, events, sessionID, "steer")
  3024. const assistantMessageID = SessionMessage.ID.create()
  3025. yield* events.publish(SessionEvent.Step.Started, {
  3026. sessionID,
  3027. assistantMessageID,
  3028. agent: AgentV2.ID.make("build"),
  3029. model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
  3030. })
  3031. yield* events.publish(SessionEvent.Tool.Input.Started, {
  3032. sessionID,
  3033. assistantMessageID,
  3034. callID: "call-pending-interrupted",
  3035. name: "echo",
  3036. })
  3037. requests.length = 0
  3038. response = []
  3039. yield* session.resume(sessionID)
  3040. expect(requests).toHaveLength(1)
  3041. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
  3042. expect(yield* session.context(sessionID)).toMatchObject([
  3043. { type: "user", text: "Recover interrupted tool input" },
  3044. { type: "assistant", content: [{ type: "tool", id: "call-pending-interrupted", state: { status: "error" } }] },
  3045. ])
  3046. }),
  3047. )
  3048. it.effect("promotes the first queued input when woken while idle", () =>
  3049. Effect.gen(function* () {
  3050. const session = yield* setup
  3051. yield* session.prompt({
  3052. sessionID,
  3053. text: "Wait in queue",
  3054. delivery: "queue",
  3055. resume: false,
  3056. })
  3057. yield* (yield* SessionExecution.Service).wake(sessionID)
  3058. while (requests.length === 0) yield* Effect.yieldNow
  3059. expect(requests).toHaveLength(1)
  3060. expect(userTexts(requests[0]!)).toEqual(["Wait in queue"])
  3061. }),
  3062. )
  3063. it.effect("retries inbox input after prompt projection rolls back", () =>
  3064. Effect.gen(function* () {
  3065. const session = yield* setup
  3066. const events = yield* EventV2.Service
  3067. const defect = new Error("fail after prompt promotion")
  3068. let fail = true
  3069. yield* events.project(SessionEvent.InputPromoted, () => (fail ? Effect.die(defect) : Effect.void))
  3070. yield* admit(session, "Recover promoted input")
  3071. expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect)
  3072. fail = false
  3073. requests.length = 0
  3074. response = reply.stop()
  3075. yield* (yield* SessionExecution.Service).wake(sessionID)
  3076. while (requests.length === 0) yield* Effect.yieldNow
  3077. expect(userTexts(requests[0]!)).toEqual(["Recover promoted input"])
  3078. }),
  3079. )
  3080. it.effect("does not strand a committed promotion when a post-commit listener defects", () =>
  3081. Effect.gen(function* () {
  3082. const session = yield* setup
  3083. const events = yield* EventV2.Service
  3084. yield* events.listen((event) =>
  3085. event.type === SessionEvent.InputPromoted.type
  3086. ? Effect.die("fail after prompt promotion commits")
  3087. : Effect.void,
  3088. )
  3089. yield* admit(session, "Run committed promotion")
  3090. yield* session.resume(sessionID)
  3091. expect(requests).toHaveLength(1)
  3092. expect(userTexts(requests[0]!)).toEqual(["Run committed promotion"])
  3093. }),
  3094. )
  3095. it.effect("adds session correlation headers to model requests", () =>
  3096. Effect.gen(function* () {
  3097. const session = yield* setup
  3098. yield* admit(session, "Run correlated request")
  3099. yield* session.resume(sessionID)
  3100. expect(requests[0]?.http?.headers).toEqual({
  3101. "x-session-affinity": sessionID,
  3102. "X-Session-Id": sessionID,
  3103. "User-Agent": App.useragent(App.make()),
  3104. "x-opencode-project": Project.ID.global,
  3105. "x-opencode-session": sessionID,
  3106. "x-opencode-client": "opencode",
  3107. })
  3108. }),
  3109. )
  3110. it.effect("adds the parent session header to child model requests", () =>
  3111. Effect.gen(function* () {
  3112. const session = yield* setup
  3113. const parentID = SessionV2.ID.make("ses_runner_parent")
  3114. const { db } = yield* Database.Service
  3115. yield* db
  3116. .update(SessionTable)
  3117. .set({ parent_id: parentID })
  3118. .where(eq(SessionTable.id, sessionID))
  3119. .run()
  3120. .pipe(Effect.orDie)
  3121. yield* admit(session, "Run child request")
  3122. yield* session.resume(sessionID)
  3123. expect(requests[0]?.http?.headers?.["x-parent-session-id"]).toBe(parentID)
  3124. }),
  3125. )
  3126. it.effect("runs different sessions concurrently", () =>
  3127. Effect.gen(function* () {
  3128. const session = yield* setup
  3129. yield* insertSession(otherSessionID)
  3130. yield* admit(session, "Run first")
  3131. yield* session.prompt({
  3132. sessionID: otherSessionID,
  3133. text: "Run second",
  3134. resume: false,
  3135. })
  3136. streamGate = yield* Deferred.make<void>()
  3137. streamStarted = yield* Deferred.make<void>()
  3138. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3139. yield* Deferred.await(streamStarted)
  3140. streamStarted = yield* Deferred.make<void>()
  3141. const second = yield* session.resume(otherSessionID).pipe(Effect.forkChild)
  3142. yield* Deferred.await(streamStarted)
  3143. expect(requests).toHaveLength(2)
  3144. expect(requests.map((request) => request.providerOptions?.openai?.promptCacheKey)).toEqual([
  3145. sessionID,
  3146. otherSessionID,
  3147. ])
  3148. yield* Deferred.succeed(streamGate, undefined)
  3149. yield* Fiber.join(first)
  3150. yield* Fiber.join(second)
  3151. streamGate = undefined
  3152. streamStarted = undefined
  3153. }),
  3154. )
  3155. it.effect("bounds 64-character session prompt cache keys", () =>
  3156. Effect.gen(function* () {
  3157. const session = yield* setup
  3158. const longSessionID = SessionV2.ID.make(`ses_${"a".repeat(64)}`)
  3159. const otherLongSessionID = SessionV2.ID.make(`ses_${"b".repeat(64)}`)
  3160. yield* insertSession(longSessionID)
  3161. yield* insertSession(otherLongSessionID)
  3162. yield* session.prompt({
  3163. sessionID: longSessionID,
  3164. text: "Run long session",
  3165. resume: false,
  3166. })
  3167. yield* session.prompt({
  3168. sessionID: otherLongSessionID,
  3169. text: "Run other long session",
  3170. resume: false,
  3171. })
  3172. yield* session.resume(longSessionID)
  3173. yield* session.resume(otherLongSessionID)
  3174. const keys = requests.map((request) => request.providerOptions?.openai?.promptCacheKey)
  3175. expect(keys).toEqual([longSessionID.slice(4), otherLongSessionID.slice(4)])
  3176. expect(keys.every((key) => typeof key === "string" && key.length === 64)).toBe(true)
  3177. expect(keys[0]).not.toBe(keys[1])
  3178. }),
  3179. )
  3180. it.effect("fans out one failed run and allows a later retry", () =>
  3181. Effect.gen(function* () {
  3182. const session = yield* setup
  3183. yield* admit(session, "Retry after failure")
  3184. streamFailure = invalidRequest()
  3185. streamGate = yield* Deferred.make<void>()
  3186. streamStarted = yield* Deferred.make<void>()
  3187. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3188. yield* Deferred.await(streamStarted)
  3189. const second = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3190. yield* Effect.yieldNow
  3191. expect(requests).toHaveLength(1)
  3192. yield* Deferred.succeed(streamGate, undefined)
  3193. const [firstExit, secondExit] = yield* Effect.all([Fiber.await(first), Fiber.await(second)])
  3194. expect(secondExit).toEqual(firstExit)
  3195. streamFailure = undefined
  3196. streamGate = undefined
  3197. streamStarted = undefined
  3198. yield* session.resume(sessionID)
  3199. expect(requests).toHaveLength(2)
  3200. }),
  3201. )
  3202. it.effect("durably settles local tool failures before continuing", () =>
  3203. Effect.gen(function* () {
  3204. const session = yield* setup
  3205. yield* admit(session, "Call missing")
  3206. responses = [reply.tool("call-missing", "missing", {}), reply.text("Recovered", "text-after-error")]
  3207. yield* session.resume(sessionID)
  3208. expect(requests).toHaveLength(2)
  3209. expect(yield* session.context(sessionID)).toMatchObject([
  3210. { type: "user", text: "Call missing" },
  3211. {
  3212. type: "assistant",
  3213. content: [
  3214. {
  3215. type: "tool",
  3216. id: "call-missing",
  3217. state: {
  3218. status: "error",
  3219. error: { type: "tool.unknown", message: "Unknown tool: missing" },
  3220. },
  3221. },
  3222. ],
  3223. },
  3224. { type: "assistant", finish: "stop", content: [{ type: "text", text: "Recovered" }] },
  3225. ])
  3226. }),
  3227. )
  3228. it.effect("returns unexpected local tool defects to the model and continues", () =>
  3229. Effect.gen(function* () {
  3230. const session = yield* setup
  3231. yield* admit(session, "Call defect")
  3232. responses = [reply.tool("call-defect", "defect", {}), reply.text("Recovered", "text-after-defect")]
  3233. yield* session.resume(sessionID)
  3234. expect(requests).toHaveLength(2)
  3235. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
  3236. const context = yield* session.context(sessionID)
  3237. expect(context).toMatchObject([
  3238. { type: "user", text: "Call defect" },
  3239. {
  3240. type: "assistant",
  3241. content: [
  3242. {
  3243. type: "tool",
  3244. id: "call-defect",
  3245. state: {
  3246. status: "error",
  3247. error: { type: "unknown", message: "unexpected tool defect" },
  3248. },
  3249. },
  3250. ],
  3251. },
  3252. { type: "assistant", finish: "stop", content: [{ type: "text", text: "Recovered" }] },
  3253. ])
  3254. const assistant = requireAssistant(context)
  3255. expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([
  3256. "session.step.started.1",
  3257. "session.tool.called.1",
  3258. "session.tool.failed.2",
  3259. "session.step.ended.1",
  3260. ])
  3261. }),
  3262. )
  3263. it.effect("returns tool-wrapped policy blocks to the model and continues", () =>
  3264. Effect.gen(function* () {
  3265. const session = yield* setup
  3266. const registry = yield* ToolRegistry.Service
  3267. yield* registry.register(
  3268. {
  3269. blocked: Tool.make({
  3270. description: "Fail because policy blocked execution",
  3271. input: Schema.Struct({}),
  3272. output: Schema.Struct({}),
  3273. execute: () =>
  3274. Effect.fail(new PermissionV2.BlockedError({ rules: [], permission: "blocked", resources: ["*"] })).pipe(
  3275. Effect.mapError(() => new Tool.Failure({ message: "Permission blocked" })),
  3276. ),
  3277. }),
  3278. },
  3279. { codemode: false },
  3280. )
  3281. yield* admit(session, "Call blocked")
  3282. responses = [reply.tool("call-blocked", "blocked", {}), reply.stop()]
  3283. yield* session.resume(sessionID)
  3284. expect(requests).toHaveLength(2)
  3285. expect(yield* session.context(sessionID)).toMatchObject([
  3286. { type: "user", text: "Call blocked" },
  3287. {
  3288. type: "assistant",
  3289. content: [
  3290. { type: "tool", id: "call-blocked", state: { status: "error", error: { message: "Permission blocked" } } },
  3291. ],
  3292. },
  3293. { type: "assistant", finish: "stop" },
  3294. ])
  3295. }),
  3296. )
  3297. it.effect("interrupts runner continuation when permission approval is declined", () =>
  3298. Effect.gen(function* () {
  3299. const session = yield* setup
  3300. const registry = yield* ToolRegistry.Service
  3301. yield* registry.register(
  3302. {
  3303. declined: Tool.make({
  3304. description: "Fail because the user declined approval",
  3305. input: Schema.Struct({}),
  3306. output: Schema.Struct({}),
  3307. execute: () => Effect.die(new PermissionV2.DeclinedError()),
  3308. }),
  3309. },
  3310. { codemode: false },
  3311. )
  3312. yield* admit(session, "Call declined")
  3313. response = reply.tool("call-declined", "declined", {})
  3314. const exit = yield* session.resume(sessionID).pipe(Effect.exit)
  3315. expect(exit._tag).toBe("Failure")
  3316. if (exit._tag === "Failure") expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
  3317. expect(requests).toHaveLength(1)
  3318. expect(yield* session.context(sessionID)).toMatchObject([
  3319. { type: "user", text: "Call declined" },
  3320. {
  3321. type: "assistant",
  3322. content: [
  3323. {
  3324. type: "tool",
  3325. id: "call-declined",
  3326. state: { status: "error", error: { type: "aborted", message: "The user declined this tool call" } },
  3327. },
  3328. ],
  3329. },
  3330. ])
  3331. }),
  3332. )
  3333. it.effect("returns permission corrections to the model and continues", () =>
  3334. Effect.gen(function* () {
  3335. const session = yield* setup
  3336. const registry = yield* ToolRegistry.Service
  3337. yield* registry.register(
  3338. {
  3339. corrected: Tool.make({
  3340. description: "Fail with user correction feedback",
  3341. input: Schema.Struct({}),
  3342. output: Schema.Struct({}),
  3343. execute: () =>
  3344. Effect.fail(new PermissionV2.CorrectedError({ feedback: "Use another tool" })).pipe(
  3345. Effect.mapError(() => new Tool.Failure({ message: "Use another tool" })),
  3346. ),
  3347. }),
  3348. },
  3349. { codemode: false },
  3350. )
  3351. yield* admit(session, "Call corrected")
  3352. responses = [reply.tool("call-corrected", "corrected", {}), reply.stop()]
  3353. yield* session.resume(sessionID)
  3354. expect(requests).toHaveLength(2)
  3355. expect(yield* session.context(sessionID)).toMatchObject([
  3356. { type: "user", text: "Call corrected" },
  3357. {
  3358. type: "assistant",
  3359. content: [
  3360. { type: "tool", id: "call-corrected", state: { status: "error", error: { message: "Use another tool" } } },
  3361. ],
  3362. },
  3363. { type: "assistant", finish: "stop" },
  3364. ])
  3365. }),
  3366. )
  3367. it.effect("fails the drain when tool output persistence fails", () =>
  3368. Effect.gen(function* () {
  3369. const session = yield* setup
  3370. yield* admit(session, "Call storefail")
  3371. responses = [reply.tool("call-storefail", "storefail", {}), []]
  3372. const exit = yield* session.resume(sessionID).pipe(Effect.exit)
  3373. expect(Exit.isFailure(exit)).toBe(true)
  3374. expect(requests).toHaveLength(1)
  3375. expect(yield* session.context(sessionID)).toMatchObject([
  3376. { type: "user", text: "Call storefail" },
  3377. {
  3378. type: "assistant",
  3379. content: [
  3380. {
  3381. type: "tool",
  3382. id: "call-storefail",
  3383. state: {
  3384. status: "error",
  3385. error: {
  3386. type: "unknown",
  3387. message: expect.stringContaining("Failed to write tool output"),
  3388. },
  3389. },
  3390. },
  3391. ],
  3392. finish: "error",
  3393. error: { type: "unknown", message: expect.stringContaining("Failed to write tool output") },
  3394. },
  3395. ])
  3396. }),
  3397. )
  3398. it.effect("returns configured permission denials to the model and continues", () =>
  3399. Effect.gen(function* () {
  3400. const session = yield* setup
  3401. const registry = yield* ToolRegistry.Service
  3402. yield* registry.register({ permissionfail: permissionFail }, { codemode: false })
  3403. yield* admit(session, "Reject permission")
  3404. responses = [
  3405. reply.tool("call-permission", "permissionfail", {}),
  3406. [LLMEvent.stepStart({ index: 0 }), LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } })],
  3407. ]
  3408. yield* session.resume(sessionID)
  3409. expect(requests).toHaveLength(2)
  3410. expect(yield* session.context(sessionID)).toMatchObject([
  3411. { type: "user" },
  3412. {
  3413. type: "assistant",
  3414. content: [
  3415. {
  3416. type: "tool",
  3417. id: "call-permission",
  3418. state: {
  3419. status: "error",
  3420. error: {
  3421. type: "permission.rejected",
  3422. message: "Permission denied: edit",
  3423. },
  3424. },
  3425. },
  3426. ],
  3427. },
  3428. { type: "assistant", finish: "stop" },
  3429. ])
  3430. expect(yield* recordedEventTypes(sessionID)).not.toContain("session.step.failed.1")
  3431. }),
  3432. )
  3433. it.effect("interrupts runner continuation when a question is cancelled", () =>
  3434. Effect.gen(function* () {
  3435. const session = yield* setup
  3436. const registry = yield* ToolRegistry.Service
  3437. yield* registry.register(
  3438. {
  3439. question: Tool.make({
  3440. description: "Ask the user",
  3441. input: Schema.Struct({}),
  3442. output: Schema.Struct({}),
  3443. execute: () => Effect.die(new QuestionTool.CancelledError()),
  3444. }),
  3445. },
  3446. { codemode: false },
  3447. )
  3448. yield* admit(session, "Ask then stop")
  3449. responses = [reply.tool("call-question", "question", {}), []]
  3450. const run = yield* session.resume(sessionID).pipe(Effect.exit, Effect.forkChild)
  3451. const exit = yield* Fiber.join(run)
  3452. expect(exit._tag).toBe("Failure")
  3453. if (exit._tag === "Failure") expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
  3454. expect(requests).toHaveLength(1)
  3455. expect(yield* session.context(sessionID)).toMatchObject([
  3456. { type: "user", text: "Ask then stop" },
  3457. {
  3458. type: "assistant",
  3459. content: [
  3460. {
  3461. type: "tool",
  3462. id: "call-question",
  3463. state: { status: "error", error: { type: "aborted", message: "The user dismissed this question" } },
  3464. },
  3465. ],
  3466. },
  3467. ])
  3468. }),
  3469. )
  3470. it.effect("awaits started local tools before surfacing provider stream failure", () =>
  3471. Effect.gen(function* () {
  3472. const session = yield* setup
  3473. yield* admit(session, "Settle before failing")
  3474. const failure = providerUnavailable()
  3475. toolExecutionGate = yield* Deferred.make<void>()
  3476. responseStream = Stream.concat(
  3477. Stream.fromIterable([
  3478. LLMEvent.stepStart({ index: 0 }),
  3479. LLMEvent.toolCall({ id: "call-before-failure", name: "echo", input: { text: "settle" } }),
  3480. ]),
  3481. Stream.fail(failure),
  3482. )
  3483. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3484. while (executions.length === 0) yield* Effect.yieldNow
  3485. yield* Effect.yieldNow
  3486. yield* Deferred.succeed(toolExecutionGate, undefined)
  3487. expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure)
  3488. toolExecutionGate = undefined
  3489. const context = yield* session.context(sessionID)
  3490. expect(context).toMatchObject([
  3491. { type: "user", text: "Settle before failing" },
  3492. {
  3493. type: "assistant",
  3494. content: [
  3495. {
  3496. type: "tool",
  3497. id: "call-before-failure",
  3498. state: { status: "completed", content: [{ type: "text", text: "settle" }] },
  3499. },
  3500. ],
  3501. },
  3502. ])
  3503. const assistant = requireAssistant(context)
  3504. expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([
  3505. "session.step.started.1",
  3506. "session.tool.called.1",
  3507. "session.tool.success.2",
  3508. "session.step.failed.1",
  3509. ])
  3510. }),
  3511. )
  3512. it.effect("durably fails blocked local tools when a step is interrupted", () =>
  3513. Effect.gen(function* () {
  3514. const session = yield* setup
  3515. yield* admit(session, "Interrupt blocked tool")
  3516. toolExecutionGate = yield* Deferred.make<void>()
  3517. responseStream = Stream.concat(
  3518. Stream.fromIterable([
  3519. LLMEvent.stepStart({ index: 0 }),
  3520. LLMEvent.toolCall({ id: "call-before-interrupt", name: "echo", input: { text: "blocked" } }),
  3521. ]),
  3522. Stream.never,
  3523. )
  3524. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3525. while (executions.length === 0) yield* Effect.yieldNow
  3526. yield* session.interrupt(sessionID)
  3527. toolExecutionGate = undefined
  3528. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  3529. yield* session.interrupt(sessionID)
  3530. const context = yield* session.context(sessionID)
  3531. expect(context).toMatchObject([
  3532. { type: "user", text: "Interrupt blocked tool" },
  3533. {
  3534. type: "assistant",
  3535. content: [
  3536. {
  3537. type: "tool",
  3538. id: "call-before-interrupt",
  3539. state: { status: "error", error: { type: "aborted", message: "Tool execution interrupted" } },
  3540. },
  3541. ],
  3542. },
  3543. ])
  3544. const assistant = requireAssistant(context)
  3545. expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([
  3546. "session.step.started.1",
  3547. "session.tool.called.1",
  3548. "session.tool.failed.2",
  3549. "session.step.failed.1",
  3550. ])
  3551. yield* replaySessionProjection(sessionID)
  3552. expect(yield* session.context(sessionID)).toMatchObject([
  3553. { type: "user", text: "Interrupt blocked tool" },
  3554. { type: "assistant", content: [{ type: "tool", id: "call-before-interrupt", state: { status: "error" } }] },
  3555. ])
  3556. requests.length = 0
  3557. responseStream = undefined
  3558. response = []
  3559. yield* session.resume(sessionID)
  3560. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
  3561. }),
  3562. )
  3563. it.effect("interrupts a blocked step without local tool execution", () =>
  3564. Effect.gen(function* () {
  3565. const session = yield* setup
  3566. yield* admit(session, "Interrupt provider")
  3567. streamGate = yield* Deferred.make<void>()
  3568. streamStarted = yield* Deferred.make<void>()
  3569. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3570. yield* Deferred.await(streamStarted)
  3571. yield* session.interrupt(sessionID)
  3572. const exit = yield* Fiber.await(run)
  3573. streamGate = undefined
  3574. streamStarted = undefined
  3575. expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBeTrue()
  3576. expect(requests).toHaveLength(1)
  3577. expect(yield* session.context(sessionID)).toMatchObject([
  3578. { type: "user", text: "Interrupt provider" },
  3579. { type: "assistant", finish: "error", error: { type: "aborted", message: "Step interrupted" } },
  3580. ])
  3581. expect(yield* recordedEventTypes(sessionID)).toContain("session.step.failed.1")
  3582. yield* session.interrupt(sessionID)
  3583. }),
  3584. )
  3585. it.effect("durably fails blocked local tools when interrupted while awaiting settlement", () =>
  3586. Effect.gen(function* () {
  3587. const session = yield* setup
  3588. yield* admit(session, "Interrupt tool settlement")
  3589. toolExecutionGate = yield* Deferred.make<void>()
  3590. toolExecutionsStarted = yield* Deferred.make<void>()
  3591. toolExecutionsReady = 1
  3592. response = reply.tool("call-await-interrupt", "echo", { text: "blocked" })
  3593. const runner = yield* SessionRunner.Service
  3594. const run = yield* runner.drain({ sessionID, force: true }).pipe(Effect.forkChild)
  3595. yield* Deferred.await(toolExecutionsStarted)
  3596. yield* Fiber.interrupt(run)
  3597. toolExecutionGate = undefined
  3598. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  3599. expect(yield* session.context(sessionID)).toMatchObject([
  3600. { type: "user", text: "Interrupt tool settlement" },
  3601. {
  3602. type: "assistant",
  3603. finish: "error",
  3604. error: { type: "aborted", message: "Step interrupted" },
  3605. content: [
  3606. {
  3607. type: "tool",
  3608. id: "call-await-interrupt",
  3609. state: { status: "error", error: { type: "aborted", message: "Tool execution interrupted" } },
  3610. },
  3611. ],
  3612. },
  3613. ])
  3614. const eventTypes = yield* recordedEventTypes(sessionID)
  3615. expect(eventTypes).toContain("session.step.failed.1")
  3616. expect(eventTypes).not.toContain("session.step.ended.1")
  3617. }),
  3618. )
  3619. it.effect("forces a text response on an agent's configured final step", () =>
  3620. Effect.gen(function* () {
  3621. const session = yield* setup
  3622. const agents = yield* AgentV2.Service
  3623. yield* agents.transform((editor) =>
  3624. editor.update(AgentV2.ID.make("build"), (agent) => {
  3625. agent.steps = 2
  3626. }),
  3627. )
  3628. yield* admit(session, "Finish at the limit")
  3629. responses = [
  3630. reply.tool("call-terminal", "echo", { text: "done" }),
  3631. reply.tool("call-forbidden", "echo", { text: "forbidden" }),
  3632. ]
  3633. yield* session.resume(sessionID)
  3634. expect(requests).toHaveLength(2)
  3635. expect(requests[0]?.toolChoice).toBeUndefined()
  3636. expect(requests[1]?.toolChoice).toMatchObject({ type: "none" })
  3637. // Protocols with native "none" keep these definitions for prompt caching.
  3638. expect(requests[1]?.tools.map((tool) => tool.name)).toContain("echo")
  3639. expect(requests[1]?.messages.at(-1)).toMatchObject({
  3640. role: "assistant",
  3641. content: [{ type: "text", text: expect.stringContaining("MAXIMUM STEPS REACHED") }],
  3642. })
  3643. expect(executions).toEqual(["done"])
  3644. expect(yield* session.context(sessionID)).toMatchObject([
  3645. { type: "user", text: "Finish at the limit" },
  3646. { type: "assistant", content: [{ type: "tool", id: "call-terminal", state: { status: "completed" } }] },
  3647. { type: "assistant", content: [{ type: "tool", id: "call-forbidden", state: { status: "error" } }] },
  3648. ])
  3649. }),
  3650. )
  3651. it.effect("resets the configured step allowance when steering input promotes", () =>
  3652. Effect.gen(function* () {
  3653. const session = yield* setup
  3654. const agents = yield* AgentV2.Service
  3655. yield* agents.transform((editor) =>
  3656. editor.update(AgentV2.ID.make("build"), (agent) => {
  3657. agent.steps = 2
  3658. }),
  3659. )
  3660. yield* admit(session, "Start work")
  3661. responses = [
  3662. reply.tool("call-before-steer", "echo", { text: "before" }),
  3663. reply.tool("call-after-steer", "echo", { text: "after" }),
  3664. reply.stop(),
  3665. ]
  3666. streamGate = yield* Deferred.make<void>()
  3667. streamStarted = yield* Deferred.make<void>()
  3668. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3669. yield* Deferred.await(streamStarted)
  3670. yield* session.prompt({ sessionID, text: "Change direction" })
  3671. yield* Deferred.succeed(streamGate, undefined)
  3672. yield* Fiber.join(run)
  3673. streamGate = undefined
  3674. streamStarted = undefined
  3675. expect(requests).toHaveLength(3)
  3676. expect(requests[1]?.toolChoice).toBeUndefined()
  3677. expect(requests[1]?.tools).not.toEqual([])
  3678. expect(requests[2]?.toolChoice).toMatchObject({ type: "none" })
  3679. expect(executions).toEqual(["before", "after"])
  3680. }),
  3681. )
  3682. it.effect("projects provider errors as terminal assistant step failures", () =>
  3683. Effect.gen(function* () {
  3684. const session = yield* setup
  3685. yield* admit(session, "Fail durably")
  3686. response = [LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "Provider unavailable" })]
  3687. expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable")
  3688. expect(requests).toHaveLength(1)
  3689. expect(yield* session.context(sessionID)).toMatchObject([
  3690. { type: "user", text: "Fail durably" },
  3691. { type: "assistant", finish: "error", error: { type: "provider.unknown", message: "Provider unavailable" } },
  3692. ])
  3693. }),
  3694. )
  3695. it.effect("projects provider errors emitted before assistant step start", () =>
  3696. Effect.gen(function* () {
  3697. const session = yield* setup
  3698. yield* admit(session, "Fail before step")
  3699. response = [LLMEvent.providerError({ message: "Provider unavailable" })]
  3700. expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable")
  3701. expect(requests).toHaveLength(1)
  3702. expect(yield* session.context(sessionID)).toMatchObject([
  3703. { type: "user", text: "Fail before step" },
  3704. { type: "assistant", finish: "error", error: { type: "provider.unknown", message: "Provider unavailable" } },
  3705. ])
  3706. }),
  3707. )
  3708. it.effect("projects content-filter finishes as visible terminal failures", () =>
  3709. Effect.gen(function* () {
  3710. const session = yield* setup
  3711. yield* admit(session, "Blocked response")
  3712. response = [
  3713. LLMEvent.stepStart({ index: 0 }),
  3714. LLMEvent.textStart({ id: "partial" }),
  3715. LLMEvent.textDelta({ id: "partial", text: "Partial" }),
  3716. LLMEvent.stepFinish({
  3717. index: 0,
  3718. reason: { normalized: "content-filter" },
  3719. usage: { nonCachedInputTokens: 8, outputTokens: 3, reasoningTokens: 1 },
  3720. }),
  3721. LLMEvent.finish({ reason: { normalized: "content-filter" } }),
  3722. ]
  3723. expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider blocked the response")
  3724. expect(yield* session.context(sessionID)).toMatchObject([
  3725. { type: "user" },
  3726. {
  3727. type: "assistant",
  3728. finish: "error",
  3729. error: { type: "provider.content-filter" },
  3730. cost: 0,
  3731. tokens: { input: 8, output: 2, reasoning: 1, cache: { read: 0, write: 0 } },
  3732. content: [{ type: "text", text: "Partial" }],
  3733. },
  3734. ])
  3735. expect(yield* session.get(sessionID)).toMatchObject({
  3736. cost: 0,
  3737. tokens: { input: 8, output: 2, reasoning: 1, cache: { read: 0, write: 0 } },
  3738. })
  3739. expect(yield* recordedEventTypes(sessionID)).not.toContain("session.step.ended.1")
  3740. }),
  3741. )
  3742. it.effect("settles a local tool before one content-filter step failure", () =>
  3743. Effect.gen(function* () {
  3744. const session = yield* setup
  3745. yield* admit(session, "Tool before blocked response")
  3746. toolExecutionGate = yield* Deferred.make<void>()
  3747. toolExecutionsStarted = yield* Deferred.make<void>()
  3748. toolExecutionsReady = 1
  3749. response = [
  3750. LLMEvent.stepStart({ index: 0 }),
  3751. LLMEvent.toolCall({ id: "call-before-content-filter", name: "echo", input: { text: "settled" } }),
  3752. LLMEvent.stepFinish({ index: 0, reason: { normalized: "content-filter" } }),
  3753. LLMEvent.finish({ reason: { normalized: "content-filter" } }),
  3754. ]
  3755. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3756. yield* Deferred.await(toolExecutionsStarted)
  3757. yield* Deferred.succeed(toolExecutionGate, undefined)
  3758. expect((yield* Fiber.join(run).pipe(Effect.flip)).message).toBe("Provider blocked the response")
  3759. toolExecutionGate = undefined
  3760. toolExecutionsStarted = undefined
  3761. const assistant = requireAssistant(yield* session.context(sessionID))
  3762. const events = yield* recordedStepSettlementEvents(sessionID, assistant.id)
  3763. expect(events.map((event) => event.type)).toEqual([
  3764. "session.step.started.1",
  3765. "session.tool.called.1",
  3766. "session.tool.success.2",
  3767. "session.step.failed.1",
  3768. ])
  3769. expect(
  3770. events.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
  3771. ).toHaveLength(1)
  3772. }),
  3773. )
  3774. it.effect("does not recover context overflow after durable assistant output", () =>
  3775. Effect.gen(function* () {
  3776. const session = yield* setup
  3777. yield* admit(session, "Fail after output")
  3778. response = [
  3779. LLMEvent.stepStart({ index: 0 }),
  3780. LLMEvent.textStart({ id: "text-partial" }),
  3781. LLMEvent.textDelta({ id: "text-partial", text: "Partial" }),
  3782. LLMEvent.textEnd({ id: "text-partial" }),
  3783. LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
  3784. ]
  3785. expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long")
  3786. expect(requests).toHaveLength(1)
  3787. expect(yield* session.context(sessionID)).toMatchObject([
  3788. { type: "user", text: "Fail after output" },
  3789. {
  3790. type: "assistant",
  3791. finish: "error",
  3792. error: { message: "prompt too long" },
  3793. content: [{ type: "text", text: "Partial" }],
  3794. },
  3795. ])
  3796. }),
  3797. )
  3798. it.effect("projects raw provider stream failures as terminal assistant step failures", () =>
  3799. Effect.gen(function* () {
  3800. const session = yield* setup
  3801. yield* admit(session, "Fail raw stream durably")
  3802. const failure = invalidRequest()
  3803. responseStream = Stream.fail(failure)
  3804. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  3805. yield* replaySessionProjection(sessionID)
  3806. expect(yield* session.context(sessionID)).toMatchObject([
  3807. { type: "user", text: "Fail raw stream durably" },
  3808. { type: "assistant", finish: "error", error: { type: "provider.invalid-request", message: "Invalid request" } },
  3809. ])
  3810. }),
  3811. )
  3812. it.effect("retries eligible pre-output failures after exponential backoff", () =>
  3813. Effect.gen(function* () {
  3814. const session = yield* setup
  3815. yield* admit(session, "Retry transport")
  3816. responseStream = Stream.fail(providerUnavailable())
  3817. response = reply.text("Recovered", "retry-success")
  3818. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3819. while (requests.length < 1) yield* Effect.yieldNow
  3820. yield* TestClock.adjust("1999 millis")
  3821. expect(requests).toHaveLength(1)
  3822. yield* TestClock.adjust("1 millis")
  3823. yield* Fiber.join(run)
  3824. expect(requests).toHaveLength(2)
  3825. const eventTypes = yield* recordedEventTypes(sessionID)
  3826. expect(eventTypes).toContain("session.retry.scheduled.1")
  3827. expect(eventTypes.filter((type) => type === "session.step.started.1")).toHaveLength(2)
  3828. expect(yield* session.context(sessionID)).toMatchObject([
  3829. { type: "user" },
  3830. { type: "assistant", finish: "stop", content: [{ type: "text", text: "Recovered" }] },
  3831. ])
  3832. yield* replaySessionProjection(sessionID)
  3833. expect((yield* session.context(sessionID)).filter((message) => message.type === "assistant")).toHaveLength(1)
  3834. }),
  3835. )
  3836. it.effect("uses a larger provider retry-after delay", () =>
  3837. Effect.gen(function* () {
  3838. const session = yield* setup
  3839. yield* admit(session, "Retry rate limit")
  3840. responseStream = Stream.fail(rateLimited(5_000))
  3841. response = reply.text("Recovered", "retry-after-success")
  3842. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3843. while (requests.length < 1) yield* Effect.yieldNow
  3844. yield* TestClock.adjust("4999 millis")
  3845. expect(requests).toHaveLength(1)
  3846. yield* TestClock.adjust("1 millis")
  3847. yield* Fiber.join(run)
  3848. expect(requests).toHaveLength(2)
  3849. }),
  3850. )
  3851. it.effect("does not retry eligible failures after observable output", () =>
  3852. Effect.gen(function* () {
  3853. const session = yield* setup
  3854. yield* admit(session, "Do not replay partial output")
  3855. const failure = rateLimited()
  3856. responseStream = Stream.fromIterable([
  3857. LLMEvent.stepStart({ index: 0 }),
  3858. LLMEvent.textStart({ id: "partial-rate-limit" }),
  3859. LLMEvent.textDelta({ id: "partial-rate-limit", text: "Partial" }),
  3860. ]).pipe(Stream.concat(Stream.fail(failure)))
  3861. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  3862. expect(requests).toHaveLength(1)
  3863. expect(yield* recordedEventTypes(sessionID)).not.toContain("session.retry.scheduled.1")
  3864. expect(yield* session.context(sessionID)).toMatchObject([
  3865. { type: "user" },
  3866. {
  3867. type: "assistant",
  3868. finish: "error",
  3869. error: { type: "provider.rate-limit" },
  3870. content: [{ type: "text", text: "Partial" }],
  3871. },
  3872. ])
  3873. }),
  3874. )
  3875. it.effect("stops after five total retry attempts", () =>
  3876. Effect.gen(function* () {
  3877. const session = yield* setup
  3878. yield* admit(session, "Exhaust retries")
  3879. streamFailure = providerUnavailable()
  3880. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3881. while (requests.length < 1) yield* Effect.yieldNow
  3882. for (const [index, delay] of [2_000, 4_000, 8_000, 16_000].entries()) {
  3883. yield* TestClock.adjust(delay)
  3884. while (requests.length < index + 2) yield* Effect.yieldNow
  3885. }
  3886. expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(streamFailure)
  3887. expect(requests).toHaveLength(5)
  3888. const database = (yield* Database.Service).db
  3889. const retries = yield* database
  3890. .select({ data: EventTable.data })
  3891. .from(EventTable)
  3892. .where(eq(EventTable.type, "session.retry.scheduled.1"))
  3893. .orderBy(asc(EventTable.seq))
  3894. .all()
  3895. .pipe(Effect.orDie)
  3896. expect(retries.map((event) => event.data)).toMatchObject([
  3897. { attempt: 2, at: 2_000 },
  3898. { attempt: 3, at: 6_000 },
  3899. { attempt: 4, at: 14_000 },
  3900. { attempt: 5, at: 30_000 },
  3901. ])
  3902. expect((yield* recordedEventTypes(sessionID)).filter((type) => type === "session.step.started.1")).toHaveLength(5)
  3903. const assistant = requireAssistant(yield* session.context(sessionID))
  3904. expect(yield* recordedStepSettlementEvents(sessionID, assistant.id)).toMatchObject([
  3905. { type: "session.step.started.1" },
  3906. { type: "session.step.started.1" },
  3907. { type: "session.step.started.1" },
  3908. { type: "session.step.started.1" },
  3909. { type: "session.step.started.1" },
  3910. { type: "session.step.failed.1" },
  3911. ])
  3912. }),
  3913. )
  3914. it.effect("retries a model call without consuming the logical agent step", () =>
  3915. Effect.gen(function* () {
  3916. const session = yield* setup
  3917. const agents = yield* AgentV2.Service
  3918. yield* agents.transform((editor) =>
  3919. editor.update(AgentV2.ID.make("build"), (agent) => {
  3920. agent.steps = 2
  3921. }),
  3922. )
  3923. yield* admit(session, "Retry without consuming a step")
  3924. const failure = providerUnavailable()
  3925. responseStream = Stream.fail(failure)
  3926. responses = [reply.tool("call-after-retry", "echo", { text: "recovered" }), reply.stop()]
  3927. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  3928. while (requests.length < 1) yield* Effect.yieldNow
  3929. yield* TestClock.adjust("2 seconds")
  3930. yield* Fiber.join(run)
  3931. expect(requests).toHaveLength(3)
  3932. expect(requests[0]?.toolChoice).toBeUndefined()
  3933. expect(requests[0]?.tools.map((tool) => tool.name)).toContain("echo")
  3934. expect(requests[1]?.toolChoice).toBeUndefined()
  3935. expect(requests[1]?.tools.map((tool) => tool.name)).toContain("echo")
  3936. expect(requests[1]?.messages.at(-1)).not.toMatchObject({
  3937. role: "assistant",
  3938. content: [{ type: "text", text: expect.stringContaining("MAXIMUM STEPS REACHED") }],
  3939. })
  3940. expect(requests[2]?.toolChoice).toMatchObject({ type: "none" })
  3941. // The final step keeps tool definitions to preserve provider prompt caching.
  3942. expect(requests[2]?.tools.map((tool) => tool.name)).toContain("echo")
  3943. expect(requests[2]?.messages.at(-1)).toMatchObject({
  3944. role: "assistant",
  3945. content: [{ type: "text", text: expect.stringContaining("MAXIMUM STEPS REACHED") }],
  3946. })
  3947. expect(executions).toEqual(["recovered"])
  3948. const eventTypes = yield* recordedEventTypes(sessionID)
  3949. expect(eventTypes.filter((type) => type === "session.step.started.1")).toHaveLength(3)
  3950. expect(eventTypes.filter((type) => type === "session.retry.scheduled.1")).toHaveLength(1)
  3951. expect((yield* session.context(sessionID)).filter((message) => message.type === "assistant")).toHaveLength(2)
  3952. }),
  3953. )
  3954. it.effect("does not retry non-eligible provider failures", () =>
  3955. Effect.gen(function* () {
  3956. const session = yield* setup
  3957. yield* admit(session, "Do not retry")
  3958. const failure = invalidRequest()
  3959. streamFailure = failure
  3960. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  3961. expect(requests).toHaveLength(1)
  3962. expect(yield* recordedEventTypes(sessionID)).not.toContain("session.retry.scheduled.1")
  3963. }),
  3964. )
  3965. it.effect("settles malformed streamed tool input before the provider failure", () =>
  3966. Effect.gen(function* () {
  3967. const session = yield* setup
  3968. yield* admit(session, "Call a malformed tool")
  3969. const failure = new LLMError({
  3970. module: "test",
  3971. method: "stream",
  3972. reason: new InvalidProviderOutputReason({ message: "Invalid JSON input for tool call echo" }),
  3973. })
  3974. responseStream = Stream.fromIterable([
  3975. LLMEvent.stepStart({ index: 0 }),
  3976. LLMEvent.toolInputStart({ id: "call-malformed", name: "echo" }),
  3977. LLMEvent.toolInputDelta({ id: "call-malformed", name: "echo", text: '{"text":"partial' }),
  3978. ]).pipe(Stream.concat(Stream.fail(failure)))
  3979. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  3980. const assistant = requireAssistant(yield* session.context(sessionID))
  3981. response = reply.stop()
  3982. yield* admit(session, "Continue")
  3983. yield* session.resume(sessionID)
  3984. expect(yield* recordedStepSettlementEvents(sessionID, assistant.id)).toMatchObject([
  3985. { type: "session.step.started.1" },
  3986. {
  3987. type: "session.tool.failed.2",
  3988. data: {
  3989. callID: "call-malformed",
  3990. error: { type: "provider.invalid-output", message: "Invalid JSON input for tool call echo" },
  3991. },
  3992. },
  3993. {
  3994. type: "session.step.failed.1",
  3995. data: { error: { type: "provider.invalid-output", message: "Invalid JSON input for tool call echo" } },
  3996. },
  3997. ])
  3998. }),
  3999. )
  4000. it.effect("continues after malformed local tool input without exposing raw arguments", () =>
  4001. Effect.gen(function* () {
  4002. const session = yield* setup
  4003. yield* admit(session, "Recover malformed tool input")
  4004. const marker = "raw-malformed-marker"
  4005. const raw = `{"text":"${marker}`
  4006. responses = [
  4007. [
  4008. LLMEvent.stepStart({ index: 0 }),
  4009. LLMEvent.toolInputStart({ id: "call-malformed", name: "echo" }),
  4010. LLMEvent.toolInputDelta({ id: "call-malformed", name: "echo", text: raw }),
  4011. LLMEvent.toolInputEnd({ id: "call-malformed", name: "echo" }),
  4012. LLMEvent.toolInputError({
  4013. id: "call-malformed",
  4014. name: "echo",
  4015. raw,
  4016. }),
  4017. LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }),
  4018. LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
  4019. ],
  4020. reply.stop(),
  4021. ]
  4022. yield* session.resume(sessionID)
  4023. expect(requests).toHaveLength(2)
  4024. expect(executions).toEqual([])
  4025. expect(JSON.stringify(requests[1])).not.toContain(marker)
  4026. expect(requests[1]?.messages).toEqual(
  4027. expect.arrayContaining([
  4028. expect.objectContaining({
  4029. role: "assistant",
  4030. content: expect.arrayContaining([
  4031. expect.objectContaining({ type: "tool-call", id: "call-malformed", name: "echo", input: {} }),
  4032. ]),
  4033. }),
  4034. expect.objectContaining({
  4035. role: "tool",
  4036. content: expect.arrayContaining([
  4037. expect.objectContaining({
  4038. type: "tool-result",
  4039. id: "call-malformed",
  4040. result: expect.objectContaining({
  4041. type: "error",
  4042. value: expect.objectContaining({
  4043. error: expect.objectContaining({
  4044. message: "Tool call arguments were malformed JSON and were not executed. Retry with valid JSON.",
  4045. }),
  4046. }),
  4047. }),
  4048. }),
  4049. ]),
  4050. }),
  4051. ]),
  4052. )
  4053. const context = yield* session.context(sessionID)
  4054. const failed = context.find(
  4055. (message): message is SessionMessage.Assistant =>
  4056. message.type === "assistant" && message.content.some((item) => item.type === "tool"),
  4057. )
  4058. expect(failed).toMatchObject({
  4059. content: [
  4060. {
  4061. type: "tool",
  4062. id: "call-malformed",
  4063. executed: false,
  4064. state: {
  4065. status: "error",
  4066. input: {},
  4067. error: {
  4068. type: "tool.input-json",
  4069. message: "Tool call arguments were malformed JSON and were not executed. Retry with valid JSON.",
  4070. },
  4071. },
  4072. },
  4073. ],
  4074. })
  4075. if (!failed) throw new Error("Malformed tool assistant missing")
  4076. expect(failed.error).toBeUndefined()
  4077. expect((yield* recordedStepSettlementEvents(sessionID, failed.id)).map((event) => event.type)).toEqual([
  4078. "session.step.started.1",
  4079. "session.tool.failed.2",
  4080. "session.step.ended.1",
  4081. ])
  4082. const database = (yield* Database.Service).db
  4083. const durable = yield* database
  4084. .select({ type: EventTable.type, data: EventTable.data })
  4085. .from(EventTable)
  4086. .where(eq(EventTable.aggregate_id, sessionID))
  4087. .all()
  4088. .pipe(Effect.orDie)
  4089. expect(durable.find((event) => event.type === "session.tool.input.ended.1")?.data).toMatchObject({
  4090. callID: "call-malformed",
  4091. text: raw,
  4092. })
  4093. }),
  4094. )
  4095. it.effect("settles a valid sibling before recovering malformed tool input", () =>
  4096. Effect.gen(function* () {
  4097. const session = yield* setup
  4098. yield* admit(session, "Run parallel tools")
  4099. toolExecutionGate = yield* Deferred.make<void>()
  4100. toolExecutionsStarted = yield* Deferred.make<void>()
  4101. toolExecutionsReady = 1
  4102. responses = [
  4103. [
  4104. LLMEvent.stepStart({ index: 0 }),
  4105. LLMEvent.toolCall({ id: "call-valid", name: "echo", input: { text: "valid" } }),
  4106. LLMEvent.toolInputError({
  4107. id: "call-malformed",
  4108. name: "echo",
  4109. raw: '{"text":"partial',
  4110. }),
  4111. LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }),
  4112. LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
  4113. ],
  4114. reply.stop(),
  4115. ]
  4116. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  4117. yield* Deferred.await(toolExecutionsStarted)
  4118. expect(requests).toHaveLength(1)
  4119. yield* Deferred.succeed(toolExecutionGate, undefined)
  4120. yield* Fiber.join(run)
  4121. toolExecutionGate = undefined
  4122. toolExecutionsStarted = undefined
  4123. expect(requests).toHaveLength(2)
  4124. expect(executions).toEqual(["valid"])
  4125. const request = requests[1]
  4126. if (!request) throw new Error("Malformed recovery request missing")
  4127. expect(request.messages.flatMap((message) => (message.role === "tool" ? message.content : []))).toEqual(
  4128. expect.arrayContaining([
  4129. expect.objectContaining({ id: "call-valid", type: "tool-result" }),
  4130. expect.objectContaining({ id: "call-malformed", type: "tool-result" }),
  4131. ]),
  4132. )
  4133. }),
  4134. )
  4135. it.effect("does not recover malformed input after sibling execution is interrupted", () =>
  4136. Effect.gen(function* () {
  4137. const session = yield* setup
  4138. yield* admit(session, "Interrupt malformed recovery")
  4139. toolExecutionGate = yield* Deferred.make<void>()
  4140. toolExecutionsStarted = yield* Deferred.make<void>()
  4141. toolExecutionsReady = 1
  4142. response = [
  4143. LLMEvent.stepStart({ index: 0 }),
  4144. LLMEvent.toolCall({ id: "call-valid", name: "echo", input: { text: "blocked" } }),
  4145. LLMEvent.toolInputError({
  4146. id: "call-malformed",
  4147. name: "echo",
  4148. raw: '{"text":"partial',
  4149. }),
  4150. LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }),
  4151. LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
  4152. ]
  4153. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  4154. yield* Deferred.await(toolExecutionsStarted)
  4155. while (
  4156. !(yield* session.context(sessionID)).some(
  4157. (message) =>
  4158. message.type === "assistant" &&
  4159. message.content.some((item) => item.type === "tool" && item.id === "call-malformed"),
  4160. )
  4161. )
  4162. yield* Effect.yieldNow
  4163. yield* session.interrupt(sessionID)
  4164. toolExecutionGate = undefined
  4165. toolExecutionsStarted = undefined
  4166. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  4167. expect(requests).toHaveLength(1)
  4168. expect(yield* session.context(sessionID)).toMatchObject([
  4169. { type: "user", text: "Interrupt malformed recovery" },
  4170. {
  4171. type: "assistant",
  4172. error: { type: "aborted", message: "Step interrupted" },
  4173. content: [
  4174. { type: "tool", id: "call-valid", state: { status: "error", error: { type: "aborted" } } },
  4175. { type: "tool", id: "call-malformed", state: { status: "error" } },
  4176. ],
  4177. },
  4178. ])
  4179. }),
  4180. )
  4181. it.effect("records malformed provider-executed input as executed", () =>
  4182. Effect.gen(function* () {
  4183. const session = yield* setup
  4184. yield* admit(session, "Fail malformed hosted input")
  4185. const failure = new LLMError({
  4186. module: "test",
  4187. method: "stream",
  4188. reason: new InvalidProviderOutputReason({ message: "Invalid hosted tool input" }),
  4189. })
  4190. responseStream = Stream.fromIterable([
  4191. LLMEvent.stepStart({ index: 0 }),
  4192. LLMEvent.toolInputStart({ id: "call-hosted", name: "web_search", providerExecuted: true }),
  4193. LLMEvent.toolInputDelta({ id: "call-hosted", name: "web_search", text: '{"query":"partial' }),
  4194. ]).pipe(Stream.concat(Stream.fail(failure)))
  4195. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  4196. expect(requireAssistant(yield* session.context(sessionID))).toMatchObject({
  4197. error: { type: "provider.invalid-output", message: "Invalid hosted tool input" },
  4198. content: [
  4199. {
  4200. type: "tool",
  4201. id: "call-hosted",
  4202. executed: true,
  4203. state: { status: "error", error: { type: "provider.invalid-output" } },
  4204. },
  4205. ],
  4206. })
  4207. }),
  4208. )
  4209. it.effect("records a provider failure after malformed input", () =>
  4210. Effect.gen(function* () {
  4211. const session = yield* setup
  4212. yield* admit(session, "Fail after malformed input")
  4213. const failure = new LLMError({
  4214. module: "test",
  4215. method: "stream",
  4216. reason: new InvalidProviderOutputReason({ message: "Provider failed after malformed input" }),
  4217. })
  4218. responseStream = Stream.fromIterable([
  4219. LLMEvent.stepStart({ index: 0 }),
  4220. LLMEvent.toolInputError({
  4221. id: "call-malformed",
  4222. name: "echo",
  4223. raw: '{"text":"partial',
  4224. }),
  4225. ]).pipe(Stream.concat(Stream.fail(failure)))
  4226. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  4227. expect(requireAssistant(yield* session.context(sessionID))).toMatchObject({
  4228. error: { type: "provider.invalid-output", message: "Provider failed after malformed input" },
  4229. content: [
  4230. {
  4231. type: "tool",
  4232. id: "call-malformed",
  4233. executed: false,
  4234. state: { status: "error", error: { type: "tool.input-json" } },
  4235. },
  4236. ],
  4237. })
  4238. expect(requests).toHaveLength(1)
  4239. }),
  4240. )
  4241. it.effect("continues after repeated malformed tool input", () =>
  4242. Effect.gen(function* () {
  4243. const session = yield* setup
  4244. yield* admit(session, "Keep producing malformed tools")
  4245. const malformed = (id: string) => [
  4246. LLMEvent.stepStart({ index: 0 }),
  4247. LLMEvent.toolInputError({
  4248. id,
  4249. name: "echo",
  4250. raw: '{"text":"partial',
  4251. }),
  4252. LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }),
  4253. LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
  4254. ]
  4255. responses = [
  4256. malformed("call-first"),
  4257. reply.tool("call-valid-between", "echo", { text: "valid" }),
  4258. malformed("call-second"),
  4259. reply.stop(),
  4260. ]
  4261. yield* session.resume(sessionID)
  4262. expect(requests).toHaveLength(4)
  4263. expect(executions).toEqual(["valid"])
  4264. expect((yield* recordedEventTypes(sessionID)).filter((type) => type === "session.step.failed.1")).toHaveLength(0)
  4265. }),
  4266. )
  4267. it.effect("does not continue malformed tool input past the agent step limit", () =>
  4268. Effect.gen(function* () {
  4269. const session = yield* setup
  4270. const agents = yield* AgentV2.Service
  4271. yield* agents.transform((editor) =>
  4272. editor.update(AgentV2.ID.make("build"), (agent) => {
  4273. agent.steps = 2
  4274. }),
  4275. )
  4276. yield* admit(session, "Stop malformed tools at the step limit")
  4277. const malformed = (id: string) => [
  4278. LLMEvent.stepStart({ index: 0 }),
  4279. LLMEvent.toolInputError({
  4280. id,
  4281. name: "echo",
  4282. raw: '{"text":"partial',
  4283. }),
  4284. LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }),
  4285. LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
  4286. ]
  4287. responses = [malformed("call-first"), malformed("call-at-limit")]
  4288. yield* session.resume(sessionID)
  4289. expect(requests).toHaveLength(2)
  4290. expect(requests[0]?.toolChoice).toBeUndefined()
  4291. expect(requests[1]?.toolChoice).toMatchObject({ type: "none" })
  4292. expect((yield* recordedEventTypes(sessionID)).filter((type) => type === "session.tool.failed.2")).toHaveLength(2)
  4293. }),
  4294. )
  4295. it.effect("does not continue automatically after a provider error follows a local tool call", () =>
  4296. Effect.gen(function* () {
  4297. const session = yield* setup
  4298. yield* admit(session, "Do not continue failed provider")
  4299. toolExecutionGate = yield* Deferred.make<void>()
  4300. toolExecutionsStarted = yield* Deferred.make<void>()
  4301. toolExecutionsReady = 1
  4302. response = [
  4303. LLMEvent.stepStart({ index: 0 }),
  4304. LLMEvent.toolCall({ id: "call-before-provider-error", name: "echo", input: { text: "settled" } }),
  4305. LLMEvent.providerError({ message: "Provider unavailable" }),
  4306. ]
  4307. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  4308. yield* Deferred.await(toolExecutionsStarted)
  4309. yield* Deferred.succeed(toolExecutionGate, undefined)
  4310. expect((yield* Fiber.join(run).pipe(Effect.flip)).message).toBe("Provider unavailable")
  4311. toolExecutionGate = undefined
  4312. toolExecutionsStarted = undefined
  4313. expect(requests).toHaveLength(1)
  4314. expect(executions).toEqual(["settled"])
  4315. const context = yield* session.context(sessionID)
  4316. const assistant = requireAssistant(context)
  4317. expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([
  4318. "session.step.started.1",
  4319. "session.tool.called.1",
  4320. "session.tool.success.2",
  4321. "session.step.failed.1",
  4322. ])
  4323. }),
  4324. )
  4325. it.effect("durably fails a hosted tool when its provider errors before returning a result", () =>
  4326. Effect.gen(function* () {
  4327. const session = yield* setup
  4328. yield* admit(session, "Fail hosted tool durably")
  4329. response = [
  4330. LLMEvent.stepStart({ index: 0 }),
  4331. hostedCall("call-hosted-provider-error", "effect"),
  4332. LLMEvent.providerError({ message: "Provider unavailable" }),
  4333. ]
  4334. expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable")
  4335. expect(requests).toHaveLength(1)
  4336. const context = yield* session.context(sessionID)
  4337. expect(context).toMatchObject([
  4338. { type: "user", text: "Fail hosted tool durably" },
  4339. {
  4340. type: "assistant",
  4341. content: [{ type: "tool", id: "call-hosted-provider-error", state: { status: "error" } }],
  4342. },
  4343. ])
  4344. const assistant = requireAssistant(context)
  4345. expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([
  4346. "session.step.started.1",
  4347. "session.tool.called.1",
  4348. "session.tool.failed.2",
  4349. "session.step.failed.1",
  4350. ])
  4351. }),
  4352. )
  4353. it.effect("preserves a tool defect before provider failure settlement", () =>
  4354. Effect.gen(function* () {
  4355. const session = yield* setup
  4356. yield* admit(session, "Defect while provider fails")
  4357. response = [
  4358. LLMEvent.stepStart({ index: 0 }),
  4359. LLMEvent.toolCall({ id: "call-defect-provider-error", name: "defect", input: {} }),
  4360. LLMEvent.providerError({ message: "Provider unavailable" }),
  4361. ]
  4362. expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable")
  4363. const context = yield* session.context(sessionID)
  4364. const assistant = requireAssistant(context)
  4365. const events = yield* recordedStepSettlementEvents(sessionID, assistant.id)
  4366. expect(events.map((event) => event.type)).toEqual([
  4367. "session.step.started.1",
  4368. "session.tool.called.1",
  4369. "session.tool.failed.2",
  4370. "session.step.failed.1",
  4371. ])
  4372. expect(events[2]?.data.error).toMatchObject({ type: "unknown", message: "unexpected tool defect" })
  4373. }),
  4374. )
  4375. it.effect("preserves the provider failure when tool output persistence also fails", () =>
  4376. Effect.gen(function* () {
  4377. const session = yield* setup
  4378. yield* admit(session, "Storage fails while provider fails")
  4379. response = [
  4380. LLMEvent.stepStart({ index: 0 }),
  4381. LLMEvent.toolCall({ id: "call-store-provider-error", name: "storefail", input: {} }),
  4382. LLMEvent.providerError({ message: "Provider unavailable" }),
  4383. ]
  4384. expect(yield* session.resume(sessionID).pipe(Effect.exit)).toMatchObject({ _tag: "Failure" })
  4385. expect(requireAssistant(yield* session.context(sessionID))).toMatchObject({
  4386. error: { type: "provider.unknown", message: "Provider unavailable" },
  4387. })
  4388. }),
  4389. )
  4390. it.effect("durably fails a hosted tool left unresolved at normal provider EOF", () =>
  4391. Effect.gen(function* () {
  4392. const session = yield* setup
  4393. yield* admit(session, "Fail hosted tool at EOF")
  4394. response = [LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-eof", "effect")]
  4395. expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider did not return a tool result")
  4396. const assistant = requireAssistant(yield* session.context(sessionID))
  4397. const events = yield* recordedStepSettlementEvents(sessionID, assistant.id)
  4398. expect(events.map((event) => event.type)).toEqual([
  4399. "session.step.started.1",
  4400. "session.tool.called.1",
  4401. "session.tool.failed.2",
  4402. "session.step.failed.1",
  4403. ])
  4404. expect(
  4405. events.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
  4406. ).toHaveLength(1)
  4407. yield* replaySessionProjection(sessionID)
  4408. expect(yield* session.context(sessionID)).toMatchObject([
  4409. { type: "user", text: "Fail hosted tool at EOF" },
  4410. {
  4411. type: "assistant",
  4412. finish: "error",
  4413. error: { type: "tool.result-missing" },
  4414. content: [{ type: "tool", id: "call-hosted-eof", state: { status: "error" } }],
  4415. },
  4416. ])
  4417. }),
  4418. )
  4419. it.effect("fails an unresolved hosted tool before one clean step end", () =>
  4420. Effect.gen(function* () {
  4421. const session = yield* setup
  4422. yield* admit(session, "Settle hosted tool before ending")
  4423. response = [
  4424. LLMEvent.stepStart({ index: 0 }),
  4425. hostedCall("call-hosted-clean-end", "effect"),
  4426. LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
  4427. LLMEvent.finish({ reason: { normalized: "stop" } }),
  4428. ]
  4429. yield* session.resume(sessionID)
  4430. const assistant = requireAssistant(yield* session.context(sessionID))
  4431. const events = yield* recordedStepSettlementEvents(sessionID, assistant.id)
  4432. expect(events.map((event) => event.type)).toEqual([
  4433. "session.step.started.1",
  4434. "session.tool.called.1",
  4435. "session.tool.failed.2",
  4436. "session.step.ended.1",
  4437. ])
  4438. expect(
  4439. events.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
  4440. ).toHaveLength(1)
  4441. }),
  4442. )
  4443. it.effect("settles unresolved local and hosted tools before one raw provider failure", () =>
  4444. Effect.gen(function* () {
  4445. const session = yield* setup
  4446. yield* admit(session, "Fail unresolved tools")
  4447. const failure = invalidRequest()
  4448. const providerFailed = yield* Deferred.make<void>()
  4449. toolExecutionGate = yield* Deferred.make<void>()
  4450. responseStream = Stream.concat(
  4451. Stream.fromIterable([
  4452. LLMEvent.stepStart({ index: 0 }),
  4453. LLMEvent.toolCall({ id: "call-local-raw-failure", name: "defect", input: {} }),
  4454. hostedCall("call-hosted-raw-failure-pair", "effect"),
  4455. ]),
  4456. Stream.fromEffect(Deferred.succeed(providerFailed, undefined)).pipe(Stream.flatMap(() => Stream.fail(failure))),
  4457. )
  4458. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  4459. yield* Deferred.await(providerFailed)
  4460. yield* Deferred.succeed(toolExecutionGate, undefined)
  4461. expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure)
  4462. toolExecutionGate = undefined
  4463. const assistant = requireAssistant(yield* session.context(sessionID))
  4464. const events = yield* recordedStepSettlementEvents(sessionID, assistant.id)
  4465. expect(events.map((event) => ({ type: event.type, callID: event.data.callID }))).toEqual([
  4466. { type: "session.step.started.1", callID: undefined },
  4467. { type: "session.tool.called.1", callID: "call-local-raw-failure" },
  4468. { type: "session.tool.called.1", callID: "call-hosted-raw-failure-pair" },
  4469. { type: "session.tool.failed.2", callID: "call-local-raw-failure" },
  4470. { type: "session.tool.failed.2", callID: "call-hosted-raw-failure-pair" },
  4471. { type: "session.step.failed.1", callID: undefined },
  4472. ])
  4473. expect(
  4474. events.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
  4475. ).toHaveLength(1)
  4476. }),
  4477. )
  4478. it.effect("durably fails a hosted tool left unresolved by a raw provider stream failure", () =>
  4479. Effect.gen(function* () {
  4480. const session = yield* setup
  4481. yield* admit(session, "Fail hosted tool on raw failure")
  4482. const failure = providerUnavailable()
  4483. responseStream = Stream.concat(
  4484. Stream.fromIterable([LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-raw-failure", "effect")]),
  4485. Stream.fail(failure),
  4486. )
  4487. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  4488. expect(requests).toHaveLength(1)
  4489. const assistant = requireAssistant(yield* session.context(sessionID))
  4490. const events = yield* recordedStepSettlementEvents(sessionID, assistant.id)
  4491. expect(events.map((event) => event.type)).toEqual([
  4492. "session.step.started.1",
  4493. "session.tool.called.1",
  4494. "session.tool.failed.2",
  4495. "session.step.failed.1",
  4496. ])
  4497. expect(
  4498. events.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
  4499. ).toHaveLength(1)
  4500. yield* replaySessionProjection(sessionID)
  4501. expect(yield* session.context(sessionID)).toMatchObject([
  4502. { type: "user", text: "Fail hosted tool on raw failure" },
  4503. {
  4504. type: "assistant",
  4505. finish: "error",
  4506. error: { type: "provider.transport", message: "Provider unavailable" },
  4507. content: [{ type: "tool", id: "call-hosted-raw-failure", state: { status: "error" } }],
  4508. },
  4509. ])
  4510. }),
  4511. )
  4512. it.effect("rejects a second text start before the open fragment ends", () =>
  4513. Effect.gen(function* () {
  4514. const session = yield* setup
  4515. yield* admit(session, "Two blocks")
  4516. response = [
  4517. LLMEvent.stepStart({ index: 0 }),
  4518. LLMEvent.textStart({ id: "text-1" }),
  4519. LLMEvent.textStart({ id: "text-2" }),
  4520. ]
  4521. const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))
  4522. expect(defect).toBeInstanceOf(Error)
  4523. if (!(defect instanceof Error)) return
  4524. expect(defect.message).toBe("text start before end: text-2")
  4525. }),
  4526. )
  4527. it.effect("projects sequential text fragments as separate content parts", () =>
  4528. Effect.gen(function* () {
  4529. const session = yield* setup
  4530. yield* admit(session, "Two blocks")
  4531. response = [
  4532. LLMEvent.stepStart({ index: 0 }),
  4533. LLMEvent.textStart({ id: "text-1" }),
  4534. LLMEvent.textDelta({ id: "text-1", text: "First" }),
  4535. LLMEvent.textEnd({ id: "text-1" }),
  4536. LLMEvent.textStart({ id: "text-2" }),
  4537. LLMEvent.textDelta({ id: "text-2", text: "Second" }),
  4538. LLMEvent.textEnd({ id: "text-2" }),
  4539. LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
  4540. LLMEvent.finish({ reason: { normalized: "stop" } }),
  4541. ]
  4542. yield* session.resume(sessionID)
  4543. expect(yield* session.context(sessionID)).toMatchObject([
  4544. { type: "user", text: "Two blocks" },
  4545. {
  4546. type: "assistant",
  4547. content: [
  4548. { type: "text", text: "First" },
  4549. { type: "text", text: "Second" },
  4550. ],
  4551. },
  4552. ])
  4553. }),
  4554. )
  4555. for (const kind of fragmentKinds) {
  4556. it.effect(`broadcasts provider ${kind} deltas without storing projection rewrites`, () =>
  4557. verifyEphemeralDeltas(kind),
  4558. )
  4559. it.effect(`durably closes partial ${kind} when the provider stream fails`, () => verifyPartialFlushOnFailure(kind))
  4560. it.effect(`durably closes partial ${kind} when the provider stream is interrupted`, () =>
  4561. verifyPartialFlushOnInterruption(kind),
  4562. )
  4563. }
  4564. it.effect("rejects duplicate streamed text starts", () =>
  4565. Effect.gen(function* () {
  4566. const session = yield* setup
  4567. response = [LLMEvent.textStart({ id: "text-1" }), LLMEvent.textStart({ id: "text-1" })]
  4568. const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))
  4569. expect(defect).toBeInstanceOf(Error)
  4570. if (!(defect instanceof Error)) return
  4571. expect(defect.message).toBe("Duplicate text start: text-1")
  4572. }),
  4573. )
  4574. it.effect("transitions streamed raw tool input to parsed called input", () =>
  4575. Effect.gen(function* () {
  4576. const session = yield* setup
  4577. yield* admit(session, "Call provider tool")
  4578. response = [
  4579. LLMEvent.stepStart({ index: 0 }),
  4580. LLMEvent.toolInputStart({ id: "call-parsed", name: "web_search" }),
  4581. LLMEvent.toolInputDelta({ id: "call-parsed", name: "web_search", text: '{"query":"hello"}' }),
  4582. LLMEvent.toolInputEnd({ id: "call-parsed", name: "web_search" }),
  4583. hostedCall("call-parsed", "hello"),
  4584. LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
  4585. LLMEvent.finish({ reason: { normalized: "stop" } }),
  4586. ]
  4587. yield* session.resume(sessionID)
  4588. expect(yield* session.context(sessionID)).toMatchObject([
  4589. { type: "user", text: "Call provider tool" },
  4590. {
  4591. type: "assistant",
  4592. content: [{ type: "tool", id: "call-parsed", state: { status: "error", input: { query: "hello" } } }],
  4593. },
  4594. ])
  4595. }),
  4596. )
  4597. it.effect("rejects malformed streamed tool input ordering", () =>
  4598. Effect.gen(function* () {
  4599. const session = yield* setup
  4600. response = [LLMEvent.toolInputDelta({ id: "call-1", name: "read", text: "{}" })]
  4601. const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))
  4602. expect(defect).toBeInstanceOf(Error)
  4603. if (!(defect instanceof Error)) return
  4604. expect(defect.message).toBe("Tool input delta before start: call-1")
  4605. }),
  4606. )
  4607. })