Compare commits

...

3 Commits

Author SHA1 Message Date
Kit Langton 12c7dc1bfd Merge branch 'v2' into fix-stale-tools 2026-07-05 19:48:02 -04:00
Kit Langton df56253799 fix(tui): reconcile messages by event watermark 2026-07-05 13:45:33 -04:00
Kit Langton 87eb062238 fix(tui): clear stale tool preparation state 2026-07-05 13:24:32 -04:00
2 changed files with 327 additions and 11 deletions
+27 -10
View File
@@ -100,6 +100,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
directory: process.cwd(),
})
const messageIndex = new Map<string, Map<string, number>>()
const appliedEventSeq = new Map<string, number>()
let connectionGeneration = 0
let statusChanges: Set<string> | undefined
let bootstrapping: Promise<void> | undefined
@@ -208,6 +209,11 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
}
function handleEvent(event: V2Event) {
if ("durable" in event && event.durable)
appliedEventSeq.set(
event.durable.aggregateID,
Math.max(appliedEventSeq.get(event.durable.aggregateID) ?? -1, event.durable.seq),
)
switch (event.type) {
case "session.created":
void result.session.refresh(event.data.sessionID)
@@ -382,7 +388,9 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
break
case "session.text.started":
message.update(event.data.sessionID, (draft, index) => {
message.assistant(draft, index, event.data.assistantMessageID)?.content.push({
const assistant = message.assistant(draft, index, event.data.assistantMessageID)
if (!assistant || message.latestText(assistant, event.data.textID)) return
assistant.content.push({
type: "text",
id: event.data.textID,
text: "",
@@ -409,7 +417,9 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
break
case "session.tool.input.started":
message.update(event.data.sessionID, (draft, index) => {
message.assistant(draft, index, event.data.assistantMessageID)?.content.push({
const assistant = message.assistant(draft, index, event.data.assistantMessageID)
if (!assistant || message.latestTool(assistant, event.data.callID)) return
assistant.content.push({
type: "tool",
id: event.data.callID,
name: event.data.name,
@@ -442,7 +452,7 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
message.assistant(draft, index, event.data.assistantMessageID),
event.data.callID,
)
if (!match) return
if (match?.state.status !== "pending") return
match.time.ran = event.created
match.provider = event.data.provider
match.state = { status: "running", input: event.data.input, structured: {}, content: [] }
@@ -506,7 +516,9 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
break
case "session.reasoning.started":
message.update(event.data.sessionID, (draft, index) => {
message.assistant(draft, index, event.data.assistantMessageID)?.content.push({
const assistant = message.assistant(draft, index, event.data.assistantMessageID)
if (!assistant || message.latestReasoning(assistant, event.data.reasoningID)) return
assistant.content.push({
type: "reasoning",
id: event.data.reasoningID,
text: "",
@@ -693,21 +705,24 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
return position === undefined ? undefined : messages?.[position]
},
async refresh(sessionID: string) {
const response = await sdk.api.message.list({ sessionID, limit: 200, order: "desc" })
const loaded = mutable(response.data).toReversed()
const live = [...(store.session.message[sessionID] ?? [])]
setStore("session", "message", sessionID, [])
messageIndex.set(sessionID, new Map())
const loaded = mutable(
(await sdk.api.message.list({ sessionID, limit: 200, order: "desc" })).data,
).toReversed()
const loadedIDs = new Set(loaded.map((message) => message.id))
const liveByID = new Map(live.map((message) => [message.id, message]))
const snapshotIsCurrent =
response.watermark !== undefined && response.watermark >= (appliedEventSeq.get(sessionID) ?? -1)
const messages = [
...loaded.map((message) => {
if (message.type === "user") return message
return liveByID.get(message.id) ?? message
if (snapshotIsCurrent) return message
const live = liveByID.get(message.id)
if (!live || ("completed" in message.time && message.time.completed !== undefined)) return message
return live
}),
...live.filter((message) => !loadedIDs.has(message.id)),
].toSorted((a, b) => a.time.created - b.time.created)
if (snapshotIsCurrent) appliedEventSeq.set(sessionID, response.watermark)
messageIndex.set(sessionID, new Map(messages.map((message, index) => [message.id, index])))
setStore("session", "message", sessionID, messages)
},
@@ -925,6 +940,8 @@ export const { use: useData, provider: DataProvider } = createSimpleContext({
if (details.type === "server.connected") {
refreshActive()
void bootstrap()
for (const sessionID of Object.keys(store.session.message))
void result.session.message.refresh(sessionID).catch(() => undefined)
return
}
handleEvent(details)
+300 -1
View File
@@ -108,9 +108,241 @@ test("refreshes resources into reactive getters", async () => {
}
})
test("refresh reconciles messages by durable snapshot watermark", async () => {
const events = createEventStream()
const sessionID = "session-stale-tool"
const messageID = "message-stale-tool"
const liveMessageID = "message-live-tool"
let requests = 0
const calls = createFetch((url) => {
if (url.pathname === `/api/session/${sessionID}/message` && ++requests === 1)
return json({
data: [
{
id: messageID,
type: "assistant",
agent: "build",
model: { id: "model", providerID: "provider" },
content: [
{
type: "tool",
id: "call-write",
name: "write",
provider: { executed: false },
state: {
status: "error",
input: {},
content: [],
structured: {},
error: { type: "unknown", message: "Tool execution interrupted" },
},
time: { created: 1, completed: 3 },
},
],
finish: "error",
error: { type: "unknown", message: "Step interrupted" },
time: { created: 1, completed: 3 },
},
],
watermark: 2,
cursor: {},
})
if (url.pathname === `/api/session/${sessionID}/message`)
return json({
data: [
{
id: liveMessageID,
type: "assistant",
agent: "build",
model: { id: "model", providerID: "provider" },
content: [
{
type: "tool",
id: "call-read",
name: "read",
state: { status: "pending", input: "" },
time: { created: 4 },
},
],
time: { created: 3 },
},
],
watermark: 4,
cursor: {},
})
return undefined
}, events)
let data!: ReturnType<typeof useData>
function Probe() {
data = useData()
return <box />
}
const app = await testRender(() => (
<TestTuiContexts>
<SDKProvider client={createClient(calls.fetch)} api={createApi(calls.fetch)}>
<ProjectProvider>
<DataProvider>
<Probe />
</DataProvider>
</ProjectProvider>
</SDKProvider>
</TestTuiContexts>
))
try {
emitEvent(events, {
id: "evt_step_started_stale_tool",
created: 1,
type: "session.step.started",
durable: durable(sessionID),
data: {
sessionID,
assistantMessageID: messageID,
agent: "build",
model: { id: "model", providerID: "provider" },
},
})
emitEvent(events, {
id: "evt_tool_started_stale_tool",
created: 2,
type: "session.tool.input.started",
durable: durable(sessionID, 1),
data: {
sessionID,
assistantMessageID: messageID,
callID: "call-write",
name: "write",
},
})
emitEvent(events, {
id: "evt_step_failed_stale_tool",
created: 3,
type: "session.step.failed",
durable: durable(sessionID, 2),
data: {
sessionID,
assistantMessageID: messageID,
error: { type: "unknown", message: "Step interrupted" },
},
})
emitEvent(events, {
id: "evt_prompt_after_stale_tool",
created: 4,
type: "session.prompt.admitted",
durable: durable(sessionID, 3),
data: {
sessionID,
inputID: "message-after-stale-tool",
prompt: { text: "Continue" },
delivery: "steer",
},
})
await wait(() => {
const message = data.session.message.get(sessionID, messageID)
return message?.type === "assistant" && message.time.completed === 3
})
await data.session.message.refresh(sessionID)
const message = data.session.message.get(sessionID, messageID)
expect(message?.type).toBe("assistant")
if (message?.type !== "assistant") return
expect(message.content[0]).toMatchObject({ type: "tool", state: { status: "error" } })
emitEvent(events, {
id: "evt_tool_started_stale_replay",
created: 2,
type: "session.tool.input.started",
durable: durable(sessionID, 1),
data: {
sessionID,
assistantMessageID: messageID,
callID: "call-write",
name: "write",
},
})
emitEvent(events, {
id: "evt_tool_called_stale_replay",
created: 2,
type: "session.tool.called",
durable: durable(sessionID, 2),
data: {
sessionID,
assistantMessageID: messageID,
callID: "call-write",
tool: "write",
input: { path: "README.md" },
provider: { executed: false },
},
})
await Bun.sleep(10)
const replayed = data.session.message.get(sessionID, messageID)
expect(replayed?.type).toBe("assistant")
if (replayed?.type !== "assistant") return
expect(replayed.content).toHaveLength(1)
expect(replayed.content[0]).toMatchObject({ type: "tool", state: { status: "error" } })
emitEvent(events, {
id: "evt_step_started_live_tool",
created: 5,
type: "session.step.started",
durable: durable(sessionID, 4),
data: {
sessionID,
assistantMessageID: liveMessageID,
agent: "build",
model: { id: "model", providerID: "provider" },
},
})
emitEvent(events, {
id: "evt_tool_started_live_tool",
created: 6,
type: "session.tool.input.started",
durable: durable(sessionID, 5),
data: {
sessionID,
assistantMessageID: liveMessageID,
callID: "call-read",
name: "read",
},
})
emitEvent(events, {
id: "evt_tool_called_live_tool",
created: 7,
type: "session.tool.called",
durable: durable(sessionID, 6),
data: {
sessionID,
assistantMessageID: liveMessageID,
callID: "call-read",
tool: "read",
input: { path: "README.md" },
provider: { executed: false },
},
})
await wait(() => {
const message = data.session.message.get(sessionID, liveMessageID)
return (
message?.type === "assistant" &&
message.content[0]?.type === "tool" &&
message.content[0].state.status === "running"
)
})
await data.session.message.refresh(sessionID)
const live = data.session.message.get(sessionID, liveMessageID)
expect(live?.type).toBe("assistant")
if (live?.type !== "assistant") return
expect(live.content[0]).toMatchObject({ type: "tool", state: { status: "running" } })
} finally {
app.renderer.destroy()
}
})
test("reconnects the event stream and bootstraps fresh data", async () => {
const events = createEventStream()
const requests = { active: 0, event: 0, model: 0 }
const requests = { active: 0, event: 0, message: 0, model: 0 }
let resolveActive!: (response: Response) => void
const calls = createFetch((url) => {
if (url.pathname === "/api/event") {
@@ -124,6 +356,40 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
resolveActive = resolve
})
}
if (url.pathname === "/api/session/session-reconnect/message") {
requests.message++
return json({
data: [
{
id: "message-reconnect",
type: "assistant",
agent: "build",
model: { id: "model", providerID: "provider" },
content: [
{
type: "tool",
id: "call-reconnect",
name: "write",
provider: { executed: false },
state: {
status: "error",
input: {},
content: [],
structured: {},
error: { type: "unknown", message: "Tool execution interrupted" },
},
time: { created: 1, completed: 2 },
},
],
finish: "error",
error: { type: "unknown", message: "Step interrupted" },
time: { created: 0, completed: 2 },
},
],
watermark: 2,
cursor: {},
})
}
if (url.pathname !== "/api/model") return
requests.model++
return json({
@@ -169,6 +435,30 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
await wait(() => data.session.status("session-stale") === "running")
expect(data.connection.status()).toBe("connected")
expect(data.connection.attempt()).toBe(0)
emitEvent(events, {
id: "evt_step_started_before_reconnect",
created: 0,
type: "session.step.started",
durable: durable("session-reconnect"),
data: {
sessionID: "session-reconnect",
assistantMessageID: "message-reconnect",
agent: "build",
model: { id: "model", providerID: "provider" },
},
})
emitEvent(events, {
id: "evt_tool_started_before_reconnect",
created: 1,
type: "session.tool.input.started",
durable: durable("session-reconnect", 1),
data: {
sessionID: "session-reconnect",
assistantMessageID: "message-reconnect",
callID: "call-reconnect",
name: "write",
},
})
events.disconnect()
await wait(() => data.connection.status() === "connecting")
@@ -193,8 +483,17 @@ test("reconnects the event stream and bootstraps fresh data", async () => {
await wait(() => data.location.model.list()?.[0]?.id === "model-2", 4000)
await wait(() => data.session.status("session-stale") === "idle")
await wait(() => {
const message = data.session.message.get("session-reconnect", "message-reconnect")
return (
message?.type === "assistant" &&
message.content[0]?.type === "tool" &&
message.content[0].state.status === "error"
)
})
expect(data.session.status("session-new")).toBe("running")
expect(requests.event).toBe(2)
expect(requests.message).toBe(1)
expect(data.connection.status()).toBe("connected")
expect(data.connection.attempt()).toBe(0)
expect(data.connection.error()).toBeUndefined()