测试骨架的 subscribe bootstrap 按生产口径过滤

- harness.ts 记录未完成条目集合、事件序号与最新生命周期锚点
- subscribe 只回放未完成条目与最新生命周期锚点,已完成条目与增量正文不再补发
- 去掉已下发事件的序号记录,避免序号表无限增长
This commit is contained in:
2026-09-16 23:48:46 +08:00
parent 1b3a3b003f
commit 6a4cb2bc18
@@ -627,6 +627,23 @@ function createPlanGddStateView(
};
}
/**
* 运行态事件里的条目身份:只有 `item.started` / `item.completed` 带条目。
*
* `item.delta` 只带 itemId 不带条目,属于瞬时事件,不参与 bootstrap 回放。
*/
function directThreadEventItemId(
event: Record<string, unknown>,
): string | null {
if (event.type !== 'item.started' && event.type !== 'item.completed') {
return null;
}
const item = event.item;
if (!item || typeof item !== 'object') return null;
const itemId = (item as Record<string, unknown>).itemId;
return typeof itemId === 'string' && itemId ? itemId : null;
}
function createProjectSupervisorRuntimeHarness({
projectPath = '/tmp/authorized-game',
sessionId = 'supervisor-session-active',
@@ -757,6 +774,18 @@ function createProjectSupervisorRuntimeHarness({
// 运行态事件里只有一个"最新回合是否在跑"的布尔:与 Thread Manager 只保留一条生命周期
// 锚点一致,重复 `turn.started` 在生产里不会出现。
let directThreadTurnRunning = false;
// bootstrap 的回放规则与 Rust 侧 `is_bootstrap_event` 对齐:只补"这一刻还没结束的条目"
// 与最新一条生命周期锚点;已完成条目、增量正文这类瞬时事件都不回放。
const directThreadActiveItemIds = new Set<string>();
const directThreadPendingEventSeqs = new Map<
Record<string, unknown>,
number
>();
let directThreadEventSequence = 0;
let directThreadLifecycleAnchor: {
seq: number;
event: Record<string, unknown>;
} | null = null;
const conversationRecord = (
role: 'user' | 'assistant',
@@ -954,12 +983,35 @@ function createProjectSupervisorRuntimeHarness({
if (command === 'subscribe_direct_project_thread') {
directThreadSubscriptionSequence += 1;
directThreadSubscriptionId = `direct-thread-${directThreadSubscriptionSequence}`;
const events = pendingDirectThreadEvents;
const pending = pendingDirectThreadEvents;
pendingDirectThreadEvents = [];
// 游标落在队尾:只有"还没结束的条目"与最新生命周期锚点作为 bootstrap 回放,
// 已完成条目与增量正文不再补发(Rust 侧 `is_bootstrap_event` 的同一口径)。
const replay = pending
.filter((event) => {
const itemId = directThreadEventItemId(event);
return Boolean(itemId && directThreadActiveItemIds.has(itemId));
})
.map((event) => ({
seq: directThreadPendingEventSeqs.get(event) ?? 0,
event,
}));
for (const event of pending) {
directThreadPendingEventSeqs.delete(event);
}
if (
directThreadLifecycleAnchor &&
!replay.some(
(entry) => entry.event === directThreadLifecycleAnchor?.event,
)
) {
replay.push(directThreadLifecycleAnchor);
}
replay.sort((left, right) => left.seq - right.seq);
return {
subscriptionId: directThreadSubscriptionId,
lastCompletedItemId: directThreadLastCompletedItemId,
events,
events: replay.map((entry) => entry.event),
};
}
if (command === 'consume_direct_project_thread') {
@@ -1188,8 +1240,24 @@ function createProjectSupervisorRuntimeHarness({
...events: Array<Record<string, unknown>>
) => {
for (const event of events) {
directThreadEventSequence += 1;
directThreadPendingEventSeqs.set(event, directThreadEventSequence);
if (event.type === 'turn.started') directThreadTurnRunning = true;
if (event.type === 'turn.completed') directThreadTurnRunning = false;
if (event.type === 'turn.started' || event.type === 'turn.completed') {
// 生命周期锚点只留最新一条,与 Thread Manager 的 `lifecycle_anchor` 一致。
directThreadLifecycleAnchor = {
seq: directThreadEventSequence,
event,
};
}
const itemId = directThreadEventItemId(event);
if (itemId && event.type === 'item.started') {
directThreadActiveItemIds.add(itemId);
}
if (itemId && event.type === 'item.completed') {
directThreadActiveItemIds.delete(itemId);
}
}
pendingDirectThreadEvents.push(...events);
if (directThreadSubscriptionId) {