846 lines
39 KiB
TypeScript
846 lines
39 KiB
TypeScript
import crypto from "node:crypto";
|
|
import fs from "node:fs";
|
|
import path from "node:path";
|
|
import type { AcpConfig, PermissionPolicy } from "../config.js";
|
|
import { chatKeyFor, type AssistantBinding, type DurableSessionStore } from "../core/durable-session-store.js";
|
|
import type { Proposal, ProposalPending, ProposalStore } from "../core/proposal-store.js";
|
|
import type { ResolvedBot } from "../roles/role-registry.js";
|
|
import type { ConversationRequest, ConversationResponse, ConversationRuntime, RuntimeEventSink, RuntimeStats } from "./types.js";
|
|
import { AcpSessionRestoreError, AcpWorker } from "./worker.js";
|
|
|
|
const ASSISTANT_POLICY: PermissionPolicy = { mode: "deny", allowedTools: [], allowedCommandPatterns: [] };
|
|
const ASSISTANT_ENVELOPE = /<GORI_ASSISTANT_ACTION_V1>(?<json>[\s\S]*?)<\/GORI_ASSISTANT_ACTION_V1>\s*$/;
|
|
const WORKER_ENVELOPE = /<GORI_WORKER_RESULT_V1>(?<json>[\s\S]*?)<\/GORI_WORKER_RESULT_V1>\s*$/;
|
|
|
|
export type AssistantAction =
|
|
| { type: "create_proposal"; title: string; goal: string; steps: string[] }
|
|
| { type: "confirm"; id?: string; answer?: string }
|
|
| { type: "start_next" }
|
|
| { type: "cancel"; id?: string }
|
|
| { type: "stop" };
|
|
|
|
export interface ParsedAssistantEnvelope {
|
|
reply: string;
|
|
actions: AssistantAction[];
|
|
}
|
|
|
|
export type WorkerResultStatus = "SUCCESS" | "NEEDS_CONFIRMATION" | "FAILED";
|
|
|
|
export interface ParsedWorkerResult {
|
|
status: WorkerResultStatus;
|
|
summary: string;
|
|
question?: string;
|
|
nextStep?: string;
|
|
dirty?: boolean;
|
|
}
|
|
|
|
interface ActiveWorker {
|
|
proposalId: string;
|
|
worker: AcpWorker;
|
|
settled: boolean;
|
|
run?: Promise<void>;
|
|
}
|
|
|
|
export interface AssistantManagerOptions {
|
|
assistantWorkspaceHome?: string;
|
|
kimiCodeHome?: string;
|
|
allowUnverifiedAssistantAgent?: boolean;
|
|
maxAssistantSessions?: number;
|
|
}
|
|
|
|
export class AssistantManager implements ConversationRuntime {
|
|
private readonly assistantWorkers = new Map<string, AcpWorker>();
|
|
private readonly chatChains = new Map<string, Promise<unknown>>();
|
|
private readonly assistantSalt = crypto.randomBytes(32);
|
|
private readonly maxAssistantSessions: number;
|
|
private readonly sweeper: NodeJS.Timeout;
|
|
private active?: ActiveWorker;
|
|
private eventSink?: RuntimeEventSink;
|
|
private kimiHome?: string;
|
|
private crashes = 0;
|
|
private shuttingDown = false;
|
|
private schedulerChain: Promise<void> = Promise.resolve();
|
|
|
|
constructor(
|
|
private readonly config: AcpConfig,
|
|
private readonly bot: ResolvedBot,
|
|
private readonly store: DurableSessionStore,
|
|
readonly proposals: ProposalStore,
|
|
private readonly options: AssistantManagerOptions = {}
|
|
) {
|
|
this.maxAssistantSessions = options.maxAssistantSessions || 4;
|
|
this.sweeper = setInterval(() => void this.sweep().catch(() => undefined), config.sweepIntervalMs);
|
|
this.sweeper.unref();
|
|
}
|
|
|
|
async initialize(): Promise<void> {
|
|
for (const proposal of this.proposals.list({ status: "working" })) {
|
|
if (proposal.workerProcessGroup) {
|
|
try {
|
|
await AcpWorker.terminatePersistedGroup(proposal.workerProcessGroup, this.config.cancelGraceMs);
|
|
} catch (error) {
|
|
console.error(`Persisted worker cleanup failed for proposal ${proposal.id} (${error instanceof Error ? error.name : "unknown error"})`);
|
|
}
|
|
}
|
|
const summary = "worker_lost: the runner restarted while this proposal was working; its worker was terminated and the proposal was marked failed";
|
|
await this.proposals.update(proposal.id, {
|
|
status: "failed",
|
|
pending: { kind: "failure", summary },
|
|
lastWorkerSummary: summary,
|
|
workerProcessGroup: undefined,
|
|
finishedAt: Date.now()
|
|
});
|
|
}
|
|
}
|
|
|
|
setEventSink(sink: RuntimeEventSink): void {
|
|
this.eventSink = sink;
|
|
}
|
|
|
|
async prompt(request: ConversationRequest): Promise<ConversationResponse> {
|
|
if (this.shuttingDown) throw new Error("ACP runtime is shutting down");
|
|
const conversationKey = assistantKeyFor(chatKeyFor(request.platform, request.chatId), request.userId);
|
|
const reply = await this.enqueueChat(conversationKey, () => this.runAssistantTurn(conversationKey, request));
|
|
return { text: reply, botId: this.bot.id, agentId: this.bot.agent.id, mode: "assistant" };
|
|
}
|
|
|
|
async cancel(platform: string, chatId: string, userId: string): Promise<boolean> {
|
|
const chatKey = chatKeyFor(platform, chatId);
|
|
return this.scheduler(() => this.cancelLocked(chatKey, userId));
|
|
}
|
|
|
|
async confirm(platform: string, chatId: string, userId: string): Promise<boolean> {
|
|
const chatKey = chatKeyFor(platform, chatId);
|
|
return this.scheduler(() => this.confirmLocked(chatKey, userId));
|
|
}
|
|
|
|
async stop(platform: string, chatId: string, userId: string): Promise<boolean> {
|
|
const chatKey = chatKeyFor(platform, chatId);
|
|
return this.scheduler(() => this.stopActiveLocked(chatKey, userId));
|
|
}
|
|
|
|
listProposals(platform: string, chatId: string, userId: string): Proposal[] {
|
|
const chatKey = chatKeyFor(platform, chatId);
|
|
return this.proposals.list().filter((proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId);
|
|
}
|
|
|
|
async reset(platform: string, chatId: string, userId: string): Promise<void> {
|
|
const conversationKey = assistantKeyFor(chatKeyFor(platform, chatId), userId);
|
|
const worker = this.assistantWorkers.get(conversationKey);
|
|
if (worker) {
|
|
this.assistantWorkers.delete(conversationKey);
|
|
await worker.terminate();
|
|
}
|
|
await this.store.deleteBinding(conversationKey);
|
|
}
|
|
|
|
status(platform: string, chatId: string, userId?: string): Record<string, string | number | boolean> {
|
|
const proposals = this.proposals.list();
|
|
const active = this.active;
|
|
const activeProposal = active ? this.proposals.get(active.proposalId) : undefined;
|
|
return {
|
|
bot: this.bot.id,
|
|
agent: this.bot.agent.id,
|
|
workspace: this.bot.workspace,
|
|
assistantSessions: this.assistantWorkers.size,
|
|
proposals: proposals.length,
|
|
queuedProposals: proposals.filter((proposal) => proposal.status === "queued").length,
|
|
workingProposals: proposals.filter((proposal) => proposal.status === "working").length,
|
|
awaitingConfirmation: proposals.filter((proposal) => proposal.status === "awaiting_user_confirmation").length,
|
|
workerRunning: Boolean(active?.worker.inFlight),
|
|
owner: Boolean(userId && activeProposal && activeProposal.ownerChatKey === chatKeyFor(platform, chatId) && activeProposal.requesterUserId === userId)
|
|
};
|
|
}
|
|
|
|
stats(): RuntimeStats {
|
|
const assistants = [...this.assistantWorkers.values()];
|
|
return {
|
|
activeWorkers: assistants.length + (this.active ? 1 : 0),
|
|
inFlight: assistants.filter((worker) => worker.inFlight).length + (this.active?.worker.inFlight ? 1 : 0),
|
|
crashes: this.crashes,
|
|
persistedBindings: this.store.stats().bindings
|
|
};
|
|
}
|
|
|
|
async shutdown(): Promise<void> {
|
|
if (this.shuttingDown) return;
|
|
this.shuttingDown = true;
|
|
clearInterval(this.sweeper);
|
|
const active = this.active;
|
|
this.active = undefined;
|
|
if (active) {
|
|
active.settled = true;
|
|
if (active.worker.inFlight) await active.worker.cancel().catch(() => undefined);
|
|
active.worker.terminateImmediately();
|
|
await active.worker.terminate().catch(() => undefined);
|
|
await active.run?.catch(() => undefined);
|
|
}
|
|
const assistants = [...this.assistantWorkers.values()];
|
|
this.assistantWorkers.clear();
|
|
await Promise.all(assistants.map((worker) => worker.terminate().catch(() => undefined)));
|
|
await Promise.allSettled([...this.chatChains.values()]);
|
|
}
|
|
|
|
private async runAssistantTurn(conversationKey: string, request: ConversationRequest): Promise<string> {
|
|
const worker = await this.acquireAssistantWorker(conversationKey);
|
|
try {
|
|
const reply = await worker.prompt(this.assistantUserPrompt(request));
|
|
const parsed = await this.parseAssistantReply(worker, reply);
|
|
if (worker.isolationViolated) throw new Error("Assistant session attempted forbidden tool activity");
|
|
await this.store.touchBinding(conversationKey);
|
|
await this.executeActions(chatKeyFor(request.platform, request.chatId), request, parsed.actions);
|
|
return parsed.reply;
|
|
} catch (error) {
|
|
if (worker.isolationViolated) await this.store.deleteBinding(conversationKey).catch(() => undefined);
|
|
throw error;
|
|
}
|
|
}
|
|
|
|
private async parseAssistantReply(worker: AcpWorker, reply: string): Promise<ParsedAssistantEnvelope> {
|
|
let parsed = parseAssistantActions(reply);
|
|
if (!parsed) parsed = parseAssistantActions(await worker.prompt(assistantRepairPrompt()));
|
|
if (!parsed) throw new Error("Assistant did not return a valid GORI_ASSISTANT_ACTION_V1 envelope after one repair attempt");
|
|
return parsed;
|
|
}
|
|
|
|
private async executeActions(chatKey: string, request: ConversationRequest, actions: AssistantAction[]): Promise<void> {
|
|
for (const action of actions) {
|
|
if (action.type === "create_proposal") {
|
|
await this.proposals.create({
|
|
title: action.title,
|
|
goal: action.goal,
|
|
steps: action.steps,
|
|
ownerChatKey: chatKey,
|
|
requesterUserId: request.userId
|
|
});
|
|
} else if (action.type === "confirm") {
|
|
await this.scheduler(() => this.confirmLocked(chatKey, request.userId, action.id, action.answer));
|
|
} else if (action.type === "start_next") {
|
|
await this.scheduler(() => this.tryStartLocked(chatKey, request.userId));
|
|
} else if (action.type === "cancel") {
|
|
await this.scheduler(() => this.cancelLocked(chatKey, request.userId, action.id));
|
|
} else if (action.type === "stop") {
|
|
await this.scheduler(() => this.stopActiveLocked(chatKey, request.userId));
|
|
}
|
|
}
|
|
}
|
|
|
|
private async confirmLocked(chatKey: string, userId: string, id?: string, answer?: string): Promise<boolean> {
|
|
const proposal = this.resolveTargetProposal(chatKey, userId, id, ["proposed", "awaiting_user_confirmation"]);
|
|
if (!proposal) return false;
|
|
if (proposal.status === "proposed") {
|
|
await this.proposals.update(proposal.id, { status: "queued", confirmedAt: Date.now() });
|
|
await this.tryStartLocked(chatKey, userId);
|
|
return true;
|
|
}
|
|
const pending = proposal.pending;
|
|
if (!pending) return false;
|
|
if (pending.kind === "success") {
|
|
await this.terminateActiveLocked(proposal.id);
|
|
await this.proposals.update(proposal.id, { status: "completed", pending: undefined, finishedAt: Date.now() });
|
|
return true;
|
|
}
|
|
if (pending.kind === "failure") {
|
|
await this.terminateActiveLocked(proposal.id);
|
|
if (answer && answer.trim().toLowerCase() === "retry") {
|
|
await this.proposals.update(proposal.id, { status: "working", pending: undefined, startedAt: Date.now() });
|
|
this.launchWorker(proposal.id);
|
|
} else {
|
|
await this.proposals.update(proposal.id, { status: "failed", pending: undefined, finishedAt: Date.now() });
|
|
}
|
|
return true;
|
|
}
|
|
await this.proposals.update(proposal.id, { status: "working", pending: undefined });
|
|
this.continueWorker(proposal.id, answer || "Confirmed. Please continue.", pending.question);
|
|
return true;
|
|
}
|
|
|
|
private async cancelLocked(chatKey: string, userId: string, id?: string): Promise<boolean> {
|
|
const proposal = this.resolveTargetProposal(chatKey, userId, id, ["proposed", "queued"]);
|
|
if (!proposal) return false;
|
|
await this.proposals.update(proposal.id, { status: "cancelled", finishedAt: Date.now() });
|
|
return true;
|
|
}
|
|
|
|
private async stopActiveLocked(chatKey: string, userId: string): Promise<boolean> {
|
|
const active = this.active;
|
|
const activeProposal = active ? this.proposals.get(active.proposalId) : undefined;
|
|
if (active && activeProposal && activeProposal.ownerChatKey === chatKey && activeProposal.requesterUserId === userId) {
|
|
active.settled = true;
|
|
this.active = undefined;
|
|
if (active.worker.inFlight) await active.worker.cancel().catch(() => undefined);
|
|
active.worker.terminateImmediately();
|
|
await active.worker.terminate().catch(() => undefined);
|
|
void active.run?.catch(() => undefined);
|
|
await this.proposals.update(activeProposal.id, {
|
|
status: "cancelled",
|
|
pending: undefined,
|
|
workerProcessGroup: undefined,
|
|
lastWorkerSummary: "stopped by user",
|
|
finishedAt: Date.now()
|
|
});
|
|
return true;
|
|
}
|
|
const waiting = this.resolveTargetProposal(chatKey, userId, undefined, ["awaiting_user_confirmation"]);
|
|
if (!waiting?.pending) return false;
|
|
await this.terminateActiveLocked(waiting.id);
|
|
if (waiting.pending.kind === "success") {
|
|
await this.proposals.update(waiting.id, { status: "completed", pending: undefined, finishedAt: Date.now() });
|
|
} else if (waiting.pending.kind === "failure") {
|
|
await this.proposals.update(waiting.id, { status: "failed", pending: undefined, finishedAt: Date.now() });
|
|
} else {
|
|
await this.proposals.update(waiting.id, {
|
|
status: "cancelled",
|
|
pending: undefined,
|
|
lastWorkerSummary: "stopped by user",
|
|
finishedAt: Date.now()
|
|
});
|
|
}
|
|
return true;
|
|
}
|
|
|
|
private async terminateActiveLocked(proposalId: string): Promise<void> {
|
|
const active = this.active;
|
|
if (!active || active.proposalId !== proposalId) return;
|
|
active.settled = true;
|
|
this.active = undefined;
|
|
active.worker.terminateImmediately();
|
|
await active.worker.terminate().catch(() => undefined);
|
|
void active.run?.catch(() => undefined);
|
|
await this.proposals.update(proposalId, { workerProcessGroup: undefined });
|
|
}
|
|
|
|
private async tryStartLocked(chatKey: string, userId: string): Promise<void> {
|
|
if (this.active) return;
|
|
if (this.proposals.list({ status: "working" }).length > 0) return;
|
|
if (this.proposals.list({ status: "awaiting_user_confirmation" }).length > 0) return;
|
|
const next = this.proposals.list({ status: "queued" })
|
|
.find((proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId);
|
|
if (!next) return;
|
|
await this.proposals.update(next.id, { status: "working", startedAt: Date.now() });
|
|
this.launchWorker(next.id);
|
|
}
|
|
|
|
private resolveTargetProposal(chatKey: string, userId: string, id: string | undefined, statuses: Proposal["status"][]): Proposal | undefined {
|
|
const owned = (proposal: Proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId;
|
|
if (id) {
|
|
const proposal = this.proposals.get(id);
|
|
return proposal && owned(proposal) && statuses.includes(proposal.status) ? proposal : undefined;
|
|
}
|
|
return this.proposals.list()
|
|
.filter((proposal) => owned(proposal) && statuses.includes(proposal.status))
|
|
.sort((left, right) => right.updatedAt - left.updatedAt)[0];
|
|
}
|
|
|
|
private launchWorker(proposalId: string, resumeContext?: string): void {
|
|
const proposal = this.proposals.get(proposalId);
|
|
if (!proposal || proposal.status !== "working") return;
|
|
const worker = new AcpWorker(this.bot, this.config, (crashed, error) => this.onWorkerCrash(proposalId, crashed, error), {
|
|
kind: "worker",
|
|
cwd: this.bot.workspace,
|
|
onSpawn: async (group) => {
|
|
if (this.active?.proposalId !== proposalId) throw new Error("Worker proposal changed before process identity was persisted");
|
|
await this.proposals.update(proposalId, { workerProcessGroup: group });
|
|
}
|
|
});
|
|
const active: ActiveWorker = { proposalId, worker, settled: false };
|
|
this.active = active;
|
|
const run = this.startWorkerRun(active, resumeContext);
|
|
active.run = run;
|
|
void run.catch((error) => this.handleWorkerRunFailure(active, error));
|
|
}
|
|
|
|
private continueWorker(proposalId: string, answer: string, question?: string): void {
|
|
const active = this.active;
|
|
if (active && active.proposalId === proposalId && active.worker.nativeSessionId) {
|
|
active.settled = false;
|
|
const run = this.runWorkerTurn(active, workerAnswerPrompt(answer, question));
|
|
active.run = run;
|
|
void run.catch((error) => this.handleWorkerRunFailure(active, error));
|
|
return;
|
|
}
|
|
const proposal = this.proposals.get(proposalId);
|
|
const context = proposal
|
|
? `The previous worker was lost. Its last summary: ${proposal.lastWorkerSummary || "unknown"}. Its question: ${question || "unknown"}. The user answered: ${answer}`
|
|
: answer;
|
|
this.launchWorker(proposalId, context);
|
|
}
|
|
|
|
private async startWorkerRun(active: ActiveWorker, resumeContext?: string): Promise<void> {
|
|
const nativeSessionId = await active.worker.start();
|
|
if (this.active !== active || active.settled) throw new Error("Worker was superseded before its session started");
|
|
await this.proposals.update(active.proposalId, { workerNativeSessionId: nativeSessionId });
|
|
await active.worker.prompt(this.bot.workerBootstrap, "initializing");
|
|
if (this.active !== active || active.settled) throw new Error("Worker was superseded during bootstrap");
|
|
const proposal = this.proposals.get(active.proposalId);
|
|
if (!proposal) throw new Error("Worker proposal no longer exists");
|
|
await this.runWorkerTurn(active, workerTaskPrompt(proposal, resumeContext));
|
|
}
|
|
|
|
private async runWorkerTurn(active: ActiveWorker, promptText: string): Promise<void> {
|
|
const reply = await active.worker.prompt(withWorkerResultProtocol(promptText));
|
|
if (this.active !== active || active.settled) throw new Error("Worker was superseded during its turn");
|
|
let result = parseWorkerResult(reply);
|
|
if (!result) {
|
|
const repaired = await active.worker.prompt(workerRepairPrompt());
|
|
if (this.active !== active || active.settled) throw new Error("Worker was superseded during result repair");
|
|
result = parseWorkerResult(repaired);
|
|
}
|
|
if (!result) throw new Error("Worker did not return a valid GORI_WORKER_RESULT_V1 envelope after one repair attempt");
|
|
await this.settleWorkerResult(active, result);
|
|
}
|
|
|
|
private async settleWorkerResult(active: ActiveWorker, result: ParsedWorkerResult): Promise<void> {
|
|
if (this.active !== active || active.settled || this.shuttingDown) return;
|
|
const transitioned = await this.scheduler(async () => {
|
|
if (this.active !== active || active.settled || this.shuttingDown) return false;
|
|
const proposal = this.proposals.get(active.proposalId);
|
|
if (!proposal || proposal.status !== "working") { active.settled = true; return false; }
|
|
active.settled = true;
|
|
const pending: ProposalPending = result.status === "SUCCESS"
|
|
? { kind: "success", summary: result.summary }
|
|
: result.status === "FAILED"
|
|
? { kind: "failure", summary: result.summary }
|
|
: { kind: "step", summary: result.summary, question: result.question, nextStep: result.nextStep };
|
|
await this.proposals.update(active.proposalId, {
|
|
status: "awaiting_user_confirmation",
|
|
pending,
|
|
lastWorkerSummary: result.summary
|
|
});
|
|
if (result.status !== "NEEDS_CONFIRMATION") {
|
|
this.active = undefined;
|
|
await active.worker.terminate().catch(() => undefined);
|
|
await this.proposals.update(active.proposalId, { workerProcessGroup: undefined });
|
|
}
|
|
return true;
|
|
});
|
|
if (transitioned) await this.notifyOwner(active.proposalId);
|
|
}
|
|
|
|
private async handleWorkerRunFailure(active: ActiveWorker, error: unknown): Promise<void> {
|
|
try {
|
|
if (this.shuttingDown || active.settled) return;
|
|
await active.worker.terminate().catch(() => undefined);
|
|
const transitioned = await this.scheduler(async () => {
|
|
if (this.active !== active || active.settled || this.shuttingDown) return false;
|
|
const proposal = this.proposals.get(active.proposalId);
|
|
if (!proposal || proposal.status !== "working") { active.settled = true; return false; }
|
|
active.settled = true;
|
|
this.active = undefined;
|
|
const summary = `worker_error: ${error instanceof Error ? error.message : String(error)}`;
|
|
await this.proposals.update(active.proposalId, {
|
|
status: "awaiting_user_confirmation",
|
|
pending: { kind: "failure", summary },
|
|
lastWorkerSummary: summary,
|
|
workerProcessGroup: undefined
|
|
});
|
|
return true;
|
|
});
|
|
if (transitioned) await this.notifyOwner(active.proposalId);
|
|
} catch (cleanupError) {
|
|
console.error(`Worker failure handling failed (${cleanupError instanceof Error ? cleanupError.name : "unknown error"})`);
|
|
}
|
|
}
|
|
|
|
private onWorkerCrash(proposalId: string, _worker: AcpWorker, error: Error): void {
|
|
this.crashes++;
|
|
console.error(`ACP worker crash: ${error.message}`);
|
|
const active = this.active;
|
|
if (!active || active.proposalId !== proposalId) return;
|
|
void this.handleWorkerRunFailure(active, new Error("worker_crashed: the ACP worker process exited unexpectedly"));
|
|
}
|
|
|
|
private async notifyOwner(proposalId: string): Promise<void> {
|
|
const proposal = this.proposals.get(proposalId);
|
|
const sink = this.eventSink;
|
|
if (!proposal || !sink) return;
|
|
const chatKey = proposal.ownerChatKey;
|
|
const conversationKey = assistantKeyFor(chatKey, proposal.requesterUserId);
|
|
void this.enqueueChat(conversationKey, async () => {
|
|
if (this.shuttingDown) return;
|
|
try {
|
|
const reply = await this.runAssistantEvent(conversationKey, proposal.id);
|
|
if (reply) await sink(chatKey, reply);
|
|
} catch (error) {
|
|
console.error(`Assistant event notification failed (${error instanceof Error ? error.name : "unknown error"})`);
|
|
const latest = this.proposals.get(proposalId);
|
|
if (latest?.pending) await sink(chatKey, fallbackEventText(latest)).catch(() => undefined);
|
|
}
|
|
}).catch(() => undefined);
|
|
}
|
|
|
|
private async runAssistantEvent(conversationKey: string, proposalId: string): Promise<string> {
|
|
const proposal = this.proposals.get(proposalId);
|
|
if (!proposal?.pending) return "";
|
|
const worker = await this.acquireAssistantWorker(conversationKey);
|
|
try {
|
|
const reply = await worker.prompt(assistantEventPrompt(proposal));
|
|
const parsed = await this.parseAssistantReply(worker, reply);
|
|
if (worker.isolationViolated) throw new Error("Assistant session attempted forbidden tool activity");
|
|
await this.store.touchBinding(conversationKey);
|
|
return parsed.reply;
|
|
} catch (error) {
|
|
if (worker.isolationViolated) await this.store.deleteBinding(conversationKey).catch(() => undefined);
|
|
throw error;
|
|
}
|
|
}
|
|
|
|
private async acquireAssistantWorker(conversationKey: string): Promise<AcpWorker> {
|
|
const existing = this.assistantWorkers.get(conversationKey);
|
|
if (existing) {
|
|
existing.lastUsedAt = Date.now();
|
|
return existing;
|
|
}
|
|
await this.reserveAssistantCapacity();
|
|
const persisted = this.store.getBinding(conversationKey);
|
|
const binding = persisted && this.validBinding(persisted) ? persisted : undefined;
|
|
if (persisted && !binding) await this.store.deleteBinding(conversationKey);
|
|
const cwd = this.prepareAssistantWorkspace(conversationKey, binding);
|
|
const spawnWorker = async (nativeSessionId?: string): Promise<AcpWorker> => {
|
|
const worker = new AcpWorker(this.bot, this.config, (crashed, error) => {
|
|
this.crashes++;
|
|
if (this.assistantWorkers.get(conversationKey) === crashed) this.assistantWorkers.delete(conversationKey);
|
|
console.error(`ACP assistant worker crash: ${error.message}`);
|
|
}, {
|
|
kind: "assistant",
|
|
cwd,
|
|
env: this.assistantEnv(),
|
|
policy: ASSISTANT_POLICY,
|
|
onIsolationViolation: (violating) => {
|
|
if (this.assistantWorkers.get(conversationKey) === violating) this.assistantWorkers.delete(conversationKey);
|
|
},
|
|
allowUnverifiedAssistantAgent: this.options.allowUnverifiedAssistantAgent
|
|
});
|
|
await worker.start(nativeSessionId);
|
|
return worker;
|
|
};
|
|
let worker: AcpWorker;
|
|
let resumed = Boolean(binding);
|
|
try {
|
|
worker = await spawnWorker(binding?.nativeSessionId);
|
|
} catch (error) {
|
|
if (!(error instanceof AcpSessionRestoreError) || !binding) throw error;
|
|
await this.store.deleteBinding(conversationKey);
|
|
resumed = false;
|
|
worker = await spawnWorker();
|
|
}
|
|
this.assistantWorkers.set(conversationKey, worker);
|
|
if (!resumed) {
|
|
try {
|
|
await worker.prompt(this.bot.assistantBootstrap, "initializing");
|
|
const now = Date.now();
|
|
await this.store.setBinding({
|
|
chatKey: conversationKey,
|
|
agentId: this.bot.agent.id,
|
|
nativeSessionId: worker.nativeSessionId!,
|
|
assistantWorkspace: cwd,
|
|
botFingerprint: this.bot.fingerprint,
|
|
createdAt: now,
|
|
updatedAt: now
|
|
});
|
|
} catch (error) {
|
|
if (this.assistantWorkers.get(conversationKey) === worker) this.assistantWorkers.delete(conversationKey);
|
|
await worker.terminate().catch(() => undefined);
|
|
throw error;
|
|
}
|
|
}
|
|
return worker;
|
|
}
|
|
|
|
private validBinding(binding: AssistantBinding): boolean {
|
|
return binding.botFingerprint === this.bot.fingerprint && binding.agentId === this.bot.agent.id;
|
|
}
|
|
|
|
private async reserveAssistantCapacity(): Promise<void> {
|
|
while (this.assistantWorkers.size >= this.maxAssistantSessions
|
|
|| this.assistantWorkers.size + (this.active ? 1 : 0) >= this.config.maxProcesses) {
|
|
const candidate = [...this.assistantWorkers.entries()]
|
|
.filter(([, worker]) => !worker.inFlight)
|
|
.sort((left, right) => left[1].lastUsedAt - right[1].lastUsedAt)[0];
|
|
if (!candidate) throw new Error(`Assistant session limit reached (${this.maxAssistantSessions})`);
|
|
this.assistantWorkers.delete(candidate[0]);
|
|
await candidate[1].terminate();
|
|
}
|
|
}
|
|
|
|
private async sweep(): Promise<void> {
|
|
const cutoff = Date.now() - this.config.idleTimeoutMs;
|
|
for (const [chatKey, worker] of [...this.assistantWorkers.entries()]) {
|
|
if (!worker.inFlight && worker.lastUsedAt < cutoff) {
|
|
this.assistantWorkers.delete(chatKey);
|
|
await worker.terminate().catch(() => undefined);
|
|
}
|
|
}
|
|
}
|
|
|
|
private prepareAssistantWorkspace(chatKey: string, binding?: AssistantBinding): string {
|
|
const home = path.resolve(this.options.assistantWorkspaceHome || process.env.GORI_AGENT_HOME || process.cwd());
|
|
const key = crypto.createHmac("sha256", this.assistantSalt).update(chatKey).digest("hex");
|
|
const directory = binding?.assistantWorkspace || path.join(home, "state", "assistant-workspaces", key);
|
|
return prepareAssistantWorkspace(directory, this.bot.workspace);
|
|
}
|
|
|
|
private assistantEnv(): Record<string, string> {
|
|
if (this.bot.agent.env.KIMI_CODE_HOME) return {};
|
|
if (!this.kimiHome) {
|
|
const home = this.options.kimiCodeHome
|
|
|| path.join(path.resolve(this.options.assistantWorkspaceHome || process.env.GORI_AGENT_HOME || process.cwd()), "state", "kimi", "assistant");
|
|
ensurePrivateDirectoryChain(home);
|
|
this.kimiHome = home;
|
|
}
|
|
return { KIMI_CODE_HOME: this.kimiHome };
|
|
}
|
|
|
|
private assistantUserPrompt(request: ConversationRequest): string {
|
|
const chatKey = chatKeyFor(request.platform, request.chatId);
|
|
return [
|
|
"[User message]",
|
|
request.text,
|
|
"",
|
|
"[Proposal states]",
|
|
...this.proposalSummaries(chatKey, request.userId),
|
|
"",
|
|
"[Worker state]",
|
|
this.workerStateSummary(chatKey, request.userId)
|
|
].join("\n");
|
|
}
|
|
|
|
private proposalSummaries(chatKey: string, userId: string): string[] {
|
|
const proposals = this.proposals.list()
|
|
.filter((proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId);
|
|
if (proposals.length === 0) return ["none"];
|
|
return proposals.map((proposal) => {
|
|
const pending = proposal.pending
|
|
? ` pending=${proposal.pending.kind} summary=${JSON.stringify(proposal.pending.summary)}${proposal.pending.question ? ` question=${JSON.stringify(proposal.pending.question)}` : ""}`
|
|
: "";
|
|
return `- id=${proposal.id} status=${proposal.status} title=${JSON.stringify(proposal.title)}${pending}`;
|
|
});
|
|
}
|
|
|
|
private workerStateSummary(chatKey: string, userId: string): string {
|
|
const active = this.active;
|
|
if (!active) return "no active worker";
|
|
const proposal = this.proposals.get(active.proposalId);
|
|
if (!proposal || proposal.ownerChatKey !== chatKey || proposal.requesterUserId !== userId) {
|
|
return "busy: a worker is executing another user's proposal; details hidden";
|
|
}
|
|
const worker = active.worker;
|
|
const runningSeconds = worker.turnStartedAt ? Math.max(0, Math.floor((Date.now() - worker.turnStartedAt) / 1_000)) : 0;
|
|
return `proposal=${active.proposalId} phase=${worker.phase} inFlight=${worker.inFlight} runningSeconds=${runningSeconds}`;
|
|
}
|
|
|
|
private enqueueChat<T>(chatKey: string, operation: () => Promise<T>): Promise<T> {
|
|
const tail = this.chatChains.get(chatKey) || Promise.resolve();
|
|
const result = tail.then(operation, operation);
|
|
const marker = result.then(() => undefined, () => undefined);
|
|
this.chatChains.set(chatKey, marker);
|
|
void marker.then(() => { if (this.chatChains.get(chatKey) === marker) this.chatChains.delete(chatKey); });
|
|
return result;
|
|
}
|
|
|
|
private scheduler<T>(operation: () => Promise<T>): Promise<T> {
|
|
const result = this.schedulerChain.then(operation, operation);
|
|
this.schedulerChain = result.then(() => undefined, () => undefined);
|
|
return result;
|
|
}
|
|
}
|
|
|
|
export function assistantKeyFor(chatKey: string, userId: string): string {
|
|
return `${chatKey}#${userId}`;
|
|
}
|
|
|
|
export function parseAssistantActions(text: string): ParsedAssistantEnvelope | undefined { const match = ASSISTANT_ENVELOPE.exec(text);
|
|
if (!match?.groups) return undefined;
|
|
let value: unknown;
|
|
try { value = JSON.parse(match.groups.json); } catch { return undefined; }
|
|
if (!isRecord(value) || typeof value.reply !== "string" || !Array.isArray(value.actions)) return undefined;
|
|
const actions: AssistantAction[] = [];
|
|
for (const item of value.actions) {
|
|
const action = parseAssistantAction(item);
|
|
if (!action) return undefined;
|
|
actions.push(action);
|
|
}
|
|
return { reply: value.reply, actions };
|
|
}
|
|
|
|
function parseAssistantAction(value: unknown): AssistantAction | undefined {
|
|
if (!isRecord(value) || typeof value.type !== "string") return undefined;
|
|
if (value.type === "create_proposal") {
|
|
if (typeof value.title !== "string" || value.title.length === 0
|
|
|| typeof value.goal !== "string" || value.goal.length === 0
|
|
|| !Array.isArray(value.steps) || value.steps.some((step) => typeof step !== "string" || step.length === 0)) return undefined;
|
|
return { type: "create_proposal", title: value.title, goal: value.goal, steps: value.steps as string[] };
|
|
}
|
|
if (value.type === "confirm") {
|
|
if ((value.id !== undefined && typeof value.id !== "string") || (value.answer !== undefined && typeof value.answer !== "string")) return undefined;
|
|
return { type: "confirm", id: value.id as string | undefined, answer: value.answer as string | undefined };
|
|
}
|
|
if (value.type === "start_next") return { type: "start_next" };
|
|
if (value.type === "cancel") {
|
|
if (value.id !== undefined && typeof value.id !== "string") return undefined;
|
|
return { type: "cancel", id: value.id as string | undefined };
|
|
}
|
|
if (value.type === "stop") return { type: "stop" };
|
|
return undefined;
|
|
}
|
|
|
|
export function parseWorkerResult(text: string): ParsedWorkerResult | undefined {
|
|
const match = WORKER_ENVELOPE.exec(text);
|
|
if (!match?.groups) return undefined;
|
|
let value: unknown;
|
|
try { value = JSON.parse(match.groups.json); } catch { return undefined; }
|
|
if (!isRecord(value) || typeof value.summary !== "string") return undefined;
|
|
if (value.status !== "SUCCESS" && value.status !== "NEEDS_CONFIRMATION" && value.status !== "FAILED") return undefined;
|
|
if ((value.question !== undefined && typeof value.question !== "string")
|
|
|| (value.nextStep !== undefined && typeof value.nextStep !== "string")
|
|
|| (value.dirty !== undefined && typeof value.dirty !== "boolean")) return undefined;
|
|
return {
|
|
status: value.status,
|
|
summary: value.summary,
|
|
question: value.question as string | undefined,
|
|
nextStep: value.nextStep as string | undefined,
|
|
dirty: value.dirty as boolean | undefined
|
|
};
|
|
}
|
|
|
|
function withWorkerResultProtocol(text: string): string {
|
|
return `${text}\n\nWhen this turn is finished, end your response with exactly one hidden worker result envelope. Use <GORI_WORKER_RESULT_V1>{"status":"SUCCESS","summary":"..."}</GORI_WORKER_RESULT_V1> only when the proposal is fully complete, status "FAILED" with a summary when it cannot be completed, or status "NEEDS_CONFIRMATION" with a "question" when user confirmation is required before continuing. Optional fields: "nextStep", "dirty". Do not emit any other status value or any text after the envelope.`;
|
|
}
|
|
|
|
function workerRepairPrompt(): string {
|
|
return "Your previous response did not end with a valid GORI_WORKER_RESULT_V1 envelope. Return exactly one valid envelope with status SUCCESS, NEEDS_CONFIRMATION, or FAILED and a summary. Do not perform more work.";
|
|
}
|
|
|
|
function assistantRepairPrompt(): string {
|
|
return "Your previous response did not end with a valid GORI_ASSISTANT_ACTION_V1 envelope. Reply again with the user-facing text in the JSON \"reply\" field and an \"actions\" array (empty if no state change is needed).";
|
|
}
|
|
|
|
function workerTaskPrompt(proposal: Proposal, resumeContext?: string): string {
|
|
return [
|
|
"Execute this confirmed proposal in the workspace.",
|
|
`Title: ${proposal.title}`,
|
|
`Goal: ${proposal.goal}`,
|
|
"Steps:",
|
|
...proposal.steps.map((step, index) => `${index + 1}. ${step}`),
|
|
...(resumeContext ? ["", "Additional context:", resumeContext] : []),
|
|
"",
|
|
"If you encounter a dirty or unexpected target, or anything that requires the user's decision, stop and report NEEDS_CONFIRMATION with a clear question instead of forcing the change."
|
|
].join("\n");
|
|
}
|
|
|
|
function workerAnswerPrompt(answer: string, question?: string): string {
|
|
return [
|
|
"The user answered your question about the current proposal.",
|
|
`Your question was: ${question || "unknown"}`,
|
|
`User answer: ${answer}`,
|
|
"Continue executing the same proposal accordingly."
|
|
].join("\n");
|
|
}
|
|
|
|
function assistantEventPrompt(proposal: Proposal): string {
|
|
const pending = proposal.pending!;
|
|
return [
|
|
"[Internal event from gori-agent. This is not a user message.]",
|
|
`The worker for proposal "${proposal.title}" (id ${proposal.id}) reported ${pending.kind === "step" ? "NEEDS_CONFIRMATION" : pending.kind === "success" ? "SUCCESS" : "FAILED"}.`,
|
|
`Summary: ${pending.summary}`,
|
|
pending.question ? `Question for the user: ${pending.question}` : "Question for the user: none",
|
|
pending.nextStep ? `Suggested next step: ${pending.nextStep}` : "Suggested next step: none",
|
|
"Tell the user, in your own words, what happened and what they can do next (confirm, answer the question, retry, or stop). End with the usual GORI_ASSISTANT_ACTION_V1 envelope; its actions array MUST be empty because actions are ignored for internal events."
|
|
].join("\n");
|
|
}
|
|
|
|
function fallbackEventText(proposal: Proposal): string {
|
|
const pending = proposal.pending!;
|
|
const state = pending.kind === "step" ? "等待你回答一个问题" : pending.kind === "success" ? "执行完成,待你确认" : "执行失败,待你确认或重试";
|
|
return `任务「${proposal.title}」${state}。摘要:${pending.summary}${pending.question ? ` 问题:${pending.question}` : ""}`;
|
|
}
|
|
|
|
function isRecord(value: unknown): value is Record<string, unknown> {
|
|
return typeof value === "object" && value !== null && !Array.isArray(value);
|
|
}
|
|
|
|
function prepareAssistantWorkspace(directory: string, projectWorkspace: string): string {
|
|
const project = fs.realpathSync.native(projectWorkspace);
|
|
if (pathsOverlap(canonicalProspectivePath(directory), project)) throw new Error("Assistant workspace must not overlap the project workspace");
|
|
ensurePrivateDirectoryChain(directory);
|
|
const workspace = fs.realpathSync.native(directory);
|
|
const agents = ensurePrivateSubdirectories(workspace, [".kimi-code", "agents"]);
|
|
const profile = path.join(agents, "agent.md");
|
|
const temp = path.join(agents, `.agent.${process.pid}.${crypto.randomBytes(8).toString("hex")}.tmp`);
|
|
const content = `---\nname: agent\ndescription: gori-agent no-tool assistant profile\noverride: true\ntools: []\nsubagents: []\n---\nYou are the user-facing Assistant. You have no tools and cannot delegate. Never claim to inspect files, execute commands, call Skills or MCP, or see the Worker's private context. If the information in a prompt is insufficient, say so clearly.\n`;
|
|
let fd: number | undefined;
|
|
try {
|
|
fd = fs.openSync(temp, "wx", 0o600);
|
|
fs.writeFileSync(fd, content, "utf8");
|
|
fs.fsyncSync(fd);
|
|
fs.closeSync(fd);
|
|
fd = undefined;
|
|
fs.renameSync(temp, profile);
|
|
fs.chmodSync(profile, 0o600);
|
|
} catch (error) {
|
|
if (fd !== undefined) fs.closeSync(fd);
|
|
try { fs.unlinkSync(temp); } catch { /* best effort */ }
|
|
throw error;
|
|
}
|
|
return workspace;
|
|
}
|
|
|
|
function pathsOverlap(left: string, right: string): boolean {
|
|
const relativeA = path.relative(left, right);
|
|
const relativeB = path.relative(right, left);
|
|
return relativeA === "" || (!relativeA.startsWith("..") && !path.isAbsolute(relativeA))
|
|
|| (!relativeB.startsWith("..") && !path.isAbsolute(relativeB));
|
|
}
|
|
|
|
function canonicalProspectivePath(target: string): string {
|
|
const suffix: string[] = [];
|
|
let current = path.resolve(target);
|
|
while (!fs.existsSync(current)) {
|
|
const parent = path.dirname(current);
|
|
if (parent === current) throw new Error(`Cannot resolve assistant workspace path: ${target}`);
|
|
suffix.unshift(path.basename(current));
|
|
current = parent;
|
|
}
|
|
return path.join(fs.realpathSync.native(current), ...suffix);
|
|
}
|
|
|
|
function ensurePrivateDirectoryChain(target: string): void {
|
|
const pending: string[] = [];
|
|
let current = path.resolve(target);
|
|
while (!fs.existsSync(current)) {
|
|
const parent = path.dirname(current);
|
|
if (parent === current) throw new Error(`Cannot create assistant workspace path: ${target}`);
|
|
pending.unshift(path.basename(current));
|
|
current = parent;
|
|
}
|
|
assertSafeDirectory(current);
|
|
for (const segment of pending) current = createPrivateSubdirectory(current, segment);
|
|
assertPrivateDirectory(current);
|
|
}
|
|
|
|
function ensurePrivateSubdirectories(root: string, segments: string[]): string {
|
|
let current = root;
|
|
for (const segment of segments) current = createPrivateSubdirectory(current, segment);
|
|
return current;
|
|
}
|
|
|
|
function createPrivateSubdirectory(parent: string, segment: string): string {
|
|
const directory = path.join(parent, segment);
|
|
try { fs.mkdirSync(directory, { mode: 0o700 }); }
|
|
catch (error) { if ((error as NodeJS.ErrnoException).code !== "EEXIST") throw error; }
|
|
assertPrivateDirectory(directory);
|
|
return directory;
|
|
}
|
|
|
|
function assertSafeDirectory(directory: string): fs.Stats {
|
|
const stat = fs.lstatSync(directory);
|
|
if (!stat.isDirectory() || stat.isSymbolicLink()) throw new Error(`Refusing unsafe assistant workspace directory: ${directory}`);
|
|
return stat;
|
|
}
|
|
|
|
function assertPrivateDirectory(directory: string): void {
|
|
const stat = assertSafeDirectory(directory);
|
|
if (typeof process.getuid === "function" && stat.uid !== process.getuid()) throw new Error(`Assistant workspace directory is not owned by the current user: ${directory}`);
|
|
fs.chmodSync(directory, 0o700);
|
|
}
|