7aa610294c
Proposal semantics: worker results no longer distinguish success/fail; only an explicit finish settles a pending proposal. Pending owner input follows up by resuming the original worker session. Dirty pending blocks start_next globally; clean pending only blocks its owner. QQ adapter now downloads image attachments and passes them as ACP image content blocks; video/file attachments degrade to link text.
258 lines
10 KiB
TypeScript
258 lines
10 KiB
TypeScript
import { spawn, type ChildProcessWithoutNullStreams } from "node:child_process";
|
|
import crypto from "node:crypto";
|
|
import fs from "node:fs";
|
|
import type { AcpConfig, PermissionPolicy } from "../config.js";
|
|
import type { WorkerProcessGroup } from "../core/proposal-store.js";
|
|
import type { ResolvedBot } from "../roles/role-registry.js";
|
|
import { AcpClient, type AcpPromptContent, type SafeActivityCategory } from "./client.js";
|
|
|
|
export class AcpSessionRestoreError extends Error {}
|
|
|
|
export type AcpWorkerPhase = "idle" | "initializing" | "processing";
|
|
export type AcpWorkerKind = "assistant" | "worker";
|
|
|
|
export interface AcpWorkerOptions {
|
|
kind?: AcpWorkerKind;
|
|
cwd?: string;
|
|
env?: Record<string, string>;
|
|
policy?: PermissionPolicy;
|
|
onActivity?: (category?: SafeActivityCategory) => void;
|
|
onIsolationViolation?: (worker: AcpWorker) => void;
|
|
onSpawn?: (group: WorkerProcessGroup) => Promise<void>;
|
|
allowUnverifiedAssistantAgent?: boolean;
|
|
}
|
|
|
|
export class AcpWorker {
|
|
private child?: ChildProcessWithoutNullStreams;
|
|
private client?: AcpClient;
|
|
private abort?: AbortController;
|
|
private exited = false;
|
|
private stopping = false;
|
|
private termination?: Promise<void>;
|
|
private stderrBytes = 0;
|
|
readonly kind: AcpWorkerKind;
|
|
readonly cwd: string;
|
|
nativeSessionId?: string;
|
|
processGroup?: WorkerProcessGroup;
|
|
isolationViolated = false;
|
|
supportsImages = false;
|
|
lastUsedAt = Date.now();
|
|
turnStartedAt?: number;
|
|
lastActivityAt?: number;
|
|
phase: AcpWorkerPhase = "idle";
|
|
inFlight = false;
|
|
|
|
constructor(
|
|
readonly bot: ResolvedBot,
|
|
private readonly config: AcpConfig,
|
|
private readonly onCrash: (worker: AcpWorker, error: Error) => void,
|
|
private readonly options: AcpWorkerOptions = {}
|
|
) {
|
|
this.kind = options.kind || "worker";
|
|
this.cwd = options.cwd || bot.workspace;
|
|
}
|
|
|
|
static async terminatePersistedGroup(group: WorkerProcessGroup, timeoutMs: number): Promise<void> {
|
|
if (!processGroupExists(group.pgid)) return;
|
|
if (!processGroupHasToken(group.pgid, group.token)) {
|
|
throw new Error(`Persisted ACP process group ${group.pgid} no longer matches its worker token`);
|
|
}
|
|
signalProcessGroup(group.pgid, "SIGKILL");
|
|
await waitForProcessGroupExit(group.pgid, timeoutMs);
|
|
}
|
|
|
|
async start(nativeSessionId?: string): Promise<string> {
|
|
const processToken = crypto.randomUUID();
|
|
this.child = spawn(this.bot.agent.command, this.bot.agent.args, {
|
|
cwd: this.cwd,
|
|
env: { ...process.env, ...this.bot.agent.env, ...this.options.env, GORI_AGENT_WORKER_TOKEN: processToken },
|
|
shell: false,
|
|
detached: true,
|
|
stdio: ["pipe", "pipe", "pipe"]
|
|
});
|
|
this.child.stderr.on("data", (chunk: Buffer) => this.logStderr(chunk));
|
|
this.child.once("error", (error) => this.crashed(error));
|
|
this.child.once("close", (code, signal) => {
|
|
this.exited = true;
|
|
if (!this.stopping) this.crashed(new Error(`ACP worker exited (code=${code ?? "null"}, signal=${signal ?? "null"})`));
|
|
});
|
|
this.client = new AcpClient(this.child, {
|
|
initializeTimeoutMs: this.config.initializeTimeoutMs,
|
|
policy: this.options.policy || this.bot.permissions,
|
|
forbidToolActivity: this.kind === "assistant",
|
|
onToolViolation: () => {
|
|
this.isolationViolated = true;
|
|
this.options.onIsolationViolation?.(this);
|
|
this.terminateImmediately();
|
|
},
|
|
onSessionActivity: (category) => {
|
|
this.lastActivityAt = Date.now();
|
|
this.options.onActivity?.(category);
|
|
}
|
|
});
|
|
try {
|
|
if (!this.child.pid) throw new Error("ACP worker has no process ID");
|
|
this.processGroup = { pgid: this.child.pid, token: processToken };
|
|
await this.options.onSpawn?.(this.processGroup);
|
|
const initialized = await this.client.initialize();
|
|
this.supportsImages = this.client.supportsImageInput();
|
|
if (this.kind === "assistant" && !this.options.allowUnverifiedAssistantAgent && initialized.agentInfo?.name !== "Kimi Code CLI") {
|
|
throw new Error("Assistant sessions require Kimi Code ACP with execution-layer no-tool profiles");
|
|
}
|
|
if (nativeSessionId) {
|
|
await this.client.resumeSession(nativeSessionId, this.cwd);
|
|
this.nativeSessionId = nativeSessionId;
|
|
} else {
|
|
this.nativeSessionId = await this.client.newSession(this.cwd);
|
|
}
|
|
return this.nativeSessionId;
|
|
} catch (error) {
|
|
await this.terminate();
|
|
if (nativeSessionId) throw new AcpSessionRestoreError(`Cannot restore ACP session '${nativeSessionId}': ${error instanceof Error ? error.message : String(error)}`);
|
|
throw error;
|
|
}
|
|
}
|
|
|
|
async prompt(input: string | AcpPromptContent[], phase: Exclude<AcpWorkerPhase, "idle"> = "processing"): Promise<string> {
|
|
if (!this.client || !this.nativeSessionId || this.exited) throw new Error("ACP worker is not available");
|
|
if (this.inFlight) throw new Error("ACP worker already has an in-flight turn");
|
|
const now = Date.now();
|
|
this.inFlight = true;
|
|
this.phase = phase;
|
|
this.turnStartedAt = now;
|
|
this.lastActivityAt = now;
|
|
this.lastUsedAt = now;
|
|
this.abort = new AbortController();
|
|
let timeout: NodeJS.Timeout | undefined;
|
|
const timeoutPromise = new Promise<never>((_resolve, reject) => {
|
|
timeout = setTimeout(() => {
|
|
void this.client?.cancel();
|
|
this.abort?.abort();
|
|
reject(new Error(`ACP prompt timed out after ${this.config.promptTimeoutMs}ms`));
|
|
setTimeout(() => { if (!this.exited) void this.terminate(); }, this.config.cancelGraceMs).unref();
|
|
}, this.config.promptTimeoutMs);
|
|
});
|
|
try {
|
|
const result = await Promise.race([this.client.prompt(input, this.abort.signal), timeoutPromise]);
|
|
if (this.kind === "assistant") await this.client.settleIsolation();
|
|
return result;
|
|
} finally {
|
|
if (timeout) clearTimeout(timeout);
|
|
this.inFlight = false;
|
|
this.phase = "idle";
|
|
this.turnStartedAt = undefined;
|
|
this.lastActivityAt = undefined;
|
|
this.abort = undefined;
|
|
this.lastUsedAt = Date.now();
|
|
}
|
|
}
|
|
|
|
async cancel(): Promise<boolean> {
|
|
if (!this.inFlight || !this.client) return false;
|
|
await this.client.cancel();
|
|
this.abort?.abort();
|
|
setTimeout(() => { if (this.inFlight) void this.terminate(); }, this.config.cancelGraceMs).unref();
|
|
return true;
|
|
}
|
|
|
|
terminateImmediately(): void {
|
|
this.stopping = true;
|
|
this.abort?.abort();
|
|
this.client?.close(new Error("ACP worker terminated because workspace ownership was lost"));
|
|
if (this.child) killProcessGroup(this.child, "SIGKILL");
|
|
this.termination ||= this.finishTermination("SIGKILL");
|
|
}
|
|
|
|
terminate(): Promise<void> {
|
|
if (this.termination) return this.termination;
|
|
this.stopping = true;
|
|
this.abort?.abort();
|
|
this.client?.close();
|
|
this.termination = this.finishTermination("SIGTERM");
|
|
return this.termination;
|
|
}
|
|
|
|
private async finishTermination(initialSignal: NodeJS.Signals): Promise<void> {
|
|
if (!this.child?.pid) return;
|
|
const pid = this.child.pid;
|
|
killProcessGroup(this.child, initialSignal);
|
|
await waitForExit(this.child, this.config.cancelGraceMs);
|
|
killProcessGroup(this.child, "SIGKILL");
|
|
await waitForExit(this.child, this.config.cancelGraceMs);
|
|
await waitForProcessGroupExit(pid, this.config.cancelGraceMs);
|
|
}
|
|
|
|
private logStderr(chunk: Buffer): void {
|
|
const remaining = Math.max(0, 16_384 - this.stderrBytes);
|
|
if (!remaining) return;
|
|
const text = chunk.subarray(0, remaining).toString("utf8").trimEnd();
|
|
this.stderrBytes += Buffer.byteLength(text);
|
|
if (text) console.error(`[acp:${this.kind}:${this.bot.agent.id}] ${text}`);
|
|
}
|
|
|
|
private crashed(error: Error): void {
|
|
if (this.stopping) return;
|
|
this.stopping = true;
|
|
this.abort?.abort();
|
|
this.client?.close(error);
|
|
this.termination = this.finishTermination("SIGKILL");
|
|
void this.termination.then(
|
|
() => this.onCrash(this, error),
|
|
(cleanupError) => this.onCrash(this, new Error(`${error.message}; process-group cleanup failed: ${cleanupError instanceof Error ? cleanupError.message : String(cleanupError)}`))
|
|
);
|
|
}
|
|
}
|
|
|
|
function killProcessGroup(child: ChildProcessWithoutNullStreams, signal: NodeJS.Signals): void {
|
|
if (child.pid) signalProcessGroup(child.pid, signal);
|
|
}
|
|
|
|
function signalProcessGroup(pgid: number, signal: NodeJS.Signals): void {
|
|
try { process.kill(-pgid, signal); }
|
|
catch (error) { if ((error as NodeJS.ErrnoException).code !== "ESRCH") throw error; }
|
|
}
|
|
|
|
function processGroupHasToken(pgid: number, token: string): boolean {
|
|
for (const entry of fs.readdirSync("/proc")) {
|
|
if (!/^\d+$/.test(entry)) continue;
|
|
try {
|
|
const stat = fs.readFileSync(`/proc/${entry}/stat`, "utf8");
|
|
const close = stat.lastIndexOf(")");
|
|
const fields = stat.slice(close + 2).split(" ");
|
|
if (Number(fields[2]) !== pgid) continue;
|
|
const environ = fs.readFileSync(`/proc/${entry}/environ`);
|
|
if (environ.toString("utf8").split("\0").includes(`GORI_AGENT_WORKER_TOKEN=${token}`)) return true;
|
|
} catch (error) {
|
|
const code = (error as NodeJS.ErrnoException).code;
|
|
if (code !== "ENOENT" && code !== "EACCES" && code !== "EPERM") throw error;
|
|
}
|
|
}
|
|
return false;
|
|
}
|
|
|
|
function waitForExit(child: ChildProcessWithoutNullStreams, timeoutMs: number): Promise<void> {
|
|
if (child.exitCode !== null || child.signalCode !== null) return Promise.resolve();
|
|
return new Promise((resolve) => {
|
|
const timer = setTimeout(resolve, timeoutMs);
|
|
child.once("close", () => { clearTimeout(timer); resolve(); });
|
|
});
|
|
}
|
|
|
|
async function waitForProcessGroupExit(pid: number, timeoutMs: number): Promise<void> {
|
|
const deadline = Date.now() + timeoutMs;
|
|
while (processGroupExists(pid) && Date.now() < deadline) {
|
|
await new Promise((resolve) => setTimeout(resolve, 10));
|
|
}
|
|
if (processGroupExists(pid)) throw new Error(`ACP process group ${pid} did not terminate`);
|
|
}
|
|
|
|
function processGroupExists(pid: number): boolean {
|
|
try { process.kill(-pid, 0); return true; }
|
|
catch (error) {
|
|
const code = (error as NodeJS.ErrnoException).code;
|
|
if (code === "ESRCH") return false;
|
|
if (code === "EPERM") return true;
|
|
throw error;
|
|
}
|
|
}
|