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

const startedAt = Date.now()
const childEnv: NodeJS.ProcessEnv = {
...process.env,
// Run-isolated opencode state.
Expand Down Expand Up @@ -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,
Expand All @@ -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, {
Expand All @@ -608,6 +642,7 @@ async function main() {
metrics: {
bench_run_time: benchRunTime,
opencode_exit_code: result.exitCode,
...perTurnMetrics,
patch_mode: args.patchMode,
},
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/bench/replay.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 = []
Expand Down
8 changes: 6 additions & 2 deletions packages/opencode/src/cli/cmd/run.ts
Original file line number Diff line number Diff line change
Expand Up @@ -434,7 +434,7 @@ export const RunCommand = effectCmd({

function emit(type: string, data: Record<string, unknown>) {
if (args.format === "json") {
if (benchEventTypesOnly) {
if (benchEventTypesOnly && type !== "tool_use") {
process.stdout.write(type + EOL)
return true
}
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading