| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900 |
- import { afterEach, describe, expect, mock, spyOn, test } from "bun:test"
- import { OpencodeClient, type GlobalEvent } from "@opencode-ai/sdk/v2"
- import { createSessionTransport } from "@/cli/cmd/run/stream.transport"
- import type { FooterApi, FooterEvent, RunFilePart, StreamCommit } from "@/cli/cmd/run/types"
- type EventStream = Awaited<ReturnType<OpencodeClient["event"]["subscribe"]>>["stream"]
- type GlobalEventStream = Awaited<ReturnType<OpencodeClient["global"]["event"]>>["stream"]
- type SdkEvent = EventStream extends AsyncGenerator<infer T, unknown, unknown> ? T : never
- type SessionMessage = NonNullable<Awaited<ReturnType<OpencodeClient["session"]["messages"]>>["data"]>[number]
- type SessionChild = NonNullable<Awaited<ReturnType<OpencodeClient["session"]["children"]>>["data"]>[number]
- type SessionToolPart = Extract<SessionMessage["parts"][number], { type: "tool" }>
- type SessionStatusMap = NonNullable<Awaited<ReturnType<OpencodeClient["session"]["status"]>>["data"]>
- type TextPart = Extract<SessionMessage["parts"][number], { type: "text" }>
- afterEach(() => {
- mock.restore()
- })
- function defer<T = void>() {
- let resolve!: (value: T | PromiseLike<T>) => void
- let reject!: (error?: unknown) => void
- const promise = new Promise<T>((next, fail) => {
- resolve = next
- reject = fail
- })
- return { promise, resolve, reject }
- }
- async function waitFor<T>(check: () => T | undefined, timeout = 1_000): Promise<T> {
- const end = Date.now() + timeout
- while (Date.now() < end) {
- const value = check()
- if (value !== undefined) {
- return value
- }
- await Bun.sleep(10)
- }
- throw new Error("timed out waiting for value")
- }
- function busy(sessionID = "session-1") {
- return {
- id: `evt-${sessionID}-busy`,
- type: "session.status",
- properties: {
- sessionID,
- status: {
- type: "busy",
- },
- },
- } satisfies SdkEvent
- }
- function idle(sessionID = "session-1") {
- return {
- id: `evt-${sessionID}-idle`,
- type: "session.status",
- properties: {
- sessionID,
- status: {
- type: "idle",
- },
- },
- } satisfies SdkEvent
- }
- function retry(sessionID: string, attempt: number, message: string) {
- return {
- id: `evt-${sessionID}-retry-${attempt}`,
- type: "session.status",
- properties: {
- sessionID,
- status: {
- type: "retry",
- attempt,
- message,
- next: 1,
- },
- },
- } satisfies SdkEvent
- }
- function assistant(id: string) {
- return {
- id: `evt-${id}`,
- type: "message.updated",
- properties: {
- sessionID: "session-1",
- info: assistantMessage({
- sessionID: "session-1",
- id,
- parts: [],
- }).info,
- },
- } satisfies SdkEvent
- }
- const StreamClosed = undefined as never
- function feed<T, R = never>(returnValue: R = StreamClosed) {
- const list: T[] = []
- let done = false
- let wake: (() => void) | undefined
- const wrapped = (async function* (): AsyncGenerator<T, R, unknown> {
- while (!done || list.length > 0) {
- if (list.length === 0) {
- await new Promise<void>((resolve) => {
- wake = resolve
- })
- continue
- }
- const next = list.shift()
- if (!next) {
- continue
- }
- yield next
- }
- return returnValue as R
- })()
- return {
- stream: wrapped,
- push(value: T) {
- list.push(value)
- wake?.()
- wake = undefined
- },
- close() {
- done = true
- wake?.()
- wake = undefined
- },
- }
- }
- function eventFeed() {
- return feed<SdkEvent>()
- }
- function globalFeed() {
- return feed<GlobalEvent>()
- }
- function emptyStream(): EventStream {
- return (async function* (): AsyncGenerator<SdkEvent> {})()
- }
- function ok<T>(data: T) {
- return Promise.resolve({
- data,
- error: undefined,
- request: new Request("https://opencode.test"),
- response: new Response(),
- })
- }
- function sse(stream: EventStream) {
- return Promise.resolve({ stream })
- }
- function globalSse(stream: GlobalEventStream) {
- return Promise.resolve({ stream })
- }
- function wrapGlobalStream(stream: EventStream): GlobalEventStream {
- return (async function* (): GlobalEventStream {
- for await (const event of stream) {
- yield globalEvent(event as GlobalEvent["payload"])
- }
- return StreamClosed
- })()
- }
- function statusMap(busy: boolean): SessionStatusMap {
- if (busy) {
- return { "session-1": { type: "busy" } }
- }
- return {}
- }
- function assistantMessage(input: { sessionID: string; id: string; parts: SessionMessage["parts"] }): SessionMessage {
- return {
- info: {
- id: input.id,
- sessionID: input.sessionID,
- role: "assistant",
- time: {
- created: 1,
- },
- parentID: "msg-user-1",
- modelID: "gpt-5",
- providerID: "openai",
- mode: "chat",
- agent: "build",
- path: {
- cwd: "/tmp",
- root: "/tmp",
- },
- cost: 0,
- tokens: {
- input: 1,
- output: 1,
- reasoning: 0,
- cache: {
- read: 0,
- write: 0,
- },
- },
- },
- parts: input.parts,
- }
- }
- function runningTool(input: {
- sessionID: string
- messageID: string
- id: string
- callID: string
- tool: string
- body: Record<string, unknown>
- metadata?: Record<string, unknown>
- }): SessionToolPart {
- return {
- id: input.id,
- sessionID: input.sessionID,
- messageID: input.messageID,
- type: "tool",
- callID: input.callID,
- tool: input.tool,
- state: {
- status: "running",
- input: input.body,
- ...(input.metadata ? { metadata: input.metadata } : {}),
- time: {
- start: 1,
- },
- },
- }
- }
- function completedTool(input: {
- sessionID: string
- messageID: string
- id: string
- callID: string
- tool: string
- body: Record<string, unknown>
- output?: string
- metadata?: Record<string, unknown>
- }): SessionToolPart {
- return {
- id: input.id,
- sessionID: input.sessionID,
- messageID: input.messageID,
- type: "tool",
- callID: input.callID,
- tool: input.tool,
- state: {
- status: "completed",
- input: input.body,
- output: input.output ?? "",
- title: input.tool,
- metadata: input.metadata ?? {},
- time: {
- start: 1,
- end: 2,
- },
- },
- }
- }
- function textPart(id: string, messageID: string, text: string, sessionID = "session-1"): TextPart {
- return {
- id,
- sessionID,
- messageID,
- type: "text",
- text,
- }
- }
- function textUpdated(part: TextPart): SdkEvent {
- return {
- id: `evt-${part.id}-updated`,
- type: "message.part.updated",
- properties: {
- sessionID: part.sessionID,
- part,
- time: 1,
- },
- }
- }
- function toolUpdated(part: SessionToolPart): SdkEvent {
- return {
- id: `evt-${part.id}-updated`,
- type: "message.part.updated",
- properties: {
- sessionID: part.sessionID,
- part,
- time: 1,
- },
- }
- }
- function textDelta(messageID: string, partID: string, delta: string, sessionID = "session-1"): SdkEvent {
- return {
- id: `evt-${partID}-delta`,
- type: "message.part.delta",
- properties: {
- sessionID,
- messageID,
- partID,
- field: "text",
- delta,
- },
- }
- }
- function child(id: string): SessionChild {
- return {
- id,
- slug: id,
- projectID: "project-1",
- directory: "/tmp",
- title: id,
- version: "1",
- time: {
- created: 1,
- updated: 1,
- },
- }
- }
- function globalEvent(payload: SdkEvent | GlobalEvent["payload"]): GlobalEvent {
- return {
- directory: "/tmp",
- project: "project-1",
- payload: payload as GlobalEvent["payload"],
- }
- }
- function footer(fn?: (commit: StreamCommit) => void) {
- const commits: StreamCommit[] = []
- const events: FooterEvent[] = []
- let closed = false
- let idleCalls = 0
- const api: FooterApi = {
- get isClosed() {
- return closed
- },
- onPrompt: () => () => {},
- onClose: () => () => {},
- event(next) {
- events.push(next)
- },
- append(next) {
- commits.push(next)
- fn?.(next)
- },
- idle() {
- idleCalls += 1
- return Promise.resolve()
- },
- close() {
- closed = true
- },
- destroy() {
- closed = true
- },
- }
- return {
- api,
- commits,
- events,
- get idleCalls() {
- return idleCalls
- },
- }
- }
- function sdk(
- input: {
- stream?: EventStream
- globalStream?: GlobalEventStream
- subscribe?: OpencodeClient["event"]["subscribe"]
- globalEvent?: OpencodeClient["global"]["event"]
- promptAsync?: OpencodeClient["session"]["promptAsync"]
- status?: OpencodeClient["session"]["status"]
- messages?: OpencodeClient["session"]["messages"]
- children?: OpencodeClient["session"]["children"]
- permissions?: OpencodeClient["permission"]["list"]
- questions?: OpencodeClient["question"]["list"]
- } = {},
- ) {
- const client = new OpencodeClient()
- const subscribe: OpencodeClient["event"]["subscribe"] = input.subscribe ?? (() => sse(input.stream ?? emptyStream()))
- const globalEvent: OpencodeClient["global"]["event"] =
- input.globalEvent ?? (() => globalSse(input.globalStream ?? wrapGlobalStream(input.stream ?? emptyStream())))
- const promptAsync: OpencodeClient["session"]["promptAsync"] = input.promptAsync ?? (() => ok(undefined))
- const status: OpencodeClient["session"]["status"] = input.status ?? (() => ok({}))
- const messages: OpencodeClient["session"]["messages"] = input.messages ?? (() => ok([]))
- const children: OpencodeClient["session"]["children"] = input.children ?? (() => ok([]))
- const permissions: OpencodeClient["permission"]["list"] = input.permissions ?? (() => ok([]))
- const questions: OpencodeClient["question"]["list"] = input.questions ?? (() => ok([]))
- spyOn(client.event, "subscribe").mockImplementation(subscribe)
- spyOn(client.global, "event").mockImplementation(globalEvent)
- spyOn(client.session, "promptAsync").mockImplementation(promptAsync)
- spyOn(client.session, "status").mockImplementation(status)
- spyOn(client.session, "messages").mockImplementation(messages)
- spyOn(client.session, "children").mockImplementation(children)
- spyOn(client.permission, "list").mockImplementation(permissions)
- spyOn(client.question, "list").mockImplementation(questions)
- return client
- }
- describe("run stream transport", () => {
- test("does not replay persisted main-session history during bootstrap by default", async () => {
- const src = eventFeed()
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async ({ sessionID }) =>
- sessionID === "session-1"
- ? ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [
- {
- ...textPart("text-1", "msg-1", "Hello."),
- time: {
- start: 1,
- end: 2,
- },
- },
- ],
- }),
- ])
- : ok([]),
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- expect(ui.commits).toEqual([])
- expect(ui.idleCalls).toBe(0)
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("replays persisted main-session history during bootstrap when enabled", async () => {
- const src = eventFeed()
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async ({ sessionID }) =>
- sessionID === "session-1"
- ? ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [
- {
- ...textPart("text-1", "msg-1", "Hello."),
- time: {
- start: 1,
- end: 2,
- },
- },
- ],
- }),
- ])
- : ok([]),
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await waitFor(() => ui.commits.find((item) => item.kind === "assistant" && item.text === "Hello."))
- expect(ui.idleCalls).toBeGreaterThan(0)
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("caps replayed bootstrap history to the configured number of messages", async () => {
- const src = eventFeed()
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async ({ sessionID }) =>
- ok(
- sessionID === "session-1"
- ? [
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [
- {
- ...textPart("text-1", "msg-1", "Hello."),
- time: {
- start: 1,
- end: 2,
- },
- },
- ],
- }),
- assistantMessage({
- sessionID: "session-1",
- id: "msg-2",
- parts: [
- {
- ...textPart("text-2", "msg-2", "World."),
- time: {
- start: 3,
- end: 4,
- },
- },
- ],
- }),
- ]
- : [],
- ),
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- replayLimit: 1,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await waitFor(() => (ui.commits.length > 0 ? ui.commits : undefined))
- expect(ui.commits.filter((item) => item.kind === "assistant")).toEqual([
- expect.objectContaining({
- text: "World.",
- }),
- ])
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("skips buffered pre-bootstrap deltas already covered by replay history", async () => {
- const src = eventFeed()
- const ui = footer()
- const gate = defer<void>()
- let transport: Awaited<ReturnType<typeof createSessionTransport>> | undefined
- const task = createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async ({ sessionID }) => {
- if (sessionID !== "session-1") {
- return ok([])
- }
- await gate.promise
- return ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [textPart("text-1", "msg-1", "Hello")],
- }),
- ])
- },
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await Promise.resolve()
- src.push(textDelta("msg-1", "text-1", "lo"))
- gate.resolve()
- transport = await task
- await waitFor(() => (ui.commits.length > 0 ? ui.commits : undefined))
- await Bun.sleep(20)
- expect(ui.commits.filter((item) => item.kind === "assistant")).toEqual([
- expect.objectContaining({
- text: "Hello",
- }),
- ])
- } finally {
- src.close()
- await transport?.close()
- }
- })
- test("applies buffered pre-bootstrap deltas not yet persisted", async () => {
- const src = eventFeed()
- const ui = footer()
- const gate = defer<void>()
- let transport: Awaited<ReturnType<typeof createSessionTransport>> | undefined
- const task = createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async ({ sessionID }) => {
- if (sessionID !== "session-1") {
- return ok([])
- }
- await gate.promise
- return ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [textPart("text-1", "msg-1", "")],
- }),
- ])
- },
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await Promise.resolve()
- src.push(textDelta("msg-1", "text-1", "Hello"))
- gate.resolve()
- transport = await task
- await waitFor(() => (ui.commits.length > 0 ? ui.commits : undefined))
- await Bun.sleep(20)
- expect(ui.commits.filter((item) => item.kind === "assistant")).toEqual([
- expect.objectContaining({
- text: "Hello",
- }),
- ])
- } finally {
- src.close()
- await transport?.close()
- }
- })
- test("preserves running footer state for resumed active sessions", async () => {
- const src = eventFeed()
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async ({ sessionID }) =>
- sessionID === "session-1"
- ? ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [
- runningTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "bash-1",
- callID: "call-1",
- tool: "bash",
- body: {
- command: "pwd",
- },
- }),
- ],
- }),
- ])
- : ok([]),
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- const patch = await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.patch")
- return item?.type === "stream.patch" ? item.patch : undefined
- })
- expect(patch).toEqual(
- expect.objectContaining({
- phase: "running",
- status: "running bash",
- }),
- )
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("drops completed historical subagent tabs during bootstrap", async () => {
- const src = eventFeed()
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async ({ sessionID }) => {
- if (sessionID !== "session-1") {
- return ok([])
- }
- return ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [
- completedTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "task-1",
- callID: "call-1",
- tool: "task",
- body: {
- description: "Explore run folder",
- subagent_type: "explore",
- },
- metadata: {
- sessionId: "child-1",
- },
- }),
- ],
- }),
- ])
- },
- children: async () => ok([child("child-1")]),
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- const state = await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- return item?.type === "stream.subagent" ? item.state : undefined
- })
- expect(state.tabs).toEqual([])
- expect(state.details).toEqual({})
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("bootstraps child tabs and resumed blocker input", async () => {
- const src = eventFeed()
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async ({ sessionID }) => {
- if (sessionID === "session-1") {
- return ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [
- runningTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "task-1",
- callID: "call-1",
- tool: "task",
- body: {
- description: "Explore run folder",
- subagent_type: "explore",
- },
- metadata: {
- sessionId: "child-1",
- },
- }),
- ],
- }),
- ])
- }
- return ok([
- assistantMessage({
- sessionID: "child-1",
- id: "msg-child-1",
- parts: [
- runningTool({
- sessionID: "child-1",
- messageID: "msg-child-1",
- id: "edit-1",
- callID: "call-edit-1",
- tool: "edit",
- body: {
- filePath: "src/run/subagent-data.ts",
- diff: "@@ -1 +1 @@",
- },
- }),
- ],
- }),
- ])
- },
- children: async () => ok([child("child-1")]),
- permissions: async () =>
- ok([
- {
- id: "perm-1",
- sessionID: "child-1",
- permission: "edit",
- patterns: ["src/run/subagent-data.ts"],
- metadata: {},
- always: [],
- tool: {
- messageID: "msg-child-1",
- callID: "call-edit-1",
- },
- },
- ]),
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- const boot = await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- const state = item?.type === "stream.subagent" ? item.state : undefined
- return state?.tabs.some((tab) => tab.sessionID === "child-1") &&
- state.permissions.some((req) => req.id === "perm-1")
- ? state
- : undefined
- })
- expect(boot.tabs).toEqual([
- expect.objectContaining({
- sessionID: "child-1",
- label: "Explore",
- description: "Pending permission",
- status: "running",
- }),
- ])
- expect(boot.permissions).toEqual([
- expect.objectContaining({
- id: "perm-1",
- sessionID: "child-1",
- metadata: {
- input: {
- filePath: "src/run/subagent-data.ts",
- diff: "@@ -1 +1 @@",
- },
- },
- }),
- ])
- transport.selectSubagent("child-1")
- const selected = await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- const state = item?.type === "stream.subagent" ? item.state : undefined
- const detail = state?.details["child-1"]
- return detail?.commits.some(
- (commit) => commit.kind === "tool" && commit.tool === "edit" && commit.phase === "start",
- )
- ? state
- : undefined
- })
- expect(selected.details).toEqual({
- "child-1": {
- sessionID: "child-1",
- commits: [
- expect.objectContaining({
- kind: "tool",
- tool: "edit",
- phase: "start",
- }),
- ],
- },
- })
- expect(
- await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.view")
- return item?.type === "stream.view" && item.view.type === "permission" && item.view.request.id === "perm-1"
- ? item
- : undefined
- }),
- ).toEqual({
- type: "stream.view",
- view: {
- type: "permission",
- request: expect.objectContaining({
- id: "perm-1",
- metadata: {
- input: {
- filePath: "src/run/subagent-data.ts",
- diff: "@@ -1 +1 @@",
- },
- },
- }),
- },
- })
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("bootstraps child session output before selection", async () => {
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({
- messages: async ({ sessionID }) => {
- if (sessionID === "session-1") {
- return ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [
- runningTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "task-1",
- callID: "call-1",
- tool: "task",
- body: {
- description: "Explore run.ts",
- subagent_type: "explore",
- },
- metadata: {
- sessionId: "child-1",
- },
- }),
- ],
- }),
- ])
- }
- return sessionID === "child-1"
- ? ok([
- assistantMessage({
- sessionID: "child-1",
- id: "msg-child-1",
- parts: [textPart("txt-child-1", "msg-child-1", "subagent summary", "child-1")],
- }),
- ])
- : ok([])
- },
- children: async () => ok([child("child-1")]),
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- return item?.type === "stream.subagent" && item.state.tabs.some((tab) => tab.sessionID === "child-1")
- ? item
- : undefined
- })
- transport.selectSubagent("child-1")
- expect(
- await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- const detail = item?.type === "stream.subagent" ? item.state.details["child-1"] : undefined
- return detail?.commits.some((commit) => commit.kind === "assistant" && commit.text === "subagent summary")
- ? detail
- : undefined
- }),
- ).toEqual({
- sessionID: "child-1",
- commits: [
- expect.objectContaining({
- kind: "assistant",
- text: "subagent summary",
- }),
- ],
- })
- } finally {
- await transport.close()
- }
- })
- test("does not block startup on child history bootstrap", async () => {
- const pending = defer<Awaited<ReturnType<typeof ok<SessionMessage[]>>>>()
- const ui = footer()
- let transport: Awaited<ReturnType<typeof createSessionTransport>> | undefined
- const task = createSessionTransport({
- sdk: sdk({
- messages: async ({ sessionID }) => {
- if (sessionID === "session-1") {
- return ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-1",
- parts: [
- runningTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "task-1",
- callID: "call-1",
- tool: "task",
- body: {
- description: "Explore run.ts",
- subagent_type: "explore",
- },
- metadata: {
- sessionId: "child-1",
- },
- }),
- ],
- }),
- ])
- }
- if (sessionID === "child-1") {
- return pending.promise
- }
- return ok([])
- },
- children: async () => ok([child("child-1")]),
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- }).then((item) => {
- transport = item
- return item
- })
- try {
- const state = await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- return item?.type === "stream.subagent" && item.state.tabs.some((tab) => tab.sessionID === "child-1")
- ? item.state
- : undefined
- })
- await waitFor(() => transport)
- expect(state).toEqual({
- tabs: [expect.objectContaining({ sessionID: "child-1", status: "running" })],
- details: {},
- permissions: [],
- questions: [],
- })
- } finally {
- pending.resolve(ok([]))
- await task
- await transport?.close()
- }
- })
- test("replays child events buffered during bootstrap once the tab is known", async () => {
- const global = globalFeed()
- const ui = footer()
- const gate = defer<void>()
- let transport: Awaited<ReturnType<typeof createSessionTransport>> | undefined
- const task = createSessionTransport({
- sdk: sdk({
- globalStream: global.stream,
- messages: async ({ sessionID }) => {
- if (sessionID !== "session-1") {
- return ok([])
- }
- await gate.promise
- return ok([])
- },
- children: async () => ok([]),
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await Promise.resolve()
- global.push(globalEvent(retry("child-1", 1, "retry child")))
- global.push(
- globalEvent({
- id: "evt-child-message",
- type: "message.updated",
- properties: {
- sessionID: "child-1",
- info: assistantMessage({
- sessionID: "child-1",
- id: "msg-child-1",
- parts: [],
- }).info,
- },
- }),
- )
- global.push(globalEvent(textUpdated(textPart("txt-child-1", "msg-child-1", "", "child-1"))))
- global.push(globalEvent(textDelta("msg-child-1", "txt-child-1", "Hello", "child-1")))
- global.push(
- globalEvent(
- toolUpdated(
- runningTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "task-1",
- callID: "call-1",
- tool: "task",
- body: {
- description: "Explore run.ts",
- subagent_type: "explore",
- },
- metadata: {
- sessionId: "child-1",
- },
- }),
- ),
- ),
- )
- gate.resolve()
- transport = await task
- await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- return item?.type === "stream.subagent" && item.state.tabs.some((tab) => tab.sessionID === "child-1")
- ? item
- : undefined
- })
- transport.selectSubagent("child-1")
- const detail = await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- const next = item?.type === "stream.subagent" ? item.state.details["child-1"] : undefined
- return next?.commits.some((commit) => commit.kind === "error" && commit.text === "retry child") &&
- next.commits.some((commit) => commit.kind === "assistant" && commit.text === "Hello")
- ? next
- : undefined
- })
- expect(detail).toEqual({
- sessionID: "child-1",
- commits: expect.arrayContaining([
- expect.objectContaining({
- kind: "error",
- text: "retry child",
- }),
- expect.objectContaining({
- kind: "assistant",
- text: "Hello",
- }),
- ]),
- })
- } finally {
- global.close()
- await transport?.close()
- }
- })
- test("streams selected subagent output from global events while it is running", async () => {
- const global = globalFeed()
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({
- globalStream: global.stream,
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- global.push(globalEvent(assistant("msg-1")))
- global.push(
- globalEvent(
- toolUpdated(
- runningTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "task-1",
- callID: "call-1",
- tool: "task",
- body: {
- description: "Explore run.ts",
- subagent_type: "explore",
- },
- metadata: {
- sessionId: "child-1",
- },
- }),
- ),
- ),
- )
- await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- return item?.type === "stream.subagent" && item.state.tabs.some((tab) => tab.sessionID === "child-1")
- ? item
- : undefined
- })
- transport.selectSubagent("child-1")
- global.push(
- globalEvent({
- id: "evt-child-message",
- type: "message.updated",
- properties: {
- sessionID: "child-1",
- info: assistantMessage({
- sessionID: "child-1",
- id: "msg-child-1",
- parts: [],
- }).info,
- },
- }),
- )
- global.push(globalEvent(textUpdated(textPart("txt-child-1", "msg-child-1", "hello", "child-1"))))
- expect(
- await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- const detail = item?.type === "stream.subagent" ? item.state.details["child-1"] : undefined
- return detail?.commits.some((commit) => commit.kind === "assistant" && commit.text === "hello")
- ? detail
- : undefined
- }),
- ).toEqual({
- sessionID: "child-1",
- commits: [
- expect.objectContaining({
- kind: "assistant",
- text: "hello",
- }),
- ],
- })
- global.push(globalEvent(textUpdated(textPart("txt-child-1", "msg-child-1", "hello world", "child-1"))))
- expect(
- await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.subagent")
- const detail = item?.type === "stream.subagent" ? item.state.details["child-1"] : undefined
- return detail?.commits.some((commit) => commit.kind === "assistant" && commit.text === "hello world")
- ? detail
- : undefined
- }, 2_000),
- ).toEqual({
- sessionID: "child-1",
- commits: [
- expect.objectContaining({
- kind: "assistant",
- text: "hello world",
- }),
- ],
- })
- } finally {
- global.close()
- await transport.close()
- }
- })
- test("recovers pending questions from question.list when question.asked is missed", async () => {
- const src = eventFeed()
- const ui = footer()
- let questionCalls = 0
- const request = {
- id: "question-1",
- sessionID: "session-1",
- questions: [
- {
- question: "Which area should I inspect first?",
- header: "Area",
- options: [{ label: "CLI", description: "Look at the direct run flow." }],
- multiple: false,
- },
- ],
- tool: {
- messageID: "msg-1",
- callID: "call-question-1",
- },
- }
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- questions: async () => {
- questionCalls += 1
- return ok(questionCalls > 1 ? [request] : [])
- },
- promptAsync: async () => {
- queueMicrotask(() => {
- src.push(busy())
- src.push(assistant("msg-1"))
- src.push(
- toolUpdated(
- runningTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "question-tool-1",
- callID: "call-question-1",
- tool: "question",
- body: {
- questions: request.questions,
- },
- }),
- ),
- )
- })
- return ok(undefined)
- },
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- const ctrl = new AbortController()
- try {
- const run = transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "hello", parts: [] },
- files: [],
- includeFiles: false,
- signal: ctrl.signal,
- })
- const view = await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.view")
- return item?.type === "stream.view" && item.view.type === "question" ? item.view : undefined
- })
- expect(view).toEqual({
- type: "question",
- request,
- })
- expect(ui.events).toContainEqual({
- type: "stream.patch",
- patch: {
- phase: "running",
- status: "awaiting answer",
- },
- })
- src.push(
- toolUpdated(
- completedTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "question-tool-1",
- callID: "call-question-1",
- tool: "question",
- body: {
- questions: request.questions,
- },
- output: "User has answered your questions.",
- metadata: {
- answers: [["CLI"]],
- },
- }),
- ),
- )
- expect(
- await waitFor(() => {
- const item = ui.events.findLast((event) => event.type === "stream.view")
- return item?.type === "stream.view" && item.view.type === "prompt" ? item : undefined
- }),
- ).toEqual({
- type: "stream.view",
- view: { type: "prompt" },
- })
- ctrl.abort()
- await run
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("does not resurrect questions if question.list resolves after tool completion", async () => {
- const src = eventFeed()
- const ui = footer()
- const started = defer()
- const request = {
- id: "question-race-1",
- sessionID: "session-1",
- questions: [
- {
- question: "Which area should I inspect first?",
- header: "Area",
- options: [{ label: "CLI", description: "Look at the direct run flow." }],
- multiple: false,
- },
- ],
- tool: {
- messageID: "msg-1",
- callID: "call-question-race-1",
- },
- }
- const pending = defer<Awaited<ReturnType<typeof ok<(typeof request)[]>>>>()
- let questionCalls = 0
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- questions: async () => {
- questionCalls += 1
- if (questionCalls === 1) {
- return ok([])
- }
- if (questionCalls === 2) {
- started.resolve()
- return pending.promise
- }
- return ok([])
- },
- promptAsync: async () => {
- queueMicrotask(() => {
- src.push(busy())
- src.push(assistant("msg-1"))
- src.push(
- toolUpdated(
- runningTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "question-race-tool-1",
- callID: "call-question-race-1",
- tool: "question",
- body: {
- questions: request.questions,
- },
- }),
- ),
- )
- })
- return ok(undefined)
- },
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- const ctrl = new AbortController()
- try {
- const run = transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "hello", parts: [] },
- files: [],
- includeFiles: false,
- signal: ctrl.signal,
- })
- await started.promise
- src.push(
- toolUpdated(
- completedTool({
- sessionID: "session-1",
- messageID: "msg-1",
- id: "question-race-tool-1",
- callID: "call-question-race-1",
- tool: "question",
- body: {
- questions: request.questions,
- },
- output: "User has answered your questions.",
- metadata: {
- answers: [["CLI"]],
- },
- }),
- ),
- )
- await waitFor(() => {
- const commit = ui.commits.findLast(
- (item) => item.kind === "tool" && item.partID === "question-race-tool-1" && item.toolState === "completed",
- )
- return commit ? true : undefined
- })
- pending.resolve(ok([request]))
- await Bun.sleep(50)
- expect(
- ui.events.some(
- (event) =>
- event.type === "stream.view" && event.view.type === "question" && event.view.request.id === request.id,
- ),
- ).toBe(false)
- ctrl.abort()
- await run
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("respects the includeFiles flag when building prompt payloads", async () => {
- const src = eventFeed()
- const ui = footer()
- const seen: unknown[] = []
- const file: RunFilePart = {
- type: "file",
- url: "file:///tmp/a.ts",
- filename: "a.ts",
- mime: "text/plain",
- }
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- promptAsync: async (input) => {
- seen.push(input)
- queueMicrotask(() => {
- src.push(busy())
- src.push(idle())
- })
- return ok(undefined)
- },
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "hello", parts: [] },
- files: [file],
- includeFiles: true,
- })
- await transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "again", parts: [] },
- files: [file],
- includeFiles: false,
- })
- expect(seen).toEqual([
- expect.objectContaining({
- parts: [file, { type: "text", text: "hello" }],
- }),
- expect.objectContaining({
- parts: [{ type: "text", text: "again" }],
- }),
- ])
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("falls back to session status polling when idle events are missing", async () => {
- const src = eventFeed()
- const ui = footer()
- let busy = true
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- promptAsync: async () => {
- queueMicrotask(() => {
- src.push(assistant("msg-1"))
- busy = false
- })
- return ok(undefined)
- },
- status: async () => ok(statusMap(busy)),
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await Promise.race([
- transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "hello", parts: [] },
- files: [],
- includeFiles: false,
- }),
- new Promise((_, reject) => setTimeout(() => reject(new Error("turn timed out")), 1_000)),
- ])
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("flushes interrupted output when the active turn aborts", async () => {
- const src = eventFeed()
- const seen = defer()
- const ui = footer((commit) => {
- if (commit.kind === "assistant" && commit.phase === "progress") {
- seen.resolve()
- }
- })
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- promptAsync: async () => {
- queueMicrotask(() => {
- src.push(busy())
- src.push(assistant("msg-1"))
- src.push(textUpdated(textPart("txt-1", "msg-1", "")))
- src.push(textDelta("msg-1", "txt-1", "unfinished"))
- })
- return ok(undefined)
- },
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- const ctrl = new AbortController()
- try {
- const task = transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "hello", parts: [] },
- files: [],
- includeFiles: false,
- signal: ctrl.signal,
- })
- await seen.promise
- ctrl.abort()
- await task
- expect(ui.commits).toEqual([
- {
- kind: "assistant",
- text: "unfinished",
- phase: "progress",
- source: "assistant",
- messageID: "msg-1",
- partID: "txt-1",
- },
- {
- kind: "assistant",
- text: "",
- phase: "final",
- source: "assistant",
- messageID: "msg-1",
- partID: "txt-1",
- interrupted: true,
- },
- ])
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("closes an active turn without rejecting it", async () => {
- const src = eventFeed()
- const ui = footer()
- const ready = defer()
- let aborted = false
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- promptAsync: async (_input, opt) => {
- ready.resolve()
- await new Promise<void>((resolve) => {
- const onAbort = () => {
- aborted = true
- opt?.signal?.removeEventListener("abort", onAbort)
- resolve()
- }
- opt?.signal?.addEventListener("abort", onAbort, { once: true })
- })
- return ok(undefined)
- },
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- const task = transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "hello", parts: [] },
- files: [],
- includeFiles: false,
- })
- await ready.promise
- await transport.close()
- await task
- expect(aborted).toBe(true)
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("rejects the active turn when the event stream faults", async () => {
- const ui = footer()
- const ready = defer()
- const transport = await createSessionTransport({
- sdk: sdk({
- globalEvent: () =>
- globalSse(
- (async function* (): AsyncGenerator<GlobalEvent> {
- await ready.promise
- yield globalEvent(busy())
- throw new Error("boom")
- })(),
- ),
- promptAsync: async () => {
- ready.resolve()
- return ok(undefined)
- },
- status: async () => ok({ "session-1": { type: "busy" } }),
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await expect(
- transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "hello", parts: [] },
- files: [],
- includeFiles: false,
- }),
- ).rejects.toThrow("boom")
- } finally {
- await transport.close()
- }
- })
- test("rejects the active turn when the backing instance is disposed", async () => {
- const ui = footer()
- const ready = defer()
- const transport = await createSessionTransport({
- sdk: sdk({
- globalEvent: () =>
- globalSse(
- (async function* (): AsyncGenerator<GlobalEvent> {
- await ready.promise
- yield globalEvent({
- id: "evt-disposed",
- type: "server.instance.disposed",
- properties: {
- directory: "/tmp",
- },
- })
- })(),
- ),
- promptAsync: async () => {
- ready.resolve()
- return ok(undefined)
- },
- status: async () => ok({}),
- }),
- directory: "/tmp",
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- await expect(
- transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "hello", parts: [] },
- files: [],
- includeFiles: false,
- }),
- ).rejects.toThrow("instance disposed")
- } finally {
- await transport.close()
- }
- })
- test("rejects concurrent turns", async () => {
- const src = eventFeed()
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- }),
- sessionID: "session-1",
- thinking: true,
- limits: () => ({}),
- footer: ui.api,
- })
- const ctrl = new AbortController()
- try {
- const task = transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "one", parts: [] },
- files: [],
- includeFiles: false,
- signal: ctrl.signal,
- })
- await expect(
- transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "two", parts: [] },
- files: [],
- includeFiles: false,
- }),
- ).rejects.toThrow("prompt already running")
- ctrl.abort()
- await task
- } finally {
- src.close()
- await transport.close()
- }
- })
- })
|