session-runner.test.ts 124 KB

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