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
2 changes: 1 addition & 1 deletion docs/console.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ SameSite=Strict、Max-Age 取浏览器上限 400 天;请求经 HTTPS 或反向

| 路由 | 页 | 内容 |
|---|---|---|
| `live` | 终端 | 与 bot 对话、时间线、上下文圈、fork;全新部署上多一组开场引导 |
| `live` | 终端 | 与 bot 对话、时间线、主循环运行阶段、上下文圈、fork;全新部署上多一组开场引导 |
| `core` | 运行诊断 | run、session、事件、运行日志,以及 Core 自己的数据与配置 |
| `usage` | 用量与成本 | 按 session、按天的 token 与费用 |
| `providers` | 模型供应商 | 端点表(见 [providers.md](providers.md)) |
Expand Down
1 change: 1 addition & 0 deletions docs/world-compatibility.md
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@
- `outputTap()`:把主 session 的模型输出流逐段交给 World;`externalizes` 返回 true 后,该轮不再被 `preempt` 取消。
- `cognition`:World 向 Persona 请求后台认知计算。
- `llmStalls(withinMs)`:最近一段时间内模型调用失败或流中断的次数。
- `onRunPhase(phase)`:主循环进入投递、模型调用、工具执行、重试等待、交接或空闲时,以及工具开始或结束时同步通知。

控制台面板与配置组不属于等级;没有控制台的宿主由 World 的配置文件提供配置。

Expand Down
1 change: 1 addition & 0 deletions docs/worlds.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ World 不直接访问 Memory 或调用 Persona 的工具。
| `outputTap?()` | 主 session 输出流的接收器(演出、字幕);这一刻没有接收器时返回 `undefined` |
| `onHandoffEnded?()`、`onTurnEnded?()` | 交接结束与主循环一轮结束的通知;隐藏的 World 不接收 |
| `onEventsSettled?(events, outcome)` | 本 World 的事件写入主 session(`delivered`)或被操作者清空队列丢弃(`discarded`) |
| `onRunPhase?(phase)` | 主循环的 `RunPhase`(`idle`、`delivering`、`model`、`tools`、`backoff`、`handoff`,轮序号,执行中的工具名)改变时同步调用;隐藏的 World 不接收 |
| `shutdownVerification?()` | 关机前要核对的外部状态,同步只读快照 |

Core 通过 `WorldHost` 向 World 提供以下能力:
Expand Down
1 change: 1 addition & 0 deletions src/bot.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1243,6 +1243,7 @@ export function createBot<C extends CoreConfig>(
onSessionReset: (cb) => core.session.onReset(cb),
onEvent: (cb) => core.store.onAppend(cb),
onRunlog: (cb) => core.runlog.onWrite(cb),
onRunPhase: (cb) => core.loop.onRunPhase(cb),
runId: () => core.run.id,
toolSchemas: () => core.loop.getToolSchemas(),
},
Expand Down
6 changes: 6 additions & 0 deletions src/core/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,12 @@ piggyback 只入队,随后续唤醒一起投递。
默认允许连续重试 2 次、每批最多 4 次,退避为 2 秒、10 秒;上下文超限、抢占、关机或轮数
达到硬上限时不重试。

`MainLoop` 维护 `RunPhase`:`delivering`(写入一批事件,含等待 `onDelivery`)、`model`、`tools`、
`backoff`(带 `retryAt`)、`handoff`,其余时刻是 `idle`;一批的轮次结束即回到 `idle`,批末钩子在
`idle` 下运行。`running` 按开始顺序列出执行中的工具,含流式提前执行的调用。state 或 round 改变、
工具开始或结束时同步通知可见 World 的 `onRunPhase` 与控制台的订阅,异常记 warn;`getStatus().phase`
返回当前值。暂停与投递闸门不进入 `RunPhase`,由 `paused`、`scheduleBlocked` 报告。

