1
0

stream-v2.transport.test.ts 115 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355235623572358235923602361236223632364236523662367236823692370237123722373237423752376237723782379238023812382238323842385238623872388238923902391239223932394239523962397239823992400240124022403240424052406240724082409241024112412241324142415241624172418241924202421242224232424242524262427242824292430243124322433243424352436243724382439244024412442244324442445244624472448244924502451245224532454245524562457245824592460246124622463246424652466246724682469247024712472247324742475247624772478247924802481248224832484248524862487248824892490249124922493249424952496249724982499250025012502250325042505250625072508250925102511251225132514251525162517251825192520252125222523252425252526252725282529253025312532253325342535253625372538253925402541254225432544254525462547254825492550255125522553255425552556255725582559256025612562256325642565256625672568256925702571257225732574257525762577257825792580258125822583258425852586258725882589259025912592259325942595259625972598259926002601260226032604260526062607260826092610261126122613261426152616261726182619262026212622262326242625262626272628262926302631263226332634263526362637263826392640264126422643264426452646264726482649265026512652265326542655265626572658265926602661266226632664266526662667266826692670267126722673267426752676267726782679268026812682268326842685268626872688268926902691269226932694269526962697269826992700270127022703270427052706270727082709271027112712271327142715271627172718271927202721272227232724272527262727272827292730273127322733273427352736273727382739274027412742274327442745274627472748274927502751275227532754275527562757275827592760276127622763276427652766276727682769277027712772277327742775277627772778277927802781278227832784278527862787278827892790279127922793279427952796279727982799280028012802280328042805280628072808280928102811281228132814281528162817281828192820282128222823282428252826282728282829283028312832283328342835283628372838283928402841284228432844284528462847284828492850285128522853285428552856285728582859286028612862286328642865286628672868286928702871287228732874287528762877287828792880288128822883288428852886288728882889289028912892289328942895289628972898289929002901290229032904290529062907290829092910291129122913291429152916291729182919292029212922292329242925292629272928292929302931293229332934293529362937293829392940294129422943294429452946294729482949295029512952295329542955295629572958295929602961296229632964296529662967296829692970297129722973297429752976297729782979298029812982298329842985298629872988298929902991299229932994299529962997299829993000300130023003300430053006300730083009301030113012301330143015301630173018301930203021302230233024302530263027302830293030303130323033303430353036303730383039304030413042304330443045304630473048304930503051305230533054305530563057305830593060306130623063306430653066306730683069307030713072307330743075307630773078307930803081308230833084308530863087308830893090309130923093309430953096309730983099310031013102310331043105310631073108310931103111311231133114311531163117311831193120312131223123312431253126312731283129313031313132313331343135313631373138313931403141314231433144314531463147314831493150315131523153315431553156315731583159316031613162316331643165316631673168316931703171317231733174317531763177317831793180318131823183318431853186318731883189319031913192319331943195319631973198319932003201320232033204320532063207320832093210321132123213321432153216321732183219322032213222322332243225322632273228322932303231323232333234323532363237323832393240324132423243324432453246324732483249325032513252325332543255325632573258325932603261326232633264326532663267326832693270327132723273327432753276327732783279328032813282328332843285328632873288328932903291329232933294329532963297329832993300330133023303330433053306330733083309331033113312331333143315331633173318331933203321332233233324332533263327332833293330333133323333333433353336333733383339334033413342334333443345334633473348334933503351335233533354335533563357335833593360336133623363336433653366336733683369337033713372337333743375337633773378337933803381338233833384338533863387338833893390339133923393339433953396339733983399340034013402340334043405340634073408340934103411341234133414341534163417341834193420342134223423342434253426342734283429343034313432343334343435343634373438343934403441344234433444344534463447344834493450345134523453345434553456345734583459346034613462346334643465346634673468346934703471347234733474347534763477347834793480348134823483348434853486348734883489349034913492349334943495349634973498349935003501350235033504350535063507350835093510351135123513351435153516351735183519352035213522352335243525352635273528352935303531353235333534353535363537353835393540354135423543354435453546354735483549355035513552355335543555355635573558355935603561356235633564356535663567356835693570357135723573357435753576357735783579358035813582358335843585358635873588358935903591359235933594359535963597359835993600360136023603360436053606360736083609361036113612361336143615361636173618361936203621362236233624362536263627362836293630363136323633363436353636363736383639364036413642364336443645364636473648364936503651365236533654365536563657365836593660366136623663366436653666366736683669367036713672367336743675367636773678367936803681368236833684368536863687368836893690369136923693369436953696369736983699370037013702370337043705370637073708370937103711371237133714371537163717371837193720372137223723372437253726372737283729373037313732373337343735373637373738373937403741374237433744374537463747374837493750375137523753375437553756375737583759376037613762376337643765376637673768376937703771377237733774377537763777377837793780378137823783378437853786378737883789379037913792379337943795379637973798379938003801380238033804380538063807380838093810381138123813381438153816381738183819382038213822382338243825382638273828382938303831383238333834383538363837383838393840384138423843
  1. import { afterEach, describe, expect, mock, spyOn, test } from "bun:test"
  2. import fs from "fs/promises"
  3. import path from "path"
  4. import { pathToFileURL } from "node:url"
  5. import {
  6. OpenCode,
  7. type EventSubscribeOutput,
  8. type FormInfo,
  9. type MessageListOutput,
  10. type OpenCodeClient,
  11. type PermissionRequest,
  12. } from "@opencode-ai/client/promise"
  13. import { createSessionTransport } from "../../src/mini/stream-v2.transport"
  14. import type { StreamCommit } from "../../src/mini/types"
  15. import { createFooterApiFixture } from "./fixture/footer-api"
  16. import { canonicalToolPart } from "./fixture/tool-part"
  17. import { tmpdir } from "../fixture/fixture"
  18. type RunV2Event = EventSubscribeOutput
  19. function feed() {
  20. const values: RunV2Event[] = []
  21. let closed = false
  22. let wake: (() => void) | undefined
  23. const stream = (async function* (): AsyncGenerator<RunV2Event, void, unknown> {
  24. while (!closed || values.length > 0) {
  25. if (values.length === 0) {
  26. await new Promise<void>((resolve) => {
  27. wake = resolve
  28. })
  29. continue
  30. }
  31. const value = values.shift()
  32. if (value) yield value
  33. }
  34. })()
  35. return {
  36. stream,
  37. push(value: RunV2Event) {
  38. values.push(value)
  39. wake?.()
  40. wake = undefined
  41. },
  42. close() {
  43. closed = true
  44. wake?.()
  45. wake = undefined
  46. },
  47. }
  48. }
  49. function ok<T>(data: T) {
  50. return Promise.resolve(data)
  51. }
  52. function defer<T = void>() {
  53. let resolve!: (value: T | PromiseLike<T>) => void
  54. const promise = new Promise<T>((done) => {
  55. resolve = done
  56. })
  57. return { promise, resolve }
  58. }
  59. function connected(id = "evt_connected") {
  60. return { id, type: "server.connected", data: {} } satisfies RunV2Event
  61. }
  62. function durable(sessionID: string, seq?: number): { aggregateID: string; seq: number; version: 1 }
  63. function durable<const Version extends 1 | 2>(
  64. sessionID: string,
  65. seq: number,
  66. version: Version,
  67. ): { aggregateID: string; seq: number; version: Version }
  68. function durable(sessionID: string, seq = 0, version: 1 | 2 = 1) {
  69. return { aggregateID: sessionID, seq, version }
  70. }
  71. function promptAdmission(input: Parameters<OpenCodeClient["session"]["prompt"]>[0], sessionID = "ses_1") {
  72. return {
  73. id: input.id ?? "msg_prompt",
  74. sessionID,
  75. type: "user" as const,
  76. data: {
  77. text: input.text,
  78. files: input.files,
  79. agents: input.agents,
  80. metadata: input.metadata,
  81. },
  82. delivery: input.delivery ?? ("steer" as const),
  83. timeCreated: 2,
  84. }
  85. }
  86. function footer() {
  87. return createFooterApiFixture()
  88. }
  89. type SessionMessages = MessageListOutput["data"]
  90. function compaction(status: "running" | "completed", summary: string): SessionMessages[number] {
  91. const message = {
  92. id: "msg_compaction",
  93. type: "compaction" as const,
  94. reason: "auto" as const,
  95. summary,
  96. recent: "",
  97. time: { created: 1 },
  98. }
  99. if (status === "running") return { ...message, status }
  100. return { ...message, status }
  101. }
  102. function form(id: string, sessionID: string, title = id): FormInfo {
  103. return {
  104. id,
  105. sessionID,
  106. title,
  107. fields: [
  108. {
  109. key: "answer",
  110. type: "string",
  111. options: [{ value: "yes", label: "Yes" }],
  112. custom: true,
  113. },
  114. ],
  115. }
  116. }
  117. function eventForm(info: FormInfo): Extract<RunV2Event, { type: "form.created" }>["data"]["form"] {
  118. return info as Extract<RunV2Event, { type: "form.created" }>["data"]["form"]
  119. }
  120. function sdk(input: {
  121. streams: ReturnType<typeof feed>[]
  122. active?: () => Record<string, { type: "running" }>
  123. messages?: Record<string, SessionMessages>
  124. sessions?: Array<{ id: string; parentID?: string; title?: string; agent?: string; time: { updated: number } }>
  125. forms?: Record<string, FormInfo[]>
  126. globals?: FormInfo[]
  127. globalLocation?: { directory: string; workspaceID?: string }
  128. permissions?: Record<string, PermissionRequest[]>
  129. pending?: Record<string, Awaited<ReturnType<OpenCodeClient["session"]["pending"]["list"]>>>
  130. wait?: () => Promise<void>
  131. }) {
  132. const client = OpenCode.make({ baseUrl: "https://opencode.test" })
  133. let subscription = 0
  134. spyOn(client.event, "subscribe").mockImplementation(() => input.streams[subscription++]?.stream ?? feed().stream)
  135. spyOn(client.message, "list").mockImplementation((request) =>
  136. ok({
  137. data: input.messages?.[request.sessionID] ?? [
  138. {
  139. id: "msg_old",
  140. type: "user" as const,
  141. text: "previous prompt",
  142. files: [],
  143. agents: [],
  144. time: { created: 1 },
  145. },
  146. ],
  147. cursor: {},
  148. }),
  149. )
  150. spyOn(client.permission, "list").mockImplementation((request) => ok(input.permissions?.[request.sessionID] ?? []))
  151. spyOn(client.form, "list").mockImplementation((request) => ok(input.forms?.[request.sessionID] ?? []))
  152. spyOn(client.form.request, "list").mockImplementation(() =>
  153. ok({
  154. location: {
  155. directory: input.globalLocation?.directory ?? "/tmp",
  156. workspaceID: input.globalLocation?.workspaceID,
  157. project: {
  158. id: "proj_1",
  159. directory: input.globalLocation?.directory ?? "/tmp",
  160. canonical: input.globalLocation?.directory ?? "/tmp",
  161. },
  162. },
  163. data: input.globals ?? [],
  164. }),
  165. )
  166. spyOn(client.session, "active").mockImplementation(() => ok(input.active?.() ?? {}))
  167. spyOn(client.session.pending, "list").mockImplementation((request) => ok(input.pending?.[request.sessionID] ?? []))
  168. spyOn(client.session, "wait").mockImplementation(() => input.wait?.() ?? ok(undefined))
  169. spyOn(client.session, "message").mockImplementation((request) => {
  170. const message = input.messages?.[request.sessionID]?.find((item) => item.id === request.messageID)
  171. return message ? (ok(message) as never) : Promise.reject(new Error(`message not found: ${request.messageID}`))
  172. })
  173. spyOn(client.session, "switchAgent").mockImplementation(() => ok(undefined))
  174. spyOn(client.session, "switchModel").mockImplementation(() => ok(undefined))
  175. // The generated methods have conditional return types for throwOnError; the
  176. // minimal shapes below are enough for family discovery and model fallback.
  177. spyOn(client.session, "list").mockImplementation((request) => {
  178. const parentID = request?.parentID
  179. return ok({
  180. location: { directory: "/tmp", project: { id: "proj_1", directory: "/tmp" } },
  181. data:
  182. input.sessions?.filter((session) =>
  183. parentID === undefined
  184. ? true
  185. : parentID === null
  186. ? session.parentID === undefined
  187. : session.parentID === parentID,
  188. ) ?? [],
  189. }) as never
  190. })
  191. spyOn(client.model, "default").mockImplementation(
  192. () =>
  193. ok({
  194. location: { directory: "/tmp", project: { id: "proj_1", directory: "/tmp" } },
  195. data: undefined,
  196. }) as never,
  197. )
  198. return client
  199. }
  200. afterEach(() => {
  201. mock.restore()
  202. })
  203. describe("V2 mini transport", () => {
  204. test("renders projected compactions as labeled transcript boundaries", async () => {
  205. const events = feed()
  206. events.push(connected())
  207. const ui = footer()
  208. const transport = await createSessionTransport({
  209. sdk: sdk({
  210. streams: [events],
  211. messages: {
  212. ses_1: [compaction("completed", "## Transport")],
  213. },
  214. }),
  215. sessionID: "ses_1",
  216. thinking: false,
  217. replay: true,
  218. footer: ui.api,
  219. })
  220. expect(ui.commits).toMatchObject([
  221. { text: "Compaction", compaction: true, messageID: "msg_compaction" },
  222. { text: "## Transport", phase: "progress", messageID: "msg_compaction" },
  223. { text: "", phase: "final", messageID: "msg_compaction" },
  224. ])
  225. await transport.close()
  226. })
  227. test("shows an active compaction boundary before live summary output without history replay", async () => {
  228. const events = feed()
  229. events.push(connected())
  230. const ui = footer()
  231. const transport = await createSessionTransport({
  232. sdk: sdk({
  233. streams: [events],
  234. active: () => ({ ses_1: { type: "running" } }),
  235. messages: {
  236. ses_1: [compaction("running", "")],
  237. },
  238. }),
  239. sessionID: "ses_1",
  240. thinking: false,
  241. footer: ui.api,
  242. })
  243. events.push({
  244. id: "evt_compaction_delta",
  245. created: 2,
  246. type: "session.compaction.delta",
  247. data: { sessionID: "ses_1", text: "Transport" },
  248. })
  249. events.push({
  250. id: "evt_compaction_ended",
  251. created: 3,
  252. type: "session.compaction.ended",
  253. durable: durable("ses_1", 3),
  254. data: { sessionID: "ses_1", reason: "auto", text: "Transport", recent: "" },
  255. })
  256. while (!ui.commits.some((commit) => commit.phase === "final")) await Bun.sleep(0)
  257. expect(ui.commits).toMatchObject([
  258. { text: "Compaction", compaction: true, messageID: "msg_compaction" },
  259. { text: "Transport", phase: "progress", messageID: "msg_compaction" },
  260. { text: "", phase: "final", messageID: "msg_compaction" },
  261. ])
  262. await transport.close()
  263. })
  264. test("reports session title changes", async () => {
  265. const events = feed()
  266. events.push(connected())
  267. const titles: string[] = []
  268. const transport = await createSessionTransport({
  269. sdk: sdk({ streams: [events] }),
  270. sessionID: "ses_1",
  271. thinking: false,
  272. footer: footer().api,
  273. onSessionTitle: (title) => titles.push(title),
  274. })
  275. events.push({
  276. id: "evt_renamed",
  277. created: 1,
  278. type: "session.renamed",
  279. durable: durable("ses_1", 1),
  280. data: { sessionID: "ses_1", title: "Greeting" },
  281. })
  282. while (titles.length === 0) await Bun.sleep(0)
  283. expect(titles).toEqual(["Greeting"])
  284. await transport.close()
  285. })
  286. test("formats footer usage with compact tokens and context percentage", async () => {
  287. const events = feed()
  288. events.push(connected())
  289. const ui = footer()
  290. const transport = await createSessionTransport({
  291. sdk: sdk({ streams: [events] }),
  292. sessionID: "ses_1",
  293. thinking: false,
  294. footer: ui.api,
  295. contextLimit: (model) => (model.providerID === "test" && model.modelID === "model" ? 160_000 : undefined),
  296. })
  297. events.push({
  298. id: "evt_step_started",
  299. created: 1,
  300. type: "session.step.started",
  301. durable: durable("ses_1", 1),
  302. data: {
  303. sessionID: "ses_1",
  304. assistantMessageID: "msg_assistant",
  305. agent: "build",
  306. model: { providerID: "test", id: "model" },
  307. },
  308. })
  309. events.push({
  310. id: "evt_step_ended",
  311. created: 2,
  312. type: "session.step.ended",
  313. durable: durable("ses_1", 2),
  314. data: {
  315. sessionID: "ses_1",
  316. assistantMessageID: "msg_assistant",
  317. finish: "stop",
  318. cost: 0,
  319. tokens: { input: 7_000, output: 500, reasoning: 8, cache: { read: 0, write: 0 } },
  320. },
  321. })
  322. while (!ui.events.some((event) => event.type === "stream.patch" && event.patch.usage)) await Bun.sleep(0)
  323. expect(ui.events).toContainEqual({ type: "stream.patch", patch: { usage: "7.5K (5%)" } })
  324. await transport.close()
  325. })
  326. test("recursively hydrates blockers for direct and transitive descendants", async () => {
  327. const events = feed()
  328. events.push(connected())
  329. const client = sdk({
  330. streams: [events],
  331. sessions: [
  332. { id: "ses_child", parentID: "ses_1", title: "Child", time: { updated: 2 } },
  333. { id: "ses_grandchild", parentID: "ses_child", title: "Grandchild", time: { updated: 1 } },
  334. ],
  335. forms: {
  336. ses_child: [form("frm_child", "ses_child")],
  337. ses_grandchild: [form("frm_grandchild", "ses_grandchild")],
  338. },
  339. })
  340. const ui = footer()
  341. const transport = await createSessionTransport({
  342. sdk: client,
  343. sessionID: "ses_1",
  344. thinking: false,
  345. footer: ui.api,
  346. })
  347. const snapshots = ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  348. expect(snapshots.at(-1)?.tabs.map((item) => item.sessionID)).toEqual(["ses_child", "ses_grandchild"])
  349. expect(snapshots.at(-1)?.forms.map((item) => item.id)).toEqual(["frm_child", "frm_grandchild"])
  350. expect(
  351. ui.events.find(
  352. (event) => event.type === "stream.view" && event.view.type === "form" && event.view.request.id === "frm_child",
  353. ),
  354. ).toMatchObject({
  355. type: "stream.view",
  356. view: { type: "form", request: { id: "frm_child", sessionID: "ses_child" } },
  357. })
  358. transport.settleForm?.("ses_child", "frm_child")
  359. expect(ui.events.at(-1)).toMatchObject({
  360. type: "stream.view",
  361. view: { type: "form", request: { id: "frm_grandchild", sessionID: "ses_grandchild" } },
  362. })
  363. await transport.close()
  364. })
  365. test("resolves a pre-existing child permission from its exact source message at startup", async () => {
  366. const events = feed()
  367. events.push(connected())
  368. const sourceMessage = {
  369. id: "msg_child_source",
  370. type: "assistant" as const,
  371. agent: "build",
  372. model: { providerID: "test", id: "model" },
  373. content: [
  374. canonicalToolPart(
  375. "shell",
  376. {
  377. status: "running" as const,
  378. input: { command: "git status --short" },
  379. metadata: {},
  380. },
  381. "call_child_source",
  382. ),
  383. ],
  384. time: { created: 1 },
  385. }
  386. const permission: PermissionRequest = {
  387. id: "per_child_startup",
  388. sessionID: "ses_child",
  389. action: "shell",
  390. resources: ["git status --short"],
  391. source: { type: "tool", messageID: "msg_child_source", id: "call_child_source" },
  392. }
  393. const client = sdk({
  394. streams: [events],
  395. sessions: [{ id: "ses_child", parentID: "ses_1", title: "Child", time: { updated: 1 } }],
  396. permissions: { ses_child: [permission] },
  397. messages: {
  398. ses_child: [sourceMessage],
  399. },
  400. })
  401. const releaseSource = defer<void>()
  402. let sourceLookups = 0
  403. spyOn(client.session, "message").mockImplementation(async () => {
  404. sourceLookups++
  405. if (sourceLookups === 1) throw new Error("source temporarily unavailable")
  406. await releaseSource.promise
  407. return sourceMessage as never
  408. })
  409. const ui = footer()
  410. const transport = await createSessionTransport({
  411. sdk: client,
  412. sessionID: "ses_1",
  413. thinking: false,
  414. footer: ui.api,
  415. })
  416. while (sourceLookups < 2) await Bun.sleep(0)
  417. expect(
  418. ui.events.some(
  419. (event) =>
  420. event.type === "stream.view" && event.view.type === "permission" && event.view.request.id === permission.id,
  421. ),
  422. ).toBe(false)
  423. releaseSource.resolve()
  424. while (
  425. !ui.events.some(
  426. (event) =>
  427. event.type === "stream.view" && event.view.type === "permission" && event.view.request.id === permission.id,
  428. )
  429. )
  430. await Bun.sleep(0)
  431. expect(client.session.message).toHaveBeenCalledWith(
  432. { sessionID: "ses_child", messageID: "msg_child_source" },
  433. { signal: expect.any(AbortSignal) },
  434. )
  435. expect(
  436. ui.events.find(
  437. (event) =>
  438. event.type === "stream.view" && event.view.type === "permission" && event.view.request.id === permission.id,
  439. ),
  440. ).toMatchObject({
  441. view: {
  442. request: {
  443. tool: {
  444. id: "call_child_source",
  445. name: "shell",
  446. state: { status: "running", input: { command: "git status --short" } },
  447. },
  448. },
  449. },
  450. })
  451. expect(client.message.list).not.toHaveBeenCalledWith(
  452. expect.objectContaining({ sessionID: "ses_child" }),
  453. expect.anything(),
  454. )
  455. await transport.close()
  456. })
  457. test("reduces nested form owners idempotently and filters global events by complete location", async () => {
  458. const events = feed()
  459. events.push(connected())
  460. const client = sdk({
  461. streams: [events],
  462. sessions: [{ id: "ses_child", parentID: "ses_1", title: "Child", time: { updated: 1 } }],
  463. globalLocation: { directory: "/work", workspaceID: "wrk_1" },
  464. })
  465. const ui = footer()
  466. const transport = await createSessionTransport({
  467. sdk: client,
  468. location: { directory: "/work", workspaceID: "wrk_1" },
  469. sessionID: "ses_1",
  470. thinking: false,
  471. footer: ui.api,
  472. })
  473. const child = form("frm_child_live", "ses_child")
  474. events.push({ id: "evt_child_form", created: 1, type: "form.created", data: { form: eventForm(child) } })
  475. events.push({ id: "evt_child_form_retry", created: 2, type: "form.created", data: { form: eventForm(child) } })
  476. while (
  477. !ui.events.some(
  478. (event) => event.type === "stream.view" && event.view.type === "form" && event.view.request.id === child.id,
  479. )
  480. )
  481. await Bun.sleep(0)
  482. const childSnapshots = ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  483. expect(childSnapshots.at(-1)?.forms.filter((item) => item.id === child.id)).toHaveLength(1)
  484. events.push({
  485. id: "evt_child_form_done",
  486. created: 3,
  487. type: "form.replied",
  488. data: { id: child.id, sessionID: "ses_child", answer: { answer: "yes" } },
  489. })
  490. const global = form("frm_global_live", "global")
  491. events.push({
  492. id: "evt_global_wrong",
  493. created: 4,
  494. type: "form.created",
  495. location: { directory: "/work", workspaceID: "wrk_other" },
  496. data: { form: eventForm(global) },
  497. })
  498. await Bun.sleep(0)
  499. expect(
  500. ui.events.some(
  501. (event) => event.type === "stream.view" && event.view.type === "form" && event.view.request.id === global.id,
  502. ),
  503. ).toBe(false)
  504. events.push({
  505. id: "evt_global_right",
  506. created: 5,
  507. type: "form.created",
  508. location: { directory: "/work", workspaceID: "wrk_1" },
  509. data: { form: eventForm(global) },
  510. })
  511. while (
  512. !ui.events.some(
  513. (event) => event.type === "stream.view" && event.view.type === "form" && event.view.request.id === global.id,
  514. )
  515. )
  516. await Bun.sleep(0)
  517. expect(ui.events.at(-1)).toMatchObject({
  518. type: "stream.view",
  519. view: {
  520. type: "form",
  521. request: { id: "frm_global_live", location: { directory: "/work", workspaceID: "wrk_1" } },
  522. },
  523. })
  524. const beforeCancel = ui.events.filter((event) => event.type === "stream.view").length
  525. events.push({
  526. id: "evt_global_done",
  527. created: 6,
  528. type: "form.cancelled",
  529. location: { directory: "/work", workspaceID: "wrk_1" },
  530. data: { id: global.id, sessionID: "global" },
  531. })
  532. while (ui.events.filter((event) => event.type === "stream.view").length === beforeCancel) await Bun.sleep(0)
  533. expect(ui.events.filter((event) => event.type === "stream.view").at(-1)).toEqual({
  534. type: "stream.view",
  535. view: { type: "prompt" },
  536. })
  537. await transport.close()
  538. })
  539. test("waits authoritatively and reconciles the projected terminal suffix", async () => {
  540. const events = feed()
  541. events.push(connected())
  542. const settled = defer()
  543. const messages: SessionMessages = []
  544. const client = sdk({
  545. streams: [events],
  546. messages: { ses_1: messages },
  547. wait: () => settled.promise,
  548. })
  549. const ui = footer()
  550. const transport = await createSessionTransport({
  551. sdk: client,
  552. sessionID: "ses_1",
  553. thinking: false,
  554. footer: ui.api,
  555. })
  556. let admitted = false
  557. spyOn(client.session, "prompt").mockImplementation((request) => {
  558. admitted = true
  559. return ok({ data: promptAdmission(request) }) as never
  560. })
  561. const turn = transport.runPromptTurn({
  562. agent: undefined,
  563. model: undefined,
  564. variant: undefined,
  565. prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
  566. files: [],
  567. includeFiles: true,
  568. })
  569. while (!admitted) await Bun.sleep(0)
  570. events.push({
  571. id: "evt_prompted",
  572. created: 0,
  573. type: "session.input.promoted",
  574. durable: durable("ses_1"),
  575. data: {
  576. sessionID: "ses_1",
  577. inputID: "msg_prompt",
  578. },
  579. })
  580. events.push({
  581. id: "evt_text",
  582. created: 0,
  583. type: "session.text.delta",
  584. data: {
  585. sessionID: "ses_1",
  586. assistantMessageID: "msg_assistant",
  587. ordinal: 0,
  588. delta: "ans",
  589. },
  590. })
  591. events.push({
  592. id: "evt_settled",
  593. created: 0,
  594. type: "session.execution.succeeded",
  595. durable: durable("ses_1"),
  596. data: { sessionID: "ses_1" },
  597. })
  598. let done = false
  599. void turn.then(() => {
  600. done = true
  601. })
  602. await Bun.sleep(0)
  603. expect(done).toBe(false)
  604. messages.push(
  605. { id: "msg_prompt", type: "user", text: "hello", time: { created: 2 } },
  606. {
  607. id: "msg_assistant",
  608. type: "assistant",
  609. agent: "build",
  610. model: { providerID: "test", id: "model" },
  611. content: [{ type: "text", text: "answer" }],
  612. time: { created: 3, completed: 4 },
  613. },
  614. )
  615. settled.resolve()
  616. await turn
  617. expect(ui.commits.map((item) => item.text)).toEqual(["ans", "wer"])
  618. await transport.close()
  619. })
  620. test("shows durable pending delivery and appends queued input on promotion", async () => {
  621. const events = feed()
  622. events.push(connected())
  623. const client = sdk({
  624. streams: [events],
  625. pending: {
  626. ses_1: [
  627. {
  628. id: "msg_queued",
  629. sessionID: "ses_1",
  630. timeCreated: 1,
  631. type: "user",
  632. data: { text: "follow up" },
  633. delivery: "queue",
  634. },
  635. ],
  636. },
  637. })
  638. const ui = footer()
  639. const transport = await createSessionTransport({
  640. sdk: client,
  641. sessionID: "ses_1",
  642. thinking: false,
  643. footer: ui.api,
  644. })
  645. const pending = () =>
  646. ui.events
  647. .findLast((item) => item.type === "queued.prompts")
  648. ?.prompts.map((item) => [item.messageID, item.delivery])
  649. expect(pending()).toEqual([["msg_queued", "queue"]])
  650. events.push({
  651. id: "evt_promoted",
  652. created: 2,
  653. type: "session.input.promoted",
  654. durable: durable("ses_1", 2),
  655. data: { sessionID: "ses_1", inputID: "msg_queued" },
  656. })
  657. while (!ui.commits.some((item) => item.messageID === "msg_queued")) await Bun.sleep(0)
  658. expect(ui.commits).toContainEqual(
  659. expect.objectContaining({ kind: "user", messageID: "msg_queued", text: "follow up" }),
  660. )
  661. expect(pending()).toEqual([])
  662. const prompt = spyOn(client.session, "prompt").mockImplementation(
  663. (request) => ok(promptAdmission(request)) as never,
  664. )
  665. await transport.queuePromptTurn({
  666. agent: "review",
  667. model: undefined,
  668. variant: undefined,
  669. prompt: { messageID: "msg_next", text: "another", parts: [] },
  670. files: [],
  671. includeFiles: false,
  672. })
  673. expect(client.session.switchAgent).toHaveBeenCalledWith({ sessionID: "ses_1", agent: "review" }, expect.anything())
  674. expect(prompt).toHaveBeenCalledWith(expect.objectContaining({ delivery: "queue" }), expect.anything())
  675. events.push({
  676. id: "evt_earlier_admission",
  677. created: 3,
  678. type: "session.input.admitted",
  679. durable: durable("ses_1", 1),
  680. data: {
  681. sessionID: "ses_1",
  682. inputID: "msg_earlier",
  683. input: { type: "user", data: { text: "earlier" }, delivery: "steer" },
  684. },
  685. })
  686. while (true) {
  687. const pending = ui.events.findLast((item) => item.type === "queued.prompts")
  688. if (pending?.type === "queued.prompts" && pending.prompts.length >= 2) break
  689. await Bun.sleep(0)
  690. }
  691. expect(pending()).toEqual([
  692. ["msg_next", "queue"],
  693. ["msg_earlier", "steer"],
  694. ])
  695. await transport.close()
  696. })
  697. test("reports an observed execution failure before prompt promotion", async () => {
  698. const events = feed()
  699. events.push(connected())
  700. const idle = defer()
  701. const client = sdk({ streams: [events], messages: { ses_1: [] }, wait: () => idle.promise })
  702. const ui = footer()
  703. const transport = await createSessionTransport({
  704. sdk: client,
  705. sessionID: "ses_1",
  706. thinking: false,
  707. footer: ui.api,
  708. })
  709. let admitted = false
  710. spyOn(client.session, "prompt").mockImplementation((request) => {
  711. admitted = true
  712. return ok(promptAdmission(request)) as never
  713. })
  714. const turn = transport.runPromptTurn({
  715. agent: undefined,
  716. model: undefined,
  717. variant: undefined,
  718. prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
  719. files: [],
  720. includeFiles: false,
  721. })
  722. while (!admitted) await Bun.sleep(0)
  723. events.push({
  724. id: "evt_failed",
  725. created: 2,
  726. type: "session.execution.failed",
  727. durable: durable("ses_1", 2),
  728. data: { sessionID: "ses_1", error: { type: "unknown", message: "instructions unavailable" } },
  729. })
  730. await Bun.sleep(0)
  731. idle.resolve()
  732. await turn
  733. expect(ui.commits).toContainEqual(
  734. expect.objectContaining({ kind: "error", messageID: "msg_prompt", text: "instructions unavailable" }),
  735. )
  736. await transport.close()
  737. })
  738. test("attributes an execution-only failure to the latest promoted prompt", async () => {
  739. const events = feed()
  740. events.push(connected())
  741. const idle = defer()
  742. const messages: SessionMessages = []
  743. const client = sdk({ streams: [events], messages: { ses_1: messages }, wait: () => idle.promise })
  744. const ui = footer()
  745. const transport = await createSessionTransport({
  746. sdk: client,
  747. sessionID: "ses_1",
  748. thinking: false,
  749. footer: ui.api,
  750. })
  751. let admitted = false
  752. spyOn(client.session, "prompt").mockImplementation((request) => {
  753. admitted = true
  754. return ok(promptAdmission(request)) as never
  755. })
  756. const turn = transport.runPromptTurn({
  757. agent: undefined,
  758. model: undefined,
  759. variant: undefined,
  760. prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
  761. files: [],
  762. includeFiles: false,
  763. })
  764. while (!admitted) await Bun.sleep(0)
  765. events.push({
  766. id: "evt_prompt_promoted",
  767. created: 2,
  768. type: "session.input.promoted",
  769. durable: durable("ses_1", 2),
  770. data: { sessionID: "ses_1", inputID: "msg_prompt" },
  771. })
  772. await transport.queuePromptTurn({
  773. agent: undefined,
  774. model: undefined,
  775. variant: undefined,
  776. prompt: { messageID: "msg_queued", text: "follow up", parts: [] },
  777. files: [],
  778. includeFiles: false,
  779. })
  780. events.push({
  781. id: "evt_queued_promoted",
  782. created: 3,
  783. type: "session.input.promoted",
  784. durable: durable("ses_1", 3),
  785. data: { sessionID: "ses_1", inputID: "msg_queued" },
  786. })
  787. events.push({
  788. id: "evt_failed",
  789. created: 4,
  790. type: "session.execution.failed",
  791. durable: durable("ses_1", 4),
  792. data: { sessionID: "ses_1", error: { type: "unknown", message: "model unavailable" } },
  793. })
  794. await Bun.sleep(0)
  795. messages.push(
  796. { id: "msg_prompt", type: "user", text: "hello", time: { created: 2 } },
  797. { id: "msg_queued", type: "user", text: "follow up", time: { created: 3 } },
  798. )
  799. idle.resolve()
  800. await turn
  801. expect(ui.commits).toContainEqual(
  802. expect.objectContaining({ kind: "error", messageID: "msg_queued", text: "model unavailable" }),
  803. )
  804. await transport.close()
  805. })
  806. test("sends local file and directory mentions as structured prompt files", async () => {
  807. await using tmp = await tmpdir()
  808. const filePath = path.join(tmp.path, "note.ts")
  809. const contextPath = path.join(tmp.path, "context.txt")
  810. const directoryPath = path.join(tmp.path, "docs")
  811. await Bun.write(filePath, "export const answer = 42\n")
  812. await Bun.write(contextPath, "context body")
  813. await fs.mkdir(directoryPath)
  814. await Bun.write(path.join(directoryPath, "README.md"), "# hello\n")
  815. const events = feed()
  816. events.push(connected())
  817. const client = sdk({ streams: [events] })
  818. const ui = footer()
  819. const transport = await createSessionTransport({
  820. sdk: client,
  821. readTextFile: (url) => fs.readFile(new URL(url), "utf8"),
  822. sessionID: "ses_1",
  823. thinking: false,
  824. footer: ui.api,
  825. })
  826. let request: Parameters<OpenCodeClient["session"]["prompt"]>[0] | undefined
  827. spyOn(client.session, "prompt").mockImplementation((input) => {
  828. request = input
  829. queueMicrotask(() => {
  830. events.push({
  831. id: "evt_prompted",
  832. created: 0,
  833. type: "session.input.promoted",
  834. durable: durable("ses_1"),
  835. data: {
  836. sessionID: "ses_1",
  837. inputID: "msg_prompt",
  838. },
  839. })
  840. events.push({
  841. id: "evt_settled",
  842. created: 0,
  843. type: "session.execution.succeeded",
  844. durable: durable("ses_1"),
  845. data: { sessionID: "ses_1" },
  846. })
  847. })
  848. return ok({ data: promptAdmission(input) }) as never
  849. })
  850. await transport.runPromptTurn({
  851. agent: undefined,
  852. model: undefined,
  853. variant: undefined,
  854. prompt: {
  855. messageID: "msg_prompt",
  856. text: "Review @note.ts and @docs",
  857. parts: [
  858. {
  859. type: "file",
  860. url: pathToFileURL(filePath).href,
  861. mime: "text/plain",
  862. filename: "note.ts",
  863. source: { type: "file", path: "note.ts", text: { start: 7, end: 15, value: "@note.ts" } },
  864. },
  865. {
  866. type: "file",
  867. url: pathToFileURL(`${directoryPath}${path.sep}`).href,
  868. mime: "application/x-directory",
  869. filename: "docs",
  870. source: { type: "file", path: "docs/", text: { start: 20, end: 25, value: "@docs" } },
  871. },
  872. ],
  873. },
  874. files: [
  875. { type: "file", url: pathToFileURL(contextPath).href, filename: "context.txt", mime: "text/plain" },
  876. { type: "file", url: "file:///tmp/image.png", filename: "image.png", mime: "image/png" },
  877. ],
  878. includeFiles: true,
  879. })
  880. expect(request?.text).toBe('Review @note.ts and @docs\n\n<file name="context.txt">\ncontext body\n</file>')
  881. expect(request?.files).toEqual([
  882. { uri: "file:///tmp/image.png", name: "image.png" },
  883. {
  884. uri: pathToFileURL(filePath).href,
  885. name: "note.ts",
  886. mention: { start: 7, end: 15, text: "@note.ts" },
  887. },
  888. {
  889. uri: pathToFileURL(`${directoryPath}${path.sep}`).href,
  890. name: "docs",
  891. mention: { start: 20, end: 25, text: "@docs" },
  892. },
  893. ])
  894. await transport.close()
  895. })
  896. test("sends attached file mentions as structured prompt files without reading them", async () => {
  897. const events = feed()
  898. events.push(connected())
  899. const client = sdk({ streams: [events] })
  900. const ui = footer()
  901. const remoteRead = spyOn(client.file, "read")
  902. const remoteList = spyOn(client.file, "list")
  903. const transport = await createSessionTransport({
  904. sdk: client,
  905. location: { directory: "/remote/project" },
  906. sessionID: "ses_1",
  907. thinking: false,
  908. footer: ui.api,
  909. })
  910. let request: Parameters<OpenCodeClient["session"]["prompt"]>[0] | undefined
  911. // The generated method has conditional return types for throwOnError; this mock represents the successful branch.
  912. // @ts-expect-error successful SDK response is valid for both modes at runtime
  913. spyOn(client.session, "prompt").mockImplementation((input) => {
  914. request = input
  915. queueMicrotask(() => {
  916. events.push({
  917. id: "evt_prompted",
  918. created: 0,
  919. type: "session.input.promoted",
  920. durable: durable("ses_1"),
  921. data: {
  922. sessionID: "ses_1",
  923. inputID: "msg_prompt",
  924. },
  925. })
  926. events.push({
  927. id: "evt_settled",
  928. created: 0,
  929. type: "session.execution.succeeded",
  930. durable: durable("ses_1"),
  931. data: { sessionID: "ses_1" },
  932. })
  933. })
  934. return ok({ data: promptAdmission(input) })
  935. })
  936. await transport.runPromptTurn({
  937. agent: undefined,
  938. model: undefined,
  939. variant: undefined,
  940. prompt: {
  941. messageID: "msg_prompt",
  942. text: "Review @note.ts and @docs",
  943. parts: [
  944. {
  945. type: "file",
  946. url: "file:///remote/project/note.ts",
  947. mime: "text/plain",
  948. filename: "note.ts",
  949. source: { type: "file", path: "note.ts", text: { start: 7, end: 15, value: "@note.ts" } },
  950. },
  951. {
  952. type: "file",
  953. url: "file:///remote/project/docs",
  954. mime: "application/x-directory",
  955. filename: "docs",
  956. source: { type: "file", path: "docs", text: { start: 20, end: 25, value: "@docs" } },
  957. },
  958. ],
  959. },
  960. files: [],
  961. includeFiles: true,
  962. })
  963. expect(remoteRead).not.toHaveBeenCalled()
  964. expect(remoteList).not.toHaveBeenCalled()
  965. expect(request?.text).toBe("Review @note.ts and @docs")
  966. expect(request?.files).toEqual([
  967. {
  968. uri: "file:///remote/project/note.ts",
  969. name: "note.ts",
  970. mention: { start: 7, end: 15, text: "@note.ts" },
  971. },
  972. {
  973. uri: "file:///remote/project/docs",
  974. name: "docs",
  975. mention: { start: 20, end: 25, text: "@docs" },
  976. },
  977. ])
  978. await transport.close()
  979. })
  980. test("sends local media mentions as structured prompt files", async () => {
  981. await using tmp = await tmpdir()
  982. const filePath = path.join(tmp.path, "diagram.png")
  983. await Bun.write(filePath, Uint8Array.of(0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a, 0x00))
  984. const events = feed()
  985. events.push(connected())
  986. const client = sdk({ streams: [events] })
  987. const ui = footer()
  988. const transport = await createSessionTransport({
  989. sdk: client,
  990. sessionID: "ses_1",
  991. thinking: false,
  992. footer: ui.api,
  993. })
  994. let request: Parameters<OpenCodeClient["session"]["prompt"]>[0] | undefined
  995. // The generated method has conditional return types for throwOnError; this mock represents the successful branch.
  996. // @ts-expect-error successful SDK response is valid for both modes at runtime
  997. spyOn(client.session, "prompt").mockImplementation((input) => {
  998. request = input
  999. queueMicrotask(() => {
  1000. events.push({
  1001. id: "evt_prompted",
  1002. created: 0,
  1003. type: "session.input.promoted",
  1004. durable: durable("ses_1"),
  1005. data: {
  1006. sessionID: "ses_1",
  1007. inputID: "msg_prompt",
  1008. },
  1009. })
  1010. events.push({
  1011. id: "evt_settled",
  1012. created: 0,
  1013. type: "session.execution.succeeded",
  1014. durable: durable("ses_1"),
  1015. data: { sessionID: "ses_1" },
  1016. })
  1017. })
  1018. return ok({ data: promptAdmission(input) })
  1019. })
  1020. await transport.runPromptTurn({
  1021. agent: undefined,
  1022. model: undefined,
  1023. variant: undefined,
  1024. prompt: {
  1025. messageID: "msg_prompt",
  1026. text: "Review @diagram.png",
  1027. parts: [
  1028. {
  1029. type: "file",
  1030. url: pathToFileURL(filePath).href,
  1031. mime: "text/plain",
  1032. filename: "diagram.png",
  1033. source: { type: "file", path: "diagram.png", text: { start: 7, end: 19, value: "@diagram.png" } },
  1034. },
  1035. ],
  1036. },
  1037. files: [],
  1038. includeFiles: true,
  1039. })
  1040. expect(request?.text).toBe("Review @diagram.png")
  1041. expect(request?.files).toEqual([
  1042. {
  1043. name: "diagram.png",
  1044. uri: pathToFileURL(filePath).href,
  1045. mention: { start: 7, end: 19, text: "@diagram.png" },
  1046. },
  1047. ])
  1048. await transport.close()
  1049. })
  1050. test("shows V2 blockers and replies through the runtime-owned session API", async () => {
  1051. const events = feed()
  1052. events.push(connected())
  1053. const client = sdk({ streams: [events] })
  1054. const ui = footer()
  1055. const transport = await createSessionTransport({
  1056. sdk: client,
  1057. sessionID: "ses_1",
  1058. thinking: false,
  1059. footer: ui.api,
  1060. })
  1061. events.push({
  1062. id: "evt_permission",
  1063. created: 0,
  1064. type: "permission.asked",
  1065. data: { id: "per_1", sessionID: "ses_1", action: "read", resources: ["/tmp/file"] },
  1066. })
  1067. await Bun.sleep(0)
  1068. expect(ui.events).toContainEqual({
  1069. type: "stream.view",
  1070. view: {
  1071. type: "permission",
  1072. request: {
  1073. id: "per_1",
  1074. sessionID: "ses_1",
  1075. action: "read",
  1076. resources: ["/tmp/file"],
  1077. },
  1078. },
  1079. })
  1080. await transport.close()
  1081. })
  1082. test("reconnects and hydrates without completing before session.wait", async () => {
  1083. const first = feed()
  1084. const second = feed()
  1085. first.push(connected("evt_connected_1"))
  1086. second.push(connected("evt_connected_2"))
  1087. const idle = defer()
  1088. let running = true
  1089. const client = sdk({
  1090. streams: [first, second],
  1091. active: () => {
  1092. const active: Record<string, { type: "running" }> = {}
  1093. if (running) active.ses_1 = { type: "running" }
  1094. return active
  1095. },
  1096. wait: () => idle.promise,
  1097. })
  1098. let projected = false
  1099. spyOn(client.message, "list").mockImplementation(() =>
  1100. ok({
  1101. data: projected
  1102. ? [
  1103. {
  1104. id: "msg_prompt",
  1105. type: "user",
  1106. text: "hello",
  1107. files: [],
  1108. agents: [],
  1109. time: { created: 2 },
  1110. },
  1111. ]
  1112. : [],
  1113. cursor: {},
  1114. }),
  1115. )
  1116. const ui = footer()
  1117. const transport = await createSessionTransport({
  1118. sdk: client,
  1119. sessionID: "ses_1",
  1120. thinking: false,
  1121. footer: ui.api,
  1122. })
  1123. let admitted = false
  1124. // The generated method has conditional return types for throwOnError; this mock represents the successful branch.
  1125. // @ts-expect-error successful SDK response is valid for both modes at runtime
  1126. spyOn(client.session, "prompt").mockImplementation((request) => {
  1127. admitted = true
  1128. return ok({ data: promptAdmission(request) })
  1129. })
  1130. const turn = transport.runPromptTurn({
  1131. agent: undefined,
  1132. model: undefined,
  1133. variant: undefined,
  1134. prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
  1135. files: [],
  1136. includeFiles: true,
  1137. })
  1138. while (!admitted) await Bun.sleep(0)
  1139. projected = true
  1140. running = false
  1141. second.push({
  1142. id: "evt_prior_failed",
  1143. created: 1,
  1144. type: "session.execution.failed",
  1145. durable: durable("ses_1", 1),
  1146. data: { sessionID: "ses_1", error: { type: "unknown", message: "prior execution failed" } },
  1147. })
  1148. second.push({
  1149. id: "evt_prompted",
  1150. created: 2,
  1151. type: "session.input.promoted",
  1152. durable: durable("ses_1", 2),
  1153. data: { sessionID: "ses_1", inputID: "msg_prompt" },
  1154. })
  1155. first.close()
  1156. while (!ui.events.some((event) => event.type === "stream.patch" && event.patch.status === "reconnecting"))
  1157. await Bun.sleep(0)
  1158. idle.resolve()
  1159. await turn
  1160. await transport.close()
  1161. })
  1162. test("does not duplicate the optimistic user row when reconnect hydration recovers a missed prompt", async () => {
  1163. const first = feed()
  1164. const second = feed()
  1165. first.push(connected("evt_connected_1"))
  1166. second.push(connected("evt_connected_2"))
  1167. let running = true
  1168. let projected = false
  1169. const client = sdk({
  1170. streams: [first, second],
  1171. active: () => {
  1172. const active: Record<string, { type: "running" }> = {}
  1173. if (running) active.ses_1 = { type: "running" }
  1174. return active
  1175. },
  1176. })
  1177. spyOn(client.message, "list").mockImplementation(() =>
  1178. ok({
  1179. data: projected
  1180. ? [
  1181. {
  1182. id: "msg_prompt",
  1183. type: "user",
  1184. text: "hello",
  1185. files: [],
  1186. agents: [],
  1187. time: { created: 2 },
  1188. },
  1189. ]
  1190. : [],
  1191. cursor: {},
  1192. }),
  1193. )
  1194. const ui = footer()
  1195. ui.commits.push({ kind: "user", source: "system", text: "hello", phase: "start", messageID: "msg_prompt" })
  1196. const transport = await createSessionTransport({
  1197. sdk: client,
  1198. sessionID: "ses_1",
  1199. thinking: false,
  1200. footer: ui.api,
  1201. })
  1202. let admitted = false
  1203. // The generated method has conditional return types for throwOnError; this mock represents the successful branch.
  1204. // @ts-expect-error successful SDK response is valid for both modes at runtime
  1205. spyOn(client.session, "prompt").mockImplementation((request) => {
  1206. admitted = true
  1207. return ok({ data: promptAdmission(request) })
  1208. })
  1209. const turn = transport.runPromptTurn({
  1210. agent: undefined,
  1211. model: undefined,
  1212. variant: undefined,
  1213. prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
  1214. files: [],
  1215. includeFiles: true,
  1216. })
  1217. while (!admitted) await Bun.sleep(0)
  1218. projected = true
  1219. running = false
  1220. first.close()
  1221. await turn
  1222. expect(ui.commits.filter((item) => item.kind === "user" && item.messageID === "msg_prompt")).toHaveLength(1)
  1223. await transport.close()
  1224. })
  1225. test("replaces the client for buffered hydration, descendants, turns, and interrupts", async () => {
  1226. const firstEvents = feed()
  1227. const secondEvents = feed()
  1228. firstEvents.push(connected("evt_connected_1"))
  1229. secondEvents.push(connected("evt_connected_2"))
  1230. const first = sdk({ streams: [firstEvents] })
  1231. const second = sdk({
  1232. streams: [secondEvents],
  1233. sessions: [{ id: "ses_child", parentID: "ses_1", title: "Child", time: { updated: 2 } }],
  1234. forms: { ses_child: [form("frm_child", "ses_child")] },
  1235. })
  1236. const firstPrompt = spyOn(first.session, "prompt")
  1237. const firstInterrupt = spyOn(first.session, "interrupt")
  1238. spyOn(first.message, "list").mockImplementation(() =>
  1239. ok({
  1240. data: [
  1241. {
  1242. id: "msg_assistant",
  1243. type: "assistant",
  1244. agent: "build",
  1245. model: { providerID: "test", id: "model" },
  1246. content: [{ type: "text", text: "partial" }],
  1247. time: { created: 1 },
  1248. },
  1249. ],
  1250. cursor: {},
  1251. }),
  1252. )
  1253. let releaseHydration!: () => void
  1254. let replacementHydrating = false
  1255. const hydration = new Promise<void>((resolve) => {
  1256. releaseHydration = resolve
  1257. })
  1258. let releaseCatalog!: () => void
  1259. let refreshes = 0
  1260. const catalog = new Promise<void>((resolve) => {
  1261. releaseCatalog = resolve
  1262. })
  1263. spyOn(second.message, "list").mockImplementation(async (request) => {
  1264. if (request.sessionID !== "ses_1") return ok({ data: [], cursor: {} })
  1265. replacementHydrating = true
  1266. await hydration
  1267. return ok({
  1268. data: [
  1269. {
  1270. id: "msg_assistant",
  1271. type: "assistant",
  1272. agent: "build",
  1273. model: { providerID: "test", id: "model" },
  1274. content: [{ type: "text", text: "partial replacement" }],
  1275. time: { created: 1 },
  1276. },
  1277. ],
  1278. cursor: {},
  1279. })
  1280. })
  1281. const current: OpenCodeClient[] = []
  1282. const ui = footer()
  1283. const transport = await createSessionTransport({
  1284. sdk: first,
  1285. reconnect: async () => second,
  1286. onClient: (client) => current.push(client),
  1287. sessionID: "ses_1",
  1288. thinking: false,
  1289. replay: true,
  1290. footer: ui.api,
  1291. onCatalogRefresh: () => {
  1292. refreshes++
  1293. if (refreshes === 2) return catalog
  1294. },
  1295. })
  1296. firstEvents.close()
  1297. while (!replacementHydrating) await Bun.sleep(0)
  1298. await expect(
  1299. transport.runPromptTurn({
  1300. agent: undefined,
  1301. model: undefined,
  1302. variant: undefined,
  1303. prompt: { messageID: "msg_blocked", text: "blocked", parts: [] },
  1304. files: [],
  1305. includeFiles: true,
  1306. }),
  1307. ).rejects.toThrow("Event stream is reconnecting")
  1308. secondEvents.push({
  1309. id: "evt_buffered_text",
  1310. created: 2,
  1311. type: "session.text.delta",
  1312. data: {
  1313. sessionID: "ses_1",
  1314. assistantMessageID: "msg_assistant",
  1315. ordinal: 0,
  1316. delta: " replacement",
  1317. },
  1318. })
  1319. let resized = false
  1320. const resize = transport.replayOnResize({
  1321. localRows: () => [],
  1322. reset: async () => {
  1323. resized = true
  1324. },
  1325. })
  1326. releaseHydration()
  1327. while (
  1328. !ui.events.some(
  1329. (event) => event.type === "stream.view" && event.view.type === "form" && event.view.request.id === "frm_child",
  1330. )
  1331. )
  1332. await Bun.sleep(0)
  1333. while (refreshes < 2) await Bun.sleep(0)
  1334. await resize
  1335. expect(resized).toBe(false)
  1336. await expect(
  1337. transport.runPromptTurn({
  1338. agent: undefined,
  1339. model: undefined,
  1340. variant: undefined,
  1341. prompt: { messageID: "msg_catalog_blocked", text: "blocked", parts: [] },
  1342. files: [],
  1343. includeFiles: true,
  1344. }),
  1345. ).rejects.toThrow("Event stream is reconnecting")
  1346. releaseCatalog()
  1347. await Bun.sleep(0)
  1348. expect(current).toEqual([second])
  1349. expect(first.event.subscribe).toHaveBeenCalledTimes(1)
  1350. expect(second.event.subscribe).toHaveBeenCalledTimes(1)
  1351. expect(second.session.list).toHaveBeenCalled()
  1352. expect(second.form.list).toHaveBeenCalledWith({ sessionID: "ses_child" }, { signal: expect.any(AbortSignal) })
  1353. expect(ui.commits.filter((commit) => commit.messageID === "msg_assistant").map((commit) => commit.text)).toEqual([
  1354. "partial",
  1355. " replacement",
  1356. ])
  1357. const prompt = spyOn(second.session, "prompt").mockImplementation((request) => {
  1358. queueMicrotask(() => {
  1359. secondEvents.push({
  1360. id: "evt_replacement_prompt",
  1361. created: 3,
  1362. type: "session.input.promoted",
  1363. durable: durable("ses_1", 1),
  1364. data: { sessionID: "ses_1", inputID: "msg_replacement" },
  1365. })
  1366. secondEvents.push({
  1367. id: "evt_replacement_settled",
  1368. created: 4,
  1369. type: "session.execution.succeeded",
  1370. durable: durable("ses_1", 2),
  1371. data: { sessionID: "ses_1" },
  1372. })
  1373. })
  1374. return ok({ data: promptAdmission(request) }) as never
  1375. })
  1376. await transport.runPromptTurn({
  1377. agent: undefined,
  1378. model: undefined,
  1379. variant: undefined,
  1380. prompt: { messageID: "msg_replacement", text: "replacement prompt", parts: [] },
  1381. files: [],
  1382. includeFiles: true,
  1383. })
  1384. const interrupt = spyOn(second.session, "interrupt").mockImplementation(() => ok(undefined))
  1385. await transport.interruptActiveTurn()
  1386. expect(prompt).toHaveBeenCalled()
  1387. expect(interrupt).toHaveBeenCalledWith({ sessionID: "ses_1" })
  1388. expect(firstPrompt).not.toHaveBeenCalled()
  1389. expect(firstInterrupt).not.toHaveBeenCalled()
  1390. await transport.close()
  1391. })
  1392. test("reconciles buffered deltas already present in a resize snapshot", async () => {
  1393. const events = feed()
  1394. events.push(connected())
  1395. const client = sdk({ streams: [events] })
  1396. const ui = footer()
  1397. const transport = await createSessionTransport({
  1398. sdk: client,
  1399. sessionID: "ses_1",
  1400. thinking: false,
  1401. replay: true,
  1402. footer: ui.api,
  1403. })
  1404. spyOn(client.message, "list").mockImplementation(() =>
  1405. ok({
  1406. data: [
  1407. {
  1408. id: "msg_assistant",
  1409. type: "assistant",
  1410. agent: "build",
  1411. model: { providerID: "test", id: "model" },
  1412. content: [{ type: "text", text: "the answer" }],
  1413. time: { created: 2, completed: 3 },
  1414. },
  1415. ],
  1416. cursor: {},
  1417. }),
  1418. )
  1419. let reset!: () => void
  1420. const resetting = new Promise<void>((resolve) => {
  1421. reset = resolve
  1422. })
  1423. const replay = transport.replayOnResize({ localRows: () => [], reset: () => resetting })
  1424. events.push({
  1425. id: "evt_text_started",
  1426. created: 0,
  1427. type: "session.text.started",
  1428. durable: durable("ses_1"),
  1429. data: {
  1430. sessionID: "ses_1",
  1431. assistantMessageID: "msg_assistant",
  1432. ordinal: 0,
  1433. },
  1434. })
  1435. events.push({
  1436. id: "evt_text",
  1437. created: 0,
  1438. type: "session.text.delta",
  1439. data: {
  1440. sessionID: "ses_1",
  1441. assistantMessageID: "msg_assistant",
  1442. ordinal: 0,
  1443. delta: "answer",
  1444. },
  1445. })
  1446. await Bun.sleep(0)
  1447. reset()
  1448. await replay
  1449. expect(ui.commits.filter((item) => item.text === "the answer")).toHaveLength(1)
  1450. expect(ui.commits.some((item) => item.text === "answer")).toBe(false)
  1451. await transport.close()
  1452. })
  1453. test("replays live assistant text missing from the resize projection", async () => {
  1454. const events = feed()
  1455. events.push(connected())
  1456. const client = sdk({ streams: [events] })
  1457. spyOn(client.message, "list").mockImplementation(() =>
  1458. ok({
  1459. data: [
  1460. {
  1461. id: "msg_assistant",
  1462. type: "assistant",
  1463. agent: "build",
  1464. model: { providerID: "test", id: "model" },
  1465. content: [{ type: "text", text: "partial" }],
  1466. time: { created: 2, completed: 3 },
  1467. },
  1468. ],
  1469. cursor: {},
  1470. }),
  1471. )
  1472. const ui = footer()
  1473. const live: StreamCommit[] = []
  1474. const transport = await createSessionTransport({
  1475. sdk: client,
  1476. sessionID: "ses_1",
  1477. thinking: false,
  1478. replay: true,
  1479. footer: ui.api,
  1480. onCommit: (commit) => live.push(commit),
  1481. })
  1482. events.push({
  1483. id: "evt_text",
  1484. created: 0,
  1485. type: "session.text.delta",
  1486. data: {
  1487. sessionID: "ses_1",
  1488. assistantMessageID: "msg_assistant",
  1489. ordinal: 0,
  1490. delta: " suffix",
  1491. },
  1492. })
  1493. await Bun.sleep(0)
  1494. expect(live.map((commit) => commit.text)).toEqual(["partial suffix"])
  1495. await transport.replayOnResize({
  1496. localRows: () => [
  1497. { commit: live[0]! },
  1498. {
  1499. commit: {
  1500. ...live[0]!,
  1501. partID: "text:1",
  1502. text: "entirely local",
  1503. },
  1504. },
  1505. ],
  1506. reset: async () => {},
  1507. })
  1508. expect(ui.commits.filter((commit) => commit.messageID === "msg_assistant").map((commit) => commit.text)).toEqual([
  1509. "partial",
  1510. " suffix",
  1511. "partial",
  1512. " suffix",
  1513. "entirely local",
  1514. ])
  1515. await transport.close()
  1516. })
  1517. test("does not replay a resize-buffered suffix twice", async () => {
  1518. const events = feed()
  1519. events.push(connected())
  1520. const client = sdk({ streams: [events] })
  1521. spyOn(client.message, "list").mockImplementation(() =>
  1522. ok({
  1523. data: [
  1524. {
  1525. id: "msg_assistant",
  1526. type: "assistant",
  1527. agent: "build",
  1528. model: { providerID: "test", id: "model" },
  1529. content: [{ type: "text", text: "partial" }],
  1530. time: { created: 2, completed: 3 },
  1531. },
  1532. ],
  1533. cursor: {},
  1534. }),
  1535. )
  1536. const ui = footer()
  1537. const live: StreamCommit[] = []
  1538. const transport = await createSessionTransport({
  1539. sdk: client,
  1540. sessionID: "ses_1",
  1541. thinking: false,
  1542. replay: true,
  1543. footer: ui.api,
  1544. onCommit: (commit) => live.push(commit),
  1545. })
  1546. let reset!: () => void
  1547. const resetting = new Promise<void>((resolve) => {
  1548. reset = resolve
  1549. })
  1550. const replay = transport.replayOnResize({
  1551. localRows: () => live.map((commit) => ({ commit })),
  1552. reset: () => resetting,
  1553. })
  1554. events.push({
  1555. id: "evt_text",
  1556. created: 0,
  1557. type: "session.text.delta",
  1558. data: {
  1559. sessionID: "ses_1",
  1560. assistantMessageID: "msg_assistant",
  1561. ordinal: 0,
  1562. delta: " suffix",
  1563. },
  1564. })
  1565. await Bun.sleep(0)
  1566. reset()
  1567. await replay
  1568. expect(ui.commits.filter((commit) => commit.messageID === "msg_assistant").map((commit) => commit.text)).toEqual([
  1569. "partial",
  1570. "partial",
  1571. " suffix",
  1572. ])
  1573. expect(live.map((commit) => commit.text)).toEqual(["partial suffix"])
  1574. await transport.close()
  1575. })
  1576. test("preserves active text and reasoning across resize before terminal projection", async () => {
  1577. const events = feed()
  1578. events.push(connected())
  1579. const client = sdk({ streams: [events] })
  1580. spyOn(client.message, "list").mockImplementation(() => ok({ data: [], cursor: {} }))
  1581. const ui = footer()
  1582. const live: StreamCommit[] = []
  1583. const transport = await createSessionTransport({
  1584. sdk: client,
  1585. sessionID: "ses_1",
  1586. thinking: true,
  1587. replay: true,
  1588. footer: ui.api,
  1589. onCommit: (commit) => live.push(commit),
  1590. })
  1591. events.push({
  1592. id: "evt_text",
  1593. created: 0,
  1594. type: "session.text.delta",
  1595. data: {
  1596. sessionID: "ses_1",
  1597. assistantMessageID: "msg_assistant",
  1598. ordinal: 0,
  1599. delta: "hello",
  1600. },
  1601. })
  1602. events.push({
  1603. id: "evt_reasoning",
  1604. created: 0,
  1605. type: "session.reasoning.delta",
  1606. data: {
  1607. sessionID: "ses_1",
  1608. assistantMessageID: "msg_assistant",
  1609. ordinal: 0,
  1610. delta: "thought",
  1611. },
  1612. })
  1613. await Bun.sleep(0)
  1614. expect(live.map((commit) => commit.text)).toEqual(["hello", "Thinking: thought"])
  1615. await transport.replayOnResize({
  1616. localRows: () => live.map((commit) => ({ commit })),
  1617. reset: async () => {},
  1618. })
  1619. expect(ui.commits.slice(-2).map((commit) => commit.text)).toEqual(["hello", "Thinking: thought"])
  1620. await transport.close()
  1621. })
  1622. test("serializes and coalesces overlapping resize replays", async () => {
  1623. const events = feed()
  1624. events.push(connected())
  1625. const client = sdk({ streams: [events] })
  1626. const ui = footer()
  1627. const transport = await createSessionTransport({
  1628. sdk: client,
  1629. sessionID: "ses_1",
  1630. thinking: false,
  1631. replay: true,
  1632. footer: ui.api,
  1633. })
  1634. let release!: () => void
  1635. const blocked = new Promise<void>((resolve) => {
  1636. release = resolve
  1637. })
  1638. const order: string[] = []
  1639. const first = transport.replayOnResize({
  1640. localRows: () => [],
  1641. reset: async () => {
  1642. order.push("first:start")
  1643. await blocked
  1644. order.push("first:end")
  1645. },
  1646. })
  1647. await Bun.sleep(0)
  1648. const second = transport.replayOnResize({
  1649. localRows: () => [],
  1650. reset: async () => {
  1651. order.push("second")
  1652. },
  1653. })
  1654. release()
  1655. await Promise.all([first, second])
  1656. expect(second).toBe(first)
  1657. expect(order).toEqual(["first:start", "first:end", "second"])
  1658. await transport.close()
  1659. })
  1660. test("restores local output and drains buffered events when resize hydration fails", async () => {
  1661. const events = feed()
  1662. events.push(connected())
  1663. const client = sdk({ streams: [events] })
  1664. const ui = footer()
  1665. const live: StreamCommit[] = []
  1666. const transport = await createSessionTransport({
  1667. sdk: client,
  1668. sessionID: "ses_1",
  1669. thinking: false,
  1670. replay: true,
  1671. footer: ui.api,
  1672. onCommit: (commit) => live.push(commit),
  1673. })
  1674. events.push({
  1675. id: "evt_text_1",
  1676. created: 0,
  1677. type: "session.text.delta",
  1678. data: {
  1679. sessionID: "ses_1",
  1680. assistantMessageID: "msg_assistant",
  1681. ordinal: 0,
  1682. delta: "hello",
  1683. },
  1684. })
  1685. await Bun.sleep(0)
  1686. spyOn(client.message, "list").mockImplementation(() => Promise.reject(new Error("projection failed")))
  1687. const replay = transport.replayOnResize({
  1688. localRows: () => live.map((commit) => ({ commit })),
  1689. reset: async () => {},
  1690. })
  1691. events.push({
  1692. id: "evt_text_2",
  1693. created: 0,
  1694. type: "session.text.delta",
  1695. data: {
  1696. sessionID: "ses_1",
  1697. assistantMessageID: "msg_assistant",
  1698. ordinal: 0,
  1699. delta: " world",
  1700. },
  1701. })
  1702. await expect(replay).rejects.toThrow("projection failed")
  1703. expect(ui.commits.slice(-2).map((commit) => commit.text)).toEqual(["hello", " world"])
  1704. expect(live.at(-1)?.text).toBe("hello world")
  1705. await transport.close()
  1706. })
  1707. test("dedupes a projected step failure from live redelivery", async () => {
  1708. const events = feed()
  1709. events.push(connected())
  1710. const client = sdk({ streams: [events] })
  1711. spyOn(client.message, "list").mockImplementation(() =>
  1712. ok({
  1713. data: [
  1714. {
  1715. id: "msg_assistant",
  1716. type: "assistant",
  1717. agent: "build",
  1718. model: { providerID: "test", id: "model" },
  1719. content: [],
  1720. error: { type: "provider.transport", message: "provider failed" },
  1721. time: { created: 2, completed: 3 },
  1722. },
  1723. ],
  1724. cursor: {},
  1725. }),
  1726. )
  1727. const ui = footer()
  1728. const transport = await createSessionTransport({
  1729. sdk: client,
  1730. sessionID: "ses_1",
  1731. thinking: false,
  1732. replay: true,
  1733. footer: ui.api,
  1734. })
  1735. events.push({
  1736. id: "evt_step_failed",
  1737. created: 2,
  1738. type: "session.step.failed",
  1739. durable: durable("ses_1", 1),
  1740. data: {
  1741. sessionID: "ses_1",
  1742. assistantMessageID: "msg_assistant",
  1743. error: { type: "provider.transport", message: "provider failed" },
  1744. },
  1745. })
  1746. await Bun.sleep(0)
  1747. expect(ui.commits.filter((commit) => commit.kind === "error" && commit.text === "provider failed")).toHaveLength(1)
  1748. await transport.close()
  1749. })
  1750. test("dedupes a retained live step failure from resize projection", async () => {
  1751. const events = feed()
  1752. events.push(connected())
  1753. const client = sdk({ streams: [events] })
  1754. const ui = footer()
  1755. const live: StreamCommit[] = []
  1756. const transport = await createSessionTransport({
  1757. sdk: client,
  1758. sessionID: "ses_1",
  1759. thinking: false,
  1760. replay: true,
  1761. footer: ui.api,
  1762. onCommit: (commit) => live.push(commit),
  1763. })
  1764. events.push({
  1765. id: "evt_step_failed",
  1766. created: 2,
  1767. type: "session.step.failed",
  1768. durable: durable("ses_1", 1),
  1769. data: {
  1770. sessionID: "ses_1",
  1771. assistantMessageID: "msg_assistant",
  1772. error: { type: "provider.transport", message: "provider failed" },
  1773. },
  1774. })
  1775. await Bun.sleep(0)
  1776. expect(live[0]?.messageID).toBe("msg_assistant")
  1777. spyOn(client.message, "list").mockImplementation(() =>
  1778. ok({
  1779. data: [
  1780. {
  1781. id: "msg_assistant",
  1782. type: "assistant",
  1783. agent: "build",
  1784. model: { providerID: "test", id: "model" },
  1785. content: [],
  1786. error: { type: "provider.transport", message: "provider failed" },
  1787. time: { created: 2, completed: 3 },
  1788. },
  1789. ],
  1790. cursor: {},
  1791. }),
  1792. )
  1793. await transport.replayOnResize({
  1794. localRows: () => live.map((commit) => ({ commit })),
  1795. reset: async () => {},
  1796. })
  1797. expect(ui.commits.filter((commit) => commit.kind === "error" && commit.text === "provider failed")).toHaveLength(2)
  1798. await transport.close()
  1799. })
  1800. test("preserves an execution-only local error beside its projected prompt", async () => {
  1801. const events = feed()
  1802. events.push(connected())
  1803. const client = sdk({ streams: [events] })
  1804. const ui = footer()
  1805. const transport = await createSessionTransport({
  1806. sdk: client,
  1807. sessionID: "ses_1",
  1808. thinking: false,
  1809. replay: true,
  1810. footer: ui.api,
  1811. })
  1812. spyOn(client.message, "list").mockImplementation(() =>
  1813. ok({
  1814. data: [
  1815. {
  1816. id: "msg_prompt",
  1817. type: "user",
  1818. text: "hello",
  1819. files: [],
  1820. agents: [],
  1821. time: { created: 2 },
  1822. },
  1823. ],
  1824. cursor: {},
  1825. }),
  1826. )
  1827. await transport.replayOnResize({
  1828. localRows: () => [
  1829. {
  1830. commit: {
  1831. kind: "error",
  1832. source: "system",
  1833. text: "model unavailable",
  1834. phase: "start",
  1835. messageID: "msg_prompt",
  1836. },
  1837. },
  1838. ],
  1839. reset: async () => {},
  1840. })
  1841. expect(ui.commits.some((commit) => commit.kind === "error" && commit.text === "model unavailable")).toBe(true)
  1842. await transport.close()
  1843. })
  1844. test("scopes text and reasoning ordinals by assistant message", async () => {
  1845. const events = feed()
  1846. events.push(connected())
  1847. const client = sdk({ streams: [events] })
  1848. spyOn(client.message, "list").mockImplementation(() =>
  1849. ok({
  1850. data: [
  1851. {
  1852. id: "msg_b",
  1853. type: "assistant",
  1854. agent: "build",
  1855. model: { providerID: "test", id: "model" },
  1856. content: [
  1857. { type: "reasoning", text: "second thought" },
  1858. { type: "text", text: "second answer" },
  1859. ],
  1860. time: { created: 4, completed: 5 },
  1861. },
  1862. {
  1863. id: "msg_a",
  1864. type: "assistant",
  1865. agent: "build",
  1866. model: { providerID: "test", id: "model" },
  1867. content: [
  1868. { type: "reasoning", text: "first thought" },
  1869. { type: "text", text: "first answer" },
  1870. ],
  1871. time: { created: 2, completed: 3 },
  1872. },
  1873. ],
  1874. cursor: {},
  1875. }),
  1876. )
  1877. const ui = footer()
  1878. const transport = await createSessionTransport({
  1879. sdk: client,
  1880. sessionID: "ses_1",
  1881. thinking: true,
  1882. replay: true,
  1883. footer: ui.api,
  1884. })
  1885. expect(ui.commits.map((item) => item.text)).toEqual([
  1886. "Thinking: first thought",
  1887. "first answer",
  1888. "Thinking: second thought",
  1889. "second answer",
  1890. ])
  1891. await transport.close()
  1892. })
  1893. test("renders full reasoning when only the ended event is observed", async () => {
  1894. const events = feed()
  1895. events.push(connected())
  1896. const client = sdk({ streams: [events] })
  1897. const ui = footer()
  1898. const transport = await createSessionTransport({
  1899. sdk: client,
  1900. sessionID: "ses_1",
  1901. thinking: true,
  1902. footer: ui.api,
  1903. })
  1904. events.push({
  1905. id: "evt_reasoning",
  1906. created: 0,
  1907. type: "session.reasoning.ended",
  1908. durable: durable("ses_1"),
  1909. data: {
  1910. sessionID: "ses_1",
  1911. assistantMessageID: "msg_assistant",
  1912. ordinal: 0,
  1913. text: "considering",
  1914. },
  1915. })
  1916. await Bun.sleep(0)
  1917. expect(ui.commits.at(-1)?.text).toBe("Thinking: considering")
  1918. await transport.close()
  1919. })
  1920. test("tracks repeated root call IDs independently across assistant messages", async () => {
  1921. const events = feed()
  1922. events.push(connected())
  1923. const client = sdk({ streams: [events] })
  1924. const ui = footer()
  1925. const transport = await createSessionTransport({
  1926. sdk: client,
  1927. sessionID: "ses_1",
  1928. thinking: false,
  1929. footer: ui.api,
  1930. })
  1931. for (const [index, messageID] of ["msg_tool_one", "msg_tool_two"].entries()) {
  1932. events.push({
  1933. id: `evt_repeated_input_${index}`,
  1934. created: index * 3 + 1,
  1935. type: "session.tool.input.started",
  1936. durable: durable("ses_1", index * 3),
  1937. data: { sessionID: "ses_1", assistantMessageID: messageID, id: "call_repeated", name: "read" },
  1938. })
  1939. events.push({
  1940. id: `evt_repeated_called_${index}`,
  1941. created: index * 3 + 2,
  1942. type: "session.tool.called",
  1943. durable: durable("ses_1", index * 3 + 1),
  1944. data: {
  1945. sessionID: "ses_1",
  1946. assistantMessageID: messageID,
  1947. id: "call_repeated",
  1948. input: { path: `${index + 1}.txt` },
  1949. executed: true,
  1950. },
  1951. })
  1952. events.push({
  1953. id: `evt_repeated_success_${index}`,
  1954. created: index * 3 + 3,
  1955. type: "session.tool.success",
  1956. durable: durable("ses_1", index * 3 + 2, 2),
  1957. data: {
  1958. sessionID: "ses_1",
  1959. assistantMessageID: messageID,
  1960. id: "call_repeated",
  1961. metadata: {},
  1962. content: [{ type: "text", text: "" }],
  1963. executed: true,
  1964. },
  1965. })
  1966. }
  1967. await Bun.sleep(0)
  1968. const commits = ui.commits.filter((item) => item.part?.id === "call_repeated")
  1969. expect(commits.map((item) => [item.messageID, item.phase])).toEqual([
  1970. ["msg_tool_one", "start"],
  1971. ["msg_tool_one", "final"],
  1972. ["msg_tool_two", "start"],
  1973. ["msg_tool_two", "final"],
  1974. ])
  1975. expect(
  1976. commits
  1977. .filter((item) => item.phase === "final")
  1978. .map((item) => (item.part?.state.status === "streaming" ? undefined : item.part?.state.input)),
  1979. ).toEqual([{ path: "1.txt" }, { path: "2.txt" }])
  1980. await transport.close()
  1981. })
  1982. test("reduces root tool progress and preserves it on failure", async () => {
  1983. const events = feed()
  1984. events.push(connected())
  1985. const client = sdk({ streams: [events] })
  1986. const ui = footer()
  1987. const transport = await createSessionTransport({
  1988. sdk: client,
  1989. sessionID: "ses_1",
  1990. thinking: false,
  1991. footer: ui.api,
  1992. })
  1993. events.push({
  1994. id: "evt_progress_input",
  1995. created: 1,
  1996. type: "session.tool.input.started",
  1997. durable: durable("ses_1"),
  1998. data: {
  1999. sessionID: "ses_1",
  2000. assistantMessageID: "msg_progress",
  2001. id: "call_progress",
  2002. name: "shell",
  2003. },
  2004. })
  2005. events.push({
  2006. id: "evt_progress_called",
  2007. created: 2,
  2008. type: "session.tool.called",
  2009. durable: durable("ses_1", 1),
  2010. data: {
  2011. sessionID: "ses_1",
  2012. assistantMessageID: "msg_progress",
  2013. id: "call_progress",
  2014. input: { command: "printf partial && false" },
  2015. executed: true,
  2016. },
  2017. })
  2018. events.push({
  2019. id: "evt_progress",
  2020. created: 3,
  2021. type: "session.tool.progress",
  2022. data: {
  2023. sessionID: "ses_1",
  2024. assistantMessageID: "msg_progress",
  2025. id: "call_progress",
  2026. metadata: { checkpoint: 1 },
  2027. },
  2028. })
  2029. events.push({
  2030. id: "evt_progress_failed",
  2031. created: 4,
  2032. type: "session.tool.failed",
  2033. durable: durable("ses_1", 3, 2),
  2034. data: {
  2035. sessionID: "ses_1",
  2036. assistantMessageID: "msg_progress",
  2037. id: "call_progress",
  2038. error: { type: "unknown", message: "boom" },
  2039. metadata: { checkpoint: 1 },
  2040. content: [{ type: "text", text: "partial" }],
  2041. executed: true,
  2042. },
  2043. })
  2044. await Bun.sleep(0)
  2045. const commits = ui.commits.filter((item) => item.part?.id === "call_progress")
  2046. expect(commits.map((item) => [item.phase, item.text, item.toolState])).toEqual([
  2047. ["start", "running shell", "running"],
  2048. ["progress", "partial", "running"],
  2049. ["final", "boom", "error"],
  2050. ])
  2051. expect(commits.at(-1)?.part?.state).toMatchObject({
  2052. status: "error",
  2053. metadata: { checkpoint: 1 },
  2054. content: [{ type: "text", text: "partial" }],
  2055. })
  2056. await transport.close()
  2057. })
  2058. test("falls back to the default model when selecting a variant on a fresh session", async () => {
  2059. const events = feed()
  2060. events.push(connected())
  2061. const client = sdk({ streams: [events] })
  2062. const ui = footer()
  2063. const transport = await createSessionTransport({
  2064. sdk: client,
  2065. location: { directory: "/project", workspaceID: "wrk_1" },
  2066. sessionID: "ses_1",
  2067. thinking: false,
  2068. footer: ui.api,
  2069. })
  2070. spyOn(client.session, "get").mockImplementation(() => ok({ model: undefined }) as never)
  2071. const defaultModel = spyOn(client.model, "default").mockImplementation(
  2072. () =>
  2073. ok({
  2074. location: { directory: "/tmp", project: { id: "proj_1", directory: "/tmp" } },
  2075. data: { id: "gpt-5", providerID: "openai" },
  2076. }) as never,
  2077. )
  2078. const switched = spyOn(client.session, "switchModel").mockImplementation(() => ok(undefined))
  2079. let admitted = false
  2080. // The generated method has conditional return types for throwOnError; this mock represents the successful branch.
  2081. // @ts-expect-error successful SDK response is valid for both modes at runtime
  2082. spyOn(client.session, "prompt").mockImplementation((request) => {
  2083. admitted = true
  2084. return ok({ data: promptAdmission(request) })
  2085. })
  2086. const turn = transport.runPromptTurn({
  2087. agent: undefined,
  2088. model: undefined,
  2089. variant: "high",
  2090. prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
  2091. files: [],
  2092. includeFiles: true,
  2093. })
  2094. while (!admitted) await Bun.sleep(0)
  2095. events.push({
  2096. id: "evt_prompted",
  2097. created: 0,
  2098. type: "session.input.promoted",
  2099. durable: durable("ses_1"),
  2100. data: {
  2101. sessionID: "ses_1",
  2102. inputID: "msg_prompt",
  2103. },
  2104. })
  2105. events.push({
  2106. id: "evt_settled",
  2107. created: 0,
  2108. type: "session.execution.succeeded",
  2109. durable: durable("ses_1"),
  2110. data: { sessionID: "ses_1" },
  2111. })
  2112. await turn
  2113. expect(switched).toHaveBeenCalledWith(
  2114. { sessionID: "ses_1", model: { providerID: "openai", id: "gpt-5", variant: "high" } },
  2115. { signal: undefined },
  2116. )
  2117. expect(defaultModel).toHaveBeenCalledWith(
  2118. { location: { directory: "/project", workspace: "wrk_1" } },
  2119. { signal: undefined },
  2120. )
  2121. await transport.close()
  2122. })
  2123. test("interrupts the current Session when an active turn is aborted", async () => {
  2124. const events = feed()
  2125. events.push(connected())
  2126. const idle = defer()
  2127. const client = sdk({ streams: [events], wait: () => idle.promise })
  2128. const ui = footer()
  2129. const transport = await createSessionTransport({
  2130. sdk: client,
  2131. sessionID: "ses_1",
  2132. thinking: false,
  2133. footer: ui.api,
  2134. })
  2135. let admitted = false
  2136. // The generated method has conditional return types for throwOnError; this mock represents the successful branch.
  2137. // @ts-expect-error successful SDK response is valid for both modes at runtime
  2138. spyOn(client.session, "prompt").mockImplementation((request) => {
  2139. admitted = true
  2140. return ok({ data: promptAdmission(request) })
  2141. })
  2142. const interrupted = spyOn(client.session, "interrupt").mockImplementation(() => ok(undefined))
  2143. const controller = new AbortController()
  2144. const turn = transport.runPromptTurn({
  2145. agent: undefined,
  2146. model: undefined,
  2147. variant: undefined,
  2148. prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
  2149. files: [],
  2150. includeFiles: true,
  2151. signal: controller.signal,
  2152. })
  2153. while (!admitted) await Bun.sleep(0)
  2154. events.push({
  2155. id: "evt_prompted",
  2156. created: 0,
  2157. type: "session.input.promoted",
  2158. durable: durable("ses_1"),
  2159. data: {
  2160. sessionID: "ses_1",
  2161. inputID: "msg_prompt",
  2162. },
  2163. })
  2164. await Bun.sleep(0)
  2165. controller.abort()
  2166. idle.resolve()
  2167. await turn
  2168. expect(interrupted).toHaveBeenCalledWith({ sessionID: "ses_1" })
  2169. await transport.close()
  2170. })
  2171. test("runs a shell turn through v2.session.shell and renders live output", async () => {
  2172. const events = feed()
  2173. events.push(connected())
  2174. const client = sdk({ streams: [events] })
  2175. const ui = footer()
  2176. const transport = await createSessionTransport({
  2177. sdk: client,
  2178. sessionID: "ses_1",
  2179. thinking: false,
  2180. footer: ui.api,
  2181. })
  2182. let request: Parameters<OpenCodeClient["session"]["shell"]>[0] | undefined
  2183. spyOn(client.session, "shell").mockImplementation((input) => {
  2184. request = input
  2185. queueMicrotask(() => {
  2186. events.push({
  2187. id: input.id ?? "evt_missing",
  2188. created: 0,
  2189. type: "session.shell.started",
  2190. durable: durable("ses_1"),
  2191. data: {
  2192. sessionID: "ses_1",
  2193. shell: {
  2194. id: "sh_shell",
  2195. status: "running",
  2196. command: "ls",
  2197. cwd: "/tmp",
  2198. shell: "/bin/sh",
  2199. file: "/tmp/opencode-shell",
  2200. metadata: {},
  2201. time: { started: 0 },
  2202. },
  2203. },
  2204. })
  2205. events.push({
  2206. id: "evt_shell_end",
  2207. created: 0,
  2208. type: "session.shell.ended",
  2209. durable: durable("ses_1", 1),
  2210. data: {
  2211. sessionID: "ses_1",
  2212. shell: {
  2213. id: "sh_shell",
  2214. status: "exited",
  2215. command: "ls",
  2216. cwd: "/tmp",
  2217. shell: "/bin/sh",
  2218. file: "/tmp/opencode-shell",
  2219. exit: 0,
  2220. metadata: {},
  2221. time: { started: 0, completed: 1 },
  2222. },
  2223. output: { output: "file.txt", cursor: 8, size: 8, truncated: false },
  2224. },
  2225. })
  2226. })
  2227. return ok(undefined) as never
  2228. })
  2229. await transport.runPromptTurn({
  2230. agent: undefined,
  2231. model: undefined,
  2232. variant: undefined,
  2233. prompt: { text: "ls", parts: [], mode: "shell" },
  2234. files: [],
  2235. includeFiles: true,
  2236. })
  2237. expect(request).toMatchObject({ sessionID: "ses_1", command: "ls", id: expect.stringMatching(/^evt_/) })
  2238. expect(ui.commits.filter((item) => item.shell)).toMatchObject([
  2239. { phase: "start", partID: "shell:sh_shell", tool: "shell", toolState: "running", shell: { command: "ls" } },
  2240. {
  2241. phase: "progress",
  2242. partID: "shell:sh_shell",
  2243. text: "file.txt",
  2244. toolState: "completed",
  2245. shell: { command: "ls" },
  2246. },
  2247. ])
  2248. expect(ui.events).toContainEqual({ type: "stream.patch", patch: { phase: "running", status: "running shell" } })
  2249. await transport.close()
  2250. })
  2251. test("aborts an active shell turn without interrupting the session", async () => {
  2252. const events = feed()
  2253. events.push(connected())
  2254. const client = sdk({ streams: [events] })
  2255. const ui = footer()
  2256. const transport = await createSessionTransport({
  2257. sdk: client,
  2258. sessionID: "ses_1",
  2259. thinking: false,
  2260. footer: ui.api,
  2261. })
  2262. let started = false
  2263. let aborted = false
  2264. spyOn(client.session, "shell").mockImplementation(
  2265. (_input, options) =>
  2266. new Promise((_, reject) => {
  2267. started = true
  2268. options?.signal?.addEventListener("abort", () => {
  2269. aborted = true
  2270. reject(new Error("aborted"))
  2271. })
  2272. }) as never,
  2273. )
  2274. const interrupted = spyOn(client.session, "interrupt").mockImplementation(() => ok(undefined))
  2275. const turn = transport.runPromptTurn({
  2276. agent: undefined,
  2277. model: undefined,
  2278. variant: undefined,
  2279. prompt: { text: "sleep 100", parts: [], mode: "shell" },
  2280. files: [],
  2281. includeFiles: true,
  2282. })
  2283. while (!started) await Bun.sleep(0)
  2284. await transport.interruptActiveTurn()
  2285. await turn
  2286. expect(aborted).toBe(true)
  2287. expect(interrupted).not.toHaveBeenCalled()
  2288. await transport.close()
  2289. })
  2290. test("does not resolve an owned shell output wait from an unrelated shell", async () => {
  2291. const events = feed()
  2292. events.push(connected())
  2293. const client = sdk({ streams: [events] })
  2294. const ui = footer()
  2295. const transport = await createSessionTransport({
  2296. sdk: client,
  2297. sessionID: "ses_1",
  2298. thinking: false,
  2299. footer: ui.api,
  2300. })
  2301. let request: Parameters<OpenCodeClient["session"]["shell"]>[0] | undefined
  2302. let complete!: () => void
  2303. spyOn(client.session, "shell").mockImplementation((input) => {
  2304. request = input
  2305. return new Promise<void>((resolve) => {
  2306. complete = resolve
  2307. }) as never
  2308. })
  2309. let done = false
  2310. const turn = transport
  2311. .runPromptTurn({
  2312. agent: undefined,
  2313. model: undefined,
  2314. variant: undefined,
  2315. prompt: { text: "pwd", parts: [], mode: "shell" },
  2316. files: [],
  2317. includeFiles: true,
  2318. })
  2319. .then(() => {
  2320. done = true
  2321. })
  2322. while (!request) await Bun.sleep(0)
  2323. events.push({
  2324. id: "evt_unrelated_shell",
  2325. created: 0,
  2326. type: "session.shell.started",
  2327. durable: durable("ses_1"),
  2328. data: {
  2329. sessionID: "ses_1",
  2330. shell: {
  2331. id: "sh_unrelated",
  2332. status: "running",
  2333. command: "other",
  2334. cwd: "/tmp",
  2335. shell: "/bin/sh",
  2336. file: "/tmp/unrelated",
  2337. metadata: {},
  2338. time: { started: 0 },
  2339. },
  2340. },
  2341. })
  2342. events.push({
  2343. id: "evt_unrelated_end",
  2344. created: 0,
  2345. type: "session.shell.ended",
  2346. durable: durable("ses_1", 1),
  2347. data: {
  2348. sessionID: "ses_1",
  2349. shell: {
  2350. id: "sh_unrelated",
  2351. status: "exited",
  2352. command: "other",
  2353. cwd: "/tmp",
  2354. shell: "/bin/sh",
  2355. file: "/tmp/unrelated",
  2356. exit: 0,
  2357. metadata: {},
  2358. time: { started: 0, completed: 1 },
  2359. },
  2360. output: { output: "wrong", cursor: 5, size: 5, truncated: false },
  2361. },
  2362. })
  2363. await Bun.sleep(0)
  2364. complete()
  2365. await Bun.sleep(0)
  2366. expect(done).toBe(false)
  2367. events.push({
  2368. id: request.id ?? "evt_missing",
  2369. created: 0,
  2370. type: "session.shell.started",
  2371. durable: durable("ses_1", 2),
  2372. data: {
  2373. sessionID: "ses_1",
  2374. shell: {
  2375. id: "sh_owned",
  2376. status: "running",
  2377. command: "pwd",
  2378. cwd: "/tmp",
  2379. shell: "/bin/sh",
  2380. file: "/tmp/owned",
  2381. metadata: {},
  2382. time: { started: 0 },
  2383. },
  2384. },
  2385. })
  2386. events.push({
  2387. id: "evt_owned_end",
  2388. created: 0,
  2389. type: "session.shell.ended",
  2390. durable: durable("ses_1", 3),
  2391. data: {
  2392. sessionID: "ses_1",
  2393. shell: {
  2394. id: "sh_owned",
  2395. status: "exited",
  2396. command: "pwd",
  2397. cwd: "/tmp",
  2398. shell: "/bin/sh",
  2399. file: "/tmp/owned",
  2400. exit: 0,
  2401. metadata: {},
  2402. time: { started: 0, completed: 1 },
  2403. },
  2404. output: { output: "/tmp", cursor: 4, size: 4, truncated: false },
  2405. },
  2406. })
  2407. await turn
  2408. expect(request.id).toMatch(/^evt_/)
  2409. expect(ui.commits.some((item) => item.partID === "shell:sh_owned" && item.text === "/tmp")).toBe(true)
  2410. await transport.close()
  2411. })
  2412. test("hydrates projected shell transcripts once and dedupes live redelivery", async () => {
  2413. const events = feed()
  2414. events.push(connected())
  2415. const client = sdk({
  2416. streams: [events],
  2417. messages: {
  2418. ses_1: [
  2419. {
  2420. id: "msg_shell",
  2421. type: "shell" as const,
  2422. shellID: "sh_1",
  2423. status: "exited",
  2424. command: "ls",
  2425. exit: 0,
  2426. output: { output: "file.txt", cursor: 8, size: 8, truncated: false },
  2427. time: { created: 1, completed: 2 },
  2428. },
  2429. ],
  2430. },
  2431. })
  2432. const ui = footer()
  2433. const transport = await createSessionTransport({
  2434. sdk: client,
  2435. sessionID: "ses_1",
  2436. thinking: false,
  2437. replay: true,
  2438. footer: ui.api,
  2439. })
  2440. events.push({
  2441. id: "evt_shell_end",
  2442. created: 0,
  2443. type: "session.shell.ended",
  2444. durable: durable("ses_1", 1),
  2445. data: {
  2446. sessionID: "ses_1",
  2447. shell: {
  2448. id: "sh_1",
  2449. status: "exited",
  2450. command: "ls",
  2451. cwd: "/tmp",
  2452. shell: "/bin/sh",
  2453. file: "/tmp/opencode-shell",
  2454. exit: 0,
  2455. metadata: {},
  2456. time: { started: 0, completed: 1 },
  2457. },
  2458. output: { output: "file.txt", cursor: 8, size: 8, truncated: false },
  2459. },
  2460. })
  2461. await Bun.sleep(0)
  2462. await Bun.sleep(0)
  2463. expect(ui.commits.filter((item) => item.shell)).toMatchObject([
  2464. { phase: "start", partID: "shell:sh_1", shell: { command: "ls" } },
  2465. { phase: "progress", partID: "shell:sh_1", text: "file.txt", toolState: "completed" },
  2466. ])
  2467. await transport.close()
  2468. })
  2469. test("renders failed projected shells as errors and marks truncated live output", async () => {
  2470. const events = feed()
  2471. events.push(connected())
  2472. const client = sdk({
  2473. streams: [events],
  2474. messages: {
  2475. ses_1: [
  2476. {
  2477. id: "msg_failed_shell",
  2478. type: "shell" as const,
  2479. shellID: "sh_failed",
  2480. status: "exited",
  2481. command: "false",
  2482. exit: 7,
  2483. output: { output: "failure output", cursor: 14, size: 14, truncated: false },
  2484. time: { created: 1, completed: 2 },
  2485. },
  2486. ],
  2487. },
  2488. })
  2489. const ui = footer()
  2490. const transport = await createSessionTransport({
  2491. sdk: client,
  2492. sessionID: "ses_1",
  2493. thinking: false,
  2494. replay: true,
  2495. footer: ui.api,
  2496. })
  2497. events.push({
  2498. id: "evt_truncated_start",
  2499. created: 0,
  2500. type: "session.shell.started",
  2501. durable: durable("ses_1"),
  2502. data: {
  2503. sessionID: "ses_1",
  2504. shell: {
  2505. id: "sh_truncated",
  2506. status: "running",
  2507. command: "long",
  2508. cwd: "/tmp",
  2509. shell: "/bin/sh",
  2510. file: "/tmp/truncated",
  2511. metadata: {},
  2512. time: { started: 0 },
  2513. },
  2514. },
  2515. })
  2516. events.push({
  2517. id: "evt_truncated_end",
  2518. created: 0,
  2519. type: "session.shell.ended",
  2520. durable: durable("ses_1", 1),
  2521. data: {
  2522. sessionID: "ses_1",
  2523. shell: {
  2524. id: "sh_truncated",
  2525. status: "exited",
  2526. command: "long",
  2527. cwd: "/tmp",
  2528. shell: "/bin/sh",
  2529. file: "/tmp/truncated",
  2530. exit: 0,
  2531. metadata: {},
  2532. time: { started: 0, completed: 1 },
  2533. },
  2534. output: { output: "partial", cursor: 7, size: 20, truncated: false },
  2535. },
  2536. })
  2537. await Bun.sleep(0)
  2538. expect(ui.commits).toContainEqual(
  2539. expect.objectContaining({ toolState: "error", toolError: "Shell exited with code 7" }),
  2540. )
  2541. expect(ui.commits).toContainEqual(expect.objectContaining({ text: "partial\n[output truncated]" }))
  2542. await transport.close()
  2543. })
  2544. test("routes command prompts through v2.session.command", async () => {
  2545. const events = feed()
  2546. events.push(connected())
  2547. const client = sdk({ streams: [events] })
  2548. const ui = footer()
  2549. const transport = await createSessionTransport({
  2550. sdk: client,
  2551. sessionID: "ses_1",
  2552. thinking: false,
  2553. footer: ui.api,
  2554. })
  2555. let request: Parameters<OpenCodeClient["session"]["command"]>[0] | undefined
  2556. spyOn(client.session, "command").mockImplementation((input) => {
  2557. request = input
  2558. queueMicrotask(() => {
  2559. events.push({
  2560. id: "evt_prompted",
  2561. created: 0,
  2562. type: "session.input.promoted",
  2563. durable: durable("ses_1"),
  2564. data: {
  2565. sessionID: "ses_1",
  2566. inputID: "msg_cmd",
  2567. },
  2568. })
  2569. events.push({
  2570. id: "evt_settled",
  2571. created: 0,
  2572. type: "session.execution.succeeded",
  2573. durable: durable("ses_1"),
  2574. data: { sessionID: "ses_1" },
  2575. })
  2576. })
  2577. return ok({
  2578. id: input.id ?? "msg_cmd",
  2579. sessionID: "ses_1",
  2580. type: "user" as const,
  2581. data: { text: "evaluated template" },
  2582. delivery: "steer" as const,
  2583. timeCreated: 2,
  2584. })
  2585. })
  2586. await transport.runPromptTurn({
  2587. agent: "build",
  2588. model: { providerID: "test", modelID: "model" },
  2589. variant: undefined,
  2590. prompt: {
  2591. messageID: "msg_cmd",
  2592. text: "/deploy prod",
  2593. parts: [
  2594. {
  2595. type: "file",
  2596. url: "file:///tmp/mentioned.txt",
  2597. filename: "mentioned.txt",
  2598. source: { type: "file", text: { start: 8, end: 12, value: "prod" } },
  2599. },
  2600. ],
  2601. command: { name: "deploy", arguments: "prod" },
  2602. },
  2603. files: [{ type: "file", url: "file:///tmp/context.txt", filename: "context.txt", mime: "text/plain" }],
  2604. includeFiles: true,
  2605. })
  2606. expect(request).toMatchObject({
  2607. sessionID: "ses_1",
  2608. id: "msg_cmd",
  2609. command: "deploy",
  2610. arguments: "prod",
  2611. agent: "build",
  2612. model: { providerID: "test", id: "model" },
  2613. files: [
  2614. { uri: "file:///tmp/context.txt", name: "context.txt" },
  2615. {
  2616. uri: "file:///tmp/mentioned.txt",
  2617. name: "mentioned.txt",
  2618. mention: { start: 8, end: 12, text: "prod" },
  2619. },
  2620. ],
  2621. delivery: "steer",
  2622. })
  2623. // Selection rides the command payload; no separate client-side switch.
  2624. expect(client.session.switchAgent).not.toHaveBeenCalled()
  2625. expect(client.session.switchModel).not.toHaveBeenCalled()
  2626. await transport.close()
  2627. })
  2628. test("routes skill prompts through v2.session.skill and settles without promotion", async () => {
  2629. const events = feed()
  2630. events.push(connected())
  2631. const client = sdk({ streams: [events] })
  2632. const ui = footer()
  2633. const transport = await createSessionTransport({
  2634. sdk: client,
  2635. sessionID: "ses_1",
  2636. thinking: false,
  2637. footer: ui.api,
  2638. })
  2639. let request: Parameters<OpenCodeClient["session"]["skill"]>[0] | undefined
  2640. const command = spyOn(client.session, "command")
  2641. const prompt = spyOn(client.session, "prompt")
  2642. spyOn(client.session, "skill").mockImplementation((input) => {
  2643. request = input
  2644. queueMicrotask(() => {
  2645. events.push({
  2646. id: "evt_skill",
  2647. created: 0,
  2648. type: "session.skill.activated",
  2649. durable: durable("ses_1"),
  2650. data: {
  2651. sessionID: "ses_1",
  2652. id: input.skill ?? "tigerstyle",
  2653. name: input.skill ?? "tigerstyle",
  2654. text: "skill instructions",
  2655. },
  2656. })
  2657. events.push({
  2658. id: "evt_settled",
  2659. created: 0,
  2660. type: "session.execution.succeeded",
  2661. durable: durable("ses_1"),
  2662. data: { sessionID: "ses_1" },
  2663. })
  2664. })
  2665. return ok(undefined) as never
  2666. })
  2667. await transport.runPromptTurn({
  2668. agent: "review",
  2669. model: undefined,
  2670. variant: undefined,
  2671. prompt: {
  2672. messageID: "msg_skill",
  2673. text: "/tigerstyle",
  2674. parts: [],
  2675. command: { name: "tigerstyle", arguments: "", source: "skill" },
  2676. },
  2677. files: [],
  2678. includeFiles: true,
  2679. })
  2680. expect(client.session.switchAgent).toHaveBeenCalledWith({ sessionID: "ses_1", agent: "review" }, expect.anything())
  2681. expect(request).toMatchObject({ sessionID: "ses_1", id: "msg_skill", skill: "tigerstyle" })
  2682. expect(command).not.toHaveBeenCalled()
  2683. expect(prompt).not.toHaveBeenCalled()
  2684. expect(ui.commits).toContainEqual(
  2685. expect.objectContaining({ kind: "system", text: '→ Skill "tigerstyle"', messageID: "msg_skill" }),
  2686. )
  2687. await transport.close()
  2688. })
  2689. test("refreshes catalogs on connection and location-scoped invalidations", async () => {
  2690. const events = feed()
  2691. events.push(connected())
  2692. const client = sdk({ streams: [events] })
  2693. const ui = footer()
  2694. let refreshes = 0
  2695. const transport = await createSessionTransport({
  2696. sdk: client,
  2697. location: {
  2698. directory: "/project",
  2699. workspaceID: "work-1",
  2700. },
  2701. sessionID: "ses_1",
  2702. thinking: false,
  2703. footer: ui.api,
  2704. onCatalogRefresh: () => refreshes++,
  2705. })
  2706. expect(refreshes).toBe(1)
  2707. for (const type of [
  2708. "catalog.updated",
  2709. "integration.updated",
  2710. "agent.updated",
  2711. "command.updated",
  2712. "skill.updated",
  2713. "reference.updated",
  2714. ] as const)
  2715. events.push({
  2716. id: `evt_${type}`,
  2717. created: 0,
  2718. type,
  2719. location: { directory: "/project", workspaceID: "work-1" },
  2720. data: {},
  2721. })
  2722. events.push({
  2723. id: "evt_foreign_catalog",
  2724. created: 0,
  2725. type: "catalog.updated",
  2726. location: { directory: "/other" },
  2727. data: {},
  2728. })
  2729. events.push({
  2730. id: "evt_foreign_workspace_catalog",
  2731. created: 0,
  2732. type: "catalog.updated",
  2733. location: { directory: "/project", workspaceID: "work-2" },
  2734. data: {},
  2735. })
  2736. while (refreshes < 7) await Bun.sleep(0)
  2737. await Bun.sleep(0)
  2738. expect(refreshes).toBe(7)
  2739. await transport.close()
  2740. })
  2741. test("hydrates skill activation messages once and dedupes live redelivery", async () => {
  2742. const events = feed()
  2743. events.push(connected())
  2744. const client = sdk({
  2745. streams: [events],
  2746. messages: {
  2747. ses_1: [
  2748. {
  2749. id: "msg_skill",
  2750. type: "skill" as const,
  2751. skill: "tigerstyle",
  2752. name: "tigerstyle",
  2753. text: "skill instructions",
  2754. time: { created: 2 },
  2755. },
  2756. ],
  2757. },
  2758. })
  2759. const ui = footer()
  2760. const transport = await createSessionTransport({
  2761. sdk: client,
  2762. sessionID: "ses_1",
  2763. thinking: false,
  2764. replay: true,
  2765. footer: ui.api,
  2766. })
  2767. events.push({
  2768. id: "evt_skill",
  2769. created: 0,
  2770. type: "session.skill.activated",
  2771. durable: durable("ses_1"),
  2772. data: {
  2773. sessionID: "ses_1",
  2774. id: "tigerstyle",
  2775. name: "tigerstyle",
  2776. text: "skill instructions",
  2777. },
  2778. })
  2779. await Bun.sleep(0)
  2780. await Bun.sleep(0)
  2781. expect(ui.commits.filter((item) => item.text === '→ Skill "tigerstyle"')).toHaveLength(1)
  2782. await transport.close()
  2783. })
  2784. test("discovers a subagent from its terminal failure snapshot", async () => {
  2785. const events = feed()
  2786. events.push(connected())
  2787. const client = sdk({ streams: [events], messages: { ses_child_failed: [] } })
  2788. const ui = footer()
  2789. const transport = await createSessionTransport({
  2790. sdk: client,
  2791. sessionID: "ses_1",
  2792. thinking: false,
  2793. footer: ui.api,
  2794. })
  2795. const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  2796. events.push({
  2797. id: "evt_failed_subagent_input",
  2798. created: 1,
  2799. type: "session.tool.input.started",
  2800. durable: durable("ses_1"),
  2801. data: {
  2802. sessionID: "ses_1",
  2803. assistantMessageID: "msg_failed_subagent",
  2804. id: "call_failed_subagent",
  2805. name: "subagent",
  2806. },
  2807. })
  2808. events.push({
  2809. id: "evt_failed_subagent_called",
  2810. created: 2,
  2811. type: "session.tool.called",
  2812. durable: durable("ses_1", 1),
  2813. data: {
  2814. sessionID: "ses_1",
  2815. assistantMessageID: "msg_failed_subagent",
  2816. id: "call_failed_subagent",
  2817. input: { agent: "explore", description: "Inspect failure", prompt: "inspect" },
  2818. executed: true,
  2819. },
  2820. })
  2821. events.push({
  2822. id: "evt_failed_subagent",
  2823. created: 3,
  2824. type: "session.tool.failed",
  2825. durable: durable("ses_1", 2, 2),
  2826. data: {
  2827. sessionID: "ses_1",
  2828. assistantMessageID: "msg_failed_subagent",
  2829. id: "call_failed_subagent",
  2830. error: { type: "unknown", message: "subagent failed" },
  2831. metadata: { sessionID: "ses_child_failed", status: "running" },
  2832. executed: true,
  2833. },
  2834. })
  2835. while (!states().some((state) => state.tabs.some((tab) => tab.sessionID === "ses_child_failed"))) await Bun.sleep(0)
  2836. expect(states().at(-1)?.tabs).toMatchObject([
  2837. {
  2838. sessionID: "ses_child_failed",
  2839. label: "Explore",
  2840. description: "Inspect failure",
  2841. },
  2842. ])
  2843. await transport.close()
  2844. })
  2845. test("discovers current subagents from progress and reduces descendant tool state", async () => {
  2846. const events = feed()
  2847. events.push(connected())
  2848. const client = sdk({ streams: [events], messages: { ses_child_progress: [] } })
  2849. const ui = footer()
  2850. const transport = await createSessionTransport({
  2851. sdk: client,
  2852. sessionID: "ses_1",
  2853. thinking: false,
  2854. footer: ui.api,
  2855. })
  2856. const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  2857. events.push({
  2858. id: "evt_subagent_input",
  2859. created: 1,
  2860. type: "session.tool.input.started",
  2861. durable: durable("ses_1"),
  2862. data: {
  2863. sessionID: "ses_1",
  2864. assistantMessageID: "msg_subagent",
  2865. id: "call_subagent",
  2866. name: "subagent",
  2867. },
  2868. })
  2869. events.push({
  2870. id: "evt_subagent_called",
  2871. created: 2,
  2872. type: "session.tool.called",
  2873. durable: durable("ses_1", 1),
  2874. data: {
  2875. sessionID: "ses_1",
  2876. assistantMessageID: "msg_subagent",
  2877. id: "call_subagent",
  2878. input: { agent: "explore", description: "Inspect progress", prompt: "inspect" },
  2879. executed: true,
  2880. },
  2881. })
  2882. events.push({
  2883. id: "evt_subagent_progress",
  2884. created: 3,
  2885. type: "session.tool.progress",
  2886. data: {
  2887. sessionID: "ses_1",
  2888. assistantMessageID: "msg_subagent",
  2889. id: "call_subagent",
  2890. metadata: { sessionID: "ses_child_progress", status: "running" },
  2891. },
  2892. })
  2893. while (!states().some((state) => state.tabs.some((tab) => tab.sessionID === "ses_child_progress")))
  2894. await Bun.sleep(0)
  2895. expect(states().at(-1)?.tabs).toMatchObject([
  2896. {
  2897. sessionID: "ses_child_progress",
  2898. label: "Explore",
  2899. description: "Inspect progress",
  2900. status: "running",
  2901. background: undefined,
  2902. },
  2903. ])
  2904. transport.selectSubagent("ses_child_progress")
  2905. while (!states().at(-1)?.details.ses_child_progress) await Bun.sleep(0)
  2906. events.push({
  2907. id: "evt_child_tool_input",
  2908. created: 4,
  2909. type: "session.tool.input.started",
  2910. durable: durable("ses_child_progress"),
  2911. data: {
  2912. sessionID: "ses_child_progress",
  2913. assistantMessageID: "msg_child_tool",
  2914. id: "call_child_shell",
  2915. name: "shell",
  2916. },
  2917. })
  2918. events.push({
  2919. id: "evt_child_tool_called",
  2920. created: 5,
  2921. type: "session.tool.called",
  2922. durable: durable("ses_child_progress", 1),
  2923. data: {
  2924. sessionID: "ses_child_progress",
  2925. assistantMessageID: "msg_child_tool",
  2926. id: "call_child_shell",
  2927. input: { command: "printf child && false" },
  2928. executed: true,
  2929. },
  2930. })
  2931. events.push({
  2932. id: "evt_child_tool_progress",
  2933. created: 6,
  2934. type: "session.tool.progress",
  2935. data: {
  2936. sessionID: "ses_child_progress",
  2937. assistantMessageID: "msg_child_tool",
  2938. id: "call_child_shell",
  2939. metadata: { checkpoint: "child" },
  2940. },
  2941. })
  2942. events.push({
  2943. id: "evt_child_permission",
  2944. created: 7,
  2945. type: "permission.asked",
  2946. data: {
  2947. id: "per_child",
  2948. sessionID: "ses_child_progress",
  2949. action: "shell",
  2950. resources: ["printf child && false"],
  2951. source: { type: "tool", messageID: "msg_child_tool", id: "call_child_shell" },
  2952. },
  2953. })
  2954. events.push({
  2955. id: "evt_child_tool_failed",
  2956. created: 8,
  2957. type: "session.tool.failed",
  2958. durable: durable("ses_child_progress", 3, 2),
  2959. data: {
  2960. sessionID: "ses_child_progress",
  2961. assistantMessageID: "msg_child_tool",
  2962. id: "call_child_shell",
  2963. error: { type: "unknown", message: "child boom" },
  2964. metadata: { checkpoint: "child" },
  2965. content: [{ type: "text", text: "child partial" }],
  2966. executed: true,
  2967. },
  2968. })
  2969. while (
  2970. !states()
  2971. .at(-1)
  2972. ?.details.ses_child_progress?.commits.some(
  2973. (item) => item.part?.id === "call_child_shell" && item.toolState === "error",
  2974. )
  2975. )
  2976. await Bun.sleep(0)
  2977. const commits = states().at(-1)?.details.ses_child_progress?.commits ?? []
  2978. expect(
  2979. commits
  2980. .filter((item) => item.part?.id === "call_child_shell")
  2981. .map((item) => [item.phase, item.text, item.toolState]),
  2982. ).toEqual([
  2983. ["progress", "child partial", "running"],
  2984. ["final", "child boom", "error"],
  2985. ])
  2986. expect(
  2987. commits.find((item) => item.part?.id === "call_child_shell" && item.toolState === "error")?.part?.state,
  2988. ).toMatchObject({
  2989. status: "error",
  2990. metadata: { checkpoint: "child" },
  2991. content: [{ type: "text", text: "child partial" }],
  2992. })
  2993. expect(
  2994. ui.events.find(
  2995. (event) =>
  2996. event.type === "stream.view" && event.view.type === "permission" && event.view.request.id === "per_child",
  2997. ),
  2998. ).toMatchObject({
  2999. view: {
  3000. request: {
  3001. sessionID: "ses_child_progress",
  3002. tool: { id: "call_child_shell", state: { input: { command: "printf child && false" } } },
  3003. },
  3004. },
  3005. })
  3006. await transport.close()
  3007. })
  3008. test("discovers a live child session and tracks its tab and selected detail", async () => {
  3009. const events = feed()
  3010. events.push(connected())
  3011. const client = sdk({
  3012. streams: [events],
  3013. messages: {
  3014. ses_child: [
  3015. {
  3016. id: "msg_task",
  3017. type: "user" as const,
  3018. text: "task prompt",
  3019. files: [],
  3020. agents: [],
  3021. time: { created: 1 },
  3022. },
  3023. {
  3024. id: "msg_child_a",
  3025. type: "assistant" as const,
  3026. agent: "explore",
  3027. model: { providerID: "test", id: "model" },
  3028. content: [{ type: "text" as const, text: "child answer" }],
  3029. time: { created: 2 },
  3030. },
  3031. ],
  3032. },
  3033. })
  3034. spyOn(client.session, "get").mockImplementation(
  3035. () =>
  3036. ok({
  3037. id: "ses_child",
  3038. parentID: "ses_1",
  3039. projectID: "proj_1",
  3040. agent: "explore",
  3041. cost: 0,
  3042. tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
  3043. time: { created: 1, updated: 1 },
  3044. title: "Find files",
  3045. location: { directory: "/tmp" },
  3046. }) as never,
  3047. )
  3048. const ui = footer()
  3049. const transport = await createSessionTransport({
  3050. sdk: client,
  3051. sessionID: "ses_1",
  3052. thinking: false,
  3053. footer: ui.api,
  3054. })
  3055. const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  3056. transport.selectSubagent("ses_child")
  3057. events.push({
  3058. id: "evt_child_step",
  3059. created: 0,
  3060. type: "session.step.started",
  3061. durable: durable("ses_child"),
  3062. data: {
  3063. sessionID: "ses_child",
  3064. assistantMessageID: "msg_child_a",
  3065. agent: "explore",
  3066. model: { providerID: "test", id: "model" },
  3067. },
  3068. })
  3069. while (!states().some((state) => state.details.ses_child?.commits.some((item) => item.text === "task prompt")))
  3070. await Bun.sleep(0)
  3071. expect(states().at(-1)?.tabs).toMatchObject([
  3072. { sessionID: "ses_child", label: "Explore", title: "Find files", status: "running" },
  3073. ])
  3074. expect(
  3075. states()
  3076. .at(-1)
  3077. ?.details.ses_child?.commits.filter((item) => item.text === "child answer"),
  3078. ).toHaveLength(1)
  3079. events.push({
  3080. id: "evt_child_text_replayed",
  3081. created: 0,
  3082. type: "session.text.delta",
  3083. data: {
  3084. sessionID: "ses_child",
  3085. assistantMessageID: "msg_child_a",
  3086. ordinal: 0,
  3087. delta: "answer",
  3088. },
  3089. })
  3090. await Bun.sleep(0)
  3091. expect(
  3092. states()
  3093. .at(-1)
  3094. ?.details.ses_child?.commits.filter((item) => item.text === "child answer"),
  3095. ).toHaveLength(1)
  3096. events.push({
  3097. id: "evt_child_text_suffix",
  3098. created: 0,
  3099. type: "session.text.delta",
  3100. data: {
  3101. sessionID: "ses_child",
  3102. assistantMessageID: "msg_child_a",
  3103. ordinal: 0,
  3104. delta: " suffix",
  3105. },
  3106. })
  3107. while (
  3108. !states().some((state) => state.details.ses_child?.commits.some((item) => item.text === "child answer suffix"))
  3109. )
  3110. await Bun.sleep(0)
  3111. events.push({
  3112. id: "evt_child_settled",
  3113. created: 0,
  3114. type: "session.execution.succeeded",
  3115. durable: durable("ses_child"),
  3116. data: { sessionID: "ses_child" },
  3117. })
  3118. while (!states().some((state) => state.tabs.some((tab) => tab.status === "completed"))) await Bun.sleep(0)
  3119. await transport.close()
  3120. })
  3121. test("reveals an admitted child prompt only when it is promoted after hydration", async () => {
  3122. const events = feed()
  3123. events.push(connected())
  3124. const client = sdk({
  3125. streams: [events],
  3126. messages: { ses_child: [] },
  3127. sessions: [{ id: "ses_child", parentID: "ses_1", time: { updated: 1 } }],
  3128. })
  3129. const ui = footer()
  3130. const transport = await createSessionTransport({
  3131. sdk: client,
  3132. sessionID: "ses_1",
  3133. thinking: false,
  3134. footer: ui.api,
  3135. })
  3136. const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  3137. transport.selectSubagent("ses_child")
  3138. while (!states().some((state) => state.details.ses_child)) await Bun.sleep(0)
  3139. events.push({
  3140. id: "evt_child_admitted",
  3141. created: 1,
  3142. type: "session.input.admitted",
  3143. durable: durable("ses_child"),
  3144. data: {
  3145. sessionID: "ses_child",
  3146. inputID: "msg_child_prompt",
  3147. input: { type: "user", data: { text: "actual child prompt" }, delivery: "steer" },
  3148. },
  3149. })
  3150. await Bun.sleep(0)
  3151. expect(
  3152. states()
  3153. .at(-1)
  3154. ?.details.ses_child?.commits.some((item) => item.messageID === "msg_child_prompt"),
  3155. ).toBe(false)
  3156. events.push({
  3157. id: "evt_child_promoted",
  3158. created: 2,
  3159. type: "session.input.promoted",
  3160. durable: durable("ses_child", 1),
  3161. data: { sessionID: "ses_child", inputID: "msg_child_prompt" },
  3162. })
  3163. while (
  3164. !states()
  3165. .at(-1)
  3166. ?.details.ses_child?.commits.some(
  3167. (item) => item.messageID === "msg_child_prompt" && item.text === "actual child prompt",
  3168. )
  3169. )
  3170. await Bun.sleep(0)
  3171. await transport.close()
  3172. })
  3173. test("preserves a pre-hydration admission promoted during stale hydration", async () => {
  3174. const events = feed()
  3175. events.push(connected())
  3176. const client = sdk({
  3177. streams: [events],
  3178. sessions: [{ id: "ses_child", parentID: "ses_1", time: { updated: 1 } }],
  3179. })
  3180. let childHydrating = false
  3181. let releaseHydration!: () => void
  3182. const hydration = new Promise<void>((resolve) => {
  3183. releaseHydration = resolve
  3184. })
  3185. spyOn(client.message, "list").mockImplementation(async (request) => {
  3186. if (request.sessionID === "ses_child") {
  3187. childHydrating = true
  3188. await hydration
  3189. }
  3190. return ok({ data: [], cursor: {} })
  3191. })
  3192. const ui = footer()
  3193. const transport = await createSessionTransport({
  3194. sdk: client,
  3195. sessionID: "ses_1",
  3196. thinking: false,
  3197. footer: ui.api,
  3198. })
  3199. const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  3200. events.push({
  3201. id: "evt_child_admitted_race",
  3202. created: 1,
  3203. type: "session.input.admitted",
  3204. durable: durable("ses_child"),
  3205. data: {
  3206. sessionID: "ses_child",
  3207. inputID: "msg_child_race",
  3208. input: { type: "user", data: { text: "prompt admitted before hydration" }, delivery: "steer" },
  3209. },
  3210. })
  3211. await Bun.sleep(0)
  3212. transport.selectSubagent("ses_child")
  3213. while (!childHydrating) await Bun.sleep(0)
  3214. events.push({
  3215. id: "evt_child_promoted_race",
  3216. created: 2,
  3217. type: "session.input.promoted",
  3218. durable: durable("ses_child", 1),
  3219. data: { sessionID: "ses_child", inputID: "msg_child_race" },
  3220. })
  3221. await Bun.sleep(0)
  3222. releaseHydration()
  3223. await Bun.sleep(0)
  3224. await Bun.sleep(0)
  3225. while (
  3226. !states()
  3227. .at(-1)
  3228. ?.details.ses_child?.commits.some(
  3229. (item) => item.messageID === "msg_child_race" && item.text === "prompt admitted before hydration",
  3230. )
  3231. )
  3232. await Bun.sleep(0)
  3233. await transport.close()
  3234. })
  3235. test("retries child hydration after a bounded live-event overflow", async () => {
  3236. const events = feed()
  3237. events.push(connected())
  3238. const client = sdk({
  3239. streams: [events],
  3240. sessions: [{ id: "ses_child", parentID: "ses_1", time: { updated: 1 } }],
  3241. })
  3242. let childRequests = 0
  3243. let releaseStale!: () => void
  3244. let releaseRetry!: () => void
  3245. const stale = new Promise<void>((resolve) => {
  3246. releaseStale = resolve
  3247. })
  3248. const retry = new Promise<void>((resolve) => {
  3249. releaseRetry = resolve
  3250. })
  3251. spyOn(client.message, "list").mockImplementation(async (request) => {
  3252. if (request.sessionID !== "ses_child") return ok({ data: [], cursor: {} })
  3253. childRequests++
  3254. if (childRequests === 1) {
  3255. await stale
  3256. return ok({ data: [], cursor: {} })
  3257. }
  3258. await retry
  3259. return ok({
  3260. data: [
  3261. {
  3262. id: "msg_overflow_assistant",
  3263. type: "assistant" as const,
  3264. agent: "explore",
  3265. model: { providerID: "test", id: "model" },
  3266. content: [{ type: "text" as const, id: "txt_overflow_64", text: "live 64" }],
  3267. time: { created: 2, completed: 3 },
  3268. },
  3269. {
  3270. id: "msg_overflow_baseline",
  3271. type: "user" as const,
  3272. text: "baseline history",
  3273. files: [],
  3274. agents: [],
  3275. time: { created: 1 },
  3276. },
  3277. ],
  3278. cursor: {},
  3279. })
  3280. })
  3281. const ui = footer()
  3282. const transport = await createSessionTransport({
  3283. sdk: client,
  3284. sessionID: "ses_1",
  3285. thinking: false,
  3286. footer: ui.api,
  3287. })
  3288. const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  3289. transport.selectSubagent("ses_child")
  3290. while (childRequests < 1) await Bun.sleep(0)
  3291. for (let index = 0; index < 65; index++)
  3292. events.push({
  3293. id: `evt_overflow_${index}`,
  3294. created: index,
  3295. type: "session.text.delta",
  3296. data: {
  3297. sessionID: "ses_child",
  3298. assistantMessageID: "msg_overflow_assistant",
  3299. ordinal: index,
  3300. delta: `live ${index}`,
  3301. },
  3302. })
  3303. while (
  3304. !states()
  3305. .at(-1)
  3306. ?.details.ses_child?.commits.some((item) => item.text === "live 64")
  3307. )
  3308. await Bun.sleep(0)
  3309. releaseStale()
  3310. while (childRequests < 2) await Bun.sleep(0)
  3311. expect(
  3312. states()
  3313. .at(-1)
  3314. ?.details.ses_child?.commits.some((item) => item.text === "live 64"),
  3315. ).toBe(true)
  3316. releaseRetry()
  3317. while (
  3318. !states()
  3319. .at(-1)
  3320. ?.details.ses_child?.commits.some((item) => item.text === "baseline history")
  3321. )
  3322. await Bun.sleep(0)
  3323. expect(
  3324. states()
  3325. .at(-1)
  3326. ?.details.ses_child?.commits.some((item) => item.text === "live 64"),
  3327. ).toBe(true)
  3328. expect(childRequests).toBe(2)
  3329. await transport.close()
  3330. })
  3331. test("reconciles pre-hydration tool metadata without downgrading projected completion", async () => {
  3332. const events = feed()
  3333. events.push(connected())
  3334. const client = sdk({
  3335. streams: [events],
  3336. sessions: [{ id: "ses_child", parentID: "ses_1", time: { updated: 1 } }],
  3337. })
  3338. let childHydrating = false
  3339. let releaseHydration!: () => void
  3340. const hydration = new Promise<void>((resolve) => {
  3341. releaseHydration = resolve
  3342. })
  3343. spyOn(client.message, "list").mockImplementation(async (request) => {
  3344. if (request.sessionID !== "ses_child") return ok({ data: [], cursor: {} })
  3345. childHydrating = true
  3346. await hydration
  3347. return ok({
  3348. data: [
  3349. {
  3350. id: "msg_tool_projected",
  3351. type: "assistant" as const,
  3352. agent: "explore",
  3353. model: { providerID: "test", id: "model" },
  3354. content: [
  3355. {
  3356. type: "tool" as const,
  3357. id: "call_overlap",
  3358. name: "shell",
  3359. state: {
  3360. status: "completed" as const,
  3361. input: { command: "projected" },
  3362. content: [{ type: "text" as const, text: "projected result" }],
  3363. metadata: {},
  3364. },
  3365. time: { created: 1, ran: 1, completed: 2 },
  3366. },
  3367. ],
  3368. time: { created: 1, completed: 2 },
  3369. },
  3370. ],
  3371. cursor: {},
  3372. })
  3373. })
  3374. const ui = footer()
  3375. const transport = await createSessionTransport({
  3376. sdk: client,
  3377. sessionID: "ses_1",
  3378. thinking: false,
  3379. footer: ui.api,
  3380. })
  3381. const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  3382. const inputStarted = (id: string, name: string, seq: number) =>
  3383. events.push({
  3384. id: `evt_started_${id}`,
  3385. created: seq,
  3386. type: "session.tool.input.started",
  3387. durable: durable("ses_child", seq),
  3388. data: { sessionID: "ses_child", assistantMessageID: "msg_tool_projected", id, name },
  3389. })
  3390. const called = (id: string, input: Record<string, unknown>, seq: number) =>
  3391. events.push({
  3392. id: `evt_called_${id}`,
  3393. created: seq,
  3394. type: "session.tool.called",
  3395. durable: durable("ses_child", seq),
  3396. data: {
  3397. sessionID: "ses_child",
  3398. assistantMessageID: "msg_tool_projected",
  3399. id,
  3400. input,
  3401. executed: true,
  3402. },
  3403. })
  3404. inputStarted("call_terminal", "grep", 0)
  3405. called("call_terminal", { pattern: "needle" }, 1)
  3406. await Bun.sleep(0)
  3407. transport.selectSubagent("ses_child")
  3408. while (!childHydrating) await Bun.sleep(0)
  3409. events.push({
  3410. id: "evt_success_terminal",
  3411. created: 2,
  3412. type: "session.tool.success",
  3413. durable: durable("ses_child", 2, 2),
  3414. data: {
  3415. sessionID: "ses_child",
  3416. assistantMessageID: "msg_tool_projected",
  3417. id: "call_terminal",
  3418. metadata: {},
  3419. content: [{ type: "text", text: "found" }],
  3420. executed: true,
  3421. },
  3422. })
  3423. inputStarted("call_overlap", "shell", 3)
  3424. called("call_overlap", { command: "stale" }, 4)
  3425. await Bun.sleep(0)
  3426. const beforeHydration = states().length
  3427. releaseHydration()
  3428. while (states().length === beforeHydration) await Bun.sleep(0)
  3429. await Bun.sleep(0)
  3430. const commits = states().at(-1)?.details.ses_child?.commits ?? []
  3431. expect(commits.find((item) => item.partID === "prt_call_terminal")).toMatchObject({
  3432. tool: "grep",
  3433. toolState: "completed",
  3434. part: { state: { input: { pattern: "needle" } } },
  3435. })
  3436. expect(commits.find((item) => item.partID === "prt_call_overlap")).toMatchObject({
  3437. tool: "shell",
  3438. toolState: "completed",
  3439. part: { state: { input: { command: "projected" } } },
  3440. })
  3441. await transport.close()
  3442. })
  3443. test("keeps child terminal state observed during discovery", async () => {
  3444. const events = feed()
  3445. events.push(connected())
  3446. const client = sdk({ streams: [events] })
  3447. let resolveGet: (() => void) | undefined
  3448. const gate = new Promise<void>((resolve) => {
  3449. resolveGet = resolve
  3450. })
  3451. spyOn(client.session, "get").mockImplementation(async () => {
  3452. await gate
  3453. return ok({
  3454. id: "ses_child",
  3455. parentID: "ses_1",
  3456. projectID: "proj_1",
  3457. agent: "explore",
  3458. cost: 0,
  3459. tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
  3460. time: { created: 1, updated: 1 },
  3461. title: "Find files",
  3462. location: { directory: "/tmp" },
  3463. }) as never
  3464. })
  3465. const ui = footer()
  3466. const transport = await createSessionTransport({
  3467. sdk: client,
  3468. sessionID: "ses_1",
  3469. thinking: false,
  3470. footer: ui.api,
  3471. })
  3472. const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  3473. // Both events arrive while session.get is still in flight.
  3474. events.push({
  3475. id: "evt_child_step",
  3476. created: 0,
  3477. type: "session.step.started",
  3478. durable: durable("ses_child"),
  3479. data: {
  3480. sessionID: "ses_child",
  3481. assistantMessageID: "msg_child_a",
  3482. agent: "explore",
  3483. model: { providerID: "test", id: "model" },
  3484. },
  3485. })
  3486. events.push({
  3487. id: "evt_child_settled",
  3488. created: 0,
  3489. type: "session.execution.interrupted",
  3490. durable: durable("ses_child"),
  3491. data: { sessionID: "ses_child", reason: "user" },
  3492. })
  3493. await Bun.sleep(0)
  3494. resolveGet?.()
  3495. while (!states().some((state) => state.tabs.some((tab) => tab.status === "cancelled"))) await Bun.sleep(0)
  3496. await transport.close()
  3497. })
  3498. test("does not resurrect a settled child from stale discovery buffer", async () => {
  3499. const events = feed()
  3500. events.push(connected())
  3501. const client = sdk({ streams: [events] })
  3502. let resolveGet: (() => void) | undefined
  3503. const gate = new Promise<void>((resolve) => {
  3504. resolveGet = resolve
  3505. })
  3506. spyOn(client.session, "get").mockImplementation(async () => {
  3507. await gate
  3508. return ok({
  3509. id: "ses_child",
  3510. parentID: "ses_1",
  3511. projectID: "proj_1",
  3512. agent: "explore",
  3513. cost: 0,
  3514. tokens: { input: 0, output: 0, reasoning: 0, cache: { read: 0, write: 0 } },
  3515. time: { created: 1, updated: 1 },
  3516. title: "Find files",
  3517. location: { directory: "/tmp" },
  3518. }) as never
  3519. })
  3520. const ui = footer()
  3521. const transport = await createSessionTransport({
  3522. sdk: client,
  3523. sessionID: "ses_1",
  3524. thinking: false,
  3525. footer: ui.api,
  3526. })
  3527. const states = () => ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  3528. // Child event arrives first and gets buffered behind the gated session.get.
  3529. events.push({
  3530. id: "evt_child_step",
  3531. created: 0,
  3532. type: "session.step.started",
  3533. durable: durable("ses_child"),
  3534. data: {
  3535. sessionID: "ses_child",
  3536. assistantMessageID: "msg_child_a",
  3537. agent: "explore",
  3538. model: { providerID: "test", id: "model" },
  3539. },
  3540. })
  3541. // Parent's background subagent tool.success adopts the child mid-discovery.
  3542. events.push({
  3543. id: "evt_parent_input",
  3544. created: 0,
  3545. type: "session.tool.input.started",
  3546. durable: durable("ses_1"),
  3547. data: {
  3548. sessionID: "ses_1",
  3549. assistantMessageID: "msg_parent_a",
  3550. id: "call_sub",
  3551. name: "subagent",
  3552. },
  3553. })
  3554. events.push({
  3555. id: "evt_parent_call",
  3556. created: 0,
  3557. type: "session.tool.called",
  3558. durable: durable("ses_1"),
  3559. data: {
  3560. sessionID: "ses_1",
  3561. assistantMessageID: "msg_parent_a",
  3562. id: "call_sub",
  3563. input: { agent: "explore", description: "Find things", prompt: "go", background: true },
  3564. executed: true,
  3565. },
  3566. })
  3567. events.push({
  3568. id: "evt_parent_success",
  3569. created: 0,
  3570. type: "session.tool.success",
  3571. durable: durable("ses_1", 1, 2),
  3572. data: {
  3573. sessionID: "ses_1",
  3574. assistantMessageID: "msg_parent_a",
  3575. id: "call_sub",
  3576. metadata: { sessionID: "ses_child", status: "running", output: "" },
  3577. content: [{ type: "text", text: "" }],
  3578. executed: true,
  3579. },
  3580. })
  3581. // The settled event arrives after adoption, so it applies directly.
  3582. events.push({
  3583. id: "evt_child_settled",
  3584. created: 0,
  3585. type: "session.execution.interrupted",
  3586. durable: durable("ses_child"),
  3587. data: { sessionID: "ses_child", reason: "shutdown" },
  3588. })
  3589. while (!states().some((state) => state.tabs.some((tab) => tab.status === "cancelled"))) await Bun.sleep(0)
  3590. // Resolving discovery must not replay the buffered step.started over the
  3591. // terminal status.
  3592. const before = states().length
  3593. resolveGet?.()
  3594. while (states().length === before) await Bun.sleep(0)
  3595. await Bun.sleep(0)
  3596. await Bun.sleep(0)
  3597. expect(states().at(-1)?.tabs).toMatchObject([{ sessionID: "ses_child", status: "cancelled" }])
  3598. await transport.close()
  3599. })
  3600. test("adopts historical children from the session family list", async () => {
  3601. const events = feed()
  3602. events.push(connected())
  3603. const client = sdk({
  3604. streams: [events],
  3605. sessions: [
  3606. { id: "ses_child_old", parentID: "ses_1", title: "Earlier subagent", agent: "explore", time: { updated: 9 } },
  3607. { id: "ses_unrelated", title: "Different session", time: { updated: 5 } },
  3608. { id: "ses_sibling", parentID: "ses_2", title: "Someone else's child", time: { updated: 4 } },
  3609. ],
  3610. })
  3611. const ui = footer()
  3612. const transport = await createSessionTransport({
  3613. sdk: client,
  3614. sessionID: "ses_1",
  3615. thinking: false,
  3616. footer: ui.api,
  3617. })
  3618. const states = ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  3619. expect(client.session.list).toHaveBeenCalledWith(
  3620. { parentID: "ses_1", limit: 100, order: "desc" },
  3621. { signal: expect.any(AbortSignal) },
  3622. )
  3623. expect(states.at(-1)?.tabs).toMatchObject([
  3624. {
  3625. sessionID: "ses_child_old",
  3626. label: "Explore",
  3627. title: "Earlier subagent",
  3628. status: "completed",
  3629. },
  3630. ])
  3631. await transport.close()
  3632. })
  3633. test("hydrates completed subagent children from projected tool output", async () => {
  3634. const events = feed()
  3635. events.push(connected())
  3636. const client = sdk({
  3637. streams: [events],
  3638. messages: {
  3639. ses_1: [
  3640. {
  3641. id: "msg_parent",
  3642. type: "assistant" as const,
  3643. agent: "build",
  3644. model: { providerID: "test", id: "model" },
  3645. time: { created: 1, completed: 3 },
  3646. content: [
  3647. {
  3648. type: "tool" as const,
  3649. id: "call_sub",
  3650. name: "subagent",
  3651. state: {
  3652. status: "completed" as const,
  3653. input: { agent: "explore", description: "Find things", prompt: "go" },
  3654. content: [{ type: "text" as const, text: "done" }],
  3655. metadata: { sessionID: "ses_child", status: "completed", output: "done" },
  3656. },
  3657. time: { created: 1, ran: 1, completed: 2 },
  3658. },
  3659. ],
  3660. },
  3661. ],
  3662. },
  3663. })
  3664. const ui = footer()
  3665. const transport = await createSessionTransport({
  3666. sdk: client,
  3667. sessionID: "ses_1",
  3668. thinking: false,
  3669. footer: ui.api,
  3670. })
  3671. const states = ui.events.flatMap((event) => (event.type === "stream.subagent" ? [event.state] : []))
  3672. expect(states.at(-1)?.tabs).toMatchObject([
  3673. {
  3674. sessionID: "ses_child",
  3675. label: "Explore",
  3676. description: "Find things",
  3677. status: "completed",
  3678. },
  3679. ])
  3680. await transport.close()
  3681. })
  3682. })