session-runner.test.ts 99 KB

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