参数不是合法 JSON 时返回 `TOOL_FAILED_BAD_ARGS`,不执行工具;未知工具返回 `UNKNOWN_TOOL`;
handler 异常转为失败回执。流式生成时 `EagerDispatch` 可提前执行完整的工具调用,遵守
`barrierAfter` 顺序,并按 call id 配对结果。回执超过 8000 字符时记录 warn。
Expand Down
107 changes: 97 additions & 10 deletions src/core/loop.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ import type {
ModelSpec,
OutputTap,
Persona,
RunPhase,
SessionDecl,
SessionOpeningReason,
ToolCallContext,
Expand Down Expand Up @@ -172,6 +173,7 @@ export interface LoopStatus {
lastDeliveredCursor: number;
/** 水位之后该进上下文却还没投递的外部事件数。 */
behind: number;
phase: RunPhase;
}

/**
Expand Down Expand Up @@ -331,12 +333,69 @@ export class MainLoop {
private backoffWake: (() => void) | null = null;
/** 关机或循环换代时取消全部工具;interruptible 工具另外合并当前轮的 interrupt 信号。 */
private readonly shutdown = new AbortController();
private phase: RunPhase;
/** 执行中的工具调用,键为每次调用的登记凭据,按开始顺序。 */
private readonly runningTools = new Map<object, string>();
private readonly phaseListeners: Array<(phase: RunPhase) => void> = [];

constructor(deps: MainLoopDeps) {
this.d = deps;
this.phase = { state: 'idle', running: [], enteredAt: nowIso(deps.cfg.timezone) };
deps.session.onReset(() => { this.anchor = null; });
}

/** 控制台订阅 RunPhase 变化;与 World.onRunPhase 同时、同步调用。 */
onRunPhase(listener: (phase: RunPhase) => void): void {
this.phaseListeners.push(listener);
}

/** state、round 或 retryAt 不变时不通知,enteredAt 保持进入时的时刻。 */
private enterPhase(state: RunPhase['state'], detail: { round?: number; retryAt?: string } = {}): void {
const cur = this.phase;
if (cur.state === state && cur.round === detail.round && cur.retryAt === detail.retryAt) return;
this.phase = {
state,
...(detail.round !== undefined ? { round: detail.round } : {}),
running: [...this.runningTools.values()],
...(detail.retryAt !== undefined ? { retryAt: detail.retryAt } : {}),
enteredAt: nowIso(this.d.cfg.timezone),
};
this.emitPhase();
}

/** 登记一次开始执行的工具调用,返回它结束时调用的函数。 */
private readonly toolStarted = (name: string): (() => void) => {
const token = {};
this.runningTools.set(token, name);
this.phase = { ...this.phase, running: [...this.runningTools.values()] };
this.emitPhase();
return () => {
if (!this.runningTools.delete(token)) return;
this.phase = { ...this.phase, running: [...this.runningTools.values()] };
this.emitPhase();
};
};

private emitPhase(): void {
if (!this.activeNow()) return;
const { phase } = this;
for (const world of this.d.worlds.visible()) {
if (!world.onRunPhase) continue;
try {
world.onRunPhase(phase);
} catch (e) {
this.d.log.warn('World 运行阶段钩子异常', { id: world.id, err: e });
}
}
for (const listener of this.phaseListeners) {
try {
listener(phase);
} catch (e) {
this.d.log.warn('运行阶段订阅异常', { err: e });
}
}
}

private active(generation: number): boolean {
return !this.stopped && !this.sealed && this.generation === generation;
}
Expand Down Expand Up @@ -531,9 +590,11 @@ export class MainLoop {
* 候选按 source、origin 和处理函数分组,再按来源项在批次中的顺序生成正文。
* 正文归档后调用并等待 onDelivery;钩子完成前注入的内部项追加到内部行末尾、外部正文之前。
* 内部行合成一条 user 消息;外部正文按 eventDelivery 进入合成工具回执或同一条 user 消息。
* round 是批内轮间投递时的当前轮序号。
*/
private async deliverBatch(batch: WakeItem[], generation: number): Promise<boolean> {
private async deliverBatch(batch: WakeItem[], generation: number, round?: number): Promise<boolean> {
if (!this.active(generation)) return false;
this.enterPhase('delivering', { round });
const { session, persona, store, cfg, log } = this.d;
const projections = this.prepareCandidateProjections(batch, generation);
if (!this.active(generation)) return false;
Expand Down Expand Up @@ -943,11 +1004,14 @@ export class MainLoop {
if (changed) {
await this.rounds(generation);
if (!this.active(generation)) break;
this.enterPhase('idle');
// 先执行本批登记的交接请求,再运行批末钩子与容量检查。
await this.flushRequestedHandoff(generation);
if (!this.active(generation)) break;
await this.batchEndCheck(generation);
if (!this.active(generation)) break;
} else {
this.enterPhase('idle');
}

// 延迟渲染和 piggyback 项不阻止进入空闲钩子。
Expand Down Expand Up @@ -1016,7 +1080,7 @@ export class MainLoop {
this.takeReadyBeforeRequest = false;
const ready = bus.takeIfReady();
if (ready) {
await this.deliverBatch(ready, generation);
await this.deliverBatch(ready, generation, round);
if (!this.active(generation)) return;
}
}
Expand All @@ -1039,6 +1103,7 @@ export class MainLoop {
externalized: false, interrupted: false, abortReason: null,
};
this.currentRound = flight;
this.enterPhase('model', { round });
const ctx: ToolCallContext = {
role: decl.id,
log,
Expand All @@ -1058,6 +1123,7 @@ export class MainLoop {
(name) => this.d.toolOwner?.(name),
flight.interrupt.signal,
() => (flight.interrupted ? NOT_EXECUTED_INTERRUPTED : null),
this.toolStarted,
)
: null;
/**
Expand All @@ -1073,7 +1139,7 @@ export class MainLoop {
this.finishTurn();
return false;
}
await this.deliverBatch(arrived, generation);
await this.deliverBatch(arrived, generation, round);
return this.active(generation);
};
// 轮级观测:首个内容事件的延迟、模型往返、工具阻塞,一轮一条 debug 记录(event=round)。
Expand Down Expand Up @@ -1216,12 +1282,13 @@ export class MainLoop {
const delayMs = resubmit.backoffMs[Math.min(consecutiveFailures, resubmit.backoffMs.length) - 1] ?? 0;
log.warn('LLM 调用失败,退避后在本批内重试', { ...detail, resubmits, delayMs });
noteRound('failed', { resubmit: true, delayMs });
this.enterPhase('backoff', { round, retryAt: nowIso(this.d.cfg.timezone, new Date(Date.now() + delayMs)) });
await this.backoff(delayMs);
if (!this.active(generation)) return;
// 重试前先投递等待期间已就绪的事件。
const ready = bus.takeIfReady();
if (ready) {
await this.deliverBatch(ready, generation);
await this.deliverBatch(ready, generation, round);
if (!this.active(generation)) return;
}
continue;
Expand Down Expand Up @@ -1256,6 +1323,7 @@ export class MainLoop {
return;
}

this.enterPhase('tools', { round });
const results: ContextRecord[] = [];
let barrierHit = false;
// endsTurn 工具真正执行过(没被屏障跳过、参数合法)才算数
Expand Down Expand Up @@ -1301,7 +1369,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), flight.interrupt.signal,
() => this.active(generation), this.d.toolOwner?.(def.name), flight.interrupt.signal, this.toolStarted,
);
if (!this.active(generation)) return;
}
Expand Down Expand Up @@ -1335,7 +1403,7 @@ export class MainLoop {
// 不再执行模型请求时,事件退回总线随下一批投递。
for (const item of arrived) bus.push(item, { trigger: 'flush' });
} else {
await this.deliverBatch(arrived, generation);
await this.deliverBatch(arrived, generation, round);
if (!this.active(generation)) return;
}
}
Expand Down Expand Up @@ -1382,6 +1450,7 @@ export class MainLoop {
const calls = partial.flatMap(entry => entry.item.type === 'function_call' ? [entry.item] : []);
this.pendingToolCalls = new Set(calls.map((call) => call.call_id));
for (const entry of partial) session.append(entry);
if (calls.some((call) => eager?.has(call.call_id))) this.enterPhase('tools', { round: this.phase.round });
for (const call of calls) {
const ran = eager?.take(call.call_id);
const out = ran !== undefined ? await ran : { text: notExecuted };
Expand Down Expand Up @@ -1488,7 +1557,14 @@ export class MainLoop {
if (!this.active(generation)) return Promise.resolve();
if (this.truncatePromise) return this.truncatePromise;
let tracked: Promise<void>;
tracked = this.enqueueMaintenance(() => this.performHandoff(generation), generation).finally(() => {
tracked = this.enqueueMaintenance(async () => {
this.enterPhase('handoff');
try {
await this.performHandoff(generation);
} finally {
this.enterPhase('idle');
}
}, generation).finally(() => {
if (this.truncatePromise === tracked) this.truncatePromise = null;
});
this.truncatePromise = tracked;
Expand Down Expand Up @@ -1894,6 +1970,7 @@ export class MainLoop {
lastUsage: this.lastUsage,
lastDeliveredCursor: this.d.state.data.lastDeliveredCursor,
behind: this.watermarkBacklog().behind,
phase: this.phase,
};
}

Expand Down Expand Up @@ -1948,7 +2025,10 @@ function parseToolArgs(raw: string): Record<string, unknown> | null {
}
}

/** 普通执行与流式提前执行共用工具处理及日志记录;异常转换为失败回执。 */
/**
* 普通执行与流式提前执行共用工具处理及日志记录;异常转换为失败回执。
* started 在 handler 开始时登记调用,返回的函数在回执确定时调用。
*/
function runToolHandler(
def: ToolDef,
args: Record<string, unknown>,
Expand All @@ -1959,16 +2039,22 @@ function runToolHandler(
canRecord: () => boolean = () => true,
mod?: string,
interrupt?: AbortSignal,
started?: (name: string) => () => void,
): Promise<ToolOutcome> {
const startedAt = Date.now();
const signal = def.interruptible && interrupt && ctx.signal ? AbortSignal.any([ctx.signal, interrupt]) : ctx.signal;
let finished: (() => void) | undefined;
return withAnchors({ call: callId }, () => Promise.resolve()
.then(() => def.handler(args, { ...ctx, callId, signal }))
.then(() => {
finished = started?.(def.name);
return 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,
}))
.finally(() => finished?.())
.then((out): ToolOutcome => (def.interruptible && interrupt?.aborted
? { ...out, text: `${out.text}
${INTERRUPTED_WHILE_RUNNING}` }
Expand Down Expand Up @@ -2001,6 +2087,7 @@ class EagerDispatch {
private readonly interrupt?: AbortSignal,
/** 轮到执行时返回未执行标记则跳过该调用。 */
private readonly skip: () => string | null = () => null,
private readonly started?: (name: string) => () => void,
) {}


Expand Down Expand Up @@ -2039,7 +2126,7 @@ class EagerDispatch {
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)
? runToolHandler(def, args, this.ctx, call.id, this.role, this.toolLog, this.active, this.owner(def.name), this.interrupt, this.started)
: { text: skipped };
});
this.chain = run.then(() => undefined);
Expand Down
23 changes: 23 additions & 0 deletions src/core/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -856,6 +856,24 @@ export interface ShutdownExternalCheck {
manualAction: string;
}

/** 主 session 循环此刻在做什么;每个字段都是 Core 直接观察到的事实。 */
export interface RunPhase {
/**
* idle:等待下一批事件,批末钩子也在此状态运行;delivering:把一批事件写入主 session,含等待
* Persona.onDelivery;model:模型调用进行中;tools:本轮模型输出已结束,执行其中的工具调用;
* backoff:模型调用失败,等待重试;handoff:上下文交接进行中。
*/
state: 'idle' | 'delivering' | 'model' | 'tools' | 'backoff' | 'handoff';
/** 本批的轮序号,从 1 起;轮间的 delivering 带当前轮。idle、handoff 与一批的首次 delivering 时缺省。 */
round?: number;
/** 正在执行的工具名,按开始顺序;包含模型输出期间提前执行的调用。 */
running: readonly string[];
/** backoff 时下一次请求的时刻(ISO)。 */
retryAt?: string;
/** 进入当前 state 的时刻(ISO,部署时区)。 */
enteredAt: string;
}

/** World 的环境描述、事件和工具契约。 */
export interface World {
id: string;
Expand Down Expand Up @@ -896,6 +914,11 @@ export interface World {
* 之后经 queueExternalEvents 写入 session 时通知。
*/
onEventsSettled?(events: readonly EventEnvelope[], outcome: 'delivered' | 'discarded'): void;
/**
* 主循环的 RunPhase 变化时同步调用:state 或 round 改变、工具开始或结束。在主循环的调用栈上执行,
* 只记录状态,不做耗时操作。隐藏 World 不接收通知。
*/
onRunPhase?(phase: RunPhase): void;
/**
* stop() 完成后返回外部状态检查的同步只读快照。
* 网络检查须在 stop() 的既有期限内完成并缓存。
Expand Down
2 changes: 2 additions & 0 deletions src/web/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,8 @@ manifest 的 `CONSOLE_PROTOCOL_VERSION` 不匹配时,浏览器拒绝加载。

`WebApp` 默认监听 `127.0.0.1`,支持由依赖配置指定监听地址。从首选端口起最多尝试五个端口;
端口为 0 时仅申请一次系统分配。WebSocket 使用 `noServer` 分派 `/ws/debug`、`/ws/sessions` 和面板流。
`/ws/debug` 在主循环的 `RunPhase` 改变时推 `phase` 帧,值与 `status` 帧的 `loop.phase` 相同;
`status` 帧只在 session 追加、暂停与继续、端点变更时推送。
所有请求与 upgrade 先校验 Host 头:只接受回环名、显式绑定的地址和 `allowedHosts` 里的名字,绑到通配地址时不校验。
写请求与 upgrade 再按 Host 校验 Origin,拒绝不匹配或无效的 Origin;缺少 Origin 时放行,Origin 是
`allowedHosts` 里的名字时放行。监听非回环地址而没有设访问密码时启动记一条 warn。
Expand Down
Loading
Loading