diff --git a/packages/opencode/src/bench/cli.ts b/packages/opencode/src/bench/cli.ts index 6a566612a6be..c8bbb7b875ee 100644 --- a/packages/opencode/src/bench/cli.ts +++ b/packages/opencode/src/bench/cli.ts @@ -26,6 +26,12 @@ import os from "node:os" import { spawn } from "node:child_process" import { runDeepReset } from "./deep_reset" import { bootstrapRepoIfMissing } from "./bootstrap_repo" +import { + collectCompletionMetrics, + parseToolExecutionMetric, + updateNemoGymMetrics, + type ActionExecutionLatencyMetric, +} from "./metrics" import { capturePatch, ensureCommitIdentity, parsePatchMode, recordBaselineCommit, type PatchMode } from "./patch" import * as BenchTerminalError from "./terminal_error" // opencode's built-in anthropic system prompt — Bun bundles .txt as a string. @@ -371,7 +377,13 @@ function runOpencode(args: { env: NodeJS.ProcessEnv opencodeBin: string agent: string -}): Promise<{ exitCode: number; stdout: string; stderr: string; terminalError?: BenchTerminalError.Kind }> { +}): Promise<{ + exitCode: number + stdout: string + stderr: string + actionExecutionLatencies: ActionExecutionLatencyMetric[] + terminalError?: BenchTerminalError.Kind +}> { // Use the same bun binary that's currently running — guaranteed to exist // and avoids PATH lookup quirks under Bun's posix_spawn. const bunPath = process.execPath @@ -411,15 +423,23 @@ function runOpencode(args: { } const MAX_KEEP = 256 * 1024 // keep only a bounded tail for error reporting let lineBuf = "" + const actionExecutionLatencies = new Map() child.stdout?.on("data", (b) => { const chunk = b.toString("utf8") observeTerminalSignal(chunk) lineBuf += chunk let idx: number while ((idx = lineBuf.indexOf("\n")) >= 0) { - const line = lineBuf.slice(0, idx) + const rawLine = lineBuf.slice(0, idx) lineBuf = lineBuf.slice(idx + 1) - // Forward the event type so the gym log captures progress cheaply. + const actionMetric = parseToolExecutionMetric(rawLine) + if (actionMetric) { + const metricID = `${actionMetric.session_id}:${actionMetric.observation_id}` + actionExecutionLatencies.set(metricID, actionMetric) + } + // Keep full tool details in metrics, but expose only the event type to + // the outer Gym log so large inputs and outputs are not duplicated. + const line = actionMetric ? "tool_use" : rawLine process.stdout.write(line + "\n") stdout = (stdout + line + "\n").slice(-MAX_KEEP) } @@ -430,10 +450,13 @@ function runOpencode(args: { stderr = (stderr + chunk).slice(-MAX_KEEP) process.stderr.write(chunk) }) - child.on("close", (code) => resolve({ exitCode: code ?? 0, stdout, stderr, terminalError })) + const metrics = () => [...actionExecutionLatencies.values()].sort((a, b) => a.timestamp.localeCompare(b.timestamp)) + child.on("close", (code) => + resolve({ exitCode: code ?? 0, stdout, stderr, actionExecutionLatencies: metrics(), terminalError }), + ) child.on("error", (err) => { stderr += String(err) - resolve({ exitCode: 999, stdout, stderr, terminalError }) + resolve({ exitCode: 999, stdout, stderr, actionExecutionLatencies: metrics(), terminalError }) }) }) } @@ -478,6 +501,7 @@ function detectOpencodeBin(): string { } async function main() { + const initializeStartedAt = Date.now() const args = parseArgs(process.argv.slice(2)) const instance = await readInstance(args.instanceDictPath, args.selectedId) // workspaceRoot is decided gym-side based on dataset_name; we use it verbatim. @@ -539,7 +563,6 @@ async function main() { replayManifest, }) - const startedAt = Date.now() const childEnv: NodeJS.ProcessEnv = { ...process.env, // Run-isolated opencode state. @@ -588,6 +611,8 @@ async function main() { } const opencodeBin = detectOpencodeBin() + const initializeRuntimeTime = (Date.now() - initializeStartedAt) / 1000 + const startedAt = Date.now() const result = await runOpencode({ workspaceRoot, modelName, @@ -599,6 +624,15 @@ async function main() { const patch = await capturePatch(workspaceRoot, args.patchMode, baselineCommit) const benchRunTime = (Date.now() - startedAt) / 1000 + const completionMetrics = await collectCompletionMetrics(completionsDir) + const perTurnMetrics = { + response_latencies: completionMetrics.responseLatencies, + action_execution_latencies: result.actionExecutionLatencies, + token_usages: completionMetrics.tokenUsages, + } + await updateNemoGymMetrics(process.env.NEMO_GYM_METRICS_FPATH, { + initialize_runtime_time: initializeRuntimeTime, + }) const error = BenchTerminalError.toGymError(result.exitCode, result.terminalError) const outPath = await writeOutputJsonl(args.outputDir, instance.instance_id, { @@ -608,6 +642,7 @@ async function main() { metrics: { bench_run_time: benchRunTime, opencode_exit_code: result.exitCode, + ...perTurnMetrics, patch_mode: args.patchMode, }, error, diff --git a/packages/opencode/src/bench/metrics.ts b/packages/opencode/src/bench/metrics.ts new file mode 100644 index 000000000000..147191e616b5 --- /dev/null +++ b/packages/opencode/src/bench/metrics.ts @@ -0,0 +1,218 @@ +import { promises as fs } from "node:fs" +import path from "node:path" + +interface ResponseLatencyMetric { + latency: number + response_id: string + request_kind: "agent" | "title" | "subagent" + session_id: string + parent_session_id: string | null + session_turn: number + start_timestamp: string + timestamp: string +} + +export interface ActionExecutionLatencyMetric { + observation_type: string + observation_id: string + session_id: string + child_session_id?: string + input?: Record + output?: string + latency: number + message: string + start_timestamp: string + timestamp: string +} + +interface TokenUsageMetric { + prompt_tokens: number + completion_tokens: number + reasoning_tokens?: number + response_id: string +} + +interface CompletionMetrics { + responseLatencies: ResponseLatencyMetric[] + tokenUsages: TokenUsageMetric[] +} + +interface CompletionDump { + response?: { + id?: unknown + usage?: { + prompt_tokens?: unknown + completion_tokens?: unknown + completion_tokens_details?: { + reasoning_tokens?: unknown + } | null + } + } + latency?: unknown + request_kind?: unknown + request_started_at?: unknown + session_id?: unknown + parent_session_id?: unknown + turn?: unknown + timestamp?: unknown +} + +function finiteNumber(value: unknown): number | undefined { + return typeof value === "number" && Number.isFinite(value) ? value : undefined +} + +function nonNegativeInteger(value: unknown): number { + const number = finiteNumber(value) + return number === undefined ? 0 : Math.max(0, Math.trunc(number)) +} + +function optionalNonNegativeInteger(value: unknown): number | undefined { + const number = finiteNumber(value) + return number === undefined || number < 0 ? undefined : Math.trunc(number) +} + +function isoTimestamp(seconds: number): string { + return new Date(seconds * 1000).toISOString() +} + +export function parseToolExecutionMetric(line: string): ActionExecutionLatencyMetric | undefined { + let event: Record + try { + event = JSON.parse(line) + } catch { + return undefined + } + + const part = event.part + const state = part?.state + if ( + event.type !== "tool_use" || + part?.type !== "tool" || + (state?.status !== "completed" && state?.status !== "error") + ) { + return undefined + } + + const recordedStart = finiteNumber(state.time?.start) + const end = finiteNumber(state.time?.end) + const callID = typeof part.callID === "string" ? part.callID : typeof part.id === "string" ? part.id : undefined + if (recordedStart === undefined || end === undefined || end < recordedStart || !callID) return undefined + + const title = typeof state.title === "string" ? state.title : "" + const error = typeof state.error === "string" ? state.error : "" + const metadata = + "metadata" in state && state.metadata && typeof state.metadata === "object" + ? (state.metadata as Record) + : undefined + const childSessionID = + part.tool === "task" && typeof metadata?.sessionId === "string" ? metadata.sessionId : undefined + const input = + state.input && typeof state.input === "object" && !Array.isArray(state.input) + ? (state.input as Record) + : undefined + const output = typeof state.output === "string" ? state.output : undefined + return { + observation_type: typeof part.tool === "string" ? part.tool : "opencode_tool", + observation_id: callID, + session_id: typeof event.sessionID === "string" ? event.sessionID : "", + ...(childSessionID ? { child_session_id: childSessionID } : {}), + ...(input ? { input } : {}), + ...(output === undefined ? {} : { output }), + latency: (end - recordedStart) / 1000, + message: title || error, + start_timestamp: new Date(recordedStart).toISOString(), + timestamp: new Date(end).toISOString(), + } +} + +export async function collectCompletionMetrics(completionsDir: string): Promise { + const records: Array<{ + startedAtSeconds: number + timestampSeconds: number + responseLatency: ResponseLatencyMetric + tokenUsage: TokenUsageMetric + }> = [] + + for (const name of await fs.readdir(completionsDir)) { + if (!name.endsWith(".json")) continue + + let dump: CompletionDump + try { + dump = JSON.parse(await fs.readFile(path.join(completionsDir, name), "utf8")) + } catch { + continue + } + + const latency = finiteNumber(dump.latency) + const requestStartedAtSeconds = finiteNumber(dump.request_started_at) + const timestampSeconds = finiteNumber(dump.timestamp) + const responseID = typeof dump.response?.id === "string" ? dump.response.id : "" + if ( + latency === undefined || + latency < 0 || + requestStartedAtSeconds === undefined || + timestampSeconds === undefined || + timestampSeconds < requestStartedAtSeconds || + !responseID + ) + continue + + const requestKind = dump.request_kind === "title" || dump.request_kind === "subagent" ? dump.request_kind : "agent" + const sessionID = typeof dump.session_id === "string" ? dump.session_id : "" + const parentSessionID = typeof dump.parent_session_id === "string" ? dump.parent_session_id : null + const sessionTurn = nonNegativeInteger(dump.turn) + const promptTokens = nonNegativeInteger(dump.response?.usage?.prompt_tokens) + const completionTokens = nonNegativeInteger(dump.response?.usage?.completion_tokens) + const reasoningTokens = optionalNonNegativeInteger( + dump.response?.usage?.completion_tokens_details?.reasoning_tokens, + ) + records.push({ + startedAtSeconds: requestStartedAtSeconds, + timestampSeconds, + responseLatency: { + latency, + response_id: responseID, + request_kind: requestKind, + session_id: sessionID, + parent_session_id: parentSessionID, + session_turn: sessionTurn, + start_timestamp: isoTimestamp(requestStartedAtSeconds), + timestamp: isoTimestamp(timestampSeconds), + }, + tokenUsage: { + prompt_tokens: promptTokens, + completion_tokens: completionTokens, + ...(reasoningTokens === undefined ? {} : { reasoning_tokens: reasoningTokens }), + response_id: responseID, + }, + }) + } + + // Subagent requests can overlap, so completion order is not turn-start order. + records.sort( + (a, b) => + a.startedAtSeconds - b.startedAtSeconds || + a.timestampSeconds - b.timestampSeconds || + a.responseLatency.response_id.localeCompare(b.responseLatency.response_id), + ) + return { + responseLatencies: records.map((record) => record.responseLatency), + tokenUsages: records.map((record) => record.tokenUsage), + } +} + +export async function updateNemoGymMetrics( + metricsPath: string | undefined, + update: Record, +): Promise { + if (!metricsPath) return + + let existing: Record = {} + try { + existing = JSON.parse(await fs.readFile(metricsPath, "utf8")) + } catch {} + + const tmpPath = `${metricsPath}.tmp.${process.pid}.${Date.now()}` + await fs.writeFile(tmpPath, JSON.stringify({ ...existing, ...update })) + await fs.rename(tmpPath, metricsPath) +} diff --git a/packages/opencode/src/bench/replay.ts b/packages/opencode/src/bench/replay.ts index 45aab7fa8e82..68331d3d77d6 100644 --- a/packages/opencode/src/bench/replay.ts +++ b/packages/opencode/src/bench/replay.ts @@ -73,7 +73,11 @@ export function parseReplayMessages(raw: string): ParsedReplay { if (msg.role === "assistant") { replayTurns.push({ content: typeof msg.content === "string" ? msg.content : replayMessageText(msg.content) || null, - toolCalls: msg.tool_calls?.map((tc) => ({ id: tc.id, name: tc.function.name, arguments: tc.function.arguments })), + toolCalls: msg.tool_calls?.map((tc) => ({ + id: tc.id, + name: tc.function.name, + arguments: tc.function.arguments, + })), ...(pendingUserTexts.length ? { precedingUserTexts: pendingUserTexts } : {}), }) pendingUserTexts = [] diff --git a/packages/opencode/src/cli/cmd/run.ts b/packages/opencode/src/cli/cmd/run.ts index 4ab9f3be6868..9b47ca93ed48 100644 --- a/packages/opencode/src/cli/cmd/run.ts +++ b/packages/opencode/src/cli/cmd/run.ts @@ -434,7 +434,7 @@ export const RunCommand = effectCmd({ function emit(type: string, data: Record) { if (args.format === "json") { - if (benchEventTypesOnly) { + if (benchEventTypesOnly && type !== "tool_use") { process.stdout.write(type + EOL) return true } @@ -465,10 +465,14 @@ export const RunCommand = effectCmd({ if (event.type === "message.part.updated") { const part = event.properties.part + + if (part.type === "tool" && (part.state.status === "completed" || part.state.status === "error")) { + if (emit("tool_use", { sessionID: part.sessionID, part })) continue + } + if (part.sessionID !== sessionID) continue if (part.type === "tool" && (part.state.status === "completed" || part.state.status === "error")) { - if (emit("tool_use", { part })) continue if (part.state.status === "completed") { tool(part) continue diff --git a/packages/opencode/src/provider/sdk/nemo-gym/language-model.ts b/packages/opencode/src/provider/sdk/nemo-gym/language-model.ts index 07203dcb21fc..9d3dcc176dce 100644 --- a/packages/opencode/src/provider/sdk/nemo-gym/language-model.ts +++ b/packages/opencode/src/provider/sdk/nemo-gym/language-model.ts @@ -104,6 +104,9 @@ interface ChatResponseUsage { prompt_tokens?: number | null completion_tokens?: number | null total_tokens?: number | null + completion_tokens_details?: { + reasoning_tokens?: number | null + } | null } interface ChatResponse { @@ -276,6 +279,17 @@ export class NemoGymLanguageModel implements LanguageModelV3 { private globalTurn = 0 private readonly sessionStartGlobalTurn = new Map() + private _requestKind( + messages: ChatRequestMessage[], + parentSessionID: string | undefined, + ): "agent" | "title" | "subagent" { + if (parentSessionID) return "subagent" + const titlePrompt = "Generate a title for this conversation:" + return messages.some((message) => message.role === "user" && JSON.stringify(message.content).includes(titlePrompt)) + ? "title" + : "agent" + } + constructor(modelId: string, cfg: NemoGymLanguageModelConfig) { this.modelId = modelId this.provider = cfg.provider @@ -566,7 +580,9 @@ export class NemoGymLanguageModel implements LanguageModelV3 { const { warnings, loggedMessages, requestParams, globalTurn, sessionStartGlobalTurn } = await this._buildRequestParams(options, session) + const requestStartedAt = Date.now() const { responseJson } = await this._postChat(requestParams) + const responseCompletedAt = Date.now() const choice = responseJson.choices[0] if (!choice) throw new Error("nemo-gym: empty choices in response") @@ -602,6 +618,8 @@ export class NemoGymLanguageModel implements LanguageModelV3 { providerSpecificFields, requestParams, session, + requestStartedAt, + responseCompletedAt, globalTurn, sessionStartGlobalTurn, }) @@ -657,7 +675,9 @@ export class NemoGymLanguageModel implements LanguageModelV3 { controller.enqueue({ type: "stream-start", warnings }) try { + const requestStartedAt = Date.now() const { responseJson } = await self._postChat(requestParams) + const responseCompletedAt = Date.now() const choice = responseJson.choices[0] if (!choice) throw new Error("nemo-gym: empty choices in response") @@ -690,6 +710,8 @@ export class NemoGymLanguageModel implements LanguageModelV3 { providerSpecificFields, requestParams, session, + requestStartedAt, + responseCompletedAt, globalTurn, sessionStartGlobalTurn, }) @@ -1017,6 +1039,8 @@ export class NemoGymLanguageModel implements LanguageModelV3 { providerSpecificFields: Record requestParams: Record session: SessionHeaders + requestStartedAt: number + responseCompletedAt: number globalTurn: number sessionStartGlobalTurn: number }) { @@ -1027,7 +1051,13 @@ export class NemoGymLanguageModel implements LanguageModelV3 { (args.session.parentSessionID ? this.liveToRecordedSession.get(args.session.parentSessionID) : undefined) if (this.cfg.onCompletion) { try { - await this.cfg.onCompletion({ turn, ...args }) + await this.cfg.onCompletion({ + turn, + messages: args.messages, + response: args.response, + providerSpecificFields: args.providerSpecificFields, + requestParams: args.requestParams, + }) } catch (err) { console.warn(`[nemo-gym] onCompletion hook threw: ${String(err)}`) } @@ -1064,7 +1094,10 @@ export class NemoGymLanguageModel implements LanguageModelV3 { turn, global_turn: args.globalTurn, session_start_global_turn: args.sessionStartGlobalTurn, - timestamp: Date.now() / 1000, + request_kind: this._requestKind(args.messages, args.session.parentSessionID), + request_started_at: args.requestStartedAt / 1000, + latency: (args.responseCompletedAt - args.requestStartedAt) / 1000, + timestamp: args.responseCompletedAt / 1000, } const tmp = `${fpath}.tmp` await fs.writeFile(tmp, JSON.stringify(payload)) diff --git a/packages/opencode/test/bench/replay.test.ts b/packages/opencode/test/bench/replay.test.ts index 25ab4c23b67b..097c6ea95154 100644 --- a/packages/opencode/test/bench/replay.test.ts +++ b/packages/opencode/test/bench/replay.test.ts @@ -131,7 +131,13 @@ describe("parseReplayMessages", () => { test("joins array-of-parts user content for the initial instruction", () => { const raw = JSON.stringify([ - { role: "user", content: [{ type: "text", text: "part one" }, { type: "text", text: "part two" }] }, + { + role: "user", + content: [ + { type: "text", text: "part one" }, + { type: "text", text: "part two" }, + ], + }, ]) const { initialUserText } = parseReplayMessages(raw) expect(initialUserText).toBe("part one\npart two") @@ -186,10 +192,7 @@ describe("parseReplayManifest", () => { ) expect(manifest.rootSessionId).toBe("recorded-root") - expect(manifest.sessions.map((session) => session.sessionId)).toEqual([ - "recorded-grandchild", - "recorded-child", - ]) + expect(manifest.sessions.map((session) => session.sessionId)).toEqual(["recorded-grandchild", "recorded-child"]) expect(manifest.sessions[0]).toMatchObject({ parentSessionId: "recorded-child", spawnCallId: "call_nested",