From 4f7b4a9115b6015f182b6f83022b7c298b630465 Mon Sep 17 00:00:00 2001 From: zenord Date: Wed, 19 Aug 2026 16:39:42 +0800 Subject: [PATCH] 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. --- AGENTS.md | 1 + README.md | 1 + src/acp/assistant-manager.ts | 48 ++++++++++++++++++++++++++++---- test/assistant-manager.test.ts | 47 +++++++++++++++++++++++++++++++ test/fixtures/fake-acp-agent.mjs | 21 ++++++++++++++ 5 files changed, 112 insertions(+), 6 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index 0ba35f5..ba66e35 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -220,6 +220,7 @@ Assistant 每条回复以隐藏 `GORI_ASSISTANT_ACTION_V2` envelope 结尾(`re - 阻塞规则:任何 `working` 全局拒绝 start_next/follow_up;owner 自己有 pending 时拒绝该 owner 的 start_next(提示先 finish 或 follow_up);任何 `workspaceDirty: true` 的 pending 全局拒绝 start_next(提示先处理 dirty 工作区);非 dirty 的 pending 不挡其他用户。 - `stop` 只对 working 生效:停 worker(现有进程组清理)后 Proposal 转为 pending(summary「被用户中止」、workspaceDirty)。 - Assistant 有 pending 时每轮 prompt 注入 pending card;回复没提到 pending 时 runtime 在 reply 末尾追加人话兜底提醒。 +- Worker working 期间 runtime 按类别聚合 worker 的 `session/update` 工具活动(turn 开始重置计数;只留类别/时间,不留参数与正文);owner 的 Assistant prompt 注入紧凑状态卡(title、mm:ss 时长、最后活动类别与距今、各类别计数、快照指引),非 owner 只给脱敏 busy。 - Gateway 无固定时间提醒;worker 落定后由 Assistant 生成事件文案发给 owner chat,QQ 无主动消息权限时优先使用被动回复窗口(群 5 分钟/私聊 60 分钟,按 4.5/55 分钟保守判定),过期或失败则不主动发、记录 lastEventDelivery 并在下次入站时补发;`/status` 输出当前 chat 的事件投递状态与 schedulerState/blockedReason/nextAction。 Assistant 隔离:cwd 在实例私有 `state/assistant-workspaces/`,目录 key 使用进程随机 salt 的 HMAC,不得与项目 workspace 重叠;Kimi 项目级 agent override 必须设置 `tools: []`、`subagents: []`,ACP `mcpServers: []`,未显式配置时 `KIMI_CODE_HOME` 指向实例私有 `state/kimi/assistant/`;任何 tool update 或 permission request 都视为隔离违约,fail closed、终止进程组并删除 binding。每个 ACP worker 独立进程组 + 随机 token;runner 重启时 working Proposal 校验 token 清理旧进程组后标记为 pending(worker_lost、workspaceDirty: true)。 diff --git a/README.md b/README.md index 36cddd6..42ebc29 100644 --- a/README.md +++ b/README.md @@ -206,6 +206,7 @@ QQ 出站图片:Worker/Assistant 报告的 workspace 内图片(png/jpg、≤ - 阻塞规则:任何 `working` 全局拒绝 `start_next`/`follow_up`;owner 自己有 `pending` 时拒绝该 owner 的 `start_next`(提示先 finish 或 follow_up);任何 `workspaceDirty` 的 `pending` 全局拒绝 `start_next`(提示先处理 dirty 工作区)。非 dirty 的 pending 不挡其他用户。 - `stop` 只对 `working` 生效:停掉 Worker(含进程组清理)并把 Proposal 标记为 `pending`(summary 为「被用户中止」、`workspaceDirty: true`)。 - Assistant 有 pending 时每轮 prompt 注入 pending card(id/title/summary/question);Assistant 回复没提到 pending 时,runtime 在回复末尾追加一条人话兜底提醒。 +- Worker working 期间,runtime 把 `session/update` 的工具活动按类别(read/search/write/execute/delegate/other)聚合成实时快照(不含工具参数、输出或正文);owner 的下一轮 Assistant prompt 注入紧凑状态卡(title、运行时长、最后活动类别、各类别计数),Assistant 据此用人话描述进度;非 owner 仍只见脱敏 busy。 - Gateway 不再有 15/60/180/480 秒的固定时间提醒;Worker 落定结果时通过 Assistant 生成一条事件说明,由平台 adapter 作为**新消息**发给 owner chat(QQ 同样是新消息,不引用原消息)。 Assistant 隔离与旧 side Session 一致且更严格: diff --git a/src/acp/assistant-manager.ts b/src/acp/assistant-manager.ts index 9db5b7e..99852fb 100644 --- a/src/acp/assistant-manager.ts +++ b/src/acp/assistant-manager.ts @@ -8,7 +8,7 @@ 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 } from "./client.js"; +import type { AcpPromptContent, SafeActivityCategory } from "./client.js"; import { AcpSessionRestoreError, AcpWorker } from "./worker.js"; const ASSISTANT_POLICY: PermissionPolicy = { mode: "deny", allowedTools: [], allowedCommandPatterns: [] }; @@ -48,8 +48,20 @@ interface ActiveWorker { worker: AcpWorker; settled: boolean; run?: Promise; + 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>; +} + +const ACTIVITY_CATEGORIES: readonly SafeActivityCategory[] = ["read", "search", "write", "execute", "delegate", "other"]; + interface WorkerLaunchOptions { resumeSessionId?: string; resumeContext?: string; @@ -453,17 +465,24 @@ export class AssistantManager implements ConversationRuntime { private launchWorker(proposalId: string, options: WorkerLaunchOptions = {}): void { const proposal = this.proposals.get(proposalId); if (!proposal || proposal.status !== "working") return; - const active: ActiveWorker = { proposalId, worker: this.spawnProposalWorker(proposalId), settled: false }; + 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): AcpWorker { + 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 }); @@ -480,7 +499,7 @@ export class AssistantManager implements ConversationRuntime { } 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.worker = this.spawnProposalWorker(active.proposalId, active.activity); await active.worker.start(); options = { ...options, @@ -503,6 +522,10 @@ export class AssistantManager implements ConversationRuntime { } private async runWorkerTurn(active: ActiveWorker, promptText: string, images?: IncomingAttachment[]): Promise { + 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); @@ -831,8 +854,21 @@ export class AssistantManager implements ConversationRuntime { 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}`; + 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(chatKey: string, operation: () => Promise): Promise { diff --git a/test/assistant-manager.test.ts b/test/assistant-manager.test.ts index 4d693bf..2d2eaa9 100644 --- a/test/assistant-manager.test.ts +++ b/test/assistant-manager.test.ts @@ -722,6 +722,53 @@ test("image attachments degrade to a text note when the agent has no image capab } }); +test("the owner sees a live worker activity card while others get only a sanitized busy state", async () => { + const harness = await createHarness(); + try { + await harness.manager.prompt(request("create proposal: toolhang", "user-a")); + const proposal = harness.proposals.list()[0]!; + await harness.manager.prompt(request("confirm", "user-a")); + const internals = harness.manager as unknown as { + active?: { activity: { turnStartedAt?: number; lastActivityAt?: number; lastCategory?: string; counts: Record } }; + }; + await waitFor(() => (internals.active?.activity.counts.read || 0) === 2); + + const activity = internals.active!.activity; + assert.equal(activity.lastCategory, "execute"); + assert.equal(activity.counts.read, 2); + assert.equal(activity.counts.execute, 1); + assert.ok(activity.turnStartedAt); + assert.ok(activity.lastActivityAt); + + await harness.manager.prompt(request("进度如何", "user-a")); + await harness.manager.prompt(request("进度如何", "user-b")); + const progressPrompts = readLog(harness.logFile).filter((entry) => entry.method === "session/prompt" + && entry.text?.startsWith("[User message]\n进度如何")); + assert.equal(progressPrompts.length, 2); + const ownerPrompt = progressPrompts[0]!.text!; + assert.match(ownerPrompt, /title="Test proposal" running=\d{2}:\d{2}/); + assert.match(ownerPrompt, /lastActivity=execute \(\d+s ago\)/); + assert.match(ownerPrompt, /counts: read×2 execute×1/); + assert.match(ownerPrompt, /real-time snapshot of the worker/); + const otherPrompt = progressPrompts[1]!.text!; + assert.match(otherPrompt, /details hidden/); + assert.doesNotMatch(otherPrompt, /lastActivity=/); + assert.doesNotMatch(otherPrompt, /Test proposal/); + + // Once the worker is no longer working, the card disappears. + assert.equal(await harness.manager.stop("webhook", "chat-1", "user-a"), true); + assert.equal(harness.proposals.get(proposal.id)!.status, "pending"); + await harness.manager.prompt(request("hello", "user-a")); + const idlePrompts = readLog(harness.logFile).filter((entry) => entry.method === "session/prompt" + && entry.text?.startsWith("[User message]\nhello")); + assert.equal(idlePrompts.length, 1); + assert.match(idlePrompts[0]!.text!, /no active worker/); + assert.doesNotMatch(idlePrompts[0]!.text!, /lastActivity=/); + } finally { + await closeHarness(harness); + } +}); + test("an assistant session idle beyond the reset threshold is dropped when the owner has no unfinished proposal", async () => { const harness = await createHarness(); try { diff --git a/test/fixtures/fake-acp-agent.mjs b/test/fixtures/fake-acp-agent.mjs index d3b5158..6ccc05e 100644 --- a/test/fixtures/fake-acp-agent.mjs +++ b/test/fixtures/fake-acp-agent.mjs @@ -107,6 +107,27 @@ const app = acp.agent({ name: "fake-acp-agent" }) } if (text.startsWith("Execute this confirmed proposal")) { const goal = ((text.split("Goal: ")[1] || "").split("\n")[0] || "").trim(); + if (goal.includes("toolhang")) { + await client.notify(acp.methods.client.session.update, { + sessionId: params.sessionId, + update: { sessionUpdate: "tool_call", toolCallId: "read-1", title: "Read package.json", kind: "read", status: "in_progress" } + }); + await client.notify(acp.methods.client.session.update, { + sessionId: params.sessionId, + update: { sessionUpdate: "tool_call", toolCallId: "read-2", title: "Read README.md", kind: "read", status: "in_progress" } + }); + await client.notify(acp.methods.client.session.update, { + sessionId: params.sessionId, + update: { sessionUpdate: "tool_call", toolCallId: "exec-1", title: "bash npm test", kind: "execute", status: "in_progress" } + }); + await new Promise((resolve) => { + const done = () => resolve(undefined); + pending.set(params.sessionId, done); + signal.addEventListener("abort", done, { once: true }); + }); + pending.delete(params.sessionId); + return { stopReason: "cancelled" }; + } if (goal.includes("hang")) { await new Promise((resolve) => { const done = () => resolve(undefined);