Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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
Expand Down
8 changes: 8 additions & 0 deletions history.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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). */
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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" ? "*" : " "
Expand Down
8 changes: 8 additions & 0 deletions test/session.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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`)
35 changes: 34 additions & 1 deletion test/universal.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand Down Expand Up @@ -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`)
105 changes: 95 additions & 10 deletions tui.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -311,7 +311,9 @@ export default Plugin.define({
const turns = new Map<string, Turn>()
// 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<string, number>()
// 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<string, number>()
// 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
Expand Down Expand Up @@ -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<void> {
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<string>()
/**
* 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<void> {
// `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.
Expand All @@ -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}` : ""}`)
Expand Down Expand Up @@ -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) {
Expand All @@ -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")
}

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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)}`)
}
Expand Down
39 changes: 35 additions & 4 deletions universal.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -160,7 +186,12 @@ export function lastModel(
export function aggregateTurn(
steps: readonly SessionMessageAssistant[],
marks: ReadonlyMap<string, Turn>,
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]
Expand Down Expand Up @@ -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,
Expand Down
Loading