diff --git a/CHANGELOG.md b/CHANGELOG.md index 7ebae41..172c441 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,13 @@ +## [Unreleased] +### Fixed +- A reply is recorded when it ends, not when OpenCode's execution does. A + message sent while a reply is running is queued into the same execution, + so the first reply was never shown: a 19m 53s reply left the previous + turn in the sidebar. +- A reply that was interrupted or failed is shown, marked `interrupted` or + `failed`, with OpenCode's figures up to that point, instead of leaving the + previous turn in place. The history panel marks it too. + ## [0.3.1] – 2026-09-24 ### Changed - The first turn after OpenCode starts can show engine figures on diff --git a/history.ts b/history.ts index 46d9f72..ce9bd80 100644 --- a/history.ts +++ b/history.ts @@ -76,6 +76,8 @@ export interface TurnRecord { * parallel, so not a sum). Rates are never combined across them. */ subagents?: { count: number; tokens: number; spanS: number; cost?: number; steps?: number } + /** Set when the reply did not finish: stopped by the user, or failed. */ + outcome?: "interrupted" | "failed" engine?: { prefillTokS?: number /** Tokens committed per verify pass (MTPLX's multi-token prediction). */ @@ -122,6 +124,11 @@ export interface Summary { /** A turn's streaming time: recorded, or derived from an older row's rate. */ export function streamOf(t: TurnRecord): number | undefined { + // A reply that did not finish has no tokens recorded for its last step, so + // its tokens over its streaming time would understate the speed (measured: + // an interrupted turn, 0 tokens over 6s of streaming, halved a session's + // average). It counts toward nothing that divides by streaming time. + if (t.outcome) return undefined if (t.streamS !== undefined && t.streamS > 0) return t.streamS // Rows recorded before streamS existed carry a generation rate whose // window is tokens / rate. A whole-turn rate is not generation, so no. @@ -180,6 +187,7 @@ export function formatRow(t: TurnRecord, modelWidth = 18): string { t.ttft !== undefined ? `ttft ${nn(t.ttft, 2)}s` : "", money(t.cost), t.cached !== undefined && t.cached > 0 ? `${ni(t.cached)} cached` : "", + t.outcome ?? "", ].filter(Boolean) // A leading marker rather than a column, so a row is readable at any width. const mark = t.source === "engine" ? "*" : " " diff --git a/test/session.test.mjs b/test/session.test.mjs index 512ec2a..a0f4d18 100644 --- a/test/session.test.mjs +++ b/test/session.test.mjs @@ -249,4 +249,12 @@ test("the roll-up carries its sub-agents' steps -- one engine request each", () assert.equal(rollupSubagents([child({ steps: undefined })], ["ses_child"], 0, 40_000).steps, 1) }) +test("an unfinished reply counts toward no speed or time split", () => { + // OpenCode records no tokens for an interrupted step: 0 tokens over 6s. + const s = summariseSession([row({ outcome: "interrupted", tokens: 0, streamS: 6 }), row({ tokens: 100, streamS: 2 })], SID) + assert.equal(s.genTokS, 50) + assert.equal(s.trend.length, 1) + assert.equal(s.turns, 2) +}) + console.log(`\n${passed} passed`) diff --git a/test/universal.test.mjs b/test/universal.test.mjs index c4bca30..e49aaaa 100644 --- a/test/universal.test.mjs +++ b/test/universal.test.mjs @@ -8,7 +8,7 @@ // counts never disagreed — only the rate's numerator did. // Run with: bun test/universal.test.mjs import { strict as assert } from "node:assert" -import { turnRate, universalLine, universalView, turnSteps, lastModel, aggregateTurn, DEFAULT_DISPLAY } from "../universal.ts" +import { turnRate, universalLine, universalView, turnSteps, turnUserAt, lastModel, aggregateTurn, DEFAULT_DISPLAY } from "../universal.ts" let passed = 0 function test(name, fn) { @@ -379,4 +379,37 @@ test("a session with no model named yet gives none, not a guess", () => { assert.equal(lastModel([]), undefined) }) +// ---- a reply's boundaries, when one execution holds several ----------------- +// A message sent while a reply runs is queued into the same execution +// (measured: a 9-step reply ended `stop`, and the queued message's step began +// 2s later, in one execution that ended interrupted). Each reply is a turn. +const queued = [ + { type: "user", id: "u1", time: { created: 1_000 } }, + { type: "assistant", id: "a1", time: { created: 2_000 } }, + { type: "assistant", id: "a2", time: { created: 3_000 } }, + { type: "user", id: "u2", time: { created: 2_500 } }, + { type: "assistant", id: "b1", time: { created: 9_000 } }, +] + +test("a reply ending at a given step excludes the queued message after it", () => { + assert.deepEqual(turnSteps(queued, "a2").map((m) => m.id), ["a1", "a2"]) + assert.equal(turnUserAt(queued, "a2"), 1_000) +}) + +test("without an end step, the turn is the latest reply", () => { + assert.deepEqual(turnSteps(queued).map((m) => m.id), ["b1"]) + assert.equal(turnUserAt(queued), 2_500) +}) + +test("an unknown end step falls back to the latest reply", () => { + assert.deepEqual(turnSteps(queued, "gone").map((m) => m.id), ["b1"]) +}) + +test("an interrupted reply's total runs to when it stopped", () => { + const partial = [{ type: "assistant", id: "x", time: { created: 5_000 }, tokens: { output: 10, reasoning: 0 } }] + const { info } = aggregateTurn(partial, new Map(), { execStart: 4_000, endAt: 34_000 }) + assert.equal(info.time.completed, 34_000) + assert.equal(info.time.created, 4_000) +}) + console.log(`\n${passed} passed`) diff --git a/tui.tsx b/tui.tsx index 50f9fa1..88f9595 100644 --- a/tui.tsx +++ b/tui.tsx @@ -29,7 +29,7 @@ import { appendFileSync } from "node:fs" import { short } from "./format" import type { HttpOptions } from "./http" -import { universalView, turnRate, turnSteps, lastModel, aggregateTurn, type Turn, type Display, DEFAULT_DISPLAY } from "./universal" +import { universalView, turnRate, turnSteps, turnUserAt, lastModel, aggregateTurn, type Turn, type Display, DEFAULT_DISPLAY } from "./universal" import { record, historyLines, type History, type TurnRecord } from "./history" import { emptyPanels, lineFor, keyFor, setLine, LatestPerKey, PLACEHOLDER, type Panels } from "./panels" import { encodeView, decodeView, LABEL_WIDTH, type TurnView } from "./rows" @@ -311,7 +311,9 @@ export default Plugin.define({ const turns = new Map() // When each session's current execution started: the start of what the // user waits for, which the turn's total runs from (retries included). - const execStart = new Map() + // When each session's current reply started: the execution starting, + // then each reply ending, since one execution can hold several replies. + const replyStart = new Map() // Per step: its provider and model (from session.step.started), and for // engines that report their latest request -- MTPLX, KoboldCpp, mlx-serve // -- a read taken at that step's end. A tool-using turn is one request @@ -772,8 +774,19 @@ export default Plugin.define({ // a turn in one tab suppress another tab's line (see panels.ts). const latest = new LatestPerKey() - async function report(sessionID: string): Promise { - const seq = latest.begin(sessionID) + // Replies already reported, by their last step's id. A reply is reported + // when its last step ends, and again asked for when the execution ends; + // it must render once. + const reported = new Set() + /** + * One reply: the user message's steps, ending at `endID` (the step that + * ended the reply) or at the latest step. `outcome` is set when the + * execution was interrupted or failed before the reply finished. + */ + async function report( + sessionID: string, + opts: { endID?: string; outcome?: "interrupted" | "failed" } = {} + ): Promise { // `message.list()` is a union of message kinds and its tail after a turn // is an "idle" marker, not the reply — measured, see audit P1. Scan // backwards for the assistant message. @@ -782,12 +795,30 @@ export default Plugin.define({ // The turn is every assistant message since the last user message: one // per step when it calls tools. Reading only the last one showed // `140 tok 10.20s` for a 316-token, 37s turn (measured, vllm-mlx). - const steps = turnSteps(msgs) - const agg = aggregateTurn(steps, turns, { execStart: execStart.get(sessionID) }) - execStart.delete(sessionID) + const steps = turnSteps(msgs, opts.endID) + const lastID = steps[steps.length - 1]?.id + if (!lastID || reported.has(lastID)) return + reported.add(lastID) + if (reported.size > 64) { + const oldest = reported.values().next().value + if (oldest !== undefined) reported.delete(oldest) + } + const seq = latest.begin(sessionID) + // The reply started when the execution did -- or, for a message queued + // behind an earlier reply in the same execution, when that reply ended, + // or when the message was sent if that is later. + const userAt = turnUserAt(msgs, opts.endID) + const since = replyStart.get(sessionID) + const start = since !== undefined && userAt !== undefined ? Math.max(since, userAt) : (since ?? userAt) + const agg = aggregateTurn(steps, turns, { + execStart: start, + endAt: opts.outcome ? Date.now() : undefined, + }) + replyStart.set(sessionID, Date.now()) const info = agg.info const turn = agg.turn if (!info) return + dbg(`report: ${sessionID} ending ${lastID}${opts.outcome ? ` (${opts.outcome})` : ""}; ${steps.length} step(s)`) dbg( `turn: ${steps.length} assistant message(s) [${steps .map((m) => `${(m.tokens?.output ?? 0) + (m.tokens?.reasoning ?? 0)}${m.finish ? `/${m.finish}` : ""}`) @@ -860,7 +891,9 @@ export default Plugin.define({ } let line: TurnView | null = null try { - line = await enrich(provider, model, info, turn, http, tier2, steps, sameEngine) + // An unfinished reply's last step never completed, so the engine has + // no reading of it to check against: OpenCode's figures only. + if (!opts.outcome) line = await enrich(provider, model, info, turn, http, tier2, steps, sameEngine) // Recorded before the fallback overwrites it, so history knows which // tier the figures actually came from. } catch (e: unknown) { @@ -884,7 +917,15 @@ export default Plugin.define({ // above are measured and complete; only their SOURCE changes once a // baseline exists, and the rate in particular can move an order of // magnitude when it does. Split to fit the box's 32 cells. - if (tier2.pendingBaseline) line.notes.push("engine telemetry", "from the next turn") + if (opts.outcome) { + // OpenCode records no tokens for a step it stopped mid-stream + // (measured: an interrupted reply's step came back 0/error after 7s + // of thinking), so a 0 here is unknown, not none. + const out0 = (info.tokens?.output ?? 0) + (info.tokens?.reasoning ?? 0) + if (out0 === 0) line.rows = line.rows.filter(([label]) => label !== "tokens" && label !== "speed") + line.notes.push(opts.outcome) + } + else if (tier2.pendingBaseline) line.notes.push("engine telemetry", "from the next turn") else if (tier2.sharedWindow) line.notes.push("engine data skipped:", "overlapping requests") } @@ -919,6 +960,7 @@ export default Plugin.define({ steps: turn?.steps, engine: enriched ? tier2.engine : undefined, subagents, + outcome: opts.outcome, } setHistory((d) => { d.turns = record({ turns: d.turns }, rec).turns @@ -992,7 +1034,7 @@ export default Plugin.define({ ctx.data.on("session.execution.started", (evt) => { const sid = (evt as { data?: { sessionID?: string } }).data?.sessionID if (typeof sid !== "string") return - execStart.set(sid, Date.now()) + replyStart.set(sid, Date.now()) // The engine this turn will use: the session's last model, or one // just selected. A new session on the default model has neither, // and its first turn goes unprimed, as before. @@ -1131,6 +1173,49 @@ export default Plugin.define({ } }) ) + // A reply ends when a step ends with a final answer. That is + // its own turn even if the execution goes on: a message sent while it + // ran is queued into the same execution, which then only ends -- or is + // interrupted -- after the next reply (measured: a 19m 53s reply never + // reported, because its execution ended interrupted 20m in). + off.push( + ctx.data.on("session.step.ended", (evt) => { + const d = (evt as { data?: { sessionID?: string; assistantMessageID?: string; finish?: string } }).data + if (typeof d?.sessionID !== "string" || typeof d.assistantMessageID !== "string") return + // Not "error" or "unknown": OpenCode retries a step under the same + // message id, and a failed attempt must not end the reply before + // the retry does. The execution's own end catches those. + if (d.finish !== "stop" && d.finish !== "length" && d.finish !== "content-filter") return + const sid = d.sessionID + const id = d.assistantMessageID + // After the host has applied the step to the message it ends. + setTimeout(() => { + const m = ctx.data.session.message.get(sid, id) as SessionMessageAssistant | undefined + dbg(`reply ended ${id} (${d.finish}); completed ${m?.time?.completed !== undefined}, tokens ${m?.tokens?.output ?? "?"}`) + report(sid, { endID: id }).catch((e: unknown) => dbg(`report threw: ${String(e)}`)) + }, 50) + }) + ) + // An execution stopped before its reply finished still used the engine; + // the reply is shown, marked, rather than leaving the previous turn up. + off.push( + ctx.data.on("session.execution.interrupted", (evt) => { + const sid = (evt as { data?: { sessionID?: string } }).data?.sessionID + if (typeof sid === "string") { + dbg(`event execution.interrupted ${sid}`) + report(sid, { outcome: "interrupted" }).catch((e: unknown) => dbg(`report threw: ${String(e)}`)) + } + }) + ) + off.push( + ctx.data.on("session.execution.failed", (evt) => { + const sid = (evt as { data?: { sessionID?: string } }).data?.sessionID + if (typeof sid === "string") { + dbg(`event execution.failed ${sid}`) + report(sid, { outcome: "failed" }).catch((e: unknown) => dbg(`report threw: ${String(e)}`)) + } + }) + ) } catch (e: unknown) { dbg(`subscribe failed: ${String(e)}`) } diff --git a/universal.ts b/universal.ts index 5b80cf1..a72d8d4 100644 --- a/universal.ts +++ b/universal.ts @@ -112,10 +112,18 @@ export function turnRate( * Reading only the last one showed `140 tok 10.20s` for a 316-token, 37s turn. */ export function turnSteps( - msgs: readonly ({ type?: string } | undefined)[] + msgs: readonly ({ type?: string; id?: string } | undefined)[], + /** + * The reply's last step. Given, the turn ends there rather than at the end + * of the list: a message sent while a reply is running is queued into the + * same execution, so by the time the first reply ends the list can already + * hold the next user message after it (measured: a 9-step reply ending in + * `stop`, then a queued message's step 2s later, in one execution). + */ + endID?: string ): SessionMessageAssistant[] { const steps: SessionMessageAssistant[] = [] - for (let i = msgs.length - 1; i >= 0; i--) { + for (let i = lastIndex(msgs, endID); i >= 0; i--) { const m = msgs[i] if (!m) continue if (m.type === "user") break @@ -124,6 +132,24 @@ export function turnSteps( return steps } +/** When the user message that opened the turn ending at `endID` was sent. */ +export function turnUserAt( + msgs: readonly ({ type?: string; id?: string; time?: { created?: number } } | undefined)[], + endID?: string +): number | undefined { + for (let i = lastIndex(msgs, endID); i >= 0; i--) { + const m = msgs[i] + if (m?.type === "user") return m.time?.created + } + return undefined +} + +function lastIndex(msgs: readonly ({ id?: string } | undefined)[], endID?: string): number { + if (endID === undefined) return msgs.length - 1 + const i = msgs.findIndex((m) => m?.id === endID) + return i >= 0 ? i : msgs.length - 1 +} + /** * The model a session last used or selected, newest first: an assistant * message's model, or a `model-switched` entry. Undefined for a session with @@ -160,7 +186,12 @@ export function lastModel( export function aggregateTurn( steps: readonly SessionMessageAssistant[], marks: ReadonlyMap, - opts: { execStart?: number } = {} + /** + * `execStart`: when the reply started (the execution, or the previous reply + * in it ending). `endAt`: when an interrupted reply stopped, since its last + * step never completed. + */ + opts: { execStart?: number; endAt?: number } = {} ): { info: SessionMessageAssistant | undefined; turn: Turn | undefined } { const first = steps[0] const last = steps[steps.length - 1] @@ -198,7 +229,7 @@ export function aggregateTurn( const info: SessionMessageAssistant = { ...last, - time: { created: opts.execStart ?? first.time.created, completed: last.time?.completed }, + time: { created: opts.execStart ?? first.time.created, completed: last.time?.completed ?? opts.endAt }, tokens: { input: last.tokens?.input ?? 0, output,