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
9 changes: 9 additions & 0 deletions docs/issues/2026-09-10-chat-relay-prompt-timeout.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
# P0:Chat relay 误判长任务超时

## 时间

2026-09-10

## 事故简述

Chat relay 在 `send_prompt` 后连续 5 分钟未收到业务帧或 JSON-RPC 结果时,主动将请求标记为超时。长时间且中途无输出的 Bash 仍在正常执行,却被错误收敛为 `turn_failed`,导致用户看到虚假失败。该超时越权判断 Agent 生命周期,现已完整删除,终态仅以 Agent 返回、断连或显式取消为准。
Original file line number Diff line number Diff line change
Expand Up @@ -88,7 +88,6 @@ function createHarness(overrides: Partial<GatewayDependencies> = {}): Harness {
createMessageHandler: () => () => undefined,
bindInstanceSession: () => undefined,
openReplayWindow: () => undefined,
convergeStuckPrompt: () => undefined,
} as unknown as RelayEventHandler;
const dependencies: GatewayDependencies = {
registry,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -187,21 +187,18 @@ describe("relay 与会话频道的内存边界", () => {
expect(sent[0]?.method).toBe("session/list");
});

// 无效 sessionUpdate 与无 id 的 result 不得生成投影;保活帧不改变业务入站时间,但未知事件仍会留下活跃信号
test("ignores unknown events and idless results without treating keepalive as activity", async () => {
// 无效 sessionUpdate、无 id 的 result 与保活帧均不得生成投影
test("ignores unknown events, idless results, and keepalive frames", async () => {
const { manager, handler } = await setupRelay();
const shared = relay("rcs-1");
shared.lastInboundAt = 42;
const consume = handler.createMessageHandler(shared);

await consume(message({ type: "ping" }));
expect(shared.lastInboundAt).toBe(42);
await consume(message({ jsonrpc: "2.0", result: { stopReason: "end_turn" } }));
await consume(
message({ jsonrpc: "2.0", method: "session/update", params: { update: { sessionUpdate: "not_supported" } } }),
);

expect(shared.lastInboundAt).toBeGreaterThan(42);
expect(getEntriesMap(manager.getChatYdoc("rcs-1")!).size).toBe(0);
});

Expand Down
29 changes: 0 additions & 29 deletions packages/chat-channel/src/channel/connection-types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,14 +18,6 @@ export type RelayMessage = EngineRelayMessage;
*/
export const REPLAY_WINDOW_MS = 10_000;

/**
* prompt 静默超时阈值(ms)。send_prompt 登记后若 agent 全程无任何业务帧
* (流式输出/事件/JSON-RPC 响应均刷新 lastInboundAt,保活帧除外)持续超过该阈值,
* 判定 prompt 已卡死(如被路由到错误 session 后 agent 无响应),收敛 turn_failed,
* 防止前端 loading 永久卡死。期间有任何业务帧则超时判定顺延,不误杀长输出。
*/
export const PROMPT_TIMEOUT_MS = 5 * 60_000;

/** 最小 WebSocket 连接抽象:与传输框架解耦(Elysia WS / Hono WSContext 适配为同一形状) */
export interface WsConnection {
/**
Expand Down Expand Up @@ -107,14 +99,6 @@ export interface SharedRelay {
* 聚合层能区分终态归属,stale turn 的迟到终态不再误伤新 turn。
*/
pendingPromptTurns?: Map<number | string, string>;
/**
* 在途 prompt 的超时定时器(rpcId → timer)。registerPendingPrompt 登记时启动,
* 到点时若 agent 全程静默(lastInboundAt 距今 ≥ PROMPT_TIMEOUT_MS)收敛 turn_failed,
* 有业务帧则重排等待;JSON-RPC result/error 消费登记时清除(见 clearPendingPromptTimeout)。
*/
pendingPromptTimeouts?: Map<number | string, ReturnType<typeof setTimeout>>;
/** 最近一次收到 Agent 业务入站消息的时间戳(ms),流式输出期间持续刷新;保活帧不更新 */
lastInboundAt?: number;
/** session/list 轮询因 status 门禁未置位而连续跳过的次数(连续 3 次告警,成功后清零) */
sessionListSkipCount?: number;
/**
Expand Down Expand Up @@ -144,16 +128,3 @@ export interface SharedRelay {
/** 未结束 callback 的独立 assistant 历史 entry,禁止无头 chunk 回退到主 active turn。 */
callbackAssistantEntryId?: string | null;
}

/**
* 清除指定 rpcId 的 prompt 超时定时器。
* 在 JSON-RPC result(成功)/error(拒绝)消费登记时调用,防止定时器到点后
* 对已正常完成的 prompt 误收敛;超时收敛路径自身也会消费登记并清除定时器。
*/
export function clearPendingPromptTimeout(shared: SharedRelay, rpcId: number | string): void {
const timer = shared.pendingPromptTimeouts?.get(rpcId);
if (timer) {
clearTimeout(timer);
shared.pendingPromptTimeouts?.delete(rpcId);
}
}
23 changes: 2 additions & 21 deletions packages/chat-channel/src/channel/gateway.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import type { DocManager } from "../state";
import { flushPendingYjsActions, forwardYjsAction } from "./action-forward";
import type { YjsBroadcaster } from "./broadcaster";
import type { ConnectionRegistry } from "./connection-registry";
import { type ClientConnection, PROMPT_TIMEOUT_MS, type SharedRelay, type WsConnection } from "./connection-types";
import type { ClientConnection, SharedRelay, WsConnection } from "./connection-types";
import { type PendingInitialSync, synchronizeInitialDocs } from "./gateway-sync";
import type { RelayEventHandler } from "./relay-event-handler";
import type { SessionChannel, SessionConnection } from "./session-channel";
Expand Down Expand Up @@ -386,14 +386,9 @@ export class Gateway {
}
// 在途会话同步请求登记随 relay 释放一并清空,避免残留条目无界增长
shared.pendingSessionSyncIds?.clear();
// 在途 prompt 登记与超时定时器一并清空:relay 释放后不再需要收敛,
// 残留定时器到点会 dispatch 到已销毁的 doc(且引用泄漏)
// 在途 prompt 登记随 relay 释放一并清空
shared.pendingPromptIds?.clear();
shared.pendingPromptTurns?.clear();
if (shared.pendingPromptTimeouts) {
for (const timer of shared.pendingPromptTimeouts.values()) clearTimeout(timer);
shared.pendingPromptTimeouts.clear();
}
// 回放窗口定时器一并清理(同泄漏语义),窗口判定缓存随之重置
if (shared.replayWindowTimer) {
clearTimeout(shared.replayWindowTimer);
Expand Down Expand Up @@ -469,20 +464,6 @@ export class Gateway {
if (!shared.pendingPromptTurns) shared.pendingPromptTurns = new Map();
shared.pendingPromptTurns.set(rpcId, turnId);
}
if (!shared.pendingPromptTimeouts) shared.pendingPromptTimeouts = new Map();
const schedule = () => {
const timer = setTimeout(() => {
if (!shared.pendingPromptIds?.has(rpcId)) return;
const lastInboundAt = shared.lastInboundAt ?? 0;
if (Date.now() - lastInboundAt < PROMPT_TIMEOUT_MS) {
schedule();
return;
}
this.dependencies.relayEvents.convergeStuckPrompt(shared, rpcId);
}, PROMPT_TIMEOUT_MS);
shared.pendingPromptTimeouts?.set(rpcId, timer);
};
schedule();
},
};
}
Expand Down
91 changes: 1 addition & 90 deletions packages/chat-channel/src/channel/relay-event-handler.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -358,55 +358,15 @@ describe("RelayEventHandler", () => {
expect(processed).not.toContain("turn_failed");
});

// prompt 静默超时(gateway 定时器到点且 agent 全程无业务帧)时,convergeStuckPrompt
// 必须收敛 turn_failed 并消费登记、清除定时器(B 方案:防前端 loading 永久卡死)。
test("convergeStuckPrompt converges the turn and clears pending registration", async () => {
const registry = new ConnectionRegistry();
const broadcaster = new YjsBroadcaster(registry);
const processed: string[] = [];
const reports: Array<[string, unknown]> = [];
const handler = createRelayEvents(registry, broadcaster, processed, {
reportError: (message, error) => reports.push([message, error]),
});
const shared = relayOn("rcs-1");
const timer = setTimeout(() => {}, 1000);
shared.pendingPromptIds = new Set([1]);
shared.pendingPromptTimeouts = new Map([[1, timer]]);

handler.convergeStuckPrompt(shared, 1);

expect(processed).toContain("turn_failed");
expect(shared.pendingPromptIds?.size).toBe(0);
expect(shared.pendingPromptTimeouts?.size).toBe(0);
expect(reports).toHaveLength(1);
expect(reports[0]?.[0]).toContain("prompt timed out");
clearTimeout(timer);
});

// 未登记的超时 id 不得收敛(与 error 路径的登记匹配语义一致),幂等保护。
test("convergeStuckPrompt ignores unregistered prompt ids", async () => {
const registry = new ConnectionRegistry();
const broadcaster = new YjsBroadcaster(registry);
const processed: string[] = [];
const handler = createRelayEvents(registry, broadcaster, processed);
const shared = relayOn("rcs-1");

handler.convergeStuckPrompt(shared, 99);

expect(processed).not.toContain("turn_failed");
});

// prompt 成功路径的 JSON-RPC result 响应(acp-link 的 session/prompt 返回 turnId)
// 必须消费在途 prompt 登记并清除定时器,否则登记永久残留(成功路径残留修复)
// 必须消费在途 prompt 登记,避免成功路径永久残留
test("prompt success result consumes the pending prompt registration", async () => {
const registry = new ConnectionRegistry();
const broadcaster = new YjsBroadcaster(registry);
const processed: string[] = [];
const handler = createRelayEvents(registry, broadcaster, processed);
const shared = relayOn("rcs-1");
const timer = setTimeout(() => {}, 1000);
shared.pendingPromptIds = new Set([1]);
shared.pendingPromptTimeouts = new Map([[1, timer]]);

await handler.createMessageHandler(shared)({
jsonrpc: "2.0",
Expand All @@ -415,55 +375,6 @@ describe("RelayEventHandler", () => {
} as unknown as RelayMessage);

expect(shared.pendingPromptIds?.size).toBe(0);
expect(shared.pendingPromptTimeouts?.size).toBe(0);
clearTimeout(timer);
});

// lastInboundAt 只被 agent 输出/事件类帧刷新(session/update 通知与私有帧),
// 保活帧与 JSON-RPC 响应帧(result/error)不得刷新——否则 10s 一次的
// list_sessions 轮询响应持续刷新时间戳,卡死的 prompt(agent 全程静默)
// 永远等不到超时收敛(判定被无限重排,loading 永久)。
test("lastInboundAt is refreshed by business frames but not keepalive or JSON-RPC response frames", async () => {
const registry = new ConnectionRegistry();
const broadcaster = new YjsBroadcaster(registry);
const processed: string[] = [];
const handler = createRelayEvents(registry, broadcaster, processed);
const shared = relayOn("rcs-1");
const before = Date.now();

// 保活帧:不刷新
await handler.createMessageHandler(shared)({ type: "keep_alive" } as unknown as RelayMessage);
expect(shared.lastInboundAt).toBeUndefined();

// JSON-RPC 响应帧(list_sessions 轮询 result / prompt error):不刷新
await handler.createMessageHandler(shared)({
jsonrpc: "2.0",
id: 42,
result: { sessions: [] },
} as unknown as RelayMessage);
expect(shared.lastInboundAt).toBeUndefined();
await handler.createMessageHandler(shared)({
jsonrpc: "2.0",
id: 43,
error: { code: -32000, message: "No active session" },
} as unknown as RelayMessage);
expect(shared.lastInboundAt).toBeUndefined();

// JSON-RPC 通知(session/update 流式增量):刷新
await handler.createMessageHandler(shared)({
jsonrpc: "2.0",
method: "session/update",
params: { sessionId: "ses-1", update: { sessionUpdate: "agent_message_chunk" } },
} as unknown as RelayMessage);
expect(shared.lastInboundAt).toBeGreaterThanOrEqual(before);

// 私有帧(非保活):刷新
const after = shared.lastInboundAt ?? 0;
await handler.createMessageHandler(shared)({
type: "agent_message_chunk",
payload: { type: "text", text: "hi" },
} as unknown as RelayMessage);
expect(shared.lastInboundAt).toBeGreaterThanOrEqual(after);
});

// relay 错误只发送给当前 RCS 会话,其他用户或会话不得收到 Agent 错误。
Expand Down
50 changes: 1 addition & 49 deletions packages/chat-channel/src/channel/relay-event-handler.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ import {
import type { DocManager } from "../state";
import type { YjsBroadcaster } from "./broadcaster";
import type { ConnectionRegistry } from "./connection-registry";
import { clearPendingPromptTimeout, REPLAY_WINDOW_MS, type RelayMessage, type SharedRelay } from "./connection-types";
import { REPLAY_WINDOW_MS, type RelayMessage, type SharedRelay } from "./connection-types";

/** 运行时错误只暴露稳定分类和安全文案,并以同一 ID 写入安全诊断日志。 */
function agentRuntimeError(
Expand Down Expand Up @@ -104,11 +104,6 @@ function createReplayTurnId(): string {
return `turn_replay_${Date.now()}_${Math.random().toString(36).slice(2, 10)}`;
}

/** 保活类消息类型(与 acp-idle-monitor 的 isIgnoredActivityMessageType 规则一致),不计入业务帧 */
function isKeepaliveMsgType(type: string | undefined): boolean {
return type === "keep_alive" || type === "heartbeat" || type === "ping" || type === "pong";
}

/** 读取聚合层活动 turn(Session Doc root.session.activeTurnId/Status 为权威,与 chat-writer 一致) */
function readActiveTurn(
docManager: DocManager,
Expand Down Expand Up @@ -192,18 +187,6 @@ export class RelayEventHandler {
// touchInstanceActivity 内部已过滤 keep_alive/ heartbeat/ping/pong 等保活消息
this.dependencies.touchInstanceActivity(shared.instanceId, raw);

// 刷新"最后业务入站帧"时间戳:prompt 超时收敛(gateway 定时器)依赖它判断
// agent 是否仍在活跃输出。**仅 agent 输出/事件类帧刷新**:session/update 通知
// 与私有帧(流式输出/工具/权限)代表 agent 活跃;JSON-RPC 响应帧(result/error)
// 不刷新——否则 10s 一次的 list_sessions 轮询响应会持续刷新时间戳,卡死的
// prompt(agent 全程静默)永远等不到超时收敛(判定被无限重排,loading 永久)。
// prompt 自身的 result/error 由 pendingPromptIds 消费分支收敛,不依赖此时间戳。
const isRpcResponse =
rpcCheck != null && rpcCheck.id !== undefined && ("result" in rpcCheck || "error" in rpcCheck);
if (!isKeepaliveMsgType(msgType) && !isRpcResponse) {
shared.lastInboundAt = Date.now();
}

// binding 校验:session-bound ACP 通知(session/update、peri/agent_event、
// peri/unstable_event)携带的 sessionId 必须与当前实例绑定的 ACP session 一致,
// 不一致(过期会话/串流)直接丢弃,不得写入 Y.Doc——扩展自原 session/update
Expand Down Expand Up @@ -548,7 +531,6 @@ export class RelayEventHandler {
const rpcError = rpcCheck.error as Record<string, unknown> | undefined;
if (rpcId !== undefined && rpcId !== null && shared.pendingPromptIds?.has(rpcId) === true) {
shared.pendingPromptIds?.delete(rpcId);
clearPendingPromptTimeout(shared, rpcId);
// 回传 turnId:聚合层按归属终结对应 turn(stale turn 的迟到终态不误伤新 turn)
const turnId = shared.pendingPromptTurns?.get(rpcId);
shared.pendingPromptTurns?.delete(rpcId);
Expand Down Expand Up @@ -606,7 +588,6 @@ export class RelayEventHandler {
if (rpcId !== undefined && rpcId !== null && shared.pendingPromptIds?.has(rpcId) === true) {
shared.pendingPromptIds.delete(rpcId);
shared.pendingPromptTurns?.delete(rpcId);
clearPendingPromptTimeout(shared, rpcId);
}
const syncRequested = rpcId !== undefined && rpcId !== null && shared.pendingSessionSyncIds?.has(rpcId) === true;
if (syncRequested) shared.pendingSessionSyncIds?.delete(rpcId);
Expand Down Expand Up @@ -774,35 +755,6 @@ export class RelayEventHandler {
}
}

/**
* 收敛卡死的在途 prompt(gateway 超时定时器到点且 agent 全程静默时调用)。
* 消费登记并清除定时器后收敛 turn_failed——与 error 拒绝路径相同的终态语义,
* 使前端 loading 不会永久卡死。错误内容脱敏(通用文案),只记录实例上下文。
*/
convergeStuckPrompt(shared: SharedRelay, rpcId: number | string): void {
if (!shared.pendingPromptIds?.has(rpcId)) return;
shared.pendingPromptIds.delete(rpcId);
clearPendingPromptTimeout(shared, rpcId);
// 回传 turnId:聚合层按归属终结对应 turn(stale turn 的迟到终态不误伤新 turn)
const turnId = shared.pendingPromptTurns?.get(rpcId);
shared.pendingPromptTurns?.delete(rpcId);
this.dependencies.reportError("[YJS-FE] prompt timed out (no agent response)", {
instanceId: shared.instanceId,
});
const publicError = agentRuntimeError(
"AGENT_RUNTIME.PROMPT_TIMEOUT",
"relay.prompt_timeout",
this.dependencies.log,
);
// Prompt 超时同样是 turn 终态,只进入会话时间线,不生成顶部连接错误。
this.dispatch(shared, {
type: "turn_failed",
update: { publicError },
content: null,
turnId,
});
}

private sendSafeErrorToRcsSession(shared: SharedRelay, error: PublicError): void {
this.dependencies.registry.forEachByRcsSession(shared.rcsSessionId, (entry) => {
this.dependencies.broadcaster.sendToYjsWs(entry.ws, {
Expand Down
Loading