mini: settle on wait and durable pending (#37984)
This commit is contained in:
parent
fd97d789ef
commit
592ef7433a
16 changed files with 908 additions and 657 deletions
|
|
@ -2,7 +2,6 @@ import type { FooterApi, FooterEvent, RunPrompt, StreamCommit } from "../../../s
|
|||
|
||||
export function createFooterApiFixture(input: { events?: FooterEvent[]; commits?: StreamCommit[] } = {}) {
|
||||
const prompts = new Set<(input: RunPrompt) => void>()
|
||||
const queuedRemoves = new Set<(messageID: string) => boolean | Promise<boolean>>()
|
||||
const closes = new Set<() => void>()
|
||||
const events = input.events ?? []
|
||||
const commits = input.commits ?? []
|
||||
|
|
@ -17,10 +16,6 @@ export function createFooterApiFixture(input: { events?: FooterEvent[]; commits?
|
|||
prompts.add(fn)
|
||||
return () => prompts.delete(fn)
|
||||
},
|
||||
onQueuedRemove(fn) {
|
||||
queuedRemoves.add(fn)
|
||||
return () => queuedRemoves.delete(fn)
|
||||
},
|
||||
onClose(fn) {
|
||||
if (closed) {
|
||||
fn()
|
||||
|
|
@ -46,7 +41,6 @@ export function createFooterApiFixture(input: { events?: FooterEvent[]; commits?
|
|||
destroy() {
|
||||
api.close()
|
||||
prompts.clear()
|
||||
queuedRemoves.clear()
|
||||
closes.clear()
|
||||
},
|
||||
}
|
||||
|
|
@ -60,8 +54,5 @@ export function createFooterApiFixture(input: { events?: FooterEvent[]; commits?
|
|||
const prompt: RunPrompt = mode ? { text, parts: [], mode } : { text, parts: [] }
|
||||
for (const fn of [...prompts]) fn(prompt)
|
||||
},
|
||||
removeQueued(messageID: string) {
|
||||
for (const fn of [...queuedRemoves]) void fn(messageID)
|
||||
},
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -12,7 +12,6 @@ test("down opens subagents from an empty prompt", async () => {
|
|||
const [state] = createSignal<FooterState>({
|
||||
phase: "idle",
|
||||
status: "",
|
||||
queue: 0,
|
||||
model: "gpt-5",
|
||||
usage: "",
|
||||
first: false,
|
||||
|
|
@ -70,7 +69,6 @@ test("down opens subagents from an empty prompt", async () => {
|
|||
onRows={() => {}}
|
||||
onLayout={() => {}}
|
||||
onStatus={() => {}}
|
||||
onQueuedRemove={async () => true}
|
||||
/>
|
||||
</Keymap.Provider>
|
||||
)
|
||||
|
|
|
|||
|
|
@ -92,7 +92,6 @@ function footerState(input: Partial<FooterState> = {}) {
|
|||
return createSignal<FooterState>({
|
||||
phase: "idle",
|
||||
status: "",
|
||||
queue: 0,
|
||||
model: "gpt-5",
|
||||
usage: "",
|
||||
first: false,
|
||||
|
|
@ -158,7 +157,6 @@ async function renderFooter(
|
|||
onRows={() => {}}
|
||||
onLayout={() => {}}
|
||||
onStatus={() => {}}
|
||||
onQueuedRemove={async () => true}
|
||||
/>
|
||||
</Keymap.Provider>
|
||||
)
|
||||
|
|
@ -680,10 +678,15 @@ test("direct subagent panel closes when moving up from the first item", async ()
|
|||
}
|
||||
})
|
||||
|
||||
test("direct queued prompt panel renders pending prompt actions", async () => {
|
||||
const [prompts] = createSignal([{ messageID: "m-1", prompt: { text: "fix the auth test", parts: [] } }])
|
||||
const edited: string[] = []
|
||||
const deleted: string[] = []
|
||||
test("direct pending panel shows durable delivery without edit actions", async () => {
|
||||
const [prompts] = createSignal([
|
||||
{
|
||||
messageID: "m-1",
|
||||
prompt: { text: "fix the auth test", parts: [] },
|
||||
delivery: "queue" as const,
|
||||
admittedSeq: 1,
|
||||
},
|
||||
])
|
||||
|
||||
const app = await testRender(
|
||||
() => (
|
||||
|
|
@ -692,12 +695,6 @@ test("direct queued prompt panel renders pending prompt actions", async () => {
|
|||
theme={() => RUN_THEME_FALLBACK.footer}
|
||||
prompts={prompts}
|
||||
onClose={() => {}}
|
||||
onEdit={(prompt) => {
|
||||
edited.push(prompt.messageID)
|
||||
}}
|
||||
onDelete={(prompt) => {
|
||||
deleted.push(prompt.messageID)
|
||||
}}
|
||||
/>
|
||||
</box>
|
||||
),
|
||||
|
|
@ -709,16 +706,14 @@ test("direct queued prompt panel renders pending prompt actions", async () => {
|
|||
const frame = app.captureCharFrame()
|
||||
const list = panelMenu(app.renderer.root)
|
||||
|
||||
expect(frame).toContain("Queued prompts")
|
||||
expect(frame).toContain("Pending work")
|
||||
expect(frame).toContain("fix the auth test")
|
||||
expect(frame).toContain("queued")
|
||||
expect(frame).toContain("queue")
|
||||
expect(frame).not.toContain("┌")
|
||||
expect(frame).not.toContain("┃")
|
||||
expectPaletteList(list, 0)
|
||||
app.mockInput.pressKey("e", { ctrl: true })
|
||||
app.mockInput.pressKey("DELETE")
|
||||
expect(edited).toEqual(["m-1"])
|
||||
expect(deleted).toEqual(["m-1"])
|
||||
expect(frame).not.toContain("edit")
|
||||
expect(frame).not.toContain("remove")
|
||||
} finally {
|
||||
app.renderer.destroy()
|
||||
}
|
||||
|
|
@ -999,11 +994,10 @@ test.skip("direct footer clears the synthetic skills draft when the panel closes
|
|||
}
|
||||
})
|
||||
|
||||
test("direct footer shows editable prompts and additional queued work while running", async () => {
|
||||
test("direct footer shows authoritative pending work while running", async () => {
|
||||
const [state] = createSignal<FooterState>({
|
||||
phase: "running",
|
||||
status: "",
|
||||
queue: 3,
|
||||
model: "gpt-5",
|
||||
usage: "",
|
||||
first: false,
|
||||
|
|
@ -1036,7 +1030,14 @@ test("direct footer shows editable prompts and additional queued work while runn
|
|||
state={state}
|
||||
view={view}
|
||||
subagent={subagents}
|
||||
queuedPrompts={() => [{ messageID: "m-queued", prompt: { text: "follow up", parts: [] } }]}
|
||||
queuedPrompts={() => [
|
||||
{
|
||||
messageID: "m-queued",
|
||||
prompt: { text: "follow up", parts: [] },
|
||||
delivery: "queue",
|
||||
admittedSeq: 1,
|
||||
},
|
||||
]}
|
||||
theme={() => RUN_THEME_FALLBACK}
|
||||
tuiConfig={tuiConfig}
|
||||
onSubmit={() => true}
|
||||
|
|
@ -1053,7 +1054,6 @@ test("direct footer shows editable prompts and additional queued work while runn
|
|||
onRows={() => {}}
|
||||
onLayout={() => {}}
|
||||
onStatus={() => {}}
|
||||
onQueuedRemove={async () => true}
|
||||
/>
|
||||
</Keymap.Provider>
|
||||
)
|
||||
|
|
@ -1088,9 +1088,9 @@ test("direct footer shows editable prompts and additional queued work while runn
|
|||
|
||||
expect(spinner).toBeDefined()
|
||||
expect(frame).toContain("a-model-name-long-enough-to-force-responsive-truncation")
|
||||
expect(frame).toContain("3 queued")
|
||||
expect(frame).toContain("1 pending")
|
||||
expect(frame).toContain("ctrl+b background")
|
||||
expect(frame).toContain("ctrl+x q 3 queued")
|
||||
expect(frame).toContain("ctrl+x q 1 pending")
|
||||
expect(frame).toContain("↓ subagents")
|
||||
expect(frame).toContain("ctrl+p cmd")
|
||||
expect(frame).toContain("a-model-name-long-enough-to-force-responsive-truncation")
|
||||
|
|
|
|||
|
|
@ -1,8 +1,16 @@
|
|||
import { describe, expect, test } from "bun:test"
|
||||
import { runPromptQueue } from "../../src/mini/runtime.queue"
|
||||
import { runPromptQueue as runPromptQueueBase, type QueueInput } from "../../src/mini/runtime.queue"
|
||||
import type { RunPrompt } from "../../src/mini/types"
|
||||
import { createFooterApiFixture } from "./fixture/footer-api"
|
||||
|
||||
function runPromptQueue(input: Omit<QueueInput, "admit" | "settle"> & Partial<Pick<QueueInput, "admit" | "settle">>) {
|
||||
return runPromptQueueBase({
|
||||
admit: async () => {},
|
||||
settle: async () => {},
|
||||
...input,
|
||||
})
|
||||
}
|
||||
|
||||
describe("run runtime queue", () => {
|
||||
test("ignores empty prompts", async () => {
|
||||
const ui = createFooterApiFixture()
|
||||
|
|
@ -171,206 +179,90 @@ describe("run runtime queue", () => {
|
|||
])
|
||||
})
|
||||
|
||||
test("passes prompts to onSend", async () => {
|
||||
test("durably admits in-flight follow-ups in submission order", async () => {
|
||||
const ui = createFooterApiFixture()
|
||||
const seen: string[] = []
|
||||
|
||||
await runPromptQueue({
|
||||
footer: ui.api,
|
||||
initialInput: " hello ",
|
||||
onSend: (input) => {
|
||||
seen.push(input.text)
|
||||
},
|
||||
run: async () => {
|
||||
ui.api.close()
|
||||
},
|
||||
})
|
||||
|
||||
expect(seen).toEqual([" hello "])
|
||||
})
|
||||
|
||||
test("appends the user row before the turn starts", async () => {
|
||||
const ui = createFooterApiFixture()
|
||||
|
||||
await runPromptQueue({
|
||||
footer: ui.api,
|
||||
initialInput: "/fmt bash",
|
||||
run: async () => {
|
||||
expect(ui.commits).toEqual([
|
||||
{
|
||||
kind: "user",
|
||||
text: "/fmt bash",
|
||||
phase: "start",
|
||||
source: "system",
|
||||
messageID: expect.any(String),
|
||||
},
|
||||
])
|
||||
ui.api.close()
|
||||
},
|
||||
})
|
||||
})
|
||||
|
||||
test("runs queued prompts in order", async () => {
|
||||
const ui = createFooterApiFixture()
|
||||
const seen: string[] = []
|
||||
let wake: (() => void) | undefined
|
||||
const gate = new Promise<void>((resolve) => {
|
||||
wake = resolve
|
||||
})
|
||||
const admitted: string[] = []
|
||||
const gate = Promise.withResolvers<void>()
|
||||
|
||||
const task = runPromptQueue({
|
||||
footer: ui.api,
|
||||
run: async (input) => {
|
||||
seen.push(input.text)
|
||||
if (seen.length === 1) {
|
||||
await gate
|
||||
return
|
||||
}
|
||||
|
||||
ui.api.close()
|
||||
run: async (input, _signal, onAdmitted) => {
|
||||
admitted.push(`${input.text}:steer`)
|
||||
onAdmitted()
|
||||
await gate.promise
|
||||
},
|
||||
admit: async (input) => {
|
||||
admitted.push(`${input.text}:queue`)
|
||||
},
|
||||
settle: async () => ui.api.close(),
|
||||
})
|
||||
|
||||
ui.submit("one")
|
||||
ui.submit("two")
|
||||
await Promise.resolve()
|
||||
expect(seen).toEqual(["one"])
|
||||
|
||||
wake?.()
|
||||
await task
|
||||
|
||||
expect(seen).toEqual(["one", "two"])
|
||||
})
|
||||
|
||||
test("exposes ordinary in-flight prompts for removal before sending", async () => {
|
||||
const ui = createFooterApiFixture()
|
||||
const turns: RunPrompt[] = []
|
||||
let wake: (() => void) | undefined
|
||||
const gate = new Promise<void>((resolve) => {
|
||||
wake = resolve
|
||||
})
|
||||
|
||||
const task = runPromptQueue({
|
||||
footer: ui.api,
|
||||
run: async (input) => {
|
||||
turns.push(input)
|
||||
await gate
|
||||
},
|
||||
})
|
||||
|
||||
ui.submit("one")
|
||||
ui.submit("two")
|
||||
await Promise.resolve()
|
||||
await Promise.resolve()
|
||||
|
||||
expect(turns.map((item) => item.text)).toEqual(["one"])
|
||||
expect(turns[0]?.messageID).toEqual(expect.any(String))
|
||||
ui.submit("three")
|
||||
while (admitted.length < 3) await Bun.sleep(0)
|
||||
expect(admitted).toEqual(["one:steer", "two:queue", "three:queue"])
|
||||
expect(ui.commits.map((item) => item.text)).toEqual(["one"])
|
||||
const first = ui.events.find((item) => item.type === "queued.prompts")
|
||||
const event = ui.events.findLast((item) => item.type === "queued.prompts")
|
||||
expect(first?.type === "queued.prompts" ? first.prompts : []).toEqual([])
|
||||
expect(
|
||||
first?.type === "queued.prompts" && event?.type === "queued.prompts" ? first.prompts === event.prompts : true,
|
||||
).toBe(false)
|
||||
expect(ui.events.findLast((item) => item.type === "queue")).toEqual({ type: "queue", queue: 1 })
|
||||
expect(event?.type === "queued.prompts" ? event.prompts.map((item) => item.prompt.text) : []).toEqual(["two"])
|
||||
if (event?.type === "queued.prompts") ui.removeQueued(event.prompts[0]!.messageID)
|
||||
await Promise.resolve()
|
||||
|
||||
wake?.()
|
||||
ui.api.close()
|
||||
gate.resolve()
|
||||
await task
|
||||
expect(turns.map((item) => item.text)).toEqual(["one"])
|
||||
})
|
||||
|
||||
test("removing one managed queued prompt preserves the others", async () => {
|
||||
test("continues durable admission after one fails", async () => {
|
||||
const ui = createFooterApiFixture()
|
||||
const turns: string[] = []
|
||||
let wake: (() => void) | undefined
|
||||
const gate = new Promise<void>((resolve) => {
|
||||
wake = resolve
|
||||
})
|
||||
|
||||
const admitted: string[] = []
|
||||
const errors: string[] = []
|
||||
const gate = Promise.withResolvers<void>()
|
||||
const task = runPromptQueue({
|
||||
footer: ui.api,
|
||||
run: async (input) => {
|
||||
turns.push(input.text)
|
||||
if (input.text === "active") await gate
|
||||
if (input.text === "queued three") ui.api.close()
|
||||
run: async (_input, _signal, admitted) => {
|
||||
admitted()
|
||||
await gate.promise
|
||||
},
|
||||
})
|
||||
|
||||
ui.submit("active")
|
||||
ui.submit("queued one")
|
||||
ui.submit("queued two")
|
||||
ui.submit("queued three")
|
||||
await Promise.resolve()
|
||||
await Promise.resolve()
|
||||
|
||||
const event = ui.events.findLast((item) => item.type === "queued.prompts")
|
||||
if (event?.type === "queued.prompts") {
|
||||
const second = event.prompts.find((item) => item.prompt.text === "queued two")
|
||||
if (second) ui.removeQueued(second.messageID)
|
||||
}
|
||||
|
||||
wake?.()
|
||||
await task
|
||||
expect(turns).toEqual(["active", "queued one", "queued three"])
|
||||
})
|
||||
|
||||
test("drains a prompt queued during an in-flight turn", async () => {
|
||||
const ui = createFooterApiFixture()
|
||||
const seen: string[] = []
|
||||
let wake: (() => void) | undefined
|
||||
const gate = new Promise<void>((resolve) => {
|
||||
wake = resolve
|
||||
})
|
||||
|
||||
const task = runPromptQueue({
|
||||
footer: ui.api,
|
||||
run: async (input) => {
|
||||
seen.push(input.text)
|
||||
if (seen.length === 1) {
|
||||
await gate
|
||||
return
|
||||
}
|
||||
|
||||
ui.api.close()
|
||||
admit: async (input) => {
|
||||
if (input.text === "two") throw new Error("admission failed")
|
||||
admitted.push(input.text)
|
||||
},
|
||||
onAdmissionError: (_prompt, error) => {
|
||||
errors.push(error instanceof Error ? error.message : String(error))
|
||||
},
|
||||
settle: async () => ui.api.close(),
|
||||
})
|
||||
|
||||
ui.submit("one")
|
||||
await Promise.resolve()
|
||||
expect(seen).toEqual(["one"])
|
||||
|
||||
wake?.()
|
||||
await Promise.resolve()
|
||||
ui.submit("two")
|
||||
ui.submit("three")
|
||||
while (admitted.length === 0) await Bun.sleep(0)
|
||||
gate.resolve()
|
||||
await task
|
||||
|
||||
expect(seen).toEqual(["one", "two"])
|
||||
expect(errors).toEqual(["admission failed"])
|
||||
expect(admitted).toEqual(["three"])
|
||||
})
|
||||
|
||||
test("close aborts the active run and drops pending queued work", async () => {
|
||||
test("close aborts an in-flight durable admission", async () => {
|
||||
const ui = createFooterApiFixture()
|
||||
const seen: string[] = []
|
||||
let hit = false
|
||||
let admissionHit = false
|
||||
const admissionStarted = Promise.withResolvers<void>()
|
||||
|
||||
const task = runPromptQueue({
|
||||
footer: ui.api,
|
||||
run: async (input, signal) => {
|
||||
seen.push(input.text)
|
||||
run: async (_input, signal, admitted) => {
|
||||
admitted()
|
||||
await new Promise<void>((resolve) => signal.addEventListener("abort", () => resolve(), { once: true }))
|
||||
},
|
||||
admit: async (_prompt, signal) => {
|
||||
admissionStarted.resolve()
|
||||
await new Promise<void>((resolve) => {
|
||||
if (signal.aborted) {
|
||||
hit = true
|
||||
admissionHit = true
|
||||
resolve()
|
||||
return
|
||||
}
|
||||
|
||||
signal.addEventListener(
|
||||
"abort",
|
||||
() => {
|
||||
hit = true
|
||||
admissionHit = true
|
||||
resolve()
|
||||
},
|
||||
{ once: true },
|
||||
|
|
@ -382,11 +274,11 @@ describe("run runtime queue", () => {
|
|||
ui.submit("one")
|
||||
await Promise.resolve()
|
||||
ui.submit("two")
|
||||
await admissionStarted.promise
|
||||
ui.api.close()
|
||||
await task
|
||||
|
||||
expect(hit).toBe(true)
|
||||
expect(seen).toEqual(["one"])
|
||||
expect(admissionHit).toBe(true)
|
||||
})
|
||||
|
||||
test("propagates run errors", async () => {
|
||||
|
|
|
|||
|
|
@ -93,6 +93,8 @@ describe("run interactive runtime", () => {
|
|||
streamStarted.resolve()
|
||||
return {
|
||||
runPromptTurn: async () => {},
|
||||
queuePromptTurn: async () => {},
|
||||
waitForIdle: async () => {},
|
||||
interruptActiveTurn: async () => {},
|
||||
selectSubagent: () => {},
|
||||
settleForm: (sessionID: string, formID: string) => settled.push({ sessionID, formID }),
|
||||
|
|
@ -432,6 +434,8 @@ describe("run interactive runtime", () => {
|
|||
setTimeout(() => input.footer.close(), 0)
|
||||
return {
|
||||
runPromptTurn: async () => {},
|
||||
queuePromptTurn: async () => {},
|
||||
waitForIdle: async () => {},
|
||||
interruptActiveTurn: async () => {},
|
||||
selectSubagent: () => {},
|
||||
replayOnResize: async () => false,
|
||||
|
|
|
|||
|
|
@ -127,6 +127,8 @@ function sdk(input: {
|
|||
globals?: FormInfo[]
|
||||
globalLocation?: { directory: string; workspaceID?: string }
|
||||
permissions?: Record<string, PermissionV2Request[]>
|
||||
pending?: Record<string, Awaited<ReturnType<OpenCodeClient["session"]["pending"]["list"]>>>
|
||||
wait?: () => Promise<void>
|
||||
}) {
|
||||
const client = OpenCode.make({ baseUrl: "https://opencode.test" })
|
||||
let subscription = 0
|
||||
|
|
@ -159,6 +161,8 @@ function sdk(input: {
|
|||
}),
|
||||
)
|
||||
spyOn(client.session, "active").mockImplementation(() => ok(input.active?.() ?? {}))
|
||||
spyOn(client.session.pending, "list").mockImplementation((request) => ok(input.pending?.[request.sessionID] ?? []))
|
||||
spyOn(client.session, "wait").mockImplementation(() => input.wait?.() ?? ok(undefined))
|
||||
spyOn(client.session, "message").mockImplementation((request) => {
|
||||
const message = input.messages?.[request.sessionID]?.find((item) => item.id === request.messageID)
|
||||
return message ? (ok(message) as never) : Promise.reject(new Error(`message not found: ${request.messageID}`))
|
||||
|
|
@ -417,35 +421,23 @@ describe("V2 mini transport", () => {
|
|||
await transport.close()
|
||||
})
|
||||
|
||||
test("finalizes an idle projection before reducing live output", async () => {
|
||||
test("waits authoritatively and reconciles the projected terminal suffix", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
const settled = defer()
|
||||
const messages: SessionMessages = []
|
||||
const client = sdk({
|
||||
streams: [events],
|
||||
messages: {
|
||||
ses_1: [
|
||||
{
|
||||
id: "msg_old",
|
||||
type: "assistant",
|
||||
agent: "build",
|
||||
model: { providerID: "test", id: "model" },
|
||||
content: [{ type: "text", text: "[Link](https://example.com)" }],
|
||||
time: { created: 1 },
|
||||
},
|
||||
],
|
||||
},
|
||||
messages: { ses_1: messages },
|
||||
wait: () => settled.promise,
|
||||
})
|
||||
const ui = footer()
|
||||
const idle = spyOn(ui.api, "idle")
|
||||
const transport = await createSessionTransport({
|
||||
sdk: client,
|
||||
sessionID: "ses_1",
|
||||
thinking: true,
|
||||
replay: true,
|
||||
thinking: false,
|
||||
footer: ui.api,
|
||||
})
|
||||
expect(ui.commits.map((item) => item.text)).toEqual(["[Link](https://example.com)"])
|
||||
expect(idle).toHaveBeenCalledTimes(1)
|
||||
|
||||
let admitted = false
|
||||
spyOn(client.session, "prompt").mockImplementation((request) => {
|
||||
|
|
@ -480,7 +472,7 @@ describe("V2 mini transport", () => {
|
|||
sessionID: "ses_1",
|
||||
assistantMessageID: "msg_assistant",
|
||||
ordinal: 0,
|
||||
delta: "answer",
|
||||
delta: "ans",
|
||||
},
|
||||
})
|
||||
events.push({
|
||||
|
|
@ -490,10 +482,222 @@ describe("V2 mini transport", () => {
|
|||
durable: durable("ses_1"),
|
||||
data: { sessionID: "ses_1" },
|
||||
})
|
||||
let done = false
|
||||
void turn.then(() => {
|
||||
done = true
|
||||
})
|
||||
await Bun.sleep(0)
|
||||
expect(done).toBe(false)
|
||||
messages.push(
|
||||
{ id: "msg_prompt", type: "user", text: "hello", time: { created: 2 } },
|
||||
{
|
||||
id: "msg_assistant",
|
||||
type: "assistant",
|
||||
agent: "build",
|
||||
model: { providerID: "test", id: "model" },
|
||||
content: [{ type: "text", text: "answer" }],
|
||||
time: { created: 3, completed: 4 },
|
||||
},
|
||||
)
|
||||
settled.resolve()
|
||||
await turn
|
||||
|
||||
expect(ui.commits.map((item) => item.text)).toEqual(["[Link](https://example.com)", "answer"])
|
||||
expect(ui.events).toContainEqual({ type: "stream.patch", patch: { phase: "idle", status: "" } })
|
||||
expect(ui.commits.map((item) => item.text)).toEqual(["ans", "wer"])
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
test("shows durable pending delivery and appends queued input on promotion", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
const client = sdk({
|
||||
streams: [events],
|
||||
pending: {
|
||||
ses_1: [
|
||||
{
|
||||
admittedSeq: 1,
|
||||
id: "msg_queued",
|
||||
sessionID: "ses_1",
|
||||
timeCreated: 1,
|
||||
type: "user",
|
||||
data: { text: "follow up" },
|
||||
delivery: "queue",
|
||||
},
|
||||
],
|
||||
},
|
||||
})
|
||||
const ui = footer()
|
||||
const transport = await createSessionTransport({
|
||||
sdk: client,
|
||||
sessionID: "ses_1",
|
||||
thinking: false,
|
||||
footer: ui.api,
|
||||
})
|
||||
const pending = () =>
|
||||
ui.events
|
||||
.findLast((item) => item.type === "queued.prompts")
|
||||
?.prompts.map((item) => [item.messageID, item.delivery, item.admittedSeq])
|
||||
|
||||
expect(pending()).toEqual([["msg_queued", "queue", 1]])
|
||||
events.push({
|
||||
id: "evt_promoted",
|
||||
created: 2,
|
||||
type: "session.input.promoted",
|
||||
durable: durable("ses_1", 2),
|
||||
data: { sessionID: "ses_1", inputID: "msg_queued" },
|
||||
})
|
||||
while (!ui.commits.some((item) => item.messageID === "msg_queued")) await Bun.sleep(0)
|
||||
|
||||
expect(ui.commits).toContainEqual(
|
||||
expect.objectContaining({ kind: "user", messageID: "msg_queued", text: "follow up" }),
|
||||
)
|
||||
expect(pending()).toEqual([])
|
||||
const prompt = spyOn(client.session, "prompt").mockImplementation((request) =>
|
||||
ok({ ...promptAdmission(request), admittedSeq: 2 }) as never,
|
||||
)
|
||||
await transport.queuePromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: { messageID: "msg_next", text: "another", parts: [] },
|
||||
files: [],
|
||||
includeFiles: false,
|
||||
})
|
||||
expect(prompt).toHaveBeenCalledWith(expect.objectContaining({ delivery: "queue" }), expect.anything())
|
||||
events.push({
|
||||
id: "evt_earlier_admission",
|
||||
created: 3,
|
||||
type: "session.input.admitted",
|
||||
durable: durable("ses_1", 1),
|
||||
data: {
|
||||
sessionID: "ses_1",
|
||||
inputID: "msg_earlier",
|
||||
input: { type: "user", data: { text: "earlier" }, delivery: "steer" },
|
||||
},
|
||||
})
|
||||
while (true) {
|
||||
const pending = ui.events.findLast((item) => item.type === "queued.prompts")
|
||||
if (pending?.type === "queued.prompts" && pending.prompts.length >= 2) break
|
||||
await Bun.sleep(0)
|
||||
}
|
||||
expect(pending()).toEqual([
|
||||
["msg_earlier", "steer", 1],
|
||||
["msg_next", "queue", 2],
|
||||
])
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
test("reports an observed execution failure before prompt promotion", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
const idle = defer()
|
||||
const client = sdk({ streams: [events], messages: { ses_1: [] }, wait: () => idle.promise })
|
||||
const ui = footer()
|
||||
const transport = await createSessionTransport({
|
||||
sdk: client,
|
||||
sessionID: "ses_1",
|
||||
thinking: false,
|
||||
footer: ui.api,
|
||||
})
|
||||
let admitted = false
|
||||
spyOn(client.session, "prompt").mockImplementation((request) => {
|
||||
admitted = true
|
||||
return ok(promptAdmission(request)) as never
|
||||
})
|
||||
|
||||
const turn = transport.runPromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
|
||||
files: [],
|
||||
includeFiles: false,
|
||||
})
|
||||
while (!admitted) await Bun.sleep(0)
|
||||
events.push({
|
||||
id: "evt_failed",
|
||||
created: 2,
|
||||
type: "session.execution.failed",
|
||||
durable: durable("ses_1", 2),
|
||||
data: { sessionID: "ses_1", error: { type: "unknown", message: "instructions unavailable" } },
|
||||
})
|
||||
await Bun.sleep(0)
|
||||
idle.resolve()
|
||||
|
||||
await turn
|
||||
expect(ui.commits).toContainEqual(
|
||||
expect.objectContaining({ kind: "error", messageID: "msg_prompt", text: "instructions unavailable" }),
|
||||
)
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
test("attributes an execution-only failure to the latest promoted prompt", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
const idle = defer()
|
||||
const messages: SessionMessages = []
|
||||
const client = sdk({ streams: [events], messages: { ses_1: messages }, wait: () => idle.promise })
|
||||
const ui = footer()
|
||||
const transport = await createSessionTransport({
|
||||
sdk: client,
|
||||
sessionID: "ses_1",
|
||||
thinking: false,
|
||||
footer: ui.api,
|
||||
})
|
||||
let admitted = false
|
||||
spyOn(client.session, "prompt").mockImplementation((request) => {
|
||||
admitted = true
|
||||
return ok(promptAdmission(request)) as never
|
||||
})
|
||||
|
||||
const turn = transport.runPromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
|
||||
files: [],
|
||||
includeFiles: false,
|
||||
})
|
||||
while (!admitted) await Bun.sleep(0)
|
||||
events.push({
|
||||
id: "evt_prompt_promoted",
|
||||
created: 2,
|
||||
type: "session.input.promoted",
|
||||
durable: durable("ses_1", 2),
|
||||
data: { sessionID: "ses_1", inputID: "msg_prompt" },
|
||||
})
|
||||
await transport.queuePromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: { messageID: "msg_queued", text: "follow up", parts: [] },
|
||||
files: [],
|
||||
includeFiles: false,
|
||||
})
|
||||
events.push({
|
||||
id: "evt_queued_promoted",
|
||||
created: 3,
|
||||
type: "session.input.promoted",
|
||||
durable: durable("ses_1", 3),
|
||||
data: { sessionID: "ses_1", inputID: "msg_queued" },
|
||||
})
|
||||
events.push({
|
||||
id: "evt_failed",
|
||||
created: 4,
|
||||
type: "session.execution.failed",
|
||||
durable: durable("ses_1", 4),
|
||||
data: { sessionID: "ses_1", error: { type: "unknown", message: "model unavailable" } },
|
||||
})
|
||||
await Bun.sleep(0)
|
||||
messages.push(
|
||||
{ id: "msg_prompt", type: "user", text: "hello", time: { created: 2 } },
|
||||
{ id: "msg_queued", type: "user", text: "follow up", time: { created: 3 } },
|
||||
)
|
||||
idle.resolve()
|
||||
|
||||
await turn
|
||||
expect(ui.commits).toContainEqual(
|
||||
expect.objectContaining({ kind: "error", messageID: "msg_queued", text: "model unavailable" }),
|
||||
)
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
|
|
@ -786,11 +990,12 @@ describe("V2 mini transport", () => {
|
|||
await transport.close()
|
||||
})
|
||||
|
||||
test("rebootstraps after disconnect and completes a promoted turn from idle active state", async () => {
|
||||
test("reconnects and hydrates without completing before session.wait", async () => {
|
||||
const first = feed()
|
||||
const second = feed()
|
||||
first.push(connected("evt_connected_1"))
|
||||
second.push(connected("evt_connected_2"))
|
||||
const idle = defer()
|
||||
let running = true
|
||||
const client = sdk({
|
||||
streams: [first, second],
|
||||
|
|
@ -799,6 +1004,7 @@ describe("V2 mini transport", () => {
|
|||
if (running) active.ses_1 = { type: "running" }
|
||||
return active
|
||||
},
|
||||
wait: () => idle.promise,
|
||||
})
|
||||
let projected = false
|
||||
spyOn(client.message, "list").mockImplementation(() =>
|
||||
|
|
@ -844,11 +1050,25 @@ describe("V2 mini transport", () => {
|
|||
while (!admitted) await Bun.sleep(0)
|
||||
projected = true
|
||||
running = false
|
||||
second.push({
|
||||
id: "evt_prior_failed",
|
||||
created: 1,
|
||||
type: "session.execution.failed",
|
||||
durable: durable("ses_1", 1),
|
||||
data: { sessionID: "ses_1", error: { type: "unknown", message: "prior execution failed" } },
|
||||
})
|
||||
second.push({
|
||||
id: "evt_prompted",
|
||||
created: 2,
|
||||
type: "session.input.promoted",
|
||||
durable: durable("ses_1", 2),
|
||||
data: { sessionID: "ses_1", inputID: "msg_prompt" },
|
||||
})
|
||||
first.close()
|
||||
while (!ui.events.some((event) => event.type === "stream.patch" && event.patch.status === "reconnecting"))
|
||||
await Bun.sleep(0)
|
||||
idle.resolve()
|
||||
await turn
|
||||
|
||||
expect(ui.events).toContainEqual({ type: "stream.patch", patch: { phase: "running", status: "reconnecting" } })
|
||||
expect(ui.events).toContainEqual({ type: "stream.patch", patch: { phase: "idle", status: "" } })
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
|
|
@ -1789,52 +2009,6 @@ describe("V2 mini transport", () => {
|
|||
await transport.close()
|
||||
})
|
||||
|
||||
test("resolves an interrupted turn even when promotion never arrived", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
const client = sdk({
|
||||
streams: [events],
|
||||
active: () => ({ ses_1: { type: "running" } }),
|
||||
})
|
||||
const ui = footer()
|
||||
const transport = await createSessionTransport({
|
||||
sdk: client,
|
||||
sessionID: "ses_1",
|
||||
thinking: false,
|
||||
footer: ui.api,
|
||||
})
|
||||
let admitted = false
|
||||
// The generated method has conditional return types for throwOnError; this mock represents the successful branch.
|
||||
// @ts-expect-error successful SDK response is valid for both modes at runtime
|
||||
spyOn(client.session, "prompt").mockImplementation((request) => {
|
||||
admitted = true
|
||||
return ok({ data: promptAdmission(request) })
|
||||
})
|
||||
const interrupted = spyOn(client.session, "interrupt").mockImplementation(() => ok(undefined))
|
||||
|
||||
const turn = transport.runPromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: { messageID: "msg_prompt", text: "hello", parts: [] },
|
||||
files: [],
|
||||
includeFiles: true,
|
||||
})
|
||||
while (!admitted) await Bun.sleep(0)
|
||||
await transport.interruptActiveTurn()
|
||||
events.push({
|
||||
id: "evt_settled",
|
||||
created: 0,
|
||||
type: "session.execution.interrupted",
|
||||
durable: durable("ses_1"),
|
||||
data: { sessionID: "ses_1", reason: "user" },
|
||||
})
|
||||
await turn
|
||||
|
||||
expect(interrupted).toHaveBeenCalledWith({ sessionID: "ses_1" })
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
test("falls back to the default model when selecting a variant on a fresh session", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
|
|
@ -1901,7 +2075,8 @@ describe("V2 mini transport", () => {
|
|||
test("interrupts the current Session when an active turn is aborted", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
const client = sdk({ streams: [events] })
|
||||
const idle = defer()
|
||||
const client = sdk({ streams: [events], wait: () => idle.promise })
|
||||
const ui = footer()
|
||||
const transport = await createSessionTransport({
|
||||
sdk: client,
|
||||
|
|
@ -1940,13 +2115,7 @@ describe("V2 mini transport", () => {
|
|||
})
|
||||
await Bun.sleep(0)
|
||||
controller.abort()
|
||||
events.push({
|
||||
id: "evt_settled",
|
||||
created: 0,
|
||||
type: "session.execution.interrupted",
|
||||
durable: durable("ses_1"),
|
||||
data: { sessionID: "ses_1", reason: "user" },
|
||||
})
|
||||
idle.resolve()
|
||||
await turn
|
||||
|
||||
expect(interrupted).toHaveBeenCalledWith({ sessionID: "ses_1" })
|
||||
|
|
@ -2493,90 +2662,6 @@ describe("V2 mini transport", () => {
|
|||
await transport.close()
|
||||
})
|
||||
|
||||
test("does not resolve a skill turn before the matching activation is observed", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
const client = sdk({ streams: [events] })
|
||||
const ui = footer()
|
||||
const transport = await createSessionTransport({
|
||||
sdk: client,
|
||||
sessionID: "ses_1",
|
||||
thinking: false,
|
||||
footer: ui.api,
|
||||
})
|
||||
let sent = false
|
||||
spyOn(client.session, "skill").mockImplementation(() => {
|
||||
sent = true
|
||||
return ok(undefined) as never
|
||||
})
|
||||
|
||||
let done = false
|
||||
const turn = transport
|
||||
.runPromptTurn({
|
||||
agent: undefined,
|
||||
model: undefined,
|
||||
variant: undefined,
|
||||
prompt: {
|
||||
messageID: "msg_skill",
|
||||
text: "/tigerstyle",
|
||||
parts: [],
|
||||
command: { name: "tigerstyle", arguments: "", source: "skill" },
|
||||
},
|
||||
files: [],
|
||||
includeFiles: true,
|
||||
})
|
||||
.then(() => {
|
||||
done = true
|
||||
})
|
||||
while (!sent) await Bun.sleep(0)
|
||||
events.push({
|
||||
id: "evt_other",
|
||||
created: 0,
|
||||
type: "session.skill.activated",
|
||||
durable: durable("ses_1"),
|
||||
data: {
|
||||
sessionID: "ses_1",
|
||||
id: "other",
|
||||
name: "other",
|
||||
text: "other instructions",
|
||||
},
|
||||
})
|
||||
events.push({
|
||||
id: "evt_unrelated_settled",
|
||||
created: 0,
|
||||
type: "session.execution.succeeded",
|
||||
durable: durable("ses_1"),
|
||||
data: { sessionID: "ses_1" },
|
||||
})
|
||||
await Bun.sleep(0)
|
||||
await Bun.sleep(0)
|
||||
expect(done).toBe(false)
|
||||
|
||||
events.push({
|
||||
id: "evt_skill",
|
||||
created: 0,
|
||||
type: "session.skill.activated",
|
||||
durable: durable("ses_1"),
|
||||
data: {
|
||||
sessionID: "ses_1",
|
||||
id: "tigerstyle",
|
||||
name: "tigerstyle",
|
||||
text: "skill instructions",
|
||||
},
|
||||
})
|
||||
events.push({
|
||||
id: "evt_skill_settled",
|
||||
created: 0,
|
||||
type: "session.execution.succeeded",
|
||||
durable: durable("ses_1"),
|
||||
data: { sessionID: "ses_1" },
|
||||
})
|
||||
await turn
|
||||
|
||||
expect(done).toBe(true)
|
||||
await transport.close()
|
||||
})
|
||||
|
||||
test("refreshes catalogs on connection and location-scoped invalidations", async () => {
|
||||
const events = feed()
|
||||
events.push(connected())
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue