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.
This commit is contained in:
zenord
2026-08-19 16:17:54 +08:00
parent 9c3139380d
commit bd96595ba8
7 changed files with 128 additions and 8 deletions
+2 -1
View File
@@ -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 位,不含正文)。 - 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 事件通知。 - capacity(maxAssistantSessions/maxProcesses)、idle sweep、cancel/confirm/finish/stop、worker_lost 恢复、owner 事件通知。
- `src/core/durable-session-store.ts` - `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)拒绝并保留原文件。 - 单 writer lock、串行持久化、临时文件、fsync、原子 rename;version/identity 不匹配(含旧 v1/v2)拒绝并保留原文件。
- `src/core/proposal-store.ts` - `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-<timestamp>.bak`(0600)备份再按固定映射迁移。 - 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-<timestamp>.bak`(0600)备份再按固定映射迁移。
@@ -198,6 +198,7 @@ QQ 群 chat ID 为 `group:<group_openid>`;用户 ID 按 QQ 官方语义取值
- `promptTimeoutMs`: 默认 14400000(4 小时)。 - `promptTimeoutMs`: 默认 14400000(4 小时)。
- `cancelGraceMs`: 默认 5000。 - `cancelGraceMs`: 默认 5000。
- `idleTimeoutMs`: 默认 1800000。 - `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。 - `sweepIntervalMs`: 默认 60000。
- `maxProcesses`: 默认 8,包含 assistant + worker。 - `maxProcesses`: 默认 8,包含 assistant + worker。
- `maxAssistantSessions`: 默认 4。 - `maxAssistantSessions`: 默认 4。
+1
View File
@@ -184,6 +184,7 @@ QQ 出站图片:Worker/Assistant 报告的 workspace 内图片(png/jpg、≤
- initialize:10 秒 - initialize:10 秒
- cancel grace:5 秒 - cancel grace:5 秒
- idle assistant:30 分钟 - idle assistant:30 分钟
- assistant session reset idle:1 小时(`assistantSessionResetIdleMs`,`0` 禁用;Assistant 会话空闲超过该阈值且 owner 没有未完成 Proposal 时,下一条消息开新 session 而不是 resume)
- sweep:60 秒 - sweep:60 秒
- max processes:8(Assistant + Worker 总和) - max processes:8(Assistant + Worker 总和)
- max assistant sessions:4 - max assistant sessions:4
+1
View File
@@ -48,6 +48,7 @@
"promptTimeoutMs": 14400000, "promptTimeoutMs": 14400000,
"cancelGraceMs": 5000, "cancelGraceMs": 5000,
"idleTimeoutMs": 1800000, "idleTimeoutMs": 1800000,
"assistantSessionResetIdleMs": 3600000,
"sweepIntervalMs": 60000, "sweepIntervalMs": 60000,
"maxProcesses": 8, "maxProcesses": 8,
"maxAssistantSessions": 4 "maxAssistantSessions": 4
+27 -2
View File
@@ -645,8 +645,12 @@ export class AssistantManager implements ConversationRuntime {
} }
await this.reserveAssistantCapacity(); await this.reserveAssistantCapacity();
const persisted = this.store.getBinding(conversationKey); 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 (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 cwd = this.prepareAssistantWorkspace(conversationKey, binding);
const spawnWorker = async (nativeSessionId?: string): Promise<AcpWorker> => { const spawnWorker = async (nativeSessionId?: string): Promise<AcpWorker> => {
const worker = new AcpWorker(this.bot, this.config, (crashed, error) => { const worker = new AcpWorker(this.bot, this.config, (crashed, error) => {
@@ -688,7 +692,8 @@ export class AssistantManager implements ConversationRuntime {
assistantWorkspace: cwd, assistantWorkspace: cwd,
botFingerprint: this.bot.fingerprint, botFingerprint: this.bot.fingerprint,
createdAt: now, createdAt: now,
updatedAt: now updatedAt: now,
lastActiveAt: now
}); });
} catch (error) { } catch (error) {
if (this.assistantWorkers.get(conversationKey) === worker) this.assistantWorkers.delete(conversationKey); 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; 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> { private async reserveAssistantCapacity(): Promise<void> {
while (this.assistantWorkers.size >= this.maxAssistantSessions while (this.assistantWorkers.size >= this.maxAssistantSessions
|| this.assistantWorkers.size + (this.active ? 1 : 0) >= this.config.maxProcesses) { || this.assistantWorkers.size + (this.active ? 1 : 0) >= this.config.maxProcesses) {
+1
View File
@@ -87,6 +87,7 @@ const acpSchema = z.object({
promptTimeoutMs: z.number().int().positive().default(14_400_000), promptTimeoutMs: z.number().int().positive().default(14_400_000),
cancelGraceMs: z.number().int().positive().default(5_000), cancelGraceMs: z.number().int().positive().default(5_000),
idleTimeoutMs: z.number().int().positive().default(1_800_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), sweepIntervalMs: z.number().int().positive().default(60_000),
maxProcesses: z.number().int().positive().default(8), maxProcesses: z.number().int().positive().default(8),
maxAssistantSessions: z.number().int().positive().default(4) maxAssistantSessions: z.number().int().positive().default(4)
+5 -1
View File
@@ -14,6 +14,8 @@ export interface AssistantBinding {
botFingerprint: string; botFingerprint: string;
createdAt: number; createdAt: number;
updatedAt: number; updatedAt: number;
// Refreshed on every assistant turn; drives the idle reset of long-inactive sessions.
lastActiveAt?: number;
} }
interface StoreData { interface StoreData {
@@ -90,6 +92,7 @@ export class DurableSessionStore {
const binding = this.data.bindings[chatKey]; const binding = this.data.bindings[chatKey];
if (!binding) return; if (!binding) return;
binding.updatedAt = Date.now(); binding.updatedAt = Date.now();
binding.lastActiveAt = binding.updatedAt;
await this.persist(); await this.persist();
} }
@@ -148,7 +151,8 @@ function parseStoreData(value: unknown): StoreData {
|| typeof binding.assistantWorkspace !== "string" || typeof binding.assistantWorkspace !== "string"
|| typeof binding.botFingerprint !== "string" || typeof binding.botFingerprint !== "string"
|| typeof binding.createdAt !== "number" || 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}'`); throw new Error(`invalid state v3 binding '${chatKey}'`);
} }
} }
+91 -4
View File
@@ -26,7 +26,7 @@ interface Harness {
logFile: string; logFile: string;
} }
async function createHarness(options: { initialize?: boolean; agentEnv?: Record<string, string> } = {}): Promise<Harness> { async function createHarness(options: { initialize?: boolean; agentEnv?: Record<string, string>; acp?: Record<string, unknown> } = {}): Promise<Harness> {
const home = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-assistant-home-")); 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 workspace = await fs.promises.mkdtemp(path.join(os.tmpdir(), "gori-assistant-ws-"));
const logFile = path.join(home, "fake-acp.log"); const logFile = path.join(home, "fake-acp.log");
@@ -34,7 +34,7 @@ async function createHarness(options: { initialize?: boolean; agentEnv?: Record<
configVersion: 3, 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" } }, 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" } }, 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 identity = { botId: "test-bot", platform: "webhook" };
const store = new DurableSessionStore(path.join(home, "state", "acp-sessions.json"), identity); 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 }; return { home, workspace, store, proposals, manager, events, logFile };
} }
async function reopenManager(harness: Harness): Promise<void> { async function reopenManager(harness: Harness, options: { acp?: Record<string, unknown> } = {}): Promise<void> {
await harness.manager.shutdown().catch(() => undefined); await harness.manager.shutdown().catch(() => undefined);
await harness.proposals.close().catch(() => undefined); await harness.proposals.close().catch(() => undefined);
await harness.store.close().catch(() => undefined); await harness.store.close().catch(() => undefined);
@@ -64,7 +64,7 @@ async function reopenManager(harness: Harness): Promise<void> {
configVersion: 3, 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" } }, 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" } }, 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, { harness.manager = new AssistantManager(config.runtime.acp, new BotProfileResolver(config).bot, harness.store, harness.proposals, {
assistantWorkspaceHome: harness.home, 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 () => { test("a worker envelope truncated at the closing tag is salvaged without a repair round", async () => {
const harness = await createHarness(); const harness = await createHarness();
try { try {