Skip to content
Open
41 changes: 36 additions & 5 deletions packages/opencode/src/bench/cli.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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<string, ActionExecutionLatencyMetric>()
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)
Expand All @@ -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() })
})
})
}
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -522,7 +542,6 @@ async function main() {
maxTokens: forcedMaxTokens,
})

const startedAt = Date.now()
const childEnv: NodeJS.ProcessEnv = {
...process.env,
// Run-isolated opencode state.
Expand Down Expand Up @@ -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,
Expand All @@ -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, {
Expand All @@ -572,6 +602,7 @@ async function main() {
metrics: {
bench_run_time: benchRunTime,
opencode_exit_code: result.exitCode,
...perTurnMetrics,
},
error,
})
Expand Down
218 changes: 218 additions & 0 deletions packages/opencode/src/bench/metrics.ts
Original file line number Diff line number Diff line change
@@ -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<string, unknown>
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<string, any>
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<string, unknown>)
: 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<string, unknown>)
: 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<CompletionMetrics> {
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<string, unknown>,
): Promise<void> {
if (!metricsPath) return

let existing: Record<string, unknown> = {}
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)
}
6 changes: 5 additions & 1 deletion packages/opencode/src/cli/cmd/run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading