diff --git a/bots/cormini/persona/persona.ts b/bots/cormini/persona/persona.ts index 99920d20..2cd30c03 100644 --- a/bots/cormini/persona/persona.ts +++ b/bots/cormini/persona/persona.ts @@ -32,6 +32,7 @@ import type { ToolDef, ToolTag, } from 'cortico/core/types.ts'; +import { defaultExternalizes } from 'cortico/core/types.ts'; /** 默认的存在方式自述源文件;bot 不覆盖 orientation 时读它,控制台可编辑同一份。 */ /** 前缀装配模板与 World 段模板住在软件包里(跟着代码走),不在人格工作区。 */ @@ -698,7 +699,7 @@ export class Cormini implements Persona { if (taps.length === 1) return taps[0]; return { onEvent: (event) => { for (const t of taps) t.onEvent(event); }, - externalizes: (event) => taps.some((t) => t.externalizes?.(event) ?? false), + externalizes: (event) => taps.some((t) => (t.externalizes?.(event) ?? defaultExternalizes(event))), onRoundEnd: () => { for (const t of taps) t.onRoundEnd?.(); }, onAbort: (reason) => { for (const t of taps) t.onAbort?.(reason); }, }; diff --git a/docs/sessions.md b/docs/sessions.md index 44fbf25f..bb7111fb 100644 --- a/docs/sessions.md +++ b/docs/sessions.md @@ -26,6 +26,12 @@ Persona 声明,Core 只把 `role` 当分类标签。 模型返回的工具调用只有 `status` 为 `completed` 才执行。其余的回执为未执行,主循环记一条 warn (工具名、`call_id`、原始 `status`),同一轮排在它后面的调用也不执行。 +一批之内,新事件在轮次边界进入上下文:本轮工具全部返回之后,或 `preempt`、`interrupt` 取消模型轮 +之后。`preempt` 只取消 `outputTap` 尚未报告外部输出的模型轮,丢弃其输出;`interrupt` 取消模型轮时 +保留已外化的输出,执行工具期间则停止 `interruptible` 工具、跳过本轮尚未开始的调用。两者取消后都在 +同一批内带着新事件开始下一轮,被取消的轮计入 `rounds()` 的轮数,不调用 `onTurnEnded`;到达时没有 +可取消的轮,事件在下一次模型请求前送入。 + ## 上下文与交接 Core 的输入上限为 `hardTokens = max(0, 生效窗口 − (spec.maxTokens ?? 0))`。生效窗口取服务探测值与配置的 diff --git a/docs/world-compatibility.md b/docs/world-compatibility.md index bd5beb75..e9c73258 100644 --- a/docs/world-compatibility.md +++ b/docs/world-compatibility.md @@ -28,17 +28,20 @@ ## L3 投递语义 -- 按 `trigger` 安排投递:`preempt` 请求中断在途模型调用并立即投递,`flush` 立即投递并带上积压, +- 按 `trigger` 安排投递:`preempt` 取消尚未产生外部输出的模型调用并立即投递,`interrupt` 停止当前轮 + (模型调用与 `interruptible` 工具)并立即投递,`flush` 立即投递并带上积压, `debounce` 参与合批,`piggyback` 只排队、随其他批次投递;`deliver: false` 只存储。 - `pushDeferred`:投递时才调用 `render` 生成正文;返回 null、抛错或超时时不存储、不投递。 - `pushCandidate`:先归档原始事件,投递时整批交给 World 选取并生成正文。 - `ephemeral` 事件在下一批投递时从 session 移除,事件库保留。 - 提供事件库读取(`store`)、`drainPendingEvents(filter)` 与工具执行期间的 `queueExternalEvents`。 +- `withdrawPending` 撤回、`promotePending` 提级未投递的事件;事件写入 session 或被丢弃时调用 + `onEventsSettled`。 - 上下文交接后、恢复投递前调用 `onHandoffEnded()`。 ## L4 实时 -- `outputTap()`:把主 session 的模型输出流逐段交给 World;`externalizes` 返回 true 后,该轮不再被抢占。 +- `outputTap()`:把主 session 的模型输出流逐段交给 World;`externalizes` 返回 true 后,该轮不再被 `preempt` 取消。 - `cognition`:World 向 Persona 请求后台认知计算。 - `llmStalls(withinMs)`:最近一段时间内模型调用失败或流中断的次数。 diff --git a/docs/worlds.md b/docs/worlds.md index 9104a060..7ecec973 100644 --- a/docs/worlds.md +++ b/docs/worlds.md @@ -19,6 +19,7 @@ World 不直接访问 Memory 或调用 Persona 的工具。 | `console?()` | 控制台页声明(见 [console.md](console.md)) | | `outputTap?()` | 主 session 输出流的接收器(演出、字幕);这一刻没有接收器时返回 `undefined` | | `onHandoffEnded?()`、`onTurnEnded?()` | 交接结束与主循环一轮结束的通知;隐藏的 World 不接收 | +| `onEventsSettled?(events, outcome)` | 本 World 的事件写入主 session(`delivered`)或被操作者清空队列丢弃(`discarded`) | | `shutdownVerification?()` | 关机前要核对的外部状态,同步只读快照 | Core 通过 `WorldHost` 向 World 提供以下能力: @@ -29,17 +30,28 @@ Core 通过 `WorldHost` 向 World 提供以下能力: | `pushDeferred(e, { trigger })` | 投递时生成正文;`render` 返回 null、抛错或超时时不存储、不投递 | | `pushCandidate?(spec, { trigger })` | 先归档原始事件,在投递时选择内容并生成正文 | | `store`、`drainPendingEvents(filter)` | 读事件库;消费待投递事件(一次性) | +| `withdrawPending?(cursor)` | 撤回本 World 一条未投递的事件;事件库追加撤回记录,重启不补投 | +| `promotePending?(cursor, trigger)` | 让本 World 一条未投递的事件按 `flush`、`preempt` 或 `interrupt` 立即触发;暂停或闸门挡着时,放行后随整批投递,不打断 | | `modelFacts` | 当前端点的模型名、接受的 MIME 与上下文窗口,每次调用按当前端点读取;未选端点或端点未选模型时 `accepts` 为 false | | `blob(handle)`、`reportUsage()`、`llmStalls?()` | 附件、用量上报、模型停滞查询 | | `cognition?` | 向 Persona 请求后台认知计算;Persona 未提供时该成员不存在 | | `log` | 包含 World 区域和调用关联字段的 Logger | -`trigger` 控制投递时机:`preempt` 请求中断当前模型调用并立即投递,已提交不可逆输出时不取消调用; -`flush` 立即投递并包含积压事件;`debounce` 参与合批;`piggyback` 仅排队,随其他触发产生的批次投递。 +`trigger` 控制投递时机: +- `preempt`:立即投递,取消尚未产生外部输出的模型调用并丢弃其输出,在同一批内带着新事件重新请求; + 已有外部输出或正在执行工具时不取消。 +- `interrupt`:立即投递并停止当前轮。模型调用一律取消,已外化的输出保留;执行中的 + `interruptible` 工具收到 `signal`,回执末尾注明执行期间发出过打断信号;本轮尚未开始的调用不执行。 +- `flush`:立即投递并包含积压事件。 +- `debounce`:参与合批。 +- `piggyback`:仅排队,随其他触发产生的批次投递。 + +主循环在轮次边界接收新事件:工具全部返回后,或 `preempt`、`interrupt` 取消模型调用后。 外部事件默认 `debounce`,内部事件默认 `flush`。`deliver: false` 仅存储事件。 被动变化通过事件报告,主动查询通过工具提供。当前状态快照使用 `pushDeferred` 在投递时读取, -配合 `piggyback` 随其他事件投递。只有需要立即中止当前模型轮的事件使用 `preempt`。 +配合 `piggyback` 随其他事件投递。新信息让还没出口的回答作废时用 `preempt`;用户要求停下正在做的事时 +用 `interrupt`。 World 应通过事件报告服务器起停、存档切换、连接变化等状态变更。 事件是 `EventEnvelope`:`cursor`、`run`、`type`、`ts`、`source`(World id)、`origin`、`tags`、 @@ -48,6 +60,8 @@ World 应通过事件报告服务器起停、存档切换、连接变化等状 外部消息使用 `origin: 'external'`。 工具回执 `ToolOutcome { text, blobs?, failed? }`;handler 抛错由 Core 转成失败回执。 +`interruptible` 的工具在 `interrupt` 到达时停止并返回已完成部分的回执,Core 等它返回; +未声明的工具执行到结束。 `endsTurn` 让一个工具结束本轮,`barrierAfter` 让流式提前派发在它之后停下。 工具名在一个 bot 内全局唯一:模型按名字调用,Core 按名字归属与隐藏。用自家短名做前缀 (`mc_`、`qq_`);工具名与已挂载 World、Persona 工具或 Core 保留帧名冲突时,装配层拒绝挂载并报告原因。 diff --git a/src/core/README.md b/src/core/README.md index 69df2390..75d2d497 100644 --- a/src/core/README.md +++ b/src/core/README.md @@ -43,7 +43,8 @@ Persona 通过 `CoreApi` 访问:`injectInternal` / `injectDeferred` / `injectE ## 总线与唤醒 -`WakeBus` 提供 `preempt`、`flush`、`debounce`、`piggyback` 四种触发模式。 +`WakeBus` 提供 `preempt`、`interrupt`、`flush`、`debounce`、`piggyback` 五种触发模式。 +`preempt` 与 `interrupt` 投递后通知主循环,取消范围由主循环裁决;`promote` 把排队项改为立即触发。 debounce 的计划投递时刻为 `min(首件时刻 + maxBatchAgeMs, max(首件时刻 + minBatchAgeMs, 末件时刻 + quietGapMs))`;计数达到 `maxBatchSize` 时立即投递。 计数包含外部即时事件与候选,不包含内部事件、延迟渲染项或 piggyback 项。 @@ -53,7 +54,8 @@ piggyback 只入队,随后续唤醒一起投递。 `nextBatch()` 仅支持一个消费者,每次按 FIFO 顺序取走整批。`batching` 使用共享配置引用, 更新后的值在下一次入队时参与计算。 -投递水位 `lastDeliveredCursor` 持久化,队列不持久化。重启时补投水位之后的外部事件; +投递水位 `lastDeliveredCursor` 持久化,队列不持久化。重启时补投水位之后的外部事件,跳过有 +`core.withdrawal` 撤回记录的事件; 已被候选处理结果引用的原始归档不重复投递。内部事件仅在当次运行投递。清空分片或跳过损坏行 留下的游标空位没有可投递内容,不阻止水位推进。 @@ -65,6 +67,10 @@ piggyback 只入队,随后续唤醒一起投递。 退回总线,进入下一批。控制台的前缀重载与清空 session 在批次边界执行:正在处理批次时等该批 结束,空闲时立即;并发请求复用同一事务。 +`preempt` 取消尚未外化的模型轮,`interrupt` 取消模型轮(已外化的输出保留)或停止 `interruptible` +工具、跳过本轮尚未开始的调用。两者取消后在同一批内接收新事件并开始下一轮,被取消的轮计入轮数, +不调用 `onTurnEnded`;到达时没有可取消的轮,事件在下一次模型请求前送入。 + 一批正文归档后,Core 调用 `Persona.onDelivery` 并等待它返回的 Promise,不设期限。完成前 `injectInternal` 的内容排在这批的内部行末尾、外部正文之前;完成后的注入进入总线。 diff --git a/src/core/bus.ts b/src/core/bus.ts index 4e19007d..d6eb588e 100644 --- a/src/core/bus.ts +++ b/src/core/bus.ts @@ -5,6 +5,7 @@ * 队列包含即时事件、延迟渲染项与候选项。延迟渲染项不参与关键词匹配或外部事件计数; * 候选项按外部事件计数,关键词只检查 gateText。 * piggyback 不触发计时、关键词或溢出,也不单独获得投递许可;需要其他项触发投递。 + * preempt 与 interrupt 投递后通知主循环;取消哪一部分由主循环裁决。 * 人工暂停阻止所有投递。DeliveryGate 解除、关键词回调授权或溢出授权均放行整批。 */ import type { DeliveryGate, EventOrigin, Logger, TriggerMode, WakeItem } from './types.ts'; @@ -45,8 +46,8 @@ export class WakeBus { private waiter: ((batch: WakeItem[]) => void) | null = null; /** 已到投递条件但没有消费者在等:置真,消费者一来就取走 */ private ready = false; - /** preempt 只报告机械时机;是否仍可安全取消由主循环裁决。 */ - private preemptHandler: (() => void) | null = null; + /** preempt 与 interrupt 只报告机械时机;取消范围由主循环裁决。 */ + private preemptHandler: ((trigger: 'preempt' | 'interrupt') => void) | null = null; private readonly log: Logger; @@ -55,7 +56,7 @@ export class WakeBus { this.log = log; } - setPreemptHandler(handler: () => void): void { + setPreemptHandler(handler: (trigger: 'preempt' | 'interrupt') => void): void { this.preemptHandler = handler; } @@ -97,7 +98,7 @@ export class WakeBus { return; } - if (trigger === 'preempt') { + if (trigger === 'preempt' || trigger === 'interrupt') { const permitted = !this.blocked(); this.deliver(); this.notifyPreempt(trigger, permitted); @@ -136,14 +137,29 @@ export class WakeBus { return this.paused || (this.gate !== null && !this.bypassGateOnce); } - /** 已获准投递的 preempt 才能取消当前模型轮;暂停与闸门继续拥有更高优先级。 */ + /** 已获准投递的 preempt 与 interrupt 才通知主循环;暂停与闸门继续拥有更高优先级。 */ private notifyPreempt(trigger: TriggerMode, permitted: boolean): void { - if (trigger === 'preempt' && permitted && !this.paused) { - this.log.emit('debug', '抢占:请求取消尚未外化的在途模型轮', { event: 'preempt' }); - this.preemptHandler?.(); + if ((trigger === 'preempt' || trigger === 'interrupt') && permitted && !this.paused) { + this.log.emit('debug', trigger === 'preempt' ? '抢占:请求取消尚未外化的在途模型轮' : '打断:请求停止当前轮', { event: trigger }); + this.preemptHandler?.(trigger); } } + /** + * 把首个匹配的排队项改为按 trigger 立即投递;piggyback 项随之可单独触发投递。 + * 暂停或闸门未放行时只改标志,放行时随整批投递,不再通知主循环。没有匹配项时返回 false。 + */ + promote(pred: (item: WakeItem) => boolean, trigger: 'flush' | 'preempt' | 'interrupt'): boolean { + const queued = this.queue.find((q) => pred(q.item)); + if (!queued) return false; + queued.piggyback = false; + if (this.gate !== null && !this.bypassGateOnce) return true; + const permitted = !this.blocked(); + this.deliver(this.bypassGateOnce); + this.notifyPreempt(trigger, permitted); + return true; + } + /** 人工暂停/继续(控制台);继续时积压一次性投递 */ setPaused(v: boolean): void { if (v !== this.paused) this.log.emit('info', v ? '总线已暂停:事件照常落库,不投递' : '总线继续', { event: v ? 'paused' : 'resumed', data: { queued: this.queue.length } }); diff --git a/src/core/core.ts b/src/core/core.ts index a390a593..b28065ed 100644 --- a/src/core/core.ts +++ b/src/core/core.ts @@ -192,8 +192,9 @@ export class Core { transcript: this.transcript, toolOwner: (name) => this.worlds.find((m) => m.tools().some((t) => t.name === name))?.id, }); - this.bus.setPreemptHandler(() => { - this.loop.abortCurrentRound(); + this.bus.setPreemptHandler((trigger) => { + if (trigger === 'interrupt') this.loop.interruptCurrentRound(); + else this.loop.abortCurrentRound(); }); } @@ -255,9 +256,9 @@ export class Core { discardPendingEvents(): number { const dropped = this.bus.drainPending((item) => item.event !== undefined || item.candidate !== undefined); - this.loop.acknowledgeDiscarded( - dropped.flatMap((item) => item.event ? [item.event] : []), - ); + const events = dropped.flatMap((item) => item.event ? [item.event] : []); + this.loop.acknowledgeDiscarded(events); + this.loop.notifySettled(events, 'discarded'); return dropped.length; } @@ -557,6 +558,15 @@ export class Core { if (taken.length > 0) this.loop.acknowledgeDiscarded(taken); return taken; }, + withdrawPending: async (cursor) => { + if (!lease.active) return false; + const [taken] = this.bus.drainPending((it) => it.event?.cursor === cursor && it.event.source === mod.id); + if (!taken?.event) return false; + this.loop.recordWithdrawn(taken.event); + return true; + }, + promotePending: async (cursor, trigger) => lease.active + && this.bus.promote((it) => it.event?.cursor === cursor && it.event.source === mod.id, trigger), modelFacts: this.modelFacts(), blob: (handle) => this.resolveBlob(handle), reportUsage: (usage, opts) => { diff --git a/src/core/loop.ts b/src/core/loop.ts index ded1fa0e..6c769a05 100644 --- a/src/core/loop.ts +++ b/src/core/loop.ts @@ -27,6 +27,7 @@ import type { LLMUsage, Logger, ModelSpec, + OutputTap, Persona, SessionDecl, SessionOpeningReason, @@ -37,6 +38,7 @@ import type { ToolTag, WakeItem, } from './types.ts'; +import { defaultExternalizes } from './types.ts'; import type { WakeBus } from './bus.ts'; import type { SessionLog } from './session.ts'; import type { CoreState } from './state.ts'; @@ -47,9 +49,11 @@ import type { Transcript } from './transcript.ts'; import { setAnchors, withAnchors } from './log-context.ts'; import { withBlobLines } from './blobs.ts'; import { + INTERRUPTED_WHILE_RUNNING, MISSING_RESULT_RESTART, NOT_EXECUTED_BARRIER, NOT_EXECUTED_INCOMPLETE, + NOT_EXECUTED_INTERRUPTED, NOT_EXECUTED_LOOP_STOPPED, NOT_EXECUTED_STREAM_ABORTED, SHUTDOWN_INTERRUPTED, @@ -200,6 +204,9 @@ function eventBlobs(events: readonly EventEnvelope[]): BlobRef[] { return events.flatMap((e) => e.blobs ?? []); } +/** 撤回记录的事件类型;meta.cursor 指向被撤回的事件。 */ +const WITHDRAWAL_TYPE = 'core.withdrawal'; + /** 重启补投的数量上限;仅补投最近事件,更早的事件标记已处理并推进水位。 */ const MAX_REQUEUE = 200; @@ -228,6 +235,20 @@ interface PreparedCandidateProjection { sourceCursors: number[]; } +/** 主循环的一轮:模型请求与随后的工具执行。 */ +interface Flight { + /** 取消模型请求。 */ + controller: AbortController; + /** 合并进 interruptible 工具的 signal。 */ + interrupt: AbortController; + phase: 'model' | 'tools' | 'done'; + /** 非 reasoning 增量已经外流;此后 preempt 不再取消本轮。 */ + externalized: boolean; + /** 本轮收到 interrupt:尚未开始的工具调用不执行。 */ + interrupted: boolean; + abortReason: 'preempt' | 'interrupt' | 'shutdown' | null; +} + export class MainLoop { private d: MainLoopDeps; private running = false; @@ -292,12 +313,13 @@ export class MainLoop { private appliedVisibleWorlds: Set | null = null; /** 当前前缀中各可见 World 的工具签名,用于检测工具表漂移。 */ private appliedWorldTools = new Map(); - /** 当前模型轮;非 reasoning 增量一旦外流,本轮不再接受自动抢占。 */ - private currentRound: { - controller: AbortController; - externalized: boolean; - abortReason: 'preempt' | 'shutdown' | null; - } | null = null; + /** + * 当前轮,从模型请求到工具全部返回。model 阶段可被 preempt 取消(尚未外化时)或被 interrupt + * 取消;tools 阶段只有 interrupt 生效,停止 interruptible 工具并跳过尚未开始的调用。 + */ + private currentRound: Flight | null = null; + /** preempt 或 interrupt 到达时没有可取消的轮:下一次模型请求前先接收已就绪的事件。 */ + private takeReadyBeforeRequest = false; /** 每次 stop 都使此前捕获的异步 continuation 永久失效。 */ private generation = 0; private stopped = false; @@ -307,10 +329,7 @@ export class MainLoop { private pendingToolCalls = new Set(); /** stop() 提前结束重试等待,由 active() 决定退出。 */ private backoffWake: (() => void) | null = null; - /** - * 工具信号合并关机信号与当前模型调用信号。 - * 自动抢占仅发生在没有外部输出、尚未执行工具时;工具执行期间仅关机或循环换代会取消。 - */ + /** 关机或循环换代时取消全部工具;interruptible 工具另外合并当前轮的 interrupt 信号。 */ private readonly shutdown = new AbortController(); constructor(deps: MainLoopDeps) { @@ -326,14 +345,40 @@ export class MainLoop { return this.active(this.generation); } - /** 尝试取消尚未输出的当前模型调用;无可取消调用时返回 false。 */ + /** preempt:取消尚未输出的当前模型调用;无可取消调用时返回 false。 */ abortCurrentRound(): boolean { if (!this.activeNow()) return false; const round = this.currentRound; - if (!round || round.externalized) return false; - round.abortReason = 'preempt'; - round.controller.abort(new Error('模型轮被新输入抢占')); - return true; + if (round?.phase === 'model' && !round.externalized) { + round.abortReason = 'preempt'; + round.controller.abort(new Error('模型轮被新输入抢占')); + return true; + } + if (this.processingBatch) this.takeReadyBeforeRequest = true; + return false; + } + + /** + * interrupt:取消当前模型调用(已外化的输出保留),或停止执行中的 interruptible 工具并跳过 + * 本轮尚未开始的调用。没有进行中的轮时返回 false。 + */ + interruptCurrentRound(): boolean { + if (!this.activeNow()) return false; + const round = this.currentRound; + if (round?.phase === 'model') { + round.abortReason = round.externalized ? 'interrupt' : 'preempt'; + round.interrupted = true; + round.interrupt.abort(new Error('被新事件打断')); + round.controller.abort(new Error('模型轮被新事件打断')); + return true; + } + if (round?.phase === 'tools' && !round.interrupted) { + round.interrupted = true; + round.interrupt.abort(new Error('被新事件打断')); + return true; + } + if (this.processingBatch) this.takeReadyBeforeRequest = true; + return false; } /** @@ -582,6 +627,7 @@ export class MainLoop { } if (events.length > 0 && !inUser) this.appendEventFrame(events, generation); this.noteHandled(delivered, generation); + this.notifySettled(delivered, 'delivered'); const changed = lines.length > 0 || events.length > 0; if (changed) this.batchesHandled++; return changed; @@ -768,15 +814,23 @@ export class MainLoop { } const after = store.range({ fromCursor: from }); const referenced = new Set(); + const withdrawn = new Set(); for (const event of after) { + if (event.source === 'core' && event.type === WITHDRAWAL_TYPE) { + const cursor = (event.meta as { cursor?: unknown } | undefined)?.cursor; + if (typeof cursor === 'number') withdrawn.add(cursor); + continue; + } if (event.contextDelivery !== 'deliver') continue; const cursors = (event.meta as { sourceCursors?: unknown } | undefined)?.sourceCursors; if (!Array.isArray(cursors)) continue; for (const c of cursors) if (typeof c === 'number') referenced.add(c); } for (const c of referenced) this.settledArchives.add(c); + // 撤回的事件按已了结处理,水位可越过。 + for (const c of withdrawn) this.deliveredCursors.add(c); this.noteHandled([], this.generation); - const pending = after.filter((event) => event.origin === 'external' + const pending = after.filter((event) => event.origin === 'external' && !withdrawn.has(event.cursor) && (event.contextDelivery !== 'archive-only' || !referenced.has(event.cursor))); if (pending.length === 0) return; @@ -882,6 +936,7 @@ export class MainLoop { await this.maintenanceChain; if (!this.active(generation)) break; this.processingBatch = true; + this.takeReadyBeforeRequest = false; try { const changed = await this.deliverBatch(batch, generation); if (!this.active(generation)) break; @@ -920,7 +975,21 @@ export class MainLoop { /** 本批模型调用共享 sess 关联字段;每轮更新 round、resp 和 call。 */ private rounds(generation: number): Promise { - return withAnchors({ sess: this.d.decl.id }, () => this.roundsInScope(generation)); + return withAnchors({ sess: this.d.decl.id }, async () => { + try { + await this.roundsInScope(generation); + } finally { + this.currentRound = null; + } + }); + } + + private tapAbort(tap: OutputTap | undefined, reason: string): void { + try { + tap?.onAbort?.(reason); + } catch (tapErr) { + this.d.log.warn('outputTap.onAbort异常', { err: tapErr }); + } } private async roundsInScope(generation: number): Promise { @@ -943,6 +1012,14 @@ export class MainLoop { log.warn('模型配置不可用', { error: String(error) }); return; } + if (this.takeReadyBeforeRequest) { + this.takeReadyBeforeRequest = false; + const ready = bus.takeIfReady(); + if (ready) { + await this.deliverBatch(ready, generation); + if (!this.active(generation)) return; + } + } // 后续轮输入超限时结束本批,由批末检查执行交接。 // 首轮仍处理本批新投递的事件;上一批的容量检查已在批末执行。 if (round > 1) { @@ -957,16 +1034,17 @@ export class MainLoop { const queuedEvents: EventEnvelope[] = []; const roundNo = ++this.roundSeq; setAnchors({ round: roundNo, resp: undefined, call: undefined }); - const flight: NonNullable = { - controller: new AbortController(), externalized: false, abortReason: null, + const flight: Flight = { + controller: new AbortController(), interrupt: new AbortController(), phase: 'model', + externalized: false, interrupted: false, abortReason: null, }; this.currentRound = flight; const ctx: ToolCallContext = { role: decl.id, log, round: roundNo, - // 关机或循环换代取消工具;自动抢占不取消工具。 - signal: AbortSignal.any([flight.controller.signal, this.shutdown.signal]), + // interruptible 工具在 runToolHandler 里再合并本轮的 interrupt 信号。 + signal: this.shutdown.signal, queueExternalEvents: (events) => { if (this.active(generation)) queuedEvents.push(...events); }, @@ -976,10 +1054,28 @@ export class MainLoop { const eager = tap ? new EagerDispatch( () => this.toolDefs, ctx, log, decl.id, this.d.toolLog, - () => this.active(generation) && !flight.controller.signal.aborted, + () => this.active(generation), (name) => this.d.toolOwner?.(name), + flight.interrupt.signal, + () => (flight.interrupted ? NOT_EXECUTED_INTERRUPTED : null), ) : null; + /** + * preempt 或 interrupt 取消模型轮后,在同一批内接收新事件并开始下一轮;被取消的轮计入轮数。 + * 没有就绪事件或已到硬上限时结束本次唤醒,就绪事件留给下一批。 + */ + const resumeAfterCancel = async (arrivedEarly: EventEnvelope[]): Promise => { + flight.phase = 'done'; + const ready = round < caps.hard ? bus.takeIfReady() : null; + const arrived: WakeItem[] = [...arrivedEarly.map((event) => ({ event })), ...(ready ?? [])]; + if (!ready) { + for (const item of arrived) bus.push(item, { trigger: 'flush' }); + this.finishTurn(); + return false; + } + await this.deliverBatch(arrived, generation); + return this.active(generation); + }; // 轮级观测:首个内容事件的延迟、模型往返、工具阻塞,一轮一条 debug 记录(event=round)。 const roundStart = Date.now(); let ttftMs: number | null = null; @@ -995,9 +1091,7 @@ export class MainLoop { if (!this.active(generation) || flight.controller.signal.aborted) return; if (event.type === 'response.created') setAnchors({ resp: event.response.id }); if (ttftMs === null && ('delta' in event || event.type === 'response.output_item.added')) ttftMs = Date.now() - roundStart; - const tappedEffect = tap.externalizes ? tap.externalizes(event) - : event.type === 'response.output_text.delta' || event.type === 'response.refusal.delta' - || (event.type === 'response.output_item.added' && event.item?.type === 'function_call'); + const tappedEffect = tap.externalizes ? tap.externalizes(event) : defaultExternalizes(event); if (tappedEffect || (event.type === 'response.output_item.done' && event.item?.type === 'function_call' && event.item.status === 'completed')) flight.externalized = true; eager?.onEvent(event); try { tap.onEvent(event); } catch (error) { log.warn('outputTap.onEvent异常', { err: error }); } @@ -1025,13 +1119,25 @@ export class MainLoop { assistant = responseRecords(res.response, res.origin); setAnchors({ resp: res.response.id }); if (!this.active(generation) || flight.controller.signal.aborted) { - // 关机或换代丢弃已成功返回的结果时,仍记录这次调用的实际用量。 + if (this.active(generation) && flight.abortReason === 'interrupt') { + // 打断与响应完成同时到达:完整响应作为已外化的输出保存,尚未执行的调用不执行。 + this.mainTrack?.recordAttempts(res.attempts, undefined, { prefixHash }); + await this.recordAbortedStream(assistant, eager, generation, NOT_EXECUTED_INTERRUPTED); + if (!this.active(generation)) return; + this.tapAbort(tap, '模型轮被新事件打断'); + noteRound('interrupted'); + if (await resumeAfterCancel(queuedEvents)) continue; + return; + } + // 关机、换代或抢占丢弃已成功返回的结果时,仍记录这次调用的实际用量。 this.mainTrack?.recordAttempts(res.attempts, undefined, { outcome: 'discarded', prefixHash }); - try { - tap?.onAbort?.('core 正在关机'); - } catch (tapErr) { - log.warn('outputTap.onAbort异常', { err: tapErr }); + if (this.active(generation) && flight.abortReason === 'preempt') { + this.tapAbort(tap, '模型轮被新输入抢占'); + noteRound('preempted'); + if (await resumeAfterCancel([])) continue; + return; } + this.tapAbort(tap, 'core 正在关机'); noteRound('discarded'); this.finishTurn(); return; @@ -1044,14 +1150,23 @@ export class MainLoop { } catch (e) { if (flight.controller.signal.aborted) { this.recordFailedUsage(e, prefixHash); - try { - tap?.onAbort?.(flight.abortReason === 'shutdown' ? 'core 正在关机' : '模型轮被新输入抢占'); - } catch (tapErr) { - log.warn('outputTap.onAbort异常', { err: tapErr }); + if (!this.active(generation) || flight.abortReason === 'shutdown') { + this.tapAbort(tap, 'core 正在关机'); + log.info('模型轮随关机终止'); + noteRound('shutdown'); + this.finishTurn(); + return; } - log.info(flight.abortReason === 'shutdown' ? '模型轮随关机终止' : '尚未外化的模型轮已被新输入抢占'); - noteRound(flight.abortReason === 'shutdown' ? 'shutdown' : 'preempted'); - this.finishTurn(); + const interrupted = flight.abortReason === 'interrupt'; + if (interrupted && e instanceof GenerationError && e.partial) { + // 已外化的部分输出保存;提前执行的调用使用真实结果,其余调用不执行。 + await this.recordAbortedStream(responseRecords(e.partial, e.origin), eager, generation, NOT_EXECUTED_INTERRUPTED); + if (!this.active(generation)) return; + } + this.tapAbort(tap, interrupted ? '模型轮被新事件打断' : '模型轮被新输入抢占'); + log.info(interrupted ? '模型轮被新事件打断,已外化的输出保留' : '尚未外化的模型轮已被新输入抢占'); + noteRound(interrupted ? 'interrupted' : 'preempted'); + if (await resumeAfterCancel(queuedEvents)) continue; return; } this.recordFailedUsage(e, prefixHash); @@ -1116,7 +1231,7 @@ export class MainLoop { this.finishTurn(); return; } finally { - if (this.currentRound === flight) this.currentRound = null; + if (flight.phase === 'model') flight.phase = 'tools'; } if (!this.active(generation)) return; assistant = this.dropReservedCalls(assistant); @@ -1159,6 +1274,10 @@ export class MainLoop { results.push(functionResult(call.call_id, NOT_EXECUTED_BARRIER)); continue; } + if (flight.interrupted && !eager?.has(call.call_id)) { + results.push(functionResult(call.call_id, NOT_EXECUTED_INTERRUPTED)); + continue; + } let out: ToolOutcome; const def = this.toolDefs.find((t) => t.name === call.name); @@ -1182,7 +1301,7 @@ export class MainLoop { } out = await runToolHandler( def, args, ctx, call.call_id, decl.id, this.d.toolLog, - () => this.active(generation), this.d.toolOwner?.(def.name), + () => this.active(generation), this.d.toolOwner?.(def.name), flight.interrupt.signal, ); if (!this.active(generation)) return; } @@ -1195,6 +1314,7 @@ export class MainLoop { results.push(this.toolResult(call.call_id, out)); } toolMs = Date.now() - toolsStart; + flight.phase = 'done'; if (!this.active(generation)) return; if (round === caps.soft && results.length > 0) { @@ -1253,6 +1373,7 @@ export class MainLoop { partial: ContextRecord[], eager: EagerDispatch | null, generation: number, + notExecuted = NOT_EXECUTED_STREAM_ABORTED, ): Promise { if (!this.active(generation)) return; const { session } = this.d; @@ -1263,7 +1384,7 @@ export class MainLoop { for (const entry of partial) session.append(entry); for (const call of calls) { const ran = eager?.take(call.call_id); - const out = ran !== undefined ? await ran : { text: NOT_EXECUTED_STREAM_ABORTED }; + const out = ran !== undefined ? await ran : { text: notExecuted }; if (!this.active(generation)) return; session.append(this.toolResult(call.call_id, out)); this.pendingToolCalls.delete(call.call_id); @@ -1276,6 +1397,38 @@ export class MainLoop { this.noteHandled(events, this.generation); } + /** + * 事件库追加撤回记录并结清水位;重启补投按记录跳过该事件。事件此时已离开队列,循环停止后 + * 仍写记录,水位由重启时按记录结清。 + */ + recordWithdrawn(event: EventEnvelope): void { + const { store, cfg } = this.d; + store.append({ + type: WITHDRAWAL_TYPE, + ts: nowIso(cfg.timezone), + source: 'core', + origin: 'internal', + text: `withdrawn #${event.cursor}`, + meta: { cursor: event.cursor }, + }); + this.noteHandled([event], this.generation); + } + + /** 按事件来源通知产生它们的 World;钩子异常记 warn。 */ + notifySettled(events: readonly EventEnvelope[], outcome: 'delivered' | 'discarded'): void { + if (!this.activeNow() || events.length === 0) return; + for (const world of this.d.worlds.all()) { + if (!world.onEventsSettled) continue; + const own = events.filter((e) => e.source === world.id); + if (own.length === 0) continue; + try { + world.onEventsSettled(own, outcome); + } catch (e) { + this.d.log.warn('World 事件结清钩子异常', { id: world.id, err: e }); + } + } + } + /** 按请求内容估算,排除不会回传的历史推理。 */ private estimateOutbound(msgs: readonly ContextRecord[]): number { const view = this.d.cfg.context.keepPastThinking ? msgs : withoutPastReasoning(msgs); @@ -1805,15 +1958,21 @@ function runToolHandler( toolLog?: ToolCallLog, canRecord: () => boolean = () => true, mod?: string, + interrupt?: AbortSignal, ): Promise { const startedAt = Date.now(); + const signal = def.interruptible && interrupt && ctx.signal ? AbortSignal.any([ctx.signal, interrupt]) : ctx.signal; return withAnchors({ call: callId }, () => Promise.resolve() - .then(() => def.handler(args, { ...ctx, callId })) + .then(() => def.handler(args, { ...ctx, callId, signal })) .then((out): ToolOutcome => (typeof out === 'string' ? { text: out } : out)) .catch((e: unknown): ToolOutcome => ({ text: toolFailed(e instanceof Error ? e.message : String(e)), failed: true, })) + .then((out): ToolOutcome => (def.interruptible && interrupt?.aborted + ? { ...out, text: `${out.text} +${INTERRUPTED_WHILE_RUNNING}` } + : out)) .then((out): ToolOutcome => { if (canRecord()) recordToolCall(toolLog, role, def.name, args, startedAt, out, mod); return out; @@ -1839,6 +1998,9 @@ class EagerDispatch { private readonly toolLog?: ToolCallLog, private readonly active: () => boolean = () => true, private readonly owner: (name: string) => string | undefined = () => undefined, + private readonly interrupt?: AbortSignal, + /** 轮到执行时返回未执行标记则跳过该调用。 */ + private readonly skip: () => string | null = () => null, ) {} @@ -1857,7 +2019,7 @@ class EagerDispatch { } private dispatch(call: { id: string; name: string; args: string }): void { - if (!this.active()) return; + if (!this.active() || this.skip() !== null) return; if (this.barrierHit) return; // 保留帧调用不执行,也不写入 session。 if (RESERVED_FRAME_NAMES.has(call.name)) return; @@ -1874,13 +2036,20 @@ class EagerDispatch { const args = parseToolArgs(call.args); if (args === null) return; // 落回非流式路径的"arguments are not valid JSON"回执 // handler 按调用顺序串行执行。 - const run = this.chain.then(() => this.active() - ? runToolHandler(def, args, this.ctx, call.id, this.role, this.toolLog, this.active, this.owner(def.name)) - : { text: NOT_EXECUTED_LOOP_STOPPED }); + const run = this.chain.then(() => { + const skipped = this.active() ? this.skip() : NOT_EXECUTED_LOOP_STOPPED; + return skipped === null + ? runToolHandler(def, args, this.ctx, call.id, this.role, this.toolLog, this.active, this.owner(def.name), this.interrupt) + : { text: skipped }; + }); this.chain = run.then(() => undefined); this.results.set(call.id, run); } + has(callId: string): boolean { + return this.results.has(callId); + } + /** 取走某次调用的执行结果;没提前派发过返回 undefined(一次性,防重复配对) */ take(callId: string): Promise | undefined { const p = this.results.get(callId); diff --git a/src/core/markers.ts b/src/core/markers.ts index 0bc9c6eb..e11d0742 100644 --- a/src/core/markers.ts +++ b/src/core/markers.ts @@ -31,6 +31,10 @@ export const NOT_EXECUTED_BARRIER = '[not executed: review the preceding tool re export const NOT_EXECUTED_THREAD_ENDED = '[not executed: this thread already ended]'; /** 流在响应中途断了,提前派发的调用没跑完 */ export const NOT_EXECUTED_STREAM_ABORTED = '[not executed: stream aborted mid-response]'; +/** interrupt 事件到达时本轮尚未开始的调用 */ +export const NOT_EXECUTED_INTERRUPTED = '[not executed: interrupted by a new event]'; +/** interruptible 工具执行期间发出过 interrupt 信号,接在它的回执之后;做到哪一步由回执正文说明 */ +export const INTERRUPTED_WHILE_RUNNING = '[an interrupt signal was sent while this tool ran]'; /** 主循环已停 */ export const NOT_EXECUTED_LOOP_STOPPED = '[not executed: main loop stopped]'; /** 关机打断了这一轮,回执取不回来了 */ diff --git a/src/core/types.ts b/src/core/types.ts index 317759ab..edac46ac 100644 --- a/src/core/types.ts +++ b/src/core/types.ts @@ -215,8 +215,8 @@ export interface ToolCallContext { */ round?: number; /** - * 宿主放弃本次调用时触发,跨进程代理超时也会触发。 - * 耗时工具必须在提交外部副作用前检查此信号。 + * 宿主放弃本次调用时触发:关机、循环换代、跨进程代理超时;ToolDef.interruptible 的工具 + * 还会在 interrupt 事件到达时触发。耗时工具必须在提交外部副作用前检查此信号。 */ signal?: AbortSignal; } @@ -246,6 +246,11 @@ export interface ToolDef extends ToolSchema { * 通常同时设置 barrierAfter,使同一条输出中的后续调用得到未执行回执。 */ endsTurn?: boolean; + /** + * interrupt 事件到达时 ctx.signal 触发,handler 停止并返回已完成部分的回执;Core 等它返回, + * 并在回执末尾注明执行期间发出过打断信号。缺省时工具执行到结束,interrupt 只跳过本轮尚未开始的调用。 + */ + interruptible?: boolean; handler: (args: Record, ctx: ToolCallContext) => Promise; } @@ -416,12 +421,21 @@ export interface SessionDecl { /** 接收模型输出流;Core 不解释增量中的语义。 */ export interface OutputTap { onEvent(event: StreamEvent): void; - /** 返回 true 表示增量已产生外部输出,此后主循环不再允许抢占该轮。 */ + /** + * 返回 true 表示增量已产生外部输出,此后 preempt 不再取消该轮。未实现时按 + * defaultExternalizes 判定;合并多个接收器时逐个按同一规则判定。 + */ externalizes?(event: StreamEvent): boolean; onRoundEnd?(): void; onAbort?(reason: string): void; } +/** 接收器未声明 externalizes 时的判定:正文或拒答增量,或出现函数调用。 */ +export function defaultExternalizes(event: StreamEvent): boolean { + return event.type === 'response.output_text.delta' || event.type === 'response.refusal.delta' + || (event.type === 'response.output_item.added' && event.item?.type === 'function_call'); +} + export interface ForkOptions { /** 已声明的 session id;默认工具和轮数上限来自该声明,模型来自当前 provider。 */ id: string; @@ -557,10 +571,22 @@ export interface ModelFacts { contextWindow(): number | undefined; } -/** 事件的投递触发方式,由产生方选择;与事件的渲染时机独立。 */ +/** + * 事件的投递触发方式,由产生方选择;与事件的渲染时机独立。 + * 主循环在一批之内的轮次边界接收新事件:工具全部返回后,或 preempt、interrupt 取消模型轮后。 + * 两者到达时没有可取消的模型轮,事件在下一次模型请求前送入。 + */ export type TriggerMode = - /** 立即投递整批,并请求取消尚未产生外部输出的在途模型轮;是否可取消由主循环判断。 */ + /** + * 立即投递整批,取消尚未产生外部输出的在途模型轮,丢弃该轮输出,在同一批内带着新事件重新请求。 + * 已有外部输出或正在执行工具时不取消。 + */ | 'preempt' + /** + * 立即投递整批,并停止当前轮:模型调用无论有无外部输出都取消,已有输出保留;执行中的 + * interruptible 工具收到 signal;本轮尚未开始的工具调用不执行。新事件接在本轮回执之后。 + */ + | 'interrupt' /** 立即投递整批积压。 */ | 'flush' /** 按 quietGapMs、minBatchAgeMs、maxBatchAgeMs 和 maxBatchSize 合批。 */ @@ -656,6 +682,17 @@ export interface WorldHost { * 取走后不会在后续批次重复投递。 */ drainPendingEvents(filter: (e: EventEnvelope) => boolean): Promise; + /** + * 撤回本 World 一条尚未投递的事件:移出队列,事件库追加撤回记录,重启不再补投。 + * 事件已进入投递或不在队列中时返回 false。 + */ + withdrawPending?(cursor: number): Promise; + /** + * 让本 World 一条尚未投递的事件按 flush、preempt 或 interrupt 立即触发投递。人工暂停或投递闸门 + * 未放行时只取消它的 piggyback,放行后随整批投递,不按 trigger 取消模型轮或工具。事件已进入投递 + * 或不在队列中时返回 false。 + */ + promotePending?(cursor: number, trigger: 'flush' | 'preempt' | 'interrupt'): Promise; /** 当前模型的能力;渲染层据 MIME 支持决定是否附加二进制内容。 */ modelFacts: ModelFacts; /** @@ -853,6 +890,12 @@ export interface World { * 隐藏 World 不接收通知。 */ onTurnEnded?(): void; + /** + * 本 World 产生的事件离开队列时调用:delivered 表示已写入主 session,discarded 表示被操作者 + * 清空队列丢弃。经 withdrawPending 撤回的事件不通知;drainPendingEvents 取走的事件只在 + * 之后经 queueExternalEvents 写入 session 时通知。 + */ + onEventsSettled?(events: readonly EventEnvelope[], outcome: 'delivered' | 'discarded'): void; /** * stop() 完成后返回外部状态检查的同步只读快照。 * 网络检查须在 stop() 的既有期限内完成并缓存。 diff --git a/tests/core/bus.test.ts b/tests/core/bus.test.ts index 7bc82291..b1379315 100644 --- a/tests/core/bus.test.ts +++ b/tests/core/bus.test.ts @@ -276,6 +276,17 @@ describe('WakeBus', () => { expect(preempts).toBe(1); }); + it('promote 把排队中的 piggyback 项改为 interrupt:立即投递并通知主循环,已取走的项返回 false', async () => { + const bus = new WakeBus({ quietGapMs: 10_000, minBatchAgeMs: 0, maxBatchAgeMs: 60_000, maxBatchSize: 100 }); + const notified: string[] = []; + bus.setPreemptHandler((trigger) => { notified.push(trigger); }); + bus.push(evt(1), { trigger: 'piggyback' }); + expect(bus.promote((it) => it.event?.cursor === 1, 'interrupt')).toBe(true); + expect(notified).toEqual(['interrupt']); + expect(await bus.nextBatch()).toEqual([evt(1)]); + expect(bus.promote((it) => it.event?.cursor === 1, 'interrupt')).toBe(false); + }); + it('preempt 被关键词闸门放行后通知主循环,整批顺序不变', async () => { const bus = new WakeBus({ quietGapMs: 10_000, minBatchAgeMs: 0, maxBatchAgeMs: 60_000, maxBatchSize: 100 }); let preempts = 0; diff --git a/tests/core/loop.test.ts b/tests/core/loop.test.ts index 7dfce04f..b80c8ab4 100644 --- a/tests/core/loop.test.ts +++ b/tests/core/loop.test.ts @@ -16,6 +16,7 @@ import { ToolCallLog } from '../../src/core/tool-log.ts'; import { Transcript } from '../../src/core/transcript.ts'; import { LLMError, LLMStreamAborted } from './fixture-errors.ts'; import { SessionTracker } from '../../src/core/sessions.ts'; +import { INTERRUPTED_WHILE_RUNNING, NOT_EXECUTED_INTERRUPTED } from '../../src/core/markers.ts'; import type { UsageRecord, Persona } from '../../src/core/types.ts'; import type { ChatMessage, LLMDelta } from './fixture-types.ts'; import type { Logger, CandidateProjector, EventEnvelope, World, WorldHost, ToolDef } from '../../src/core/types.ts'; @@ -395,6 +396,119 @@ describe('MainLoop preempt', () => { expect(cancelled).toBe(false); await rig.cleanup(); }); + + it('被抢占的轮不结束本次唤醒:不调用 onTurnEnded,新输入在同一批内送入下一次请求', async () => { + let turnEnds = 0; + const rig = makeRig({ outputTap: { onDelta: () => {} }, hooks: { onTurnEnded: () => { turnEnds++; } } }); + rig.bus.setPreemptHandler(() => { rig.loop.abortCurrentRound(); }); + const chat = rig.llm.chat.bind(rig.llm); + rig.llm.chat = async (spec, messages, tools, opts) => { + if (rig.llm.calls.length === 1) rig.llm.blockUntilAbort = true; + return chat(spec, messages, tools, opts); + }; + rig.llm.script(toolReply([{ name: 'noop' }])); + try { + rig.start(); + rig.pushEvent('先做事'); + await until(() => rig.llm.calls.length === 2); + const cut = rig.store.append({ type: 'qq.message', ts: '2026-07-17T10:01:00+08:00', source: 'qq', origin: 'external', text: '插一句' }); + rig.bus.push({ event: cut }, { trigger: 'preempt' }); + await until(() => rig.llm.calls.length === 3 && rig.loop.getStatus().batchesHandled > 0 && turnEnds > 0); + await sleep(30); + expect(turnEnds).toBe(1); + expect(rig.loop.getStatus().roundsLastBatch).toBe(3); + expect(JSON.stringify(rig.llm.calls[2].messages)).toContain('插一句'); + } finally { await rig.cleanup(); } + }); +}); + +describe('MainLoop interrupt', () => { + function wireTriggers(rig: ReturnType): void { + rig.bus.setPreemptHandler((trigger) => { + if (trigger === 'interrupt') rig.loop.interruptCurrentRound(); + else rig.loop.abortCurrentRound(); + }); + } + + it('停下执行中的 interruptible 工具,同轮尚未开始的调用不执行,新事件接在回执之后', async () => { + const ran: string[] = []; + const settled: string[] = []; + const walk: ToolDef = { + ...makeTool('walk', ''), + interruptible: true, + handler: async (_args, ctx) => { + ran.push('walk'); + await new Promise((resolve) => ctx.signal!.addEventListener('abort', () => resolve(), { once: true })); + return '走到一半停下'; + }, + }; + const wave = makeTool('wave', () => { ran.push('wave'); return 'ok'; }); + const pet: World = { + ...makeFakeIO('pet', [walk, wave]), + onEventsSettled: (events, outcome) => { settled.push(...events.map((e) => `${outcome}:${e.text}`)); }, + }; + const rig = makeRig({ worlds: [pet] }); + wireTriggers(rig); + rig.llm.script(toolReply([{ name: 'walk', id: 'c_walk' }, { name: 'wave', id: 'c_wave' }])); + try { + rig.start(); + rig.pushEvent('出发'); + await until(() => ran.includes('walk')); + const stop = rig.store.append({ type: 'pet.message', ts: '2026-07-17T10:01:00+08:00', source: 'pet', origin: 'external', text: '别走了' }); + rig.bus.push({ event: stop }, { trigger: 'interrupt' }); + await until(() => rig.llm.calls.length >= 2); + const sent = rig.llm.calls[1].messages; + const receipt = (id: string) => sent.findIndex((m) => m.role === 'tool' && m.tool_call_id === id); + expect(ran).toEqual(['walk']); + expect(sent[receipt('c_walk')].content).toBe(`走到一半停下\n${INTERRUPTED_WHILE_RUNNING}`); + expect(sent[receipt('c_wave')].content).toBe(NOT_EXECUTED_INTERRUPTED); + expect(sent.findIndex((m) => String(m.content ?? '').includes('别走了'))).toBeGreaterThan(receipt('c_wave')); + expect(settled).toEqual(['delivered:别走了']); + } finally { await rig.cleanup(); } + }); + + it('取消已外化的模型轮:已输出的正文保留,新事件在同一批内送入下一次请求', async () => { + const rig = makeRig({ outputTap: { onDelta: () => {} } }); + wireTriggers(rig); + const chat = rig.llm.chat.bind(rig.llm); + rig.llm.chat = async (spec, messages, tools, opts) => { + if (rig.llm.calls.length > 0) return chat(spec, messages, tools, opts); + rig.llm.calls.push({ spec, messages, tools }); + opts?.onDelta?.({ type: 'content', text: '我先说一半' }); + return new Promise((_resolve, reject) => { + opts!.signal!.addEventListener('abort', () => reject(opts!.signal!.reason), { once: true }); + }); + }; + try { + rig.start(); + rig.pushEvent('讲个故事'); + await until(() => rig.llm.calls.length === 1); + expect(rig.loop.abortCurrentRound()).toBe(false); + const stop = rig.store.append({ type: 'qq.message', ts: '2026-07-17T10:01:00+08:00', source: 'qq', origin: 'external', text: '换个话题' }); + rig.bus.push({ event: stop }, { trigger: 'interrupt' }); + await until(() => rig.llm.calls.length >= 2); + const sent = rig.llm.calls[1].messages; + const partial = sent.findIndex((m) => m.role === 'assistant' && m.content === '我先说一半'); + expect(partial).toBeGreaterThan(0); + expect(sent.findIndex((m) => JSON.stringify(m).includes('换个话题'))).toBeGreaterThan(partial); + } finally { await rig.cleanup(); } + }); + + it('撤回记录让重启补投跳过被撤回的事件,水位越过它', async () => { + const rig = makeRig(); + const withdrawn = rig.store.append({ type: 'qq.message', ts: '2026-07-17T10:00:00+08:00', source: 'qq', origin: 'external', contextDelivery: 'deliver', text: '撤回的话' }); + rig.store.append({ type: 'qq.message', ts: '2026-07-17T10:00:01+08:00', source: 'qq', origin: 'external', contextDelivery: 'deliver', text: '留下的话' }); + rig.loop.recordWithdrawn(withdrawn); + try { + rig.start(); + rig.pushEvent('新消息'); + await until(() => rig.llm.calls.length >= 1); + const sent = JSON.stringify(rig.llm.calls[0].messages); + expect(sent).toContain('留下的话'); + expect(sent).not.toContain('撤回的话'); + await until(() => rig.state.data.lastDeliveredCursor === rig.store.latestCursor()); + } finally { await rig.cleanup(); } + }); }); describe('MainLoop 异常退出', () => { diff --git a/tests/cormini/output-tap.test.ts b/tests/cormini/output-tap.test.ts index d8fdc464..698c31f9 100644 --- a/tests/cormini/output-tap.test.ts +++ b/tests/cormini/output-tap.test.ts @@ -43,6 +43,15 @@ describe('outputTap 扇入', () => { expect(seen).toEqual(['甲']); }); + it('扇入多个接收器时,没声明 externalizes 的接收器按单接收器时的默认规则判定', () => { + const persona = new Cormini({ + memoryDir: dir, + worlds: [world('a', () => ({ onEvent: () => {} })), world('b', () => ({ onEvent: () => {} }))], + }); + + expect(persona.declareSessions()[0].outputTap?.externalizes?.(delta('甲'))).toBe(true); + }); + it('没有任何 World 给接收器时不声明 tap', () => { const persona = new Cormini({ memoryDir: dir,