stream.transport.test.ts 50 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900
  1. import { afterEach, describe, expect, mock, spyOn, test } from "bun:test"
  2. import { OpencodeClient, type GlobalEvent } from "@opencode-ai/sdk/v2"
  3. import { createSessionTransport } from "@/cli/cmd/run/stream.transport"
  4. import type { FooterApi, FooterEvent, RunFilePart, StreamCommit } from "@/cli/cmd/run/types"
  5. type EventStream = Awaited<ReturnType<OpencodeClient["event"]["subscribe"]>>["stream"]
  6. type GlobalEventStream = Awaited<ReturnType<OpencodeClient["global"]["event"]>>["stream"]
  7. type SdkEvent = EventStream extends AsyncGenerator<infer T, unknown, unknown> ? T : never
  8. type SessionMessage = NonNullable<Awaited<ReturnType<OpencodeClient["session"]["messages"]>>["data"]>[number]
  9. type SessionChild = NonNullable<Awaited<ReturnType<OpencodeClient["session"]["children"]>>["data"]>[number]
  10. type SessionToolPart = Extract<SessionMessage["parts"][number], { type: "tool" }>
  11. type SessionStatusMap = NonNullable<Awaited<ReturnType<OpencodeClient["session"]["status"]>>["data"]>
  12. type TextPart = Extract<SessionMessage["parts"][number], { type: "text" }>
  13. afterEach(() => {
  14. mock.restore()
  15. })
  16. function defer<T = void>() {
  17. let resolve!: (value: T | PromiseLike<T>) => void
  18. let reject!: (error?: unknown) => void
  19. const promise = new Promise<T>((next, fail) => {
  20. resolve = next
  21. reject = fail
  22. })
  23. return { promise, resolve, reject }
  24. }
  25. async function waitFor<T>(check: () => T | undefined, timeout = 1_000): Promise<T> {
  26. const end = Date.now() + timeout
  27. while (Date.now() < end) {
  28. const value = check()
  29. if (value !== undefined) {
  30. return value
  31. }
  32. await Bun.sleep(10)
  33. }
  34. throw new Error("timed out waiting for value")
  35. }
  36. function busy(sessionID = "session-1") {
  37. return {
  38. id: `evt-${sessionID}-busy`,
  39. type: "session.status",
  40. properties: {
  41. sessionID,
  42. status: {
  43. type: "busy",
  44. },
  45. },
  46. } satisfies SdkEvent
  47. }
  48. function idle(sessionID = "session-1") {
  49. return {
  50. id: `evt-${sessionID}-idle`,
  51. type: "session.status",
  52. properties: {
  53. sessionID,
  54. status: {
  55. type: "idle",
  56. },
  57. },
  58. } satisfies SdkEvent
  59. }
  60. function retry(sessionID: string, attempt: number, message: string) {
  61. return {
  62. id: `evt-${sessionID}-retry-${attempt}`,
  63. type: "session.status",
  64. properties: {
  65. sessionID,
  66. status: {
  67. type: "retry",
  68. attempt,
  69. message,
  70. next: 1,
  71. },
  72. },
  73. } satisfies SdkEvent
  74. }
  75. function assistant(id: string) {
  76. return {
  77. id: `evt-${id}`,
  78. type: "message.updated",
  79. properties: {
  80. sessionID: "session-1",
  81. info: assistantMessage({
  82. sessionID: "session-1",
  83. id,
  84. parts: [],
  85. }).info,
  86. },
  87. } satisfies SdkEvent
  88. }
  89. const StreamClosed = undefined as never
  90. function feed<T, R = never>(returnValue: R = StreamClosed) {
  91. const list: T[] = []
  92. let done = false
  93. let wake: (() => void) | undefined
  94. const wrapped = (async function* (): AsyncGenerator<T, R, unknown> {
  95. while (!done || list.length > 0) {
  96. if (list.length === 0) {
  97. await new Promise<void>((resolve) => {
  98. wake = resolve
  99. })
  100. continue
  101. }
  102. const next = list.shift()
  103. if (!next) {
  104. continue
  105. }
  106. yield next
  107. }
  108. return returnValue as R
  109. })()
  110. return {
  111. stream: wrapped,
  112. push(value: T) {
  113. list.push(value)
  114. wake?.()
  115. wake = undefined
  116. },
  117. close() {
  118. done = true
  119. wake?.()
  120. wake = undefined
  121. },
  122. }
  123. }
  124. function eventFeed() {
  125. return feed<SdkEvent>()
  126. }
  127. function globalFeed() {
  128. return feed<GlobalEvent>()
  129. }
  130. function emptyStream(): EventStream {
  131. return (async function* (): AsyncGenerator<SdkEvent> {})()
  132. }
  133. function ok<T>(data: T) {
  134. return Promise.resolve({
  135. data,
  136. error: undefined,
  137. request: new Request("https://opencode.test"),
  138. response: new Response(),
  139. })
  140. }
  141. function sse(stream: EventStream) {
  142. return Promise.resolve({ stream })
  143. }
  144. function globalSse(stream: GlobalEventStream) {
  145. return Promise.resolve({ stream })
  146. }
  147. function wrapGlobalStream(stream: EventStream): GlobalEventStream {
  148. return (async function* (): GlobalEventStream {
  149. for await (const event of stream) {
  150. yield globalEvent(event as GlobalEvent["payload"])
  151. }
  152. return StreamClosed
  153. })()
  154. }
  155. function statusMap(busy: boolean): SessionStatusMap {
  156. if (busy) {
  157. return { "session-1": { type: "busy" } }
  158. }
  159. return {}
  160. }
  161. function assistantMessage(input: { sessionID: string; id: string; parts: SessionMessage["parts"] }): SessionMessage {
  162. return {
  163. info: {
  164. id: input.id,
  165. sessionID: input.sessionID,
  166. role: "assistant",
  167. time: {
  168. created: 1,
  169. },
  170. parentID: "msg-user-1",
  171. modelID: "gpt-5",
  172. providerID: "openai",
  173. mode: "chat",
  174. agent: "build",
  175. path: {
  176. cwd: "/tmp",
  177. root: "/tmp",
  178. },
  179. cost: 0,
  180. tokens: {
  181. input: 1,
  182. output: 1,
  183. reasoning: 0,
  184. cache: {
  185. read: 0,
  186. write: 0,
  187. },
  188. },
  189. },
  190. parts: input.parts,
  191. }
  192. }
  193. function runningTool(input: {
  194. sessionID: string
  195. messageID: string
  196. id: string
  197. callID: string
  198. tool: string
  199. body: Record<string, unknown>
  200. metadata?: Record<string, unknown>
  201. }): SessionToolPart {
  202. return {
  203. id: input.id,
  204. sessionID: input.sessionID,
  205. messageID: input.messageID,
  206. type: "tool",
  207. callID: input.callID,
  208. tool: input.tool,
  209. state: {
  210. status: "running",
  211. input: input.body,
  212. ...(input.metadata ? { metadata: input.metadata } : {}),
  213. time: {
  214. start: 1,
  215. },
  216. },
  217. }
  218. }
  219. function completedTool(input: {
  220. sessionID: string
  221. messageID: string
  222. id: string
  223. callID: string
  224. tool: string
  225. body: Record<string, unknown>
  226. output?: string
  227. metadata?: Record<string, unknown>
  228. }): SessionToolPart {
  229. return {
  230. id: input.id,
  231. sessionID: input.sessionID,
  232. messageID: input.messageID,
  233. type: "tool",
  234. callID: input.callID,
  235. tool: input.tool,
  236. state: {
  237. status: "completed",
  238. input: input.body,
  239. output: input.output ?? "",
  240. title: input.tool,
  241. metadata: input.metadata ?? {},
  242. time: {
  243. start: 1,
  244. end: 2,
  245. },
  246. },
  247. }
  248. }
  249. function textPart(id: string, messageID: string, text: string, sessionID = "session-1"): TextPart {
  250. return {
  251. id,
  252. sessionID,
  253. messageID,
  254. type: "text",
  255. text,
  256. }
  257. }
  258. function textUpdated(part: TextPart): SdkEvent {
  259. return {
  260. id: `evt-${part.id}-updated`,
  261. type: "message.part.updated",
  262. properties: {
  263. sessionID: part.sessionID,
  264. part,
  265. time: 1,
  266. },
  267. }
  268. }
  269. function toolUpdated(part: SessionToolPart): SdkEvent {
  270. return {
  271. id: `evt-${part.id}-updated`,
  272. type: "message.part.updated",
  273. properties: {
  274. sessionID: part.sessionID,
  275. part,
  276. time: 1,
  277. },
  278. }
  279. }
  280. function textDelta(messageID: string, partID: string, delta: string, sessionID = "session-1"): SdkEvent {
  281. return {
  282. id: `evt-${partID}-delta`,
  283. type: "message.part.delta",
  284. properties: {
  285. sessionID,
  286. messageID,
  287. partID,
  288. field: "text",
  289. delta,
  290. },
  291. }
  292. }
  293. function child(id: string): SessionChild {
  294. return {
  295. id,
  296. slug: id,
  297. projectID: "project-1",
  298. directory: "/tmp",
  299. title: id,
  300. version: "1",
  301. time: {
  302. created: 1,
  303. updated: 1,
  304. },
  305. }
  306. }
  307. function globalEvent(payload: SdkEvent | GlobalEvent["payload"]): GlobalEvent {
  308. return {
  309. directory: "/tmp",
  310. project: "project-1",
  311. payload: payload as GlobalEvent["payload"],
  312. }
  313. }
  314. function footer(fn?: (commit: StreamCommit) => void) {
  315. const commits: StreamCommit[] = []
  316. const events: FooterEvent[] = []
  317. let closed = false
  318. let idleCalls = 0
  319. const api: FooterApi = {
  320. get isClosed() {
  321. return closed
  322. },
  323. onPrompt: () => () => {},
  324. onClose: () => () => {},
  325. event(next) {
  326. events.push(next)
  327. },
  328. append(next) {
  329. commits.push(next)
  330. fn?.(next)
  331. },
  332. idle() {
  333. idleCalls += 1
  334. return Promise.resolve()
  335. },
  336. close() {
  337. closed = true
  338. },
  339. destroy() {
  340. closed = true
  341. },
  342. }
  343. return {
  344. api,
  345. commits,
  346. events,
  347. get idleCalls() {
  348. return idleCalls
  349. },
  350. }
  351. }
  352. function sdk(
  353. input: {
  354. stream?: EventStream
  355. globalStream?: GlobalEventStream
  356. subscribe?: OpencodeClient["event"]["subscribe"]
  357. globalEvent?: OpencodeClient["global"]["event"]
  358. promptAsync?: OpencodeClient["session"]["promptAsync"]
  359. status?: OpencodeClient["session"]["status"]
  360. messages?: OpencodeClient["session"]["messages"]
  361. children?: OpencodeClient["session"]["children"]
  362. permissions?: OpencodeClient["permission"]["list"]
  363. questions?: OpencodeClient["question"]["list"]
  364. } = {},
  365. ) {
  366. const client = new OpencodeClient()
  367. const subscribe: OpencodeClient["event"]["subscribe"] = input.subscribe ?? (() => sse(input.stream ?? emptyStream()))
  368. const globalEvent: OpencodeClient["global"]["event"] =
  369. input.globalEvent ?? (() => globalSse(input.globalStream ?? wrapGlobalStream(input.stream ?? emptyStream())))
  370. const promptAsync: OpencodeClient["session"]["promptAsync"] = input.promptAsync ?? (() => ok(undefined))
  371. const status: OpencodeClient["session"]["status"] = input.status ?? (() => ok({}))
  372. const messages: OpencodeClient["session"]["messages"] = input.messages ?? (() => ok([]))
  373. const children: OpencodeClient["session"]["children"] = input.children ?? (() => ok([]))
  374. const permissions: OpencodeClient["permission"]["list"] = input.permissions ?? (() => ok([]))
  375. const questions: OpencodeClient["question"]["list"] = input.questions ?? (() => ok([]))
  376. spyOn(client.event, "subscribe").mockImplementation(subscribe)
  377. spyOn(client.global, "event").mockImplementation(globalEvent)
  378. spyOn(client.session, "promptAsync").mockImplementation(promptAsync)
  379. spyOn(client.session, "status").mockImplementation(status)
  380. spyOn(client.session, "messages").mockImplementation(messages)
  381. spyOn(client.session, "children").mockImplementation(children)
  382. spyOn(client.permission, "list").mockImplementation(permissions)
  383. spyOn(client.question, "list").mockImplementation(questions)
  384. return client
  385. }
  386. describe("run stream transport", () => {
  387. test("does not replay persisted main-session history during bootstrap by default", async () => {
  388. const src = eventFeed()
  389. const ui = footer()
  390. const transport = await createSessionTransport({
  391. sdk: sdk({
  392. stream: src.stream,
  393. messages: async ({ sessionID }) =>
  394. sessionID === "session-1"
  395. ? ok([
  396. assistantMessage({
  397. sessionID: "session-1",
  398. id: "msg-1",
  399. parts: [
  400. {
  401. ...textPart("text-1", "msg-1", "Hello."),
  402. time: {
  403. start: 1,
  404. end: 2,
  405. },
  406. },
  407. ],
  408. }),
  409. ])
  410. : ok([]),
  411. }),
  412. sessionID: "session-1",
  413. thinking: true,
  414. limits: () => ({}),
  415. footer: ui.api,
  416. })
  417. try {
  418. expect(ui.commits).toEqual([])
  419. expect(ui.idleCalls).toBe(0)
  420. } finally {
  421. src.close()
  422. await transport.close()
  423. }
  424. })
  425. test("replays persisted main-session history during bootstrap when enabled", async () => {
  426. const src = eventFeed()
  427. const ui = footer()
  428. const transport = await createSessionTransport({
  429. sdk: sdk({
  430. stream: src.stream,
  431. messages: async ({ sessionID }) =>
  432. sessionID === "session-1"
  433. ? ok([
  434. assistantMessage({
  435. sessionID: "session-1",
  436. id: "msg-1",
  437. parts: [
  438. {
  439. ...textPart("text-1", "msg-1", "Hello."),
  440. time: {
  441. start: 1,
  442. end: 2,
  443. },
  444. },
  445. ],
  446. }),
  447. ])
  448. : ok([]),
  449. }),
  450. sessionID: "session-1",
  451. thinking: true,
  452. replay: true,
  453. limits: () => ({}),
  454. footer: ui.api,
  455. })
  456. try {
  457. await waitFor(() => ui.commits.find((item) => item.kind === "assistant" && item.text === "Hello."))
  458. expect(ui.idleCalls).toBeGreaterThan(0)
  459. } finally {
  460. src.close()
  461. await transport.close()
  462. }
  463. })
  464. test("caps replayed bootstrap history to the configured number of messages", async () => {
  465. const src = eventFeed()
  466. const ui = footer()
  467. const transport = await createSessionTransport({
  468. sdk: sdk({
  469. stream: src.stream,
  470. messages: async ({ sessionID }) =>
  471. ok(
  472. sessionID === "session-1"
  473. ? [
  474. assistantMessage({
  475. sessionID: "session-1",
  476. id: "msg-1",
  477. parts: [
  478. {
  479. ...textPart("text-1", "msg-1", "Hello."),
  480. time: {
  481. start: 1,
  482. end: 2,
  483. },
  484. },
  485. ],
  486. }),
  487. assistantMessage({
  488. sessionID: "session-1",
  489. id: "msg-2",
  490. parts: [
  491. {
  492. ...textPart("text-2", "msg-2", "World."),
  493. time: {
  494. start: 3,
  495. end: 4,
  496. },
  497. },
  498. ],
  499. }),
  500. ]
  501. : [],
  502. ),
  503. }),
  504. sessionID: "session-1",
  505. thinking: true,
  506. replay: true,
  507. replayLimit: 1,
  508. limits: () => ({}),
  509. footer: ui.api,
  510. })
  511. try {
  512. await waitFor(() => (ui.commits.length > 0 ? ui.commits : undefined))
  513. expect(ui.commits.filter((item) => item.kind === "assistant")).toEqual([
  514. expect.objectContaining({
  515. text: "World.",
  516. }),
  517. ])
  518. } finally {
  519. src.close()
  520. await transport.close()
  521. }
  522. })
  523. test("skips buffered pre-bootstrap deltas already covered by replay history", async () => {
  524. const src = eventFeed()
  525. const ui = footer()
  526. const gate = defer<void>()
  527. let transport: Awaited<ReturnType<typeof createSessionTransport>> | undefined
  528. const task = createSessionTransport({
  529. sdk: sdk({
  530. stream: src.stream,
  531. messages: async ({ sessionID }) => {
  532. if (sessionID !== "session-1") {
  533. return ok([])
  534. }
  535. await gate.promise
  536. return ok([
  537. assistantMessage({
  538. sessionID: "session-1",
  539. id: "msg-1",
  540. parts: [textPart("text-1", "msg-1", "Hello")],
  541. }),
  542. ])
  543. },
  544. }),
  545. sessionID: "session-1",
  546. thinking: true,
  547. replay: true,
  548. limits: () => ({}),
  549. footer: ui.api,
  550. })
  551. try {
  552. await Promise.resolve()
  553. src.push(textDelta("msg-1", "text-1", "lo"))
  554. gate.resolve()
  555. transport = await task
  556. await waitFor(() => (ui.commits.length > 0 ? ui.commits : undefined))
  557. await Bun.sleep(20)
  558. expect(ui.commits.filter((item) => item.kind === "assistant")).toEqual([
  559. expect.objectContaining({
  560. text: "Hello",
  561. }),
  562. ])
  563. } finally {
  564. src.close()
  565. await transport?.close()
  566. }
  567. })
  568. test("applies buffered pre-bootstrap deltas not yet persisted", async () => {
  569. const src = eventFeed()
  570. const ui = footer()
  571. const gate = defer<void>()
  572. let transport: Awaited<ReturnType<typeof createSessionTransport>> | undefined
  573. const task = createSessionTransport({
  574. sdk: sdk({
  575. stream: src.stream,
  576. messages: async ({ sessionID }) => {
  577. if (sessionID !== "session-1") {
  578. return ok([])
  579. }
  580. await gate.promise
  581. return ok([
  582. assistantMessage({
  583. sessionID: "session-1",
  584. id: "msg-1",
  585. parts: [textPart("text-1", "msg-1", "")],
  586. }),
  587. ])
  588. },
  589. }),
  590. sessionID: "session-1",
  591. thinking: true,
  592. replay: true,
  593. limits: () => ({}),
  594. footer: ui.api,
  595. })
  596. try {
  597. await Promise.resolve()
  598. src.push(textDelta("msg-1", "text-1", "Hello"))
  599. gate.resolve()
  600. transport = await task
  601. await waitFor(() => (ui.commits.length > 0 ? ui.commits : undefined))
  602. await Bun.sleep(20)
  603. expect(ui.commits.filter((item) => item.kind === "assistant")).toEqual([
  604. expect.objectContaining({
  605. text: "Hello",
  606. }),
  607. ])
  608. } finally {
  609. src.close()
  610. await transport?.close()
  611. }
  612. })
  613. test("preserves running footer state for resumed active sessions", async () => {
  614. const src = eventFeed()
  615. const ui = footer()
  616. const transport = await createSessionTransport({
  617. sdk: sdk({
  618. stream: src.stream,
  619. messages: async ({ sessionID }) =>
  620. sessionID === "session-1"
  621. ? ok([
  622. assistantMessage({
  623. sessionID: "session-1",
  624. id: "msg-1",
  625. parts: [
  626. runningTool({
  627. sessionID: "session-1",
  628. messageID: "msg-1",
  629. id: "bash-1",
  630. callID: "call-1",
  631. tool: "bash",
  632. body: {
  633. command: "pwd",
  634. },
  635. }),
  636. ],
  637. }),
  638. ])
  639. : ok([]),
  640. }),
  641. sessionID: "session-1",
  642. thinking: true,
  643. replay: true,
  644. limits: () => ({}),
  645. footer: ui.api,
  646. })
  647. try {
  648. const patch = await waitFor(() => {
  649. const item = ui.events.findLast((event) => event.type === "stream.patch")
  650. return item?.type === "stream.patch" ? item.patch : undefined
  651. })
  652. expect(patch).toEqual(
  653. expect.objectContaining({
  654. phase: "running",
  655. status: "running bash",
  656. }),
  657. )
  658. } finally {
  659. src.close()
  660. await transport.close()
  661. }
  662. })
  663. test("drops completed historical subagent tabs during bootstrap", async () => {
  664. const src = eventFeed()
  665. const ui = footer()
  666. const transport = await createSessionTransport({
  667. sdk: sdk({
  668. stream: src.stream,
  669. messages: async ({ sessionID }) => {
  670. if (sessionID !== "session-1") {
  671. return ok([])
  672. }
  673. return ok([
  674. assistantMessage({
  675. sessionID: "session-1",
  676. id: "msg-1",
  677. parts: [
  678. completedTool({
  679. sessionID: "session-1",
  680. messageID: "msg-1",
  681. id: "task-1",
  682. callID: "call-1",
  683. tool: "task",
  684. body: {
  685. description: "Explore run folder",
  686. subagent_type: "explore",
  687. },
  688. metadata: {
  689. sessionId: "child-1",
  690. },
  691. }),
  692. ],
  693. }),
  694. ])
  695. },
  696. children: async () => ok([child("child-1")]),
  697. }),
  698. sessionID: "session-1",
  699. thinking: true,
  700. limits: () => ({}),
  701. footer: ui.api,
  702. })
  703. try {
  704. const state = await waitFor(() => {
  705. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  706. return item?.type === "stream.subagent" ? item.state : undefined
  707. })
  708. expect(state.tabs).toEqual([])
  709. expect(state.details).toEqual({})
  710. } finally {
  711. src.close()
  712. await transport.close()
  713. }
  714. })
  715. test("bootstraps child tabs and resumed blocker input", async () => {
  716. const src = eventFeed()
  717. const ui = footer()
  718. const transport = await createSessionTransport({
  719. sdk: sdk({
  720. stream: src.stream,
  721. messages: async ({ sessionID }) => {
  722. if (sessionID === "session-1") {
  723. return ok([
  724. assistantMessage({
  725. sessionID: "session-1",
  726. id: "msg-1",
  727. parts: [
  728. runningTool({
  729. sessionID: "session-1",
  730. messageID: "msg-1",
  731. id: "task-1",
  732. callID: "call-1",
  733. tool: "task",
  734. body: {
  735. description: "Explore run folder",
  736. subagent_type: "explore",
  737. },
  738. metadata: {
  739. sessionId: "child-1",
  740. },
  741. }),
  742. ],
  743. }),
  744. ])
  745. }
  746. return ok([
  747. assistantMessage({
  748. sessionID: "child-1",
  749. id: "msg-child-1",
  750. parts: [
  751. runningTool({
  752. sessionID: "child-1",
  753. messageID: "msg-child-1",
  754. id: "edit-1",
  755. callID: "call-edit-1",
  756. tool: "edit",
  757. body: {
  758. filePath: "src/run/subagent-data.ts",
  759. diff: "@@ -1 +1 @@",
  760. },
  761. }),
  762. ],
  763. }),
  764. ])
  765. },
  766. children: async () => ok([child("child-1")]),
  767. permissions: async () =>
  768. ok([
  769. {
  770. id: "perm-1",
  771. sessionID: "child-1",
  772. permission: "edit",
  773. patterns: ["src/run/subagent-data.ts"],
  774. metadata: {},
  775. always: [],
  776. tool: {
  777. messageID: "msg-child-1",
  778. callID: "call-edit-1",
  779. },
  780. },
  781. ]),
  782. }),
  783. sessionID: "session-1",
  784. thinking: true,
  785. limits: () => ({}),
  786. footer: ui.api,
  787. })
  788. try {
  789. const boot = await waitFor(() => {
  790. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  791. const state = item?.type === "stream.subagent" ? item.state : undefined
  792. return state?.tabs.some((tab) => tab.sessionID === "child-1") &&
  793. state.permissions.some((req) => req.id === "perm-1")
  794. ? state
  795. : undefined
  796. })
  797. expect(boot.tabs).toEqual([
  798. expect.objectContaining({
  799. sessionID: "child-1",
  800. label: "Explore",
  801. description: "Pending permission",
  802. status: "running",
  803. }),
  804. ])
  805. expect(boot.permissions).toEqual([
  806. expect.objectContaining({
  807. id: "perm-1",
  808. sessionID: "child-1",
  809. metadata: {
  810. input: {
  811. filePath: "src/run/subagent-data.ts",
  812. diff: "@@ -1 +1 @@",
  813. },
  814. },
  815. }),
  816. ])
  817. transport.selectSubagent("child-1")
  818. const selected = await waitFor(() => {
  819. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  820. const state = item?.type === "stream.subagent" ? item.state : undefined
  821. const detail = state?.details["child-1"]
  822. return detail?.commits.some(
  823. (commit) => commit.kind === "tool" && commit.tool === "edit" && commit.phase === "start",
  824. )
  825. ? state
  826. : undefined
  827. })
  828. expect(selected.details).toEqual({
  829. "child-1": {
  830. sessionID: "child-1",
  831. commits: [
  832. expect.objectContaining({
  833. kind: "tool",
  834. tool: "edit",
  835. phase: "start",
  836. }),
  837. ],
  838. },
  839. })
  840. expect(
  841. await waitFor(() => {
  842. const item = ui.events.findLast((event) => event.type === "stream.view")
  843. return item?.type === "stream.view" && item.view.type === "permission" && item.view.request.id === "perm-1"
  844. ? item
  845. : undefined
  846. }),
  847. ).toEqual({
  848. type: "stream.view",
  849. view: {
  850. type: "permission",
  851. request: expect.objectContaining({
  852. id: "perm-1",
  853. metadata: {
  854. input: {
  855. filePath: "src/run/subagent-data.ts",
  856. diff: "@@ -1 +1 @@",
  857. },
  858. },
  859. }),
  860. },
  861. })
  862. } finally {
  863. src.close()
  864. await transport.close()
  865. }
  866. })
  867. test("bootstraps child session output before selection", async () => {
  868. const ui = footer()
  869. const transport = await createSessionTransport({
  870. sdk: sdk({
  871. messages: async ({ sessionID }) => {
  872. if (sessionID === "session-1") {
  873. return ok([
  874. assistantMessage({
  875. sessionID: "session-1",
  876. id: "msg-1",
  877. parts: [
  878. runningTool({
  879. sessionID: "session-1",
  880. messageID: "msg-1",
  881. id: "task-1",
  882. callID: "call-1",
  883. tool: "task",
  884. body: {
  885. description: "Explore run.ts",
  886. subagent_type: "explore",
  887. },
  888. metadata: {
  889. sessionId: "child-1",
  890. },
  891. }),
  892. ],
  893. }),
  894. ])
  895. }
  896. return sessionID === "child-1"
  897. ? ok([
  898. assistantMessage({
  899. sessionID: "child-1",
  900. id: "msg-child-1",
  901. parts: [textPart("txt-child-1", "msg-child-1", "subagent summary", "child-1")],
  902. }),
  903. ])
  904. : ok([])
  905. },
  906. children: async () => ok([child("child-1")]),
  907. }),
  908. sessionID: "session-1",
  909. thinking: true,
  910. limits: () => ({}),
  911. footer: ui.api,
  912. })
  913. try {
  914. await waitFor(() => {
  915. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  916. return item?.type === "stream.subagent" && item.state.tabs.some((tab) => tab.sessionID === "child-1")
  917. ? item
  918. : undefined
  919. })
  920. transport.selectSubagent("child-1")
  921. expect(
  922. await waitFor(() => {
  923. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  924. const detail = item?.type === "stream.subagent" ? item.state.details["child-1"] : undefined
  925. return detail?.commits.some((commit) => commit.kind === "assistant" && commit.text === "subagent summary")
  926. ? detail
  927. : undefined
  928. }),
  929. ).toEqual({
  930. sessionID: "child-1",
  931. commits: [
  932. expect.objectContaining({
  933. kind: "assistant",
  934. text: "subagent summary",
  935. }),
  936. ],
  937. })
  938. } finally {
  939. await transport.close()
  940. }
  941. })
  942. test("does not block startup on child history bootstrap", async () => {
  943. const pending = defer<Awaited<ReturnType<typeof ok<SessionMessage[]>>>>()
  944. const ui = footer()
  945. let transport: Awaited<ReturnType<typeof createSessionTransport>> | undefined
  946. const task = createSessionTransport({
  947. sdk: sdk({
  948. messages: async ({ sessionID }) => {
  949. if (sessionID === "session-1") {
  950. return ok([
  951. assistantMessage({
  952. sessionID: "session-1",
  953. id: "msg-1",
  954. parts: [
  955. runningTool({
  956. sessionID: "session-1",
  957. messageID: "msg-1",
  958. id: "task-1",
  959. callID: "call-1",
  960. tool: "task",
  961. body: {
  962. description: "Explore run.ts",
  963. subagent_type: "explore",
  964. },
  965. metadata: {
  966. sessionId: "child-1",
  967. },
  968. }),
  969. ],
  970. }),
  971. ])
  972. }
  973. if (sessionID === "child-1") {
  974. return pending.promise
  975. }
  976. return ok([])
  977. },
  978. children: async () => ok([child("child-1")]),
  979. }),
  980. sessionID: "session-1",
  981. thinking: true,
  982. limits: () => ({}),
  983. footer: ui.api,
  984. }).then((item) => {
  985. transport = item
  986. return item
  987. })
  988. try {
  989. const state = await waitFor(() => {
  990. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  991. return item?.type === "stream.subagent" && item.state.tabs.some((tab) => tab.sessionID === "child-1")
  992. ? item.state
  993. : undefined
  994. })
  995. await waitFor(() => transport)
  996. expect(state).toEqual({
  997. tabs: [expect.objectContaining({ sessionID: "child-1", status: "running" })],
  998. details: {},
  999. permissions: [],
  1000. questions: [],
  1001. })
  1002. } finally {
  1003. pending.resolve(ok([]))
  1004. await task
  1005. await transport?.close()
  1006. }
  1007. })
  1008. test("replays child events buffered during bootstrap once the tab is known", async () => {
  1009. const global = globalFeed()
  1010. const ui = footer()
  1011. const gate = defer<void>()
  1012. let transport: Awaited<ReturnType<typeof createSessionTransport>> | undefined
  1013. const task = createSessionTransport({
  1014. sdk: sdk({
  1015. globalStream: global.stream,
  1016. messages: async ({ sessionID }) => {
  1017. if (sessionID !== "session-1") {
  1018. return ok([])
  1019. }
  1020. await gate.promise
  1021. return ok([])
  1022. },
  1023. children: async () => ok([]),
  1024. }),
  1025. sessionID: "session-1",
  1026. thinking: true,
  1027. limits: () => ({}),
  1028. footer: ui.api,
  1029. })
  1030. try {
  1031. await Promise.resolve()
  1032. global.push(globalEvent(retry("child-1", 1, "retry child")))
  1033. global.push(
  1034. globalEvent({
  1035. id: "evt-child-message",
  1036. type: "message.updated",
  1037. properties: {
  1038. sessionID: "child-1",
  1039. info: assistantMessage({
  1040. sessionID: "child-1",
  1041. id: "msg-child-1",
  1042. parts: [],
  1043. }).info,
  1044. },
  1045. }),
  1046. )
  1047. global.push(globalEvent(textUpdated(textPart("txt-child-1", "msg-child-1", "", "child-1"))))
  1048. global.push(globalEvent(textDelta("msg-child-1", "txt-child-1", "Hello", "child-1")))
  1049. global.push(
  1050. globalEvent(
  1051. toolUpdated(
  1052. runningTool({
  1053. sessionID: "session-1",
  1054. messageID: "msg-1",
  1055. id: "task-1",
  1056. callID: "call-1",
  1057. tool: "task",
  1058. body: {
  1059. description: "Explore run.ts",
  1060. subagent_type: "explore",
  1061. },
  1062. metadata: {
  1063. sessionId: "child-1",
  1064. },
  1065. }),
  1066. ),
  1067. ),
  1068. )
  1069. gate.resolve()
  1070. transport = await task
  1071. await waitFor(() => {
  1072. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  1073. return item?.type === "stream.subagent" && item.state.tabs.some((tab) => tab.sessionID === "child-1")
  1074. ? item
  1075. : undefined
  1076. })
  1077. transport.selectSubagent("child-1")
  1078. const detail = await waitFor(() => {
  1079. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  1080. const next = item?.type === "stream.subagent" ? item.state.details["child-1"] : undefined
  1081. return next?.commits.some((commit) => commit.kind === "error" && commit.text === "retry child") &&
  1082. next.commits.some((commit) => commit.kind === "assistant" && commit.text === "Hello")
  1083. ? next
  1084. : undefined
  1085. })
  1086. expect(detail).toEqual({
  1087. sessionID: "child-1",
  1088. commits: expect.arrayContaining([
  1089. expect.objectContaining({
  1090. kind: "error",
  1091. text: "retry child",
  1092. }),
  1093. expect.objectContaining({
  1094. kind: "assistant",
  1095. text: "Hello",
  1096. }),
  1097. ]),
  1098. })
  1099. } finally {
  1100. global.close()
  1101. await transport?.close()
  1102. }
  1103. })
  1104. test("streams selected subagent output from global events while it is running", async () => {
  1105. const global = globalFeed()
  1106. const ui = footer()
  1107. const transport = await createSessionTransport({
  1108. sdk: sdk({
  1109. globalStream: global.stream,
  1110. }),
  1111. sessionID: "session-1",
  1112. thinking: true,
  1113. limits: () => ({}),
  1114. footer: ui.api,
  1115. })
  1116. try {
  1117. global.push(globalEvent(assistant("msg-1")))
  1118. global.push(
  1119. globalEvent(
  1120. toolUpdated(
  1121. runningTool({
  1122. sessionID: "session-1",
  1123. messageID: "msg-1",
  1124. id: "task-1",
  1125. callID: "call-1",
  1126. tool: "task",
  1127. body: {
  1128. description: "Explore run.ts",
  1129. subagent_type: "explore",
  1130. },
  1131. metadata: {
  1132. sessionId: "child-1",
  1133. },
  1134. }),
  1135. ),
  1136. ),
  1137. )
  1138. await waitFor(() => {
  1139. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  1140. return item?.type === "stream.subagent" && item.state.tabs.some((tab) => tab.sessionID === "child-1")
  1141. ? item
  1142. : undefined
  1143. })
  1144. transport.selectSubagent("child-1")
  1145. global.push(
  1146. globalEvent({
  1147. id: "evt-child-message",
  1148. type: "message.updated",
  1149. properties: {
  1150. sessionID: "child-1",
  1151. info: assistantMessage({
  1152. sessionID: "child-1",
  1153. id: "msg-child-1",
  1154. parts: [],
  1155. }).info,
  1156. },
  1157. }),
  1158. )
  1159. global.push(globalEvent(textUpdated(textPart("txt-child-1", "msg-child-1", "hello", "child-1"))))
  1160. expect(
  1161. await waitFor(() => {
  1162. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  1163. const detail = item?.type === "stream.subagent" ? item.state.details["child-1"] : undefined
  1164. return detail?.commits.some((commit) => commit.kind === "assistant" && commit.text === "hello")
  1165. ? detail
  1166. : undefined
  1167. }),
  1168. ).toEqual({
  1169. sessionID: "child-1",
  1170. commits: [
  1171. expect.objectContaining({
  1172. kind: "assistant",
  1173. text: "hello",
  1174. }),
  1175. ],
  1176. })
  1177. global.push(globalEvent(textUpdated(textPart("txt-child-1", "msg-child-1", "hello world", "child-1"))))
  1178. expect(
  1179. await waitFor(() => {
  1180. const item = ui.events.findLast((event) => event.type === "stream.subagent")
  1181. const detail = item?.type === "stream.subagent" ? item.state.details["child-1"] : undefined
  1182. return detail?.commits.some((commit) => commit.kind === "assistant" && commit.text === "hello world")
  1183. ? detail
  1184. : undefined
  1185. }, 2_000),
  1186. ).toEqual({
  1187. sessionID: "child-1",
  1188. commits: [
  1189. expect.objectContaining({
  1190. kind: "assistant",
  1191. text: "hello world",
  1192. }),
  1193. ],
  1194. })
  1195. } finally {
  1196. global.close()
  1197. await transport.close()
  1198. }
  1199. })
  1200. test("recovers pending questions from question.list when question.asked is missed", async () => {
  1201. const src = eventFeed()
  1202. const ui = footer()
  1203. let questionCalls = 0
  1204. const request = {
  1205. id: "question-1",
  1206. sessionID: "session-1",
  1207. questions: [
  1208. {
  1209. question: "Which area should I inspect first?",
  1210. header: "Area",
  1211. options: [{ label: "CLI", description: "Look at the direct run flow." }],
  1212. multiple: false,
  1213. },
  1214. ],
  1215. tool: {
  1216. messageID: "msg-1",
  1217. callID: "call-question-1",
  1218. },
  1219. }
  1220. const transport = await createSessionTransport({
  1221. sdk: sdk({
  1222. stream: src.stream,
  1223. questions: async () => {
  1224. questionCalls += 1
  1225. return ok(questionCalls > 1 ? [request] : [])
  1226. },
  1227. promptAsync: async () => {
  1228. queueMicrotask(() => {
  1229. src.push(busy())
  1230. src.push(assistant("msg-1"))
  1231. src.push(
  1232. toolUpdated(
  1233. runningTool({
  1234. sessionID: "session-1",
  1235. messageID: "msg-1",
  1236. id: "question-tool-1",
  1237. callID: "call-question-1",
  1238. tool: "question",
  1239. body: {
  1240. questions: request.questions,
  1241. },
  1242. }),
  1243. ),
  1244. )
  1245. })
  1246. return ok(undefined)
  1247. },
  1248. }),
  1249. sessionID: "session-1",
  1250. thinking: true,
  1251. limits: () => ({}),
  1252. footer: ui.api,
  1253. })
  1254. const ctrl = new AbortController()
  1255. try {
  1256. const run = transport.runPromptTurn({
  1257. agent: undefined,
  1258. model: undefined,
  1259. variant: undefined,
  1260. prompt: { text: "hello", parts: [] },
  1261. files: [],
  1262. includeFiles: false,
  1263. signal: ctrl.signal,
  1264. })
  1265. const view = await waitFor(() => {
  1266. const item = ui.events.findLast((event) => event.type === "stream.view")
  1267. return item?.type === "stream.view" && item.view.type === "question" ? item.view : undefined
  1268. })
  1269. expect(view).toEqual({
  1270. type: "question",
  1271. request,
  1272. })
  1273. expect(ui.events).toContainEqual({
  1274. type: "stream.patch",
  1275. patch: {
  1276. phase: "running",
  1277. status: "awaiting answer",
  1278. },
  1279. })
  1280. src.push(
  1281. toolUpdated(
  1282. completedTool({
  1283. sessionID: "session-1",
  1284. messageID: "msg-1",
  1285. id: "question-tool-1",
  1286. callID: "call-question-1",
  1287. tool: "question",
  1288. body: {
  1289. questions: request.questions,
  1290. },
  1291. output: "User has answered your questions.",
  1292. metadata: {
  1293. answers: [["CLI"]],
  1294. },
  1295. }),
  1296. ),
  1297. )
  1298. expect(
  1299. await waitFor(() => {
  1300. const item = ui.events.findLast((event) => event.type === "stream.view")
  1301. return item?.type === "stream.view" && item.view.type === "prompt" ? item : undefined
  1302. }),
  1303. ).toEqual({
  1304. type: "stream.view",
  1305. view: { type: "prompt" },
  1306. })
  1307. ctrl.abort()
  1308. await run
  1309. } finally {
  1310. src.close()
  1311. await transport.close()
  1312. }
  1313. })
  1314. test("does not resurrect questions if question.list resolves after tool completion", async () => {
  1315. const src = eventFeed()
  1316. const ui = footer()
  1317. const started = defer()
  1318. const request = {
  1319. id: "question-race-1",
  1320. sessionID: "session-1",
  1321. questions: [
  1322. {
  1323. question: "Which area should I inspect first?",
  1324. header: "Area",
  1325. options: [{ label: "CLI", description: "Look at the direct run flow." }],
  1326. multiple: false,
  1327. },
  1328. ],
  1329. tool: {
  1330. messageID: "msg-1",
  1331. callID: "call-question-race-1",
  1332. },
  1333. }
  1334. const pending = defer<Awaited<ReturnType<typeof ok<(typeof request)[]>>>>()
  1335. let questionCalls = 0
  1336. const transport = await createSessionTransport({
  1337. sdk: sdk({
  1338. stream: src.stream,
  1339. questions: async () => {
  1340. questionCalls += 1
  1341. if (questionCalls === 1) {
  1342. return ok([])
  1343. }
  1344. if (questionCalls === 2) {
  1345. started.resolve()
  1346. return pending.promise
  1347. }
  1348. return ok([])
  1349. },
  1350. promptAsync: async () => {
  1351. queueMicrotask(() => {
  1352. src.push(busy())
  1353. src.push(assistant("msg-1"))
  1354. src.push(
  1355. toolUpdated(
  1356. runningTool({
  1357. sessionID: "session-1",
  1358. messageID: "msg-1",
  1359. id: "question-race-tool-1",
  1360. callID: "call-question-race-1",
  1361. tool: "question",
  1362. body: {
  1363. questions: request.questions,
  1364. },
  1365. }),
  1366. ),
  1367. )
  1368. })
  1369. return ok(undefined)
  1370. },
  1371. }),
  1372. sessionID: "session-1",
  1373. thinking: true,
  1374. limits: () => ({}),
  1375. footer: ui.api,
  1376. })
  1377. const ctrl = new AbortController()
  1378. try {
  1379. const run = transport.runPromptTurn({
  1380. agent: undefined,
  1381. model: undefined,
  1382. variant: undefined,
  1383. prompt: { text: "hello", parts: [] },
  1384. files: [],
  1385. includeFiles: false,
  1386. signal: ctrl.signal,
  1387. })
  1388. await started.promise
  1389. src.push(
  1390. toolUpdated(
  1391. completedTool({
  1392. sessionID: "session-1",
  1393. messageID: "msg-1",
  1394. id: "question-race-tool-1",
  1395. callID: "call-question-race-1",
  1396. tool: "question",
  1397. body: {
  1398. questions: request.questions,
  1399. },
  1400. output: "User has answered your questions.",
  1401. metadata: {
  1402. answers: [["CLI"]],
  1403. },
  1404. }),
  1405. ),
  1406. )
  1407. await waitFor(() => {
  1408. const commit = ui.commits.findLast(
  1409. (item) => item.kind === "tool" && item.partID === "question-race-tool-1" && item.toolState === "completed",
  1410. )
  1411. return commit ? true : undefined
  1412. })
  1413. pending.resolve(ok([request]))
  1414. await Bun.sleep(50)
  1415. expect(
  1416. ui.events.some(
  1417. (event) =>
  1418. event.type === "stream.view" && event.view.type === "question" && event.view.request.id === request.id,
  1419. ),
  1420. ).toBe(false)
  1421. ctrl.abort()
  1422. await run
  1423. } finally {
  1424. src.close()
  1425. await transport.close()
  1426. }
  1427. })
  1428. test("respects the includeFiles flag when building prompt payloads", async () => {
  1429. const src = eventFeed()
  1430. const ui = footer()
  1431. const seen: unknown[] = []
  1432. const file: RunFilePart = {
  1433. type: "file",
  1434. url: "file:///tmp/a.ts",
  1435. filename: "a.ts",
  1436. mime: "text/plain",
  1437. }
  1438. const transport = await createSessionTransport({
  1439. sdk: sdk({
  1440. stream: src.stream,
  1441. promptAsync: async (input) => {
  1442. seen.push(input)
  1443. queueMicrotask(() => {
  1444. src.push(busy())
  1445. src.push(idle())
  1446. })
  1447. return ok(undefined)
  1448. },
  1449. }),
  1450. sessionID: "session-1",
  1451. thinking: true,
  1452. limits: () => ({}),
  1453. footer: ui.api,
  1454. })
  1455. try {
  1456. await transport.runPromptTurn({
  1457. agent: undefined,
  1458. model: undefined,
  1459. variant: undefined,
  1460. prompt: { text: "hello", parts: [] },
  1461. files: [file],
  1462. includeFiles: true,
  1463. })
  1464. await transport.runPromptTurn({
  1465. agent: undefined,
  1466. model: undefined,
  1467. variant: undefined,
  1468. prompt: { text: "again", parts: [] },
  1469. files: [file],
  1470. includeFiles: false,
  1471. })
  1472. expect(seen).toEqual([
  1473. expect.objectContaining({
  1474. parts: [file, { type: "text", text: "hello" }],
  1475. }),
  1476. expect.objectContaining({
  1477. parts: [{ type: "text", text: "again" }],
  1478. }),
  1479. ])
  1480. } finally {
  1481. src.close()
  1482. await transport.close()
  1483. }
  1484. })
  1485. test("falls back to session status polling when idle events are missing", async () => {
  1486. const src = eventFeed()
  1487. const ui = footer()
  1488. let busy = true
  1489. const transport = await createSessionTransport({
  1490. sdk: sdk({
  1491. stream: src.stream,
  1492. promptAsync: async () => {
  1493. queueMicrotask(() => {
  1494. src.push(assistant("msg-1"))
  1495. busy = false
  1496. })
  1497. return ok(undefined)
  1498. },
  1499. status: async () => ok(statusMap(busy)),
  1500. }),
  1501. sessionID: "session-1",
  1502. thinking: true,
  1503. limits: () => ({}),
  1504. footer: ui.api,
  1505. })
  1506. try {
  1507. await Promise.race([
  1508. transport.runPromptTurn({
  1509. agent: undefined,
  1510. model: undefined,
  1511. variant: undefined,
  1512. prompt: { text: "hello", parts: [] },
  1513. files: [],
  1514. includeFiles: false,
  1515. }),
  1516. new Promise((_, reject) => setTimeout(() => reject(new Error("turn timed out")), 1_000)),
  1517. ])
  1518. } finally {
  1519. src.close()
  1520. await transport.close()
  1521. }
  1522. })
  1523. test("flushes interrupted output when the active turn aborts", async () => {
  1524. const src = eventFeed()
  1525. const seen = defer()
  1526. const ui = footer((commit) => {
  1527. if (commit.kind === "assistant" && commit.phase === "progress") {
  1528. seen.resolve()
  1529. }
  1530. })
  1531. const transport = await createSessionTransport({
  1532. sdk: sdk({
  1533. stream: src.stream,
  1534. promptAsync: async () => {
  1535. queueMicrotask(() => {
  1536. src.push(busy())
  1537. src.push(assistant("msg-1"))
  1538. src.push(textUpdated(textPart("txt-1", "msg-1", "")))
  1539. src.push(textDelta("msg-1", "txt-1", "unfinished"))
  1540. })
  1541. return ok(undefined)
  1542. },
  1543. }),
  1544. sessionID: "session-1",
  1545. thinking: true,
  1546. limits: () => ({}),
  1547. footer: ui.api,
  1548. })
  1549. const ctrl = new AbortController()
  1550. try {
  1551. const task = transport.runPromptTurn({
  1552. agent: undefined,
  1553. model: undefined,
  1554. variant: undefined,
  1555. prompt: { text: "hello", parts: [] },
  1556. files: [],
  1557. includeFiles: false,
  1558. signal: ctrl.signal,
  1559. })
  1560. await seen.promise
  1561. ctrl.abort()
  1562. await task
  1563. expect(ui.commits).toEqual([
  1564. {
  1565. kind: "assistant",
  1566. text: "unfinished",
  1567. phase: "progress",
  1568. source: "assistant",
  1569. messageID: "msg-1",
  1570. partID: "txt-1",
  1571. },
  1572. {
  1573. kind: "assistant",
  1574. text: "",
  1575. phase: "final",
  1576. source: "assistant",
  1577. messageID: "msg-1",
  1578. partID: "txt-1",
  1579. interrupted: true,
  1580. },
  1581. ])
  1582. } finally {
  1583. src.close()
  1584. await transport.close()
  1585. }
  1586. })
  1587. test("closes an active turn without rejecting it", async () => {
  1588. const src = eventFeed()
  1589. const ui = footer()
  1590. const ready = defer()
  1591. let aborted = false
  1592. const transport = await createSessionTransport({
  1593. sdk: sdk({
  1594. stream: src.stream,
  1595. promptAsync: async (_input, opt) => {
  1596. ready.resolve()
  1597. await new Promise<void>((resolve) => {
  1598. const onAbort = () => {
  1599. aborted = true
  1600. opt?.signal?.removeEventListener("abort", onAbort)
  1601. resolve()
  1602. }
  1603. opt?.signal?.addEventListener("abort", onAbort, { once: true })
  1604. })
  1605. return ok(undefined)
  1606. },
  1607. }),
  1608. sessionID: "session-1",
  1609. thinking: true,
  1610. limits: () => ({}),
  1611. footer: ui.api,
  1612. })
  1613. try {
  1614. const task = transport.runPromptTurn({
  1615. agent: undefined,
  1616. model: undefined,
  1617. variant: undefined,
  1618. prompt: { text: "hello", parts: [] },
  1619. files: [],
  1620. includeFiles: false,
  1621. })
  1622. await ready.promise
  1623. await transport.close()
  1624. await task
  1625. expect(aborted).toBe(true)
  1626. } finally {
  1627. src.close()
  1628. await transport.close()
  1629. }
  1630. })
  1631. test("rejects the active turn when the event stream faults", async () => {
  1632. const ui = footer()
  1633. const ready = defer()
  1634. const transport = await createSessionTransport({
  1635. sdk: sdk({
  1636. globalEvent: () =>
  1637. globalSse(
  1638. (async function* (): AsyncGenerator<GlobalEvent> {
  1639. await ready.promise
  1640. yield globalEvent(busy())
  1641. throw new Error("boom")
  1642. })(),
  1643. ),
  1644. promptAsync: async () => {
  1645. ready.resolve()
  1646. return ok(undefined)
  1647. },
  1648. status: async () => ok({ "session-1": { type: "busy" } }),
  1649. }),
  1650. sessionID: "session-1",
  1651. thinking: true,
  1652. limits: () => ({}),
  1653. footer: ui.api,
  1654. })
  1655. try {
  1656. await expect(
  1657. transport.runPromptTurn({
  1658. agent: undefined,
  1659. model: undefined,
  1660. variant: undefined,
  1661. prompt: { text: "hello", parts: [] },
  1662. files: [],
  1663. includeFiles: false,
  1664. }),
  1665. ).rejects.toThrow("boom")
  1666. } finally {
  1667. await transport.close()
  1668. }
  1669. })
  1670. test("rejects the active turn when the backing instance is disposed", async () => {
  1671. const ui = footer()
  1672. const ready = defer()
  1673. const transport = await createSessionTransport({
  1674. sdk: sdk({
  1675. globalEvent: () =>
  1676. globalSse(
  1677. (async function* (): AsyncGenerator<GlobalEvent> {
  1678. await ready.promise
  1679. yield globalEvent({
  1680. id: "evt-disposed",
  1681. type: "server.instance.disposed",
  1682. properties: {
  1683. directory: "/tmp",
  1684. },
  1685. })
  1686. })(),
  1687. ),
  1688. promptAsync: async () => {
  1689. ready.resolve()
  1690. return ok(undefined)
  1691. },
  1692. status: async () => ok({}),
  1693. }),
  1694. directory: "/tmp",
  1695. sessionID: "session-1",
  1696. thinking: true,
  1697. limits: () => ({}),
  1698. footer: ui.api,
  1699. })
  1700. try {
  1701. await expect(
  1702. transport.runPromptTurn({
  1703. agent: undefined,
  1704. model: undefined,
  1705. variant: undefined,
  1706. prompt: { text: "hello", parts: [] },
  1707. files: [],
  1708. includeFiles: false,
  1709. }),
  1710. ).rejects.toThrow("instance disposed")
  1711. } finally {
  1712. await transport.close()
  1713. }
  1714. })
  1715. test("rejects concurrent turns", async () => {
  1716. const src = eventFeed()
  1717. const ui = footer()
  1718. const transport = await createSessionTransport({
  1719. sdk: sdk({
  1720. stream: src.stream,
  1721. }),
  1722. sessionID: "session-1",
  1723. thinking: true,
  1724. limits: () => ({}),
  1725. footer: ui.api,
  1726. })
  1727. const ctrl = new AbortController()
  1728. try {
  1729. const task = transport.runPromptTurn({
  1730. agent: undefined,
  1731. model: undefined,
  1732. variant: undefined,
  1733. prompt: { text: "one", parts: [] },
  1734. files: [],
  1735. includeFiles: false,
  1736. signal: ctrl.signal,
  1737. })
  1738. await expect(
  1739. transport.runPromptTurn({
  1740. agent: undefined,
  1741. model: undefined,
  1742. variant: undefined,
  1743. prompt: { text: "two", parts: [] },
  1744. files: [],
  1745. includeFiles: false,
  1746. }),
  1747. ).rejects.toThrow("prompt already running")
  1748. ctrl.abort()
  1749. await task
  1750. } finally {
  1751. src.close()
  1752. await transport.close()
  1753. }
  1754. })
  1755. })