From bd96595ba8c93fbd3dfaf1121c8fc997ac243961 Mon Sep 17 00:00:00 2001 From: zenord Date: Wed, 19 Aug 2026 16:17:54 +0800 Subject: [PATCH] feat: reset idle assistant sessions after one hour When a conversation's assistant session has been idle longer than runtime.acp.assistantSessionResetIdleMs (default 1h, 0 disables) and its owner has no unfinished proposals, the next inbound message starts a fresh session instead of resuming. Binding lastActiveAt is persisted per turn so the decision survives restarts. --- AGENTS.md | 3 +- README.md | 1 + config.example.json | 1 + src/acp/assistant-manager.ts | 29 +++++++++- src/config.ts | 1 + src/core/durable-session-store.ts | 6 +- test/assistant-manager.test.ts | 95 +++++++++++++++++++++++++++++-- 7 files changed, 128 insertions(+), 8 deletions(-) diff --git a/AGENTS.md b/AGENTS.md index aff5831..0ba35f5 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -92,7 +92,7 @@ gori-agent list - envelope 解析先严格匹配(闭合标签 + 文本末尾 anchor);缺闭合标签时从起始标签后做花括号配平(字符串/转义感知)salvage,配平点必须在文本末尾才接受,否则仍判无效走修复。worker 启动(new/resume session)、envelope 无效、repair 成败、settle(attachments/dropped 数)、worker_error、follow_up resume/fresh 兜底均有单行安全日志(proposal id 前 8 位,不含正文)。 - capacity(maxAssistantSessions/maxProcesses)、idle sweep、cancel/confirm/finish/stop、worker_lost 恢复、owner 事件通知。 - `src/core/durable-session-store.ts` - - state v3 只保存 bot/platform identity 和 assistant binding(conversation key,即 chat + user、agent ID、native session ID、assistant workspace、fingerprint、时间戳);不保存消息正文。 + - state v3 只保存 bot/platform identity 和 assistant binding(conversation key,即 chat + user、agent ID、native session ID、assistant workspace、fingerprint、时间戳、`lastActiveAt`);不保存消息正文。 - 单 writer lock、串行持久化、临时文件、fsync、原子 rename;version/identity 不匹配(含旧 v1/v2)拒绝并保留原文件。 - `src/core/proposal-store.ts` - proposals.json(version 2)保存 Proposal:title/goal/steps、owner chat、发起用户 requesterUserId、状态流(proposed/queued/working/pending/finished)、pending(summary/question?/workspaceDirty?/receivedAt)、finished(finishKind done|cancelled、finishNote?)、worker session 与进程组 PGID/token;同样的 lock 与原子写入纪律。打开 v1 文件时先写 `proposals.json.v1-.bak`(0600)备份再按固定映射迁移。 @@ -198,6 +198,7 @@ QQ 群 chat ID 为 `group:`;用户 ID 按 QQ 官方语义取值 - `promptTimeoutMs`: 默认 14400000(4 小时)。 - `cancelGraceMs`: 默认 5000。 - `idleTimeoutMs`: 默认 1800000。 +- `assistantSessionResetIdleMs`: 默认 3600000;`0` 禁用。Assistant 会话(按 binding 的 `lastActiveAt`,每轮刷新并持久化)空闲超过该阈值且该 owner(chat + user)没有任何 `proposed/queued/working/pending` Proposal 时,下一条消息删除 binding 开新 session(新 HMAC workspace 目录),而不是 resume 旧 session;有未完成 Proposal 则照常 resume。重置只在下一条入站消息时 lazy 发生,sweeper 不主动删 binding。 - `sweepIntervalMs`: 默认 60000。 - `maxProcesses`: 默认 8,包含 assistant + worker。 - `maxAssistantSessions`: 默认 4。 diff --git a/README.md b/README.md index ce5f72c..36cddd6 100644 --- a/README.md +++ b/README.md @@ -184,6 +184,7 @@ QQ 出站图片:Worker/Assistant 报告的 workspace 内图片(png/jpg、≤ - initialize:10 秒 - cancel grace:5 秒 - idle assistant:30 分钟 +- assistant session reset idle:1 小时(`assistantSessionResetIdleMs`,`0` 禁用;Assistant 会话空闲超过该阈值且 owner 没有未完成 Proposal 时,下一条消息开新 session 而不是 resume) - sweep:60 秒 - max processes:8(Assistant + Worker 总和) - max assistant sessions:4 diff --git a/config.example.json b/config.example.json index 2c9ad31..36770bf 100644 --- a/config.example.json +++ b/config.example.json @@ -48,6 +48,7 @@ "promptTimeoutMs": 14400000, "cancelGraceMs": 5000, "idleTimeoutMs": 1800000, + "assistantSessionResetIdleMs": 3600000, "sweepIntervalMs": 60000, "maxProcesses": 8, "maxAssistantSessions": 4 diff --git a/src/acp/assistant-manager.ts b/src/acp/assistant-manager.ts index 4ae63a8..9db5b7e 100644 --- a/src/acp/assistant-manager.ts +++ b/src/acp/assistant-manager.ts @@ -645,8 +645,12 @@ export class AssistantManager implements ConversationRuntime { } await this.reserveAssistantCapacity(); const persisted = this.store.getBinding(conversationKey); - const binding = persisted && this.validBinding(persisted) ? persisted : undefined; + 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 => { const worker = new AcpWorker(this.bot, this.config, (crashed, error) => { @@ -688,7 +692,8 @@ export class AssistantManager implements ConversationRuntime { assistantWorkspace: cwd, botFingerprint: this.bot.fingerprint, createdAt: now, - updatedAt: now + updatedAt: now, + lastActiveAt: now }); } catch (error) { if (this.assistantWorkers.get(conversationKey) === worker) this.assistantWorkers.delete(conversationKey); @@ -703,6 +708,26 @@ export class AssistantManager implements ConversationRuntime { 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 { while (this.assistantWorkers.size >= this.maxAssistantSessions || this.assistantWorkers.size + (this.active ? 1 : 0) >= this.config.maxProcesses) { diff --git a/src/config.ts b/src/config.ts index d6dc56f..9fd7f19 100644 --- a/src/config.ts +++ b/src/config.ts @@ -87,6 +87,7 @@ const acpSchema = z.object({ promptTimeoutMs: z.number().int().positive().default(14_400_000), cancelGraceMs: z.number().int().positive().default(5_000), idleTimeoutMs: z.number().int().positive().default(1_800_000), + assistantSessionResetIdleMs: z.number().int().nonnegative().default(3_600_000), sweepIntervalMs: z.number().int().positive().default(60_000), maxProcesses: z.number().int().positive().default(8), maxAssistantSessions: z.number().int().positive().default(4) diff --git a/src/core/durable-session-store.ts b/src/core/durable-session-store.ts index 8b4acbc..6355b01 100644 --- a/src/core/durable-session-store.ts +++ b/src/core/durable-session-store.ts @@ -14,6 +14,8 @@ export interface AssistantBinding { botFingerprint: string; createdAt: number; updatedAt: number; + // Refreshed on every assistant turn; drives the idle reset of long-inactive sessions. + lastActiveAt?: number; } interface StoreData { @@ -90,6 +92,7 @@ export class DurableSessionStore { const binding = this.data.bindings[chatKey]; if (!binding) return; binding.updatedAt = Date.now(); + binding.lastActiveAt = binding.updatedAt; await this.persist(); } @@ -148,7 +151,8 @@ function parseStoreData(value: unknown): StoreData { || typeof binding.assistantWorkspace !== "string" || typeof binding.botFingerprint !== "string" || typeof binding.createdAt !== "number" - || typeof binding.updatedAt !== "number") { + || typeof binding.updatedAt !== "number" + || (binding.lastActiveAt !== undefined && typeof binding.lastActiveAt !== "number")) { throw new Error(`invalid state v3 binding '${chatKey}'`); } } diff --git a/test/assistant-manager.test.ts b/test/assistant-manager.test.ts index a3cfc6c..4d693bf 100644 --- a/test/assistant-manager.test.ts +++ b/test/assistant-manager.test.ts @@ -26,7 +26,7 @@ interface Harness { logFile: string; } -async function createHarness(options: { initialize?: boolean; agentEnv?: Record } = {}): Promise { +async function createHarness(options: { initialize?: boolean; agentEnv?: Record; acp?: Record } = {}): Promise { const home = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-assistant-home-")); const workspace = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-assistant-ws-")); const logFile = path.join(home, "fake-acp.log"); @@ -34,7 +34,7 @@ async function createHarness(options: { initialize?: boolean; agentEnv?: Record< configVersion: 3, bot: { id: "test-bot", workspace, persona: "", agent: { id: "fake", command: process.execPath, args: [fixture], env: { FAKE_ACP_LOG: logFile, ...options.agentEnv } }, permissions: { mode: "deny" } }, gateway: { platform: { type: "webhook", secret: "secret" } }, - runtime: { acp: { promptTimeoutMs: 10_000, cancelGraceMs: 200 } } + runtime: { acp: { promptTimeoutMs: 10_000, cancelGraceMs: 200, ...options.acp } } }); const identity = { botId: "test-bot", platform: "webhook" }; const store = new DurableSessionStore(path.join(home, "state", "acp-sessions.json"), identity); @@ -51,7 +51,7 @@ async function createHarness(options: { initialize?: boolean; agentEnv?: Record< return { home, workspace, store, proposals, manager, events, logFile }; } -async function reopenManager(harness: Harness): Promise { +async function reopenManager(harness: Harness, options: { acp?: Record } = {}): Promise { await harness.manager.shutdown().catch(() => undefined); await harness.proposals.close().catch(() => undefined); await harness.store.close().catch(() => undefined); @@ -64,7 +64,7 @@ async function reopenManager(harness: Harness): Promise { configVersion: 3, bot: { id: "test-bot", workspace: harness.workspace, persona: "", agent: { id: "fake", command: process.execPath, args: [fixture], env: { FAKE_ACP_LOG: harness.logFile } }, permissions: { mode: "deny" } }, gateway: { platform: { type: "webhook", secret: "secret" } }, - runtime: { acp: { promptTimeoutMs: 10_000, cancelGraceMs: 200 } } + runtime: { acp: { promptTimeoutMs: 10_000, cancelGraceMs: 200, ...options.acp } } }); harness.manager = new AssistantManager(config.runtime.acp, new BotProfileResolver(config).bot, harness.store, harness.proposals, { assistantWorkspaceHome: harness.home, @@ -722,6 +722,93 @@ test("image attachments degrade to a text note when the agent has no image capab } }); +test("an assistant session idle beyond the reset threshold is dropped when the owner has no unfinished proposal", async () => { + const harness = await createHarness(); + try { + await harness.manager.prompt(request("hello")); + const first = harness.store.getBinding(CONVERSATION_KEY)!; + assert.ok(first.lastActiveAt); + await harness.store.setBinding({ ...first, lastActiveAt: Date.now() - 3_700_000 }); + await reopenManager(harness); + + const logs: string[] = []; + const originalLog = console.log; + console.log = (...args: unknown[]) => { logs.push(args.map(String).join(" ")); }; + try { + await harness.manager.prompt(request("hello again")); + } finally { + console.log = originalLog; + } + const reset = harness.store.getBinding(CONVERSATION_KEY)!; + assert.notEqual(reset.nativeSessionId, first.nativeSessionId); + assert.ok(logs.some((line) => /Assistant session reset for conversation [0-9a-f]{8} after \d+s idle/.test(line)), logs.join("\n")); + const restores = readLog(harness.logFile).filter((entry) => (entry.method === "session/resume" || entry.method === "session/load") + && entry.sessionId === first.nativeSessionId); + assert.equal(restores.length, 0); + } finally { + await closeHarness(harness); + } +}); + +test("an idle assistant session is resumed while the owner has an unfinished proposal", async () => { + const harness = await createHarness(); + try { + // Proposed (never confirmed) is enough to keep the session. + await harness.manager.prompt(request("create proposal: succeed")); + const first = harness.store.getBinding(CONVERSATION_KEY)!; + await harness.store.setBinding({ ...first, lastActiveAt: Date.now() - 3_700_000 }); + await reopenManager(harness); + await harness.manager.prompt(request("hello")); + assert.equal(harness.store.getBinding(CONVERSATION_KEY)!.nativeSessionId, first.nativeSessionId); + + // Pending keeps it too. + const proposal = harness.proposals.list()[0]!; + await harness.manager.prompt(request("confirm")); + await waitFor(() => harness.proposals.get(proposal.id)!.status === "pending"); + const second = harness.store.getBinding(CONVERSATION_KEY)!; + await harness.store.setBinding({ ...second, lastActiveAt: Date.now() - 3_700_000 }); + await reopenManager(harness); + await harness.manager.prompt(request("hello again")); + assert.equal(harness.store.getBinding(CONVERSATION_KEY)!.nativeSessionId, first.nativeSessionId); + const restores = readLog(harness.logFile).filter((entry) => (entry.method === "session/resume" || entry.method === "session/load") + && entry.sessionId === first.nativeSessionId); + assert.ok(restores.length > 0); + } finally { + await closeHarness(harness); + } +}); + +test("a recently active assistant session resumes and lastActiveAt persists across restarts", async () => { + const harness = await createHarness(); + try { + await harness.manager.prompt(request("hello")); + const first = harness.store.getBinding(CONVERSATION_KEY)!; + await reopenManager(harness); + await harness.manager.prompt(request("hello again")); + const touched = harness.store.getBinding(CONVERSATION_KEY)!; + assert.equal(touched.nativeSessionId, first.nativeSessionId); + assert.ok(touched.lastActiveAt! >= first.lastActiveAt!); + await reopenManager(harness); + assert.equal(harness.store.getBinding(CONVERSATION_KEY)!.lastActiveAt, touched.lastActiveAt); + } finally { + await closeHarness(harness); + } +}); + +test("assistantSessionResetIdleMs 0 disables the idle reset", async () => { + const harness = await createHarness({ acp: { assistantSessionResetIdleMs: 0 } }); + try { + await harness.manager.prompt(request("hello")); + const first = harness.store.getBinding(CONVERSATION_KEY)!; + await harness.store.setBinding({ ...first, lastActiveAt: Date.now() - 86_400_000 }); + await reopenManager(harness, { acp: { assistantSessionResetIdleMs: 0 } }); + await harness.manager.prompt(request("hello again")); + assert.equal(harness.store.getBinding(CONVERSATION_KEY)!.nativeSessionId, first.nativeSessionId); + } finally { + await closeHarness(harness); + } +}); + test("a worker envelope truncated at the closing tag is salvaged without a repair round", async () => { const harness = await createHarness(); try {