Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 2 additions & 1 deletion bots/cormini/persona/persona.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ import type {
ToolDef,
ToolTag,
} from 'cortico/core/types.ts';
import { defaultExternalizes } from 'cortico/core/types.ts';

/** 默认的存在方式自述源文件;bot 不覆盖 orientation 时读它,控制台可编辑同一份。 */
/** 前缀装配模板与 World 段模板住在软件包里(跟着代码走),不在人格工作区。 */
Expand Down Expand Up @@ -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); },
};
Expand Down
6 changes: 6 additions & 0 deletions docs/sessions.md
Original file line number Diff line number Diff line change
Expand Up @@ -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))`。生效窗口取服务探测值与配置的
Expand Down
7 changes: 5 additions & 2 deletions docs/world-compatibility.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)`:最近一段时间内模型调用失败或流中断的次数。

Expand Down
20 changes: 17 additions & 3 deletions docs/worlds.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 提供以下能力:
Expand All @@ -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`、
Expand All @@ -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 保留帧名冲突时,装配层拒绝挂载并报告原因。
Expand Down
10 changes: 8 additions & 2 deletions src/core/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 项。
Expand All @@ -53,7 +54,8 @@ piggyback 只入队,随后续唤醒一起投递。
`nextBatch()` 仅支持一个消费者,每次按 FIFO 顺序取走整批。`batching` 使用共享配置引用,
更新后的值在下一次入队时参与计算。

投递水位 `lastDeliveredCursor` 持久化,队列不持久化。重启时补投水位之后的外部事件;
投递水位 `lastDeliveredCursor` 持久化,队列不持久化。重启时补投水位之后的外部事件,跳过有
`core.withdrawal` 撤回记录的事件;
已被候选处理结果引用的原始归档不重复投递。内部事件仅在当次运行投递。清空分片或跳过损坏行
留下的游标空位没有可投递内容,不阻止水位推进。

Expand All @@ -65,6 +67,10 @@ piggyback 只入队,随后续唤醒一起投递。
退回总线,进入下一批。控制台的前缀重载与清空 session 在批次边界执行:正在处理批次时等该批
结束,空闲时立即;并发请求复用同一事务。

`preempt` 取消尚未外化的模型轮,`interrupt` 取消模型轮(已外化的输出保留)或停止 `interruptible`
工具、跳过本轮尚未开始的调用。两者取消后在同一批内接收新事件并开始下一轮,被取消的轮计入轮数,
不调用 `onTurnEnded`;到达时没有可取消的轮,事件在下一次模型请求前送入。

一批正文归档后,Core 调用 `Persona.onDelivery` 并等待它返回的 Promise,不设期限。完成前
`injectInternal` 的内容排在这批的内部行末尾、外部正文之前;完成后的注入进入总线。

Expand Down
32 changes: 24 additions & 8 deletions src/core/bus.ts
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
* 队列包含即时事件、延迟渲染项与候选项。延迟渲染项不参与关键词匹配或外部事件计数;
* 候选项按外部事件计数,关键词只检查 gateText。
* piggyback 不触发计时、关键词或溢出,也不单独获得投递许可;需要其他项触发投递。
* preempt 与 interrupt 投递后通知主循环;取消哪一部分由主循环裁决。
* 人工暂停阻止所有投递。DeliveryGate 解除、关键词回调授权或溢出授权均放行整批。
*/
import type { DeliveryGate, EventOrigin, Logger, TriggerMode, WakeItem } from './types.ts';
Expand Down Expand Up @@ -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;

Expand All @@ -55,7 +56,7 @@ export class WakeBus {
this.log = log;
}

setPreemptHandler(handler: () => void): void {
setPreemptHandler(handler: (trigger: 'preempt' | 'interrupt') => void): void {
this.preemptHandler = handler;
}

Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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;
Comment thread
cursor[bot] marked this conversation as resolved.
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 } });
Expand Down
20 changes: 15 additions & 5 deletions src/core/core.ts
Original file line number Diff line number Diff line change
Expand Up @@ -192,8 +192,9 @@ export class Core<C extends CoreConfig = CoreConfig> {
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();
});
}

Expand Down Expand Up @@ -255,9 +256,9 @@ export class Core<C extends CoreConfig = CoreConfig> {
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;
}

Expand Down Expand Up @@ -557,6 +558,15 @@ export class Core<C extends CoreConfig = CoreConfig> {
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) => {
Expand Down
Loading
Loading