修复: 打通设计回复流式投递与结构化确认
问题:Run 完成会抢在 SSE 流式事件之前清空 pending,导致回复整块出现;无 Quote 的确认按钮又会降级为普通聊天,因此无法创建任务。 实现:为 Main 事件队列增加有界投递屏障,严格识别生成确认短语并只调用结构化 Quote Action,接入 generation capability 提示,同时补齐竞态、误触和任务对账测试。
This commit is contained in:
@@ -184,11 +184,16 @@ type AgentWebSocketFactory = (
|
||||
|
||||
type TaskEventQueue = {
|
||||
events: AsyncIterable<DesignWorkspaceEvent>;
|
||||
push(event: DesignWorkspaceEvent): void;
|
||||
push(event: DesignWorkspaceEvent): Promise<void>;
|
||||
finish(): void;
|
||||
fail(error: unknown): void;
|
||||
};
|
||||
|
||||
type QueuedTaskEvent = {
|
||||
event: DesignWorkspaceEvent;
|
||||
acknowledge(): void;
|
||||
};
|
||||
|
||||
type AgentRunEventWaiter = (run: ServerAgentRun) => void;
|
||||
|
||||
const AGENT_WEBSOCKET_OPEN = 1;
|
||||
@@ -196,6 +201,7 @@ const AGENT_WEBSOCKET_PING_INTERVAL_MS = 20_000;
|
||||
const AGENT_RUN_INITIAL_POLL_INTERVAL_MS = 1_000;
|
||||
const AGENT_RUN_MAX_POLL_INTERVAL_MS = 5_000;
|
||||
const AGENT_RUN_TIMEOUT_MS = 10 * 60_000;
|
||||
const DESIGN_EVENT_DELIVERY_BARRIER_TIMEOUT_MS = 1_000;
|
||||
|
||||
function mapBrief(brief: ServerBrief): DesignBrief {
|
||||
return {
|
||||
@@ -475,7 +481,8 @@ function normalizeAgentRunEvent(value: unknown, sessionId: string): ServerAgentR
|
||||
}
|
||||
|
||||
function createTaskEventQueue(): TaskEventQueue {
|
||||
const queued: DesignWorkspaceEvent[] = [];
|
||||
const queued: QueuedTaskEvent[] = [];
|
||||
const pendingDeliveries = new Set<QueuedTaskEvent>();
|
||||
const waiters: Array<() => void> = [];
|
||||
let finished = false;
|
||||
let failed = false;
|
||||
@@ -488,9 +495,13 @@ function createTaskEventQueue(): TaskEventQueue {
|
||||
events: {
|
||||
async *[Symbol.asyncIterator]() {
|
||||
while (true) {
|
||||
const event = queued.shift();
|
||||
if (event) {
|
||||
yield event;
|
||||
const entry = queued.shift();
|
||||
if (entry) {
|
||||
try {
|
||||
yield entry.event;
|
||||
} finally {
|
||||
entry.acknowledge();
|
||||
}
|
||||
continue;
|
||||
}
|
||||
if (failed) throw failure;
|
||||
@@ -500,24 +511,52 @@ function createTaskEventQueue(): TaskEventQueue {
|
||||
},
|
||||
},
|
||||
push(event) {
|
||||
if (finished || failed) return;
|
||||
queued.push(event);
|
||||
if (finished || failed) return Promise.resolve();
|
||||
let resolveDelivery!: () => void;
|
||||
const delivered = new Promise<void>((resolve) => {
|
||||
resolveDelivery = resolve;
|
||||
});
|
||||
let acknowledged = false;
|
||||
const entry: QueuedTaskEvent = {
|
||||
event,
|
||||
acknowledge() {
|
||||
if (acknowledged) return;
|
||||
acknowledged = true;
|
||||
pendingDeliveries.delete(entry);
|
||||
resolveDelivery();
|
||||
},
|
||||
};
|
||||
pendingDeliveries.add(entry);
|
||||
queued.push(entry);
|
||||
wake();
|
||||
return delivered;
|
||||
},
|
||||
finish() {
|
||||
if (finished || failed) return;
|
||||
finished = true;
|
||||
for (const entry of [...pendingDeliveries]) entry.acknowledge();
|
||||
wake();
|
||||
},
|
||||
fail(error) {
|
||||
if (finished || failed) return;
|
||||
failed = true;
|
||||
failure = error;
|
||||
for (const entry of [...pendingDeliveries]) entry.acknowledge();
|
||||
wake();
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function boundedTaskEventDelivery(delivered: Promise<void>): Promise<void> {
|
||||
return new Promise((resolve) => {
|
||||
const timeout = setTimeout(resolve, DESIGN_EVENT_DELIVERY_BARRIER_TIMEOUT_MS);
|
||||
void delivered.then(() => {
|
||||
clearTimeout(timeout);
|
||||
resolve();
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
function webSocketCloseError(code: number): DesignWorkspaceModuleError | null {
|
||||
if (code === 1000 || code === 1001) return null;
|
||||
if (code === 4401) {
|
||||
@@ -924,6 +963,7 @@ export class WorksSquareDesignWorkspace implements DesignWorkspaceModule {
|
||||
|
||||
const { socket } = connection;
|
||||
const queue = createTaskEventQueue();
|
||||
let latestWorkspaceEventDelivery = Promise.resolve();
|
||||
let didOpen = false;
|
||||
let ending = false;
|
||||
let settled = false;
|
||||
@@ -1010,14 +1050,19 @@ export class WorksSquareDesignWorkspace implements DesignWorkspaceModule {
|
||||
socket.onmessage = ({ data }) => {
|
||||
if (ending) return;
|
||||
const agentEvent = agentEventFromWebSocketFrame(data);
|
||||
const run = normalizeAgentRunEvent(agentEvent, session.session_id);
|
||||
if (run) this.publishAgentRun(session.session_id, run);
|
||||
const event = normalizeWorkspaceEvent(
|
||||
agentEvent,
|
||||
session.session_id,
|
||||
input.workspaceId,
|
||||
);
|
||||
if (event) queue.push(event);
|
||||
if (event) {
|
||||
latestWorkspaceEventDelivery = boundedTaskEventDelivery(queue.push(event));
|
||||
}
|
||||
const run = normalizeAgentRunEvent(agentEvent, session.session_id);
|
||||
if (run) {
|
||||
const deliveryBarrier = latestWorkspaceEventDelivery;
|
||||
void deliveryBarrier.then(() => this.publishAgentRun(session.session_id, run));
|
||||
}
|
||||
};
|
||||
socket.onerror = () => {
|
||||
setTimeout(() => {
|
||||
|
||||
Reference in New Issue
Block a user