session-runner.test.ts 119 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363236423652366236723682369237023712372237323742375237623772378237923802381238223832384238523862387238823892390239123922393239423952396239723982399240024012402240324042405240624072408240924102411241224132414241524162417241824192420242124222423242424252426242724282429243024312432243324342435243624372438243924402441244224432444244524462447244824492450245124522453245424552456245724582459246024612462246324642465246624672468246924702471247224732474247524762477247824792480248124822483248424852486248724882489249024912492249324942495249624972498249925002501250225032504250525062507250825092510251125122513251425152516251725182519252025212522252325242525252625272528252925302531253225332534253525362537253825392540254125422543254425452546254725482549255025512552255325542555255625572558255925602561256225632564256525662567256825692570257125722573257425752576257725782579258025812582258325842585258625872588258925902591259225932594259525962597259825992600260126022603260426052606260726082609261026112612261326142615261626172618261926202621262226232624262526262627262826292630263126322633263426352636263726382639264026412642264326442645264626472648264926502651265226532654265526562657265826592660266126622663266426652666266726682669267026712672267326742675267626772678267926802681268226832684268526862687268826892690269126922693269426952696269726982699270027012702270327042705270627072708270927102711271227132714271527162717271827192720272127222723272427252726272727282729273027312732273327342735273627372738273927402741274227432744274527462747274827492750275127522753275427552756275727582759276027612762276327642765276627672768276927702771277227732774277527762777277827792780278127822783278427852786278727882789279027912792279327942795279627972798279928002801280228032804280528062807280828092810281128122813281428152816281728182819282028212822282328242825282628272828282928302831283228332834283528362837283828392840284128422843284428452846284728482849285028512852285328542855285628572858285928602861286228632864286528662867286828692870287128722873287428752876287728782879288028812882288328842885288628872888288928902891289228932894289528962897289828992900290129022903290429052906290729082909291029112912291329142915291629172918291929202921292229232924292529262927292829292930293129322933293429352936293729382939294029412942294329442945294629472948294929502951295229532954295529562957295829592960296129622963296429652966296729682969297029712972297329742975297629772978297929802981298229832984298529862987298829892990299129922993299429952996299729982999300030013002300330043005300630073008300930103011301230133014301530163017301830193020302130223023302430253026302730283029303030313032303330343035303630373038303930403041304230433044304530463047304830493050305130523053305430553056305730583059306030613062306330643065306630673068306930703071307230733074307530763077307830793080308130823083308430853086308730883089309030913092309330943095309630973098309931003101310231033104310531063107310831093110311131123113311431153116311731183119312031213122312331243125312631273128312931303131313231333134313531363137313831393140314131423143314431453146314731483149315031513152315331543155315631573158315931603161316231633164316531663167316831693170317131723173317431753176317731783179318031813182318331843185318631873188318931903191319231933194319531963197319831993200320132023203320432053206320732083209321032113212321332143215321632173218321932203221322232233224322532263227322832293230323132323233323432353236323732383239324032413242324332443245324632473248324932503251325232533254
  1. import { describe, expect } from "bun:test"
  2. import {
  3. LLMClient,
  4. LLMError,
  5. LLMEvent,
  6. Model,
  7. Tool,
  8. TransportReason,
  9. type LLMClientShape,
  10. type LLMRequest,
  11. } from "@opencode-ai/llm"
  12. import * as OpenAIChat from "@opencode-ai/llm/protocols/openai-chat"
  13. import { Database } from "@opencode-ai/core/database/database"
  14. import { EventV2 } from "@opencode-ai/core/event"
  15. import { PermissionV2 } from "@opencode-ai/core/permission"
  16. import { EventTable } from "@opencode-ai/core/event/sql"
  17. import { Project } from "@opencode-ai/core/project"
  18. import { ProjectTable } from "@opencode-ai/core/project/sql"
  19. import { QuestionV2 } from "@opencode-ai/core/question"
  20. import { AbsolutePath } from "@opencode-ai/core/schema"
  21. import { SessionV2 } from "@opencode-ai/core/session"
  22. import { ContextSnapshotDecodeError } from "@opencode-ai/core/session/error"
  23. import { SessionEvent } from "@opencode-ai/core/session/event"
  24. import { SessionInput } from "@opencode-ai/core/session/input"
  25. import { SessionMessage } from "@opencode-ai/core/session/message"
  26. import { Prompt } from "@opencode-ai/core/session/prompt"
  27. import { SessionProjector } from "@opencode-ai/core/session/projector"
  28. import { SessionExecution } from "@opencode-ai/core/session/execution"
  29. import { SessionContextEpoch } from "@opencode-ai/core/session/context-epoch"
  30. import { SessionRunCoordinator } from "@opencode-ai/core/session/run-coordinator"
  31. import { SessionRunner } from "@opencode-ai/core/session/runner"
  32. import * as SessionRunnerLLM from "@opencode-ai/core/session/runner/llm"
  33. import { SessionRunnerModel } from "@opencode-ai/core/session/runner/model"
  34. import { ToolRegistry } from "@opencode-ai/core/tool/registry"
  35. import { ApplicationTools } from "@opencode-ai/core/tool/application-tools"
  36. import { AgentV2 } from "@opencode-ai/core/agent"
  37. import { NativeTool } from "@opencode-ai/core/tool/native"
  38. import {
  39. SessionContextEpochTable,
  40. SessionInputTable,
  41. SessionMessageTable,
  42. SessionTable,
  43. } from "@opencode-ai/core/session/sql"
  44. import { SessionStore } from "@opencode-ai/core/session/store"
  45. import { SystemContext } from "@opencode-ai/core/system-context"
  46. import { SystemContextRegistry } from "@opencode-ai/core/system-context/registry"
  47. import { SkillGuidance } from "@opencode-ai/core/skill/guidance"
  48. import { ModelV2 } from "@opencode-ai/core/model"
  49. import { Location } from "@opencode-ai/core/location"
  50. import { ProviderV2 } from "@opencode-ai/core/provider"
  51. import { Cause, DateTime, Deferred, Effect, Exit, Fiber, Layer, Schema, Stream } from "effect"
  52. import { asc, eq } from "drizzle-orm"
  53. import { testEffect } from "./lib/effect"
  54. const database = Database.layerFromPath(":memory:")
  55. const events = EventV2.layer.pipe(Layer.provide(database))
  56. const questions = QuestionV2.layer.pipe(Layer.provide(events))
  57. const projector = SessionProjector.layer.pipe(Layer.provide(events), Layer.provide(database))
  58. const store = SessionStore.layer.pipe(Layer.provide(database))
  59. const requests: LLMRequest[] = []
  60. let response: LLMEvent[] = []
  61. let responses: LLMEvent[][] | undefined
  62. let responseStream: Stream.Stream<LLMEvent, LLMError> | undefined
  63. let streamGate: Deferred.Deferred<void> | undefined
  64. let streamStarted: Deferred.Deferred<void> | undefined
  65. let streamFailure: LLMError | undefined
  66. let toolExecutionGate: Deferred.Deferred<void> | undefined
  67. let toolExecutionsStarted: Deferred.Deferred<void> | undefined
  68. let toolExecutionsReady = 5
  69. let activeToolExecutions = 0
  70. let maxActiveToolExecutions = 0
  71. const client = Layer.succeed(
  72. LLMClient.Service,
  73. LLMClient.Service.of({
  74. prepare: () => Effect.die("unused"),
  75. stream: ((request: LLMRequest) => {
  76. requests.push(request)
  77. if (responseStream) {
  78. const stream = responseStream
  79. responseStream = undefined
  80. return stream
  81. }
  82. const events = streamFailure
  83. ? Stream.fail(streamFailure)
  84. : Stream.fromIterable(responses === undefined ? response : (responses.shift() ?? []))
  85. if (!streamGate) return events
  86. return Stream.unwrap(
  87. (streamStarted ? Deferred.succeed(streamStarted, undefined) : Effect.void).pipe(
  88. Effect.andThen(Deferred.await(streamGate)),
  89. Effect.as(events),
  90. ),
  91. )
  92. }) as unknown as LLMClientShape["stream"],
  93. generate: () => Effect.die("unused"),
  94. }),
  95. )
  96. const model = Model.make({ id: "fake-model", provider: "fake", route: OpenAIChat.route })
  97. const replacementModel = Model.make({ id: "replacement", provider: "fake", route: OpenAIChat.route })
  98. const authorizations: ToolRegistry.AuthorizeInput[] = []
  99. const executions: string[] = []
  100. const permission = Layer.succeed(
  101. PermissionV2.Service,
  102. PermissionV2.Service.of({
  103. assert: () => Effect.die("unused"),
  104. ask: () => Effect.die("unused"),
  105. reply: () => Effect.die("unused"),
  106. get: () => Effect.die("unused"),
  107. forSession: () => Effect.die("unused"),
  108. list: () => Effect.die("unused"),
  109. }),
  110. )
  111. const applications = ApplicationTools.layer
  112. const registry = ToolRegistry.layer.pipe(Layer.provide(permission), Layer.provide(applications))
  113. const agents = AgentV2.layer
  114. const echo = Layer.effectDiscard(
  115. ToolRegistry.Service.use((registry) =>
  116. registry.contribute((editor) => {
  117. ;(editor.set("echo", {
  118. authorize: (input) =>
  119. Effect.sync(() => {
  120. authorizations.push(input)
  121. }),
  122. tool: Tool.make({
  123. description: "Echo text",
  124. parameters: Schema.Struct({ text: Schema.String }),
  125. success: Schema.Struct({ text: Schema.String }),
  126. toModelOutput: ({ output }) => [{ type: "text", text: output.text }],
  127. execute: ({ text }) =>
  128. Effect.gen(function* () {
  129. executions.push(text)
  130. activeToolExecutions++
  131. maxActiveToolExecutions = Math.max(maxActiveToolExecutions, activeToolExecutions)
  132. if (activeToolExecutions === toolExecutionsReady && toolExecutionsStarted) {
  133. yield* Deferred.succeed(toolExecutionsStarted, undefined)
  134. }
  135. if (toolExecutionGate) yield* Deferred.await(toolExecutionGate)
  136. return { text }
  137. }).pipe(Effect.ensuring(Effect.sync(() => activeToolExecutions--))),
  138. }),
  139. }),
  140. editor.set("defect", {
  141. tool: Tool.make({
  142. description: "Fail unexpectedly",
  143. parameters: Schema.Struct({}),
  144. success: Schema.Struct({}),
  145. execute: () => Effect.die("unexpected tool defect"),
  146. }),
  147. }))
  148. }),
  149. ),
  150. ).pipe(Layer.provide(registry))
  151. let modelResolveHook = Effect.void
  152. const models = SessionRunnerModel.layerWith((session) =>
  153. modelResolveHook.pipe(Effect.as(session.model?.id === "replacement" ? replacementModel : model)),
  154. )
  155. const systemContextKey = SystemContext.Key.make("test/context")
  156. let systemBaseline = "Initial context"
  157. let systemRemoved = false
  158. let systemUnavailable = false
  159. let systemLoadHook = Effect.void
  160. const skillBaselines = new Map<AgentV2.ID, string>()
  161. const systemContext = Layer.effectDiscard(
  162. SystemContextRegistry.Service.pipe(
  163. Effect.flatMap((registry) =>
  164. registry.contribute({
  165. key: systemContextKey,
  166. load: Effect.sync(() =>
  167. SystemContext.combine(
  168. systemRemoved
  169. ? []
  170. : [
  171. SystemContext.make({
  172. key: systemContextKey,
  173. codec: Schema.toCodecJson(Schema.String),
  174. load: systemLoadHook.pipe(
  175. Effect.andThen(
  176. Effect.sync(() => (systemUnavailable ? SystemContext.unavailable : systemBaseline)),
  177. ),
  178. ),
  179. baseline: String,
  180. update: (_previous, current) => current,
  181. removed: () => "System context source removed: test/context",
  182. }),
  183. ],
  184. ),
  185. ),
  186. }),
  187. ),
  188. ),
  189. ).pipe(Layer.provideMerge(SystemContextRegistry.layer))
  190. const location = Location.layer({ directory: AbsolutePath.make("/project") }).pipe(Layer.provide(Project.defaultLayer))
  191. const skillGuidance = Layer.mock(SkillGuidance.Service, {
  192. load: (agent) =>
  193. Effect.succeed(
  194. skillBaselines.has(agent.id)
  195. ? SystemContext.make({
  196. key: SystemContext.Key.make("test/skill-guidance"),
  197. codec: Schema.toCodecJson(Schema.String),
  198. load: Effect.succeed(skillBaselines.get(agent.id)!),
  199. baseline: String,
  200. update: (_previous, current) => current,
  201. removed: () => "Skill guidance removed",
  202. })
  203. : SystemContext.empty,
  204. ),
  205. })
  206. const runner = SessionRunnerLLM.layer.pipe(
  207. Layer.provide(database),
  208. Layer.provide(store),
  209. Layer.provide(events),
  210. Layer.provide(client),
  211. Layer.provide(registry),
  212. Layer.provide(models),
  213. Layer.provide(systemContext),
  214. Layer.provide(location),
  215. Layer.provide(agents),
  216. Layer.provide(skillGuidance),
  217. )
  218. const coordinator = SessionRunCoordinator.layer.pipe(Layer.provide(runner))
  219. const execution = Layer.effect(
  220. SessionExecution.Service,
  221. SessionRunCoordinator.Service.pipe(
  222. Effect.map((coordinator) =>
  223. SessionExecution.Service.of({
  224. resume: coordinator.run,
  225. wake: coordinator.wake,
  226. interrupt: coordinator.interrupt,
  227. }),
  228. ),
  229. ),
  230. ).pipe(Layer.provide(coordinator))
  231. const sessions = SessionV2.layer.pipe(
  232. Layer.provide(events),
  233. Layer.provide(database),
  234. Layer.provide(store),
  235. Layer.provide(Project.defaultLayer),
  236. Layer.provide(execution),
  237. )
  238. const it = testEffect(
  239. Layer.mergeAll(
  240. database,
  241. events,
  242. questions,
  243. projector,
  244. store,
  245. client,
  246. permission,
  247. applications,
  248. agents,
  249. registry,
  250. echo,
  251. models,
  252. systemContext,
  253. location,
  254. skillGuidance,
  255. runner,
  256. coordinator,
  257. execution,
  258. sessions,
  259. ),
  260. )
  261. const sessionID = SessionV2.ID.make("ses_runner_test")
  262. const otherSessionID = SessionV2.ID.make("ses_runner_other")
  263. const insertSession = (id: SessionV2.ID) =>
  264. Effect.gen(function* () {
  265. const { db } = yield* Database.Service
  266. yield* db
  267. .insert(SessionTable)
  268. .values({
  269. id,
  270. project_id: Project.ID.global,
  271. slug: id,
  272. directory: "/project",
  273. title: "test",
  274. version: "test",
  275. })
  276. .onConflictDoNothing()
  277. .run()
  278. .pipe(Effect.orDie)
  279. })
  280. const setup = Effect.gen(function* () {
  281. const { db } = yield* Database.Service
  282. response = []
  283. systemBaseline = "Initial context"
  284. systemRemoved = false
  285. systemUnavailable = false
  286. systemLoadHook = Effect.void
  287. modelResolveHook = Effect.void
  288. skillBaselines.clear()
  289. responses = undefined
  290. streamFailure = undefined
  291. responseStream = undefined
  292. streamGate = undefined
  293. streamStarted = undefined
  294. toolExecutionGate = undefined
  295. toolExecutionsStarted = undefined
  296. toolExecutionsReady = 5
  297. activeToolExecutions = 0
  298. maxActiveToolExecutions = 0
  299. yield* db
  300. .insert(ProjectTable)
  301. .values({ id: Project.ID.global, worktree: AbsolutePath.make("/project"), sandboxes: [] })
  302. .onConflictDoNothing()
  303. .run()
  304. .pipe(Effect.orDie)
  305. yield* insertSession(sessionID)
  306. })
  307. const providerUnavailable = () =>
  308. new LLMError({
  309. module: "test",
  310. method: "stream",
  311. reason: new TransportReason({ message: "Provider unavailable" }),
  312. })
  313. const userTexts = (request: LLMRequest) =>
  314. request.messages.flatMap((message) =>
  315. message.role === "user"
  316. ? message.content.flatMap((content) => (content.type === "text" ? [content.text] : []))
  317. : [],
  318. )
  319. const replaySessionProjection = (id: SessionV2.ID) =>
  320. Effect.gen(function* () {
  321. const { db } = yield* Database.Service
  322. const events = yield* EventV2.Service
  323. const recorded = yield* db
  324. .select()
  325. .from(EventTable)
  326. .where(eq(EventTable.aggregate_id, id))
  327. .orderBy(asc(EventTable.seq))
  328. .all()
  329. .pipe(Effect.orDie)
  330. yield* events.remove(id)
  331. yield* db.delete(SessionInputTable).where(eq(SessionInputTable.session_id, id)).run().pipe(Effect.orDie)
  332. yield* db.delete(SessionMessageTable).where(eq(SessionMessageTable.session_id, id)).run().pipe(Effect.orDie)
  333. yield* events.replayAll(
  334. recorded.map((event) => ({
  335. id: event.id,
  336. aggregateID: event.aggregate_id,
  337. seq: event.seq,
  338. type: event.type,
  339. data: event.data,
  340. })),
  341. )
  342. })
  343. type FragmentKind = "text" | "reasoning" | "tool input"
  344. type FragmentFixture = {
  345. readonly delta: EventV2.Definition
  346. readonly completeEvents: LLMEvent[]
  347. readonly partialEvents: LLMEvent[]
  348. readonly expectedAssistant: unknown
  349. readonly expectedContent: unknown
  350. }
  351. const fragmentKinds: readonly FragmentKind[] = ["text", "reasoning", "tool input"]
  352. const fragmentID = (kind: FragmentKind, suffix: string) => `${kind === "tool input" ? "call" : kind}-${suffix}`
  353. const fragmentFixture = (kind: FragmentKind, id: string, chunks: readonly string[]): FragmentFixture => {
  354. const text = chunks.join("")
  355. switch (kind) {
  356. case "text": {
  357. const partialEvents = [
  358. LLMEvent.stepStart({ index: 0 }),
  359. LLMEvent.textStart({ id }),
  360. ...chunks.map((text) => LLMEvent.textDelta({ id, text })),
  361. ]
  362. const expectedContent = { type: "text", id, text }
  363. return {
  364. delta: SessionEvent.Text.Delta,
  365. partialEvents,
  366. completeEvents: [
  367. ...partialEvents,
  368. LLMEvent.textEnd({ id }),
  369. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  370. LLMEvent.finish({ reason: "stop" }),
  371. ],
  372. expectedAssistant: { type: "assistant", finish: "stop", content: [expectedContent] },
  373. expectedContent,
  374. }
  375. }
  376. case "reasoning": {
  377. const partialEvents = [
  378. LLMEvent.stepStart({ index: 0 }),
  379. LLMEvent.reasoningStart({ id }),
  380. ...chunks.map((text) => LLMEvent.reasoningDelta({ id, text })),
  381. ]
  382. const expectedContent = { type: "reasoning", id, text }
  383. return {
  384. delta: SessionEvent.Reasoning.Delta,
  385. partialEvents,
  386. completeEvents: [
  387. ...partialEvents,
  388. LLMEvent.reasoningEnd({ id }),
  389. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  390. LLMEvent.finish({ reason: "stop" }),
  391. ],
  392. expectedAssistant: { type: "assistant", finish: "stop", content: [expectedContent] },
  393. expectedContent,
  394. }
  395. }
  396. case "tool input": {
  397. const partialEvents = [
  398. LLMEvent.stepStart({ index: 0 }),
  399. LLMEvent.toolInputStart({ id, name: "echo" }),
  400. ...chunks.map((text) => LLMEvent.toolInputDelta({ id, name: "echo", text })),
  401. ]
  402. const expectedContent = { type: "tool", id, state: { status: "pending", input: text } }
  403. return {
  404. delta: SessionEvent.Tool.Input.Delta,
  405. partialEvents,
  406. completeEvents: [...partialEvents, LLMEvent.toolInputEnd({ id, name: "echo" })],
  407. expectedAssistant: { type: "assistant", content: [expectedContent] },
  408. expectedContent,
  409. }
  410. }
  411. }
  412. }
  413. const verifyEphemeralDeltas = (kind: FragmentKind) =>
  414. Effect.gen(function* () {
  415. yield* setup
  416. const session = yield* SessionV2.Service
  417. const prompt = `Stream ${kind}`
  418. const chunks = Array.from({ length: 32 }, (_, index) => `${index},`)
  419. const fixture = fragmentFixture(kind, fragmentID(kind, "many"), chunks)
  420. const expectedContext = [{ type: "user", text: prompt }, fixture.expectedAssistant]
  421. yield* session.prompt({ sessionID, prompt: new Prompt({ text: prompt }), resume: false })
  422. const events = yield* EventV2.Service
  423. const live = yield* events.subscribe(fixture.delta).pipe(Stream.take(32), Stream.runCollect, Effect.forkScoped)
  424. yield* Effect.yieldNow
  425. response = fixture.completeEvents
  426. yield* session.resume(sessionID)
  427. const { db } = yield* Database.Service
  428. const deltas = yield* db
  429. .select({ type: EventTable.type })
  430. .from(EventTable)
  431. .where(eq(EventTable.type, EventV2.versionedType(fixture.delta.type, 1)))
  432. .all()
  433. .pipe(Effect.orDie)
  434. expect(Array.from(yield* Fiber.join(live))).toHaveLength(32)
  435. expect(deltas).toHaveLength(0)
  436. expect(yield* session.context(sessionID)).toMatchObject(expectedContext)
  437. yield* replaySessionProjection(sessionID)
  438. expect(yield* session.context(sessionID)).toMatchObject(expectedContext)
  439. })
  440. const verifyPartialFlushOnFailure = (kind: FragmentKind) =>
  441. Effect.gen(function* () {
  442. yield* setup
  443. const session = yield* SessionV2.Service
  444. const prompt = `Fail after ${kind}`
  445. const fixture = fragmentFixture(kind, fragmentID(kind, "partial"), ["Partial"])
  446. const failure = providerUnavailable()
  447. yield* session.prompt({ sessionID, prompt: new Prompt({ text: prompt }), resume: false })
  448. responseStream = Stream.concat(Stream.fromIterable(fixture.partialEvents), Stream.fail(failure))
  449. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  450. expect(yield* session.context(sessionID)).toMatchObject([
  451. { type: "user", text: prompt },
  452. {
  453. type: "assistant",
  454. finish: "error",
  455. error: { type: "unknown", message: "Provider unavailable" },
  456. content: [fixture.expectedContent],
  457. },
  458. ])
  459. })
  460. const verifyPartialFlushOnInterruption = (kind: FragmentKind) =>
  461. Effect.gen(function* () {
  462. yield* setup
  463. const session = yield* SessionV2.Service
  464. const prompt = `Interrupt after ${kind}`
  465. const fixture = fragmentFixture(kind, fragmentID(kind, "interrupted"), ["Partial"])
  466. const streamed = yield* Deferred.make<void>()
  467. yield* session.prompt({ sessionID, prompt: new Prompt({ text: prompt }), resume: false })
  468. responseStream = Stream.concat(
  469. Stream.fromIterable(fixture.partialEvents),
  470. Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)),
  471. )
  472. const runner = yield* SessionRunner.Service
  473. const fiber = yield* runner.run({ sessionID, force: true }).pipe(Effect.forkChild)
  474. yield* Deferred.await(streamed)
  475. yield* Fiber.interrupt(fiber)
  476. expect(yield* session.context(sessionID)).toMatchObject([
  477. { type: "user", text: prompt },
  478. {
  479. type: "assistant",
  480. content: [
  481. kind === "tool input"
  482. ? { type: "tool", id: fragmentID(kind, "interrupted"), state: { status: "error" } }
  483. : fixture.expectedContent,
  484. ],
  485. },
  486. ])
  487. })
  488. describe("SessionRunnerLLM", () => {
  489. it.effect("advertises and executes a globally attached application tool", () =>
  490. Effect.gen(function* () {
  491. yield* setup
  492. const applicationTools = yield* ApplicationTools.Service
  493. const session = yield* SessionV2.Service
  494. const contexts: NativeTool.Context[] = []
  495. yield* applicationTools.attach({
  496. application_context: NativeTool.make({
  497. description: "Read application context",
  498. parameters: Schema.Struct({ query: Schema.String }),
  499. success: Schema.Struct({ answer: Schema.String }),
  500. execute: ({ query }, context) =>
  501. Effect.sync(() => {
  502. contexts.push(context)
  503. return { answer: query.toUpperCase() }
  504. }),
  505. }),
  506. })
  507. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Use application context" }), resume: false })
  508. responses = [
  509. [
  510. LLMEvent.stepStart({ index: 0 }),
  511. LLMEvent.toolCall({ id: "call-application", name: "application_context", input: { query: "hello" } }),
  512. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  513. LLMEvent.finish({ reason: "tool-calls" }),
  514. ],
  515. [],
  516. ]
  517. yield* session.resume(sessionID)
  518. expect(requests[0]?.tools.map((tool) => tool.name)).toContain("application_context")
  519. expect(contexts).toEqual([{ sessionID, id: "call-application", name: "application_context" }])
  520. expect(yield* session.context(sessionID)).toMatchObject([
  521. { type: "user", text: "Use application context" },
  522. {
  523. type: "assistant",
  524. content: [
  525. {
  526. type: "tool",
  527. id: "call-application",
  528. state: { status: "completed", structured: { answer: "HELLO" } },
  529. },
  530. ],
  531. },
  532. ])
  533. }),
  534. )
  535. it.effect("starts a real runner turn after default prompt recording", () =>
  536. Effect.gen(function* () {
  537. yield* setup
  538. const session = yield* SessionV2.Service
  539. requests.length = 0
  540. responses = undefined
  541. streamGate = undefined
  542. streamStarted = undefined
  543. response = []
  544. const message = yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Run automatically" }) })
  545. expect(requests).toHaveLength(1)
  546. expect(yield* session.messages({ sessionID })).toMatchObject([
  547. { id: message.id, type: "user", text: "Run automatically" },
  548. ])
  549. }),
  550. )
  551. it.effect("streams one request with registry definitions from chronological V2 user history", () =>
  552. Effect.gen(function* () {
  553. yield* setup
  554. const session = yield* SessionV2.Service
  555. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  556. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  557. requests.length = 0
  558. responses = undefined
  559. streamGate = undefined
  560. streamStarted = undefined
  561. response = []
  562. yield* session.resume(sessionID)
  563. expect(requests).toHaveLength(1)
  564. expect(requests[0]?.model).toBe(model)
  565. expect(requests[0]?.tools.map((tool) => tool.name)).toEqual(["echo", "defect"])
  566. expect(requests[0]?.messages.map((message) => ({ role: message.role, content: message.content }))).toEqual([
  567. { role: "user", content: [{ type: "text", text: "First" }] },
  568. { role: "user", content: [{ type: "text", text: "Second" }] },
  569. ])
  570. expect(yield* session.messages({ sessionID })).toHaveLength(2)
  571. }),
  572. )
  573. it.effect("retries the first provider turn after system context becomes available", () =>
  574. Effect.gen(function* () {
  575. yield* setup
  576. const session = yield* SessionV2.Service
  577. const { db } = yield* Database.Service
  578. const messageID = SessionMessage.ID.create()
  579. systemUnavailable = true
  580. yield* session.prompt({ id: messageID, sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  581. requests.length = 0
  582. const exit = yield* session.resume(sessionID).pipe(Effect.exit)
  583. expect(Exit.isFailure(exit)).toBe(true)
  584. if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(SystemContext.InitializationBlocked)
  585. expect(requests).toHaveLength(0)
  586. expect(yield* SessionInput.hasPending(db, sessionID, "steer")).toBe(true)
  587. expect(
  588. yield* db
  589. .select()
  590. .from(SessionContextEpochTable)
  591. .where(eq(SessionContextEpochTable.session_id, sessionID))
  592. .get(),
  593. ).toBeUndefined()
  594. systemUnavailable = false
  595. yield* session.prompt({ id: messageID, sessionID, prompt: new Prompt({ text: "First" }) })
  596. yield* (yield* SessionRunCoordinator.Service).awaitIdle(sessionID)
  597. expect(requests).toHaveLength(1)
  598. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user"])
  599. }),
  600. )
  601. it.effect("interrupts a source Location runner after a Session moves", () =>
  602. Effect.gen(function* () {
  603. yield* setup
  604. const session = yield* SessionV2.Service
  605. const events = yield* EventV2.Service
  606. const { db } = yield* Database.Service
  607. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  608. requests.length = 0
  609. response = []
  610. yield* session.resume(sessionID)
  611. yield* events.publish(SessionEvent.Moved, {
  612. sessionID,
  613. timestamp: DateTime.makeUnsafe(1),
  614. location: { directory: AbsolutePath.make("/moved") },
  615. })
  616. expect(
  617. yield* db
  618. .select()
  619. .from(SessionContextEpochTable)
  620. .where(eq(SessionContextEpochTable.session_id, sessionID))
  621. .get(),
  622. ).toBeUndefined()
  623. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  624. const exit = yield* session.resume(sessionID).pipe(Effect.exit)
  625. expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true)
  626. expect(requests).toHaveLength(1)
  627. expect(yield* SessionInput.hasPending(db, sessionID, "steer")).toBe(true)
  628. }),
  629. )
  630. it.effect("fails gracefully when a stored context snapshot cannot be decoded", () =>
  631. Effect.gen(function* () {
  632. yield* setup
  633. const session = yield* SessionV2.Service
  634. const { db } = yield* Database.Service
  635. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  636. response = []
  637. yield* session.resume(sessionID)
  638. yield* db
  639. .update(SessionContextEpochTable)
  640. .set({ snapshot: { invalid: { value: "bad" } } })
  641. .where(eq(SessionContextEpochTable.session_id, sessionID))
  642. .run()
  643. .pipe(Effect.orDie)
  644. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  645. requests.length = 0
  646. const exit = yield* session.resume(sessionID).pipe(Effect.exit)
  647. expect(Exit.isFailure(exit)).toBe(true)
  648. if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(ContextSnapshotDecodeError)
  649. expect(requests).toHaveLength(0)
  650. }),
  651. )
  652. it.effect("does not create a source Location epoch after a concurrent Session move", () =>
  653. Effect.gen(function* () {
  654. yield* setup
  655. const session = yield* SessionV2.Service
  656. const events = yield* EventV2.Service
  657. const { db } = yield* Database.Service
  658. let moved = false
  659. systemLoadHook = Effect.suspend(() => {
  660. if (moved) return Effect.void
  661. moved = true
  662. return events
  663. .publish(SessionEvent.Moved, {
  664. sessionID,
  665. timestamp: DateTime.makeUnsafe(1),
  666. location: { directory: AbsolutePath.make("/moved") },
  667. })
  668. .pipe(Effect.asVoid)
  669. })
  670. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  671. expect(Exit.isFailure(yield* session.resume(sessionID).pipe(Effect.exit))).toBe(true)
  672. expect(yield* SessionInput.hasPending(db, sessionID, "steer")).toBe(true)
  673. expect(
  674. yield* db
  675. .select()
  676. .from(SessionContextEpochTable)
  677. .where(eq(SessionContextEpochTable.session_id, sessionID))
  678. .get(),
  679. ).toBeUndefined()
  680. expect((yield* session.get(sessionID)).location.directory).toBe(AbsolutePath.make("/moved"))
  681. }),
  682. )
  683. it.effect("reuses one durable baseline after the context producer changes", () =>
  684. Effect.gen(function* () {
  685. yield* setup
  686. const session = yield* SessionV2.Service
  687. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  688. requests.length = 0
  689. response = []
  690. yield* session.resume(sessionID)
  691. systemBaseline = "Changed context"
  692. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  693. yield* session.resume(sessionID)
  694. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  695. ["Initial context"],
  696. ["Initial context"],
  697. ])
  698. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "user", "system"])
  699. expect(requests[1]?.messages.at(-1)?.content).toEqual([{ type: "text", text: "Changed context" }])
  700. expect(yield* session.messages({ sessionID })).toHaveLength(3)
  701. const { db } = yield* Database.Service
  702. expect(
  703. yield* db
  704. .select({ id: EventTable.id })
  705. .from(EventTable)
  706. .where(eq(EventTable.type, "session.next.context.updated.1"))
  707. .all()
  708. .pipe(Effect.orDie),
  709. ).toHaveLength(1)
  710. yield* replaySessionProjection(sessionID)
  711. expect(yield* session.messages({ sessionID })).toHaveLength(3)
  712. }),
  713. )
  714. it.effect("includes the effective default agent system before durable context", () =>
  715. Effect.gen(function* () {
  716. yield* setup
  717. const agent = yield* AgentV2.Service
  718. yield* agent.update((editor) =>
  719. editor.update(AgentV2.ID.make("build"), (agent) => {
  720. agent.system = "Build agent instructions"
  721. agent.mode = "primary"
  722. }),
  723. )
  724. const session = yield* SessionV2.Service
  725. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  726. requests.length = 0
  727. response = fragmentFixture("text", "text-build", ["Done"]).completeEvents
  728. yield* session.resume(sessionID)
  729. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Build agent instructions", "Initial context"])
  730. }),
  731. )
  732. it.effect("uses the configured default agent system for omitted-agent sessions", () =>
  733. Effect.gen(function* () {
  734. yield* setup
  735. const agent = yield* AgentV2.Service
  736. yield* agent.update((editor) => {
  737. editor.update(AgentV2.ID.make("build"), (agent) => {
  738. agent.system = "Build agent instructions"
  739. agent.mode = "primary"
  740. })
  741. editor.update(AgentV2.ID.make("reviewer"), (agent) => {
  742. agent.system = "Reviewer instructions"
  743. agent.mode = "primary"
  744. })
  745. editor.default(AgentV2.ID.make("reviewer"))
  746. })
  747. const session = yield* SessionV2.Service
  748. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  749. requests.length = 0
  750. response = fragmentFixture("text", "text-reviewer", ["Done"]).completeEvents
  751. yield* session.resume(sessionID)
  752. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"])
  753. expect((yield* session.messages({ sessionID }))[0]).toMatchObject({ type: "assistant", agent: "reviewer" })
  754. }),
  755. )
  756. it.effect("uses an explicitly selected non-build agent system", () =>
  757. Effect.gen(function* () {
  758. yield* setup
  759. const { db } = yield* Database.Service
  760. const agent = yield* AgentV2.Service
  761. yield* agent.update((editor) =>
  762. editor.update(AgentV2.ID.make("reviewer"), (agent) => {
  763. agent.system = "Reviewer instructions"
  764. agent.mode = "primary"
  765. }),
  766. )
  767. yield* db
  768. .update(SessionTable)
  769. .set({ agent: "reviewer" })
  770. .where(eq(SessionTable.id, sessionID))
  771. .run()
  772. .pipe(Effect.orDie)
  773. const session = yield* SessionV2.Service
  774. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  775. requests.length = 0
  776. response = fragmentFixture("text", "text-selected", ["Done"]).completeEvents
  777. yield* session.resume(sessionID)
  778. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"])
  779. expect((yield* session.messages({ sessionID }))[0]).toMatchObject({ type: "assistant", agent: "reviewer" })
  780. }),
  781. )
  782. it.effect("composes selected-agent skill guidance and replaces it after an agent switch", () =>
  783. Effect.gen(function* () {
  784. yield* setup
  785. const session = yield* SessionV2.Service
  786. const events = yield* EventV2.Service
  787. skillBaselines.set(AgentV2.ID.make("build"), "Build skills")
  788. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  789. requests.length = 0
  790. response = []
  791. yield* session.resume(sessionID)
  792. skillBaselines.set(AgentV2.ID.make("reviewer"), "Reviewer skills")
  793. yield* events.publish(SessionEvent.AgentSwitched, {
  794. sessionID,
  795. messageID: SessionMessage.ID.create(),
  796. timestamp: DateTime.makeUnsafe(1),
  797. agent: "reviewer",
  798. })
  799. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  800. yield* session.resume(sessionID)
  801. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  802. ["Initial context\n\nBuild skills"],
  803. ["Initial context\n\nReviewer skills"],
  804. ])
  805. }),
  806. )
  807. it.effect("retries first-epoch preparation when the selected agent changes during observation", () =>
  808. Effect.gen(function* () {
  809. yield* setup
  810. const session = yield* SessionV2.Service
  811. const events = yield* EventV2.Service
  812. skillBaselines.set(AgentV2.ID.make("build"), "Build skills")
  813. skillBaselines.set(AgentV2.ID.make("reviewer"), "Reviewer skills")
  814. let switched = false
  815. systemLoadHook = Effect.suspend(() => {
  816. if (switched) return Effect.void
  817. switched = true
  818. return events
  819. .publish(SessionEvent.AgentSwitched, {
  820. sessionID,
  821. messageID: SessionMessage.ID.create(),
  822. timestamp: DateTime.makeUnsafe(1),
  823. agent: "reviewer",
  824. })
  825. .pipe(Effect.asVoid)
  826. })
  827. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  828. requests.length = 0
  829. response = []
  830. yield* session.resume(sessionID)
  831. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  832. ["Initial context\n\nReviewer skills"],
  833. ])
  834. }),
  835. )
  836. it.effect("opens a queued activity once when the selected agent changes during observation", () =>
  837. Effect.gen(function* () {
  838. yield* setup
  839. const session = yield* SessionV2.Service
  840. const events = yield* EventV2.Service
  841. skillBaselines.set(AgentV2.ID.make("build"), "Build skills")
  842. skillBaselines.set(AgentV2.ID.make("reviewer"), "Reviewer skills")
  843. let switched = false
  844. systemLoadHook = Effect.suspend(() => {
  845. if (switched) return Effect.void
  846. switched = true
  847. return events
  848. .publish(SessionEvent.AgentSwitched, {
  849. sessionID,
  850. messageID: SessionMessage.ID.create(),
  851. timestamp: DateTime.makeUnsafe(1),
  852. agent: "reviewer",
  853. })
  854. .pipe(Effect.asVoid)
  855. })
  856. yield* session.prompt({
  857. sessionID,
  858. prompt: new Prompt({ text: "Queued" }),
  859. delivery: "queue",
  860. resume: false,
  861. })
  862. requests.length = 0
  863. response = []
  864. yield* session.resume(sessionID)
  865. expect(requests).toHaveLength(1)
  866. expect((yield* session.context(sessionID)).filter((message) => message.type === "user")).toHaveLength(1)
  867. }),
  868. )
  869. it.effect("retries an agent switch before the final provider-dispatch boundary", () =>
  870. Effect.gen(function* () {
  871. yield* setup
  872. const session = yield* SessionV2.Service
  873. const events = yield* EventV2.Service
  874. const { db } = yield* Database.Service
  875. skillBaselines.set(AgentV2.ID.make("build"), "Build skills")
  876. skillBaselines.set(AgentV2.ID.make("reviewer"), "Reviewer skills")
  877. let switched = false
  878. modelResolveHook = Effect.suspend(() => {
  879. if (switched) return Effect.void
  880. switched = true
  881. return events
  882. .publish(SessionEvent.AgentSwitched, {
  883. sessionID,
  884. messageID: SessionMessage.ID.create(),
  885. timestamp: DateTime.makeUnsafe(1),
  886. agent: "reviewer",
  887. })
  888. .pipe(Effect.asVoid)
  889. })
  890. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  891. requests.length = 0
  892. response = []
  893. yield* session.resume(sessionID)
  894. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  895. ["Initial context\n\nReviewer skills"],
  896. ])
  897. expect(
  898. yield* db
  899. .select({ replacementSeq: SessionContextEpochTable.replacement_seq })
  900. .from(SessionContextEpochTable)
  901. .where(eq(SessionContextEpochTable.session_id, sessionID))
  902. .get()
  903. .pipe(Effect.orDie),
  904. ).toEqual({ replacementSeq: null })
  905. }),
  906. )
  907. it.effect("retries a model switch before the final provider-dispatch boundary", () =>
  908. Effect.gen(function* () {
  909. yield* setup
  910. const session = yield* SessionV2.Service
  911. const events = yield* EventV2.Service
  912. let switched = false
  913. modelResolveHook = Effect.suspend(() => {
  914. if (switched) return Effect.void
  915. switched = true
  916. return events
  917. .publish(SessionEvent.ModelSwitched, {
  918. sessionID,
  919. messageID: SessionMessage.ID.create(),
  920. timestamp: DateTime.makeUnsafe(1),
  921. model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") },
  922. })
  923. .pipe(Effect.asVoid)
  924. })
  925. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  926. requests.length = 0
  927. response = []
  928. yield* session.resume(sessionID)
  929. expect(requests.map((request) => request.model)).toEqual([replacementModel])
  930. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([["Initial context"]])
  931. }),
  932. )
  933. it.effect("fences an unchanged epoch read across an agent ABA replacement request", () =>
  934. Effect.gen(function* () {
  935. yield* setup
  936. const session = yield* SessionV2.Service
  937. const events = yield* EventV2.Service
  938. const { db } = yield* Database.Service
  939. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  940. response = []
  941. yield* session.resume(sessionID)
  942. let switched = false
  943. systemLoadHook = Effect.suspend(() => {
  944. if (switched) return Effect.void
  945. switched = true
  946. return events
  947. .publish(SessionEvent.AgentSwitched, {
  948. sessionID,
  949. messageID: SessionMessage.ID.create(),
  950. timestamp: DateTime.makeUnsafe(1),
  951. agent: AgentV2.ID.make("reviewer"),
  952. })
  953. .pipe(
  954. Effect.andThen(
  955. events.publish(SessionEvent.AgentSwitched, {
  956. sessionID,
  957. messageID: SessionMessage.ID.create(),
  958. timestamp: DateTime.makeUnsafe(2),
  959. agent: AgentV2.defaultID,
  960. }),
  961. ),
  962. Effect.asVoid,
  963. )
  964. })
  965. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  966. requests.length = 0
  967. yield* session.resume(sessionID)
  968. expect(requests).toHaveLength(1)
  969. expect(
  970. yield* db
  971. .select({ replacementSeq: SessionContextEpochTable.replacement_seq })
  972. .from(SessionContextEpochTable)
  973. .where(eq(SessionContextEpochTable.session_id, sessionID))
  974. .get()
  975. .pipe(Effect.orDie),
  976. ).toEqual({ replacementSeq: null })
  977. }),
  978. )
  979. it.effect("rejects stale agent guidance when committing an existing-epoch replacement", () =>
  980. Effect.gen(function* () {
  981. yield* setup
  982. const session = yield* SessionV2.Service
  983. const events = yield* EventV2.Service
  984. const { db } = yield* Database.Service
  985. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  986. response = []
  987. yield* session.resume(sessionID)
  988. yield* events.publish(SessionEvent.AgentSwitched, {
  989. sessionID,
  990. messageID: SessionMessage.ID.create(),
  991. timestamp: DateTime.makeUnsafe(1),
  992. agent: AgentV2.ID.make("reviewer"),
  993. })
  994. const context = (text: string) =>
  995. Effect.succeed(
  996. SystemContext.make({
  997. key: systemContextKey,
  998. codec: Schema.toCodecJson(Schema.String),
  999. load: Effect.succeed(text),
  1000. baseline: String,
  1001. update: (_previous, current) => current,
  1002. }),
  1003. )
  1004. const location = (yield* session.get(sessionID)).location
  1005. expect(
  1006. yield* SessionContextEpoch.prepare(
  1007. db,
  1008. events,
  1009. context("Stale build context"),
  1010. sessionID,
  1011. location,
  1012. AgentV2.defaultID,
  1013. ).pipe(Effect.catchDefect(Effect.succeed)),
  1014. ).toBeInstanceOf(SessionContextEpoch.AgentMismatch)
  1015. expect(
  1016. yield* SessionContextEpoch.prepare(
  1017. db,
  1018. events,
  1019. context("Reviewer context"),
  1020. sessionID,
  1021. location,
  1022. AgentV2.ID.make("reviewer"),
  1023. ),
  1024. ).toMatchObject({ baseline: "Reviewer context" })
  1025. }),
  1026. )
  1027. it.effect("blocks a cross-agent provider turn while replacement context is unavailable", () =>
  1028. Effect.gen(function* () {
  1029. yield* setup
  1030. const session = yield* SessionV2.Service
  1031. const events = yield* EventV2.Service
  1032. skillBaselines.set(AgentV2.defaultID, "Build skills")
  1033. skillBaselines.set(AgentV2.ID.make("reviewer"), "Reviewer skills")
  1034. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  1035. response = []
  1036. yield* session.resume(sessionID)
  1037. yield* events.publish(SessionEvent.AgentSwitched, {
  1038. sessionID,
  1039. messageID: SessionMessage.ID.create(),
  1040. timestamp: DateTime.makeUnsafe(1),
  1041. agent: AgentV2.ID.make("reviewer"),
  1042. })
  1043. systemUnavailable = true
  1044. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  1045. requests.length = 0
  1046. const blocked = yield* session.resume(sessionID).pipe(Effect.exit)
  1047. expect(Exit.isFailure(blocked)).toBe(true)
  1048. if (Exit.isFailure(blocked))
  1049. expect(Cause.squash(blocked.cause)).toBeInstanceOf(SessionContextEpoch.AgentReplacementBlocked)
  1050. expect(requests).toHaveLength(0)
  1051. systemUnavailable = false
  1052. yield* session.resume(sessionID)
  1053. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1054. ["Initial context\n\nReviewer skills"],
  1055. ])
  1056. }),
  1057. )
  1058. it.effect("admits removed context as a chronological System message", () =>
  1059. Effect.gen(function* () {
  1060. yield* setup
  1061. const session = yield* SessionV2.Service
  1062. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  1063. requests.length = 0
  1064. response = []
  1065. yield* session.resume(sessionID)
  1066. systemRemoved = true
  1067. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  1068. yield* session.resume(sessionID)
  1069. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "user", "system"])
  1070. expect(requests[1]?.messages.at(-1)?.content).toEqual([
  1071. { type: "text", text: "System context source removed: test/context" },
  1072. ])
  1073. expect(yield* session.messages({ sessionID })).toHaveLength(3)
  1074. }),
  1075. )
  1076. it.effect("replaces the baseline lazily after a model switch and drops prior System updates", () =>
  1077. Effect.gen(function* () {
  1078. yield* setup
  1079. const session = yield* SessionV2.Service
  1080. const events = yield* EventV2.Service
  1081. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  1082. requests.length = 0
  1083. response = []
  1084. yield* session.resume(sessionID)
  1085. systemBaseline = "Changed context"
  1086. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  1087. yield* session.resume(sessionID)
  1088. yield* events.publish(SessionEvent.ModelSwitched, {
  1089. sessionID,
  1090. messageID: SessionMessage.ID.create(),
  1091. timestamp: DateTime.makeUnsafe(1),
  1092. model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") },
  1093. })
  1094. systemBaseline = "Replacement context"
  1095. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Third" }), resume: false })
  1096. yield* session.resume(sessionID)
  1097. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1098. ["Initial context"],
  1099. ["Initial context"],
  1100. ["Replacement context"],
  1101. ])
  1102. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "user", "system"])
  1103. expect(requests[2]?.messages.map((message) => message.role)).toEqual(["user", "user", "user"])
  1104. expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([
  1105. "user",
  1106. "user",
  1107. "model-switched",
  1108. "user",
  1109. ])
  1110. yield* replaySessionProjection(sessionID)
  1111. expect(yield* session.messages({ sessionID })).toHaveLength(5)
  1112. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Fourth" }), resume: false })
  1113. yield* session.resume(sessionID)
  1114. }),
  1115. )
  1116. it.effect("defers replacement while admitted context is temporarily unavailable", () =>
  1117. Effect.gen(function* () {
  1118. yield* setup
  1119. const session = yield* SessionV2.Service
  1120. const events = yield* EventV2.Service
  1121. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  1122. requests.length = 0
  1123. response = []
  1124. yield* session.resume(sessionID)
  1125. yield* events.publish(SessionEvent.ModelSwitched, {
  1126. sessionID,
  1127. messageID: SessionMessage.ID.create(),
  1128. timestamp: DateTime.makeUnsafe(1),
  1129. model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") },
  1130. })
  1131. systemUnavailable = true
  1132. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  1133. yield* session.resume(sessionID)
  1134. systemUnavailable = false
  1135. systemBaseline = "Replacement context"
  1136. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Third" }), resume: false })
  1137. yield* session.resume(sessionID)
  1138. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1139. ["Initial context"],
  1140. ["Initial context"],
  1141. ["Replacement context"],
  1142. ])
  1143. }),
  1144. )
  1145. it.effect("advances a pending replacement to the latest invalidation boundary", () =>
  1146. Effect.gen(function* () {
  1147. yield* setup
  1148. const session = yield* SessionV2.Service
  1149. const events = yield* EventV2.Service
  1150. const { db } = yield* Database.Service
  1151. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  1152. response = []
  1153. yield* session.resume(sessionID)
  1154. yield* events.publish(SessionEvent.ModelSwitched, {
  1155. sessionID,
  1156. messageID: SessionMessage.ID.create(),
  1157. timestamp: DateTime.makeUnsafe(1),
  1158. model: { id: ModelV2.ID.make("replacement-1"), providerID: ProviderV2.ID.make("fake") },
  1159. })
  1160. yield* events.publish(SessionEvent.ModelSwitched, {
  1161. sessionID,
  1162. messageID: SessionMessage.ID.create(),
  1163. timestamp: DateTime.makeUnsafe(2),
  1164. model: { id: ModelV2.ID.make("replacement-2"), providerID: ProviderV2.ID.make("fake") },
  1165. })
  1166. const latest = yield* SessionInput.latestSeq(db, sessionID)
  1167. expect(
  1168. yield* db
  1169. .select({ replacementSeq: SessionContextEpochTable.replacement_seq })
  1170. .from(SessionContextEpochTable)
  1171. .where(eq(SessionContextEpochTable.session_id, sessionID))
  1172. .get()
  1173. .pipe(Effect.orDie),
  1174. ).toEqual({ replacementSeq: latest })
  1175. }),
  1176. )
  1177. it.effect("retries epoch preparation until observation-time invalidations settle", () =>
  1178. Effect.gen(function* () {
  1179. yield* setup
  1180. const session = yield* SessionV2.Service
  1181. const events = yield* EventV2.Service
  1182. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  1183. response = []
  1184. yield* session.resume(sessionID)
  1185. requests.length = 0
  1186. systemBaseline = "Changed context"
  1187. let invalidations = 0
  1188. systemLoadHook = Effect.suspend(() => {
  1189. if (invalidations === 4) return Effect.void
  1190. invalidations++
  1191. return events
  1192. .publish(SessionEvent.ModelSwitched, {
  1193. sessionID,
  1194. messageID: SessionMessage.ID.create(),
  1195. timestamp: DateTime.makeUnsafe(invalidations),
  1196. model: { id: ModelV2.ID.make(`replacement-${invalidations}`), providerID: ProviderV2.ID.make("fake") },
  1197. })
  1198. .pipe(Effect.asVoid)
  1199. })
  1200. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  1201. yield* session.resume(sessionID)
  1202. expect(invalidations).toBe(4)
  1203. expect(requests).toHaveLength(1)
  1204. expect(requests[0]?.system.map((part) => part.text)).toEqual(["Changed context"])
  1205. }),
  1206. )
  1207. it.effect("replays retained context projections while replacement is pending", () =>
  1208. Effect.gen(function* () {
  1209. yield* setup
  1210. const session = yield* SessionV2.Service
  1211. const events = yield* EventV2.Service
  1212. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  1213. requests.length = 0
  1214. response = []
  1215. yield* session.resume(sessionID)
  1216. systemBaseline = "Changed context"
  1217. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  1218. yield* session.resume(sessionID)
  1219. yield* events.publish(SessionEvent.ModelSwitched, {
  1220. sessionID,
  1221. messageID: SessionMessage.ID.create(),
  1222. timestamp: DateTime.makeUnsafe(1),
  1223. model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") },
  1224. })
  1225. yield* replaySessionProjection(sessionID)
  1226. systemBaseline = "Replacement context"
  1227. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Third" }), resume: false })
  1228. yield* session.resume(sessionID)
  1229. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Replacement context"])
  1230. }),
  1231. )
  1232. it.effect("replaces the baseline lazily after completed compaction without reopening replacement on replay", () =>
  1233. Effect.gen(function* () {
  1234. yield* setup
  1235. const session = yield* SessionV2.Service
  1236. const events = yield* EventV2.Service
  1237. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  1238. requests.length = 0
  1239. response = []
  1240. yield* session.resume(sessionID)
  1241. yield* events.publish(SessionEvent.Compaction.Started, {
  1242. sessionID,
  1243. messageID: SessionMessage.ID.create(),
  1244. timestamp: DateTime.makeUnsafe(1),
  1245. reason: "manual",
  1246. })
  1247. yield* events.publish(SessionEvent.Compaction.Ended, {
  1248. sessionID,
  1249. timestamp: DateTime.makeUnsafe(2),
  1250. text: "summary",
  1251. })
  1252. systemBaseline = "Replacement context"
  1253. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  1254. yield* session.resume(sessionID)
  1255. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1256. ["Initial context"],
  1257. ["Replacement context"],
  1258. ])
  1259. yield* replaySessionProjection(sessionID)
  1260. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Third" }), resume: false })
  1261. yield* session.resume(sessionID)
  1262. }),
  1263. )
  1264. it.effect("preserves effective System updates while compaction replacement is blocked", () =>
  1265. Effect.gen(function* () {
  1266. yield* setup
  1267. const session = yield* SessionV2.Service
  1268. const events = yield* EventV2.Service
  1269. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First" }), resume: false })
  1270. requests.length = 0
  1271. response = []
  1272. yield* session.resume(sessionID)
  1273. systemBaseline = "Changed context"
  1274. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second" }), resume: false })
  1275. yield* session.resume(sessionID)
  1276. yield* events.publish(SessionEvent.Compaction.Started, {
  1277. sessionID,
  1278. messageID: SessionMessage.ID.create(),
  1279. timestamp: DateTime.makeUnsafe(1),
  1280. reason: "manual",
  1281. })
  1282. yield* events.publish(SessionEvent.Compaction.Ended, {
  1283. sessionID,
  1284. timestamp: DateTime.makeUnsafe(2),
  1285. text: "summary",
  1286. })
  1287. systemUnavailable = true
  1288. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Third" }), resume: false })
  1289. yield* session.resume(sessionID)
  1290. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Initial context"])
  1291. expect(
  1292. requests
  1293. .at(-1)
  1294. ?.messages.some(
  1295. (message) =>
  1296. message.role === "system" &&
  1297. message.content[0]?.type === "text" &&
  1298. message.content[0].text === "Changed context",
  1299. ),
  1300. ).toBe(true)
  1301. }),
  1302. )
  1303. it.effect("projects reasoning and tool events without executing or continuing tools", () =>
  1304. Effect.gen(function* () {
  1305. yield* setup
  1306. const session = yield* SessionV2.Service
  1307. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Use tools" }), resume: false })
  1308. requests.length = 0
  1309. responses = undefined
  1310. streamGate = undefined
  1311. streamStarted = undefined
  1312. response = [
  1313. LLMEvent.stepStart({ index: 0 }),
  1314. LLMEvent.reasoningStart({ id: "reasoning-1" }),
  1315. LLMEvent.reasoningDelta({ id: "reasoning-1", text: "Think" }),
  1316. LLMEvent.reasoningEnd({ id: "reasoning-1" }),
  1317. LLMEvent.toolInputStart({ id: "call-error", name: "write" }),
  1318. LLMEvent.toolInputDelta({ id: "call-error", name: "write", text: '{"path":"README.md"}' }),
  1319. LLMEvent.toolInputEnd({ id: "call-error", name: "write" }),
  1320. LLMEvent.toolCall({ id: "call-error", name: "write", input: { path: "README.md" }, providerExecuted: true }),
  1321. LLMEvent.toolError({ id: "call-error", name: "write", message: "Denied" }),
  1322. LLMEvent.toolResult({ id: "call-error", name: "write", result: { type: "error", value: "Denied" } }),
  1323. LLMEvent.toolCall({
  1324. id: "call-provider",
  1325. name: "web_search",
  1326. input: { query: "hello" },
  1327. providerExecuted: true,
  1328. providerMetadata: { fake: { source: "provider" } },
  1329. }),
  1330. LLMEvent.toolResult({
  1331. id: "call-provider",
  1332. name: "web_search",
  1333. result: {
  1334. type: "content",
  1335. value: [
  1336. { type: "text", text: "Hello" },
  1337. { type: "media", mediaType: "image/png", data: "data:image/png;base64,aGVsbG8=", filename: "hello.png" },
  1338. ],
  1339. },
  1340. providerExecuted: true,
  1341. providerMetadata: { fake: { source: "provider" } },
  1342. }),
  1343. LLMEvent.stepFinish({
  1344. index: 0,
  1345. reason: "tool-calls",
  1346. usage: {
  1347. inputTokens: 10,
  1348. nonCachedInputTokens: 8,
  1349. outputTokens: 4,
  1350. reasoningTokens: 1,
  1351. cacheReadInputTokens: 2,
  1352. },
  1353. }),
  1354. LLMEvent.finish({ reason: "tool-calls" }),
  1355. ]
  1356. yield* session.resume(sessionID)
  1357. expect(requests).toHaveLength(1)
  1358. expect(requests[0]?.tools.map((tool) => tool.name)).toEqual(["echo", "defect"])
  1359. expect(yield* session.context(sessionID)).toMatchObject([
  1360. { type: "user", text: "Use tools" },
  1361. {
  1362. type: "assistant",
  1363. finish: "tool-calls",
  1364. tokens: { input: 8, output: 3, reasoning: 1, cache: { read: 2, write: 0 } },
  1365. content: [
  1366. { type: "reasoning", id: "reasoning-1", text: "Think" },
  1367. {
  1368. type: "tool",
  1369. id: "call-error",
  1370. name: "write",
  1371. state: {
  1372. status: "error",
  1373. input: { path: "README.md" },
  1374. error: { type: "unknown", message: "Denied" },
  1375. },
  1376. },
  1377. {
  1378. type: "tool",
  1379. id: "call-provider",
  1380. name: "web_search",
  1381. provider: { executed: true, metadata: { fake: { source: "provider" } } },
  1382. state: {
  1383. status: "completed",
  1384. input: { query: "hello" },
  1385. structured: {},
  1386. content: [
  1387. { type: "text", text: "Hello" },
  1388. { type: "file", mime: "image/png", source: { type: "data", data: "aGVsbG8=" }, name: "hello.png" },
  1389. ],
  1390. },
  1391. },
  1392. ],
  1393. },
  1394. ])
  1395. }),
  1396. )
  1397. it.effect("continues with reloaded history after durably settling one local tool call", () =>
  1398. Effect.gen(function* () {
  1399. yield* setup
  1400. const session = yield* SessionV2.Service
  1401. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Echo this" }), resume: false })
  1402. requests.length = 0
  1403. authorizations.length = 0
  1404. executions.length = 0
  1405. streamGate = undefined
  1406. streamStarted = undefined
  1407. responses = [
  1408. [
  1409. LLMEvent.stepStart({ index: 0 }),
  1410. LLMEvent.toolCall({ id: "call-echo", name: "echo", input: { text: "hello" } }),
  1411. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  1412. LLMEvent.finish({ reason: "tool-calls" }),
  1413. ],
  1414. [
  1415. LLMEvent.stepStart({ index: 0 }),
  1416. LLMEvent.textStart({ id: "text-final" }),
  1417. LLMEvent.textDelta({ id: "text-final", text: "Done" }),
  1418. LLMEvent.textEnd({ id: "text-final" }),
  1419. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1420. LLMEvent.finish({ reason: "stop" }),
  1421. ],
  1422. ]
  1423. yield* session.resume(sessionID)
  1424. expect(requests).toHaveLength(2)
  1425. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
  1426. expect(authorizations).toMatchObject([{ sessionID, call: { id: "call-echo", name: "echo" } }])
  1427. expect(executions).toEqual(["hello"])
  1428. expect(yield* session.context(sessionID)).toMatchObject([
  1429. { type: "user", text: "Echo this" },
  1430. {
  1431. type: "assistant",
  1432. finish: "tool-calls",
  1433. content: [
  1434. {
  1435. type: "tool",
  1436. id: "call-echo",
  1437. name: "echo",
  1438. state: {
  1439. status: "completed",
  1440. input: { text: "hello" },
  1441. structured: { text: "hello" },
  1442. content: [{ type: "text", text: "hello" }],
  1443. },
  1444. },
  1445. ],
  1446. },
  1447. { type: "assistant", finish: "stop", content: [{ type: "text", id: "text-final", text: "Done" }] },
  1448. ])
  1449. }),
  1450. )
  1451. it.effect("reloads a model switch before a tool-driven continuation turn", () =>
  1452. Effect.gen(function* () {
  1453. yield* setup
  1454. const session = yield* SessionV2.Service
  1455. const events = yield* EventV2.Service
  1456. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Echo this" }), resume: false })
  1457. requests.length = 0
  1458. responses = [
  1459. [
  1460. LLMEvent.stepStart({ index: 0 }),
  1461. LLMEvent.toolCall({ id: "call-echo", name: "echo", input: { text: "hello" } }),
  1462. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  1463. LLMEvent.finish({ reason: "tool-calls" }),
  1464. ],
  1465. [
  1466. LLMEvent.stepStart({ index: 0 }),
  1467. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1468. LLMEvent.finish({ reason: "stop" }),
  1469. ],
  1470. ]
  1471. toolExecutionGate = yield* Deferred.make<void>()
  1472. toolExecutionsStarted = yield* Deferred.make<void>()
  1473. toolExecutionsReady = 1
  1474. const run = yield* Effect.forkChild(session.resume(sessionID))
  1475. yield* Deferred.await(toolExecutionsStarted)
  1476. yield* events.publish(SessionEvent.ModelSwitched, {
  1477. sessionID,
  1478. messageID: SessionMessage.ID.create(),
  1479. timestamp: DateTime.makeUnsafe(1),
  1480. model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") },
  1481. })
  1482. systemBaseline = "Replacement context"
  1483. yield* Deferred.succeed(toolExecutionGate, undefined)
  1484. yield* Fiber.join(run)
  1485. expect(requests.map((request) => request.model)).toEqual([model, replacementModel])
  1486. expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([
  1487. ["Initial context"],
  1488. ["Replacement context"],
  1489. ])
  1490. }),
  1491. )
  1492. it.effect("restores durable reasoning provider metadata in a second-turn request", () =>
  1493. Effect.gen(function* () {
  1494. yield* setup
  1495. const session = yield* SessionV2.Service
  1496. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Think first" }), resume: false })
  1497. requests.length = 0
  1498. response = [
  1499. LLMEvent.stepStart({ index: 0 }),
  1500. LLMEvent.reasoningStart({ id: "reasoning-anthropic" }),
  1501. LLMEvent.reasoningDelta({ id: "reasoning-anthropic", text: "Signed thought" }),
  1502. LLMEvent.reasoningEnd({ id: "reasoning-anthropic", providerMetadata: { anthropic: { signature: "sig_1" } } }),
  1503. LLMEvent.reasoningStart({
  1504. id: "reasoning-openai",
  1505. providerMetadata: { openai: { itemId: "rs_1", reasoningEncryptedContent: null } },
  1506. }),
  1507. LLMEvent.reasoningDelta({ id: "reasoning-openai", text: "Encrypted thought" }),
  1508. LLMEvent.reasoningEnd({
  1509. id: "reasoning-openai",
  1510. providerMetadata: { openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" } },
  1511. }),
  1512. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1513. LLMEvent.finish({ reason: "stop" }),
  1514. ]
  1515. yield* session.resume(sessionID)
  1516. yield* replaySessionProjection(sessionID)
  1517. expect(yield* session.context(sessionID)).toMatchObject([
  1518. { type: "user", text: "Think first" },
  1519. {
  1520. type: "assistant",
  1521. content: [
  1522. { type: "reasoning", text: "Signed thought", providerMetadata: { anthropic: { signature: "sig_1" } } },
  1523. {
  1524. type: "reasoning",
  1525. text: "Encrypted thought",
  1526. providerMetadata: { openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" } },
  1527. },
  1528. ],
  1529. },
  1530. ])
  1531. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Continue" }), resume: false })
  1532. response = []
  1533. yield* session.resume(sessionID)
  1534. expect(requests[1]?.messages[1]?.content).toEqual([
  1535. { type: "reasoning", text: "Signed thought", providerMetadata: { anthropic: { signature: "sig_1" } } },
  1536. {
  1537. type: "reasoning",
  1538. text: "Encrypted thought",
  1539. providerMetadata: { openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" } },
  1540. },
  1541. ])
  1542. }),
  1543. )
  1544. it.effect("replays durable provider-executed tool results inline in a second-turn request", () =>
  1545. Effect.gen(function* () {
  1546. yield* setup
  1547. const session = yield* SessionV2.Service
  1548. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Search first" }), resume: false })
  1549. requests.length = 0
  1550. response = [
  1551. LLMEvent.stepStart({ index: 0 }),
  1552. LLMEvent.toolCall({
  1553. id: "hosted-search",
  1554. name: "web_search",
  1555. input: { query: "Effect" },
  1556. providerExecuted: true,
  1557. providerMetadata: { openai: { itemId: "hosted-search" } },
  1558. }),
  1559. LLMEvent.toolResult({
  1560. id: "hosted-search",
  1561. name: "web_search",
  1562. result: { type: "json", value: [{ title: "Effect" }] },
  1563. providerExecuted: true,
  1564. providerMetadata: { anthropic: { blockType: "web_search_tool_result" } },
  1565. }),
  1566. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1567. LLMEvent.finish({ reason: "stop" }),
  1568. ]
  1569. yield* session.resume(sessionID)
  1570. yield* replaySessionProjection(sessionID)
  1571. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Continue" }), resume: false })
  1572. response = []
  1573. yield* session.resume(sessionID)
  1574. expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "user"])
  1575. expect(requests[1]?.messages[1]?.content).toMatchObject([
  1576. {
  1577. type: "tool-call",
  1578. id: "hosted-search",
  1579. name: "web_search",
  1580. input: { query: "Effect" },
  1581. providerExecuted: true,
  1582. providerMetadata: { openai: { itemId: "hosted-search" } },
  1583. },
  1584. {
  1585. type: "tool-result",
  1586. id: "hosted-search",
  1587. name: "web_search",
  1588. result: { type: "json", value: [{ title: "Effect" }] },
  1589. providerExecuted: true,
  1590. providerMetadata: { anthropic: { blockType: "web_search_tool_result" } },
  1591. },
  1592. ])
  1593. }),
  1594. )
  1595. it.effect("starts recorded local tools eagerly and awaits settlement before continuing", () =>
  1596. Effect.gen(function* () {
  1597. yield* setup
  1598. const session = yield* SessionV2.Service
  1599. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Echo five times" }), resume: false })
  1600. requests.length = 0
  1601. executions.length = 0
  1602. toolExecutionGate = yield* Deferred.make<void>()
  1603. toolExecutionsStarted = yield* Deferred.make<void>()
  1604. const providerGate = yield* Deferred.make<void>()
  1605. response = []
  1606. responses = undefined
  1607. const initial = Stream.fromIterable([
  1608. LLMEvent.stepStart({ index: 0 }),
  1609. ...Array.from({ length: 5 }, (_, index) =>
  1610. LLMEvent.toolCall({ id: `call-echo-${index}`, name: "echo", input: { text: `${index}` } }),
  1611. ),
  1612. ])
  1613. const final = Stream.fromIterable([
  1614. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  1615. LLMEvent.finish({ reason: "tool-calls" }),
  1616. ])
  1617. streamGate = undefined
  1618. responseStream = Stream.concat(
  1619. initial,
  1620. Stream.fromEffect(Deferred.await(providerGate)).pipe(Stream.flatMap(() => final)),
  1621. )
  1622. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1623. yield* Deferred.await(toolExecutionsStarted)
  1624. expect(executions).toHaveLength(5)
  1625. expect(maxActiveToolExecutions).toBe(5)
  1626. expect(yield* session.context(sessionID)).toMatchObject([
  1627. { type: "user", text: "Echo five times" },
  1628. {
  1629. type: "assistant",
  1630. content: Array.from({ length: 5 }, (_, index) => ({
  1631. type: "tool",
  1632. id: `call-echo-${index}`,
  1633. state: { status: "running", input: { text: `${index}` } },
  1634. })),
  1635. },
  1636. ])
  1637. yield* Deferred.succeed(providerGate, undefined)
  1638. yield* Effect.yieldNow
  1639. expect(requests).toHaveLength(1)
  1640. yield* Deferred.succeed(toolExecutionGate, undefined)
  1641. yield* Fiber.join(run)
  1642. toolExecutionGate = undefined
  1643. toolExecutionsStarted = undefined
  1644. expect(executions).toHaveLength(5)
  1645. expect(maxActiveToolExecutions).toBe(5)
  1646. expect(requests).toHaveLength(2)
  1647. }),
  1648. )
  1649. it.effect("settles repeated provider-local tool call IDs against their owning assistant messages", () =>
  1650. Effect.gen(function* () {
  1651. yield* setup
  1652. const session = yield* SessionV2.Service
  1653. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Echo twice" }), resume: false })
  1654. requests.length = 0
  1655. executions.length = 0
  1656. responses = [
  1657. [
  1658. LLMEvent.stepStart({ index: 0 }),
  1659. LLMEvent.toolCall({ id: "tool_0", name: "echo", input: { text: "first" } }),
  1660. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  1661. LLMEvent.finish({ reason: "tool-calls" }),
  1662. ],
  1663. [
  1664. LLMEvent.stepStart({ index: 0 }),
  1665. LLMEvent.toolCall({ id: "tool_0", name: "echo", input: { text: "second" } }),
  1666. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  1667. LLMEvent.finish({ reason: "tool-calls" }),
  1668. ],
  1669. [],
  1670. ]
  1671. yield* session.resume(sessionID)
  1672. expect(executions).toEqual(["first", "second"])
  1673. expect(requests).toHaveLength(3)
  1674. expect(yield* session.context(sessionID)).toMatchObject([
  1675. { type: "user", text: "Echo twice" },
  1676. {
  1677. type: "assistant",
  1678. content: [
  1679. {
  1680. type: "tool",
  1681. id: "tool_0",
  1682. state: { status: "completed", structured: { text: "first" }, content: [{ type: "text", text: "first" }] },
  1683. },
  1684. ],
  1685. },
  1686. {
  1687. type: "assistant",
  1688. content: [
  1689. {
  1690. type: "tool",
  1691. id: "tool_0",
  1692. state: {
  1693. status: "completed",
  1694. structured: { text: "second" },
  1695. content: [{ type: "text", text: "second" }],
  1696. },
  1697. },
  1698. ],
  1699. },
  1700. ])
  1701. yield* replaySessionProjection(sessionID)
  1702. expect(yield* session.context(sessionID)).toMatchObject([
  1703. { type: "user", text: "Echo twice" },
  1704. {
  1705. type: "assistant",
  1706. content: [
  1707. {
  1708. type: "tool",
  1709. id: "tool_0",
  1710. state: { status: "completed", structured: { text: "first" }, content: [{ type: "text", text: "first" }] },
  1711. },
  1712. ],
  1713. },
  1714. {
  1715. type: "assistant",
  1716. content: [
  1717. {
  1718. type: "tool",
  1719. id: "tool_0",
  1720. state: {
  1721. status: "completed",
  1722. structured: { text: "second" },
  1723. content: [{ type: "text", text: "second" }],
  1724. },
  1725. },
  1726. ],
  1727. },
  1728. ])
  1729. }),
  1730. )
  1731. it.effect("joins concurrent resume calls into one active provider run", () =>
  1732. Effect.gen(function* () {
  1733. yield* setup
  1734. const session = yield* SessionV2.Service
  1735. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Run once" }), resume: false })
  1736. requests.length = 0
  1737. responses = undefined
  1738. response = [
  1739. LLMEvent.stepStart({ index: 0 }),
  1740. LLMEvent.textStart({ id: "text-once" }),
  1741. LLMEvent.textDelta({ id: "text-once", text: "Once" }),
  1742. LLMEvent.textEnd({ id: "text-once" }),
  1743. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1744. LLMEvent.finish({ reason: "stop" }),
  1745. ]
  1746. streamGate = yield* Deferred.make<void>()
  1747. streamStarted = yield* Deferred.make<void>()
  1748. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1749. yield* Deferred.await(streamStarted)
  1750. const second = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1751. yield* Effect.yieldNow
  1752. expect(requests).toHaveLength(1)
  1753. yield* Deferred.succeed(streamGate, undefined)
  1754. yield* Fiber.join(first)
  1755. yield* Fiber.join(second)
  1756. streamGate = undefined
  1757. streamStarted = undefined
  1758. expect(requests).toHaveLength(1)
  1759. expect(yield* session.context(sessionID)).toMatchObject([
  1760. { type: "user", text: "Run once" },
  1761. { type: "assistant", finish: "stop", content: [{ type: "text", id: "text-once", text: "Once" }] },
  1762. ])
  1763. }),
  1764. )
  1765. it.effect("steers an active provider turn with newly recorded prompts", () =>
  1766. Effect.gen(function* () {
  1767. yield* setup
  1768. const session = yield* SessionV2.Service
  1769. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Start working" }), resume: false })
  1770. requests.length = 0
  1771. responses = [
  1772. [
  1773. LLMEvent.stepStart({ index: 0 }),
  1774. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1775. LLMEvent.finish({ reason: "stop" }),
  1776. ],
  1777. [
  1778. LLMEvent.stepStart({ index: 0 }),
  1779. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1780. LLMEvent.finish({ reason: "stop" }),
  1781. ],
  1782. ]
  1783. streamGate = yield* Deferred.make<void>()
  1784. streamStarted = yield* Deferred.make<void>()
  1785. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1786. yield* Deferred.await(streamStarted)
  1787. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Change direction" }) })
  1788. yield* Deferred.succeed(streamGate, undefined)
  1789. yield* Fiber.join(first)
  1790. streamGate = undefined
  1791. streamStarted = undefined
  1792. yield* Effect.yieldNow
  1793. expect(requests).toHaveLength(2)
  1794. expect(userTexts(requests[0]!)).toEqual(["Start working"])
  1795. expect(userTexts(requests[1]!)).toEqual(["Start working", "Change direction"])
  1796. expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([
  1797. "user",
  1798. "assistant",
  1799. "user",
  1800. "assistant",
  1801. ])
  1802. }),
  1803. )
  1804. it.effect("starts queued input after the active activity settles", () =>
  1805. Effect.gen(function* () {
  1806. yield* setup
  1807. const session = yield* SessionV2.Service
  1808. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Start working" }), resume: false })
  1809. requests.length = 0
  1810. responses = [
  1811. [
  1812. LLMEvent.stepStart({ index: 0 }),
  1813. LLMEvent.toolCall({ id: "call-echo", name: "echo", input: { text: "hello" } }),
  1814. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  1815. LLMEvent.finish({ reason: "tool-calls" }),
  1816. ],
  1817. [
  1818. LLMEvent.stepStart({ index: 0 }),
  1819. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1820. LLMEvent.finish({ reason: "stop" }),
  1821. ],
  1822. [
  1823. LLMEvent.stepStart({ index: 0 }),
  1824. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1825. LLMEvent.finish({ reason: "stop" }),
  1826. ],
  1827. ]
  1828. streamGate = yield* Deferred.make<void>()
  1829. streamStarted = yield* Deferred.make<void>()
  1830. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1831. yield* Deferred.await(streamStarted)
  1832. yield* session.prompt({
  1833. sessionID,
  1834. prompt: new Prompt({ text: "Wait until the next activity" }),
  1835. delivery: "queue",
  1836. })
  1837. yield* Deferred.succeed(streamGate, undefined)
  1838. yield* Fiber.join(first)
  1839. streamGate = undefined
  1840. streamStarted = undefined
  1841. expect(requests).toHaveLength(3)
  1842. expect(userTexts(requests[0]!)).toEqual(["Start working"])
  1843. expect(userTexts(requests[1]!)).toEqual(["Start working"])
  1844. expect(userTexts(requests[2]!)).toEqual(["Start working", "Wait until the next activity"])
  1845. }),
  1846. )
  1847. it.effect("preserves durable queued input for a later wake after interruption", () =>
  1848. Effect.gen(function* () {
  1849. yield* setup
  1850. const session = yield* SessionV2.Service
  1851. const { db } = yield* Database.Service
  1852. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Interrupt current work" }), resume: false })
  1853. requests.length = 0
  1854. responses = [
  1855. [],
  1856. [
  1857. LLMEvent.stepStart({ index: 0 }),
  1858. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1859. LLMEvent.finish({ reason: "stop" }),
  1860. ],
  1861. ]
  1862. streamGate = yield* Deferred.make<void>()
  1863. streamStarted = yield* Deferred.make<void>()
  1864. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1865. yield* Deferred.await(streamStarted)
  1866. yield* session.prompt({
  1867. sessionID,
  1868. prompt: new Prompt({ text: "Run after interrupt" }),
  1869. delivery: "queue",
  1870. })
  1871. yield* session.interrupt(sessionID)
  1872. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  1873. expect(requests).toHaveLength(1)
  1874. expect(yield* SessionInput.hasPending(db, sessionID, "queue")).toBe(true)
  1875. const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1876. while (requests.length < 2) yield* Effect.yieldNow
  1877. yield* Deferred.succeed(streamGate, undefined)
  1878. yield* Fiber.join(resumed)
  1879. streamGate = undefined
  1880. streamStarted = undefined
  1881. expect(requests).toHaveLength(2)
  1882. expect(userTexts(requests[0]!)).toEqual(["Interrupt current work"])
  1883. expect(userTexts(requests[1]!)).toEqual(["Interrupt current work", "Run after interrupt"])
  1884. }),
  1885. )
  1886. it.effect("preserves durable steering input for a later resume after interruption", () =>
  1887. Effect.gen(function* () {
  1888. yield* setup
  1889. const session = yield* SessionV2.Service
  1890. const { db } = yield* Database.Service
  1891. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Interrupt current work" }), resume: false })
  1892. requests.length = 0
  1893. responses = [
  1894. [],
  1895. [
  1896. LLMEvent.stepStart({ index: 0 }),
  1897. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1898. LLMEvent.finish({ reason: "stop" }),
  1899. ],
  1900. ]
  1901. streamGate = yield* Deferred.make<void>()
  1902. streamStarted = yield* Deferred.make<void>()
  1903. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1904. yield* Deferred.await(streamStarted)
  1905. yield* session.prompt({
  1906. sessionID,
  1907. prompt: new Prompt({ text: "Steer after interrupt" }),
  1908. })
  1909. yield* session.interrupt(sessionID)
  1910. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  1911. expect(requests).toHaveLength(1)
  1912. expect(yield* SessionInput.hasPending(db, sessionID, "steer")).toBe(true)
  1913. const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1914. while (requests.length < 2) yield* Effect.yieldNow
  1915. yield* Deferred.succeed(streamGate, undefined)
  1916. yield* Fiber.join(resumed)
  1917. streamGate = undefined
  1918. streamStarted = undefined
  1919. expect(requests).toHaveLength(2)
  1920. expect(userTexts(requests[0]!)).toEqual(["Interrupt current work"])
  1921. expect(userTexts(requests[1]!)).toEqual(["Interrupt current work", "Steer after interrupt"])
  1922. }),
  1923. )
  1924. it.effect("runs queued active inputs as separate FIFO activities", () =>
  1925. Effect.gen(function* () {
  1926. yield* setup
  1927. const session = yield* SessionV2.Service
  1928. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Start working" }), resume: false })
  1929. requests.length = 0
  1930. responses = [
  1931. [
  1932. LLMEvent.stepStart({ index: 0 }),
  1933. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1934. LLMEvent.finish({ reason: "stop" }),
  1935. ],
  1936. [
  1937. LLMEvent.stepStart({ index: 0 }),
  1938. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1939. LLMEvent.finish({ reason: "stop" }),
  1940. ],
  1941. [
  1942. LLMEvent.stepStart({ index: 0 }),
  1943. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1944. LLMEvent.finish({ reason: "stop" }),
  1945. ],
  1946. ]
  1947. streamGate = yield* Deferred.make<void>()
  1948. streamStarted = yield* Deferred.make<void>()
  1949. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  1950. yield* Deferred.await(streamStarted)
  1951. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Queue first" }), delivery: "queue" })
  1952. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Queue second" }), delivery: "queue" })
  1953. yield* Deferred.succeed(streamGate, undefined)
  1954. yield* Fiber.join(first)
  1955. streamGate = undefined
  1956. streamStarted = undefined
  1957. expect(requests).toHaveLength(3)
  1958. expect(userTexts(requests[0]!)).toEqual(["Start working"])
  1959. expect(userTexts(requests[1]!)).toEqual(["Start working", "Queue first"])
  1960. expect(userTexts(requests[2]!)).toEqual(["Start working", "Queue first", "Queue second"])
  1961. }),
  1962. )
  1963. it.effect("opens queued input after idle steering activity settles", () =>
  1964. Effect.gen(function* () {
  1965. yield* setup
  1966. const session = yield* SessionV2.Service
  1967. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Start steering activity" }), resume: false })
  1968. yield* session.prompt({
  1969. sessionID,
  1970. prompt: new Prompt({ text: "Queue later activity" }),
  1971. delivery: "queue",
  1972. resume: false,
  1973. })
  1974. requests.length = 0
  1975. responses = [
  1976. [
  1977. LLMEvent.stepStart({ index: 0 }),
  1978. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1979. LLMEvent.finish({ reason: "stop" }),
  1980. ],
  1981. [
  1982. LLMEvent.stepStart({ index: 0 }),
  1983. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  1984. LLMEvent.finish({ reason: "stop" }),
  1985. ],
  1986. ]
  1987. yield* session.resume(sessionID)
  1988. expect(requests).toHaveLength(2)
  1989. expect(userTexts(requests[0]!)).toEqual(["Start steering activity"])
  1990. expect(userTexts(requests[1]!)).toEqual(["Start steering activity", "Queue later activity"])
  1991. }),
  1992. )
  1993. it.effect("coalesces steers into the active queued activity before starting the next queued activity", () =>
  1994. Effect.gen(function* () {
  1995. yield* setup
  1996. const session = yield* SessionV2.Service
  1997. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Start working" }), resume: false })
  1998. requests.length = 0
  1999. responses = [
  2000. [
  2001. LLMEvent.stepStart({ index: 0 }),
  2002. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2003. LLMEvent.finish({ reason: "stop" }),
  2004. ],
  2005. [
  2006. LLMEvent.stepStart({ index: 0 }),
  2007. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2008. LLMEvent.finish({ reason: "stop" }),
  2009. ],
  2010. [
  2011. LLMEvent.stepStart({ index: 0 }),
  2012. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2013. LLMEvent.finish({ reason: "stop" }),
  2014. ],
  2015. [
  2016. LLMEvent.stepStart({ index: 0 }),
  2017. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2018. LLMEvent.finish({ reason: "stop" }),
  2019. ],
  2020. ]
  2021. const firstGate = yield* Deferred.make<void>()
  2022. const secondGate = yield* Deferred.make<void>()
  2023. streamGate = firstGate
  2024. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2025. while (requests.length < 1) yield* Effect.yieldNow
  2026. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Queue first" }), delivery: "queue" })
  2027. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Queue second" }), delivery: "queue" })
  2028. streamGate = secondGate
  2029. yield* Deferred.succeed(firstGate, undefined)
  2030. while (requests.length < 2) yield* Effect.yieldNow
  2031. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Steer first queued activity" }) })
  2032. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Also steer first queued activity" }) })
  2033. yield* Deferred.succeed(secondGate, undefined)
  2034. yield* Fiber.join(first)
  2035. streamGate = undefined
  2036. expect(requests).toHaveLength(4)
  2037. expect(userTexts(requests[0]!)).toEqual(["Start working"])
  2038. expect(userTexts(requests[1]!)).toEqual(["Start working", "Queue first"])
  2039. expect(userTexts(requests[2]!)).toEqual([
  2040. "Start working",
  2041. "Queue first",
  2042. "Steer first queued activity",
  2043. "Also steer first queued activity",
  2044. ])
  2045. expect(userTexts(requests[3]!)).toEqual([
  2046. "Start working",
  2047. "Queue first",
  2048. "Steer first queued activity",
  2049. "Also steer first queued activity",
  2050. "Queue second",
  2051. ])
  2052. }),
  2053. )
  2054. it.effect("coalesces multiple active steering prompts into one continuation turn", () =>
  2055. Effect.gen(function* () {
  2056. yield* setup
  2057. const session = yield* SessionV2.Service
  2058. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Start working" }), resume: false })
  2059. requests.length = 0
  2060. responses = [
  2061. [
  2062. LLMEvent.stepStart({ index: 0 }),
  2063. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2064. LLMEvent.finish({ reason: "stop" }),
  2065. ],
  2066. [
  2067. LLMEvent.stepStart({ index: 0 }),
  2068. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2069. LLMEvent.finish({ reason: "stop" }),
  2070. ],
  2071. ]
  2072. streamGate = yield* Deferred.make<void>()
  2073. streamStarted = yield* Deferred.make<void>()
  2074. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2075. yield* Deferred.await(streamStarted)
  2076. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "First steer" }) })
  2077. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Second steer" }) })
  2078. yield* Deferred.succeed(streamGate, undefined)
  2079. yield* Fiber.join(first)
  2080. streamGate = undefined
  2081. streamStarted = undefined
  2082. yield* Effect.yieldNow
  2083. expect(requests).toHaveLength(2)
  2084. expect(userTexts(requests[1]!)).toEqual(["Start working", "First steer", "Second steer"])
  2085. yield* (yield* SessionRunCoordinator.Service).wake(sessionID)
  2086. yield* Effect.yieldNow
  2087. expect(requests).toHaveLength(2)
  2088. }),
  2089. )
  2090. it.effect("runs steering input accepted while the active provider turn fails", () =>
  2091. Effect.gen(function* () {
  2092. yield* setup
  2093. const session = yield* SessionV2.Service
  2094. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Start working" }), resume: false })
  2095. requests.length = 0
  2096. responses = undefined
  2097. response = []
  2098. streamFailure = providerUnavailable()
  2099. streamGate = yield* Deferred.make<void>()
  2100. streamStarted = yield* Deferred.make<void>()
  2101. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2102. yield* Deferred.await(streamStarted)
  2103. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Recover with this" }) })
  2104. yield* Deferred.succeed(streamGate, undefined)
  2105. expect(yield* Fiber.join(first).pipe(Effect.flip)).toBe(streamFailure)
  2106. streamFailure = undefined
  2107. streamGate = undefined
  2108. streamStarted = undefined
  2109. yield* Effect.yieldNow
  2110. expect(requests).toHaveLength(2)
  2111. expect(userTexts(requests[1]!)).toEqual(["Start working", "Recover with this"])
  2112. }),
  2113. )
  2114. it.effect("durably fails local tools left running by a prior process before continuing", () =>
  2115. Effect.gen(function* () {
  2116. yield* setup
  2117. const session = yield* SessionV2.Service
  2118. const events = yield* EventV2.Service
  2119. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Recover interrupted tool" }), resume: false })
  2120. yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID, Number.MAX_SAFE_INTEGER)
  2121. const assistantMessageID = SessionMessage.ID.create()
  2122. yield* events.publish(SessionEvent.Step.Started, {
  2123. sessionID,
  2124. assistantMessageID,
  2125. timestamp: yield* DateTime.now,
  2126. agent: "build",
  2127. model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
  2128. })
  2129. yield* events.publish(SessionEvent.Tool.Input.Started, {
  2130. sessionID,
  2131. timestamp: yield* DateTime.now,
  2132. assistantMessageID,
  2133. callID: "call-interrupted",
  2134. name: "echo",
  2135. })
  2136. yield* events.publish(SessionEvent.Tool.Input.Ended, {
  2137. sessionID,
  2138. timestamp: yield* DateTime.now,
  2139. assistantMessageID,
  2140. callID: "call-interrupted",
  2141. text: '{"text":"stale"}',
  2142. })
  2143. yield* events.publish(SessionEvent.Tool.Called, {
  2144. sessionID,
  2145. timestamp: yield* DateTime.now,
  2146. assistantMessageID,
  2147. callID: "call-interrupted",
  2148. tool: "echo",
  2149. input: { text: "stale" },
  2150. provider: { executed: false },
  2151. })
  2152. requests.length = 0
  2153. response = []
  2154. yield* session.resume(sessionID)
  2155. expect(requests).toHaveLength(1)
  2156. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
  2157. expect(yield* session.context(sessionID)).toMatchObject([
  2158. { type: "user", text: "Recover interrupted tool" },
  2159. {
  2160. type: "assistant",
  2161. content: [
  2162. {
  2163. type: "tool",
  2164. id: "call-interrupted",
  2165. state: { status: "error", error: { type: "unknown", message: "Tool execution interrupted" } },
  2166. },
  2167. ],
  2168. },
  2169. ])
  2170. }),
  2171. )
  2172. it.effect("durably fails hosted tools left running by a prior process before continuing inline", () =>
  2173. Effect.gen(function* () {
  2174. yield* setup
  2175. const session = yield* SessionV2.Service
  2176. const events = yield* EventV2.Service
  2177. yield* session.prompt({
  2178. sessionID,
  2179. prompt: new Prompt({ text: "Recover interrupted hosted tool" }),
  2180. resume: false,
  2181. })
  2182. yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID, Number.MAX_SAFE_INTEGER)
  2183. const assistantMessageID = SessionMessage.ID.create()
  2184. yield* events.publish(SessionEvent.Step.Started, {
  2185. sessionID,
  2186. assistantMessageID,
  2187. timestamp: yield* DateTime.now,
  2188. agent: "build",
  2189. model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
  2190. })
  2191. yield* events.publish(SessionEvent.Tool.Input.Started, {
  2192. sessionID,
  2193. timestamp: yield* DateTime.now,
  2194. assistantMessageID,
  2195. callID: "call-hosted-interrupted",
  2196. name: "web_search",
  2197. })
  2198. yield* events.publish(SessionEvent.Tool.Input.Ended, {
  2199. sessionID,
  2200. timestamp: yield* DateTime.now,
  2201. assistantMessageID,
  2202. callID: "call-hosted-interrupted",
  2203. text: '{"query":"stale"}',
  2204. })
  2205. yield* events.publish(SessionEvent.Tool.Called, {
  2206. sessionID,
  2207. timestamp: yield* DateTime.now,
  2208. assistantMessageID,
  2209. callID: "call-hosted-interrupted",
  2210. tool: "web_search",
  2211. input: { query: "stale" },
  2212. provider: { executed: true, metadata: { openai: { itemId: "call-hosted-interrupted" } } },
  2213. })
  2214. requests.length = 0
  2215. response = []
  2216. yield* session.resume(sessionID)
  2217. expect(requests).toHaveLength(1)
  2218. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant"])
  2219. expect(requests[0]?.messages[1]?.content).toMatchObject([
  2220. {
  2221. type: "tool-call",
  2222. id: "call-hosted-interrupted",
  2223. providerExecuted: true,
  2224. providerMetadata: { openai: { itemId: "call-hosted-interrupted" } },
  2225. },
  2226. { type: "tool-result", id: "call-hosted-interrupted", providerExecuted: true, result: { type: "error" } },
  2227. ])
  2228. }),
  2229. )
  2230. it.effect("durably fails pending tool input left by a prior process before continuing", () =>
  2231. Effect.gen(function* () {
  2232. yield* setup
  2233. const session = yield* SessionV2.Service
  2234. const events = yield* EventV2.Service
  2235. yield* session.prompt({
  2236. sessionID,
  2237. prompt: new Prompt({ text: "Recover interrupted tool input" }),
  2238. resume: false,
  2239. })
  2240. yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID, Number.MAX_SAFE_INTEGER)
  2241. const assistantMessageID = SessionMessage.ID.create()
  2242. yield* events.publish(SessionEvent.Step.Started, {
  2243. sessionID,
  2244. assistantMessageID,
  2245. timestamp: yield* DateTime.now,
  2246. agent: "build",
  2247. model: { id: ModelV2.ID.make("fake-model"), providerID: ProviderV2.ID.make("fake") },
  2248. })
  2249. yield* events.publish(SessionEvent.Tool.Input.Started, {
  2250. sessionID,
  2251. timestamp: yield* DateTime.now,
  2252. assistantMessageID,
  2253. callID: "call-pending-interrupted",
  2254. name: "echo",
  2255. })
  2256. requests.length = 0
  2257. response = []
  2258. yield* session.resume(sessionID)
  2259. expect(requests).toHaveLength(1)
  2260. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
  2261. expect(yield* session.context(sessionID)).toMatchObject([
  2262. { type: "user", text: "Recover interrupted tool input" },
  2263. { type: "assistant", content: [{ type: "tool", id: "call-pending-interrupted", state: { status: "error" } }] },
  2264. ])
  2265. }),
  2266. )
  2267. it.effect("starts the first queued activity when woken while idle", () =>
  2268. Effect.gen(function* () {
  2269. yield* setup
  2270. const session = yield* SessionV2.Service
  2271. yield* session.prompt({
  2272. sessionID,
  2273. prompt: new Prompt({ text: "Wait for fresh activity" }),
  2274. delivery: "queue",
  2275. resume: false,
  2276. })
  2277. requests.length = 0
  2278. yield* (yield* SessionRunCoordinator.Service).wake(sessionID)
  2279. yield* Effect.yieldNow
  2280. expect(requests).toHaveLength(1)
  2281. expect(userTexts(requests[0]!)).toEqual(["Wait for fresh activity"])
  2282. }),
  2283. )
  2284. it.effect("does not spend one activity step budget across queued activities", () =>
  2285. Effect.gen(function* () {
  2286. yield* setup
  2287. const session = yield* SessionV2.Service
  2288. const queued = Array.from({ length: 26 }, (_, index) => `Queued activity ${index + 1}`)
  2289. for (const text of queued) {
  2290. yield* session.prompt({ sessionID, prompt: new Prompt({ text }), delivery: "queue", resume: false })
  2291. }
  2292. requests.length = 0
  2293. responses = queued.map(() => [
  2294. LLMEvent.stepStart({ index: 0 }),
  2295. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2296. LLMEvent.finish({ reason: "stop" }),
  2297. ])
  2298. yield* session.resume(sessionID)
  2299. expect(requests).toHaveLength(queued.length)
  2300. expect(userTexts(requests.at(-1)!)).toEqual(queued)
  2301. }),
  2302. )
  2303. it.effect("retries inbox input after prompt projection rolls back", () =>
  2304. Effect.gen(function* () {
  2305. yield* setup
  2306. const session = yield* SessionV2.Service
  2307. const events = yield* EventV2.Service
  2308. const defect = new Error("fail after prompt promotion")
  2309. let fail = true
  2310. yield* events.project(SessionEvent.PromptLifecycle.Promoted, () => (fail ? Effect.die(defect) : Effect.void))
  2311. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Recover promoted input" }), resume: false })
  2312. expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect)
  2313. fail = false
  2314. requests.length = 0
  2315. response = [
  2316. LLMEvent.stepStart({ index: 0 }),
  2317. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2318. LLMEvent.finish({ reason: "stop" }),
  2319. ]
  2320. yield* (yield* SessionRunCoordinator.Service).wake(sessionID)
  2321. while (requests.length === 0) yield* Effect.yieldNow
  2322. expect(userTexts(requests[0]!)).toEqual(["Recover promoted input"])
  2323. }),
  2324. )
  2325. it.effect("does not strand a committed promotion when a post-commit listener defects", () =>
  2326. Effect.gen(function* () {
  2327. yield* setup
  2328. const session = yield* SessionV2.Service
  2329. const events = yield* EventV2.Service
  2330. yield* events.listen((event) =>
  2331. event.type === SessionEvent.PromptLifecycle.Promoted.type
  2332. ? Effect.die("fail after prompt promotion commits")
  2333. : Effect.void,
  2334. )
  2335. yield* session.prompt({
  2336. sessionID,
  2337. prompt: new Prompt({ text: "Run committed promotion" }),
  2338. resume: false,
  2339. })
  2340. requests.length = 0
  2341. yield* session.resume(sessionID)
  2342. expect(requests).toHaveLength(1)
  2343. expect(userTexts(requests[0]!)).toEqual(["Run committed promotion"])
  2344. }),
  2345. )
  2346. it.effect("runs different sessions concurrently", () =>
  2347. Effect.gen(function* () {
  2348. yield* setup
  2349. yield* insertSession(otherSessionID)
  2350. const session = yield* SessionV2.Service
  2351. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Run first" }), resume: false })
  2352. yield* session.prompt({ sessionID: otherSessionID, prompt: new Prompt({ text: "Run second" }), resume: false })
  2353. requests.length = 0
  2354. responses = undefined
  2355. response = []
  2356. streamGate = yield* Deferred.make<void>()
  2357. streamStarted = yield* Deferred.make<void>()
  2358. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2359. yield* Deferred.await(streamStarted)
  2360. const second = yield* session.resume(otherSessionID).pipe(Effect.forkChild)
  2361. yield* Effect.yieldNow
  2362. expect(requests).toHaveLength(2)
  2363. yield* Deferred.succeed(streamGate, undefined)
  2364. yield* Fiber.join(first)
  2365. yield* Fiber.join(second)
  2366. streamGate = undefined
  2367. streamStarted = undefined
  2368. }),
  2369. )
  2370. it.effect("fans out one failed run and allows a later retry", () =>
  2371. Effect.gen(function* () {
  2372. yield* setup
  2373. const session = yield* SessionV2.Service
  2374. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Retry after failure" }), resume: false })
  2375. requests.length = 0
  2376. responses = undefined
  2377. response = []
  2378. streamFailure = providerUnavailable()
  2379. streamGate = yield* Deferred.make<void>()
  2380. streamStarted = yield* Deferred.make<void>()
  2381. const first = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2382. yield* Deferred.await(streamStarted)
  2383. const second = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2384. yield* Effect.yieldNow
  2385. expect(requests).toHaveLength(1)
  2386. yield* Deferred.succeed(streamGate, undefined)
  2387. const [firstExit, secondExit] = yield* Effect.all([Fiber.await(first), Fiber.await(second)])
  2388. expect(secondExit).toEqual(firstExit)
  2389. streamFailure = undefined
  2390. streamGate = undefined
  2391. streamStarted = undefined
  2392. yield* session.resume(sessionID)
  2393. expect(requests).toHaveLength(2)
  2394. }),
  2395. )
  2396. it.effect("durably settles local tool failures before continuing", () =>
  2397. Effect.gen(function* () {
  2398. yield* setup
  2399. const session = yield* SessionV2.Service
  2400. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Call missing" }), resume: false })
  2401. requests.length = 0
  2402. responses = [
  2403. [
  2404. LLMEvent.stepStart({ index: 0 }),
  2405. LLMEvent.toolCall({ id: "call-missing", name: "missing", input: {} }),
  2406. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  2407. LLMEvent.finish({ reason: "tool-calls" }),
  2408. ],
  2409. [
  2410. LLMEvent.stepStart({ index: 0 }),
  2411. LLMEvent.textStart({ id: "text-after-error" }),
  2412. LLMEvent.textDelta({ id: "text-after-error", text: "Recovered" }),
  2413. LLMEvent.textEnd({ id: "text-after-error" }),
  2414. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2415. LLMEvent.finish({ reason: "stop" }),
  2416. ],
  2417. ]
  2418. streamGate = undefined
  2419. streamStarted = undefined
  2420. yield* session.resume(sessionID)
  2421. expect(requests).toHaveLength(2)
  2422. expect(yield* session.context(sessionID)).toMatchObject([
  2423. { type: "user", text: "Call missing" },
  2424. {
  2425. type: "assistant",
  2426. content: [
  2427. {
  2428. type: "tool",
  2429. id: "call-missing",
  2430. state: { status: "error", error: { message: "Unknown tool: missing" } },
  2431. },
  2432. ],
  2433. },
  2434. { type: "assistant", finish: "stop", content: [{ type: "text", id: "text-after-error", text: "Recovered" }] },
  2435. ])
  2436. }),
  2437. )
  2438. it.effect("durably settles unexpected local tool defects before continuing", () =>
  2439. Effect.gen(function* () {
  2440. yield* setup
  2441. const session = yield* SessionV2.Service
  2442. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Call defect" }), resume: false })
  2443. requests.length = 0
  2444. responses = [
  2445. [
  2446. LLMEvent.stepStart({ index: 0 }),
  2447. LLMEvent.toolCall({ id: "call-defect", name: "defect", input: {} }),
  2448. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  2449. LLMEvent.finish({ reason: "tool-calls" }),
  2450. ],
  2451. [],
  2452. ]
  2453. yield* session.resume(sessionID)
  2454. expect(requests).toHaveLength(2)
  2455. expect(yield* session.context(sessionID)).toMatchObject([
  2456. { type: "user", text: "Call defect" },
  2457. {
  2458. type: "assistant",
  2459. content: [
  2460. {
  2461. type: "tool",
  2462. id: "call-defect",
  2463. state: { status: "error", error: { message: "unexpected tool defect" } },
  2464. },
  2465. ],
  2466. },
  2467. ])
  2468. }),
  2469. )
  2470. it.effect("interrupts runner continuation when a question is dismissed", () =>
  2471. Effect.gen(function* () {
  2472. yield* setup
  2473. const session = yield* SessionV2.Service
  2474. const registry = yield* ToolRegistry.Service
  2475. const questions = yield* QuestionV2.Service
  2476. const transform = yield* registry.transform()
  2477. yield* transform((editor) =>
  2478. editor.set("question", {
  2479. tool: Tool.make({
  2480. description: "Ask the user",
  2481. parameters: Schema.Struct({}),
  2482. success: Schema.Struct({}),
  2483. }),
  2484. execute: ({ sessionID }) => questions.ask({ sessionID, questions: [] }).pipe(Effect.as({}), Effect.orDie),
  2485. }),
  2486. )
  2487. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Ask then stop" }), resume: false })
  2488. requests.length = 0
  2489. responses = [
  2490. [
  2491. LLMEvent.stepStart({ index: 0 }),
  2492. LLMEvent.toolCall({ id: "call-question", name: "question", input: {} }),
  2493. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  2494. LLMEvent.finish({ reason: "tool-calls" }),
  2495. ],
  2496. [],
  2497. ]
  2498. const run = yield* session.resume(sessionID).pipe(Effect.exit, Effect.forkChild)
  2499. let pending = yield* questions.list()
  2500. while (pending.length === 0) {
  2501. yield* Effect.yieldNow
  2502. pending = yield* questions.list()
  2503. }
  2504. yield* questions.reject(pending[0]!.id)
  2505. const exit = yield* Fiber.join(run)
  2506. expect(exit._tag).toBe("Failure")
  2507. if (exit._tag === "Failure") expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
  2508. expect(requests).toHaveLength(1)
  2509. expect(yield* session.context(sessionID)).toMatchObject([
  2510. { type: "user", text: "Ask then stop" },
  2511. {
  2512. type: "assistant",
  2513. content: [
  2514. {
  2515. type: "tool",
  2516. id: "call-question",
  2517. state: { status: "error", error: { type: "unknown", message: "Tool execution interrupted" } },
  2518. },
  2519. ],
  2520. },
  2521. ])
  2522. }),
  2523. )
  2524. it.effect("awaits started local tools before surfacing provider stream failure", () =>
  2525. Effect.gen(function* () {
  2526. yield* setup
  2527. const session = yield* SessionV2.Service
  2528. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Settle before failing" }), resume: false })
  2529. const failure = providerUnavailable()
  2530. toolExecutionGate = yield* Deferred.make<void>()
  2531. responseStream = Stream.concat(
  2532. Stream.fromIterable([
  2533. LLMEvent.stepStart({ index: 0 }),
  2534. LLMEvent.toolCall({ id: "call-before-failure", name: "echo", input: { text: "settle" } }),
  2535. ]),
  2536. Stream.fail(failure),
  2537. )
  2538. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2539. while (executions.length === 0) yield* Effect.yieldNow
  2540. yield* Effect.yieldNow
  2541. yield* Deferred.succeed(toolExecutionGate, undefined)
  2542. expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure)
  2543. toolExecutionGate = undefined
  2544. expect(yield* session.context(sessionID)).toMatchObject([
  2545. { type: "user", text: "Settle before failing" },
  2546. {
  2547. type: "assistant",
  2548. content: [
  2549. { type: "tool", id: "call-before-failure", state: { status: "completed", structured: { text: "settle" } } },
  2550. ],
  2551. },
  2552. ])
  2553. }),
  2554. )
  2555. it.effect("durably fails blocked local tools when a provider turn is interrupted", () =>
  2556. Effect.gen(function* () {
  2557. yield* setup
  2558. const session = yield* SessionV2.Service
  2559. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Interrupt blocked tool" }), resume: false })
  2560. executions.length = 0
  2561. toolExecutionGate = yield* Deferred.make<void>()
  2562. responseStream = Stream.concat(
  2563. Stream.fromIterable([
  2564. LLMEvent.stepStart({ index: 0 }),
  2565. LLMEvent.toolCall({ id: "call-before-interrupt", name: "echo", input: { text: "blocked" } }),
  2566. ]),
  2567. Stream.never,
  2568. )
  2569. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2570. while (executions.length === 0) yield* Effect.yieldNow
  2571. yield* session.interrupt(sessionID)
  2572. toolExecutionGate = undefined
  2573. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  2574. yield* session.interrupt(sessionID)
  2575. expect(yield* session.context(sessionID)).toMatchObject([
  2576. { type: "user", text: "Interrupt blocked tool" },
  2577. {
  2578. type: "assistant",
  2579. content: [
  2580. {
  2581. type: "tool",
  2582. id: "call-before-interrupt",
  2583. state: { status: "error", error: { type: "unknown", message: "Tool execution interrupted" } },
  2584. },
  2585. ],
  2586. },
  2587. ])
  2588. yield* replaySessionProjection(sessionID)
  2589. expect(yield* session.context(sessionID)).toMatchObject([
  2590. { type: "user", text: "Interrupt blocked tool" },
  2591. { type: "assistant", content: [{ type: "tool", id: "call-before-interrupt", state: { status: "error" } }] },
  2592. ])
  2593. requests.length = 0
  2594. responseStream = undefined
  2595. response = []
  2596. yield* session.resume(sessionID)
  2597. expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
  2598. }),
  2599. )
  2600. it.effect("interrupts a blocked provider turn without local tool activity", () =>
  2601. Effect.gen(function* () {
  2602. yield* setup
  2603. const session = yield* SessionV2.Service
  2604. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Interrupt provider" }), resume: false })
  2605. requests.length = 0
  2606. response = []
  2607. streamGate = yield* Deferred.make<void>()
  2608. streamStarted = yield* Deferred.make<void>()
  2609. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2610. yield* Deferred.await(streamStarted)
  2611. yield* session.interrupt(sessionID)
  2612. const exit = yield* Fiber.await(run)
  2613. streamGate = undefined
  2614. streamStarted = undefined
  2615. expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBeTrue()
  2616. expect(requests).toHaveLength(1)
  2617. yield* session.interrupt(sessionID)
  2618. }),
  2619. )
  2620. it.effect("durably fails blocked local tools when interrupted while awaiting settlement", () =>
  2621. Effect.gen(function* () {
  2622. yield* setup
  2623. const session = yield* SessionV2.Service
  2624. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Interrupt tool settlement" }), resume: false })
  2625. executions.length = 0
  2626. toolExecutionGate = yield* Deferred.make<void>()
  2627. response = [
  2628. LLMEvent.stepStart({ index: 0 }),
  2629. LLMEvent.toolCall({ id: "call-await-interrupt", name: "echo", input: { text: "blocked" } }),
  2630. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  2631. LLMEvent.finish({ reason: "tool-calls" }),
  2632. ]
  2633. const runner = yield* SessionRunner.Service
  2634. const run = yield* runner.run({ sessionID, force: true }).pipe(Effect.forkChild)
  2635. while (executions.length === 0) yield* Effect.yieldNow
  2636. yield* Fiber.interrupt(run)
  2637. toolExecutionGate = undefined
  2638. expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" })
  2639. expect(yield* session.context(sessionID)).toMatchObject([
  2640. { type: "user", text: "Interrupt tool settlement" },
  2641. {
  2642. type: "assistant",
  2643. content: [
  2644. {
  2645. type: "tool",
  2646. id: "call-await-interrupt",
  2647. state: { status: "error", error: { type: "unknown", message: "Tool execution interrupted" } },
  2648. },
  2649. ],
  2650. },
  2651. ])
  2652. }),
  2653. )
  2654. it.effect("fails after the bounded number of local tool continuation steps", () =>
  2655. Effect.gen(function* () {
  2656. yield* setup
  2657. const session = yield* SessionV2.Service
  2658. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Loop forever" }), resume: false })
  2659. requests.length = 0
  2660. authorizations.length = 0
  2661. executions.length = 0
  2662. streamGate = undefined
  2663. streamStarted = undefined
  2664. responses = Array.from({ length: 25 }, (_, index) => [
  2665. LLMEvent.stepStart({ index: 0 }),
  2666. LLMEvent.toolCall({ id: `call-echo-${index}`, name: "echo", input: { text: `${index}` } }),
  2667. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  2668. LLMEvent.finish({ reason: "tool-calls" }),
  2669. ])
  2670. const failure = yield* session.resume(sessionID).pipe(Effect.flip)
  2671. expect(failure).toMatchObject({ _tag: "SessionRunner.StepLimitExceededError", sessionID, limit: 25 })
  2672. expect(requests).toHaveLength(25)
  2673. expect(executions).toHaveLength(25)
  2674. }),
  2675. )
  2676. it.effect("does not restart a capped tool loop for a coalesced stale wake", () =>
  2677. Effect.gen(function* () {
  2678. yield* setup
  2679. const session = yield* SessionV2.Service
  2680. const coordinator = yield* SessionRunCoordinator.Service
  2681. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Loop forever" }), resume: false })
  2682. requests.length = 0
  2683. responses = Array.from({ length: 25 }, (_, index) => [
  2684. LLMEvent.stepStart({ index: 0 }),
  2685. LLMEvent.toolCall({ id: `call-capped-${index}`, name: "echo", input: { text: `${index}` } }),
  2686. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  2687. LLMEvent.finish({ reason: "tool-calls" }),
  2688. ])
  2689. streamGate = yield* Deferred.make<void>()
  2690. streamStarted = yield* Deferred.make<void>()
  2691. const run = yield* session.resume(sessionID).pipe(Effect.forkChild)
  2692. yield* Deferred.await(streamStarted)
  2693. yield* coordinator.wake(sessionID)
  2694. yield* Deferred.succeed(streamGate, undefined)
  2695. expect(yield* Fiber.join(run).pipe(Effect.flip)).toMatchObject({ _tag: "SessionRunner.StepLimitExceededError" })
  2696. streamGate = undefined
  2697. streamStarted = undefined
  2698. yield* Effect.yieldNow
  2699. expect(requests).toHaveLength(25)
  2700. }),
  2701. )
  2702. it.effect("accepts a terminal response on the final bounded provider turn", () =>
  2703. Effect.gen(function* () {
  2704. yield* setup
  2705. const session = yield* SessionV2.Service
  2706. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Finish at the limit" }), resume: false })
  2707. requests.length = 0
  2708. responses = [
  2709. ...Array.from({ length: 24 }, (_, index) => [
  2710. LLMEvent.stepStart({ index: 0 }),
  2711. LLMEvent.toolCall({ id: `call-terminal-${index}`, name: "echo", input: { text: `${index}` } }),
  2712. LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }),
  2713. LLMEvent.finish({ reason: "tool-calls" }),
  2714. ]),
  2715. [
  2716. LLMEvent.stepStart({ index: 0 }),
  2717. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2718. LLMEvent.finish({ reason: "stop" }),
  2719. ],
  2720. ]
  2721. yield* session.resume(sessionID)
  2722. expect(requests).toHaveLength(25)
  2723. }),
  2724. )
  2725. it.effect("projects provider errors as terminal assistant step failures", () =>
  2726. Effect.gen(function* () {
  2727. yield* setup
  2728. const session = yield* SessionV2.Service
  2729. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Fail durably" }), resume: false })
  2730. requests.length = 0
  2731. responses = undefined
  2732. streamGate = undefined
  2733. streamStarted = undefined
  2734. response = [LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "Provider unavailable" })]
  2735. yield* session.resume(sessionID)
  2736. expect(requests).toHaveLength(1)
  2737. expect(yield* session.context(sessionID)).toMatchObject([
  2738. { type: "user", text: "Fail durably" },
  2739. { type: "assistant", finish: "error", error: { type: "unknown", message: "Provider unavailable" } },
  2740. ])
  2741. }),
  2742. )
  2743. it.effect("projects provider errors emitted before assistant step start", () =>
  2744. Effect.gen(function* () {
  2745. yield* setup
  2746. const session = yield* SessionV2.Service
  2747. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Fail before step" }), resume: false })
  2748. requests.length = 0
  2749. response = [LLMEvent.providerError({ message: "Provider unavailable" })]
  2750. yield* session.resume(sessionID)
  2751. expect(requests).toHaveLength(1)
  2752. expect(yield* session.context(sessionID)).toMatchObject([
  2753. { type: "user", text: "Fail before step" },
  2754. { type: "assistant", finish: "error", error: { type: "unknown", message: "Provider unavailable" } },
  2755. ])
  2756. }),
  2757. )
  2758. it.effect("projects raw provider stream failures as terminal assistant step failures", () =>
  2759. Effect.gen(function* () {
  2760. yield* setup
  2761. const session = yield* SessionV2.Service
  2762. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Fail raw stream durably" }), resume: false })
  2763. const failure = providerUnavailable()
  2764. responseStream = Stream.fail(failure)
  2765. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  2766. yield* replaySessionProjection(sessionID)
  2767. expect(yield* session.context(sessionID)).toMatchObject([
  2768. { type: "user", text: "Fail raw stream durably" },
  2769. { type: "assistant", finish: "error", error: { type: "unknown", message: "Provider unavailable" } },
  2770. ])
  2771. }),
  2772. )
  2773. it.effect("does not continue automatically after a provider error follows a local tool call", () =>
  2774. Effect.gen(function* () {
  2775. yield* setup
  2776. const session = yield* SessionV2.Service
  2777. yield* session.prompt({
  2778. sessionID,
  2779. prompt: new Prompt({ text: "Do not continue failed provider" }),
  2780. resume: false,
  2781. })
  2782. requests.length = 0
  2783. const executionCount = executions.length
  2784. response = [
  2785. LLMEvent.stepStart({ index: 0 }),
  2786. LLMEvent.toolCall({ id: "call-before-provider-error", name: "echo", input: { text: "settled" } }),
  2787. LLMEvent.providerError({ message: "Provider unavailable" }),
  2788. ]
  2789. yield* session.resume(sessionID)
  2790. expect(requests).toHaveLength(1)
  2791. expect(executions.slice(executionCount)).toEqual(["settled"])
  2792. }),
  2793. )
  2794. it.effect("durably fails a hosted tool when its provider errors before returning a result", () =>
  2795. Effect.gen(function* () {
  2796. yield* setup
  2797. const session = yield* SessionV2.Service
  2798. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Fail hosted tool durably" }), resume: false })
  2799. requests.length = 0
  2800. response = [
  2801. LLMEvent.stepStart({ index: 0 }),
  2802. LLMEvent.toolCall({
  2803. id: "call-hosted-provider-error",
  2804. name: "web_search",
  2805. input: { query: "effect" },
  2806. providerExecuted: true,
  2807. }),
  2808. LLMEvent.providerError({ message: "Provider unavailable" }),
  2809. ]
  2810. yield* session.resume(sessionID)
  2811. expect(requests).toHaveLength(1)
  2812. expect(yield* session.context(sessionID)).toMatchObject([
  2813. { type: "user", text: "Fail hosted tool durably" },
  2814. {
  2815. type: "assistant",
  2816. content: [{ type: "tool", id: "call-hosted-provider-error", state: { status: "error" } }],
  2817. },
  2818. ])
  2819. }),
  2820. )
  2821. it.effect("durably fails a hosted tool left unresolved at normal provider EOF", () =>
  2822. Effect.gen(function* () {
  2823. yield* setup
  2824. const session = yield* SessionV2.Service
  2825. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Fail hosted tool at EOF" }), resume: false })
  2826. response = [
  2827. LLMEvent.stepStart({ index: 0 }),
  2828. LLMEvent.toolCall({
  2829. id: "call-hosted-eof",
  2830. name: "web_search",
  2831. input: { query: "effect" },
  2832. providerExecuted: true,
  2833. }),
  2834. ]
  2835. yield* session.resume(sessionID)
  2836. yield* replaySessionProjection(sessionID)
  2837. expect(yield* session.context(sessionID)).toMatchObject([
  2838. { type: "user", text: "Fail hosted tool at EOF" },
  2839. { type: "assistant", content: [{ type: "tool", id: "call-hosted-eof", state: { status: "error" } }] },
  2840. ])
  2841. }),
  2842. )
  2843. it.effect("durably fails a hosted tool left unresolved by a raw provider stream failure", () =>
  2844. Effect.gen(function* () {
  2845. yield* setup
  2846. const session = yield* SessionV2.Service
  2847. yield* session.prompt({
  2848. sessionID,
  2849. prompt: new Prompt({ text: "Fail hosted tool on raw failure" }),
  2850. resume: false,
  2851. })
  2852. const failure = providerUnavailable()
  2853. responseStream = Stream.concat(
  2854. Stream.fromIterable([
  2855. LLMEvent.stepStart({ index: 0 }),
  2856. LLMEvent.toolCall({
  2857. id: "call-hosted-raw-failure",
  2858. name: "web_search",
  2859. input: { query: "effect" },
  2860. providerExecuted: true,
  2861. }),
  2862. ]),
  2863. Stream.fail(failure),
  2864. )
  2865. expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure)
  2866. yield* replaySessionProjection(sessionID)
  2867. expect(yield* session.context(sessionID)).toMatchObject([
  2868. { type: "user", text: "Fail hosted tool on raw failure" },
  2869. {
  2870. type: "assistant",
  2871. finish: "error",
  2872. error: { type: "unknown", message: "Provider unavailable" },
  2873. content: [{ type: "tool", id: "call-hosted-raw-failure", state: { status: "error" } }],
  2874. },
  2875. ])
  2876. }),
  2877. )
  2878. it.effect("keeps interleaved assistant text blocks separate", () =>
  2879. Effect.gen(function* () {
  2880. yield* setup
  2881. const session = yield* SessionV2.Service
  2882. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Two blocks" }), resume: false })
  2883. responses = undefined
  2884. streamGate = undefined
  2885. streamStarted = undefined
  2886. response = [
  2887. LLMEvent.stepStart({ index: 0 }),
  2888. LLMEvent.textStart({ id: "text-1" }),
  2889. LLMEvent.textStart({ id: "text-2" }),
  2890. LLMEvent.textDelta({ id: "text-1", text: "First" }),
  2891. LLMEvent.textDelta({ id: "text-2", text: "Second" }),
  2892. LLMEvent.textEnd({ id: "text-1" }),
  2893. LLMEvent.textEnd({ id: "text-2" }),
  2894. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  2895. LLMEvent.finish({ reason: "stop" }),
  2896. ]
  2897. yield* session.resume(sessionID)
  2898. expect(yield* session.context(sessionID)).toMatchObject([
  2899. { type: "user", text: "Two blocks" },
  2900. {
  2901. type: "assistant",
  2902. content: [
  2903. { type: "text", id: "text-1", text: "First" },
  2904. { type: "text", id: "text-2", text: "Second" },
  2905. ],
  2906. },
  2907. ])
  2908. }),
  2909. )
  2910. for (const kind of fragmentKinds) {
  2911. it.effect(`broadcasts provider ${kind} deltas without storing projection rewrites`, () =>
  2912. verifyEphemeralDeltas(kind),
  2913. )
  2914. it.effect(`durably closes partial ${kind} when the provider stream fails`, () => verifyPartialFlushOnFailure(kind))
  2915. it.effect(`durably closes partial ${kind} when the provider stream is interrupted`, () =>
  2916. verifyPartialFlushOnInterruption(kind),
  2917. )
  2918. }
  2919. it.effect("rejects duplicate streamed text starts", () =>
  2920. Effect.gen(function* () {
  2921. yield* setup
  2922. const session = yield* SessionV2.Service
  2923. responses = undefined
  2924. streamGate = undefined
  2925. streamStarted = undefined
  2926. response = [LLMEvent.textStart({ id: "text-1" }), LLMEvent.textStart({ id: "text-1" })]
  2927. expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe(
  2928. "Duplicate text start: text-1",
  2929. )
  2930. }),
  2931. )
  2932. it.effect("transitions streamed raw tool input to parsed called input", () =>
  2933. Effect.gen(function* () {
  2934. yield* setup
  2935. const session = yield* SessionV2.Service
  2936. yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Call provider tool" }), resume: false })
  2937. responses = undefined
  2938. streamGate = undefined
  2939. streamStarted = undefined
  2940. response = [
  2941. LLMEvent.stepStart({ index: 0 }),
  2942. LLMEvent.toolInputStart({ id: "call-parsed", name: "web_search" }),
  2943. LLMEvent.toolInputDelta({ id: "call-parsed", name: "web_search", text: '{"query":"hello"}' }),
  2944. LLMEvent.toolInputEnd({ id: "call-parsed", name: "web_search" }),
  2945. LLMEvent.toolCall({ id: "call-parsed", name: "web_search", input: { query: "hello" }, providerExecuted: true }),
  2946. ]
  2947. yield* session.resume(sessionID)
  2948. expect(yield* session.context(sessionID)).toMatchObject([
  2949. { type: "user", text: "Call provider tool" },
  2950. {
  2951. type: "assistant",
  2952. content: [{ type: "tool", id: "call-parsed", state: { status: "error", input: { query: "hello" } } }],
  2953. },
  2954. ])
  2955. }),
  2956. )
  2957. it.effect("rejects malformed streamed tool input ordering", () =>
  2958. Effect.gen(function* () {
  2959. yield* setup
  2960. const session = yield* SessionV2.Service
  2961. responses = undefined
  2962. streamGate = undefined
  2963. streamStarted = undefined
  2964. response = [LLMEvent.toolInputDelta({ id: "call-1", name: "read", text: "{}" })]
  2965. expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe(
  2966. "Tool input delta before start: call-1",
  2967. )
  2968. }),
  2969. )
  2970. })