From ed7d712488f94e8fa540f16cad4581496dc0b25f Mon Sep 17 00:00:00 2001 From: error2913 <2913949387@qq.com> Date: Sat, 29 Aug 2026 02:22:06 +0800 Subject: [PATCH 1/2] =?UTF-8?q?feat(session):=20=E5=BC=83=E7=94=A8=20.ai?= =?UTF-8?q?=20steer=EF=BC=8C=E7=BB=9F=E4=B8=80=E4=BC=9A=E8=AF=9D=E5=BF=99?= =?UTF-8?q?=E6=97=B6=E6=B6=88=E6=81=AF=E6=8C=82=E8=B5=B7=E9=98=9F=E5=88=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 删除 .ai steer 命令与 steerQueue/方向提示注入机制(.ai live/.ai status 保持原有展示,live 新增挂起状态) - 新增 Session.starting/pendingQueue/deferReceipt/flushPending:会话运行或启动期间收到的新消息一律挂起,由 run/runStream 每轮模型请求前统一入库,修复工具链中插入 user 导致 tool 失配 - Agent.run/runStream 增加同会话闸门(running/starting):同一会话不再并发或排队,始终只有一个会话;acquire 未获许可时不再误释放并发槽 - 触发类挂起消息在链结束后由 resumePending 串行再起一轮;记录类只入库不续跑;.ai stop 清空挂起队列并自增 stopVersion,停止后不复活 - .ai live 展示「启动中」「挂起消息: N」「挂起N」状态 - 新增单元测试:忙时挂起、flush 入库/合并/触发标记/空队列/触发原因提示、stop 清理、同会话闸门防并发 --- README.md | 1 - ...41\345\235\227\350\257\246\350\247\243.md" | 6 +- ...44\344\270\216\351\205\215\347\275\256.md" | 1 - scripts/unit-test-entry.ts | 137 ++++++++++++++++++ skills/aiplugin4-test-suite/SKILL.md | 2 +- .../references/control.md | 6 +- src/agent/agent.ts | 65 ++++++--- src/cmd/root_cmd.ts | 2 - src/cmd/sub_cmd/live.ts | 8 + src/cmd/sub_cmd/steer.ts | 27 ---- src/pipeline.ts | 22 ++- src/session/session.ts | 77 +++++++--- 12 files changed, 278 insertions(+), 76 deletions(-) delete mode 100644 src/cmd/sub_cmd/steer.ts diff --git a/README.md b/README.md index d0f081d..ca7eca5 100644 --- a/README.md +++ b/README.md @@ -358,7 +358,6 @@ max_tokens = 2048 | `.ai role [<名称>]` | - | 查看 / 切换角色设定 | | `.ai model [<模型名>]` | `.ai model deepseek-chat` | 查看 / 设置当前会话模型,`clr` 清除设置恢复默认 | | `.ai stop` | - | 完全暂停当前对话(打断流式输出/工具链/排队请求,清计时器) | -| `.ai steer <内容>` | - | 不打断对话,向工具链插入方向提示,下一轮生效 | ### 记忆管理命令 diff --git "a/docs/03-\346\240\270\345\277\203\346\250\241\345\235\227\350\257\246\350\247\243.md" "b/docs/03-\346\240\270\345\277\203\346\250\241\345\235\227\350\257\246\350\247\243.md" index 362551c..121dcac 100644 --- "a/docs/03-\346\240\270\345\277\203\346\250\241\345\235\227\350\257\246\350\247\243.md" +++ "b/docs/03-\346\240\270\345\277\203\346\250\241\345\235\227\350\257\246\350\247\243.md" @@ -91,7 +91,7 @@ - `session.ts`: - `Setting`:priv、standby、regexTrigger、counter、timer、prob、modelName、activeTimeInfo(start/end/segs)。 - - `Session`:agentName、sessionId、sessionType、state、context、memory、setting、stream、bucket、tool(state/callCount/listen)、lastCtx(不持久化);运行时字段 `running`(是否有 run/runStream 在跑)、`stopVersion`(.ai stop 自增的中止版本号)、`steerQueue`(.ai steer 方向提示队列)均不持久化。 + - `Session`:agentName、sessionId、sessionType、state、context、memory、setting、stream、bucket、tool(state/callCount/listen)、lastCtx(不持久化);运行时字段 `running`(是否有 run/runStream 在跑)、`starting`(是否正在启动 run/runStream,置位期间同会话新消息一律挂起)、`stopVersion`(.ai stop 自增的中止版本号)、`pendingQueue`(会话忙时挂起的待入库消息队列)、`activeRuns`(本会话当前在跑的请求数)均不持久化。 - `toolState` getter:按 `BLOCKED` 与 `DEFAULT_CLOSED` 同步每会话工具开关,并清理已删除工具的残留状态。 - `chat()`:配额检查 → 重置 → 流式判断 → `Agent.run/runStream` → `save()`。 - `reply()`:逐段发送并写入上下文。 @@ -138,7 +138,7 @@ - `run()`:完整编排循环——组装消息 → 请求 → 工具执行(函数调用或提示词工程)→ 复读检测重试 → 最终回复。 - `runStream()`:流式编排,轮询后端流,检测 ```function 代码块截断,执行工具后递归续流。 - `.ai stop` 中止机制:run/runStream 启动时捕获 `session.stopVersion`,在消息组装、模型请求返回、工具执行后、最终回复前设置检查点,检测到版本变化即中止;并发队列中该会话的等待项由 `requestLimiter.cancelBySession` 丢弃。 - - `.ai steer` 注入机制:方向提示进入 `session.steerQueue`,下一轮模型请求组装后以 `system` 消息「【方向提示】…」追加进请求(不入库),实现工具链中不打断地调整方向。 + - 挂起队列机制:会话运行(或正在启动)时收到的新消息经 `session.deferReceipt` 进入 `pendingQueue` 挂起,`run/runStream` 每轮模型请求前 `flushPending` 统一入库(此时上一轮工具回调已写入上下文,插入位置合法);记录类只入库不续跑,触发类在链结束后由 `resumePending` 用第一条触发消息串行再起一轮;`.ai stop` 清空挂起队列并自增 `stopVersion`,停止后不复活。 - `run_context.ts`:`AgentRunContext` 记录单次 run 的观测信息(runId/轮次/工具调用事件/耗时)。 - `api.ts`:对外 API,启动时挂载 `globalThis.aiplugin4`(智能体调用 `chat`/`run`/`getAgent` 与工具注册 `registerTool`),供其他海豹插件使用。 - `agents/`:包含 `compress_agent`(压缩)、`summarize_agent`(摘要);`samples.ts` 的 `sample_agent` 仅作参考,不在 `initAgents` 中初始化。 @@ -147,7 +147,7 @@ - `root_cmd.ts`:注册根命令 `.ai`(注册到 `cmdMap['ai']` 与 `cmdMap['AI']`),别名经 `aliasToCmd` 归一;子命令统一分发,兜住异步异常并回复错误;支持 `--page`/`--p` 翻页。 - `privilege.ts`:命令权限系统,`CmdPrivInfo { priv: [会话权限, 用户权限, 强触权限] }`,预置 `U`(user)/`M`(master)/`I`(inviter)/`S`(会话需 inviter,否则骰主);`PrivilegeManager` 管理 `cmdPriv` 存储与校验,支持 `*` 通配。 -- `sub_cmd/`:17 个文件,其中 16 个注册到根命令,`sample.ts` 未注册(仅开发参考),详见 [05-命令与配置](05-命令与配置.md)。 +- `sub_cmd/`:21 个文件,其中 20 个注册到根命令,`sample.ts` 未注册(仅开发参考),详见 [05-命令与配置](05-命令与配置.md)。 ## 定时器(src/timer.ts) diff --git "a/docs/05-\345\221\275\344\273\244\344\270\216\351\205\215\347\275\256.md" "b/docs/05-\345\221\275\344\273\244\344\270\216\351\205\215\347\275\256.md" index d18f493..92f7e36 100644 --- "a/docs/05-\345\221\275\344\273\244\344\270\216\351\205\215\347\275\256.md" +++ "b/docs/05-\345\221\275\344\273\244\344\270\216\351\205\215\347\275\256.md" @@ -29,7 +29,6 @@ | `block` | 黑名单管理(骰主) | `.ai block add <用户ID/群ID> <原因>`;`.ai block rm <用户ID/群ID>`;`.ai block list`;被拉黑对象无法触发 AI 对话,指令不受影响;AI 建议拉黑默认需骰主确认,可在工具配置关闭 | | `token` | token 用量统计 | 见下文 | | `stop` | 完全暂停当前对话(打断流式/工具链/排队,清计时器) | `.ai stop` | -| `steer` | 向当前对话插入方向提示(不打断) | `.ai steer <内容>` | | `model` | 查看/设置会话模型 | `.ai model`;`.ai model <模型名>`;`.ai model clr` | | `sample` | 示例智能体(未注册到根命令) | 仅开发参考 | diff --git a/scripts/unit-test-entry.ts b/scripts/unit-test-entry.ts index 609b62f..3125883 100644 --- a/scripts/unit-test-entry.ts +++ b/scripts/unit-test-entry.ts @@ -20,6 +20,8 @@ import { InMemoryMemoryStorage } from "../src/memory/v2/storage"; import { createMemoryEngine } from "../src/memory/v2"; import { requestLimiter } from "../src/utils/concurrency"; import { Context } from "../src/context/context"; +import Agent from "../src/agent/agent"; +import { Session } from "../src/session/session"; import Image from "../src/resource/image"; import Tool, { toolMap } from "../src/tool/tool"; import { registerDispatchTools } from "../src/tool/tools/core/tool_dispatch"; @@ -907,4 +909,139 @@ export const tests: Record void | Promise> = { TC.intConfigs['请求队列上限'] = 0; resetConfigCache(); }, + + /** 会话忙时挂起:deferReceipt 不直接入库,进入 pendingQueue(空闲时 handleReceipt 仍直接入库) */ + async testDeferReceiptWhileBusy(): Promise { + const session = new Session(); + session.sessionId = 'sess_defer'; + session.sessionType = 'user'; + session.running = true; + const ctx = makeCtx(); + const msg = { message: '测试消息' } as any; + const messageArray = [{ type: 'text', data: { text: '测试消息' } }] as any; + await session.deferReceipt(ctx, msg, messageArray, 'trigger'); + assert.equal(session.pendingQueue.length, 1, '忙时应进入挂起队列'); + assert.equal(session.context.messages.length, 0, '忙时不应直接写入上下文'); + const p = session.pendingQueue[0]; + assert.equal(p.kind, 'trigger'); + assert.equal(p.content, '测试消息'); + assert.equal(p.userId, 'QQ:10000'); + + // 空闲对照:handleReceipt 直接入库 + session.running = false; + await session.handleReceipt(ctx, msg, messageArray); + assert.equal(session.context.messages.length, 1, '空闲时 handleReceipt 应直接入库'); + assert.equal(session.pendingQueue.length, 1, '已挂起队列不受空闲路径影响'); + }, + + /** flushPending:统一入库并返回是否存在触发类消息;连续 user 消息自动合并 */ + async testFlushPendingAddsMessagesAndReturnsFlag(): Promise { + const session = new Session(); + session.sessionId = 'sess_flush'; + session.sessionType = 'user'; + const ctx = makeCtx(); + const msgA = { message: '第一条' } as any; + const msgB = { message: '第二条' } as any; + const arrA = [{ type: 'text', data: { text: '第一条' } }] as any; + const arrB = [{ type: 'text', data: { text: '第二条' } }] as any; + + // 仅记录类:flushPending 返回 false(链结束不续跑) + await session.deferReceipt(ctx, msgA, arrA, 'record'); + assert.equal(await session.flushPending(), false, '纯记录类不应续跑'); + assert.equal(session.pendingQueue.length, 0, 'flush 后队列应清空'); + + // 触发类 + 记录类混排:返回 true,且全部入库(连续 user 消息合并为同一条) + const session2 = new Session(); + session2.sessionId = 'sess_flush2'; + session2.sessionType = 'user'; + await session2.deferReceipt(ctx, msgA, arrA, 'record'); + await session2.deferReceipt(ctx, msgB, arrB, 'trigger'); + assert.equal(await session2.flushPending(), true, '含触发类应返回 true'); + assert.equal(session2.context.messages.length, 1, '连续 user 消息应合并为同一条'); + const userMsg = session2.context.messages[0] as any; + assert.equal(userMsg.role, 'user'); + const texts = userMsg.contentItems.map((i: any) => i.text); + assert.deepEqual(texts, ['第一条', '第二条'], '连续 user 消息应合并为同一条'); + }, + + /** flushPending 空队列:返回 false 且不改动上下文 */ + async testFlushPendingEmpty(): Promise { + const session = new Session(); + session.sessionId = 'sess_empty'; + session.sessionType = 'user'; + assert.equal(await session.flushPending(), false); + assert.equal(session.context.messages.length, 0); + }, + + /** AI 设定触发挂起:systemReason 在 flush 时以「触发原因提示」写入用户消息 */ + async testFlushPendingAddsSystemReason(): Promise { + const session = new Session(); + session.sessionId = 'sess_reason'; + session.sessionType = 'user'; + const ctx = makeCtx(); + const msg = { message: '关键词消息' } as any; + const messageArray = [{ type: 'text', data: { text: '关键词消息' } }] as any; + await session.deferReceipt(ctx, msg, messageArray, 'trigger', '因为你提到了关键词'); + assert.equal(await session.flushPending(), true); + const userMsg = session.context.messages[0] as any; + const items = userMsg.contentItems; + const reasonItem = items.find((i: any) => i.systemName === '触发原因提示'); + assert.ok(reasonItem, '应写入触发原因提示条目'); + assert.equal(reasonItem.text, '因为你提到了关键词'); + assert.ok(items.some((i: any) => i.text === '关键词消息'), '用户消息本身也应入库'); + }, + + /** .ai stop:清空挂起队列(停止后不复活) */ + async testStopClearsPendingQueue(): Promise { + const session = new Session(); + session.sessionId = 'sess_stop'; + session.sessionType = 'user'; + const ctx = makeCtx(); + const msg = { message: '挂起消息' } as any; + const messageArray = [{ type: 'text', data: { text: '挂起消息' } }] as any; + await session.deferReceipt(ctx, msg, messageArray, 'trigger'); + assert.equal(session.pendingQueue.length, 1); + await session.stopConversation(); + assert.equal(session.pendingQueue.length, 0, 'stop 后挂起队列应清空'); + assert.equal(session.context.messages.length, 0, 'stop 后挂起消息不应复活入库'); + }, + + /** 同会话闸门:第一条 run 挂起在 runInternal 时,第二条 run 被 starting 闸门直接拦截,activeRuns 恒 ≤1 */ + async testStartingGatePreventsConcurrentRun(): Promise { + TC.intConfigs['请求并发上限'] = 1; + TC.intConfigs['请求队列上限'] = 3; + resetConfigCache(); + + const agent = new Agent(); + const session = new Session(); + session.sessionId = 'sess_gate'; + const origInternal = (agent as any).runInternal; + let releaseRun = () => { }; + const gate = new Promise(resolve => { releaseRun = resolve; }); + (agent as any).runInternal = async () => { + assert.ok(session.activeRuns <= 1, '同一会话不应并发多个 run'); + await gate; + }; + try { + const p1 = agent.run(session, makeCtx(), { message: '1' } as any); + // 等第一条真正进入 runInternal(acquire 完成、running=true、activeRuns=1、挂起在 gate 上) + await new Promise(r => setTimeout(r, 0)); + // 第二条同刻到达:run() 同步闸门(running/starting)应直接拦截,不进入 acquire/runInternal + await agent.run(session, makeCtx(), { message: '2' } as any); + assert.equal(session.activeRuns, 1, '并发保护下 activeRuns 恒为 1'); + assert.equal(session.running, true, '第一条仍在运行'); + assert.equal(session.starting, false, '运行中不再处于启动中'); + releaseRun(); + await p1; + assert.equal(session.activeRuns, 0, '运行结束后 activeRuns 归零'); + assert.equal(session.running, false); + assert.equal(session.starting, false, '结束后 starting 复位'); + } finally { + releaseRun(); + (agent as any).runInternal = origInternal; + TC.intConfigs['请求并发上限'] = 0; + TC.intConfigs['请求队列上限'] = 0; + resetConfigCache(); + } + }, }; diff --git a/skills/aiplugin4-test-suite/SKILL.md b/skills/aiplugin4-test-suite/SKILL.md index 947ecf9..756d0f9 100644 --- a/skills/aiplugin4-test-suite/SKILL.md +++ b/skills/aiplugin4-test-suite/SKILL.md @@ -30,7 +30,7 @@ description: aiplugin4 插件综合测试与调试技能:通过 qqmcp 对测 | 用户提到的指令 | 用例文件 | |---|---| -| `status` `ctxn` `on` `off` `sb/standby` `fgt/forget` `role` `model` `stop` `steer` `ign/ignore` | [control.md](references/control.md) | +| `status` `ctxn` `on` `off` `sb/standby` `fgt/forget` `role` `model` `stop` `ign/ignore` | [control.md](references/control.md) | | `memo/memory` | [memory.md](references/memory.md) | | `tool` | [tool.md](references/tool.md) | | `priv/privilege` `prompt` `tk/token` `timer` | [admin.md](references/admin.md) | diff --git a/skills/aiplugin4-test-suite/references/control.md b/skills/aiplugin4-test-suite/references/control.md index 4b3b05b..d9e920c 100644 --- a/skills/aiplugin4-test-suite/references/control.md +++ b/skills/aiplugin4-test-suite/references/control.md @@ -2,7 +2,7 @@ 权限:U=任意成员;I=邀请者/管理/群主/白名单/骰主;S=骰主(会话权限≥1 时为邀请者)。 -## 基础控制(status/ctxn/on/standby/off/forget/role/model/stop/steer) +## 基础控制(status/ctxn/on/standby/off/forget/role/model/stop) | ID | 指令 | 权限 | 预期关键字 | 备注 | |---|---|---|---|---| @@ -32,8 +32,8 @@ | CTRL-24 | `.ai model clr` | U | `已清除当前会话的模型设置` | 结束后恢复快照模型 | | CTRL-25 | `.ai model 不存在的模型` | U | `不存在` | 错误路径 | | CTRL-26 | `.ai stop` | U | `当前没有正在进行的对话` | 无进行中对话时(原 shut 用例) | -| CTRL-27 | `.ai steer` | U | `缺少内容` | 无内容错误路径 | -| CTRL-28 | `.ai steer 测试方向` | U | `当前没有正在进行的对话` | 无进行中对话时 | +| CTRL-27 | 运行中连发普通消息 | U | 当前任务结束后自动入库、不并发回复 | 先发一个长工具链任务制造运行中会话,再连发普通消息,全程无并发/排队 | +| CTRL-28 | 运行中发触发消息 | U | 前一轮结束后串行自动再触发一轮 | 运行中发送触发消息,等待前一轮结束并自动续跑回复 | ## 忽略名单(ignore) diff --git a/src/agent/agent.ts b/src/agent/agent.ts index d4b9407..8995bb9 100644 --- a/src/agent/agent.ts +++ b/src/agent/agent.ts @@ -14,7 +14,7 @@ import { ToolName } from "../tool/tool"; import Tool from "../tool/tool"; import { ToolInfo } from "../tool/types"; import { requestLimiter } from "../utils/concurrency"; -import { buildSystemMessage, handleMessages, RequestMessage } from "../utils/message"; +import { buildSystemMessage, handleMessages } from "../utils/message"; import { checkRepeat, handleReply } from "../utils/string"; import { revive, TypeDescriptor } from "../utils/utils"; @@ -98,16 +98,38 @@ export default class Agent { async run(session: Session, ctx: seal.MsgContext, msg: seal.Message, tool_choice?: string): Promise { // 启动前捕获 stopVersion:.ai stop 发生在排队期间或拿到许可瞬间都会因版本变化而中止 const version = session.stopVersion; - if (!(await requestLimiter.acquire(session.sessionId))) return; + // 同会话闸门:已有 run 在跑或正在启动(含排队等待)时不再并发起一轮,消息已由 pipeline 挂起 + if (session.running || session.starting) return; + session.starting = true; + let acquired = false; try { + acquired = await requestLimiter.acquire(session.sessionId); + if (!acquired) return; if (session.stopVersion !== version) return; session.running = true; + session.starting = false; session.activeRuns++; await this.runInternal(session, ctx, msg, tool_choice); } finally { - session.activeRuns = Math.max(0, session.activeRuns - 1); - session.running = session.activeRuns > 0; - requestLimiter.release(session.sessionId); + // 未拿到许可(被 stop 取消/队列满/acquire 异常)时不得释放别人的并发槽 + if (acquired) { + session.activeRuns = Math.max(0, session.activeRuns - 1); + session.running = session.activeRuns > 0; + // 先释放并发槽再续跑,避免续跑 acquire 与自己持有的槽死锁 + requestLimiter.release(session.sessionId); + await this.resumePending(session, version); + } + session.starting = false; + } + } + + /** 链结束后处理挂起队列:先全部入库(记录类不续跑),若含触发类且未被 stop 则用第一条串行再起一轮 */ + private async resumePending(session: Session, version: number): Promise { + if (session.pendingQueue.length === 0) return; + const firstTrigger = session.pendingQueue.find(p => p.kind === 'trigger'); + const hasTrigger = await session.flushPending(); + if (hasTrigger && firstTrigger && session.stopVersion === version) { + await session.chat(firstTrigger.ctx, firstTrigger.msg, '挂起触发'); } } @@ -134,8 +156,9 @@ export default class Agent { // stop 中止检查点:上一轮工具执行/回调期间被 stop 则不再发下一轮请求 if (session.stopVersion !== version) return; trace.beginTurn(); + // 挂起消息入库:上一轮工具回调已写入上下文,此时插入位置合法(修复工具链中插入 user 导致 tool 失配的问题) + await session.flushPending(); const messages = await handleMessages(ctx, session, this.isMultimodalChat(session), toolInfos || [], systemMessage); - this.injectSteers(session, messages); const { content: raw_reply, tool_calls, reasoning_content } = await streamService.sendChatRequest(messages, toolInfos || [], tool_choice || 'auto', session.setting.modelName, trace.runId); // stop 中止检查点:模型请求期间被 stop,丢弃本轮输出直接中止 if (session.stopVersion !== version) return; @@ -241,27 +264,32 @@ export default class Agent { log.info(`[run] ${trace.summary()}`); } - /** 把 .ai steer 注入的方向提示追加到请求消息末尾(最新指令),并清空队列;不写入持久化上下文 */ - private injectSteers(session: Session, messages: RequestMessage[]): void { - const steers = session.drainSteers(); - for (const steer of steers) { - messages.push({ role: 'system', content: `【方向提示】${steer}` }); - } - } /** 流式编排:与 run() 同层的流式循环(start → poll → 工具调用 → 递归续流),工具轮数由配置控制 */ async runStream(session: Session, ctx: seal.MsgContext, msg: seal.Message): Promise { const version = session.stopVersion; - if (!(await requestLimiter.acquire(session.sessionId))) return; + // 同会话闸门:已有 run/runStream 在跑或正在启动(含排队等待)时不再并发起一轮,消息已由 pipeline 挂起 + if (session.running || session.starting) return; + session.starting = true; + let acquired = false; try { + acquired = await requestLimiter.acquire(session.sessionId); + if (!acquired) return; if (session.stopVersion !== version) return; session.running = true; + session.starting = false; session.activeRuns++; await this.runStreamInner(session, ctx, msg); } finally { - session.activeRuns = Math.max(0, session.activeRuns - 1); - session.running = session.activeRuns > 0; - requestLimiter.release(session.sessionId); + // 未拿到许可(被 stop 取消/队列满/acquire 异常)时不得释放别人的并发槽 + if (acquired) { + session.activeRuns = Math.max(0, session.activeRuns - 1); + session.running = session.activeRuns > 0; + // 先释放并发槽再续跑,避免续跑 acquire 与自己持有的槽死锁 + requestLimiter.release(session.sessionId); + await this.resumePending(session, version); + } + session.starting = false; } } @@ -276,10 +304,11 @@ export default class Agent { const version = session.stopVersion; await session.stopCurrentChatStream(); + // 挂起消息入库:上一轮工具回调已写入上下文,此时插入位置合法(修复工具链中插入 user 导致 tool 失配的问题) + await session.flushPending(); const sys = systemMessage ?? await buildSystemMessage(ctx, session); const messages = await handleMessages(ctx, session, this.isMultimodalChat(session), undefined, sys); - this.injectSteers(session, messages); const id = await streamService.startStream(messages, session.setting.modelName, trace.runId); if (!id) return; // stop 发生在 startStream 期间:结束刚建的新流并中止,避免轮询一个未被 stop 的流 diff --git a/src/cmd/root_cmd.ts b/src/cmd/root_cmd.ts index 4832218..4b46f41 100644 --- a/src/cmd/root_cmd.ts +++ b/src/cmd/root_cmd.ts @@ -22,7 +22,6 @@ import { registerCmdPrompt } from "./sub_cmd/prompt"; import { registerCmdRole } from "./sub_cmd/role"; import { registerCmdStandby } from "./sub_cmd/standby"; import { registerCmdStatus } from "./sub_cmd/status"; -import { registerCmdSteer } from "./sub_cmd/steer"; import { registerCmdStop } from "./sub_cmd/stop"; import { registerCmdTimer } from "./sub_cmd/timer"; import { registerCmdToken } from "./sub_cmd/token"; @@ -78,7 +77,6 @@ export class SubCmd { registerCmdIgnore(); registerCmdToken(); registerCmdStop(); - registerCmdSteer(); registerCmdModel(); registerCmdBlock(); diff --git a/src/cmd/sub_cmd/live.ts b/src/cmd/sub_cmd/live.ts index 7834c24..ef396ce 100644 --- a/src/cmd/sub_cmd/live.ts +++ b/src/cmd/sub_cmd/live.ts @@ -14,7 +14,9 @@ function stateText(session: Session, qi: QueueInfo): string { if (session.stream.id && session.stream.toolCallStatus) return '流式-工具调用中'; if (session.stream.id) return '流式输出中'; if (session.activeRuns > 0) return `运行中(${session.activeRuns}个请求)`; + if (session.starting) return '启动中'; if (qi.queuedBySession > 0) return '排队中'; + if (session.pendingQueue.length > 0) return `挂起${session.pendingQueue.length}`; return '空闲'; } @@ -23,6 +25,7 @@ function collectRuntime(session: Session): { state: string; activeRuns: number; queued: number; + pending: number; timers: { target: number; interval: number; activeTime: number }; } { const qi = requestLimiter.getQueueInfo(session.sessionId); @@ -34,6 +37,7 @@ function collectRuntime(session: Session): { state: stateText(session, qi), activeRuns: session.activeRuns, queued: qi.queuedBySession, + pending: session.pendingQueue.length, timers }; } @@ -56,6 +60,7 @@ function formatRuntime(session: Session): string { `状态: ${r.state}`, `流式: ${session.stream.id ? (session.stream.toolCallStatus ? '工具调用中' : '输出中') : '无'}`, `并发: 全局活跃 ${qi.active}/${qi.maxConcurrent} | 本会话活跃 ${r.activeRuns} | 本会话排队 ${r.queued}/${qi.maxQueue}`, + `挂起消息: ${r.pending}`, `定时器: ${timerText(r.timers)}` ].join('\n'); } @@ -64,6 +69,8 @@ function formatRuntime(session: Session): string { function isBusy(session: Session): boolean { if (session.stream.id) return true; if (session.activeRuns > 0) return true; + if (session.starting) return true; + if (session.pendingQueue.length > 0) return true; if (requestLimiter.getQueueInfo(session.sessionId).queuedBySession > 0) return true; return TimerManager.getTimers(session.sessionId).length > 0; } @@ -97,6 +104,7 @@ function formatAll(): string { const sessionType = session.sessionType === 'user' ? '私聊' : '群聊'; const timerCount = r.timers.target + r.timers.interval + r.timers.activeTime; const suffix = [ + r.pending > 0 ? `挂起${r.pending}` : '', r.queued > 0 ? `排队${r.queued}` : '', timerCount > 0 ? `定时器:${timerText(r.timers)}` : '' ].filter(Boolean).join(' '); diff --git a/src/cmd/sub_cmd/steer.ts b/src/cmd/sub_cmd/steer.ts deleted file mode 100644 index 394e40b..0000000 --- a/src/cmd/sub_cmd/steer.ts +++ /dev/null @@ -1,27 +0,0 @@ -// .ai steer:在不打断对话的前提下,向工具链插入方向提示 -import { U } from "../privilege"; -import { SubCmd, SubCmdContext } from "../root_cmd"; - -export function registerCmdSteer() { - const cmd = new SubCmd('steer'); - cmd.desc = '向当前对话插入方向提示'; - cmd.help = `帮助: -【.ai steer <内容>】不打断当前对话,把内容作为方向提示插入工具链,下一轮模型请求生效`; - cmd.priv = { priv: U }; - cmd.solve = (scc: SubCmdContext) => { - const { ctx, msg, cmdArgs, session, ret } = scc; - - const text = cmdArgs.getRestArgsFrom(2).trim(); - if (!text) { - seal.replyToSender(ctx, msg, '【.ai steer <内容>】不打断当前对话,向工具链插入方向提示\n缺少内容'); - return ret; - } - if (!session.running) { - seal.replyToSender(ctx, msg, '当前没有正在进行的对话'); - return ret; - } - session.steer(text); - seal.replyToSender(ctx, msg, '已插入方向提示,将在下一轮生效'); - return ret; - } -} diff --git a/src/pipeline.ts b/src/pipeline.ts index 0811463..9f7b3b3 100644 --- a/src/pipeline.ts +++ b/src/pipeline.ts @@ -506,6 +506,8 @@ export class MessagePipeline { const gid = ctx.isPrivate ? '' : ctx.group!.groupId; const sid = ctx.isPrivate ? uid : gid; const session = getSession(sid); + // 会话忙(正在运行或正在启动)时,新消息不直接入库/触发,改为挂起由 run 循环统一处理 + const sessionBusy = session.running || session.starting; // 检查活跃时间定时器 session.checkActiveTimer(ctx); @@ -537,6 +539,9 @@ export class MessagePipeline { if (messageArray.length === 0) return; const supplementTypes = messageArray.filter(item => item.type !== 'text').map(item => item.type); if (supplementTypes.some(type => !CQ_TYPES_ALLOW.includes(type))) return; + if (sessionBusy) { + return session.deferReceipt(ctx, msg, messageArray, 'record').then(() => session.save()); + } return session.handleReceipt(ctx, msg, messageArray).then(() => session.save()); } @@ -556,14 +561,18 @@ export class MessagePipeline { // 检查CQ码 const CQTypes = messageArray.filter(item => item.type !== 'text').map(item => item.type); if (CQTypes.length === 0 || CQTypes.every(item => CQ_TYPES_ALLOW.includes(item))) { - if (session.context.timer) clearTimeout(session.context.timer); - session.context.timer = null; + // 运行中不重置待触发计时器(待机计数/概率/计时器一律跳过,只入库不推进) + if (!sessionBusy && session.context.timer) clearTimeout(session.context.timer); + if (!sessionBusy) session.context.timer = null; // 非指令消息触发(受会话开关控制) if (session.setting.regexTrigger && triggerRegex.test(messageText)) { const fmtCondition = parseInt(seal.format(ctx, `{${triggerCondition}}`)); if (fmtCondition === 1) { markCoreMessageRecorded(coreMessageKey); + if (sessionBusy) { + return session.deferReceipt(ctx, msg, messageArray, 'trigger').then(() => session.save()); + } return session.handleReceipt(ctx, msg, messageArray) .then(() => session.chat(ctx, msg, '非指令')); } @@ -591,6 +600,11 @@ export class MessagePipeline { } markCoreMessageRecorded(coreMessageKey); + if (sessionBusy) { + // 先消费一次性触发条件再挂起,避免条件残留导致下次重复触发 + triggerConditionMap[sid].splice(i, 1); + return session.deferReceipt(ctx, msg, messageArray, 'trigger', condition.reason).then(() => session.save()); + } return session.handleReceipt(ctx, msg, messageArray) .then(() => session.context.addSystemUserMessage(condition.reason, '触发原因提示')) .then(() => triggerConditionMap[sid].splice(i, 1)) @@ -602,6 +616,10 @@ export class MessagePipeline { const setting = session.setting; if (setting.standby || Config.base.GLOBAL_STANDBY) { markCoreMessageRecorded(coreMessageKey); + if (sessionBusy) { + // 运行中待机消息只挂起入库:计数/概率/计时器一律跳过,不推进、不触发 + return session.deferReceipt(ctx, msg, messageArray, 'record').then(() => session.save()); + } return session.handleReceipt(ctx, msg, messageArray) .then((): void | Promise => { if (setting.counter > -1) { diff --git a/src/session/session.ts b/src/session/session.ts index 4eb6e52..7cbaa2d 100644 --- a/src/session/session.ts +++ b/src/session/session.ts @@ -14,7 +14,7 @@ import { ToolState } from "../tool/tool"; import { toolMap } from "../tool/tool"; import { ToolListen } from "../tool/types"; import { requestLimiter } from "../utils/concurrency"; -import { MessageSegment, normalizeRenderTags, stripInternalTags, transformArrayToContent } from "../utils/string"; +import { MessageSegment, normalizeRenderTags, transformArrayToContent } from "../utils/string"; import { TypeDescriptor } from "../utils/utils"; import { getRecordMessageId, replyToSender } from "../utils/utils"; @@ -25,8 +25,19 @@ import User from "./user"; const log = logger.withTag('session'); -/** 持久化时排除的运行时字段(监听器/运行状态/方向提示等),不写入存储、不参与 revive 恢复 */ -export const SESSION_RUNTIME_KEYS = new Set(['lastCtx', 'running', 'stopVersion', 'steerQueue', 'activeRuns']); +/** 持久化时排除的运行时字段(监听器/运行状态/挂起队列等),不写入存储、不参与 revive 恢复 */ +export const SESSION_RUNTIME_KEYS = new Set(['lastCtx', 'running', 'starting', 'stopVersion', 'pendingQueue', 'activeRuns']); + +/** 会话忙时挂起的消息:运行中收到的新消息先入队,由下一轮模型请求前统一入库(触发类可在链结束后续跑一轮) */ +export interface PendingMessage { + ctx: seal.MsgContext; + msg: seal.Message; + content: string; + userId: string; + messageId: string; + kind: 'trigger' | 'record'; + systemReason?: string; +} export class Setting { static validKeys: (keyof Setting)[] = ['priv', 'standby', 'counter', 'timer', 'prob', 'activeTimeInfo', 'modelName', 'regexTrigger']; @@ -127,8 +138,10 @@ export class Session { running = false; /** 运行时字段:stop 时自增,运行循环启动时捕获、检测到变化即中止(不持久化) */ stopVersion = 0; - /** 运行时字段:.ai steer 注入的方向提示队列,下一轮模型请求时清空注入(不持久化) */ - steerQueue: string[] = []; + /** 运行时字段:本会话正在启动 run/runStream(含在请求并发队列中等待时);置位后同会话新消息一律挂起(不持久化) */ + starting = false; + /** 运行时字段:会话忙(starting/running)时挂起的待入库消息队列,下一轮模型请求前统一入库(不持久化) */ + pendingQueue: PendingMessage[] = []; /** 运行时字段:本会话当前在跑的 run/runStream 请求数(同会话并发重叠时 >1;不持久化) */ activeRuns = 0; tool: { @@ -160,7 +173,8 @@ export class Session { this.lastCtx = null; this.running = false; this.stopVersion = 0; - this.steerQueue = []; + this.starting = false; + this.pendingQueue = []; this.activeRuns = 0; const listen = createToolListen(); this.tool = { @@ -265,6 +279,42 @@ export class Session { await this.context.addUserMessage(ctx, content, ctx.player!.userId, getRecordMessageId(ctx, msg)); } + /** 会话忙时挂起消息:与 handleReceipt 等价但不入库,进入 pendingQueue 由 flushPending 统一处理(不触发并发) */ + async deferReceipt(ctx: seal.MsgContext, msg: seal.Message, messageArray: MessageSegment[], kind: 'trigger' | 'record', systemReason?: string) { + this.lastCtx = ctx; + const { content } = await transformArrayToContent(ctx, messageArray); + this.pendingQueue.push({ + ctx, + msg, + content, + userId: ctx.player!.userId, + messageId: getRecordMessageId(ctx, msg), + kind, + systemReason + }); + } + + /** 取出并清空挂起队列 */ + drainPending(): PendingMessage[] { + const pending = this.pendingQueue; + this.pendingQueue = []; + return pending; + } + + /** 把挂起队列全部写入上下文(此时上一轮工具回调已入库,位置合法)并保存;返回是否存在触发类消息 */ + async flushPending(): Promise { + const pending = this.drainPending(); + if (pending.length === 0) return false; + let hasTrigger = false; + for (const p of pending) { + if (p.kind === 'trigger') hasTrigger = true; + await this.context.addUserMessage(p.ctx, p.content, p.userId, p.messageId); + if (p.systemReason) await this.context.addSystemUserMessage(p.systemReason, '触发原因提示'); + } + this.save(); + return hasTrigger; + } + async reply(ctx: seal.MsgContext, msg: seal.Message, contextArray: string[], replyArray: string[], _images: Image[], options: { withSegmentDelay?: boolean } = {}, reasoningContent?: string) { const { withSegmentDelay = false } = options; const { SEGMENT_DELAY_ENABLED, SEGMENT_DELAY_MS, SEGMENT_IMAGE_EXTRA_DELAY_MS } = Config.reply; @@ -390,7 +440,9 @@ export class Session { const hadTimer = this.context.timer !== null; await this.stopCurrentChatStream(); this.stopVersion++; - // 完全暂停:清掉待触发的计时器(计数器/概率/触发条件保留,需主动触发) + this.starting = false; + // 完全暂停:清掉挂起消息(停止后不复活)与待触发的计时器(计数器/概率/触发条件保留,需主动触发) + this.pendingQueue = []; if (this.context.timer) clearTimeout(this.context.timer); this.context.timer = null; const queueCleared = requestLimiter.cancelBySession(this.sessionId); @@ -399,15 +451,4 @@ export class Session { return { hadStream, hadRun, hadTimer, queueCleared }; } - /** 插入方向提示:进入 steerQueue,由下一轮模型请求注入(不打断当前对话) */ - steer(text: string): void { - this.steerQueue.push(stripInternalTags(text)); - } - - /** 取出并清空方向提示队列 */ - drainSteers(): string[] { - const steers = this.steerQueue; - this.steerQueue = []; - return steers; - } } From 43abc6a3020efcafc64c009b66a4ae0fd1fa3270 Mon Sep 17 00:00:00 2001 From: error2913 <2913949387@qq.com> Date: Sat, 29 Aug 2026 02:44:15 +0800 Subject: [PATCH 2/2] =?UTF-8?q?fix(live):=20.ai=20live=20=E5=8D=95?= =?UTF-8?q?=E4=BC=9A=E8=AF=9D=E8=A7=86=E5=9B=BE=E4=B8=8D=E5=86=8D=E5=B1=95?= =?UTF-8?q?=E7=A4=BA=E5=85=A8=E5=B1=80=E5=B9=B6=E5=8F=91=E4=BF=A1=E6=81=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 普通 .ai live 只展示本会话活跃/排队/挂起,全局活跃与全局队列上限仅保留在 .ai live all --- src/cmd/sub_cmd/live.ts | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/src/cmd/sub_cmd/live.ts b/src/cmd/sub_cmd/live.ts index ef396ce..9d467a6 100644 --- a/src/cmd/sub_cmd/live.ts +++ b/src/cmd/sub_cmd/live.ts @@ -53,13 +53,12 @@ function timerText(timers: { target: number; interval: number; activeTime: numbe /** 单会话视图:.ai live */ function formatRuntime(session: Session): string { const r = collectRuntime(session); - const qi = requestLimiter.getQueueInfo(session.sessionId); const sessionType = session.sessionType === 'user' ? '私聊' : '群聊'; return [ `【运行状态】${sessionType}会话 ${session.sessionId}`, `状态: ${r.state}`, `流式: ${session.stream.id ? (session.stream.toolCallStatus ? '工具调用中' : '输出中') : '无'}`, - `并发: 全局活跃 ${qi.active}/${qi.maxConcurrent} | 本会话活跃 ${r.activeRuns} | 本会话排队 ${r.queued}/${qi.maxQueue}`, + `并发: 本会话活跃 ${r.activeRuns} | 本会话排队 ${r.queued}`, `挂起消息: ${r.pending}`, `定时器: ${timerText(r.timers)}` ].join('\n');