| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761276227632764276527662767276827692770277127722773277427752776277727782779278027812782278327842785278627872788278927902791279227932794279527962797279827992800280128022803280428052806280728082809281028112812281328142815281628172818281928202821282228232824282528262827282828292830283128322833283428352836283728382839284028412842284328442845284628472848284928502851285228532854285528562857285828592860286128622863286428652866286728682869287028712872287328742875287628772878287928802881288228832884288528862887288828892890289128922893289428952896289728982899290029012902290329042905290629072908290929102911291229132914291529162917291829192920292129222923292429252926292729282929293029312932293329342935293629372938293929402941294229432944294529462947294829492950295129522953295429552956295729582959296029612962296329642965296629672968296929702971297229732974297529762977297829792980298129822983298429852986298729882989299029912992299329942995299629972998299930003001300230033004300530063007300830093010301130123013301430153016301730183019302030213022302330243025302630273028302930303031303230333034303530363037303830393040304130423043304430453046304730483049305030513052305330543055305630573058305930603061306230633064306530663067306830693070307130723073307430753076307730783079308030813082308330843085308630873088308930903091309230933094309530963097309830993100310131023103310431053106310731083109311031113112311331143115311631173118311931203121312231233124312531263127312831293130313131323133313431353136313731383139314031413142314331443145314631473148314931503151315231533154315531563157315831593160316131623163316431653166316731683169317031713172317331743175317631773178317931803181318231833184318531863187318831893190319131923193319431953196319731983199320032013202320332043205320632073208320932103211321232133214321532163217321832193220322132223223322432253226322732283229323032313232323332343235323632373238323932403241324232433244324532463247324832493250325132523253325432553256325732583259326032613262326332643265326632673268326932703271327232733274327532763277327832793280328132823283328432853286328732883289329032913292329332943295329632973298329933003301330233033304330533063307330833093310331133123313331433153316331733183319332033213322332333243325332633273328332933303331333233333334333533363337333833393340334133423343334433453346334733483349335033513352335333543355335633573358335933603361336233633364336533663367336833693370337133723373337433753376337733783379338033813382338333843385338633873388338933903391339233933394339533963397339833993400340134023403340434053406340734083409341034113412341334143415341634173418341934203421342234233424342534263427342834293430343134323433343434353436343734383439344034413442344334443445344634473448344934503451345234533454345534563457345834593460346134623463346434653466346734683469347034713472347334743475347634773478347934803481348234833484348534863487348834893490349134923493349434953496349734983499350035013502350335043505350635073508350935103511351235133514351535163517351835193520352135223523352435253526352735283529353035313532353335343535353635373538353935403541354235433544354535463547354835493550355135523553355435553556355735583559356035613562356335643565356635673568356935703571357235733574357535763577357835793580358135823583358435853586358735883589359035913592359335943595359635973598359936003601360236033604360536063607360836093610361136123613361436153616361736183619362036213622362336243625362636273628362936303631363236333634363536363637363836393640364136423643364436453646364736483649365036513652365336543655365636573658365936603661366236633664366536663667366836693670367136723673367436753676367736783679368036813682368336843685368636873688368936903691369236933694369536963697369836993700370137023703370437053706370737083709371037113712371337143715371637173718371937203721372237233724372537263727372837293730373137323733373437353736373737383739374037413742374337443745374637473748374937503751375237533754375537563757375837593760376137623763376437653766376737683769377037713772377337743775377637773778377937803781378237833784378537863787378837893790379137923793379437953796379737983799380038013802380338043805380638073808380938103811381238133814381538163817381838193820382138223823382438253826382738283829383038313832383338343835383638373838383938403841384238433844384538463847384838493850385138523853385438553856385738583859386038613862386338643865386638673868386938703871387238733874387538763877387838793880388138823883388438853886388738883889389038913892389338943895389638973898389939003901390239033904390539063907390839093910391139123913391439153916391739183919392039213922392339243925392639273928392939303931393239333934393539363937393839393940394139423943394439453946394739483949395039513952395339543955395639573958395939603961396239633964396539663967396839693970397139723973397439753976397739783979398039813982398339843985398639873988398939903991399239933994399539963997399839994000400140024003400440054006400740084009401040114012401340144015401640174018401940204021402240234024402540264027402840294030403140324033403440354036403740384039404040414042404340444045404640474048404940504051405240534054405540564057405840594060406140624063406440654066406740684069407040714072407340744075407640774078407940804081408240834084408540864087408840894090409140924093409440954096409740984099410041014102410341044105410641074108410941104111411241134114411541164117411841194120412141224123412441254126412741284129413041314132413341344135413641374138413941404141414241434144414541464147414841494150415141524153415441554156415741584159416041614162416341644165416641674168416941704171417241734174417541764177417841794180418141824183418441854186418741884189419041914192419341944195419641974198419942004201420242034204420542064207420842094210421142124213421442154216421742184219422042214222422342244225422642274228422942304231423242334234423542364237423842394240424142424243424442454246424742484249425042514252425342544255425642574258425942604261426242634264426542664267426842694270427142724273427442754276427742784279428042814282428342844285428642874288428942904291429242934294429542964297429842994300430143024303430443054306430743084309431043114312431343144315431643174318431943204321432243234324432543264327432843294330433143324333433443354336433743384339434043414342434343444345434643474348434943504351435243534354435543564357435843594360436143624363436443654366436743684369437043714372437343744375437643774378437943804381438243834384438543864387438843894390439143924393439443954396439743984399440044014402440344044405440644074408440944104411441244134414441544164417441844194420442144224423442444254426442744284429443044314432443344344435443644374438443944404441444244434444444544464447444844494450445144524453445444554456445744584459446044614462446344644465446644674468446944704471447244734474447544764477447844794480448144824483448444854486448744884489449044914492449344944495449644974498449945004501450245034504450545064507450845094510451145124513451445154516451745184519452045214522452345244525452645274528452945304531453245334534453545364537453845394540454145424543454445454546454745484549455045514552455345544555455645574558455945604561456245634564456545664567456845694570457145724573457445754576457745784579458045814582458345844585458645874588458945904591459245934594459545964597459845994600460146024603460446054606460746084609461046114612461346144615461646174618461946204621462246234624462546264627462846294630463146324633463446354636463746384639464046414642464346444645464646474648464946504651465246534654465546564657465846594660466146624663466446654666466746684669467046714672467346744675467646774678467946804681468246834684468546864687468846894690469146924693469446954696469746984699470047014702 |
- import { describe, expect, test } from "bun:test"
- import {
- LLMError,
- LLMEvent,
- LLMRequest,
- Message,
- Model,
- SystemPart,
- ToolFailure,
- TransportReason,
- InvalidProviderOutputReason,
- InvalidRequestReason,
- RateLimitReason,
- } from "@opencode-ai/ai"
- import * as OpenAIChat from "@opencode-ai/ai/protocols/openai-chat"
- import { TestLLM } from "@opencode-ai/ai/testing"
- import { Catalog } from "@opencode-ai/core/catalog"
- import { Database } from "@opencode-ai/core/database/database"
- import { makeLocationNode } from "@opencode-ai/util/effect/app-node"
- import { AppNodeBuilder } from "@opencode-ai/core/effect/app-node-builder"
- import { LayerNodePlatform } from "@opencode-ai/core/effect/app-node-platform"
- import { LayerNode } from "@opencode-ai/util/effect/layer-node"
- import { Bus } from "@opencode-ai/core/bus"
- import { Event } from "@opencode-ai/schema/event"
- import { App } from "@opencode-ai/core/app"
- import { Permission } from "@opencode-ai/core/permission"
- import { EventTable } from "@opencode-ai/core/event/sql"
- import { Project } from "@opencode-ai/core/project"
- import { ProjectTable } from "@opencode-ai/core/project/sql"
- import { Form } from "@opencode-ai/core/form"
- import { AbsolutePath } from "@opencode-ai/core/schema"
- import { Session } from "@opencode-ai/core/session"
- import { Snapshot } from "@opencode-ai/core/snapshot"
- import { SessionEvent } from "@opencode-ai/core/session/event"
- import { SessionPending } from "@opencode-ai/core/session/pending"
- import { SessionMessage } from "@opencode-ai/core/session/message"
- import { Money } from "@opencode-ai/schema/money"
- import { SessionProjector } from "@opencode-ai/core/session/projector"
- import { SessionExecution } from "@opencode-ai/core/session/execution"
- import { SessionRunCoordinator } from "@opencode-ai/core/session/run-coordinator"
- import { SessionRunner } from "@opencode-ai/core/session/runner"
- import * as SessionRunnerLLM from "@opencode-ai/core/session/runner/llm"
- import { SessionRunnerModel } from "@opencode-ai/core/session/runner/model"
- import { PromptCacheDiagnostics } from "@opencode-ai/core/session/prompt-cache-diagnostics"
- import { SessionUsage } from "@opencode-ai/core/session/usage"
- import { PluginSupervisor } from "@opencode-ai/core/plugin/supervisor"
- import { PluginHooks } from "@opencode-ai/core/plugin/hooks"
- import { SystemPromptPlugin } from "@opencode-ai/core/plugin/system-prompt"
- import { QuestionTool } from "@opencode-ai/core/tool/plugin/question"
- import { Agent } from "@opencode-ai/core/agent"
- import { Config } from "@opencode-ai/core/config"
- import { ConfigCompaction } from "@opencode-ai/core/config/compaction"
- import { Tool } from "@opencode-ai/core/tool"
- import type { Info } from "@opencode-ai/schema/tool"
- import {
- InstructionStateTable,
- SessionPendingTable,
- SessionMessageTable,
- SessionTable,
- } from "@opencode-ai/core/session/sql"
- import { InstructionEntry } from "@opencode-ai/core/session/instruction-entry"
- import { SessionStore } from "@opencode-ai/core/session/store"
- import { Instructions } from "@opencode-ai/core/instructions"
- import { InstructionBuiltIns } from "@opencode-ai/core/instructions/builtins"
- import { InstructionDiscovery } from "@opencode-ai/core/instruction-discovery"
- import { SkillInstructions } from "@opencode-ai/core/skill/instructions"
- import { ReferenceInstructions } from "@opencode-ai/core/reference/instructions"
- import { McpInstructions } from "@opencode-ai/core/mcp/instructions"
- import { ID } from "@opencode-ai/core/model"
- import { Location } from "@opencode-ai/core/location"
- import { Provider } from "@opencode-ai/core/provider"
- import { Cause, DateTime, Deferred, Effect, Exit, Fiber, Layer, Schema, Scope, Stream } from "effect"
- import { TestClock } from "effect/testing"
- import { asc, eq } from "drizzle-orm"
- import { testEffect } from "./lib/effect"
- import { agentHost, catalogHost, host } from "./plugin/host"
- import PROMPT_DEFAULT from "../src/session/runner/prompt/base.txt"
- import { CodeModeInstructions } from "@opencode-ai/core/codemode/instructions"
- let requests: LLMRequest[] = []
- const emptyCodeMode = `\n\n${CodeModeInstructions.render({ total: 0, shown: 0, namespaces: [] })}`
- type ToolBarrier = {
- readonly count: number
- readonly started: Deferred.Deferred<void>
- readonly release: Deferred.Deferred<void>
- active: number
- maxActive: number
- }
- let toolBarrier: ToolBarrier | undefined
- const releaseTools = (barrier: ToolBarrier) =>
- Effect.sync(() => {
- if (toolBarrier === barrier) toolBarrier = undefined
- }).pipe(Effect.andThen(Deferred.succeed(barrier.release, undefined)), Effect.asVoid)
- const blockTools = (count = 1) =>
- Effect.acquireRelease(
- Effect.all({ started: Deferred.make<void>(), release: Deferred.make<void>() }).pipe(
- Effect.map((deferreds) => {
- const barrier = { count, ...deferreds, active: 0, maxActive: 0 }
- toolBarrier = barrier
- return barrier
- }),
- ),
- releaseTools,
- ).pipe(
- Effect.map((barrier) => ({
- started: Deferred.await(barrier.started),
- release: releaseTools(barrier),
- maxActive: Effect.sync(() => barrier.maxActive),
- })),
- )
- const awaitToolBarrier = Effect.suspend(() => {
- const barrier = toolBarrier
- if (!barrier) return Effect.void
- barrier.active++
- barrier.maxActive = Math.max(barrier.maxActive, barrier.active)
- return (barrier.active === barrier.count ? Deferred.succeed(barrier.started, undefined) : Effect.void).pipe(
- Effect.andThen(Deferred.await(barrier.release)),
- Effect.ensuring(Effect.sync(() => barrier.active--)),
- )
- })
- const testLLM = TestLLM.layer({
- fallback: [],
- transformRequest: (request) =>
- LLMRequest.update(request, {
- system: request.system.map((part) => ({
- ...part,
- text: part.text.replace(emptyCodeMode, ""),
- })),
- tools: request.tools.filter((tool) => tool.name !== "execute"),
- }),
- })
- const client = TestLLM.clientLayer
- const model = Model.make({ id: "fake-model", provider: "fake", route: OpenAIChat.route })
- const defaultSystem = PROMPT_DEFAULT
- const replacementModel = Model.make({ id: "replacement", provider: "fake", route: OpenAIChat.route })
- const compactModel = Model.make({
- id: "compact",
- provider: "fake",
- route: OpenAIChat.route.with({ limits: { context: 4_000, output: 50 } }),
- })
- const fullOutputModel = Model.make({
- id: "full-output",
- provider: "fake",
- route: OpenAIChat.route.with({ limits: { context: 262_144, output: 262_144 } }),
- })
- const undersizedContextModel = Model.make({
- id: "undersized-context",
- provider: "fake",
- route: OpenAIChat.route.with({ limits: { context: 1, output: 1_000 } }),
- })
- const recoveryModel = Model.make({
- id: "recovery",
- provider: "fake",
- route: OpenAIChat.route.with({ limits: { context: 20_000, output: 1_000 } }),
- })
- test("calculates step cost using the matching context tier", () => {
- expect(
- SessionUsage.calculateCost(
- [
- {
- input: Money.USDPerMillionTokens.make(1),
- output: Money.USDPerMillionTokens.make(2),
- cache: {
- read: Money.USDPerMillionTokens.make(0.1),
- write: Money.USDPerMillionTokens.make(0.5),
- },
- },
- {
- tier: { type: "context", size: 100 },
- input: Money.USDPerMillionTokens.make(3),
- output: Money.USDPerMillionTokens.make(4),
- cache: {
- read: Money.USDPerMillionTokens.make(0.2),
- write: Money.USDPerMillionTokens.make(0.6),
- },
- },
- ],
- { input: 80, output: 10, reasoning: 2, cache: { read: 20, write: 1 } },
- ),
- ).toBeCloseTo(0.0002926)
- })
- test("does not apply an ineligible tier without base pricing", () => {
- expect(
- SessionUsage.calculateCost(
- [
- {
- tier: { type: "context", size: 100 },
- input: Money.USDPerMillionTokens.make(3),
- output: Money.USDPerMillionTokens.make(4),
- cache: {
- read: Money.USDPerMillionTokens.make(0.2),
- write: Money.USDPerMillionTokens.make(0.6),
- },
- },
- ],
- { input: 80, output: 10, reasoning: 2, cache: { read: 20, write: 0 } },
- ),
- ).toBe(Money.USD.zero)
- })
- const authorizations: Tool.Context[] = []
- const executions: string[] = []
- const permissionFail = {
- name: "permission_fail",
- description: "Reject a permission",
- input: Schema.Struct({}),
- output: Schema.Struct({}),
- execute: () =>
- new ToolFailure({
- message: "Permission denied: edit",
- error: new Permission.BlockedError({
- rules: [],
- permission: "edit",
- resources: ["src/index.ts"],
- }),
- }),
- }
- const permission = Layer.succeed(
- Permission.Service,
- Permission.Service.of({
- assert: () => Effect.die("unused"),
- ask: () => Effect.die("unused"),
- reply: () => Effect.die("unused"),
- get: () => Effect.die("unused"),
- forSession: () => Effect.die("unused"),
- list: () => Effect.die("unused"),
- }),
- )
- const transformTools = (registry: Tool.Interface, tools: Readonly<Record<string, Info>>, options?: Tool.Options) =>
- registry.transform((draft) =>
- Object.entries(tools).forEach(([name, tool]) =>
- draft.add({ ...tool, name, options: options ?? tool.options }),
- ),
- )
- const echo = Layer.effectDiscard(
- Tool.Service.use((registry) =>
- transformTools(
- registry,
- {
- echo: {
- name: "echo",
- description: "Echo text",
- input: Schema.Struct({ text: Schema.String }),
- output: Schema.Struct({ text: Schema.String }),
- execute: ({ text }, context) =>
- Effect.gen(function* () {
- authorizations.push(context)
- executions.push(text)
- yield* awaitToolBarrier
- return { output: { text }, content: text }
- }),
- },
- defect: {
- name: "defect",
- description: "Fail unexpectedly",
- input: Schema.Struct({}),
- output: Schema.Struct({}),
- execute: () => awaitToolBarrier.pipe(Effect.andThen(Effect.die("unexpected tool defect"))),
- },
- storefail: {
- name: "storefail",
- description: "Produce output that cannot be persisted",
- input: Schema.Struct({}),
- output: Schema.Struct({}),
- execute: () => Effect.succeed({ output: {} }),
- },
- },
- { codemode: false },
- ),
- ),
- )
- const echoNode = makeLocationNode({ name: "test/session-runner-tools", layer: echo, deps: [Tool.node] })
- let modelResolveHook = Effect.void
- let currentModel = model
- const models = Layer.mock(SessionRunnerModel.Service)({
- resolve: (session) =>
- modelResolveHook.pipe(
- Effect.as(
- SessionRunnerModel.resolved(session.model?.id === "replacement" ? replacementModel : currentModel, {
- capabilities: { tools: true, input: ["text", "image"], output: ["text"] },
- cost: [],
- variant: session.model?.variant,
- }),
- ),
- ),
- })
- const systemContextKey = Instructions.Key.make("test/context")
- let systemBaseline = "Initial context"
- let systemRemoved = false
- let systemUnavailable = false
- let systemLoadHook = Effect.void
- const skillBaselines = new Map<Agent.ID, string>()
- const systemContext = Layer.mock(InstructionBuiltIns.Service, {
- load: () =>
- Effect.sync(() =>
- Instructions.make({
- key: systemContextKey,
- codec: Schema.toCodecJson(Schema.String),
- read: systemLoadHook.pipe(
- Effect.andThen(
- Effect.sync(() =>
- systemUnavailable ? Instructions.unavailable : systemRemoved ? Instructions.removed : systemBaseline,
- ),
- ),
- ),
- render: {
- initial: String,
- changed: (_previous, current) => current,
- removed: () => "System context source removed: test/context",
- },
- }),
- ),
- })
- const instructionContext = Layer.mock(InstructionDiscovery.Service, { load: () => Effect.succeed(Instructions.empty) })
- const skillInstructions = Layer.mock(SkillInstructions.Service, {
- load: (agent) =>
- Effect.succeed(
- skillBaselines.has(agent.id)
- ? Instructions.make({
- key: Instructions.Key.make("test/skill-guidance"),
- codec: Schema.toCodecJson(Schema.String),
- read: Effect.succeed(skillBaselines.get(agent.id)!),
- render: {
- initial: String,
- changed: (_previous, current) => current,
- removed: () => "Skill guidance removed",
- },
- })
- : Instructions.empty,
- ),
- })
- const referenceInstructions = Layer.mock(ReferenceInstructions.Service, {
- load: () => Effect.succeed(Instructions.empty),
- })
- const mcpInstructions = Layer.mock(McpInstructions.Service, { load: () => Effect.succeed(Instructions.empty) })
- const config = Config.testLayer([
- new Config.Document({
- type: "document",
- info: new Config.Info({
- compaction: new ConfigCompaction.Info({
- buffer: 3_000,
- keep: new ConfigCompaction.Keep({ tokens: 1_000 }),
- }),
- }),
- }),
- ])
- let pluginFlushHook = Effect.void
- const pluginSupervisor = Layer.succeed(
- PluginSupervisor.Service,
- PluginSupervisor.Service.of({
- flush: Effect.suspend(() => pluginFlushHook),
- }),
- )
- const promptCatalog = Layer.mock(Catalog.Service, {
- provider: {
- get: () => Effect.succeed(undefined),
- all: () => Effect.succeed([]),
- available: () => Effect.succeed([]),
- },
- model: {
- get: () => Effect.succeed(undefined),
- all: () => Effect.succeed([]),
- available: () => Effect.succeed([]),
- default: () => Effect.succeed(undefined),
- small: () => Effect.succeed(undefined),
- },
- })
- const runnerLayer = AppNodeBuilder.build(SessionRunnerLLM.node, [
- [Snapshot.node, Snapshot.noopLayer],
- [LayerNodePlatform.llmClient, client],
- [SessionRunnerModel.node, models],
- [InstructionBuiltIns.node, systemContext],
- [InstructionDiscovery.node, instructionContext],
- [Location.node, Location.boundNode({ directory: AbsolutePath.make("/project") })],
- [SkillInstructions.node, skillInstructions],
- [ReferenceInstructions.node, referenceInstructions],
- [Permission.node, permission],
- [Config.node, config],
- [McpInstructions.node, mcpInstructions],
- [PluginSupervisor.node, pluginSupervisor],
- ])
- const execution = Layer.effect(
- SessionExecution.Service,
- Effect.gen(function* () {
- const sessionRunner = yield* SessionRunner.Service
- const coordinator = yield* SessionRunCoordinator.make<Session.ID, SessionRunner.RunError>({
- drain: (sessionID, force) => sessionRunner.drain({ sessionID, force }),
- })
- return SessionExecution.Service.of({
- active: coordinator.active,
- resume: coordinator.run,
- wake: coordinator.wake,
- interrupt: coordinator.interrupt,
- awaitIdle: coordinator.awaitIdle,
- })
- }),
- ).pipe(Layer.provide(runnerLayer))
- const it = testEffect(
- AppNodeBuilder.build(
- LayerNode.group([
- Database.node,
- Bus.node,
- Form.node,
- SessionProjector.node,
- SessionStore.node,
- Agent.node,
- Catalog.node,
- Tool.node,
- Tool.node,
- PluginHooks.node,
- PluginHooks.node,
- echoNode,
- SessionRunnerModel.node,
- InstructionBuiltIns.node,
- InstructionDiscovery.node,
- InstructionEntry.node,
- SkillInstructions.node,
- ReferenceInstructions.node,
- Config.node,
- Snapshot.node,
- SessionRunnerLLM.node,
- SessionExecution.node,
- Session.node,
- ]),
- [
- [LayerNodePlatform.llmClient, client],
- [Permission.node, permission],
- [Catalog.node, promptCatalog],
- [SessionRunnerModel.node, models],
- [InstructionBuiltIns.node, systemContext],
- [InstructionDiscovery.node, instructionContext],
- [Location.node, Location.boundNode({ directory: AbsolutePath.make("/project") })],
- [SkillInstructions.node, skillInstructions],
- [ReferenceInstructions.node, referenceInstructions],
- [Snapshot.node, Snapshot.noopLayer],
- [SessionExecution.node, execution],
- [Config.node, config],
- [PluginSupervisor.node, pluginSupervisor],
- ],
- ).pipe(Layer.provideMerge(testLLM)),
- )
- const sessionID = Session.ID.make("ses_runner_test")
- const otherSessionID = Session.ID.make("ses_runner_other")
- const admit = (session: Session.Interface, text: string) => session.prompt({ sessionID, text, resume: false })
- const runPrompt = Effect.fnUntraced(function* (session: Session.Interface, text: string) {
- const message = yield* admit(session, text)
- yield* session.resume(sessionID)
- return message
- })
- const insertSession = (id: Session.ID) =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- yield* db
- .insert(SessionTable)
- .values({
- id,
- project_id: Project.ID.global,
- slug: id,
- directory: "/project",
- title: "test",
- version: "test",
- })
- .onConflictDoNothing()
- .run()
- .pipe(Effect.orDie)
- })
- const setup = Effect.gen(function* () {
- const { db } = yield* Database.Service
- const agents = yield* Agent.Service
- const catalog = yield* Catalog.Service
- const hooks = yield* PluginHooks.Service
- const pluginHost = host({
- agent: agentHost(agents),
- catalog: catalogHost(catalog),
- session: { hook: (name, callback) => hooks.register("session", name, callback) },
- })
- yield* Effect.forEach(SystemPromptPlugin.Plugins, (plugin) => plugin.effect(pluginHost), {
- discard: true,
- })
- requests = (yield* TestLLM.Service).requests
- authorizations.length = 0
- executions.length = 0
- systemBaseline = "Initial context"
- systemRemoved = false
- systemUnavailable = false
- systemLoadHook = Effect.void
- modelResolveHook = Effect.void
- pluginFlushHook = Effect.void
- currentModel = model
- skillBaselines.clear()
- toolBarrier = undefined
- yield* agents.transform((draft) =>
- draft.update(Agent.ID.make("build"), (agent) => {
- agent.mode = "primary"
- }),
- )
- yield* db
- .insert(ProjectTable)
- .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
- .onConflictDoNothing()
- .run()
- .pipe(Effect.orDie)
- yield* insertSession(sessionID)
- return yield* Session.Service
- })
- const providerUnavailable = () =>
- new LLMError({
- module: "test",
- method: "stream",
- reason: new TransportReason({ message: "Provider unavailable" }),
- })
- const invalidRequest = () =>
- new LLMError({
- module: "test",
- method: "stream",
- reason: new InvalidRequestReason({ message: "Invalid request" }),
- })
- const rateLimited = (retryAfterMs?: number) =>
- new LLMError({
- module: "test",
- method: "stream",
- reason: new RateLimitReason({ message: "Rate limited", retryAfterMs }),
- })
- const setupOverflowRecovery = Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(TestLLM.text("Earlier answer", "text-earlier"))
- yield* runPrompt(session, "Earlier question ".repeat(700))
- currentModel = recoveryModel
- requests.length = 0
- return session
- })
- const messageTexts = (request: LLMRequest, role: "user" | "system") =>
- request.messages.flatMap((message) =>
- message.role === role ? message.content.flatMap((content) => (content.type === "text" ? [content.text] : [])) : [],
- )
- const userTexts = (request: LLMRequest) => messageTexts(request, "user")
- const systemTexts = (request: LLMRequest) => messageTexts(request, "system")
- const messageRoles = (request: LLMRequest | undefined) => request?.messages.map((message) => message.role)
- const recordedEventTypes = (id: Session.ID) =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- return yield* db
- .select({ type: EventTable.type })
- .from(EventTable)
- .where(eq(EventTable.aggregate_id, id))
- .orderBy(asc(EventTable.seq))
- .all()
- .pipe(
- Effect.orDie,
- Effect.map((rows) => rows.map((row) => row.type)),
- )
- })
- const recordedStepSettlementEvents = (id: Session.ID, assistantMessageID: SessionMessage.ID) =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- const settlementTypes = new Set([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.success.2",
- "session.tool.failed.2",
- "session.step.ended.1",
- "session.step.failed.1",
- ])
- return (yield* db
- .select({ type: EventTable.type, data: EventTable.data })
- .from(EventTable)
- .where(eq(EventTable.aggregate_id, id))
- .orderBy(asc(EventTable.seq))
- .all()
- .pipe(Effect.orDie)).filter(
- (event) => settlementTypes.has(event.type) && event.data.assistantMessageID === assistantMessageID,
- )
- })
- const recordedStepSettlementTypes = (id: Session.ID, assistantMessageID: SessionMessage.ID) =>
- recordedStepSettlementEvents(id, assistantMessageID).pipe(Effect.map((events) => events.map((event) => event.type)))
- const hostedCall = (id: string, query: string) =>
- LLMEvent.toolCall({ id, name: "web_search", input: { query }, providerExecuted: true })
- const requireAssistant = (messages: readonly SessionMessage.Info[]) => {
- const assistant = messages.find((message) => message.type === "assistant")
- if (!assistant) throw new Error("Assistant message missing")
- return assistant
- }
- const replaySessionProjection = (id: Session.ID) =>
- Effect.gen(function* () {
- const { db } = yield* Database.Service
- const bus = yield* Bus.Service
- const recorded = yield* db
- .select()
- .from(EventTable)
- .where(eq(EventTable.aggregate_id, id))
- .orderBy(asc(EventTable.seq))
- .all()
- .pipe(Effect.orDie)
- yield* bus.remove(id)
- yield* db.delete(InstructionStateTable).where(eq(InstructionStateTable.session_id, id)).run().pipe(Effect.orDie)
- yield* db.delete(SessionPendingTable).where(eq(SessionPendingTable.session_id, id)).run().pipe(Effect.orDie)
- yield* db.delete(SessionMessageTable).where(eq(SessionMessageTable.session_id, id)).run().pipe(Effect.orDie)
- yield* bus.replayAll(
- recorded.map((event) => ({
- id: event.id,
- created: DateTime.makeUnsafe(event.created),
- aggregateID: event.aggregate_id,
- seq: event.seq,
- type: event.type,
- data: event.data,
- })),
- )
- })
- type FragmentKind = "text" | "reasoning" | "tool input"
- type FragmentFixture = {
- readonly delta: Event.Definition
- readonly completeEvents: LLMEvent[]
- readonly partialEvents: LLMEvent[]
- readonly expectedAssistant: unknown
- readonly expectedContent: unknown
- }
- const fragmentKinds: readonly FragmentKind[] = ["text", "reasoning", "tool input"]
- const fragmentID = (kind: FragmentKind, suffix: string) => `${kind === "tool input" ? "call" : kind}-${suffix}`
- const fragmentFixture = (kind: FragmentKind, id: string, chunks: readonly string[]): FragmentFixture => {
- const text = chunks.join("")
- switch (kind) {
- case "text": {
- const partialEvents = [
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.textStart({ id }),
- ...chunks.map((text) => LLMEvent.textDelta({ id, text })),
- ]
- const expectedContent = { type: "text", text }
- return {
- delta: SessionEvent.Text.Delta,
- partialEvents,
- completeEvents: [
- ...partialEvents,
- LLMEvent.textEnd({ id }),
- LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
- LLMEvent.finish({ reason: { normalized: "stop" } }),
- ],
- expectedAssistant: { type: "assistant", finish: "stop", content: [expectedContent] },
- expectedContent,
- }
- }
- case "reasoning": {
- const partialEvents = [
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.reasoningStart({ id }),
- ...chunks.map((text) => LLMEvent.reasoningDelta({ id, text })),
- ]
- const expectedContent = { type: "reasoning", text }
- return {
- delta: SessionEvent.Reasoning.Delta,
- partialEvents,
- completeEvents: [
- ...partialEvents,
- LLMEvent.reasoningEnd({ id }),
- LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
- LLMEvent.finish({ reason: { normalized: "stop" } }),
- ],
- expectedAssistant: { type: "assistant", finish: "stop", content: [expectedContent] },
- expectedContent,
- }
- }
- case "tool input": {
- const partialEvents = [
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolInputStart({ id, name: "echo" }),
- ...chunks.map((text) => LLMEvent.toolInputDelta({ id, name: "echo", text })),
- ]
- const expectedContent = { type: "tool", id, state: { status: "streaming", input: text } }
- return {
- delta: SessionEvent.Tool.Input.Delta,
- partialEvents,
- completeEvents: [...partialEvents, LLMEvent.toolInputEnd({ id, name: "echo" })],
- expectedAssistant: { type: "assistant", content: [expectedContent] },
- expectedContent,
- }
- }
- }
- }
- const verifyEphemeralDeltas = (kind: FragmentKind) =>
- Effect.gen(function* () {
- const session = yield* setup
- const prompt = `Stream ${kind}`
- const chunks = Array.from({ length: 32 }, (_, index) => `${index},`)
- const fixture = fragmentFixture(kind, fragmentID(kind, "many"), chunks)
- const expectedContext = [{ type: "user", text: prompt }, fixture.expectedAssistant]
- yield* admit(session, prompt)
- const bus = yield* Bus.Service
- const live = yield* bus.subscribe(fixture.delta).pipe(Stream.take(32), Stream.runCollect, Effect.forkScoped)
- yield* Effect.yieldNow
- yield* TestLLM.push(fixture.completeEvents)
- yield* session.resume(sessionID)
- const { db } = yield* Database.Service
- const deltas = yield* db
- .select({ type: EventTable.type })
- .from(EventTable)
- .where(eq(EventTable.type, Bus.versionedType(fixture.delta.type, 1)))
- .all()
- .pipe(Effect.orDie)
- expect(Array.from(yield* Fiber.join(live))).toHaveLength(32)
- expect(deltas).toHaveLength(0)
- expect(yield* session.context(sessionID)).toMatchObject(expectedContext)
- yield* replaySessionProjection(sessionID)
- expect(yield* session.context(sessionID)).toMatchObject(expectedContext)
- })
- const verifyPartialFlushOnFailure = (kind: FragmentKind) =>
- Effect.gen(function* () {
- const session = yield* setup
- const prompt = `Fail after ${kind}`
- const fixture = fragmentFixture(kind, fragmentID(kind, "partial"), ["Partial"])
- const failure = providerUnavailable()
- yield* admit(session, prompt)
- yield* TestLLM.push(TestLLM.failAfter(failure, ...fixture.partialEvents))
- expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: prompt },
- {
- type: "assistant",
- finish: "error",
- error: { type: "provider.transport", message: "Provider unavailable" },
- content: [
- kind === "tool input"
- ? {
- type: "tool",
- id: fragmentID(kind, "partial"),
- state: {
- status: "error",
- error: { type: "provider.transport", message: "Provider unavailable" },
- },
- }
- : fixture.expectedContent,
- ],
- },
- ])
- expect(requests).toHaveLength(1)
- })
- const verifyPartialFlushOnInterruption = (kind: FragmentKind) =>
- Effect.gen(function* () {
- const session = yield* setup
- const prompt = `Interrupt after ${kind}`
- const fixture = fragmentFixture(kind, fragmentID(kind, "interrupted"), ["Partial"])
- const streamed = yield* Deferred.make<void>()
- yield* admit(session, prompt)
- yield* TestLLM.push(
- Stream.concat(
- Stream.fromIterable(fixture.partialEvents),
- Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)),
- ),
- )
- const runner = yield* SessionRunner.Service
- const fiber = yield* runner.drain({ sessionID, force: true }).pipe(Effect.forkChild)
- yield* Deferred.await(streamed)
- yield* Fiber.interrupt(fiber)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: prompt },
- {
- type: "assistant",
- finish: "error",
- error: { type: "aborted", message: "Step interrupted" },
- content: [
- kind === "tool input"
- ? { type: "tool", id: fragmentID(kind, "interrupted"), state: { status: "error" } }
- : fixture.expectedContent,
- ],
- },
- ])
- })
- describe("SessionRunnerLLM", () => {
- it.effect("applies session context hooks without exposing unavailable tools", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const hooks = yield* PluginHooks.Service
- yield* hooks.register("session", "context", (event) =>
- Effect.sync(() => {
- event.system = [SystemPart.make("Hooked system")]
- event.messages = [Message.user("Hooked message")]
- delete event.tools.echo
- event.tools.unregistered = { description: "Unavailable", input: { type: "object" } }
- }),
- )
- yield* admit(session, "Original message")
- yield* TestLLM.push(TestLLM.tool("call-removed", "echo", { text: "blocked" }))
- yield* session.resume(sessionID)
- // A hook-removed call fails independently and continues while step allowance remains.
- expect(requests).toHaveLength(2)
- expect(requests[0]?.system.map((part) => part.text)).toEqual(["Hooked system"])
- expect(requests[0]?.messages).toEqual([Message.user("Hooked message")])
- expect(requests[0]?.tools.map((tool) => tool.name)).not.toContain("echo")
- expect(requests[0]?.tools.map((tool) => tool.name)).not.toContain("unregistered")
- expect(executions).toEqual([])
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Original message" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-removed",
- state: { status: "error", error: { type: "tool.execution" } },
- },
- ],
- },
- ])
- }),
- )
- it.effect("advertises and executes a location registered tool", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const registry = yield* Tool.Service
- const contexts: Tool.Context[] = []
- yield* transformTools(
- registry,
- {
- location_context: {
- name: "location_context",
- description: "Read application context",
- input: Schema.Struct({ query: Schema.String }),
- output: Schema.Struct({ answer: Schema.String }),
- execute: ({ query }, context) =>
- Effect.gen(function* () {
- contexts.push(context)
- yield* context.progress({ phase: "reading" })
- return { output: { answer: query.toUpperCase() } }
- }),
- },
- },
- { codemode: false },
- )
- yield* admit(session, "Use application context")
- yield* TestLLM.push(TestLLM.tool("call-location", "location_context", { query: "hello" }), [])
- const bus = yield* Bus.Service
- const progressFiber = yield* bus.subscribe(SessionEvent.Tool.Progress).pipe(
- Stream.filter((event) => event.data.sessionID === sessionID && event.data.callID === "call-location"),
- Stream.take(1),
- Stream.runCollect,
- Effect.forkScoped({ startImmediately: true }),
- )
- yield* session.resume(sessionID)
- expect(requests[0]?.tools.map((tool) => tool.name)).toContain("location_context")
- expect(contexts).toEqual([
- {
- sessionID,
- agent: Agent.ID.make("build"),
- messageID: expect.stringMatching(/^msg_/),
- callID: Tool.CallID.make("call-location"),
- progress: expect.any(Function),
- },
- ])
- expect(Array.from(yield* Fiber.join(progressFiber))[0]?.data.metadata).toEqual({ phase: "reading" })
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Use application context" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-location",
- state: { status: "completed", content: [{ type: "text", text: '{"answer":"HELLO"}' }] },
- },
- ],
- },
- ])
- }),
- )
- it.effect("executes the tool advertised before a registry reload", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const registry = yield* Tool.Service
- const scope = yield* Scope.make()
- const executions: string[] = []
- yield* transformTools(
- registry,
- {
- reloaded: {
- name: "reloaded",
- description: "Record the advertised tool",
- input: Schema.Struct({}),
- output: Schema.Struct({ value: Schema.String }),
- execute: () =>
- Effect.sync(() => executions.push("advertised")).pipe(Effect.as({ output: { value: "advertised" } })),
- },
- },
- { codemode: false },
- ).pipe(Scope.provide(scope))
- yield* admit(session, "Use the reloaded tool")
- yield* TestLLM.push(TestLLM.tool("call-reloaded", "reloaded", {}), [])
- const stream = yield* TestLLM.gate
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* Scope.close(scope, Exit.void)
- yield* transformTools(
- registry,
- {
- reloaded: {
- name: "reloaded",
- description: "Record the replacement tool",
- input: Schema.Struct({}),
- output: Schema.Struct({ value: Schema.String }),
- execute: () =>
- Effect.sync(() => executions.push("replacement")).pipe(Effect.as({ output: { value: "replacement" } })),
- },
- },
- { codemode: false },
- )
- yield* stream.release
- yield* Fiber.join(run)
- expect(executions).toEqual(["advertised"])
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Use the reloaded tool" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-reloaded",
- state: { status: "completed", content: [{ type: "text", text: '{"value":"advertised"}' }] },
- },
- ],
- },
- ])
- }),
- )
- it.effect("starts a real runner step after default prompt recording", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const message = yield* session.prompt({
- sessionID,
- text: "Run automatically",
- })
- yield* session.wait(sessionID)
- expect(requests).toHaveLength(1)
- expect(yield* session.messages({ sessionID })).toMatchObject([
- { id: message.id, type: "user", text: "Run automatically" },
- ])
- }),
- )
- it.effect("runs a follow-up when synthetic input arrives during an active continuation", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const secondStarted = yield* Deferred.make<void>()
- const releaseSecond = yield* Deferred.make<void>()
- yield* TestLLM.push(
- Stream.fromIterable(TestLLM.tool("call-echo", "echo", { text: "background started" })),
- Stream.unwrap(
- Deferred.succeed(secondStarted, undefined).pipe(
- Effect.andThen(Deferred.await(releaseSecond)),
- Effect.as(Stream.fromIterable(TestLLM.stop())),
- ),
- ),
- Stream.fromIterable(TestLLM.text("Handled completion", "text-completion")),
- )
- yield* admit(session, "Start background work")
- const running = yield* session.resume(sessionID).pipe(Effect.forkChild({ startImmediately: true }))
- yield* Deferred.await(secondStarted)
- yield* session.synthetic({ sessionID, text: "Background work completed" })
- yield* Deferred.succeed(releaseSecond, undefined)
- yield* Fiber.join(running)
- expect(requests).toHaveLength(3)
- expect(userTexts(requests[2])).toContain("Background work completed")
- }),
- )
- it.effect("streams one request with registry definitions from chronological user history", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "First")
- yield* runPrompt(session, "Second")
- expect(requests).toHaveLength(1)
- expect(requests[0]?.model).toBe(model)
- expect(requests[0]?.tools.map((tool) => tool.name)).toEqual(["defect", "echo", "storefail"])
- expect(requests[0]?.messages.map((message) => ({ role: message.role, content: message.content }))).toEqual([
- { role: "user", content: [{ type: "text", text: "First" }] },
- { role: "user", content: [{ type: "text", text: "Second" }] },
- ])
- expect(yield* session.messages({ sessionID })).toHaveLength(2)
- }),
- )
- it.effect("marks the initial instruction sync as baseline metadata", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- const instructionEvents: Event.Payload[] = []
- const unsubscribe = yield* bus.listen((event) =>
- Effect.sync(() => {
- if (event.type === "session.instructions.updated") instructionEvents.push(event)
- }),
- )
- yield* runPrompt(session, "First")
- systemBaseline = "Changed context"
- yield* runPrompt(session, "Second")
- yield* unsubscribe
- expect(instructionEvents).toHaveLength(2)
- expect(instructionEvents[0]?.metadata).toEqual({ instructions: { initial: true } })
- expect(instructionEvents[1]?.metadata).toBeUndefined()
- }),
- )
- it.effect("retries the first request after system context becomes available", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const { db } = yield* Database.Service
- const messageID = SessionMessage.ID.create()
- systemUnavailable = true
- yield* session.prompt({
- id: messageID,
- sessionID,
- text: "First",
- resume: false,
- })
- const exit = yield* session.resume(sessionID).pipe(Effect.exit)
- expect(Exit.isFailure(exit)).toBe(true)
- if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(Instructions.InitializationBlocked)
- expect(requests).toHaveLength(0)
- expect(yield* SessionPending.has(db, sessionID, "steer")).toBe(true)
- expect(
- yield* db.select().from(InstructionStateTable).where(eq(InstructionStateTable.session_id, sessionID)).get(),
- ).toBeUndefined()
- systemUnavailable = false
- yield* session.prompt({ id: messageID, sessionID, text: "First" })
- yield* session.wait(sessionID)
- expect(requests).toHaveLength(1)
- expect(messageRoles(requests[0])).toEqual(["user"])
- }),
- )
- it.effect("interrupts a source Location runner after a Session moves", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- const { db } = yield* Database.Service
- yield* runPrompt(session, "First")
- yield* bus.publish(SessionEvent.Moved, {
- sessionID,
- location: Location.Ref.make({ directory: AbsolutePath.make("/moved") }),
- })
- expect(
- yield* db.select().from(InstructionStateTable).where(eq(InstructionStateTable.session_id, sessionID)).get(),
- ).toBeUndefined()
- yield* admit(session, "Second")
- const exit = yield* session.resume(sessionID).pipe(Effect.exit)
- expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true)
- expect(requests).toHaveLength(1)
- expect(yield* SessionPending.has(db, sessionID, "steer")).toBe(true)
- }),
- )
- it.effect("forks instruction values at the selected message instead of the parent's latest state", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* runPrompt(session, "First")
- systemBaseline = "Changed context"
- const second = yield* runPrompt(session, "Second")
- systemBaseline = "Latest context"
- yield* runPrompt(session, "Third")
- const forked = yield* session.fork({ sessionID, boundary: { type: "before", messageID: second.id } })
- expect(
- yield* (yield* Database.Service).db
- .select()
- .from(InstructionStateTable)
- .where(eq(InstructionStateTable.session_id, forked.id))
- .get(),
- ).toMatchObject({
- initial_values: { "test/context": Instructions.hash("Changed context") },
- current_values: { "test/context": Instructions.hash("Changed context") },
- })
- yield* session.prompt({ sessionID: forked.id, text: "Forked", resume: false })
- yield* session.resume(forked.id)
- expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([defaultSystem, "Changed context"])
- expect(systemTexts(requests.at(-1)!)).toContain("Latest context")
- const { db } = yield* Database.Service
- const bus = yield* Bus.Service
- const recorded = yield* db
- .select()
- .from(EventTable)
- .where(eq(EventTable.aggregate_id, forked.id))
- .orderBy(asc(EventTable.seq))
- .all()
- yield* bus.remove(forked.id)
- yield* db.delete(SessionTable).where(eq(SessionTable.id, forked.id)).run()
- yield* bus.replayAll(
- recorded.map((event) => ({
- id: event.id,
- created: DateTime.makeUnsafe(event.created),
- aggregateID: event.aggregate_id,
- seq: event.seq,
- type: event.type,
- data: event.data,
- })),
- )
- expect(
- yield* db.select().from(InstructionStateTable).where(eq(InstructionStateTable.session_id, forked.id)).get(),
- ).toMatchObject({ current_values: { "test/context": Instructions.hash("Latest context") } })
- }),
- )
- it.effect("keeps nested forks self-contained", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* runPrompt(session, "First")
- systemBaseline = "Changed context"
- const second = yield* runPrompt(session, "Second")
- const child = yield* session.fork({ sessionID, boundary: { type: "before", messageID: second.id } })
- const inheritedFirst = (yield* session.messages({ sessionID: child.id })).find(
- (message) => message.type === "user" && message.text === "First",
- )
- if (!inheritedFirst) return yield* Effect.die(new Error("Nested fork boundary message not found"))
- const grandchild = yield* session.fork({
- sessionID: child.id,
- boundary: { type: "before", messageID: inheritedFirst.id },
- })
- expect(
- yield* (yield* Database.Service).db
- .select()
- .from(InstructionStateTable)
- .where(eq(InstructionStateTable.session_id, grandchild.id))
- .get(),
- ).toMatchObject({
- initial_values: { "test/context": Instructions.hash("Changed context") },
- current_values: { "test/context": Instructions.hash("Changed context") },
- })
- return undefined
- }),
- )
- it.effect("rebuilds a missing instruction cache without admitting another delta", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const { db } = yield* Database.Service
- yield* runPrompt(session, "First")
- yield* db.delete(InstructionStateTable).where(eq(InstructionStateTable.session_id, sessionID)).run()
- yield* admit(session, "Second")
- requests.length = 0
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(1)
- expect(requests[0]?.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
- expect(messageRoles(requests[0])).toEqual(["user", "user"])
- expect(
- yield* db
- .select({ id: EventTable.id })
- .from(EventTable)
- .where(eq(EventTable.type, "session.instructions.updated.2"))
- .all(),
- ).toHaveLength(1)
- expect(yield* db.select().from(InstructionStateTable).get()).toMatchObject({
- initial_values: { "test/context": Instructions.hash("Initial context") },
- current_values: { "test/context": Instructions.hash("Initial context") },
- })
- }),
- )
- it.effect("keeps the initial instructions stable and derives a chronological update from values", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* runPrompt(session, "First")
- systemBaseline = "Changed context"
- yield* runPrompt(session, "Second")
- expect(
- PromptCacheDiagnostics.compare(
- PromptCacheDiagnostics.snapshot(requests[0]),
- PromptCacheDiagnostics.snapshot(requests[1]),
- ),
- ).toEqual({ status: "append-only", previousMessages: 1, currentMessages: 3 })
- expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
- [defaultSystem, "Initial context"],
- [defaultSystem, "Initial context"],
- ])
- expect(messageRoles(requests[1])).toEqual(["user", "system", "user"])
- expect(requests[1]?.messages.at(1)?.content).toEqual([{ type: "text", text: "Changed context" }])
- expect(yield* session.messages({ sessionID })).toHaveLength(2)
- const { db } = yield* Database.Service
- const updates = yield* db
- .select({ data: EventTable.data })
- .from(EventTable)
- .where(eq(EventTable.type, "session.instructions.updated.2"))
- .orderBy(asc(EventTable.seq))
- .all()
- .pipe(Effect.orDie)
- expect(updates).toHaveLength(2)
- expect(updates[0]?.data).toMatchObject({
- sessionID,
- delta: { "test/context": Instructions.hash("Initial context") },
- })
- expect(updates[1]?.data).toEqual({
- sessionID,
- delta: { "test/context": Instructions.hash("Changed context") },
- })
- yield* replaySessionProjection(sessionID)
- expect(yield* session.messages({ sessionID })).toHaveLength(2)
- }),
- )
- it.effect("uses the selected model family prompt when the agent does not override it", () =>
- Effect.gen(function* () {
- const session = yield* setup
- currentModel = Model.make({ id: "gpt-5", provider: "openai", route: OpenAIChat.route })
- yield* admit(session, "First")
- yield* TestLLM.push(TestLLM.text("Done", "text-provider-prompt"))
- yield* session.resume(sessionID)
- expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([
- expect.stringContaining("You are OpenCode, You and the user share the same workspace"),
- "Initial context",
- ])
- }),
- )
- it.effect("uses the selected model family prompt when the agent system override is empty", () =>
- Effect.gen(function* () {
- const session = yield* setup
- currentModel = Model.make({ id: "gpt-5", provider: "openai", route: OpenAIChat.route })
- const agent = yield* Agent.Service
- yield* agent.transform((editor) =>
- editor.update(Agent.ID.make("build"), (agent) => {
- agent.system = ""
- agent.mode = "primary"
- }),
- )
- yield* admit(session, "First")
- yield* TestLLM.push(TestLLM.text("Done", "text-empty-agent-system"))
- yield* session.resume(sessionID)
- expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([
- expect.stringContaining("You are OpenCode, You and the user share the same workspace"),
- "Initial context",
- ])
- }),
- )
- it.effect("includes the effective default agent system before durable context", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const agent = yield* Agent.Service
- yield* agent.transform((editor) =>
- editor.update(Agent.ID.make("build"), (agent) => {
- agent.system = "Build agent instructions"
- agent.mode = "primary"
- }),
- )
- yield* admit(session, "First")
- yield* TestLLM.push(TestLLM.text("Done", "text-build"))
- yield* session.resume(sessionID)
- expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Build agent instructions", "Initial context"])
- }),
- )
- it.effect("uses the configured default agent system for omitted-agent sessions", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const agent = yield* Agent.Service
- yield* agent.transform((editor) => {
- editor.update(Agent.ID.make("build"), (agent) => {
- agent.system = "Build agent instructions"
- agent.mode = "primary"
- })
- editor.update(Agent.ID.make("reviewer"), (agent) => {
- agent.system = "Reviewer instructions"
- agent.mode = "primary"
- })
- editor.default(Agent.ID.make("reviewer"))
- })
- yield* admit(session, "First")
- yield* TestLLM.push(TestLLM.text("Done", "text-reviewer"))
- yield* session.resume(sessionID)
- expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"])
- expect((yield* session.messages({ sessionID }))[0]).toMatchObject({ type: "assistant", agent: "reviewer" })
- }),
- )
- it.effect("uses only the agent prompt and initial instructions as system parts", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const agent = yield* Agent.Service
- yield* agent.transform((editor) =>
- editor.update(Agent.ID.make("build"), (agent) => {
- agent.system = "Build agent instructions"
- agent.mode = "primary"
- }),
- )
- yield* admit(session, "First")
- yield* TestLLM.push(TestLLM.text("Done", "text-no-system"))
- yield* session.resume(sessionID)
- expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Build agent instructions", "Initial context"])
- }),
- )
- it.effect("uses an explicitly selected non-build agent system", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const { db } = yield* Database.Service
- const agent = yield* Agent.Service
- yield* agent.transform((editor) =>
- editor.update(Agent.ID.make("reviewer"), (agent) => {
- agent.system = "Reviewer instructions"
- agent.mode = "primary"
- }),
- )
- yield* db
- .update(SessionTable)
- .set({ agent: "reviewer" })
- .where(eq(SessionTable.id, sessionID))
- .run()
- .pipe(Effect.orDie)
- yield* admit(session, "First")
- yield* TestLLM.push(TestLLM.text("Done", "text-selected"))
- yield* session.resume(sessionID)
- expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"])
- expect((yield* session.messages({ sessionID }))[0]).toMatchObject({ type: "assistant", agent: "reviewer" })
- }),
- )
- it.effect("fails before the model request when the selected agent is unavailable", () =>
- Effect.gen(function* () {
- yield* setup
- const { db } = yield* Database.Service
- yield* db
- .update(SessionTable)
- .set({ agent: "explore" })
- .where(eq(SessionTable.id, sessionID))
- .run()
- .pipe(Effect.orDie)
- const session = yield* Session.Service
- yield* session.prompt({ sessionID, text: "Inspect files", resume: false })
- requests.length = 0
- yield* TestLLM.push([])
- const failure = yield* session.resume(sessionID).pipe(Effect.flip)
- expect(failure).toMatchObject({
- _tag: "Session.AgentNotFoundError",
- sessionID,
- agent: "explore",
- })
- expect(requests).toHaveLength(0)
- }),
- )
- it.effect("waits for initial plugin readiness before constructing the model request", () =>
- Effect.gen(function* () {
- yield* setup
- const release = yield* Deferred.make<void>()
- pluginFlushHook = Deferred.await(release)
- const session = yield* Session.Service
- yield* session.prompt({ sessionID, text: "Wait for plugins", resume: false })
- requests.length = 0
- yield* TestLLM.push([])
- const running = yield* session.resume(sessionID).pipe(Effect.forkChild({ startImmediately: true }))
- yield* Effect.yieldNow
- expect(requests).toHaveLength(0)
- expect(running.pollUnsafe()).toBeUndefined()
- yield* Deferred.succeed(release, undefined)
- yield* Fiber.join(running)
- expect(requests).toHaveLength(1)
- }),
- )
- it.effect("updates selected-agent skill instructions after an agent switch", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- const agents = yield* Agent.Service
- yield* agents.transform((draft) =>
- draft.update(Agent.ID.make("reviewer"), (agent) => {
- agent.mode = "primary"
- }),
- )
- skillBaselines.set(Agent.ID.make("build"), "Build skills")
- yield* runPrompt(session, "First")
- skillBaselines.set(Agent.ID.make("reviewer"), "Reviewer skills")
- yield* bus.publish(SessionEvent.AgentSelected, {
- sessionID,
- agent: Agent.ID.make("reviewer"),
- })
- yield* runPrompt(session, "Second")
- expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
- [defaultSystem, "Initial context\n\nBuild skills"],
- [defaultSystem, "Initial context\n\nBuild skills"],
- ])
- expect(systemTexts(requests[1])).toContainEqual(expect.stringContaining("Reviewer skills"))
- }),
- )
- it.effect("keeps the sampled agent when selection changes during observation", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- skillBaselines.set(Agent.ID.make("build"), "Build skills")
- skillBaselines.set(Agent.ID.make("reviewer"), "Reviewer skills")
- let switched = false
- systemLoadHook = Effect.suspend(() => {
- if (switched) return Effect.void
- switched = true
- return bus
- .publish(SessionEvent.AgentSelected, {
- sessionID,
- agent: Agent.ID.make("reviewer"),
- })
- .pipe(Effect.asVoid)
- })
- yield* runPrompt(session, "First")
- expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
- [defaultSystem, "Initial context\n\nBuild skills"],
- ])
- }),
- )
- it.effect("keeps the sampled model when selection changes during model resolution", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- let switched = false
- modelResolveHook = Effect.suspend(() => {
- if (switched) return Effect.void
- switched = true
- return bus
- .publish(SessionEvent.ModelSelected, {
- sessionID,
- model: { id: ID.make("replacement"), providerID: Provider.ID.make("fake") },
- })
- .pipe(Effect.asVoid)
- })
- yield* runPrompt(session, "First")
- expect(requests.map((request) => request.model)).toEqual([model])
- expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
- [defaultSystem, "Initial context"],
- ])
- }),
- )
- it.effect("admits removed context as a chronological System message", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* runPrompt(session, "First")
- systemRemoved = true
- yield* runPrompt(session, "Second")
- expect(messageRoles(requests[1])).toEqual(["user", "system", "user"])
- expect(requests[1]?.messages.at(1)?.content).toEqual([
- { type: "text", text: "System context source removed: test/context" },
- ])
- expect(yield* session.messages({ sessionID })).toHaveLength(2)
- }),
- )
- it.effect("renders API context entries through add, change, and removal", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const contextEntries = yield* InstructionEntry.Service
- yield* contextEntries.put({ sessionID, key: "deploy-target", value: "production" })
- yield* runPrompt(session, "First")
- // String values render verbatim inside the initial tagged block.
- expect(requests[0]?.system.map((part) => part.text)).toEqual([
- defaultSystem,
- ["Initial context", "", '<context key="deploy-target">', "production", "</context>"].join("\n"),
- ])
- // Non-string JSON pretty-prints; the change narrates as a System update.
- yield* contextEntries.put({ sessionID, key: "deploy-target", value: { region: "us-east-1" } })
- yield* runPrompt(session, "Second")
- expect(messageRoles(requests[1])).toEqual(["user", "system", "user"])
- expect(requests[1]?.messages.at(1)?.content).toEqual([
- {
- type: "text",
- text: [
- 'The context under "deploy-target" changed and supersedes the previous value:',
- '<context key="deploy-target">',
- "{",
- ' "region": "us-east-1"',
- "}",
- "</context>",
- ].join("\n"),
- },
- ])
- expect(yield* contextEntries.list(sessionID)).toEqual([{ key: "deploy-target", value: { region: "us-east-1" } }])
- // Deleting the row announces removal through the stored removal text.
- yield* contextEntries.remove({ sessionID, key: "deploy-target" })
- yield* runPrompt(session, "Third")
- expect(messageRoles(requests[2])).toEqual(["user", "system", "user", "system", "user"])
- expect(requests[2]?.messages.at(-2)?.content).toEqual([
- { type: "text", text: 'The context under "deploy-target" no longer applies. Disregard it.' },
- ])
- expect(yield* contextEntries.list(sessionID)).toEqual([])
- }),
- )
- it.effect("retains JSON null API entries as values", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const entries = yield* InstructionEntry.Service
- yield* entries.put({ sessionID, key: "nullable", value: "present" })
- yield* runPrompt(session, "First")
- yield* entries.put({ sessionID, key: "nullable", value: null })
- yield* runPrompt(session, "Second")
- expect(requests[1]?.messages.at(1)?.content).toEqual([
- {
- type: "text",
- text: [
- 'The context under "nullable" changed and supersedes the previous value:',
- '<context key="nullable">',
- "null",
- "</context>",
- ].join("\n"),
- },
- ])
- expect(yield* entries.list(sessionID)).toEqual([{ key: "nullable", value: null }])
- }),
- )
- it.effect("rejects API instruction entries larger than 8KB", () =>
- Effect.gen(function* () {
- yield* setup
- const entries = yield* InstructionEntry.Service
- const exit = yield* entries
- .put({ sessionID, key: "oversized", value: "x".repeat(InstructionEntry.MaxValueBytes) })
- .pipe(Effect.exit)
- expect(Exit.isFailure(exit)).toBe(true)
- if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(InstructionEntry.ValueTooLargeError)
- expect(yield* entries.list(sessionID)).toEqual([])
- }),
- )
- it.effect("keeps initial instructions and chronological updates after a model switch", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- yield* runPrompt(session, "First")
- systemBaseline = "Changed context"
- yield* runPrompt(session, "Second")
- yield* bus.publish(SessionEvent.ModelSelected, {
- sessionID,
- model: { id: ID.make("replacement"), providerID: Provider.ID.make("fake") },
- })
- systemBaseline = "Replacement context"
- yield* runPrompt(session, "Third")
- expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
- [defaultSystem, "Initial context"],
- [defaultSystem, "Initial context"],
- [defaultSystem, "Initial context"],
- ])
- expect(messageRoles(requests[1])).toEqual(["user", "system", "user"])
- expect(requests[2]?.messages.filter((message) => message.role === "system")).toHaveLength(2)
- expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([
- "user",
- "user",
- "model-switched",
- "user",
- ])
- yield* replaySessionProjection(sessionID)
- expect(yield* session.messages({ sessionID })).toHaveLength(4)
- yield* runPrompt(session, "Fourth")
- }),
- )
- it.effect("preserves instruction values while a source is temporarily unavailable", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- yield* runPrompt(session, "First")
- yield* bus.publish(SessionEvent.ModelSelected, {
- sessionID,
- model: { id: ID.make("replacement"), providerID: Provider.ID.make("fake") },
- })
- systemUnavailable = true
- yield* runPrompt(session, "Second")
- systemUnavailable = false
- systemBaseline = "Replacement context"
- yield* runPrompt(session, "Third")
- expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
- [defaultSystem, "Initial context"],
- [defaultSystem, "Initial context"],
- [defaultSystem, "Initial context"],
- ])
- }),
- )
- it.effect("moves the epoch at compaction and narrates later changes", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- yield* runPrompt(session, "First")
- yield* bus.publish(SessionEvent.Compaction.Started, {
- sessionID,
- reason: "manual",
- recent: "",
- })
- yield* bus.publish(SessionEvent.Compaction.Ended, {
- sessionID,
- reason: "manual",
- text: "summary",
- recent: "",
- })
- systemBaseline = "Replacement context"
- yield* runPrompt(session, "Second")
- expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
- [defaultSystem, "Initial context"],
- [defaultSystem, "Initial context"],
- ])
- expect(messageRoles(requests[1])).toEqual(["user", "system", "user"])
- expect(requests[1]?.messages.at(1)?.content).toEqual([{ type: "text", text: "Replacement context" }])
- yield* replaySessionProjection(sessionID)
- yield* runPrompt(session, "Third")
- }),
- )
- it.effect("runs one durable compaction barrier after tool settlement and before later inputs", () =>
- Effect.gen(function* () {
- const session = yield* setup
- currentModel = recoveryModel
- const stream = yield* TestLLM.gate
- yield* TestLLM.push(
- TestLLM.tool("call-active", "echo", { text: "active" }),
- [LLMEvent.textDelta({ id: "summary", text: "durable summary" })],
- TestLLM.text("Steer complete", "text-steer"),
- TestLLM.text("Queue complete", "text-queue"),
- )
- yield* admit(session, "Active work")
- const active = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- const first = yield* session.compact({ sessionID })
- const second = yield* session.compact({ sessionID })
- expect(second.id).toBe(first.id)
- expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toMatchObject({
- id: first.id,
- })
- expect((yield* session.messages({ sessionID })).find((message) => message.id === first.id)).toBeUndefined()
- yield* admit(session, "Steer after compaction")
- yield* session.synthetic({ sessionID, text: "Completion after compaction", resume: false })
- yield* session.prompt({
- sessionID,
- text: "Queue after compaction",
- delivery: "queue",
- resume: false,
- })
- expect(yield* SessionPending.has((yield* Database.Service).db, sessionID, "steer")).toBe(false)
- yield* stream.release
- yield* Fiber.join(active)
- expect(requests).toHaveLength(4)
- expect(userTexts(requests[1])[0]).toContain("Create a new anchored summary")
- expect(userTexts(requests[2])).toContain("Steer after compaction")
- expect(userTexts(requests[2])).toContain("Completion after compaction")
- expect(userTexts(requests[3])).toContain("Queue after compaction")
- expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
- expect((yield* session.messages({ sessionID })).find((message) => message.id === first.id)).toMatchObject({
- type: "compaction",
- status: "completed",
- summary: "durable summary",
- })
- }),
- )
- it.effect("releases queued prompts when durable compaction fails", () =>
- Effect.gen(function* () {
- const session = yield* setup
- currentModel = recoveryModel
- const stream = yield* TestLLM.gate
- yield* TestLLM.push(
- TestLLM.text("Active complete", "text-active-failure"),
- [],
- TestLLM.text("Continued", "text-after-failure"),
- )
- yield* admit(session, "Active work")
- const active = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- const compaction = yield* session.compact({ sessionID })
- yield* session.prompt({
- sessionID,
- text: "Continue after failure",
- delivery: "queue",
- resume: false,
- })
- yield* stream.release
- yield* Fiber.join(active)
- expect(requests).toHaveLength(3)
- expect(userTexts(requests[2])).toContain("Continue after failure")
- expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
- expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
- type: "compaction",
- status: "failed",
- })
- expect(
- (yield* recordedEventTypes(sessionID)).filter(
- (type) => type === Bus.versionedType(SessionEvent.Compaction.Failed.type, 1),
- ),
- ).toHaveLength(1)
- }),
- )
- it.effect("explains when manual compaction has no history", () =>
- Effect.gen(function* () {
- yield* setup
- const session = yield* Session.Service
- const compaction = yield* session.compact({ sessionID })
- modelResolveHook = Effect.die("model resolution should not run")
- yield* session.resume(sessionID)
- expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
- expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
- type: "compaction",
- status: "failed",
- reason: "manual",
- error: { type: "compaction.unavailable", message: "Nothing to compact yet" },
- })
- expect(
- (yield* recordedEventTypes(sessionID)).filter(
- (type) => type === Bus.versionedType(SessionEvent.Compaction.Failed.type, 1),
- ),
- ).toHaveLength(1)
- }),
- )
- it.effect("manually compacts when the model has no context limit", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-unknown-history"))
- yield* runPrompt(session, "Earlier question")
- requests.length = 0
- yield* TestLLM.push(TestLLM.text("Manual summary", "text-manual-unknown-summary"))
- const compaction = yield* session.compact({ sessionID })
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(1)
- expect(userTexts(requests[0])[0]).toContain("Earlier question")
- expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
- type: "compaction",
- status: "completed",
- summary: "Manual summary",
- })
- }),
- )
- it.effect("preserves provider errors from manual compaction", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-provider-history"))
- yield* runPrompt(session, "Earlier question")
- yield* TestLLM.push([LLMEvent.providerError({ message: "summary unavailable" })])
- const compaction = yield* session.compact({ sessionID })
- yield* session.resume(sessionID)
- expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
- type: "compaction",
- status: "failed",
- error: { type: "provider.error", message: "summary unavailable" },
- })
- }),
- )
- it.effect("preserves typed provider failures from manual compaction", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-failure-history"))
- yield* runPrompt(session, "Earlier question")
- yield* TestLLM.push(Stream.fail(providerUnavailable()))
- const compaction = yield* session.compact({ sessionID })
- yield* session.resume(sessionID)
- expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
- type: "compaction",
- status: "failed",
- error: { type: "provider.transport", message: "Provider unavailable" },
- })
- }),
- )
- it.effect("records cancelled manual compaction without surfacing an internal failure", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-interrupt-history"))
- yield* runPrompt(session, "Earlier question")
- const streamed = yield* Deferred.make<void>()
- const partial = fragmentFixture("text", "text-manual-interrupt-summary", ["Partial summary"])
- yield* TestLLM.push(
- Stream.concat(
- Stream.fromIterable(partial.partialEvents),
- Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)),
- ),
- )
- const compaction = yield* session.compact({ sessionID })
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* Deferred.await(streamed)
- yield* session.interrupt(sessionID)
- yield* Fiber.await(run)
- expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
- expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
- type: "compaction",
- status: "failed",
- reason: "manual",
- error: { type: "aborted", message: "Compaction cancelled" },
- })
- }),
- )
- it.effect("settles an admitted manual compaction when pre-start resolution throws", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-resolution-history"))
- yield* runPrompt(session, "Earlier question")
- const compaction = yield* session.compact({ sessionID })
- modelResolveHook = Effect.die("model resolution failed")
- expect(yield* Effect.exit(session.resume(sessionID))).toMatchObject({ _tag: "Failure" })
- expect(yield* SessionPending.compaction((yield* Database.Service).db, sessionID)).toBeUndefined()
- expect((yield* session.messages({ sessionID })).find((message) => message.id === compaction.id)).toMatchObject({
- type: "compaction",
- status: "failed",
- reason: "manual",
- })
- expect(
- (yield* recordedEventTypes(sessionID)).filter(
- (type) => type === Bus.versionedType(SessionEvent.Compaction.Failed.type, 1),
- ),
- ).toHaveLength(1)
- }),
- )
- it.effect("automatically compacts into a completed summary and retained recent turn", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(TestLLM.textWithUsage("Earlier answer", "text-first", 3_950))
- yield* runPrompt(session, "Earlier question ".repeat(180))
- currentModel = compactModel
- requests.length = 0
- yield* TestLLM.push(
- TestLLM.text("## Objective\n- Preserve the task", "text-summary"),
- TestLLM.textWithUsage("Continued", "text-final", 3_950),
- )
- yield* runPrompt(session, "Recent exact request ".repeat(180))
- expect(requests).toHaveLength(2)
- expect(userTexts(requests[0])[0]).toContain("## Objective")
- expect(userTexts(requests[1])).toHaveLength(1)
- expect(userTexts(requests[1])[0]).toContain("<summary>\n## Objective\n- Preserve the task\n</summary>")
- expect(userTexts(requests[1])[0]).toContain(`[User]: ${"Recent exact request ".repeat(180)}`)
- const context = yield* (yield* SessionStore.Service).context(sessionID)
- expect(context.map((message) => message.type)).toEqual(["compaction", "assistant"])
- expect(context[0]).toMatchObject({
- type: "compaction",
- summary: "## Objective\n- Preserve the task",
- })
- requests.length = 0
- executions.length = 0
- yield* TestLLM.push(
- TestLLM.text("## Objective\n- Preserve the updated task", "text-summary-2"),
- TestLLM.text("Continued again", "text-final-2"),
- )
- yield* runPrompt(session, "Newest exact request ".repeat(180))
- expect(requests).toHaveLength(2)
- expect(userTexts(requests[0])[0]).toContain(
- "<previous-summary>\n## Objective\n- Preserve the task\n</previous-summary>",
- )
- expect(userTexts(requests[0])[0]).toContain("Recent exact request")
- expect((yield* (yield* SessionStore.Service).context(sessionID))[0]).toMatchObject({
- type: "compaction",
- summary: "## Objective\n- Preserve the updated task",
- })
- }),
- )
- it.effect("does not compact immediately when the advertised output limit fills the context", () =>
- Effect.gen(function* () {
- const session = yield* setup
- currentModel = fullOutputModel
- yield* TestLLM.push(TestLLM.textWithUsage("Earlier answer", "text-full-output-first", 9_500))
- yield* runPrompt(session, "Earlier question")
- requests.length = 0
- yield* TestLLM.push(TestLLM.text("Continued", "text-full-output-final"))
- yield* runPrompt(session, "Continue")
- expect(requests).toHaveLength(1)
- expect(userTexts(requests[0])).toContain("Continue")
- expect(yield* session.context(sessionID)).not.toContainEqual(expect.objectContaining({ type: "compaction" }))
- }),
- )
- it.effect("stops after required automatic compaction fails", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(TestLLM.textWithUsage("Earlier answer", "text-before-failed-compaction", 3_950))
- yield* runPrompt(session, "Earlier question ".repeat(180))
- currentModel = compactModel
- requests.length = 0
- yield* TestLLM.push(
- [LLMEvent.providerError({ message: "Unsupported parameter: max_output_tokens" })],
- TestLLM.text("Must not run", "text-after-failed-compaction"),
- )
- yield* admit(session, "Recent exact request ".repeat(180))
- expect(yield* Effect.exit(session.resume(sessionID))).toMatchObject({ _tag: "Failure" })
- expect(requests).toHaveLength(1)
- expect(requests[0]?.generation).toBeUndefined()
- expect(yield* session.context(sessionID)).toContainEqual(
- expect.objectContaining({
- type: "compaction",
- status: "failed",
- reason: "auto",
- error: expect.objectContaining({ message: "Unsupported parameter: max_output_tokens" }),
- }),
- )
- }),
- )
- it.effect("forces one compaction and retries after provider context overflow", () =>
- Effect.gen(function* () {
- const session = yield* setupOverflowRecovery
- yield* TestLLM.push(
- [
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
- ],
- TestLLM.text("## Objective\n- Recover overflow", "text-summary"),
- TestLLM.text("Recovered", "text-final"),
- )
- yield* runPrompt(session, "Continue")
- expect(requests).toHaveLength(3)
- expect(userTexts(requests[1])[0]).toContain("## Objective")
- expect(userTexts(requests[2])[0]).toContain("<summary>\n## Objective\n- Recover overflow\n</summary>")
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "compaction", summary: "## Objective\n- Recover overflow" },
- { type: "assistant", finish: "stop" },
- ])
- yield* replaySessionProjection(sessionID)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "compaction" },
- { type: "assistant", finish: "stop" },
- ])
- }),
- )
- it.effect("recovers from provider context overflow without a configured context limit", () =>
- Effect.gen(function* () {
- const session = yield* setupOverflowRecovery
- currentModel = model
- yield* TestLLM.push(
- [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })],
- TestLLM.text("## Objective\n- Recover unknown limit", "text-summary-unknown-limit"),
- TestLLM.text("Recovered", "text-final-unknown-limit"),
- )
- yield* runPrompt(session, "Continue")
- expect(requests).toHaveLength(3)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "compaction", summary: "## Objective\n- Recover unknown limit" },
- { type: "assistant", finish: "stop" },
- ])
- }),
- )
- it.effect("recovers from provider context overflow despite an undersized configured context limit", () =>
- Effect.gen(function* () {
- const session = yield* setupOverflowRecovery
- currentModel = undersizedContextModel
- yield* TestLLM.push(
- [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })],
- TestLLM.text("## Objective\n- Recover undersized limit", "text-summary-undersized-limit"),
- TestLLM.text("Recovered", "text-final-undersized-limit"),
- )
- yield* runPrompt(session, "Continue")
- expect(requests).toHaveLength(3)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "compaction", summary: "## Objective\n- Recover undersized limit" },
- { type: "assistant", finish: "stop" },
- ])
- }),
- )
- it.effect("persists a second context overflow after one recovery", () =>
- Effect.gen(function* () {
- const session = yield* setupOverflowRecovery
- const overflow = () => [
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
- ]
- yield* TestLLM.push(overflow(), TestLLM.text("## Objective\n- Recover once", "text-summary"), overflow())
- yield* admit(session, "Continue")
- expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long")
- expect(requests).toHaveLength(3)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "compaction" },
- { type: "assistant", finish: "error", error: { message: "prompt too long" } },
- ])
- }),
- )
- it.effect("recovers once from a raw context overflow failure", () =>
- Effect.gen(function* () {
- const session = yield* setupOverflowRecovery
- yield* TestLLM.push(
- Stream.fail(
- new LLMError({
- module: "test",
- method: "stream",
- reason: new InvalidRequestReason({
- message: "prompt too long",
- classification: "context-overflow",
- }),
- }),
- ),
- )
- yield* TestLLM.push(
- TestLLM.text("## Objective\n- Recover raw overflow", "text-summary"),
- TestLLM.text("Recovered", "text-final"),
- )
- yield* runPrompt(session, "Continue")
- expect(requests).toHaveLength(3)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "compaction", summary: "## Objective\n- Recover raw overflow" },
- { type: "assistant", finish: "stop" },
- ])
- }),
- )
- it.effect("publishes the original overflow when recovery summarization fails", () =>
- Effect.gen(function* () {
- const session = yield* setupOverflowRecovery
- yield* TestLLM.push(
- [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })],
- [LLMEvent.providerError({ message: "summary unavailable" })],
- )
- yield* admit(session, "Continue")
- expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long")
- expect(requests).toHaveLength(2)
- const context = yield* session.context(sessionID)
- expect(context).toContainEqual(
- expect.objectContaining({
- type: "compaction",
- status: "failed",
- reason: "auto",
- error: { type: "provider.error", message: "summary unavailable" },
- }),
- )
- expect(context.slice(-3)).toMatchObject([
- { type: "user", text: "Continue" },
- { type: "compaction", status: "failed", reason: "auto" },
- { type: "assistant", finish: "error", error: { message: "prompt too long" } },
- ])
- }),
- )
- it.effect("interrupts overflow recovery while the summary provider is running", () =>
- Effect.gen(function* () {
- const session = yield* setupOverflowRecovery
- yield* TestLLM.push(
- [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })],
- TestLLM.text("## Objective\n- Interrupted", "text-summary"),
- )
- const first = yield* TestLLM.gate
- yield* admit(session, "Continue")
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* first.started
- const summary = yield* TestLLM.gate
- yield* first.release
- yield* summary.started
- yield* session.interrupt(sessionID)
- const exit = yield* Fiber.await(run)
- expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBeTrue()
- expect(yield* session.context(sessionID)).toContainEqual(
- expect.objectContaining({
- type: "compaction",
- status: "failed",
- reason: "auto",
- error: { type: "compaction.interrupted", message: "Compaction was interrupted" },
- }),
- )
- }),
- )
- it.effect("uses epoch values after compaction while a source is unavailable", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- yield* runPrompt(session, "First")
- systemBaseline = "Changed context"
- yield* runPrompt(session, "Second")
- yield* bus.publish(SessionEvent.Compaction.Started, {
- sessionID,
- reason: "manual",
- recent: "",
- })
- yield* bus.publish(SessionEvent.Compaction.Ended, {
- sessionID,
- reason: "manual",
- text: "summary",
- recent: "",
- })
- systemUnavailable = true
- yield* runPrompt(session, "Third")
- // Compaction already moved current values into the new epoch before the unavailable read.
- expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([defaultSystem, "Changed context"])
- expect(systemTexts(requests.at(-1)!)).not.toContain("Changed context")
- }),
- )
- it.effect("projects reasoning and tool events without executing or continuing tools", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Use tools")
- yield* TestLLM.push(
- TestLLM.complete(
- {
- reason: { normalized: "tool-calls" },
- usage: {
- inputTokens: 10,
- nonCachedInputTokens: 8,
- outputTokens: 4,
- reasoningTokens: 1,
- cacheReadInputTokens: 2,
- },
- },
- LLMEvent.reasoningStart({ id: "reasoning-1" }),
- LLMEvent.reasoningDelta({ id: "reasoning-1", text: "Think" }),
- LLMEvent.reasoningEnd({ id: "reasoning-1" }),
- LLMEvent.toolInputStart({ id: "call-error", name: "write" }),
- LLMEvent.toolInputDelta({ id: "call-error", name: "write", text: '{"path":"README.md"}' }),
- LLMEvent.toolInputEnd({ id: "call-error", name: "write" }),
- LLMEvent.toolCall({ id: "call-error", name: "write", input: { path: "README.md" }, providerExecuted: true }),
- LLMEvent.toolError({ id: "call-error", name: "write", message: "Denied" }),
- LLMEvent.toolResult({ id: "call-error", name: "write", result: { type: "error", value: "Denied" } }),
- LLMEvent.toolCall({
- id: "call-provider",
- name: "web_search",
- input: { query: "hello" },
- providerExecuted: true,
- providerMetadata: { openai: { source: "provider" } },
- }),
- LLMEvent.toolResult({
- id: "call-provider",
- name: "web_search",
- result: {
- type: "content",
- value: [
- { type: "text", text: "Hello" },
- { type: "file", uri: "data:image/png;base64,aGVsbG8=", mime: "image/png", name: "hello.png" },
- ],
- },
- providerExecuted: true,
- providerMetadata: { openai: { source: "provider" } },
- }),
- ),
- )
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(1)
- expect(requests[0]?.tools.map((tool) => tool.name)).toEqual(["defect", "echo", "storefail"])
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Use tools" },
- {
- type: "assistant",
- finish: "tool-calls",
- cost: 0,
- tokens: { input: 8, output: 3, reasoning: 1, cache: { read: 2, write: 0 } },
- content: [
- { type: "reasoning", text: "Think" },
- {
- type: "tool",
- id: "call-error",
- name: "write",
- state: {
- status: "error",
- input: { path: "README.md" },
- error: { type: "tool.execution", message: "Denied" },
- },
- },
- {
- type: "tool",
- id: "call-provider",
- name: "web_search",
- executed: true,
- providerState: { source: "provider" },
- providerResultState: { source: "provider" },
- state: {
- status: "completed",
- input: { query: "hello" },
- content: [
- { type: "text", text: "Hello" },
- { type: "file", mime: "image/png", uri: "data:image/png;base64,aGVsbG8=", name: "hello.png" },
- ],
- },
- },
- ],
- },
- ])
- }),
- )
- it.effect("continues with reloaded history after durably settling one local tool call", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Echo this")
- yield* TestLLM.push(TestLLM.tool("call-echo", "echo", { text: "hello" }), TestLLM.text("Done", "text-final"))
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(2)
- expect(messageRoles(requests[1])).toEqual(["user", "assistant", "tool"])
- expect(authorizations).toMatchObject([{ sessionID, callID: "call-echo" }])
- expect(executions).toEqual(["hello"])
- const context = yield* session.context(sessionID)
- expect(context).toMatchObject([
- { type: "user", text: "Echo this" },
- {
- type: "assistant",
- finish: "tool-calls",
- content: [
- {
- type: "tool",
- id: "call-echo",
- name: "echo",
- state: {
- status: "completed",
- input: { text: "hello" },
- content: [{ type: "text", text: "hello" }],
- },
- },
- ],
- },
- { type: "assistant", finish: "stop", content: [{ type: "text", text: "Done" }] },
- ])
- const assistant = requireAssistant(context)
- expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.success.2",
- "session.step.ended.1",
- ])
- }),
- )
- it.effect("reloads a model switch before a tool-driven continuation step", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- yield* admit(session, "Echo this")
- yield* TestLLM.push(TestLLM.tool("call-echo", "echo", { text: "hello" }), TestLLM.stop())
- const tools = yield* blockTools()
- const run = yield* Effect.forkChild(session.resume(sessionID))
- yield* tools.started
- yield* bus.publish(SessionEvent.ModelSelected, {
- sessionID,
- model: { id: ID.make("replacement"), providerID: Provider.ID.make("fake") },
- })
- systemBaseline = "Replacement context"
- yield* tools.release
- yield* Fiber.join(run)
- expect(requests.map((request) => request.model)).toEqual([model, replacementModel])
- expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
- [defaultSystem, "Initial context"],
- [defaultSystem, "Initial context"],
- ])
- expect(systemTexts(requests[1])).toContain("Replacement context")
- }),
- )
- it.effect("restores durable reasoning provider metadata in the next request", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Think first")
- yield* TestLLM.push(
- TestLLM.stop(
- LLMEvent.reasoningStart({ id: "reasoning-anthropic" }),
- LLMEvent.reasoningDelta({ id: "reasoning-anthropic", text: "Signed thought" }),
- LLMEvent.reasoningEnd({
- id: "reasoning-anthropic",
- providerMetadata: { openai: { signature: "sig_1" }, anthropic: { ignored: true } },
- }),
- LLMEvent.reasoningStart({
- id: "reasoning-openai",
- providerMetadata: {
- openai: { itemId: "rs_1", reasoningEncryptedContent: null },
- anthropic: { ignored: true },
- },
- }),
- LLMEvent.reasoningDelta({ id: "reasoning-openai", text: "Encrypted thought" }),
- LLMEvent.reasoningEnd({
- id: "reasoning-openai",
- providerMetadata: {
- openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" },
- anthropic: { ignored: true },
- },
- }),
- ),
- )
- yield* session.resume(sessionID)
- yield* replaySessionProjection(sessionID)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Think first" },
- {
- type: "assistant",
- content: [
- {
- type: "reasoning",
- text: "Signed thought",
- state: { signature: "sig_1" },
- },
- {
- type: "reasoning",
- text: "Encrypted thought",
- state: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" },
- },
- ],
- },
- ])
- yield* admit(session, "Continue")
- yield* TestLLM.push([])
- yield* session.resume(sessionID)
- expect(requests[1]?.messages[1]?.content).toEqual([
- {
- type: "reasoning",
- text: "Signed thought",
- providerMetadata: { openai: { signature: "sig_1" } },
- },
- {
- type: "reasoning",
- text: "Encrypted thought",
- providerMetadata: { openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" } },
- },
- ])
- }),
- )
- it.effect("restores durable text provider metadata in the next request", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Check first")
- yield* TestLLM.push(
- TestLLM.stop(
- LLMEvent.textStart({ id: "commentary", providerMetadata: { openai: { phase: "commentary" } } }),
- LLMEvent.textDelta({ id: "commentary", text: "Checking." }),
- LLMEvent.textEnd({
- id: "commentary",
- providerMetadata: { openai: { phase: "commentary" }, anthropic: { ignored: true } },
- }),
- ),
- )
- yield* session.resume(sessionID)
- yield* replaySessionProjection(sessionID)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Check first" },
- {
- type: "assistant",
- content: [{ type: "text", text: "Checking.", state: { phase: "commentary" } }],
- },
- ])
- yield* admit(session, "Continue")
- yield* TestLLM.push([])
- yield* session.resume(sessionID)
- expect(requests[1]?.messages[1]?.content).toEqual([
- {
- type: "text",
- text: "Checking.",
- providerMetadata: { openai: { phase: "commentary" } },
- },
- ])
- }),
- )
- it.effect("replays durable provider-executed tool results inline in the next request", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Search first")
- yield* TestLLM.push(
- TestLLM.stop(
- LLMEvent.toolCall({
- id: "hosted-search",
- name: "web_search",
- input: { query: "Effect" },
- providerExecuted: true,
- providerMetadata: { openai: { itemId: "hosted-search" }, fake: { ignored: true } },
- }),
- LLMEvent.toolResult({
- id: "hosted-search",
- name: "web_search",
- result: { type: "json", value: [{ title: "Effect" }] },
- providerExecuted: true,
- providerMetadata: { openai: { blockType: "web_search_tool_result" }, anthropic: { ignored: true } },
- }),
- ),
- )
- yield* session.resume(sessionID)
- yield* replaySessionProjection(sessionID)
- yield* admit(session, "Continue")
- yield* TestLLM.push([])
- yield* session.resume(sessionID)
- expect(messageRoles(requests[1])).toEqual(["user", "assistant", "user"])
- expect(requests[1]?.messages[1]?.content).toMatchObject([
- {
- type: "tool-call",
- id: "hosted-search",
- name: "web_search",
- input: { query: "Effect" },
- providerExecuted: true,
- providerMetadata: { openai: { itemId: "hosted-search" } },
- },
- {
- type: "tool-result",
- id: "hosted-search",
- name: "web_search",
- // The generic replay result derives from canonical stored content.
- result: { type: "text", value: '[{"title":"Effect"}]' },
- providerExecuted: true,
- providerMetadata: { openai: { blockType: "web_search_tool_result" } },
- },
- ])
- }),
- )
- it.effect("starts recorded local tools eagerly and awaits settlement before continuing", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Echo five times")
- const tools = yield* blockTools(5)
- const providerGate = yield* Deferred.make<void>()
- const initial = Stream.fromIterable([
- LLMEvent.stepStart({ index: 0 }),
- ...Array.from({ length: 5 }, (_, index) =>
- LLMEvent.toolCall({ id: `call-echo-${index}`, name: "echo", input: { text: `${index}` } }),
- ),
- ])
- const final = Stream.fromIterable([
- LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }),
- LLMEvent.finish({ reason: { normalized: "tool-calls" } }),
- ])
- yield* TestLLM.push(
- Stream.concat(initial, Stream.fromEffect(Deferred.await(providerGate)).pipe(Stream.flatMap(() => final))),
- )
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* tools.started
- expect(executions).toHaveLength(5)
- expect(yield* tools.maxActive).toBe(5)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Echo five times" },
- {
- type: "assistant",
- content: Array.from({ length: 5 }, (_, index) => ({
- type: "tool",
- id: `call-echo-${index}`,
- state: { status: "running", input: { text: `${index}` } },
- })),
- },
- ])
- yield* Deferred.succeed(providerGate, undefined)
- yield* Effect.yieldNow
- expect(requests).toHaveLength(1)
- yield* tools.release
- yield* Fiber.join(run)
- expect(executions).toHaveLength(5)
- expect(yield* tools.maxActive).toBe(5)
- expect(requests).toHaveLength(2)
- }),
- )
- it.effect("settles repeated provider-local tool call IDs against their owning assistant messages", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Echo twice")
- yield* TestLLM.push(
- TestLLM.tool("tool_0", "echo", { text: "first" }),
- TestLLM.tool("tool_0", "echo", { text: "second" }),
- [],
- )
- yield* session.resume(sessionID)
- const expected = [
- { type: "user", text: "Echo twice" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "tool_0",
- state: { status: "completed", content: [{ type: "text", text: "first" }] },
- },
- ],
- },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "tool_0",
- state: { status: "completed", content: [{ type: "text", text: "second" }] },
- },
- ],
- },
- ]
- expect(executions).toEqual(["first", "second"])
- expect(requests).toHaveLength(3)
- expect(yield* session.context(sessionID)).toMatchObject(expected)
- yield* replaySessionProjection(sessionID)
- expect(yield* session.context(sessionID)).toMatchObject(expected)
- }),
- )
- it.effect("joins concurrent resume calls into one active provider run", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Run once")
- yield* TestLLM.push(TestLLM.text("Once", "text-once"))
- const stream = yield* TestLLM.gate
- const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- const second = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* Effect.yieldNow
- expect(requests).toHaveLength(1)
- yield* stream.release
- yield* Fiber.join(first)
- yield* Fiber.join(second)
- expect(requests).toHaveLength(1)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Run once" },
- { type: "assistant", finish: "stop", content: [{ type: "text", text: "Once" }] },
- ])
- }),
- )
- it.effect("steers an active step with newly recorded prompts", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Start working")
- yield* TestLLM.push(TestLLM.stop(), TestLLM.stop())
- const stream = yield* TestLLM.gate
- const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* session.prompt({ sessionID, text: "Change direction" })
- yield* stream.release
- yield* Fiber.join(first)
- yield* Effect.yieldNow
- expect(requests).toHaveLength(2)
- expect(userTexts(requests[0])).toEqual(["Start working"])
- expect(userTexts(requests[1])).toEqual(["Start working", "Change direction"])
- expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([
- "user",
- "assistant",
- "user",
- "assistant",
- ])
- }),
- )
- it.effect("promotes queued input after continuation ends", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Start working")
- yield* TestLLM.push(TestLLM.tool("call-echo", "echo", { text: "hello" }), TestLLM.stop(), TestLLM.stop())
- const stream = yield* TestLLM.gate
- const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* session.prompt({
- sessionID,
- text: "Wait until continuation ends",
- delivery: "queue",
- })
- yield* stream.release
- yield* Fiber.join(first)
- expect(requests).toHaveLength(3)
- expect(userTexts(requests[0])).toEqual(["Start working"])
- expect(userTexts(requests[1])).toEqual(["Start working"])
- expect(userTexts(requests[2])).toEqual(["Start working", "Wait until continuation ends"])
- }),
- )
- it.effect("preserves durable queued input for a later wake after interruption", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const { db } = yield* Database.Service
- yield* admit(session, "Interrupt current work")
- yield* TestLLM.push([], TestLLM.stop())
- const stream = yield* TestLLM.gate
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* session.prompt({
- sessionID,
- text: "Run after interrupt",
- delivery: "queue",
- })
- yield* session.interrupt(sessionID)
- expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
- expect(requests).toHaveLength(1)
- expect(yield* SessionPending.has(db, sessionID, "queue")).toBe(true)
- const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* stream.release
- yield* Fiber.join(resumed)
- expect(requests).toHaveLength(2)
- expect(userTexts(requests[0])).toEqual(["Interrupt current work"])
- expect(userTexts(requests[1])).toEqual(["Interrupt current work", "Run after interrupt"])
- }),
- )
- it.effect("preserves durable steering input for a later resume after interruption", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const { db } = yield* Database.Service
- yield* admit(session, "Interrupt current work")
- yield* TestLLM.push([], TestLLM.stop())
- const stream = yield* TestLLM.gate
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* session.prompt({
- sessionID,
- text: "Steer after interrupt",
- })
- yield* session.interrupt(sessionID)
- expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
- expect(requests).toHaveLength(1)
- expect(yield* SessionPending.has(db, sessionID, "steer")).toBe(true)
- const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* stream.release
- yield* Fiber.join(resumed)
- expect(requests).toHaveLength(2)
- expect(userTexts(requests[0])).toEqual(["Interrupt current work"])
- expect(userTexts(requests[1])).toEqual(["Interrupt current work", "Steer after interrupt"])
- }),
- )
- it.effect("promotes queued inputs one at a time in FIFO order", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Start working")
- yield* TestLLM.push(TestLLM.stop(), TestLLM.stop(), TestLLM.stop())
- const stream = yield* TestLLM.gate
- const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* session.prompt({ sessionID, text: "Queue first", delivery: "queue" })
- yield* session.prompt({ sessionID, text: "Queue second", delivery: "queue" })
- yield* stream.release
- yield* Fiber.join(first)
- expect(requests).toHaveLength(3)
- expect(userTexts(requests[0])).toEqual(["Start working"])
- expect(userTexts(requests[1])).toEqual(["Start working", "Queue first"])
- expect(userTexts(requests[2])).toEqual(["Start working", "Queue first", "Queue second"])
- }),
- )
- it.effect("promotes queued input after steering continuation ends", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Start steering")
- yield* session.prompt({
- sessionID,
- text: "Queue for later",
- delivery: "queue",
- resume: false,
- })
- yield* TestLLM.push(TestLLM.stop(), TestLLM.stop())
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(2)
- expect(userTexts(requests[0])).toEqual(["Start steering"])
- expect(userTexts(requests[1])).toEqual(["Start steering", "Queue for later"])
- }),
- )
- it.effect("promotes steers before the next queued input", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Start working")
- yield* TestLLM.push(TestLLM.stop(), TestLLM.stop(), TestLLM.stop(), TestLLM.stop())
- const firstStream = yield* TestLLM.gate
- const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* firstStream.started
- yield* session.prompt({ sessionID, text: "Queue first", delivery: "queue" })
- yield* session.prompt({ sessionID, text: "Queue second", delivery: "queue" })
- const secondStream = yield* TestLLM.gate
- yield* firstStream.release
- yield* secondStream.started
- yield* session.prompt({ sessionID, text: "Steer before next queued input" })
- yield* session.prompt({
- sessionID,
- text: "Also steer before next queued input",
- })
- yield* session.synthetic({ sessionID, text: "Background completion before next queued input" })
- yield* secondStream.release
- yield* Fiber.join(first)
- expect(requests).toHaveLength(4)
- expect(userTexts(requests[0])).toEqual(["Start working"])
- expect(userTexts(requests[1])).toEqual(["Start working", "Queue first"])
- expect(userTexts(requests[2])).toEqual([
- "Start working",
- "Queue first",
- "Steer before next queued input",
- "Also steer before next queued input",
- "Background completion before next queued input",
- ])
- expect(userTexts(requests[3])).toEqual([
- "Start working",
- "Queue first",
- "Steer before next queued input",
- "Also steer before next queued input",
- "Background completion before next queued input",
- "Queue second",
- ])
- }),
- )
- it.effect("coalesces multiple active steering prompts into one continuation step", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Start working")
- yield* TestLLM.push(TestLLM.stop(), TestLLM.stop())
- const stream = yield* TestLLM.gate
- const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* session.prompt({ sessionID, text: "First steer" })
- yield* session.prompt({ sessionID, text: "Second steer" })
- yield* stream.release
- yield* Fiber.join(first)
- yield* Effect.yieldNow
- expect(requests).toHaveLength(2)
- expect(userTexts(requests[1])).toEqual(["Start working", "First steer", "Second steer"])
- yield* (yield* SessionExecution.Service).wake(sessionID)
- yield* Effect.yieldNow
- expect(requests).toHaveLength(2)
- }),
- )
- it.effect("runs steering input accepted while the active step fails", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Start working")
- const failure = invalidRequest()
- yield* TestLLM.push(Stream.fail(failure))
- const stream = yield* TestLLM.gate
- const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* session.prompt({ sessionID, text: "Recover with this" })
- yield* stream.release
- expect(yield* Fiber.join(first).pipe(Effect.flip)).toBe(failure)
- yield* TestLLM.push([])
- yield* session.wait(sessionID)
- expect(requests).toHaveLength(2)
- expect(userTexts(requests[1])).toEqual(["Start working", "Recover with this"])
- }),
- )
- it.effect("durably fails local tools left running by a prior process before continuing", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- yield* admit(session, "Recover interrupted tool")
- yield* SessionPending.promote((yield* Database.Service).db, bus, sessionID, "steer")
- const assistantMessageID = SessionMessage.ID.create()
- yield* bus.publish(SessionEvent.Step.Started, {
- sessionID,
- assistantMessageID,
- agent: Agent.ID.make("build"),
- model: { id: ID.make("fake-model"), providerID: Provider.ID.make("fake") },
- })
- yield* bus.publish(SessionEvent.Tool.Input.Started, {
- sessionID,
- assistantMessageID,
- callID: "call-interrupted",
- name: "echo",
- })
- yield* bus.publish(SessionEvent.Tool.Input.Ended, {
- sessionID,
- assistantMessageID,
- callID: "call-interrupted",
- text: '{"text":"stale"}',
- })
- yield* bus.publish(SessionEvent.Tool.Called, {
- sessionID,
- assistantMessageID,
- callID: "call-interrupted",
- input: { text: "stale" },
- executed: false,
- })
- requests.length = 0
- yield* TestLLM.push([])
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(1)
- expect(messageRoles(requests[0])).toEqual(["user", "assistant", "tool"])
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Recover interrupted tool" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-interrupted",
- state: {
- status: "error",
- error: { type: "aborted", message: "Tool execution interrupted: echo" },
- },
- },
- ],
- },
- ])
- }),
- )
- it.effect("durably fails hosted tools left running by a prior process before continuing inline", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- yield* admit(session, "Recover interrupted hosted tool")
- yield* SessionPending.promote((yield* Database.Service).db, bus, sessionID, "steer")
- const assistantMessageID = SessionMessage.ID.create()
- yield* bus.publish(SessionEvent.Step.Started, {
- sessionID,
- assistantMessageID,
- agent: Agent.ID.make("build"),
- model: { id: ID.make("fake-model"), providerID: Provider.ID.make("fake") },
- })
- yield* bus.publish(SessionEvent.Tool.Input.Started, {
- sessionID,
- assistantMessageID,
- callID: "call-hosted-interrupted",
- name: "web_search",
- })
- yield* bus.publish(SessionEvent.Tool.Input.Ended, {
- sessionID,
- assistantMessageID,
- callID: "call-hosted-interrupted",
- text: '{"query":"stale"}',
- })
- yield* bus.publish(SessionEvent.Tool.Called, {
- sessionID,
- assistantMessageID,
- callID: "call-hosted-interrupted",
- input: { query: "stale" },
- executed: true,
- state: { itemId: "call-hosted-interrupted" },
- })
- requests.length = 0
- yield* TestLLM.push([])
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(1)
- expect(messageRoles(requests[0])).toEqual(["user", "assistant"])
- expect(requests[0]?.messages[1]?.content).toMatchObject([
- {
- type: "tool-call",
- id: "call-hosted-interrupted",
- providerExecuted: true,
- providerMetadata: { openai: { itemId: "call-hosted-interrupted" } },
- },
- { type: "tool-result", id: "call-hosted-interrupted", providerExecuted: true, result: { type: "error" } },
- ])
- }),
- )
- it.effect("durably fails pending tool input left by a prior process before continuing", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- yield* admit(session, "Recover interrupted tool input")
- yield* SessionPending.promote((yield* Database.Service).db, bus, sessionID, "steer")
- const assistantMessageID = SessionMessage.ID.create()
- yield* bus.publish(SessionEvent.Step.Started, {
- sessionID,
- assistantMessageID,
- agent: Agent.ID.make("build"),
- model: { id: ID.make("fake-model"), providerID: Provider.ID.make("fake") },
- })
- yield* bus.publish(SessionEvent.Tool.Input.Started, {
- sessionID,
- assistantMessageID,
- callID: "call-pending-interrupted",
- name: "echo",
- })
- requests.length = 0
- yield* TestLLM.push([])
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(1)
- expect(messageRoles(requests[0])).toEqual(["user", "assistant", "tool"])
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Recover interrupted tool input" },
- { type: "assistant", content: [{ type: "tool", id: "call-pending-interrupted", state: { status: "error" } }] },
- ])
- }),
- )
- it.effect("promotes the first queued input when woken while idle", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* session.prompt({
- sessionID,
- text: "Wait in queue",
- delivery: "queue",
- resume: false,
- })
- const stream = yield* TestLLM.gate
- yield* (yield* SessionExecution.Service).wake(sessionID)
- yield* stream.started
- yield* stream.release
- expect(requests).toHaveLength(1)
- expect(userTexts(requests[0])).toEqual(["Wait in queue"])
- }),
- )
- it.effect("retries inbox input after prompt projection rolls back", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- const defect = new Error("fail after prompt promotion")
- let fail = true
- yield* bus.project(SessionEvent.InputPromoted, () => (fail ? Effect.die(defect) : Effect.void))
- yield* admit(session, "Recover promoted input")
- expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect)
- fail = false
- requests.length = 0
- yield* TestLLM.push(TestLLM.stop())
- const stream = yield* TestLLM.gate
- yield* (yield* SessionExecution.Service).wake(sessionID)
- yield* stream.started
- yield* stream.release
- expect(userTexts(requests[0])).toEqual(["Recover promoted input"])
- }),
- )
- it.effect("does not strand a committed promotion when a post-commit listener defects", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const bus = yield* Bus.Service
- yield* bus.listen((event) =>
- event.type === SessionEvent.InputPromoted.type
- ? Effect.die("fail after prompt promotion commits")
- : Effect.void,
- )
- yield* runPrompt(session, "Run committed promotion")
- expect(requests).toHaveLength(1)
- expect(userTexts(requests[0])).toEqual(["Run committed promotion"])
- }),
- )
- it.effect("adds session correlation headers to model requests", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* runPrompt(session, "Run correlated request")
- expect(requests[0]?.http?.headers).toEqual({
- "x-session-affinity": sessionID,
- "X-Session-Id": sessionID,
- "User-Agent": App.useragent(App.make()),
- "x-opencode-project": Project.ID.global,
- "x-opencode-session": sessionID,
- "x-opencode-client": "opencode",
- })
- }),
- )
- it.effect("adds the parent session header to child model requests", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const parentID = Session.ID.make("ses_runner_parent")
- const { db } = yield* Database.Service
- yield* db
- .update(SessionTable)
- .set({ parent_id: parentID })
- .where(eq(SessionTable.id, sessionID))
- .run()
- .pipe(Effect.orDie)
- yield* runPrompt(session, "Run child request")
- expect(requests[0]?.http?.headers?.["x-parent-session-id"]).toBe(parentID)
- }),
- )
- it.effect("runs different sessions concurrently", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* insertSession(otherSessionID)
- yield* admit(session, "Run first")
- yield* session.prompt({
- sessionID: otherSessionID,
- text: "Run second",
- resume: false,
- })
- const stream = yield* TestLLM.gate
- const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- const second = yield* session.resume(otherSessionID).pipe(Effect.forkChild)
- yield* stream.started
- expect(requests).toHaveLength(2)
- expect(requests.map((request) => request.providerOptions?.openai?.promptCacheKey)).toEqual([
- sessionID,
- otherSessionID,
- ])
- yield* stream.release
- yield* Fiber.join(first)
- yield* Fiber.join(second)
- }),
- )
- it.effect("bounds 64-character session prompt cache keys", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const longSessionID = Session.ID.make(`ses_${"a".repeat(64)}`)
- const otherLongSessionID = Session.ID.make(`ses_${"b".repeat(64)}`)
- yield* insertSession(longSessionID)
- yield* insertSession(otherLongSessionID)
- yield* session.prompt({
- sessionID: longSessionID,
- text: "Run long session",
- resume: false,
- })
- yield* session.prompt({
- sessionID: otherLongSessionID,
- text: "Run other long session",
- resume: false,
- })
- yield* session.resume(longSessionID)
- yield* session.resume(otherLongSessionID)
- const keys = requests.map((request) => request.providerOptions?.openai?.promptCacheKey)
- expect(keys).toEqual([longSessionID.slice(4), otherLongSessionID.slice(4)])
- expect(keys.every((key) => typeof key === "string" && key.length === 64)).toBe(true)
- expect(keys[0]).not.toBe(keys[1])
- }),
- )
- it.effect("fans out one failed run and allows a later retry", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Retry after failure")
- yield* TestLLM.push(Stream.fail(invalidRequest()))
- const stream = yield* TestLLM.gate
- const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- const second = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* Effect.yieldNow
- expect(requests).toHaveLength(1)
- yield* stream.release
- const [firstExit, secondExit] = yield* Effect.all([Fiber.await(first), Fiber.await(second)])
- expect(secondExit).toEqual(firstExit)
- yield* TestLLM.push([])
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(2)
- }),
- )
- it.effect("durably settles local tool failures before continuing", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Call missing")
- yield* TestLLM.push(TestLLM.tool("call-missing", "missing", {}), TestLLM.text("Recovered", "text-after-error"))
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(2)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Call missing" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-missing",
- state: {
- status: "error",
- error: { type: "tool.execution", message: "Unknown tool: missing" },
- },
- },
- ],
- },
- { type: "assistant", finish: "stop", content: [{ type: "text", text: "Recovered" }] },
- ])
- }),
- )
- it.effect("returns unexpected local tool defects to the model and continues", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Call defect")
- yield* TestLLM.push(TestLLM.tool("call-defect", "defect", {}), TestLLM.text("Recovered", "text-after-defect"))
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(2)
- expect(messageRoles(requests[1])).toEqual(["user", "assistant", "tool"])
- const context = yield* session.context(sessionID)
- expect(context).toMatchObject([
- { type: "user", text: "Call defect" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-defect",
- state: {
- status: "error",
- error: { type: "unknown", message: "unexpected tool defect" },
- },
- },
- ],
- },
- { type: "assistant", finish: "stop", content: [{ type: "text", text: "Recovered" }] },
- ])
- const assistant = requireAssistant(context)
- expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.failed.2",
- "session.step.ended.1",
- ])
- }),
- )
- it.effect("returns tool-wrapped policy blocks to the model and continues", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const registry = yield* Tool.Service
- yield* transformTools(
- registry,
- {
- blocked: {
- name: "blocked",
- description: "Fail because policy blocked execution",
- input: Schema.Struct({}),
- output: Schema.Struct({}),
- execute: () =>
- Effect.fail(new Permission.BlockedError({ rules: [], permission: "blocked", resources: ["*"] })).pipe(
- Effect.mapError(() => new Tool.Error({ message: "Permission blocked" })),
- ),
- },
- },
- { codemode: false },
- )
- yield* admit(session, "Call blocked")
- yield* TestLLM.push(TestLLM.tool("call-blocked", "blocked", {}), TestLLM.stop())
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(2)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Call blocked" },
- {
- type: "assistant",
- content: [
- { type: "tool", id: "call-blocked", state: { status: "error", error: { message: "Permission blocked" } } },
- ],
- },
- { type: "assistant", finish: "stop" },
- ])
- }),
- )
- it.effect("interrupts runner continuation when permission approval is declined", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const registry = yield* Tool.Service
- yield* transformTools(
- registry,
- {
- declined: {
- name: "declined",
- description: "Fail because the user declined approval",
- input: Schema.Struct({}),
- output: Schema.Struct({}),
- execute: () => Effect.die(new Permission.DeclinedError()),
- },
- },
- { codemode: false },
- )
- yield* admit(session, "Call declined")
- yield* TestLLM.push(TestLLM.tool("call-declined", "declined", {}))
- const exit = yield* session.resume(sessionID).pipe(Effect.exit)
- expect(exit._tag).toBe("Failure")
- if (exit._tag === "Failure") expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
- expect(requests).toHaveLength(1)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Call declined" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-declined",
- state: { status: "error", error: { type: "aborted", message: "The user declined this tool call" } },
- },
- ],
- },
- ])
- }),
- )
- it.effect("returns permission corrections to the model and continues", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const registry = yield* Tool.Service
- yield* transformTools(
- registry,
- {
- corrected: {
- name: "corrected",
- description: "Fail with user correction feedback",
- input: Schema.Struct({}),
- output: Schema.Struct({}),
- execute: () =>
- Effect.fail(new Permission.CorrectedError({ feedback: "Use another tool" })).pipe(
- Effect.mapError(() => new Tool.Error({ message: "Use another tool" })),
- ),
- },
- },
- { codemode: false },
- )
- yield* admit(session, "Call corrected")
- yield* TestLLM.push(TestLLM.tool("call-corrected", "corrected", {}), TestLLM.stop())
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(2)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Call corrected" },
- {
- type: "assistant",
- content: [
- { type: "tool", id: "call-corrected", state: { status: "error", error: { message: "Use another tool" } } },
- ],
- },
- { type: "assistant", finish: "stop" },
- ])
- }),
- )
- it.effect("returns configured permission denials to the model and continues", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const registry = yield* Tool.Service
- yield* transformTools(registry, { permissionfail: permissionFail }, { codemode: false })
- yield* admit(session, "Reject permission")
- yield* TestLLM.push(TestLLM.tool("call-permission", "permissionfail", {}), [
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }),
- ])
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(2)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-permission",
- state: {
- status: "error",
- error: {
- type: "permission.rejected",
- message: "Permission denied: edit",
- },
- },
- },
- ],
- },
- { type: "assistant", finish: "stop" },
- ])
- expect(yield* recordedEventTypes(sessionID)).not.toContain("session.step.failed.1")
- }),
- )
- it.effect("interrupts runner continuation when a question is cancelled", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const registry = yield* Tool.Service
- yield* transformTools(
- registry,
- {
- question: {
- name: "question",
- description: "Ask the user",
- input: Schema.Struct({}),
- output: Schema.Struct({}),
- execute: () => Effect.die(new QuestionTool.CancelledError()),
- },
- },
- { codemode: false },
- )
- yield* admit(session, "Ask then stop")
- yield* TestLLM.push(TestLLM.tool("call-question", "question", {}), [])
- const run = yield* session.resume(sessionID).pipe(Effect.exit, Effect.forkChild)
- const exit = yield* Fiber.join(run)
- expect(exit._tag).toBe("Failure")
- if (exit._tag === "Failure") expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
- expect(requests).toHaveLength(1)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Ask then stop" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-question",
- state: { status: "error", error: { type: "aborted", message: "The user dismissed this question" } },
- },
- ],
- },
- ])
- }),
- )
- it.effect("awaits started local tools before surfacing provider stream failure", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Settle before failing")
- const failure = providerUnavailable()
- const tools = yield* blockTools()
- yield* TestLLM.push(
- TestLLM.failAfter(
- failure,
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolCall({ id: "call-before-failure", name: "echo", input: { text: "settle" } }),
- ),
- )
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* tools.started
- yield* tools.release
- expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure)
- const context = yield* session.context(sessionID)
- expect(context).toMatchObject([
- { type: "user", text: "Settle before failing" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-before-failure",
- state: { status: "completed", content: [{ type: "text", text: "settle" }] },
- },
- ],
- },
- ])
- const assistant = requireAssistant(context)
- expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.success.2",
- "session.step.failed.1",
- ])
- }),
- )
- it.effect("durably fails blocked local tools when a step is interrupted", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Interrupt blocked tool")
- const tools = yield* blockTools()
- yield* TestLLM.push(
- TestLLM.hangAfter(
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolCall({ id: "call-before-interrupt", name: "echo", input: { text: "blocked" } }),
- ),
- )
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* tools.started
- yield* session.interrupt(sessionID)
- expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
- yield* session.interrupt(sessionID)
- const context = yield* session.context(sessionID)
- expect(context).toMatchObject([
- { type: "user", text: "Interrupt blocked tool" },
- {
- type: "assistant",
- content: [
- {
- type: "tool",
- id: "call-before-interrupt",
- state: { status: "error", error: { type: "aborted", message: "Tool execution interrupted" } },
- },
- ],
- },
- ])
- const assistant = requireAssistant(context)
- expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.failed.2",
- "session.step.failed.1",
- ])
- yield* replaySessionProjection(sessionID)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Interrupt blocked tool" },
- { type: "assistant", content: [{ type: "tool", id: "call-before-interrupt", state: { status: "error" } }] },
- ])
- requests.length = 0
- yield* TestLLM.push([])
- yield* session.resume(sessionID)
- expect(messageRoles(requests[0])).toEqual(["user", "assistant", "tool"])
- }),
- )
- it.effect("interrupts a blocked step without local tool execution", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Interrupt provider")
- const stream = yield* TestLLM.gate
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* session.interrupt(sessionID)
- const exit = yield* Fiber.await(run)
- expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBeTrue()
- expect(requests).toHaveLength(1)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Interrupt provider" },
- { type: "assistant", finish: "error", error: { type: "aborted", message: "Step interrupted" } },
- ])
- expect(yield* recordedEventTypes(sessionID)).toContain("session.step.failed.1")
- yield* session.interrupt(sessionID)
- }),
- )
- it.effect("durably fails blocked local tools when interrupted while awaiting settlement", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Interrupt tool settlement")
- const tools = yield* blockTools()
- yield* TestLLM.push(TestLLM.tool("call-await-interrupt", "echo", { text: "blocked" }))
- const runner = yield* SessionRunner.Service
- const run = yield* runner.drain({ sessionID, force: true }).pipe(Effect.forkChild)
- yield* tools.started
- yield* Fiber.interrupt(run)
- expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Interrupt tool settlement" },
- {
- type: "assistant",
- finish: "error",
- error: { type: "aborted", message: "Step interrupted" },
- content: [
- {
- type: "tool",
- id: "call-await-interrupt",
- state: { status: "error", error: { type: "aborted", message: "Tool execution interrupted" } },
- },
- ],
- },
- ])
- const eventTypes = yield* recordedEventTypes(sessionID)
- expect(eventTypes).toContain("session.step.failed.1")
- expect(eventTypes).not.toContain("session.step.ended.1")
- }),
- )
- it.effect("forces a text response on an agent's configured final step", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const agents = yield* Agent.Service
- yield* agents.transform((editor) =>
- editor.update(Agent.ID.make("build"), (agent) => {
- agent.steps = 2
- }),
- )
- yield* admit(session, "Finish at the limit")
- yield* TestLLM.push(
- TestLLM.tool("call-terminal", "echo", { text: "done" }),
- TestLLM.tool("call-forbidden", "echo", { text: "forbidden" }),
- )
- yield* session.resume(sessionID)
- expect(requests).toHaveLength(2)
- expect(requests[0]?.toolChoice).toBeUndefined()
- expect(requests[1]?.toolChoice).toMatchObject({ type: "none" })
- // Protocols with native "none" keep these definitions for prompt caching.
- expect(requests[1]?.tools.map((tool) => tool.name)).toContain("echo")
- expect(requests[1]?.messages.at(-1)).toMatchObject({
- role: "assistant",
- content: [{ type: "text", text: expect.stringContaining("MAXIMUM STEPS REACHED") }],
- })
- expect(executions).toEqual(["done"])
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Finish at the limit" },
- { type: "assistant", content: [{ type: "tool", id: "call-terminal", state: { status: "completed" } }] },
- { type: "assistant", content: [{ type: "tool", id: "call-forbidden", state: { status: "error" } }] },
- ])
- }),
- )
- it.effect("resets the configured step allowance when steering input promotes", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const agents = yield* Agent.Service
- yield* agents.transform((editor) =>
- editor.update(Agent.ID.make("build"), (agent) => {
- agent.steps = 2
- }),
- )
- yield* admit(session, "Start work")
- yield* TestLLM.push(
- TestLLM.tool("call-before-steer", "echo", { text: "before" }),
- TestLLM.tool("call-after-steer", "echo", { text: "after" }),
- TestLLM.stop(),
- )
- const stream = yield* TestLLM.gate
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* stream.started
- yield* session.prompt({ sessionID, text: "Change direction" })
- yield* stream.release
- yield* Fiber.join(run)
- expect(requests).toHaveLength(3)
- expect(requests[1]?.toolChoice).toBeUndefined()
- expect(requests[1]?.tools).not.toEqual([])
- expect(requests[2]?.toolChoice).toMatchObject({ type: "none" })
- expect(executions).toEqual(["before", "after"])
- }),
- )
- it.effect("projects provider errors as terminal assistant step failures", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push([
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.providerError({ message: "Provider unavailable" }),
- ])
- expect((yield* runPrompt(session, "Fail durably").pipe(Effect.flip)).message).toBe("Provider unavailable")
- expect(requests).toHaveLength(1)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Fail durably" },
- { type: "assistant", finish: "error", error: { type: "provider.unknown", message: "Provider unavailable" } },
- ])
- }),
- )
- it.effect("projects provider errors emitted before assistant step start", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push([LLMEvent.providerError({ message: "Provider unavailable" })])
- expect((yield* runPrompt(session, "Fail before step").pipe(Effect.flip)).message).toBe("Provider unavailable")
- expect(requests).toHaveLength(1)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Fail before step" },
- { type: "assistant", finish: "error", error: { type: "provider.unknown", message: "Provider unavailable" } },
- ])
- }),
- )
- it.effect("projects content-filter finishes as visible terminal failures", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(
- TestLLM.complete(
- {
- reason: { normalized: "content-filter" },
- usage: { nonCachedInputTokens: 8, outputTokens: 3, reasoningTokens: 1 },
- },
- LLMEvent.textStart({ id: "partial" }),
- LLMEvent.textDelta({ id: "partial", text: "Partial" }),
- ),
- )
- expect((yield* runPrompt(session, "Blocked response").pipe(Effect.flip)).message).toBe(
- "Provider blocked the response",
- )
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user" },
- {
- type: "assistant",
- finish: "error",
- error: { type: "provider.content-filter" },
- cost: 0,
- tokens: { input: 8, output: 2, reasoning: 1, cache: { read: 0, write: 0 } },
- content: [{ type: "text", text: "Partial" }],
- },
- ])
- expect(yield* session.get(sessionID)).toMatchObject({
- cost: 0,
- tokens: { input: 8, output: 2, reasoning: 1, cache: { read: 0, write: 0 } },
- })
- expect(yield* recordedEventTypes(sessionID)).not.toContain("session.step.ended.1")
- }),
- )
- it.effect("settles a local tool before one content-filter step failure", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Tool before blocked response")
- const tools = yield* blockTools()
- yield* TestLLM.push(
- TestLLM.complete(
- { reason: { normalized: "content-filter" } },
- LLMEvent.toolCall({ id: "call-before-content-filter", name: "echo", input: { text: "settled" } }),
- ),
- )
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* tools.started
- yield* tools.release
- expect((yield* Fiber.join(run).pipe(Effect.flip)).message).toBe("Provider blocked the response")
- const assistant = requireAssistant(yield* session.context(sessionID))
- const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id)
- expect(bus.map((event) => event.type)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.success.2",
- "session.step.failed.1",
- ])
- expect(
- bus.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
- ).toHaveLength(1)
- }),
- )
- it.effect("does not recover context overflow after durable assistant output", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push([
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.textStart({ id: "text-partial" }),
- LLMEvent.textDelta({ id: "text-partial", text: "Partial" }),
- LLMEvent.textEnd({ id: "text-partial" }),
- LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
- ])
- expect((yield* runPrompt(session, "Fail after output").pipe(Effect.flip)).message).toBe("prompt too long")
- expect(requests).toHaveLength(1)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Fail after output" },
- {
- type: "assistant",
- finish: "error",
- error: { message: "prompt too long" },
- content: [{ type: "text", text: "Partial" }],
- },
- ])
- }),
- )
- it.effect("projects raw provider stream failures as terminal assistant step failures", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const failure = invalidRequest()
- yield* TestLLM.push(Stream.fail(failure))
- expect(yield* runPrompt(session, "Fail raw stream durably").pipe(Effect.flip)).toBe(failure)
- yield* replaySessionProjection(sessionID)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Fail raw stream durably" },
- { type: "assistant", finish: "error", error: { type: "provider.invalid-request", message: "Invalid request" } },
- ])
- }),
- )
- it.effect("retries eligible pre-output failures after exponential backoff", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Retry transport")
- yield* TestLLM.push(Stream.fail(providerUnavailable()))
- yield* TestLLM.push(TestLLM.text("Recovered", "retry-success"))
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* TestLLM.wait(1)
- yield* TestClock.adjust("1999 millis")
- expect(requests).toHaveLength(1)
- yield* TestClock.adjust("1 millis")
- yield* Fiber.join(run)
- expect(requests).toHaveLength(2)
- const eventTypes = yield* recordedEventTypes(sessionID)
- expect(eventTypes).toContain("session.retry.scheduled.1")
- expect(eventTypes.filter((type) => type === "session.step.started.1")).toHaveLength(2)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user" },
- { type: "assistant", finish: "stop", content: [{ type: "text", text: "Recovered" }] },
- ])
- yield* replaySessionProjection(sessionID)
- expect((yield* session.context(sessionID)).filter((message) => message.type === "assistant")).toHaveLength(1)
- }),
- )
- it.effect("uses a larger provider retry-after delay", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Retry rate limit")
- yield* TestLLM.push(Stream.fail(rateLimited(5_000)))
- yield* TestLLM.push(TestLLM.text("Recovered", "retry-after-success"))
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* TestLLM.wait(1)
- yield* TestClock.adjust("4999 millis")
- expect(requests).toHaveLength(1)
- yield* TestClock.adjust("1 millis")
- yield* Fiber.join(run)
- expect(requests).toHaveLength(2)
- }),
- )
- it.effect("does not retry eligible failures after observable output", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const failure = rateLimited()
- yield* TestLLM.push(
- TestLLM.failAfter(
- failure,
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.textStart({ id: "partial-rate-limit" }),
- LLMEvent.textDelta({ id: "partial-rate-limit", text: "Partial" }),
- ),
- )
- expect(yield* runPrompt(session, "Do not replay partial output").pipe(Effect.flip)).toBe(failure)
- expect(requests).toHaveLength(1)
- expect(yield* recordedEventTypes(sessionID)).not.toContain("session.retry.scheduled.1")
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user" },
- {
- type: "assistant",
- finish: "error",
- error: { type: "provider.rate-limit" },
- content: [{ type: "text", text: "Partial" }],
- },
- ])
- }),
- )
- it.effect("stops after five total retry attempts", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Exhaust retries")
- const failure = providerUnavailable()
- yield* TestLLM.always(Stream.fail(failure))
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* TestLLM.wait(1)
- for (const [index, delay] of [2_000, 4_000, 8_000, 16_000].entries()) {
- yield* TestClock.adjust(delay)
- yield* TestLLM.wait(index + 2)
- }
- expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure)
- expect(requests).toHaveLength(5)
- const database = (yield* Database.Service).db
- const retries = yield* database
- .select({ data: EventTable.data })
- .from(EventTable)
- .where(eq(EventTable.type, "session.retry.scheduled.1"))
- .orderBy(asc(EventTable.seq))
- .all()
- .pipe(Effect.orDie)
- expect(retries.map((event) => event.data)).toMatchObject([
- { attempt: 2, at: 2_000 },
- { attempt: 3, at: 6_000 },
- { attempt: 4, at: 14_000 },
- { attempt: 5, at: 30_000 },
- ])
- expect((yield* recordedEventTypes(sessionID)).filter((type) => type === "session.step.started.1")).toHaveLength(5)
- const assistant = requireAssistant(yield* session.context(sessionID))
- expect(yield* recordedStepSettlementEvents(sessionID, assistant.id)).toMatchObject([
- { type: "session.step.started.1" },
- { type: "session.step.started.1" },
- { type: "session.step.started.1" },
- { type: "session.step.started.1" },
- { type: "session.step.started.1" },
- { type: "session.step.failed.1" },
- ])
- }),
- )
- it.effect("retries a model call without consuming the logical agent step", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const agents = yield* Agent.Service
- yield* agents.transform((editor) =>
- editor.update(Agent.ID.make("build"), (agent) => {
- agent.steps = 2
- }),
- )
- yield* admit(session, "Retry without consuming a step")
- const failure = providerUnavailable()
- yield* TestLLM.push(Stream.fail(failure))
- yield* TestLLM.push(TestLLM.tool("call-after-retry", "echo", { text: "recovered" }), TestLLM.stop())
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* TestLLM.wait(1)
- yield* TestClock.adjust("2 seconds")
- yield* Fiber.join(run)
- expect(requests).toHaveLength(3)
- expect(requests[0]?.toolChoice).toBeUndefined()
- expect(requests[0]?.tools.map((tool) => tool.name)).toContain("echo")
- expect(requests[1]?.toolChoice).toBeUndefined()
- expect(requests[1]?.tools.map((tool) => tool.name)).toContain("echo")
- expect(requests[1]?.messages.at(-1)).not.toMatchObject({
- role: "assistant",
- content: [{ type: "text", text: expect.stringContaining("MAXIMUM STEPS REACHED") }],
- })
- expect(requests[2]?.toolChoice).toMatchObject({ type: "none" })
- // The final step keeps tool definitions to preserve provider prompt caching.
- expect(requests[2]?.tools.map((tool) => tool.name)).toContain("echo")
- expect(requests[2]?.messages.at(-1)).toMatchObject({
- role: "assistant",
- content: [{ type: "text", text: expect.stringContaining("MAXIMUM STEPS REACHED") }],
- })
- expect(executions).toEqual(["recovered"])
- const eventTypes = yield* recordedEventTypes(sessionID)
- expect(eventTypes.filter((type) => type === "session.step.started.1")).toHaveLength(3)
- expect(eventTypes.filter((type) => type === "session.retry.scheduled.1")).toHaveLength(1)
- expect((yield* session.context(sessionID)).filter((message) => message.type === "assistant")).toHaveLength(2)
- }),
- )
- it.effect("does not retry non-eligible provider failures", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const failure = invalidRequest()
- yield* TestLLM.push(Stream.fail(failure))
- expect(yield* runPrompt(session, "Do not retry").pipe(Effect.flip)).toBe(failure)
- expect(requests).toHaveLength(1)
- expect(yield* recordedEventTypes(sessionID)).not.toContain("session.retry.scheduled.1")
- }),
- )
- it.effect("settles malformed streamed tool input before the provider failure", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const failure = new LLMError({
- module: "test",
- method: "stream",
- reason: new InvalidProviderOutputReason({ message: "Invalid JSON input for tool call echo" }),
- })
- yield* TestLLM.push(
- TestLLM.failAfter(
- failure,
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolInputStart({ id: "call-malformed", name: "echo" }),
- LLMEvent.toolInputDelta({ id: "call-malformed", name: "echo", text: '{"text":"partial' }),
- ),
- )
- expect(yield* runPrompt(session, "Call a malformed tool").pipe(Effect.flip)).toBe(failure)
- const assistant = requireAssistant(yield* session.context(sessionID))
- yield* TestLLM.push(TestLLM.stop())
- yield* runPrompt(session, "Continue")
- expect(yield* recordedStepSettlementEvents(sessionID, assistant.id)).toMatchObject([
- { type: "session.step.started.1" },
- {
- type: "session.tool.failed.2",
- data: {
- callID: "call-malformed",
- error: { type: "provider.invalid-output", message: "Invalid JSON input for tool call echo" },
- },
- },
- {
- type: "session.step.failed.1",
- data: { error: { type: "provider.invalid-output", message: "Invalid JSON input for tool call echo" } },
- },
- ])
- }),
- )
- it.effect("continues after malformed local tool input without exposing raw arguments", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const marker = "raw-malformed-marker"
- const raw = `{"text":"${marker}`
- yield* TestLLM.push(
- TestLLM.toolCalls(
- LLMEvent.toolInputStart({ id: "call-malformed", name: "echo" }),
- LLMEvent.toolInputDelta({ id: "call-malformed", name: "echo", text: raw }),
- LLMEvent.toolInputEnd({ id: "call-malformed", name: "echo" }),
- LLMEvent.toolInputError({
- id: "call-malformed",
- name: "echo",
- raw,
- }),
- ),
- TestLLM.stop(),
- )
- yield* runPrompt(session, "Recover malformed tool input")
- expect(requests).toHaveLength(2)
- expect(executions).toEqual([])
- expect(JSON.stringify(requests[1])).not.toContain(marker)
- expect(requests[1]?.messages).toEqual(
- expect.arrayContaining([
- expect.objectContaining({
- role: "assistant",
- content: expect.arrayContaining([
- expect.objectContaining({ type: "tool-call", id: "call-malformed", name: "echo", input: {} }),
- ]),
- }),
- expect.objectContaining({
- role: "tool",
- content: expect.arrayContaining([
- expect.objectContaining({
- type: "tool-result",
- id: "call-malformed",
- result: expect.objectContaining({
- type: "error",
- value: expect.objectContaining({
- error: expect.objectContaining({
- message: "Tool call arguments were malformed JSON and were not executed. Retry with valid JSON.",
- }),
- }),
- }),
- }),
- ]),
- }),
- ]),
- )
- const context = yield* session.context(sessionID)
- const failed = context.find(
- (message): message is SessionMessage.Assistant =>
- message.type === "assistant" && message.content.some((item) => item.type === "tool"),
- )
- expect(failed).toMatchObject({
- content: [
- {
- type: "tool",
- id: "call-malformed",
- executed: false,
- state: {
- status: "error",
- input: {},
- error: {
- type: "tool.input-json",
- message: "Tool call arguments were malformed JSON and were not executed. Retry with valid JSON.",
- },
- },
- },
- ],
- })
- if (!failed) throw new Error("Malformed tool assistant missing")
- expect(failed.error).toBeUndefined()
- expect(yield* recordedStepSettlementTypes(sessionID, failed.id)).toEqual([
- "session.step.started.1",
- "session.tool.failed.2",
- "session.step.ended.1",
- ])
- const database = (yield* Database.Service).db
- const durable = yield* database
- .select({ type: EventTable.type, data: EventTable.data })
- .from(EventTable)
- .where(eq(EventTable.aggregate_id, sessionID))
- .all()
- .pipe(Effect.orDie)
- expect(durable.find((event) => event.type === "session.tool.input.ended.1")?.data).toMatchObject({
- callID: "call-malformed",
- text: raw,
- })
- }),
- )
- it.effect("settles a valid sibling before recovering malformed tool input", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Run parallel tools")
- const tools = yield* blockTools()
- yield* TestLLM.push(
- TestLLM.toolCalls(
- LLMEvent.toolCall({ id: "call-valid", name: "echo", input: { text: "valid" } }),
- LLMEvent.toolInputError({
- id: "call-malformed",
- name: "echo",
- raw: '{"text":"partial',
- }),
- ),
- TestLLM.stop(),
- )
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* tools.started
- expect(requests).toHaveLength(1)
- yield* tools.release
- yield* Fiber.join(run)
- expect(requests).toHaveLength(2)
- expect(executions).toEqual(["valid"])
- const request = requests[1]
- if (!request) throw new Error("Malformed recovery request missing")
- expect(request.messages.flatMap((message) => (message.role === "tool" ? message.content : []))).toEqual(
- expect.arrayContaining([
- expect.objectContaining({ id: "call-valid", type: "tool-result" }),
- expect.objectContaining({ id: "call-malformed", type: "tool-result" }),
- ]),
- )
- }),
- )
- it.effect("does not recover malformed input after sibling execution is interrupted", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Interrupt malformed recovery")
- const tools = yield* blockTools()
- yield* TestLLM.push(
- TestLLM.toolCalls(
- LLMEvent.toolCall({ id: "call-valid", name: "echo", input: { text: "blocked" } }),
- LLMEvent.toolInputError({
- id: "call-malformed",
- name: "echo",
- raw: '{"text":"partial',
- }),
- ),
- )
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* tools.started
- while (
- !(yield* session.context(sessionID)).some(
- (message) =>
- message.type === "assistant" &&
- message.content.some((item) => item.type === "tool" && item.id === "call-malformed"),
- )
- )
- yield* Effect.yieldNow
- yield* session.interrupt(sessionID)
- expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
- expect(requests).toHaveLength(1)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Interrupt malformed recovery" },
- {
- type: "assistant",
- error: { type: "aborted", message: "Step interrupted" },
- content: [
- { type: "tool", id: "call-valid", state: { status: "error", error: { type: "aborted" } } },
- { type: "tool", id: "call-malformed", state: { status: "error" } },
- ],
- },
- ])
- }),
- )
- it.effect("records malformed provider-executed input as executed", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const failure = new LLMError({
- module: "test",
- method: "stream",
- reason: new InvalidProviderOutputReason({ message: "Invalid hosted tool input" }),
- })
- yield* TestLLM.push(
- TestLLM.failAfter(
- failure,
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolInputStart({ id: "call-hosted", name: "web_search", providerExecuted: true }),
- LLMEvent.toolInputDelta({ id: "call-hosted", name: "web_search", text: '{"query":"partial' }),
- ),
- )
- expect(yield* runPrompt(session, "Fail malformed hosted input").pipe(Effect.flip)).toBe(failure)
- expect(requireAssistant(yield* session.context(sessionID))).toMatchObject({
- error: { type: "provider.invalid-output", message: "Invalid hosted tool input" },
- content: [
- {
- type: "tool",
- id: "call-hosted",
- executed: true,
- state: { status: "error", error: { type: "provider.invalid-output" } },
- },
- ],
- })
- }),
- )
- it.effect("records a provider failure after malformed input", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const failure = new LLMError({
- module: "test",
- method: "stream",
- reason: new InvalidProviderOutputReason({ message: "Provider failed after malformed input" }),
- })
- yield* TestLLM.push(
- TestLLM.failAfter(
- failure,
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolInputError({
- id: "call-malformed",
- name: "echo",
- raw: '{"text":"partial',
- }),
- ),
- )
- expect(yield* runPrompt(session, "Fail after malformed input").pipe(Effect.flip)).toBe(failure)
- expect(requireAssistant(yield* session.context(sessionID))).toMatchObject({
- error: { type: "provider.invalid-output", message: "Provider failed after malformed input" },
- content: [
- {
- type: "tool",
- id: "call-malformed",
- executed: false,
- state: { status: "error", error: { type: "tool.input-json" } },
- },
- ],
- })
- expect(requests).toHaveLength(1)
- }),
- )
- it.effect("continues after repeated malformed tool input", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const malformed = (id: string) =>
- TestLLM.toolCalls(
- LLMEvent.toolInputError({
- id,
- name: "echo",
- raw: '{"text":"partial',
- }),
- )
- yield* TestLLM.push(
- malformed("call-first"),
- TestLLM.tool("call-valid-between", "echo", { text: "valid" }),
- malformed("call-second"),
- TestLLM.stop(),
- )
- yield* runPrompt(session, "Keep producing malformed tools")
- expect(requests).toHaveLength(4)
- expect(executions).toEqual(["valid"])
- expect((yield* recordedEventTypes(sessionID)).filter((type) => type === "session.step.failed.1")).toHaveLength(0)
- }),
- )
- it.effect("does not continue malformed tool input past the agent step limit", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const agents = yield* Agent.Service
- yield* agents.transform((editor) =>
- editor.update(Agent.ID.make("build"), (agent) => {
- agent.steps = 2
- }),
- )
- const malformed = (id: string) =>
- TestLLM.toolCalls(
- LLMEvent.toolInputError({
- id,
- name: "echo",
- raw: '{"text":"partial',
- }),
- )
- yield* TestLLM.push(malformed("call-first"), malformed("call-at-limit"))
- yield* runPrompt(session, "Stop malformed tools at the step limit")
- expect(requests).toHaveLength(2)
- expect(requests[0]?.toolChoice).toBeUndefined()
- expect(requests[1]?.toolChoice).toMatchObject({ type: "none" })
- expect((yield* recordedEventTypes(sessionID)).filter((type) => type === "session.tool.failed.2")).toHaveLength(2)
- }),
- )
- it.effect("does not continue automatically after a provider error follows a local tool call", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Do not continue failed provider")
- const tools = yield* blockTools()
- yield* TestLLM.push([
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolCall({ id: "call-before-provider-error", name: "echo", input: { text: "settled" } }),
- LLMEvent.providerError({ message: "Provider unavailable" }),
- ])
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* tools.started
- yield* tools.release
- expect((yield* Fiber.join(run).pipe(Effect.flip)).message).toBe("Provider unavailable")
- expect(requests).toHaveLength(1)
- expect(executions).toEqual(["settled"])
- const context = yield* session.context(sessionID)
- const assistant = requireAssistant(context)
- expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.success.2",
- "session.step.failed.1",
- ])
- }),
- )
- it.effect("durably fails a hosted tool when its provider errors before returning a result", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push([
- LLMEvent.stepStart({ index: 0 }),
- hostedCall("call-hosted-provider-error", "effect"),
- LLMEvent.providerError({ message: "Provider unavailable" }),
- ])
- expect((yield* runPrompt(session, "Fail hosted tool durably").pipe(Effect.flip)).message).toBe(
- "Provider unavailable",
- )
- expect(requests).toHaveLength(1)
- const context = yield* session.context(sessionID)
- expect(context).toMatchObject([
- { type: "user", text: "Fail hosted tool durably" },
- {
- type: "assistant",
- content: [{ type: "tool", id: "call-hosted-provider-error", state: { status: "error" } }],
- },
- ])
- const assistant = requireAssistant(context)
- expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.failed.2",
- "session.step.failed.1",
- ])
- }),
- )
- it.effect("preserves a tool defect before provider failure settlement", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push([
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolCall({ id: "call-defect-provider-error", name: "defect", input: {} }),
- LLMEvent.providerError({ message: "Provider unavailable" }),
- ])
- expect((yield* runPrompt(session, "Defect while provider fails").pipe(Effect.flip)).message).toBe(
- "Provider unavailable",
- )
- const context = yield* session.context(sessionID)
- const assistant = requireAssistant(context)
- const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id)
- expect(bus.map((event) => event.type)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.failed.2",
- "session.step.failed.1",
- ])
- expect(bus[2]?.data.error).toMatchObject({ type: "unknown", message: "unexpected tool defect" })
- }),
- )
- it.effect("preserves the provider failure when tool output persistence also fails", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Storage fails while provider fails")
- yield* TestLLM.push([
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolCall({ id: "call-store-provider-error", name: "storefail", input: {} }),
- LLMEvent.providerError({ message: "Provider unavailable" }),
- ])
- expect(yield* session.resume(sessionID).pipe(Effect.exit)).toMatchObject({
- _tag: "Failure",
- })
- expect(requireAssistant(yield* session.context(sessionID))).toMatchObject({
- error: { type: "provider.unknown", message: "Provider unavailable" },
- })
- }),
- )
- it.effect("durably fails a hosted tool left unresolved at normal provider EOF", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push([LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-eof", "effect")])
- expect((yield* runPrompt(session, "Fail hosted tool at EOF").pipe(Effect.flip)).message).toBe(
- "Provider did not return a tool result",
- )
- const assistant = requireAssistant(yield* session.context(sessionID))
- const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id)
- expect(bus.map((event) => event.type)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.failed.2",
- "session.step.failed.1",
- ])
- expect(
- bus.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
- ).toHaveLength(1)
- yield* replaySessionProjection(sessionID)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Fail hosted tool at EOF" },
- {
- type: "assistant",
- finish: "error",
- error: { type: "tool.result-missing" },
- content: [{ type: "tool", id: "call-hosted-eof", state: { status: "error" } }],
- },
- ])
- }),
- )
- it.effect("fails an unresolved hosted tool before one clean step end", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(TestLLM.stop(hostedCall("call-hosted-clean-end", "effect")))
- yield* runPrompt(session, "Settle hosted tool before ending")
- const assistant = requireAssistant(yield* session.context(sessionID))
- const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id)
- expect(bus.map((event) => event.type)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.failed.2",
- "session.step.ended.1",
- ])
- expect(
- bus.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
- ).toHaveLength(1)
- }),
- )
- it.effect("settles unresolved local and hosted tools before one raw provider failure", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* admit(session, "Fail unresolved tools")
- const failure = invalidRequest()
- const providerFailed = yield* Deferred.make<void>()
- const tools = yield* blockTools()
- yield* TestLLM.push(
- Stream.concat(
- Stream.fromIterable([
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolCall({ id: "call-local-raw-failure", name: "defect", input: {} }),
- hostedCall("call-hosted-raw-failure-pair", "effect"),
- ]),
- Stream.fromEffect(Deferred.succeed(providerFailed, undefined)).pipe(
- Stream.flatMap(() => Stream.fail(failure)),
- ),
- ),
- )
- const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
- yield* Deferred.await(providerFailed)
- yield* tools.release
- expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure)
- const assistant = requireAssistant(yield* session.context(sessionID))
- const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id)
- expect(bus.map((event) => ({ type: event.type, callID: event.data.callID }))).toEqual([
- { type: "session.step.started.1", callID: undefined },
- { type: "session.tool.called.1", callID: "call-local-raw-failure" },
- { type: "session.tool.called.1", callID: "call-hosted-raw-failure-pair" },
- { type: "session.tool.failed.2", callID: "call-local-raw-failure" },
- { type: "session.tool.failed.2", callID: "call-hosted-raw-failure-pair" },
- { type: "session.step.failed.1", callID: undefined },
- ])
- expect(
- bus.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
- ).toHaveLength(1)
- }),
- )
- it.effect("durably fails a hosted tool left unresolved by a raw provider stream failure", () =>
- Effect.gen(function* () {
- const session = yield* setup
- const failure = providerUnavailable()
- yield* TestLLM.push(
- Stream.concat(
- Stream.fromIterable([LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-raw-failure", "effect")]),
- Stream.fail(failure),
- ),
- )
- expect(yield* runPrompt(session, "Fail hosted tool on raw failure").pipe(Effect.flip)).toBe(failure)
- expect(requests).toHaveLength(1)
- const assistant = requireAssistant(yield* session.context(sessionID))
- const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id)
- expect(bus.map((event) => event.type)).toEqual([
- "session.step.started.1",
- "session.tool.called.1",
- "session.tool.failed.2",
- "session.step.failed.1",
- ])
- expect(
- bus.filter((event) => event.type.startsWith("session.step.") && event.type !== "session.step.started.1"),
- ).toHaveLength(1)
- yield* replaySessionProjection(sessionID)
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Fail hosted tool on raw failure" },
- {
- type: "assistant",
- finish: "error",
- error: { type: "provider.transport", message: "Provider unavailable" },
- content: [{ type: "tool", id: "call-hosted-raw-failure", state: { status: "error" } }],
- },
- ])
- }),
- )
- it.effect("rejects a second text start before the open fragment ends", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push([
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.textStart({ id: "text-1" }),
- LLMEvent.textStart({ id: "text-2" }),
- ])
- const defect = yield* runPrompt(session, "Two blocks").pipe(Effect.catchDefect(Effect.succeed))
- expect(defect).toBeInstanceOf(Error)
- if (!(defect instanceof Error)) return
- expect(defect.message).toBe("text start before end: text-2")
- }),
- )
- it.effect("projects sequential text fragments as separate content parts", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(
- TestLLM.stop(
- LLMEvent.textStart({ id: "text-1" }),
- LLMEvent.textDelta({ id: "text-1", text: "First" }),
- LLMEvent.textEnd({ id: "text-1" }),
- LLMEvent.textStart({ id: "text-2" }),
- LLMEvent.textDelta({ id: "text-2", text: "Second" }),
- LLMEvent.textEnd({ id: "text-2" }),
- ),
- )
- yield* runPrompt(session, "Two blocks")
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Two blocks" },
- {
- type: "assistant",
- content: [
- { type: "text", text: "First" },
- { type: "text", text: "Second" },
- ],
- },
- ])
- }),
- )
- for (const kind of fragmentKinds) {
- it.effect(`broadcasts provider ${kind} deltas without storing projection rewrites`, () =>
- verifyEphemeralDeltas(kind),
- )
- it.effect(`durably closes partial ${kind} when the provider stream fails`, () => verifyPartialFlushOnFailure(kind))
- it.effect(`durably closes partial ${kind} when the provider stream is interrupted`, () =>
- verifyPartialFlushOnInterruption(kind),
- )
- }
- it.effect("rejects duplicate streamed text starts", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push([LLMEvent.textStart({ id: "text-1" }), LLMEvent.textStart({ id: "text-1" })])
- const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))
- expect(defect).toBeInstanceOf(Error)
- if (!(defect instanceof Error)) return
- expect(defect.message).toBe("Duplicate text start: text-1")
- }),
- )
- it.effect("transitions streamed raw tool input to parsed called input", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push(
- TestLLM.stop(
- LLMEvent.toolInputStart({ id: "call-parsed", name: "web_search" }),
- LLMEvent.toolInputDelta({ id: "call-parsed", name: "web_search", text: '{"query":"hello"}' }),
- LLMEvent.toolInputEnd({ id: "call-parsed", name: "web_search" }),
- hostedCall("call-parsed", "hello"),
- ),
- )
- yield* runPrompt(session, "Call provider tool")
- expect(yield* session.context(sessionID)).toMatchObject([
- { type: "user", text: "Call provider tool" },
- {
- type: "assistant",
- content: [{ type: "tool", id: "call-parsed", state: { status: "error", input: { query: "hello" } } }],
- },
- ])
- }),
- )
- it.effect("rejects malformed streamed tool input ordering", () =>
- Effect.gen(function* () {
- const session = yield* setup
- yield* TestLLM.push([LLMEvent.toolInputDelta({ id: "call-1", name: "read", text: "{}" })])
- const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))
- expect(defect).toBeInstanceOf(Error)
- if (!(defect instanceof Error)) return
- expect(defect.message).toBe("Tool input delta before start: call-1")
- }),
- )
- })
|