diff --git a/packages/opencode/src/bench/cli.ts b/packages/opencode/src/bench/cli.ts index c1374c03a383..4ab52dc08063 100644 --- a/packages/opencode/src/bench/cli.ts +++ b/packages/opencode/src/bench/cli.ts @@ -25,6 +25,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" // opencode's built-in anthropic system prompt — Bun bundles .txt as a string. // Used as the default when no --system-prompt override is passed. import PROMPT_ANTHROPIC from "../session/prompt/anthropic.txt" @@ -347,7 +353,12 @@ function runOpencode(args: { env: NodeJS.ProcessEnv opencodeBin: string agent: string -}): Promise<{ exitCode: number; stdout: string; stderr: string }> { +}): Promise<{ + exitCode: number + stdout: string + stderr: string + actionExecutionLatencies: ActionExecutionLatencyMetric[] +}> { // 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 @@ -401,12 +412,19 @@ 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) => { lineBuf += b.toString("utf8") let idx: number while ((idx = lineBuf.indexOf("\n")) >= 0) { - const line = scrub(lineBuf.slice(0, idx)) + const rawLine = lineBuf.slice(0, idx) lineBuf = lineBuf.slice(idx + 1) + const actionMetric = parseToolExecutionMetric(rawLine) + if (actionMetric) { + const metricID = `${actionMetric.session_id}:${actionMetric.observation_id}` + actionExecutionLatencies.set(metricID, actionMetric) + } + const line = scrub(rawLine) // Forward to our stdout so the gym log captures the event stream. process.stdout.write(line + "\n") stdout = (stdout + line + "\n").slice(-MAX_KEEP) @@ -417,10 +435,11 @@ function runOpencode(args: { stderr = (stderr + chunk).slice(-MAX_KEEP) process.stderr.write(chunk) }) - child.on("close", (code) => resolve({ exitCode: code ?? 0, stdout, stderr })) + const metrics = () => [...actionExecutionLatencies.values()].sort((a, b) => a.timestamp.localeCompare(b.timestamp)) + child.on("close", (code) => resolve({ exitCode: code ?? 0, stdout, stderr, actionExecutionLatencies: metrics() })) child.on("error", (err) => { stderr += String(err) - resolve({ exitCode: 999, stdout, stderr }) + resolve({ exitCode: 999, stdout, stderr, actionExecutionLatencies: metrics() }) }) }) } @@ -484,6 +503,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. @@ -522,7 +542,6 @@ async function main() { maxTokens: forcedMaxTokens, }) - const startedAt = Date.now() const childEnv: NodeJS.ProcessEnv = { ...process.env, // Run-isolated opencode state. @@ -552,6 +571,8 @@ async function main() { } const opencodeBin = detectOpencodeBin() + const initializeRuntimeTime = (Date.now() - initializeStartedAt) / 1000 + const startedAt = Date.now() const result = await runOpencode({ workspaceRoot, modelName, @@ -563,6 +584,15 @@ async function main() { const patch = await captureGitDiff(workspaceRoot) 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: string | null = result.exitCode === 0 ? null : `opencode_exit_${result.exitCode}` const outPath = await writeOutputJsonl(args.outputDir, instance.instance_id, { @@ -572,6 +602,7 @@ async function main() { metrics: { bench_run_time: benchRunTime, opencode_exit_code: result.exitCode, + ...perTurnMetrics, }, 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/cli/cmd/run.ts b/packages/opencode/src/cli/cmd/run.ts index a05b273e4489..d1067e02a8a1 100644 --- a/packages/opencode/src/cli/cmd/run.ts +++ b/packages/opencode/src/cli/cmd/run.ts @@ -459,10 +459,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 0a0f9e56c386..f3946519e2cd 100644 --- a/packages/opencode/src/provider/sdk/nemo-gym/language-model.ts +++ b/packages/opencode/src/provider/sdk/nemo-gym/language-model.ts @@ -75,6 +75,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 { @@ -155,6 +158,17 @@ export class NemoGymLanguageModel implements LanguageModelV3 { // a Map keeps their dump filenames from clobbering the main session's. private readonly turnCounters: Map = 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 @@ -193,11 +207,15 @@ export class NemoGymLanguageModel implements LanguageModelV3 { async doGenerate(options: LanguageModelV3CallOptions) { const { warnings, loggedMessages, requestParams } = await this._buildRequestParams(options) const session = this._sessionFromHeaders(options.headers) + const turn = this._nextTurn(session.sessionID) + 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") - const msg: ChatResponseChoice["message"] = choice.message ?? ({ role: "assistant" } as ChatResponseChoice["message"]) + const msg: ChatResponseChoice["message"] = + choice.message ?? ({ role: "assistant" } as ChatResponseChoice["message"]) const providerSpecificFields = this._extractProviderFields(msg) const providerMetadata = this._buildProviderMetadata(providerSpecificFields) @@ -223,6 +241,9 @@ export class NemoGymLanguageModel implements LanguageModelV3 { providerSpecificFields, requestParams, session, + turn, + requestStartedAt, + responseCompletedAt, }) return { @@ -249,7 +270,10 @@ export class NemoGymLanguageModel implements LanguageModelV3 { controller.enqueue({ type: "stream-start", warnings }) try { + const turn = self._nextTurn(session.sessionID) + 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") @@ -317,6 +341,9 @@ export class NemoGymLanguageModel implements LanguageModelV3 { providerSpecificFields, requestParams, session, + turn, + requestStartedAt, + responseCompletedAt, }) controller.enqueue({ @@ -575,7 +602,10 @@ export class NemoGymLanguageModel implements LanguageModelV3 { return md } - private _mapFinishReason(raw: string | null): { unified: "stop" | "length" | "tool-calls" | "error" | "other"; raw: string | undefined } { + private _mapFinishReason(raw: string | null): { + unified: "stop" | "length" | "tool-calls" | "error" | "other" + raw: string | undefined + } { if (!raw) return { unified: "other", raw: undefined } switch (raw) { case "stop": @@ -612,11 +642,19 @@ export class NemoGymLanguageModel implements LanguageModelV3 { providerSpecificFields: Record requestParams: Record session: { sessionID: string; parentSessionID: string | undefined } + turn: number + requestStartedAt: number + responseCompletedAt: number }) { - const turn = this._nextTurn(args.session.sessionID) if (this.cfg.onCompletion) { try { - await this.cfg.onCompletion({ turn, ...args }) + await this.cfg.onCompletion({ + turn: args.turn, + messages: args.messages, + response: args.response, + providerSpecificFields: args.providerSpecificFields, + requestParams: args.requestParams, + }) } catch (err) { console.warn(`[nemo-gym] onCompletion hook threw: ${String(err)}`) } @@ -626,7 +664,7 @@ export class NemoGymLanguageModel implements LanguageModelV3 { try { await fs.mkdir(this.cfg.completionsDir, { recursive: true }) - const turnStr = String(turn).padStart(4, "0") + const turnStr = String(args.turn).padStart(4, "0") const safeModel = this.modelId.replace(/\//g, "__") // sessionID is part of the filename so subagent dumps don't clobber the // main session's. Sanitized for filesystem safety. @@ -644,8 +682,11 @@ export class NemoGymLanguageModel implements LanguageModelV3 { kwargs, session_id: args.session.sessionID, parent_session_id: args.session.parentSessionID ?? null, - turn, - timestamp: Date.now() / 1000, + turn: args.turn, + 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))