Files
gori-agent/src/acp/assistant-manager.ts
T
zenord 4f7b4a9115 feat: show owners a live worker status snapshot
While a worker turn is running, the assistant prompt for the working
proposal's owner includes a compact status card (elapsed time, last
tool activity category with age, per-category counts) built from the
ACP session/update stream the runtime already receives. Other users
still see only the desensitized busy state, and raw update content
never enters the assistant session.
2026-08-19 16:39:42 +08:00

1203 lines
59 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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 { IncomingAttachment, OutgoingImageRef } from "../core/types.js";
import { MAX_OUTBOUND_IMAGES, validateWorkspaceImage } from "../core/workspace-images.js";
import type { ResolvedBot } from "../roles/role-registry.js";
import type { ConversationRequest, ConversationResponse, ConversationRuntime, RuntimeEventSink, RuntimeStats } from "./types.js";
import type { AcpPromptContent, SafeActivityCategory } from "./client.js";
import { AcpSessionRestoreError, AcpWorker } from "./worker.js";
const ASSISTANT_POLICY: PermissionPolicy = { mode: "deny", allowedTools: [], allowedCommandPatterns: [] };
const ASSISTANT_START_TAG = "<GORI_ASSISTANT_ACTION_V2>";
const WORKER_START_TAG = "<GORI_WORKER_RESULT_V2>";
const ASSISTANT_ENVELOPE = /<GORI_ASSISTANT_ACTION_V2>(?<json>[\s\S]*?)<\/GORI_ASSISTANT_ACTION_V2>\s*$/;
const WORKER_ENVELOPE = /<GORI_WORKER_RESULT_V2>(?<json>[\s\S]*?)<\/GORI_WORKER_RESULT_V2>\s*$/;
export type AssistantAction =
| { type: "create_proposal"; title: string; goal: string; steps: string[] }
| { type: "confirm"; id?: string }
| { type: "adjust_proposal"; id?: string; title?: string; goal?: string; steps?: string[] }
| { type: "follow_up"; id?: string; instruction: string }
| { type: "finish"; id?: string; note?: string }
| { type: "send_image"; path: string }
| { type: "start_next" }
| { type: "cancel"; id?: string }
| { type: "stop" };
export interface ParsedAssistantEnvelope {
reply: string;
actions: AssistantAction[];
}
export type WorkerResultStatus = "PENDING";
export interface ParsedWorkerResult {
status: WorkerResultStatus;
summary: string;
question?: string;
workspaceDirty?: boolean;
attachments?: { path: string; mimeType?: string }[];
}
interface ActiveWorker {
proposalId: string;
worker: AcpWorker;
settled: boolean;
run?: Promise<void>;
activity: WorkerActivity;
}
// Real-time snapshot of the active worker's session/update activity; only categories and
// timestamps are kept, never tool arguments, outputs, or message text.
interface WorkerActivity {
turnStartedAt?: number;
lastActivityAt?: number;
lastCategory?: SafeActivityCategory;
counts: Partial<Record<SafeActivityCategory, number>>;
}
const ACTIVITY_CATEGORIES: readonly SafeActivityCategory[] = ["read", "search", "write", "execute", "delegate", "other"];
interface WorkerLaunchOptions {
resumeSessionId?: string;
resumeContext?: string;
firstTurn?: string;
freshSessionNote?: string;
images?: IncomingAttachment[];
}
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 workspace may be dirty";
await this.proposals.update(proposal.id, {
status: "pending",
pending: { summary, workspaceDirty: true, receivedAt: Date.now() },
lastWorkerSummary: summary,
workerProcessGroup: undefined
});
}
}
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 result = await this.enqueueChat(conversationKey, () => this.runAssistantTurn(conversationKey, request));
return { text: result.text, images: result.images, 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 finish(platform: string, chatId: string, userId: string): Promise<boolean> {
const chatKey = chatKeyFor(platform, chatId);
return this.scheduler(() => this.finishLocked(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 chatKey = chatKeyFor(platform, chatId);
const proposals = this.proposals.list();
const owned = (proposal: Proposal) => Boolean(userId) && proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId;
const mine = proposals.filter(owned);
const active = this.active;
const activeProposal = active ? this.proposals.get(active.proposalId) : undefined;
const pending = proposals.filter((proposal) => proposal.status === "pending");
const working = proposals.filter((proposal) => proposal.status === "working");
const myQueued = mine.filter((proposal) => proposal.status === "queued").length;
const myPending = mine.filter((proposal) => proposal.status === "pending").length;
const myProposed = mine.filter((proposal) => proposal.status === "proposed").length;
const dirtyPending = pending.some((proposal) => proposal.pending?.workspaceDirty);
const schedulerState = active || working.length > 0 ? "working" : pending.length > 0 ? "pending" : "idle";
const blockedReason = active || working.length > 0
? (working.some(owned) || (activeProposal && owned(activeProposal)) ? "your proposal is working" : "another proposal is working")
: myPending > 0
? "your proposal is waiting for your decision"
: dirtyPending
? "a pending proposal left a dirty workspace"
: "none";
const nextAction = myPending > 0
? "finish your pending proposal or follow up with new instructions"
: myProposed > 0
? "confirm your proposed proposal to queue it"
: myQueued > 0
? (schedulerState === "idle" ? "start_next" : "wait for the current proposal to settle")
: "none";
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: working.length,
pendingProposals: pending.length,
workerRunning: Boolean(active?.worker.inFlight),
owner: Boolean(userId && activeProposal && owned(activeProposal)),
myQueuedProposals: myQueued,
myPendingProposals: myPending,
schedulerState,
blockedReason,
nextAction
};
}
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<{ text: string; images?: OutgoingImageRef[] }> {
const worker = await this.acquireAssistantWorker(conversationKey);
try {
const reply = await worker.prompt(this.assistantPromptContent(worker, 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);
const chatKey = chatKeyFor(request.platform, request.chatId);
const { corrections, images } = await this.executeActions(chatKey, request, parsed.actions);
const reminders = this.pendingReminders(chatKey, request.userId, parsed.reply);
const extras = [...corrections, ...reminders];
return { text: extras.length > 0 ? `${parsed.reply}\n${extras.join("\n")}` : parsed.reply, images };
} catch (error) {
if (worker.isolationViolated) await this.store.deleteBinding(conversationKey).catch(() => undefined);
throw error;
}
}
private pendingReminders(chatKey: string, userId: string, reply: string): string[] {
return this.proposals.list({ status: "pending" })
.filter((proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId)
.filter((proposal) => !reply.includes(proposal.title))
.map((proposal) => `提醒:你还有一个待确认的任务「${proposal.title}」,说 finish 结束,或直接继续说要求继续改。`);
}
private assistantPromptContent(worker: AcpWorker, request: ConversationRequest): string | AcpPromptContent[] {
const text = this.assistantUserPrompt(request);
const images = request.attachments || [];
if (images.length === 0) return text;
if (!worker.supportsImages) return `${text}\n\n[Attachments]\n${images.map(() => "[图片,当前 agent 不支持图片输入]").join("\n")}`;
return [
...images.map((image) => ({ type: "image" as const, data: image.data, mimeType: image.mimeType })),
{ type: "text" as const, text }
];
}
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_V2 envelope after one repair attempt");
return parsed;
}
private async executeActions(chatKey: string, request: ConversationRequest, actions: AssistantAction[]): Promise<{ corrections: string[]; images?: OutgoingImageRef[] }> {
const corrections: string[] = [];
const images: OutgoingImageRef[] = [];
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") {
const ok = await this.scheduler(() => this.confirmLocked(chatKey, request.userId, action.id));
if (!ok) corrections.push("没有可确认的提案;你只能操作自己发起的任务。");
} else if (action.type === "adjust_proposal") {
const ok = await this.scheduler(() => this.adjustProposalLocked(chatKey, request.userId, action));
if (!ok) corrections.push("没有可调整的提案;只有还没开始执行的任务可以调整,且你只能操作自己发起的任务。");
} else if (action.type === "follow_up") {
const ok = await this.scheduler(() => this.followUpLocked(chatKey, request.userId, action.id, action.instruction, request.attachments));
if (!ok) corrections.push("现在没法继续:只有待确认的任务可以继续,且你只能操作自己发起的任务;如果有别的任务正在执行,也要先等它落定。");
} else if (action.type === "finish") {
const ok = await this.scheduler(() => this.finishLocked(chatKey, request.userId, action.id, action.note));
if (!ok) corrections.push("没有可结束的任务;只有待确认的任务可以 finish,且你只能操作自己发起的任务。");
} else if (action.type === "send_image") {
const image = await validateWorkspaceImage(this.bot.workspace, action.path);
if (image) images.push({ path: image.path, mimeType: image.mimeType, filename: image.filename });
else corrections.push("这张图发不出去:路径不在工作区内,或者不是有效的 png/jpg 图片。");
} else if (action.type === "start_next") {
const started = await this.scheduler(() => this.tryStartLocked(chatKey, request.userId));
if (!started) corrections.push(this.startNextCorrection(chatKey, request.userId));
} else if (action.type === "cancel") {
const ok = await this.scheduler(() => this.cancelLocked(chatKey, request.userId, action.id));
if (!ok) corrections.push("没有可取消的提案;你只能操作自己发起的任务。");
} else if (action.type === "stop") {
const ok = await this.scheduler(() => this.stopActiveLocked(chatKey, request.userId));
if (!ok) corrections.push("没有可停止的任务;你只能操作自己发起的任务。");
}
}
return { corrections, images: images.length > 0 ? images.slice(0, MAX_OUTBOUND_IMAGES) : undefined };
}
private startNextCorrection(chatKey: string, userId: string): string {
const owned = (proposal: Proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId;
const queuedMine = this.proposals.list({ status: "queued" }).some(owned);
const kept = queuedMine ? "你的任务已保留在队列里。" : "";
if (this.active || this.proposals.list({ status: "working" }).length > 0) {
return `还不能开始:还有一个任务正在执行。${kept}`.trim();
}
if (this.proposals.list({ status: "pending" }).some(owned)) {
return `还不能开始:你有任务还在等你确认,先 finish 或直接继续说要求。${kept}`.trim();
}
if (this.proposals.list({ status: "pending" }).some((proposal) => proposal.pending?.workspaceDirty)) {
return `还不能开始:有任务留下了未整理的工作区,需要先处理 dirty 工作区。${kept}`.trim();
}
return "还不能开始:你没有已确认并排队的任务。";
}
private async confirmLocked(chatKey: string, userId: string, id?: string): Promise<boolean> {
const proposal = this.resolveTargetProposal(chatKey, userId, id, ["proposed"]);
if (!proposal) return false;
await this.proposals.update(proposal.id, { status: "queued", confirmedAt: Date.now() });
await this.tryStartLocked(chatKey, userId);
return true;
}
private async adjustProposalLocked(chatKey: string, userId: string, action: { id?: string; title?: string; goal?: string; steps?: string[] }): Promise<boolean> {
const proposal = this.resolveTargetProposal(chatKey, userId, action.id, ["proposed", "queued"]);
if (!proposal) return false;
const patch: { title?: string; goal?: string; steps?: string[] } = {};
if (action.title !== undefined) patch.title = action.title;
if (action.goal !== undefined) patch.goal = action.goal;
if (action.steps !== undefined) patch.steps = action.steps;
if (Object.keys(patch).length === 0) return false;
await this.proposals.update(proposal.id, patch);
return true;
}
private async followUpLocked(chatKey: string, userId: string, id: string | undefined, instruction: string, attachments?: IncomingAttachment[]): Promise<boolean> {
const proposal = this.resolveTargetProposal(chatKey, userId, id, ["pending"]);
if (!proposal) return false;
if (this.active || this.proposals.list({ status: "working" }).length > 0) return false;
await this.proposals.update(proposal.id, { status: "working", pending: undefined, startedAt: Date.now() });
console.log(`Worker follow-up for proposal ${proposalTag(proposal.id)}: resuming session ${proposal.workerNativeSessionId || "none"}`);
this.launchWorker(proposal.id, {
resumeSessionId: proposal.workerNativeSessionId,
firstTurn: workerFollowUpPrompt(instruction, proposal.pending?.question),
freshSessionNote: "The previous worker session could not be resumed; this is a fresh session continuing the same proposal.",
images: attachments
});
return true;
}
private async finishLocked(chatKey: string, userId: string, id?: string, note?: string): Promise<boolean> {
const proposal = this.resolveTargetProposal(chatKey, userId, id, ["pending"]);
if (!proposal) return false;
await this.terminateActiveLocked(proposal.id);
await this.proposals.update(proposal.id, {
status: "finished",
finishKind: "done",
finishNote: note,
pending: undefined,
finishedAt: Date.now()
});
return true;
}
private async cancelLocked(chatKey: string, userId: string, id?: string): Promise<boolean> {
const proposal = this.resolveTargetProposal(chatKey, userId, id, ["proposed", "queued", "pending"]);
if (!proposal) return false;
await this.terminateActiveLocked(proposal.id);
await this.proposals.update(proposal.id, {
status: "finished",
finishKind: "cancelled",
pending: undefined,
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) return false;
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);
const summary = "被用户中止";
await this.proposals.update(activeProposal.id, {
status: "pending",
pending: { summary, workspaceDirty: true, receivedAt: Date.now() },
workerProcessGroup: undefined,
lastWorkerSummary: summary
});
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<boolean> {
if (this.active) return false;
if (this.proposals.list({ status: "working" }).length > 0) return false;
const owned = (proposal: Proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId;
if (this.proposals.list({ status: "pending" }).some(owned)) return false;
if (this.proposals.list({ status: "pending" }).some((proposal) => proposal.pending?.workspaceDirty)) return false;
const next = this.proposals.list({ status: "queued" })
.find((proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId);
if (!next) return false;
await this.proposals.update(next.id, { status: "working", startedAt: Date.now() });
this.launchWorker(next.id);
return true;
}
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, options: WorkerLaunchOptions = {}): void {
const proposal = this.proposals.get(proposalId);
if (!proposal || proposal.status !== "working") return;
const activity: WorkerActivity = { counts: {} };
const active: ActiveWorker = { proposalId, worker: this.spawnProposalWorker(proposalId, activity), settled: false, activity };
this.active = active;
const run = this.startWorkerRun(active, options);
active.run = run;
void run.catch((error) => this.handleWorkerRunFailure(active, error));
}
private spawnProposalWorker(proposalId: string, activity: WorkerActivity): AcpWorker {
return new AcpWorker(this.bot, this.config, (crashed, error) => this.onWorkerCrash(proposalId, crashed, error), {
kind: "worker",
cwd: this.bot.workspace,
onActivity: (category) => {
activity.lastActivityAt = Date.now();
if (!category) return;
activity.lastCategory = category;
activity.counts[category] = (activity.counts[category] || 0) + 1;
},
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 });
}
});
}
private async startWorkerRun(active: ActiveWorker, options: WorkerLaunchOptions): Promise<void> {
let resumed = false;
if (options.resumeSessionId) {
try {
await active.worker.start(options.resumeSessionId);
resumed = true;
} catch (error) {
if (!(error instanceof AcpSessionRestoreError) || this.active !== active || active.settled) throw error;
console.log(`Worker resume failed for proposal ${proposalTag(active.proposalId)}; starting a fresh session`);
active.worker = this.spawnProposalWorker(active.proposalId, active.activity);
await active.worker.start();
options = {
...options,
resumeContext: [options.resumeContext, options.freshSessionNote].filter(Boolean).join("\n"),
firstTurn: options.firstTurn && options.freshSessionNote ? `${options.firstTurn}\n\n${options.freshSessionNote}` : options.firstTurn
};
}
} else {
await active.worker.start();
}
if (this.active !== active || active.settled) throw new Error("Worker was superseded before its session started");
console.log(`Worker started for proposal ${proposalTag(active.proposalId)} (session ${active.worker.nativeSessionId}, ${resumed ? "resumed" : "new"})`);
await this.proposals.update(active.proposalId, { workerNativeSessionId: active.worker.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");
const firstTurn = options.firstTurn || workerTaskPrompt(proposal, options.resumeContext);
await this.runWorkerTurn(active, firstTurn, options.images);
}
private async runWorkerTurn(active: ActiveWorker, promptText: string, images?: IncomingAttachment[]): Promise<void> {
active.activity.turnStartedAt = Date.now();
active.activity.lastActivityAt = undefined;
active.activity.lastCategory = undefined;
active.activity.counts = {};
const reply = await active.worker.prompt(withWorkerResultProtocol(this.workerPromptContent(active.worker, promptText, images)));
if (this.active !== active || active.settled) throw new Error("Worker was superseded during its turn");
let result = parseWorkerResult(reply);
if (!result) {
console.error(`Worker turn for proposal ${proposalTag(active.proposalId)} returned no valid envelope (strict or salvaged); requesting repair`);
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) console.log(`Worker result repair succeeded for proposal ${proposalTag(active.proposalId)}`);
}
if (!result) {
console.error(`Worker result repair failed for proposal ${proposalTag(active.proposalId)}`);
throw new Error("Worker did not return a valid GORI_WORKER_RESULT_V2 envelope after one repair attempt");
}
await this.settleWorkerResult(active, result);
}
private workerPromptContent(worker: AcpWorker, text: string, images?: IncomingAttachment[]): string | AcpPromptContent[] {
if (!images || images.length === 0) return text;
if (!worker.supportsImages) return `${text}\n\n[Attachments]\n${images.map(() => "[图片,当前 agent 不支持图片输入]").join("\n")}`;
return [
...images.map((image) => ({ type: "image" as const, data: image.data, mimeType: image.mimeType })),
{ type: "text" as const, text }
];
}
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;
this.active = undefined;
const attachments: { path: string; mimeType?: string }[] = [];
const droppedAttachments: string[] = [];
for (const ref of result.attachments || []) {
const image = await validateWorkspaceImage(this.bot.workspace, ref.path);
if (image) attachments.push({ path: image.path, mimeType: image.mimeType });
else droppedAttachments.push(ref.path);
}
if (droppedAttachments.length > 0) {
console.error(`Worker reported ${droppedAttachments.length} invalid attachment(s) for proposal ${active.proposalId}; dropped (outside workspace or not a valid png/jpg)`);
}
const pending: ProposalPending = {
summary: result.summary,
question: result.question,
workspaceDirty: result.workspaceDirty,
receivedAt: Date.now(),
attachments: attachments.length > 0 ? attachments : undefined,
droppedAttachments: droppedAttachments.length > 0 ? droppedAttachments : undefined
};
await this.proposals.update(active.proposalId, {
status: "pending",
pending,
lastWorkerSummary: result.summary
});
console.log(`Worker settled proposal ${proposalTag(active.proposalId)} as pending (attachments=${attachments.length} dropped=${droppedAttachments.length})`);
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)}`;
console.error(`Worker error for proposal ${proposalTag(active.proposalId)}: ${error instanceof Error ? error.message : String(error)}`);
await this.proposals.update(active.proposalId, {
status: "pending",
pending: { summary, workspaceDirty: true, receivedAt: Date.now() },
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, pendingImageRefs(this.proposals.get(proposalId)));
} 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), pendingImageRefs(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);
let binding = persisted && this.validBinding(persisted) ? persisted : undefined;
if (persisted && !binding) await this.store.deleteBinding(conversationKey);
if (binding && this.shouldResetAssistantSession(conversationKey, binding)) {
await this.store.deleteBinding(conversationKey);
binding = undefined;
}
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,
lastActiveAt: 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;
}
// Lazy idle reset: a long-inactive assistant session is dropped (fresh session on the next
// message) unless its owner still has an unfinished proposal that needs the old context.
private shouldResetAssistantSession(conversationKey: string, binding: AssistantBinding): boolean {
const resetMs = this.config.assistantSessionResetIdleMs;
if (!resetMs) return false;
const lastActiveAt = binding.lastActiveAt ?? binding.updatedAt;
const idleMs = Date.now() - lastActiveAt;
if (idleMs <= resetMs) return false;
const separator = conversationKey.lastIndexOf("#");
const chatKey = conversationKey.slice(0, separator);
const userId = conversationKey.slice(separator + 1);
const hasUnfinished = this.proposals.list().some((proposal) => proposal.ownerChatKey === chatKey
&& proposal.requesterUserId === userId
&& (proposal.status === "proposed" || proposal.status === "queued" || proposal.status === "working" || proposal.status === "pending"));
if (hasUnfinished) return false;
const keyHash = crypto.createHash("sha256").update(conversationKey).digest("hex").slice(0, 8);
console.log(`Assistant session reset for conversation ${keyHash} after ${Math.floor(idleMs / 1000)}s idle`);
return true;
}
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),
"",
"[Pending proposals]",
...this.pendingCards(chatKey, request.userId),
"",
"[Scheduler state]",
this.schedulerStateSummary(chatKey, request.userId),
"",
"[Worker state]",
this.workerStateSummary(chatKey, request.userId)
].join("\n");
}
private pendingCards(chatKey: string, userId: string): string[] {
const pending = this.proposals.list({ status: "pending" })
.filter((proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId);
if (pending.length === 0) return ["none"];
return pending.map((proposal) => {
const card = proposal.pending!;
return `- id=${proposal.id} title=${JSON.stringify(proposal.title)} summary=${JSON.stringify(card.summary)}${card.question ? ` question=${JSON.stringify(card.question)}` : ""}${card.workspaceDirty ? " workspaceDirty=true" : ""}`;
});
}
private schedulerStateSummary(chatKey: string, userId: string): string {
const owned = (proposal: Proposal) => proposal.ownerChatKey === chatKey && proposal.requesterUserId === userId;
const active = this.active;
const activeProposal = active ? this.proposals.get(active.proposalId) : undefined;
const workingOther = this.proposals.list({ status: "working" }).some((proposal) => !owned(proposal))
|| Boolean(activeProposal && !owned(activeProposal));
if (workingOther) return "busy: another proposal is working";
const dirtyOther = this.proposals.list({ status: "pending" })
.some((proposal) => !owned(proposal) && proposal.pending?.workspaceDirty);
if (dirtyOther) return "blocked: another proposal left a dirty workspace that needs attention";
return "idle";
}
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
? ` 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 activity = active.activity;
const runningMs = activity.turnStartedAt ? Math.max(0, Date.now() - activity.turnStartedAt) : 0;
const running = `${String(Math.floor(runningMs / 60_000)).padStart(2, "0")}:${String(Math.floor((runningMs % 60_000) / 1_000)).padStart(2, "0")}`;
const lastActivity = activity.lastCategory
? `${activity.lastCategory} (${Math.max(0, Math.floor((Date.now() - (activity.lastActivityAt || activity.turnStartedAt || Date.now())) / 1_000))}s ago)`
: "none yet";
const counts = ACTIVITY_CATEGORIES
.filter((category) => activity.counts[category])
.map((category) => `${category}×${activity.counts[category]}`)
.join(" ") || "none yet";
return [
`title=${JSON.stringify(proposal.title)} running=${running} phase=${worker.phase} inFlight=${worker.inFlight}`,
`lastActivity=${lastActivity}; counts: ${counts}`,
"guidance: this is a real-time snapshot of the worker as of this message; when the user asks about progress, paraphrase it in your own words and never invent details beyond this snapshot."
].join("\n");
}
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}`;
}
function proposalTag(proposalId: string): string {
return proposalId.slice(0, 8);
}
export function parseAssistantActions(text: string): ParsedAssistantEnvelope | undefined {
const direct = parseAssistantEnvelopeValue(strictEnvelopeJson(text, ASSISTANT_ENVELOPE));
if (direct) return direct;
const salvaged = parseAssistantEnvelopeValue(salvageEnvelopeJson(text, ASSISTANT_START_TAG));
if (salvaged) console.log("Assistant action envelope salvaged (missing closing tag)");
return salvaged;
}
function parseAssistantEnvelopeValue(json: string | undefined): ParsedAssistantEnvelope | undefined {
if (!json) return undefined;
let value: unknown;
try { value = JSON.parse(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") return undefined;
return { type: "confirm", id: value.id as string | undefined };
}
if (value.type === "adjust_proposal") {
if ((value.id !== undefined && typeof value.id !== "string")
|| (value.title !== undefined && (typeof value.title !== "string" || value.title.length === 0))
|| (value.goal !== undefined && (typeof value.goal !== "string" || value.goal.length === 0))
|| (value.steps !== undefined && (!Array.isArray(value.steps) || value.steps.some((step) => typeof step !== "string" || step.length === 0)))) return undefined;
return {
type: "adjust_proposal",
id: value.id as string | undefined,
title: value.title as string | undefined,
goal: value.goal as string | undefined,
steps: value.steps as string[] | undefined
};
}
if (value.type === "follow_up") {
if (typeof value.instruction !== "string" || value.instruction.length === 0
|| (value.id !== undefined && typeof value.id !== "string")) return undefined;
return { type: "follow_up", id: value.id as string | undefined, instruction: value.instruction };
}
if (value.type === "finish") {
if ((value.id !== undefined && typeof value.id !== "string") || (value.note !== undefined && typeof value.note !== "string")) return undefined;
return { type: "finish", id: value.id as string | undefined, note: value.note as string | undefined };
}
if (value.type === "send_image") {
if (typeof value.path !== "string" || value.path.length === 0) return undefined;
return { type: "send_image", path: value.path };
}
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 direct = parseWorkerResultValue(strictEnvelopeJson(text, WORKER_ENVELOPE));
if (direct) return direct;
const salvaged = parseWorkerResultValue(salvageEnvelopeJson(text, WORKER_START_TAG));
if (salvaged) console.log("Worker result envelope salvaged (missing closing tag)");
return salvaged;
}
function parseWorkerResultValue(json: string | undefined): ParsedWorkerResult | undefined {
if (!json) return undefined;
let value: unknown;
try { value = JSON.parse(json); } catch { return undefined; }
if (!isRecord(value) || typeof value.summary !== "string") return undefined;
if (value.status !== "PENDING") return undefined;
if ((value.question !== undefined && typeof value.question !== "string")
|| (value.workspaceDirty !== undefined && typeof value.workspaceDirty !== "boolean")) return undefined;
let attachments: { path: string; mimeType?: string }[] | undefined;
if (value.attachments !== undefined) {
if (!Array.isArray(value.attachments) || value.attachments.length > MAX_OUTBOUND_IMAGES) return undefined;
attachments = [];
for (const item of value.attachments) {
if (!isRecord(item) || typeof item.path !== "string" || item.path.length === 0
|| (item.mimeType !== undefined && typeof item.mimeType !== "string")) return undefined;
attachments.push({ path: item.path, mimeType: item.mimeType as string | undefined });
}
}
return {
status: "PENDING",
summary: value.summary,
question: value.question as string | undefined,
workspaceDirty: value.workspaceDirty as boolean | undefined,
attachments
};
}
function strictEnvelopeJson(text: string, envelope: RegExp): string | undefined {
return envelope.exec(text)?.groups?.json;
}
// Fallback for model responses truncated at the envelope's closing tag: brace-balance the JSON
// after the start tag (string- and escape-aware) and accept it only when the balanced object
// ends exactly at the end of the text. Anything unbalanced or with trailing content stays invalid.
function salvageEnvelopeJson(text: string, startTag: string): string | undefined {
const start = text.lastIndexOf(startTag);
if (start < 0) return undefined;
const after = text.slice(start + startTag.length);
if (!after.startsWith("{")) return undefined;
const end = balancedJsonObjectEnd(after);
if (end === undefined) return undefined;
if (after.slice(end).trim().length > 0) return undefined;
return after.slice(0, end);
}
function balancedJsonObjectEnd(text: string): number | undefined {
let depth = 0;
let inString = false;
let escaped = false;
for (let index = 0; index < text.length; index++) {
const char = text[index]!;
if (inString) {
if (escaped) escaped = false;
else if (char === "\\") escaped = true;
else if (char === '"') inString = false;
continue;
}
if (char === '"') inString = true;
else if (char === "{") depth++;
else if (char === "}") {
depth--;
if (depth === 0) return index + 1;
if (depth < 0) return undefined;
}
}
return undefined;
}
function withWorkerResultProtocol(input: string | AcpPromptContent[]): string | AcpPromptContent[] {
const instruction = "When this turn is finished, end your response with exactly one hidden worker result envelope: <GORI_WORKER_RESULT_V2>{\"status\":\"PENDING\",\"summary\":\"...\"}</GORI_WORKER_RESULT_V2>. The summary is a short user-readable description of what you did or what is blocking you. Add a \"question\" when you need the user's decision before continuing, and set \"workspaceDirty\": true when you left the workspace modified or are unsure about its state. When you produced image files the user should see (png/jpg only), save them inside the workspace (prefer .gori-outbox/) and report up to 3 of them as \"attachments\": [{\"path\":\"relative/or/absolute/path\"}]; never report paths outside the workspace. PENDING is the only status; do not emit any other status value or any text after the envelope.";
if (typeof input === "string") return `${input}\n\n${instruction}`;
return [...input, { type: "text" as const, text: `\n\n${instruction}` }];
}
function workerRepairPrompt(): string {
return "Your previous response did not end with a valid GORI_WORKER_RESULT_V2 envelope. Return exactly one valid envelope with status PENDING and a summary. Do not perform more work.";
}
function assistantRepairPrompt(): string {
return "Your previous response did not end with a valid GORI_ASSISTANT_ACTION_V2 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 PENDING with a clear question instead of forcing the change."
].join("\n");
}
function workerFollowUpPrompt(instruction: string, question?: string): string {
return [
"The user sent a follow-up instruction for the current proposal.",
question ? `Your last question was: ${question}` : "You had no open question.",
`User instruction: ${instruction}`,
"Continue executing the same proposal accordingly."
].join("\n");
}
function assistantEventPrompt(proposal: Proposal): string {
const pending = proposal.pending!;
const attachments = pending.attachments?.length
? `Image attachments produced by the worker (already sent to the user after your message): ${pending.attachments.map((attachment) => attachment.path).join(", ")}`
: "Image attachments: none";
const dropped = pending.droppedAttachments?.length
? `Note: ${pending.droppedAttachments.length} reported attachment(s) were invalid (outside the workspace or not png/jpg) and were dropped; mention this to the user.`
: "";
return [
"[Internal event from gori-agent. This is not a user message.]",
`The worker for proposal "${proposal.title}" (id ${proposal.id}) reported its result and the proposal is now waiting for the user's decision.`,
`Summary: ${pending.summary}`,
pending.question ? `Question for the user: ${pending.question}` : "Question for the user: none",
attachments,
...(dropped ? [dropped] : []),
"Tell the user, in your own words, what happened and what they can do next (finish to close the task, follow up with new instructions to continue it, or stop). End with the usual GORI_ASSISTANT_ACTION_V2 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 attachments = pending.attachments?.length ? `(附 ${pending.attachments.length} 张图片)` : "";
const dropped = pending.droppedAttachments?.length ? `(${pending.droppedAttachments.length} 个附件校验失败被丢弃)` : "";
return `任务「${proposal.title}」等你确认。摘要:${pending.summary}${pending.question ? ` 问题:${pending.question}` : ""}${attachments}${dropped} 说 finish 结束,或直接说要求继续改。`;
}
function pendingImageRefs(proposal: Proposal | undefined): OutgoingImageRef[] | undefined {
const attachments = proposal?.pending?.attachments;
if (!attachments || attachments.length === 0) return undefined;
return attachments.map((attachment) => ({
path: attachment.path,
mimeType: attachment.mimeType,
filename: attachment.path.split("/").pop()
}));
}
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);
}