Files
TJWaterAgent/src/runtime/opencode.ts
T

804 lines
25 KiB
TypeScript

import * as OpenCodePackage from "@opencode-ai/client";
import type { OpenCodeClient, OpenCodeEvent } from "@opencode-ai/client";
import { Service } from "@opencode-ai/client/service";
import { join, resolve } from "node:path";
import { config } from "../config.js";
import { logger } from "../logger.js";
const OPENCODE_V2_VERSION = "0.0.0-next-16741";
const isDevelopmentDebugLoggingEnabled = process.env.NODE_ENV === "development";
export const getEmbeddedServicePaths = (cwd: string, port: number) => {
const stateHome = resolve(cwd, "data", "opencode-service", String(port));
return {
stateHome,
registrationFile: join(stateHome, "opencode", "service.json"),
};
};
const logDevelopmentDebug = (
message: string,
metadata: Record<string, unknown>,
) => {
if (isDevelopmentDebugLoggingEnabled) {
logger.info(metadata, message);
}
};
export type RuntimeHealth = {
healthy: boolean;
version: string;
pid?: number;
};
export type RuntimeModelOverride = {
providerID: string;
modelID: string;
};
export type PermissionReply = "once" | "always" | "reject";
export type QuestionAnswers = string[][];
export type RuntimePart =
| { type: "text"; text: string }
| { type: "reasoning"; text: string };
export type RuntimeMessage = {
info: {
id: string;
role: string;
};
parts: RuntimePart[];
};
export type RuntimeError = {
name: string;
data?: { message?: string };
};
export type RuntimeToolPart = {
id: string;
callID: string;
messageID: string;
tool: string;
state: {
status: "pending" | "running" | "completed" | "error";
input: Record<string, unknown>;
error?: string;
};
};
type RuntimeQuestion = {
header: string;
question: string;
options: Array<{ label: string; description: string }>;
multiple?: boolean;
custom?: boolean;
};
type RuntimeQuestionTool = { messageID: string; callID: string };
export type RuntimeEvent =
| { type: "session.status"; sessionId: string; status: { type: "idle" | "busy" } | { type: "retry"; message: string } }
| { type: "session.execution.started"; sessionId: string }
| { type: "session.execution.succeeded"; sessionId: string }
| { type: "session.execution.failed"; sessionId: string; error: RuntimeError }
| { type: "session.execution.interrupted"; sessionId: string; reason: "user" | "shutdown" | "superseded" }
| { type: "permission.asked"; sessionId: string; request: { id: string; action: string; resources: string[]; save: string[]; metadata?: Record<string, unknown>; tool?: RuntimeQuestionTool } }
| { type: "permission.replied"; sessionId: string; requestId: string; reply: PermissionReply }
| { type: "question.asked"; sessionId: string; request: { id: string; questions: RuntimeQuestion[]; tool?: RuntimeQuestionTool } }
| { type: "question.replied"; sessionId: string; requestId: string; answers: QuestionAnswers }
| { type: "question.rejected"; sessionId: string; requestId: string }
| { type: "text.delta"; sessionId: string; partId: string; delta: string }
| { type: "reasoning.updated"; sessionId: string; partId: string; delta?: string; completed: boolean }
| { type: "tool.updated"; sessionId: string; part: RuntimeToolPart }
| { type: "skill.activated"; sessionId: string; name: string; payload: Record<string, unknown> }
| { type: "session.error"; sessionId?: string; error?: RuntimeError }
| { type: "session.idle"; sessionId: string };
type SessionMessageInfo = Awaited<
ReturnType<OpenCodeClient["message"]["list"]>
>["data"][number];
type AssistantContent = Extract<
SessionMessageInfo,
{ type: "assistant" }
>["content"][number];
type RuntimeFormField = {
key: string;
title?: string;
description?: string;
type: "string" | "number" | "integer" | "boolean" | "multiselect" | "external";
options?: Array<{ value: string; label: string; description?: string }>;
custom?: boolean;
};
const createClient = (options: {
baseUrl: string;
headers?: HeadersInit;
}): OpenCodeClient =>
(OpenCodePackage as unknown as {
OpenCode: { make: (value: typeof options) => OpenCodeClient };
}).OpenCode.make(options);
const toModelRef = (model: RuntimeModelOverride) => ({
providerID: model.providerID,
id: model.modelID,
});
const getConfiguredModel = (): RuntimeModelOverride => {
const [providerID, modelID] = config.OPENCODE_MODEL.split("/");
if (!providerID || !modelID) {
throw new Error(
`invalid OPENCODE_MODEL; expected provider/model, received ${config.OPENCODE_MODEL}`,
);
}
return { providerID, modelID };
};
const requireCompatibleHealth = (health: RuntimeHealth): RuntimeHealth => {
if (health.version !== OPENCODE_V2_VERSION) {
throw new Error(
`incompatible OpenCode service version: expected ${OPENCODE_V2_VERSION}, received ${health.version || "unknown"}`,
);
}
return health;
};
const toRuntimeMessage = (message: SessionMessageInfo): RuntimeMessage => {
if (message.type === "user") {
return {
info: { id: message.id, role: "user" },
parts: [{ type: "text", text: message.text }],
};
}
if (message.type === "assistant") {
return {
info: { id: message.id, role: "assistant" },
parts: message.content.flatMap((part: AssistantContent): RuntimePart[] => {
if (part.type === "text") {
return [{ type: "text", text: part.text }];
}
if (part.type === "reasoning") {
return [{ type: "reasoning", text: part.text }];
}
return [];
}),
};
}
return {
info: { id: message.id, role: message.type },
parts: [],
};
};
export class OpencodeRuntimeAdapter {
private clientPromise: Promise<OpenCodeClient> | null = null;
private ownsEmbeddedService = false;
private embeddedServiceFile: string | undefined;
private readonly pendingForms = new Map<string, RuntimeFormField[]>();
async ensureClient(): Promise<OpenCodeClient> {
if (!this.clientPromise) {
this.clientPromise = this.bootstrapClient().catch((error) => {
this.clientPromise = null;
throw error;
});
}
return this.clientPromise;
}
async health(): Promise<RuntimeHealth> {
const client = await this.ensureClient();
return requireCompatibleHealth(await client.health.get());
}
async warmup(): Promise<void> {
const client = await this.ensureClient();
requireCompatibleHealth(await client.health.get());
const session = await client.session.create({
title: "tjwater-agent-warmup",
agent: "instruction",
model: toModelRef(getConfiguredModel()),
location: { directory: process.cwd() },
});
try {
await Promise.all([
client.model.list({ location: { directory: process.cwd() } }),
client.plugin.list({ location: { directory: process.cwd() } }),
]);
} finally {
await client.session.remove({ sessionID: session.id }).catch((error: unknown) => {
logger.warn(
{ err: error, sessionId: session.id },
"failed to remove opencode warmup session",
);
});
}
}
async createSession(title?: string): Promise<{ id: string }> {
const client = await this.ensureClient();
const session = await client.session.create({
title,
agent: "instruction",
model: toModelRef(getConfiguredModel()),
location: { directory: process.cwd() },
});
if (!session.id) {
throw new Error("OpenCode V2 returned a session without an id");
}
const id = session.id;
return { id };
}
async sendPrompt(sessionId: string, text: string) {
await this.prompt(sessionId, text);
return this.messages(sessionId);
}
async prompt(sessionId: string, text: string, model?: RuntimeModelOverride) {
const client = await this.ensureClient();
const startedAt = Date.now();
logDevelopmentDebug("dispatching opencode v2 session.prompt", {
sessionId,
model: model ?? null,
textChars: text.length,
});
if (model) {
await client.session.switchModel({
sessionID: sessionId,
model: toModelRef(model),
});
}
await client.session.prompt({
sessionID: sessionId,
text,
resume: true,
});
logDevelopmentDebug("opencode v2 session.prompt returned", {
sessionId,
elapsedMs: Math.max(0, Date.now() - startedAt),
});
}
async messages(sessionId: string, limit = 20): Promise<RuntimeMessage[]> {
const client = await this.ensureClient();
const response = await client.message.list({
sessionID: sessionId,
limit,
order: "asc",
});
return response.data.map(toRuntimeMessage);
}
async revertMessage(sessionId: string, messageId: string) {
const client = await this.ensureClient();
return client.session.revert.stage({
sessionID: sessionId,
messageID: messageId,
files: false,
});
}
async revertToUserMessage(sessionId: string, options: { userOrdinal: number }) {
const messages = await this.messages(sessionId, 80);
const userMessages = messages.filter((message) => message.info.role === "user");
const targetUserMessage = userMessages[options.userOrdinal - 1];
if (!targetUserMessage) {
if (messages.length === 0 && options.userOrdinal === 1) {
logger.warn(
{ sessionId, userOrdinal: options.userOrdinal },
"skipping opencode revert because runtime session has no messages",
);
return;
}
throw new Error("target user message not found to revert");
}
const client = await this.ensureClient();
await client.session.revert.stage({
sessionID: sessionId,
messageID: targetUserMessage.info.id,
files: false,
});
await client.session.revert.commit({ sessionID: sessionId });
}
async abortSession(sessionId: string) {
const client = await this.ensureClient();
await client.session.interrupt({ sessionID: sessionId });
}
async waitForSessionIdle(sessionId: string, timeoutMs = config.OPENCODE_TIMEOUT_MS) {
const client = await this.ensureClient();
await client.session.wait(
{ sessionID: sessionId },
{ signal: AbortSignal.timeout(timeoutMs) },
);
}
async subscribeEvents(): Promise<AsyncIterable<RuntimeEvent>> {
const client = await this.ensureClient();
return normalizeEventStream(client.event.subscribe(), this.pendingForms);
}
async replyPermission(options: {
requestId: string;
sessionId: string;
reply: PermissionReply;
message?: string;
}) {
const client = await this.ensureClient();
await client.permission.reply({
sessionID: options.sessionId,
requestID: options.requestId,
reply: options.reply,
message: options.message,
});
}
async replyQuestion(options: {
requestId: string;
sessionId: string;
answers: QuestionAnswers;
}) {
const client = await this.ensureClient();
if (isFormRequestId(options.requestId)) {
const fields =
this.pendingForms.get(options.requestId) ??
(
await client.form.get({
sessionID: options.sessionId,
formID: options.requestId,
})
).fields as RuntimeFormField[];
await client.form.reply({
sessionID: options.sessionId,
formID: options.requestId,
answer: toFormAnswer(fields, options.answers),
});
this.pendingForms.delete(options.requestId);
return;
}
await client.question.reply({
sessionID: options.sessionId,
requestID: options.requestId,
answers: options.answers,
});
}
async rejectQuestion(options: { requestId: string; sessionId: string }) {
const client = await this.ensureClient();
if (isFormRequestId(options.requestId)) {
await client.form.cancel({
sessionID: options.sessionId,
formID: options.requestId,
});
this.pendingForms.delete(options.requestId);
return;
}
await client.question.reject({
sessionID: options.sessionId,
requestID: options.requestId,
});
}
async dispose(): Promise<void> {
if (this.ownsEmbeddedService && this.embeddedServiceFile) {
await Service.stop({ file: this.embeddedServiceFile }).catch((error) => {
logger.warn({ err: error }, "failed to stop opencode v2 service");
});
}
this.ownsEmbeddedService = false;
this.embeddedServiceFile = undefined;
this.clientPromise = null;
}
private async bootstrapClient(): Promise<OpenCodeClient> {
if (config.OPENCODE_MODE === "client") {
logger.info(
{ baseUrl: config.OPENCODE_CLIENT_BASE_URL, mode: config.OPENCODE_MODE },
"connecting to opencode v2 service in client mode",
);
return createClient({ baseUrl: config.OPENCODE_CLIENT_BASE_URL! });
}
process.env.TJWATER_AGENT_INTERNAL_BASE_URL = `http://127.0.0.1:${config.PORT}`;
process.env.TJWATER_AGENT_INTERNAL_TOKEN =
config.AGENT_INTERNAL_TOKEN ?? process.env.TJWATER_AGENT_INTERNAL_TOKEN ?? "";
logger.info(
{ model: config.OPENCODE_MODEL, mode: config.OPENCODE_MODE },
"starting opencode v2 service in embedded mode",
);
const startedAt = Date.now();
const { stateHome: serviceStateHome, registrationFile: serviceFile } =
getEmbeddedServicePaths(process.cwd(), config.PORT);
this.embeddedServiceFile = serviceFile;
let endpoint;
const previousStateHome = process.env.XDG_STATE_HOME;
process.env.XDG_STATE_HOME = serviceStateHome;
try {
// The service inherits callback credentials from this Agent process. Stop an
// orphaned instance before launch so a restarted Agent never reuses stale env.
await Service.stop({ file: serviceFile });
endpoint = await Service.ensure({
file: serviceFile,
command: ["opencode2", "serve", "--service"],
version: OPENCODE_V2_VERSION,
onStart: () => {
this.ownsEmbeddedService = true;
},
});
} catch (error) {
this.embeddedServiceFile = undefined;
if (isMissingOpencodeCli(error)) {
throw new Error(
"embedded mode requires the opencode2 CLI in PATH; otherwise set OPENCODE_MODE=client and provide OPENCODE_CLIENT_BASE_URL",
{ cause: error },
);
}
throw error;
} finally {
if (previousStateHome === undefined) {
delete process.env.XDG_STATE_HOME;
} else {
process.env.XDG_STATE_HOME = previousStateHome;
}
}
logger.info(
{
baseUrl: endpoint.url,
elapsedMs: Math.max(0, Date.now() - startedAt),
mode: config.OPENCODE_MODE,
},
"opencode v2 service started in embedded mode",
);
return createClient({
baseUrl: endpoint.url,
headers: Service.headers(endpoint),
});
}
}
export const opencodeRuntime = new OpencodeRuntimeAdapter();
function isMissingOpencodeCli(error: unknown): boolean {
const cause = error instanceof Error ? error.cause : undefined;
return (
(typeof error === "object" && error !== null && "code" in error && error.code === "ENOENT") ||
(typeof cause === "object" && cause !== null && "code" in cause && cause.code === "ENOENT")
);
}
async function* normalizeEventStream(
source: AsyncIterable<OpenCodeEvent>,
pendingForms: Map<string, RuntimeFormField[]>,
): AsyncIterable<RuntimeEvent> {
const tools = new Map<
string,
{ name: string; messageID: string; input: Record<string, unknown> }
>();
for await (const event of source) {
const normalized = normalizeEvent(event, tools, pendingForms);
if (normalized) {
yield normalized;
}
}
}
function normalizeEvent(
event: OpenCodeEvent,
tools: Map<string, { name: string; messageID: string; input: Record<string, unknown> }>,
pendingForms: Map<string, RuntimeFormField[]>,
): RuntimeEvent | undefined {
switch (event.type) {
case "session.status":
return { type: event.type, sessionId: event.data.sessionID, status: event.data.status };
case "session.execution.started":
case "session.execution.succeeded":
return { type: event.type, sessionId: event.data.sessionID };
case "session.execution.failed":
return {
type: event.type,
sessionId: event.data.sessionID,
error: toRuntimeError(event.data.error),
};
case "session.execution.interrupted":
return {
type: event.type,
sessionId: event.data.sessionID,
reason: event.data.reason,
};
case "permission.asked":
return {
type: event.type,
sessionId: event.data.sessionID,
request: {
id: event.data.id,
action: event.data.action,
resources: event.data.resources,
save: event.data.save ?? [],
metadata: event.data.metadata,
tool: event.data.source?.type === "tool" ? event.data.source : undefined,
},
};
case "permission.replied":
return {
type: event.type,
sessionId: event.data.sessionID,
requestId: event.data.requestID,
reply: event.data.reply,
};
case "question.asked":
return {
type: event.type,
sessionId: event.data.sessionID,
request: {
id: event.data.id,
questions: event.data.questions,
tool: event.data.tool,
},
};
case "question.replied":
return {
type: event.type,
sessionId: event.data.sessionID,
requestId: event.data.requestID,
answers: event.data.answers,
};
case "question.rejected":
return {
type: event.type,
sessionId: event.data.sessionID,
requestId: event.data.requestID,
};
case "form.created": {
const form = event.data.form;
const fields = [...form.fields] as RuntimeFormField[];
pendingForms.set(form.id, fields);
return {
type: "question.asked",
sessionId: form.sessionID,
request: {
id: form.id,
questions: fields.map((field) => toRuntimeQuestion(form.title, field)),
},
};
}
case "form.replied": {
const fields = pendingForms.get(event.data.id) ?? [];
pendingForms.delete(event.data.id);
return {
type: "question.replied",
sessionId: event.data.sessionID,
requestId: event.data.id,
answers: fields.map((field) => toQuestionAnswer(field, event.data.answer[field.key])),
};
}
case "form.cancelled":
pendingForms.delete(event.data.id);
return {
type: "question.rejected",
sessionId: event.data.sessionID,
requestId: event.data.id,
};
case "session.text.delta":
return {
type: "text.delta",
sessionId: event.data.sessionID,
partId: `text-${event.data.assistantMessageID}-${event.data.ordinal}`,
delta: event.data.delta,
};
case "session.reasoning.delta":
return {
type: "reasoning.updated",
sessionId: event.data.sessionID,
partId: `reasoning-${event.data.assistantMessageID}-${event.data.ordinal}`,
delta: event.data.delta,
completed: false,
};
case "session.reasoning.ended":
return {
type: "reasoning.updated",
sessionId: event.data.sessionID,
partId: `reasoning-${event.data.assistantMessageID}-${event.data.ordinal}`,
completed: true,
};
case "session.tool.input.started":
tools.set(event.data.callID, {
name: event.data.name,
messageID: event.data.assistantMessageID,
input: {},
});
return undefined;
case "session.tool.called": {
const tracked = tools.get(event.data.callID) ?? {
name: "unknown",
messageID: event.data.assistantMessageID,
input: {},
};
tracked.input = event.data.input;
tools.set(event.data.callID, tracked);
return {
type: "tool.updated",
sessionId: event.data.sessionID,
part: toRuntimeToolPart(event.data.callID, tracked, "running"),
};
}
case "session.tool.success": {
const tracked = tools.get(event.data.callID) ?? {
name: "unknown",
messageID: event.data.assistantMessageID,
input: {},
};
tools.delete(event.data.callID);
return {
type: "tool.updated",
sessionId: event.data.sessionID,
part: toRuntimeToolPart(event.data.callID, tracked, "completed"),
};
}
case "session.tool.failed": {
const tracked = tools.get(event.data.callID) ?? {
name: "unknown",
messageID: event.data.assistantMessageID,
input: {},
};
tools.delete(event.data.callID);
return {
type: "tool.updated",
sessionId: event.data.sessionID,
part: toRuntimeToolPart(
event.data.callID,
tracked,
"error",
getStructuredErrorMessage(event.data.error),
),
};
}
case "session.skill.activated":
return {
type: "skill.activated",
sessionId: event.data.sessionID,
name: event.data.name,
payload: { id: event.data.id, name: event.data.name },
};
case "session.error":
return {
type: event.type,
sessionId: event.data.sessionID,
error: event.data.error ? toRuntimeError(event.data.error) : undefined,
};
case "session.idle":
return { type: event.type, sessionId: event.data.sessionID };
default:
return undefined;
}
}
const isFormRequestId = (requestId: string) => requestId.startsWith("frm_");
function toRuntimeQuestion(formTitle: string, field: RuntimeFormField): RuntimeQuestion {
const options =
"options" in field && Array.isArray(field.options)
? field.options.map((option) => ({
label: option.label,
description: option.description ?? "",
}))
: field.type === "boolean"
? [
{ label: "是", description: "" },
{ label: "否", description: "" },
]
: [];
return {
header: field.title?.trim() || formTitle,
question: field.description?.trim() || field.title?.trim() || formTitle,
options,
multiple: field.type === "multiselect",
custom:
field.type === "string"
? (field.custom ?? options.length === 0)
: field.type !== "boolean" && field.type !== "multiselect",
};
}
function toFormAnswer(
fields: RuntimeFormField[],
answers: QuestionAnswers,
): Record<string, string | number | boolean | ReadonlyArray<string>> {
return Object.fromEntries(
fields.map((field, index) => [field.key, toFormValue(field, answers[index] ?? [])]),
);
}
function toFormValue(
field: RuntimeFormField,
answer: string[],
): string | number | boolean | ReadonlyArray<string> {
const mapOption = (value: string) =>
"options" in field && Array.isArray(field.options)
? (field.options.find((option) => option.label === value)?.value ?? value)
: value;
if (field.type === "multiselect") {
return answer.map(mapOption);
}
if (field.type === "boolean") {
return ["true", "yes", "是", "1"].includes((answer[0] ?? "").trim().toLowerCase());
}
if (field.type === "number" || field.type === "integer") {
const value = Number(answer[0]);
if (!Number.isFinite(value)) {
throw new Error(`invalid numeric answer for form field ${field.key}`);
}
return field.type === "integer" ? Math.trunc(value) : value;
}
return mapOption(answer[0] ?? "");
}
function toQuestionAnswer(
field: RuntimeFormField,
value: string | number | boolean | ReadonlyArray<string> | undefined,
): string[] {
const values = Array.isArray(value) ? value : value === undefined ? [] : [value];
return values.map((item) => {
const text = String(item);
if (field.type === "boolean") {
return item === true ? "是" : "否";
}
return "options" in field && Array.isArray(field.options)
? (field.options.find((option) => option.value === text)?.label ?? text)
: text;
});
}
function toRuntimeToolPart(
callID: string,
tool: { name: string; messageID: string; input: Record<string, unknown> },
status: RuntimeToolPart["state"]["status"],
error?: string,
): RuntimeToolPart {
return {
id: callID,
callID,
messageID: tool.messageID,
tool: tool.name,
state: { status, input: tool.input, error },
};
}
function toRuntimeError(error: unknown): RuntimeError {
if (typeof error === "object" && error !== null) {
const record = error as Record<string, unknown>;
const name =
typeof record._tag === "string"
? record._tag
: typeof record.name === "string"
? record.name
: "OpenCodeError";
return { name, data: { message: getStructuredErrorMessage(record) } };
}
return { name: "OpenCodeError", data: { message: String(error) } };
}
function getStructuredErrorMessage(error: unknown): string {
if (typeof error === "object" && error !== null) {
const record = error as Record<string, unknown>;
if (typeof record.message === "string") {
return record.message;
}
if (typeof record.data === "object" && record.data !== null) {
const message = (record.data as Record<string, unknown>).message;
if (typeof message === "string") {
return message;
}
}
}
return String(error);
}