| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363 |
- import { afterEach, describe, expect, mock, spyOn, test } from "bun:test"
- import { OpencodeClient, type GlobalEvent } from "@kirincode-ai/sdk/v2"
- import { createSessionTransport } from "@/cli/cmd/run/stream.transport"
- import type { FooterApi, FooterEvent, LocalReplayRow, 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" }>
- type ReasoningPart = Extract<SessionMessage["parts"][number], { type: "reasoning" }>
- 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)
- }
- 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 reasoningPart(id: string, messageID: string, text: string): ReasoningPart {
- return {
- id,
- sessionID: "session-1",
- messageID,
- type: "reasoning",
- text,
- time: { start: 1 },
- }
- }
- function reasoningUpdated(part: ReasoningPart): 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: GlobalEvent["payload"]): GlobalEvent {
- return {
- directory: "/tmp",
- project: "project-1",
- 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: () => () => {},
- onQueuedRemove: () => () => {},
- 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("rebuilds session output on resize and continues live deltas from replayed state", async () => {
- const src = eventFeed()
- const ui = footer()
- let calls = 0
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async () => {
- calls += 1
- if (calls === 1) {
- return ok([])
- }
- 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,
- })
- const localRows: LocalReplayRow[] = [
- { commit: { kind: "user", text: "pending prompt", phase: "start", source: "system", messageID: "msg-pending" } },
- ]
- const reset = mock(() => {
- localRows.push({
- commit: {
- kind: "user",
- text: "sent during reset",
- phase: "start",
- source: "system",
- messageID: "msg-during-reset",
- },
- })
- return Promise.resolve()
- })
- try {
- expect(
- await transport.replayOnResize({
- localRows: () => localRows,
- reset,
- }),
- ).toBe(true)
- expect(reset).toHaveBeenCalledTimes(1)
- expect(ui.commits).toEqual(
- expect.arrayContaining([
- expect.objectContaining({ kind: "assistant", text: "Hello" }),
- expect.objectContaining({ kind: "user", text: "sent during reset", messageID: "msg-during-reset" }),
- ]),
- )
- src.push(textUpdated(textPart("text-1", "msg-1", "Hello world")))
- await waitFor(() => ui.commits.find((commit) => commit.kind === "assistant" && commit.text === " world"))
- expect(ui.commits.filter((commit) => commit.kind === "assistant").map((commit) => commit.text)).toEqual([
- "Hello",
- " world",
- ])
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("coalesces active resize requests into one trailing replay", async () => {
- const src = eventFeed()
- const ui = footer()
- const firstReset = defer()
- const resetA = mock(() => firstReset.promise)
- const resetB = mock(() => Promise.resolve())
- const resetC = mock(() => Promise.resolve())
- const transport = await createSessionTransport({
- sdk: sdk({ stream: src.stream }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- const active = transport.replayOnResize({ localRows: () => [], reset: resetA })
- await waitFor(() => (resetA.mock.calls.length === 1 ? true : undefined))
- expect(await transport.replayOnResize({ localRows: () => [], reset: resetB })).toBe(false)
- expect(await transport.replayOnResize({ localRows: () => [], reset: resetC })).toBe(false)
- expect(resetB).not.toHaveBeenCalled()
- firstReset.resolve()
- expect(await active).toBe(true)
- expect(resetA).toHaveBeenCalledTimes(1)
- expect(resetB).not.toHaveBeenCalled()
- expect(resetC).toHaveBeenCalledTimes(1)
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("keeps coalescing resize requests while buffered events drain", async () => {
- const src = eventFeed()
- const ui = footer()
- const firstReset = defer()
- const statusGate = defer()
- const statusStarted = defer()
- let blockStatus = false
- const trace = mock((_type: string, _data?: unknown) => {})
- const resetA = mock(() => firstReset.promise)
- const resetB = mock(() => Promise.resolve())
- const resetC = mock(() => Promise.resolve())
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- status: async () => {
- if (blockStatus) {
- statusStarted.resolve()
- await statusGate.promise
- }
- return ok(statusMap(true))
- },
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- trace: { write: trace },
- })
- const turn = transport.runPromptTurn({
- agent: undefined,
- model: undefined,
- variant: undefined,
- prompt: { text: "active", parts: [] },
- files: [],
- includeFiles: false,
- })
- try {
- await waitFor(() => ui.events.find((event) => event.type === "turn.wait"))
- const active = transport.replayOnResize({ localRows: () => [], reset: resetA })
- await waitFor(() => (resetA.mock.calls.length === 1 ? true : undefined))
- blockStatus = true
- src.push(busy())
- src.push(idle())
- await waitFor(() => (trace.mock.calls.filter((call) => call[0] === "recv.event").length >= 2 ? true : undefined))
- expect(await transport.replayOnResize({ localRows: () => [], reset: resetB })).toBe(false)
- firstReset.resolve()
- await Promise.race([
- statusStarted.promise,
- Bun.sleep(1_000).then(() => {
- throw new Error("timed out waiting for buffered status drain")
- }),
- ])
- expect(await transport.replayOnResize({ localRows: () => [], reset: resetC })).toBe(false)
- expect(resetC).not.toHaveBeenCalled()
- blockStatus = false
- statusGate.resolve()
- expect(
- await Promise.race([
- active,
- Bun.sleep(1_000).then(() => {
- throw new Error("timed out waiting for trailing resize replay")
- }),
- ]),
- ).toBe(true)
- expect(resetB).not.toHaveBeenCalled()
- expect(resetC).toHaveBeenCalledTimes(1)
- } finally {
- src.close()
- await transport.close()
- await turn
- }
- })
- test("preserves assistant deltas not yet persisted when replaying during a live stream", async () => {
- const src = eventFeed()
- const ui = footer()
- let calls = 0
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async () => {
- calls += 1
- if (calls === 1) {
- return ok([])
- }
- return ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-live",
- parts: [textPart("text-live", "msg-live", "")],
- }),
- ])
- },
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- src.push(assistant("msg-live"))
- src.push(textUpdated(textPart("text-live", "msg-live", "")))
- src.push(textDelta("msg-live", "text-live", "Hello"))
- await waitFor(() => ui.commits.find((commit) => commit.kind === "assistant" && commit.text === "Hello"))
- ui.commits.length = 0
- expect(await transport.replayOnResize({ localRows: () => [], reset: () => Promise.resolve() })).toBe(true)
- src.push(textDelta("msg-live", "text-live", "Hello"))
- src.push(
- textUpdated({
- ...textPart("text-live", "msg-live", "HelloHello"),
- time: { start: 1, end: 2 },
- }),
- )
- await waitFor(() =>
- ui.commits.filter((commit) => commit.kind === "assistant" && commit.text === "Hello").length === 2
- ? true
- : undefined,
- )
- expect(
- ui.commits.filter((commit) => commit.kind === "assistant" && commit.text).map((commit) => commit.text),
- ).toEqual(["Hello", "Hello"])
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("preserves the display prefix for active reasoning restored during replay", async () => {
- const src = eventFeed()
- const ui = footer()
- let calls = 0
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async () => {
- calls += 1
- if (calls === 1) {
- return ok([])
- }
- return ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-thinking",
- parts: [reasoningPart("thinking-1", "msg-thinking", "")],
- }),
- ])
- },
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- src.push(assistant("msg-thinking"))
- src.push(reasoningUpdated(reasoningPart("thinking-1", "msg-thinking", "")))
- src.push(textDelta("msg-thinking", "thinking-1", "plan"))
- await waitFor(() => ui.commits.find((commit) => commit.kind === "reasoning" && commit.text === "Thinking: plan"))
- ui.commits.length = 0
- expect(await transport.replayOnResize({ localRows: () => [], reset: () => Promise.resolve() })).toBe(true)
- expect(ui.commits.filter((commit) => commit.kind === "reasoning").map((commit) => commit.text)).toEqual([
- "Thinking: plan",
- ])
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("does not overlay stale active text when persistence completes during replay", async () => {
- const src = eventFeed()
- const ui = footer()
- let calls = 0
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async () => {
- calls += 1
- if (calls === 1) {
- return ok([])
- }
- return ok([
- assistantMessage({
- sessionID: "session-1",
- id: "msg-finished",
- parts: [
- {
- ...textPart("text-finished", "msg-finished", "Hello"),
- time: { start: 1, end: 2 },
- },
- ],
- }),
- ])
- },
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- try {
- src.push(assistant("msg-finished"))
- src.push(textUpdated(textPart("text-finished", "msg-finished", "")))
- src.push(textDelta("msg-finished", "text-finished", "Hello"))
- await waitFor(() => ui.commits.find((commit) => commit.kind === "assistant" && commit.text === "Hello"))
- ui.commits.length = 0
- expect(await transport.replayOnResize({ localRows: () => [], reset: () => Promise.resolve() })).toBe(true)
- expect(
- ui.commits.filter((commit) => commit.kind === "assistant" && commit.text).map((commit) => commit.text),
- ).toEqual(["Hello"])
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("does not clear the terminal when resize replay snapshot fetch fails", async () => {
- const src = eventFeed()
- const ui = footer()
- let calls = 0
- const transport = await createSessionTransport({
- sdk: sdk({
- stream: src.stream,
- messages: async () => {
- calls += 1
- if (calls === 1) {
- return ok([])
- }
- throw new Error("snapshot failed")
- },
- }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- const reset = mock(() => Promise.resolve())
- try {
- expect(await transport.replayOnResize({ localRows: () => [], reset })).toBe(false)
- expect(reset).not.toHaveBeenCalled()
- expect(ui.commits).toEqual([])
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("disables resize replay for the session after terminal reset fails", async () => {
- const src = eventFeed()
- const ui = footer()
- const transport = await createSessionTransport({
- sdk: sdk({ stream: src.stream }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- const reset = mock(() => Promise.reject(new Error("clear failed")))
- try {
- expect(await transport.replayOnResize({ localRows: () => [], reset })).toBe(false)
- expect(await transport.replayOnResize({ localRows: () => [], reset })).toBe(false)
- expect(reset).toHaveBeenCalledTimes(1)
- expect(ui.commits).toContainEqual({
- kind: "error",
- text: "resize replay failed; disabled for this session",
- phase: "start",
- source: "system",
- })
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("disables resize replay when rebuilding scrollback fails after terminal reset", async () => {
- const src = eventFeed()
- const ui = footer()
- let cleared = false
- const idle = ui.api.idle
- ui.api.idle = () => (cleared ? Promise.reject(new Error("render failed")) : idle())
- const transport = await createSessionTransport({
- sdk: sdk({ stream: src.stream }),
- sessionID: "session-1",
- thinking: true,
- replay: true,
- limits: () => ({}),
- footer: ui.api,
- })
- const reset = mock(() => {
- cleared = true
- return Promise.resolve()
- })
- try {
- expect(await transport.replayOnResize({ localRows: () => [], reset })).toBe(false)
- expect(await transport.replayOnResize({ localRows: () => [], reset })).toBe(false)
- expect(reset).toHaveBeenCalledTimes(1)
- expect(ui.commits).toContainEqual({
- kind: "error",
- text: "resize replay failed; disabled for this session",
- phase: "start",
- source: "system",
- })
- } finally {
- src.close()
- await transport.close()
- }
- })
- test("keeps 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.objectContaining({ sessionID: "child-1", status: "completed" })])
- 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()
- }
- })
- })
|