chore(workbench): remove legacy trace gap-fill remnants

This commit is contained in:
root
2026-07-09 03:52:58 +02:00
parent 6999c94712
commit 8d6409dc9d
17 changed files with 57 additions and 260 deletions
+2 -2
View File
@@ -31,9 +31,9 @@ Cloud Web 的通用加载态使用 `web/hwlab-cloud-web/src/components/common/Lo
Session rail 是该规则的高频区域。`/v1/agent/conversations` 还未返回时,即使 workspace 中已有 `selectedConversationId`、sessionId、traceId 或 selected conversation snapshot,也不能把选中 session stub 渲染成单条 `.session-tab`,更不能让它占满整个 session 列表高度。加载窗口应只显示 `LoadingState`,并隐藏当前 trace 元信息、复制/删除等依赖真实 active tab 的动作;待 conversations ready 后再渲染真实 session tabs,或在真实空集合时显示空态。
Session rail 的后台恢复刷新必须有硬边界。显式用户动作或强一致操作(例如选择会话、删除当前会话)可以立即刷新会话列表;SSE error、active trace REST gap-fill、terminal refresh、trace hydration 等后台补偿路径不得绕过 session list 的冷却/合并机制去强制刷新完整列表。后台路径应优先补当前 trace、turn status、message projection 和必要的 trace events;需要刷新 session rail 时走统一的 scheduled refresh,并按 session/list key 合并已有 timer,避免网络抖动或 EventSource error storm 把 `/v1/workbench/sessions` 放大成浏览器内存和 CDP responsiveness 红灯。
Session rail 的后台恢复刷新必须有硬边界。显式用户动作或强一致操作(例如选择会话、删除当前会话)可以立即刷新会话列表;SSE error、active trace sync replay、terminal refresh、trace hydration 等后台补偿路径不得绕过 session list 的冷却/合并机制去强制刷新完整列表。后台路径应优先补当前 trace、turn status、message projection 和必要的 trace events;需要刷新 session rail 时走统一的 scheduled refresh,并按 session/list key 合并已有 timer,避免网络抖动或 EventSource error storm 把 `/v1/workbench/sessions` 放大成浏览器内存和 CDP responsiveness 红灯。
Workbench realtime 恢复与 REST gap-fill 必须统一接入已经迁移的 OpenCode-style runtime 模块。`workbench-stream-transport` 只拥有 SSE lifecycle、cursor 和 recovery reason`workbench-realtime-plan` 只把 transport action 转成纯 plan`workbench-refresh-runtime`、keyed singleflight、scheduled task runtime、trace hydration queue、server-state reducer 和 session cache 是恢复读取的唯一 substrate。`workbench.ts` 不得再持有新的 in-flight map、timer map、cursor map、REST gap-fill queue 或第二套 recovery coordinator。`force` 是用户显式操作和恢复优先级语义,不得绕过同 key 的 singleflight、cooldown、min interval 或 storm budget;这些预算、退避、并发、页数、重试和窗口参数只由 node/lane YAML-backed runtime policy 注入,SPEC 只声明字段族、责任边界和验收读取方式,不写死数值。
Workbench realtime 恢复与 sync replay 必须统一接入已经迁移的 OpenCode-style runtime 模块。`workbench-stream-transport` 只拥有 SSE lifecycle、cursor 和 recovery reason`workbench-realtime-plan` 只把 transport action 转成纯 plan`workbench-refresh-runtime`、keyed singleflight、scheduled task runtime、trace hydration queue、server-state reducer 和 session cache 是恢复读取的唯一 substrate。`workbench.ts` 不得再持有新的 in-flight map、timer map、cursor map、REST recovery queue 或第二套 recovery coordinator。`force` 是用户显式操作和恢复优先级语义,不得绕过同 key 的 singleflight、cooldown、min interval 或 storm budget;这些预算、退避、并发、页数、重试和窗口参数只由 node/lane YAML-backed runtime policy 注入,SPEC 只声明字段族、责任边界和验收读取方式,不写死数值。
Workbench Realtime Authority v2 的自动恢复只允许消费 SSE typed event 和 `/v1/workbench/sync` replay。`/workbench/sync` 返回的 durable `delta.messages``delta.turns` 等 family object 必须在前端 authority 层投影成带 `realtimeAuthority`、entity family/id/version 和 projection revision 的 `message.snapshot``turn.snapshot` 等 typed event,再进入统一 reducer;不得把 family delta 当成不可应用的普通对象,也不得用 `/v1/workbench/sessions/:id/messages``/v1/workbench/turns/:id``/v1/workbench/traces/:id/events` 自动 fan-out 补洞。跨 tab/page 的 `session-projection` signal 只能触发同一 `/workbench/sync` replay 或临时 optimistic echo,最终必须由 durable sync delta 覆盖并让 control/observer 页在同一个 session 上收敛到相同 messages/turn projection。`workbench-server-state``message.snapshot` guard 只能拒绝无 messageId、无 trace 或当前会话没有同 trace 上下文的孤儿 agent snapshotdurable user fact 之后到达的同 trace agent snapshot 必须可追加。新增修复应覆盖“observer 从空消息状态应用 sync replay 后得到 user/agent messages 与 turn status”的最小测试,并用 web-probe `observe analyze` 确认没有 persistent `cross-page-projection-divergence` 或 automatic recovery legacy fan-out 红项。
+3 -2
View File
@@ -2546,8 +2546,9 @@ function codeAgentCompatProjectionPayload(payload = {}, context = {}) {
staleMs: projection.staleMs ?? null,
blocker: projection.blocker ?? null,
workbench: traceId ? {
turnUrl: `/v1/workbench/turns/${encodeURIComponent(traceId)}`,
traceEventsUrl: `/v1/workbench/traces/${encodeURIComponent(traceId)}/events`
detailOnly: true,
turnDetailUrl: `/v1/workbench/turns/${encodeURIComponent(traceId)}`,
traceEventsDetailUrl: `/v1/workbench/traces/${encodeURIComponent(traceId)}/events`
} : null,
valuesRedacted: true,
secretMaterialStored: false
+1
View File
@@ -262,6 +262,7 @@ test("web performance summary exposes Workbench p75 and low-sample diagnostics w
assert.ok(summary.dashboard.trends.some((trend) => trend.id === "workbench-experience" && trend.points.some((point) => point.label === "提交到首个可见结果")));
assert.ok(summary.dashboard.distributions.some((distribution) => distribution.id === "samples" && distribution.buckets.some((bucket) => bucket.label === "低样本")));
assert.ok(summary.workbenchJourneys.some((row) => row.metric === "submit_to_first_visible" && row.backend === "agentrun-v01/codex" && row.p75 >= row.p50 && row.lowSample === true && row.sampleState === "low-sample"));
assert.ok(summary.workbenchJourneys.some((row) => row.metric === "session_switch_first_visible" && row.transport === "detail_history"));
assert.ok(summary.workbenchEventPhases.some((row) => row.phase === "created_to_append" && row.eventType === "backend" && row.transport === "sse"));
assert.ok(summary.workbenchBackendEvents.some((row) => row.metric === "backend_event_to_visible" && row.eventType === "terminal" && row.backend === "agentrun-v01/codex"));
assert.equal(summary.summary.problemCount, 0);
+11 -5
View File
@@ -41,7 +41,7 @@ const WORKBENCH_EVENT_PHASES = new Set([
"api_accepted_to_backend_event"
]);
const WORKBENCH_EVENT_TYPES = new Set(["assistant", "tool_call", "backend", "terminal", "error", "status", "request", "result", "unknown"]);
const WORKBENCH_TRANSPORTS = new Set(["sse", "rest_gap", "poll", "none", "unknown"]);
const WORKBENCH_TRANSPORTS = new Set(["sse", "sync_replay", "detail_history", "poll", "none", "unknown"]);
const WORKBENCH_OUTCOMES = new Set(["ok", "timeout", "error", "dropped", "stale", "partial", "empty", "network", "denied", "unknown"]);
const WORKBENCH_ENTRIES = new Set(["new", "existing", "steer", "retry", "unknown"]);
const WORKBENCH_VISIBILITY = new Set(["foreground", "background", "hidden", "unknown"]);
@@ -432,7 +432,7 @@ export function createWebPerformanceStore(options: WebPerformanceStoreOptions =
cache: optionalEnum(input.cache, WORKBENCH_CACHE_STATES, "unknown"),
auth_state: optionalEnum(input.authState ?? field(input, "auth_state"), WORKBENCH_AUTH_STATES, "unknown"),
backend: normalizeBackendLabel(input.backend),
transport: optionalEnum(input.transport, WORKBENCH_TRANSPORTS, "unknown"),
transport: normalizeWorkbenchTransport(input.transport),
visibility: optionalEnum(input.visibility, WORKBENCH_VISIBILITY, "unknown"),
outcome: optionalEnum(input.outcome, WORKBENCH_OUTCOMES, "ok")
};
@@ -450,7 +450,7 @@ export function createWebPerformanceStore(options: WebPerformanceStoreOptions =
phase,
event_type: normalizeEventType(input.eventType ?? field(input, "event_type")),
backend: normalizeBackendLabel(input.backend),
transport: optionalEnum(input.transport, WORKBENCH_TRANSPORTS, "unknown"),
transport: normalizeWorkbenchTransport(input.transport),
outcome: optionalEnum(input.outcome, WORKBENCH_OUTCOMES, "ok")
};
return { series: "workbench_event_phase", metric: phase, value, labels };
@@ -464,7 +464,7 @@ export function createWebPerformanceStore(options: WebPerformanceStoreOptions =
...baseLabels,
event_type: normalizeEventType(input.eventType ?? field(input, "event_type")),
backend: normalizeBackendLabel(input.backend),
transport: optionalEnum(input.transport, WORKBENCH_TRANSPORTS, "unknown"),
transport: normalizeWorkbenchTransport(input.transport),
outcome: optionalEnum(input.outcome, WORKBENCH_OUTCOMES, "ok")
};
return { series: "workbench_backend_event_visible", metric: "backend_event_to_visible", value, labels };
@@ -2128,7 +2128,7 @@ function eventTypeLabel(value: string) {
}
function transportLabel(value: string) {
const labels: Record<string, string> = { sse: "实时流", rest_gap: "REST 补洞", poll: "轮询", none: "无", unknown: "未知" };
const labels: Record<string, string> = { sse: "实时流", sync_replay: "同步回放", detail_history: "明细/历史", poll: "轮询", none: "无", unknown: "未知" };
return labels[value] ?? value;
}
@@ -2193,6 +2193,12 @@ function optionalEnum(value: unknown, allowed: Set<string>, fallback: string) {
return allowed.has(text) ? text : fallback;
}
function normalizeWorkbenchTransport(value: unknown) {
const text = sanitizeMetricName(value, "unknown");
const aliased = text === "rest_gap" ? "detail_history" : text;
return WORKBENCH_TRANSPORTS.has(aliased) ? aliased : "unknown";
}
function normalizeEventType(value: unknown) {
const text = sanitizeMetricName(value, "unknown");
const aliased = text === "assistant_message" ? "assistant" : text === "tool" ? "tool_call" : text;
+4 -4
View File
@@ -48,7 +48,7 @@ const requiredFiles = Object.freeze([
"src/stores/workbench-timeline-model.ts",
"src/stores/workbench-session-cache.ts",
"src/composables/useWorkbenchScrollRuntime.ts",
"src/composables/useTraceSubscription.ts",
"src/composables/workbench-trace-snapshot.ts",
"src/composables/useAutoRefresh.ts",
"src/composables/useClipboard.ts",
"src/composables/useForm.ts",
@@ -193,8 +193,8 @@ assertIncludes(workbenchColadaSource, "staleTime", "Workbench query min-interval
assertIncludes(workbenchPerformanceSource, "recordWorkbenchRuntimeDiagnostic", "Workbench performance probe must record runtime diagnostics for monitor root cause visibility");
assertIncludes(workbenchPerformanceSource, "clearResourceTimings", "Workbench performance probe must bound browser ResourceTiming retention after API enrichment");
assertIncludes(workbenchStoreSource, "recordWorkbenchRuntimeDiagnostic", "Workbench store must surface SSE recovery diagnostics to the performance probe");
assertIncludes(workbenchRealtimePlanSource, "new Set(recovery.actions)", "Realtime recovery planner must consume transport-owned actions explicitly");
assertIncludes(workbenchRealtimePlanSource, "actions.has(\"schedule-session-list\")", "Realtime stream errors must schedule bounded session list refreshes only when transport requests that action");
assertIncludes(workbenchRealtimePlanSource, "recovery.actions.includes(\"sync-replay\")", "Realtime recovery planner must consume transport-owned actions explicitly");
assert.doesNotMatch(workbenchRealtimePlanSource, /schedule-session-list/u, "Realtime recovery planner must not restore legacy session-list repair actions");
assertIncludes(workbenchRealtimePlanSource, "authority: \"automatic-recovery\"", "Realtime recovery planner must classify transport recovery as automatic recovery authority");
assert.doesNotMatch(workbenchRealtimePlanSource, /force:\s*true/u, "Realtime recovery planner must not turn transport recovery into force-refresh work");
assertIncludes(workbenchColadaSource, "const state = await queryCache.refresh(entry);", "Workbench reads must preserve Colada staleTime/min-interval governance");
@@ -204,7 +204,7 @@ assertIncludes(workbenchStoreSource, "workbenchColadaQueries.fetchSession", "Rea
assertIncludes(workbenchStoreSource, "runtimePolicy.workbenchSessionDetailMinRefreshMs", "Realtime session detail recovery budget must come from runtime policy");
assert.doesNotMatch(workbenchStoreSource, /refreshRealtimeSessionMessages[\s\S]{0,900}refreshSessionMessageProjectionPage\(id, \{ force: true \}\)/u, "Realtime session message recovery must not force-bypass the message projection refresh budget");
assert.doesNotMatch(workbenchStoreSource, /handleRealtimeStreamError[\s\S]{0,1200}refreshSessions\([^;]+force:\s*true/u, "Realtime stream errors must not force-refresh the full session list");
assert.doesNotMatch(workbenchStoreSource, /refreshActiveTraceFromRest[\s\S]{0,1200}refreshSessions\([^;]+force:\s*true/u, "Active trace REST gap-fill must not force-refresh the full session list");
assert.doesNotMatch(workbenchStoreSource, /refreshActiveTraceFromRest[\s\S]{0,1200}refreshSessions\([^;]+force:\s*true/u, "Active trace sync replay must not force-refresh the full session list");
assert.doesNotMatch(workbenchStoreSource, /refreshTerminalTraceFromRest[\s\S]{0,1200}refreshSessions\([^;]+force:\s*true/u, "Terminal trace REST refresh must not force-refresh the full session list");
assert.doesNotMatch(workbenchStoreSource, /message\.runnerTrace(?:\?\.|\.)status/u, "Workbench message lifecycle must not be inferred from runnerTrace.status");
assert.doesNotMatch(conversationPanelSource, /message\.runnerTrace(?:\?\.|\.)status/u, "ConversationPanel must not override message completion from runnerTrace.status");
@@ -264,15 +264,15 @@ test("Workbench submit first visible waits for terminal final or tool output", (
const appendedAt = new Date(wallBase - 11_500).toISOString();
const backendEvent = { label: "agentrun:run:createdAgentRun", backend: "agentrun-v01/codex", createdAt, appendedAt, message: "AgentRun created." } as TraceEvent;
startWorkbenchSubmitJourney({ traceId: "trc_secret", sessionId: "ses_secret", entry: "existing", backend: "codex", transport: "sse" });
markWorkbenchTraceEventsReceived({ traceId: "trc_secret", transport: "rest_gap", events: [backendEvent] });
markWorkbenchTraceEventsReceived({ traceId: "trc_secret", transport: "detail_history", events: [backendEvent] });
markWorkbenchTraceProjected("trc_secret");
acknowledgeWorkbenchVisible({ messages: [traceOnlyAgentMessage("ses_secret", "trc_secret", [backendEvent])], activeSessionId: "ses_secret", detailLoading: false });
const backendOnlyEvents = drainWorkbenchPerformanceEventsForTest();
assert.equal(backendOnlyEvents.some((event) => event.kind === "workbench_journey" && event.journey === "submit_to_first_visible"), false);
assert.equal(backendOnlyEvents.some((event) => event.kind === "workbench_event_phase" && event.phase === "sse_to_receive"), false);
assert.ok(backendOnlyEvents.some((event) => event.kind === "workbench_event_phase" && event.phase === "created_to_append" && event.eventType === "backend" && event.transport === "rest_gap"));
assert.ok(backendOnlyEvents.some((event) => event.kind === "workbench_event_phase" && event.phase === "receive_to_project" && event.eventType === "backend" && event.transport === "rest_gap"));
assert.ok(backendOnlyEvents.some((event) => event.kind === "workbench_event_phase" && event.phase === "created_to_append" && event.eventType === "backend" && event.transport === "detail_history"));
assert.ok(backendOnlyEvents.some((event) => event.kind === "workbench_event_phase" && event.phase === "receive_to_project" && event.eventType === "backend" && event.transport === "detail_history"));
assert.equal(backendOnlyEvents.some((event) => event.kind === "workbench_backend_event_visible" && event.eventType === "backend" && event.outcome === "stale"), false);
acknowledgeWorkbenchVisible({ messages: [traceOnlyAgentMessage("ses_secret", "trc_secret", [{ type: "assistant_message", status: "running", message: "progress only" } as TraceEvent])], activeSessionId: "ses_secret", detailLoading: false });
@@ -3,7 +3,7 @@ import test from "node:test";
import type { AgentRunProvenance, ChatMessage, TraceEvent } from "../src/types/index.ts";
import { canCancelMessage, canRetryMessage, messageTraceId, renderSafeMarkdown, traceEventBody, traceEventLabel, traceIdentityText, visibleTraceEvents } from "../src/components/workbench/message-rendering.ts";
import { mergeRunnerTrace } from "../src/composables/useTraceSubscription.ts";
import { mergeRunnerTrace } from "../src/composables/workbench-trace-snapshot.ts";
import { traceDisplayRows, traceNoiseEventCount } from "../../../tools/src/hwlab-cli/trace-renderer.ts";
test("R1 markdown rendering keeps structure and strips unsafe HTML", () => {
@@ -38,7 +38,7 @@ test("Workbench runtime policy reads injected config while preserving defaults",
workbenchSessionDetailMinRefreshMs: 1234,
workbenchSessionMessagesWindowLimit: 9,
workbenchTraceMessagesWindowLimit: 4,
workbenchRealtimeErrorGapFillMinMs: 0,
workbenchRealtimeErrorSyncReplayMinMs: 0,
workbenchRealtimeFlushMaxItemsPerChunk: 2,
workbenchRealtimeFlushMaxChunkMs: 6,
workbenchRealtimeFlushYieldMs: 5,
@@ -50,7 +50,7 @@ test("Workbench runtime policy reads injected config while preserving defaults",
assert.equal(policy.workbenchSessionDetailMinRefreshMs, 1234);
assert.equal(policy.workbenchSessionMessagesWindowLimit, 9);
assert.equal(policy.workbenchTraceMessagesWindowLimit, 4);
assert.equal(policy.workbenchRealtimeErrorGapFillMinMs, 0);
assert.equal(policy.workbenchRealtimeErrorSyncReplayMinMs, 0);
assert.equal(policy.workbenchRealtimeFlushMaxItemsPerChunk, 2);
assert.equal(policy.workbenchRealtimeFlushMaxChunkMs, 6);
assert.equal(policy.workbenchRealtimeFlushYieldMs, 5);
@@ -1,5 +1,5 @@
// SPEC: PJ2026-0106050514 Workbench实时运行面 draft-2026-06-30-p0-1297-spec-first; PJ2026-010403 API契约 draft-2026-06-18-r1; PJ2026-010401 Web工作台 draft-2026-06-18-r1; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0; PJ2026-010401080313 Workbench实时权威 draft-2026-07-08-p0-workbench-realtime-authority-v2.
// Responsibility: Workbench SSE client. Realtime events accelerate UI projection; REST snapshots remain gap-fill authority.
// Responsibility: Workbench SSE client. Realtime events and sync replay are projection authority; REST routes are detail/history reads.
import { fetchJson, type ApiRequestOptions } from "@/api/client";
import type { ApiResult, ChatMessage, ProjectionDiagnostic, TraceEvent } from "@/types";
@@ -2,7 +2,7 @@ import assert from "node:assert/strict";
import { test } from "bun:test";
import type { ChatMessage, TraceEvent } from "@/types";
import { mergeRunnerTrace } from "./useTraceSubscription";
import { mergeRunnerTrace } from "./workbench-trace-snapshot";
function event(projectedSeq: number, label = `event-${projectedSeq}`): TraceEvent {
return { projectedSeq, label, type: "event" };
@@ -1,8 +1,7 @@
// SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-18-p0-unique-projection; PJ2026-010401 Web工作台 draft-2026-06-18-r1.
// Responsibility: Trace snapshot helpers and legacy subscription adapter backed by Workbench read-model APIs.
// Responsibility: Trace snapshot helpers for Workbench realtime authority reducers.
import { workbenchAPI, type ActivityRefSource } from "@/api";
import type { AgentChatResponse, AgentChatResultResponse, AgentRunProvenance, ChatMessage, ProjectionDiagnostic, TraceEvent, WorkbenchTurnTimingProjection } from "@/types";
import type { AgentChatResultResponse, AgentRunProvenance, ChatMessage, ProjectionDiagnostic, TraceEvent, WorkbenchTurnTimingProjection } from "@/types";
import { firstNonEmptyString } from "@/utils";
export interface TraceSnapshot {
@@ -40,31 +39,6 @@ export interface TraceSnapshot {
updatedAt?: string;
}
export interface TraceSubscriptionConfig {
traceId: string;
initial: AgentChatResponse;
onActivity: () => void;
onSnapshot: (snapshot: TraceSnapshot) => void;
onComplete: (result: AgentChatResultResponse) => void;
onInfrastructureError: (error: string) => void;
signal: AbortSignal;
inactivityTimeoutMs: number;
activityRef: ActivityRefSource;
}
const TRACE_POLL_INTERVAL_MS = 1500;
const TRACE_LIVE_PAGE_LIMIT = 80;
const TRACE_LIVE_MAX_PAGES_PER_POLL = 4;
export function isTerminalStatus(status: string | undefined): boolean {
if (!status) return false;
return ["completed", "failed", "blocked", "timeout", "cancelled", "canceled"].includes(String(status));
}
export function isResultUrlStatus(response: AgentChatResponse): boolean {
return Boolean(response.resultUrl) && (response.status === "running" || response.status === "accepted" || isTerminalStatus(response.status));
}
export function snapshotToRunnerTrace(snapshot: TraceSnapshot): NonNullable<ChatMessage["runnerTrace"]> {
const events = Array.isArray(snapshot.events) ? snapshot.events : [];
return {
@@ -348,195 +322,3 @@ export function mergeTraceResults(terminal: AgentChatResultResponse, trace: Trac
lastEventLabel: mergedTrace.lastEventLabel ?? terminal.lastEventLabel ?? undefined
};
}
export async function subscribeToTrace(config: TraceSubscriptionConfig): Promise<void> {
const { traceId, initial, onActivity, onSnapshot, onComplete, onInfrastructureError, signal, inactivityTimeoutMs, activityRef } = config;
if (isTerminalStatus(initial.status)) {
const snapshot = resultToTraceSnapshot(traceId, initial as AgentChatResultResponse, "turn-api");
onSnapshot(snapshot);
onComplete(mergeTraceResults(initial as AgentChatResultResponse, snapshot));
return;
}
if (!initial.turnUrl && !traceId) {
onInfrastructureError("Code Agent initial response is missing traceId for turn status polling");
return;
}
let lastSnapshotKey = "";
let traceAfterProjectedSeq = 0;
let accumulatedTrace: TraceSnapshot | null = null;
const fetchLiveTracePages = async (): Promise<TraceSnapshot | null> => {
for (let page = 0; page < TRACE_LIVE_MAX_PAGES_PER_POLL; page += 1) {
const previousProjectedSeq = traceAfterProjectedSeq;
const tracePolled = await workbenchAPI.traceEvents(traceId, inactivityTimeoutMs, activityRef, { afterProjectedSeq: traceAfterProjectedSeq, limit: TRACE_LIVE_PAGE_LIMIT });
if (signal.aborted) return accumulatedTrace;
if (!tracePolled.ok || !tracePolled.data) return accumulatedTrace;
onActivity();
const pageSnapshot = resultToTraceSnapshot(traceId, tracePolled.data, "trace-api");
accumulatedTrace = mergeTraceSnapshots(accumulatedTrace, pageSnapshot);
const nextProjectedSeq = traceSnapshotNextProjectedSeq(pageSnapshot, traceAfterProjectedSeq);
if (nextProjectedSeq > traceAfterProjectedSeq) traceAfterProjectedSeq = nextProjectedSeq;
if (pageSnapshot.hasMore !== true || nextProjectedSeq <= previousProjectedSeq) break;
}
return accumulatedTrace;
};
for (;;) {
if (signal.aborted) return;
await sleep(TRACE_POLL_INTERVAL_MS);
if (signal.aborted) return;
const turnPolled = await workbenchAPI.turn(traceId, inactivityTimeoutMs, activityRef);
if (signal.aborted) return;
if (turnPolled.ok && turnPolled.data) {
onActivity();
const turnSnapshot = resultToTraceSnapshot(traceId, turnPolled.data, "turn-api");
const liveTrace = await fetchLiveTracePages();
if (signal.aborted) return;
const snapshot = traceSnapshotWithTurnStatus(liveTrace ?? turnSnapshot, turnSnapshot);
const snapshotKey = traceSnapshotSignalKey(snapshot);
if (snapshotKey !== lastSnapshotKey) {
lastSnapshotKey = snapshotKey;
onSnapshot(snapshot);
}
if (turnPolled.data.terminal === true || isTerminalStatus(turnPolled.data.status)) {
onComplete(mergeTraceResults(turnPolled.data, accumulatedTrace ?? snapshot));
return;
}
} else if (!turnPolled.ok && turnPolled.status >= 500) {
// Backend hiccup: keep polling and let inactivity-timeout classify a real outage.
} else {
onInfrastructureError(turnPolled.error ?? "Code Agent turn status poll failed (non-5xx)");
return;
}
}
}
function resultToTraceSnapshot(traceId: string, result: AgentChatResultResponse, eventSource = "trace-api"): TraceSnapshot {
const events = Array.isArray(result.events) ? result.events : Array.isArray(result.traceEvents) ? result.traceEvents : [];
const lastEvent = events.at(-1);
const timing = result.timing ?? null;
return {
traceId: result.traceId ?? traceId,
status: result.status,
sessionId: result.sessionId ?? null,
threadId: result.threadId ?? null,
events,
eventCount: result.eventCount ?? events.length,
eventsCompacted: result.runnerTrace?.eventsCompacted,
fullTraceLoaded: result.fullTraceLoaded,
hasMore: result.hasMore,
truncated: result.truncated,
nextProjectedSeq: typeof result.nextProjectedSeq === "number" ? result.nextProjectedSeq : null,
range: result.range as TraceSnapshot["range"],
agentRun: result.agentRun,
traceStatus: result.traceStatus,
retention: result.retention,
terminalEvidence: result.terminalEvidence,
traceSummary: result.traceSummary,
error: result.error,
timing,
startedAt: result.startedAt ?? timing?.startedAt ?? null,
lastEventAt: result.lastEventAt ?? timing?.lastEventAt ?? null,
finishedAt: result.finishedAt ?? timing?.finishedAt ?? null,
durationMs: result.durationMs ?? timing?.durationMs ?? null,
projection: result.projection ?? result.runnerTrace?.projection ?? null,
projectionStatus: result.projectionStatus ?? result.runnerTrace?.projectionStatus ?? null,
projectionHealth: result.projectionHealth ?? result.runnerTrace?.projectionHealth ?? null,
staleMs: result.staleMs ?? result.runnerTrace?.staleMs ?? null,
blocker: result.blocker ?? result.runnerTrace?.blocker ?? null,
lastEventLabel: result.lastEventLabel ?? lastEvent?.label ?? lastEvent?.type,
eventSource,
updatedAt: new Date().toISOString()
};
}
function mergeTraceSnapshots(previous: TraceSnapshot | null, next: TraceSnapshot): TraceSnapshot {
if (!previous) return next;
const previousEvents = Array.isArray(previous.events) ? previous.events : [];
const nextEvents = Array.isArray(next.events) ? next.events : [];
const events = mergeTraceEvents(previousEvents, nextEvents);
const timing = mergeTraceTimingProjection(previous, next);
return {
...previous,
...next,
events,
eventCount: next.eventCount ?? previous.eventCount ?? events.length,
timing,
startedAt: timing?.startedAt ?? next.startedAt ?? next.timing?.startedAt ?? previous.startedAt ?? previous.timing?.startedAt ?? null,
lastEventAt: timing?.lastEventAt ?? next.lastEventAt ?? next.timing?.lastEventAt ?? previous.lastEventAt ?? previous.timing?.lastEventAt ?? null,
finishedAt: timing?.finishedAt ?? next.finishedAt ?? next.timing?.finishedAt ?? previous.finishedAt ?? previous.timing?.finishedAt ?? null,
durationMs: timing?.durationMs ?? next.durationMs ?? next.timing?.durationMs ?? previous.durationMs ?? previous.timing?.durationMs ?? null,
lastEventLabel: next.lastEventLabel ?? previous.lastEventLabel ?? undefined,
updatedAt: next.updatedAt ?? previous.updatedAt ?? new Date().toISOString()
};
}
function traceSnapshotWithTurnStatus(trace: TraceSnapshot, turn: TraceSnapshot): TraceSnapshot {
const timing = mergeTraceTimingProjection(trace, turn);
return {
...trace,
traceId: trace.traceId ?? turn.traceId,
status: turn.status ?? trace.status,
sessionId: trace.sessionId ?? turn.sessionId,
threadId: trace.threadId ?? turn.threadId,
agentRun: trace.agentRun ?? turn.agentRun,
traceStatus: trace.traceStatus ?? turn.traceStatus,
terminalEvidence: trace.terminalEvidence ?? turn.terminalEvidence,
traceSummary: trace.traceSummary ?? turn.traceSummary,
error: trace.error ?? turn.error,
timing,
startedAt: timing?.startedAt ?? turn.startedAt ?? turn.timing?.startedAt ?? trace.startedAt ?? trace.timing?.startedAt ?? null,
lastEventAt: timing?.lastEventAt ?? turn.lastEventAt ?? turn.timing?.lastEventAt ?? trace.lastEventAt ?? trace.timing?.lastEventAt ?? null,
finishedAt: timing?.finishedAt ?? turn.finishedAt ?? turn.timing?.finishedAt ?? trace.finishedAt ?? trace.timing?.finishedAt ?? null,
durationMs: timing?.durationMs ?? turn.durationMs ?? turn.timing?.durationMs ?? trace.durationMs ?? trace.timing?.durationMs ?? null,
projection: trace.projection ?? turn.projection ?? null,
projectionStatus: trace.projectionStatus ?? turn.projectionStatus ?? null,
projectionHealth: trace.projectionHealth ?? turn.projectionHealth ?? null,
staleMs: trace.staleMs ?? turn.staleMs ?? null,
blocker: trace.blocker ?? turn.blocker ?? null,
waitingFor: trace.waitingFor ?? turn.waitingFor,
updatedAt: new Date().toISOString()
};
}
function traceSnapshotNextProjectedSeq(snapshot: TraceSnapshot, fallback: number): number {
const direct = Number(snapshot.nextProjectedSeq ?? snapshot.range?.toProjectedSeq);
if (Number.isFinite(direct) && direct >= 0) return Math.trunc(direct);
const events = Array.isArray(snapshot.events) ? snapshot.events : [];
return events.reduce((max, event) => {
const seq = traceEventProjectedSeq(event);
return Number.isFinite(seq) && seq > max ? Math.trunc(seq) : max;
}, fallback);
}
function traceSnapshotSignalKey(snapshot: TraceSnapshot): string {
const events = Array.isArray(snapshot.events) ? snapshot.events : [];
return [
snapshot.status ?? "",
snapshot.traceStatus ?? "",
snapshot.eventCount ?? events.length,
events.length,
snapshot.nextProjectedSeq ?? snapshot.range?.toProjectedSeq ?? "",
snapshot.lastEventAt ?? snapshot.timing?.lastEventAt ?? "",
snapshot.durationMs ?? snapshot.timing?.durationMs ?? "",
snapshot.hasMore === true ? "more" : "caught-up",
snapshot.fullTraceLoaded === true ? "loaded" : "partial",
snapshot.lastEventLabel ?? "",
traceSnapshotErrorKey(snapshot.error),
snapshot.projectionHealth ?? snapshot.projection?.projectionHealth ?? "",
snapshot.projection?.blocker?.code ?? snapshot.blocker?.code ?? ""
].join("|");
}
function traceSnapshotErrorKey(error: TraceSnapshot["error"]): string {
if (!error) return "";
if (typeof error === "string") return error;
return firstNonEmptyString(error.code, error.message) ?? "";
}
function sleep(ms: number): Promise<void> {
return new Promise((resolve) => window.setTimeout(resolve, ms));
}
@@ -25,12 +25,12 @@ export interface WorkbenchRuntimePolicy {
sessionListRealtimeRefreshDelayMs: number;
sessionListTerminalRefreshDelayMs: number;
sessionListMinRefreshIntervalMs: number;
workbenchRealtimeErrorGapFillMinMs: number;
workbenchRealtimeErrorSyncReplayMinMs: number;
workbenchRealtimeFlushMaxItemsPerChunk: number;
workbenchRealtimeFlushMaxChunkMs: number;
workbenchRealtimeFlushYieldMs: number;
workbenchActiveTraceRestGapFillInitialMs: number;
workbenchActiveTraceRestGapFillRepeatMs: number;
workbenchActiveTraceSyncReplayInitialMs: number;
workbenchActiveTraceSyncReplayRepeatMs: number;
}
const DEFAULT_WORKBENCH_RUNTIME_POLICY: WorkbenchRuntimePolicy = Object.freeze({
@@ -57,12 +57,12 @@ const DEFAULT_WORKBENCH_RUNTIME_POLICY: WorkbenchRuntimePolicy = Object.freeze({
sessionListRealtimeRefreshDelayMs: 5_000,
sessionListTerminalRefreshDelayMs: 1_500,
sessionListMinRefreshIntervalMs: 15_000,
workbenchRealtimeErrorGapFillMinMs: 2_000,
workbenchRealtimeErrorSyncReplayMinMs: 2_000,
workbenchRealtimeFlushMaxItemsPerChunk: 4,
workbenchRealtimeFlushMaxChunkMs: 8,
workbenchRealtimeFlushYieldMs: 0,
workbenchActiveTraceRestGapFillInitialMs: 2_500,
workbenchActiveTraceRestGapFillRepeatMs: 5_000
workbenchActiveTraceSyncReplayInitialMs: 2_500,
workbenchActiveTraceSyncReplayRepeatMs: 5_000
});
export function workbenchRuntimePolicy(input: unknown = runtimePolicyConfig()): WorkbenchRuntimePolicy {
@@ -91,12 +91,12 @@ export function workbenchRuntimePolicy(input: unknown = runtimePolicyConfig()):
sessionListRealtimeRefreshDelayMs: nonNegativeNumber(source.sessionListRealtimeRefreshDelayMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.sessionListRealtimeRefreshDelayMs),
sessionListTerminalRefreshDelayMs: nonNegativeNumber(source.sessionListTerminalRefreshDelayMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.sessionListTerminalRefreshDelayMs),
sessionListMinRefreshIntervalMs: nonNegativeNumber(source.sessionListMinRefreshIntervalMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.sessionListMinRefreshIntervalMs),
workbenchRealtimeErrorGapFillMinMs: nonNegativeNumber(source.workbenchRealtimeErrorGapFillMinMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeErrorGapFillMinMs),
workbenchRealtimeErrorSyncReplayMinMs: nonNegativeNumber(source.workbenchRealtimeErrorSyncReplayMinMs ?? source.workbenchRealtimeErrorGapFillMinMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeErrorSyncReplayMinMs),
workbenchRealtimeFlushMaxItemsPerChunk: positiveInteger(source.workbenchRealtimeFlushMaxItemsPerChunk, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeFlushMaxItemsPerChunk),
workbenchRealtimeFlushMaxChunkMs: nonNegativeNumber(source.workbenchRealtimeFlushMaxChunkMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeFlushMaxChunkMs),
workbenchRealtimeFlushYieldMs: nonNegativeNumber(source.workbenchRealtimeFlushYieldMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeFlushYieldMs),
workbenchActiveTraceRestGapFillInitialMs: nonNegativeNumber(source.workbenchActiveTraceRestGapFillInitialMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchActiveTraceRestGapFillInitialMs),
workbenchActiveTraceRestGapFillRepeatMs: nonNegativeNumber(source.workbenchActiveTraceRestGapFillRepeatMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchActiveTraceRestGapFillRepeatMs)
workbenchActiveTraceSyncReplayInitialMs: nonNegativeNumber(source.workbenchActiveTraceSyncReplayInitialMs ?? source.workbenchActiveTraceRestGapFillInitialMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchActiveTraceSyncReplayInitialMs),
workbenchActiveTraceSyncReplayRepeatMs: nonNegativeNumber(source.workbenchActiveTraceSyncReplayRepeatMs ?? source.workbenchActiveTraceRestGapFillRepeatMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchActiveTraceSyncReplayRepeatMs)
};
}
@@ -1,7 +1,7 @@
// SPEC: PJ2026-0106050514 Workbench实时运行面 draft-2026-06-30-p0-1297-spec-first; PJ2026-0104010803 Workbench唯一投影 draft-2026-06-18-p0-unique-projection.
// Responsibility: Pure Workbench trace/message/projection merge helpers consumed by the store orchestration layer.
import { mergeRunnerTrace, type TraceSnapshot } from "@/composables/useTraceSubscription";
import { mergeRunnerTrace, type TraceSnapshot } from "@/composables/workbench-trace-snapshot";
import type { AgentChatResultResponse, AgentRunProvenance, ApiResult, ChatMessage, ProjectionDiagnostic, TraceEvent, WorkbenchTurnTimingProjection } from "@/types";
import { firstNonEmptyString } from "@/utils";
import { normalizeErrorDiagnostic, normalizeProjectionDiagnostic } from "@/utils/workbench-error-runtime";
+6 -6
View File
@@ -10,7 +10,7 @@ import { createWorkbenchHealthProbeCache } from "@/utils/workbench-health";
import { agentErrorFromProjection, normalizeApiErrorRecord, normalizeErrorDiagnostic, normalizeProjectionDiagnostic, projectionDiagnosticFromApiFailure, projectionDiagnosticFromFailure } from "@/utils/workbench-error-runtime";
import { readWorkbenchJson, readWorkbenchNumber, readWorkbenchString, removeWorkbenchStorageKey, writeWorkbenchJson, writeWorkbenchString } from "@/utils/workbench-storage-runtime";
import { createWorkbenchStreamTransportRuntime, type WorkbenchRealtimeEvent, type WorkbenchStreamTransportRecovery } from "@/utils/workbench-realtime-runtime";
import { mergeRunnerTrace, snapshotToRunnerTrace, type TraceSnapshot } from "@/composables/useTraceSubscription";
import { mergeRunnerTrace, snapshotToRunnerTrace, type TraceSnapshot } from "@/composables/workbench-trace-snapshot";
import type { WorkbenchMessagePageResponse, WorkbenchSessionDetailResponse } from "@/api/workbench";
import type { AgentChatResponse, AgentChatResultResponse, AgentRunProvenance, ApiError, ApiResult, ChatMessage, ErrorDiagnostic, LiveSurface, ProjectionBlocker, ProjectionDiagnostic, ProviderProfile, TraceEvent, WorkbenchSessionRecord, WorkbenchTurnTimingProjection } from "@/types";
import { firstNonEmptyString, nextProtocolId, normalizeWorkbenchSessionId, normalizeWorkbenchSessionRouteId } from "@/utils";
@@ -996,7 +996,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
const events = Array.isArray(result.events) ? result.events : Array.isArray(result.traceEvents) ? result.traceEvents : [];
const activityLabel = firstNonEmptyString(result.lastEventLabel, result.status);
if (ownerSessionId === activeSessionId.value && (events.length > 0 || turnResultIsTerminalForMerge(result))) recordActivity(`trace:${activityLabel ?? "hydrated"}`);
markWorkbenchTraceEventsReceived({ traceId, events, transport: "rest_gap" });
markWorkbenchTraceEventsReceived({ traceId, events, transport: "detail_history" });
if (turnResultIsTerminalForMerge(result)) {
const terminalSeal = terminalSealResultWithoutTraceEvents(result);
rememberTurnStatus(traceId, terminalSeal);
@@ -1097,7 +1097,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
}
function reattachTrace(traceId: string): void {
const initial: AgentChatResponse = { status: "running", traceId, turnUrl: `/v1/workbench/turns/${encodeURIComponent(traceId)}` };
const initial: AgentChatResponse = { status: "running", traceId };
if (!messages.value.some((message) => message.traceId === traceId)) appendActiveMessages(makeMessage("agent", "", "running", { traceId, sessionId: selectedSessionId.value ?? undefined, threadId: selectedThreadId.value ?? undefined, title: "Code Agent", traceAutoLifecycle: "running" }));
currentRequest.value = { traceId, sessionId: selectedSessionId.value ?? null, threadId: selectedThreadId.value ?? null, status: initial.status };
scheduleActiveTraceSyncReplay(traceId, "reattach-sync-replay");
@@ -1111,7 +1111,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
realtimeTransport.restart({
sessionId,
traceId,
errorRecoveryMinMs: runtimePolicy.workbenchRealtimeErrorGapFillMinMs,
errorRecoveryMinMs: runtimePolicy.workbenchRealtimeErrorSyncReplayMinMs,
flushMaxItemsPerChunk: runtimePolicy.workbenchRealtimeFlushMaxItemsPerChunk,
flushMaxChunkMs: runtimePolicy.workbenchRealtimeFlushMaxChunkMs,
flushYieldMs: runtimePolicy.workbenchRealtimeFlushYieldMs,
@@ -1154,7 +1154,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
recordWorkbenchRuntimeDiagnostic({ module: "workbench-sync-replay", sessionId, traceId, outcome: "ok", diagnostic: workbenchSyncReplayDiagnostic(result.data, events, { reason, sinceOutboxSeq }) });
}
function scheduleActiveTraceSyncReplay(traceId: string | null | undefined, reason: string, delayMs = runtimePolicy.workbenchActiveTraceRestGapFillInitialMs): void {
function scheduleActiveTraceSyncReplay(traceId: string | null | undefined, reason: string, delayMs = runtimePolicy.workbenchActiveTraceSyncReplayInitialMs): void {
const id = firstNonEmptyString(traceId);
if (!id) return;
if (typeof window === "undefined") {
@@ -1192,7 +1192,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
return;
}
const stillActive = currentRequest.value?.traceId === id || isTraceActiveStatus(turn?.status) || isTraceActiveStatus(message?.status);
if (stillActive) scheduleActiveTraceSyncReplay(id, "active-sync-replay:repeat", runtimePolicy.workbenchActiveTraceRestGapFillRepeatMs);
if (stillActive) scheduleActiveTraceSyncReplay(id, "active-sync-replay:repeat", runtimePolicy.workbenchActiveTraceSyncReplayRepeatMs);
}
function stopRealtime(): void {
+6
View File
@@ -37,8 +37,14 @@ declare global {
sessionListRealtimeRefreshDelayMs?: number;
sessionListTerminalRefreshDelayMs?: number;
sessionListMinRefreshIntervalMs?: number;
workbenchRealtimeErrorSyncReplayMinMs?: number;
workbenchActiveTraceSyncReplayInitialMs?: number;
workbenchActiveTraceSyncReplayRepeatMs?: number;
/** @deprecated Use workbenchRealtimeErrorSyncReplayMinMs. */
workbenchRealtimeErrorGapFillMinMs?: number;
/** @deprecated Use workbenchActiveTraceSyncReplayInitialMs. */
workbenchActiveTraceRestGapFillInitialMs?: number;
/** @deprecated Use workbenchActiveTraceSyncReplayRepeatMs. */
workbenchActiveTraceRestGapFillRepeatMs?: number;
};
};
@@ -83,7 +83,7 @@ interface WorkbenchPerformanceEvent {
interface TraceEventTimingInput {
traceId: string | null | undefined;
events: TraceEvent[];
transport: "sse" | "rest_gap" | "poll";
transport: "sse" | "sync_replay" | "detail_history" | "poll";
serverSentAt?: string | null | undefined;
eventCreatedAt?: string | null | undefined;
traceSeq?: number | string | null | undefined;
@@ -771,7 +771,8 @@ function normalizeBackend(value: unknown): string {
function normalizeTransport(value: unknown): string {
const text = safeText(value).toLowerCase();
return text === "rest_gap" || text === "poll" || text === "sse" ? text : "unknown";
if (text === "rest_gap") return "detail_history";
return text === "detail_history" || text === "sync_replay" || text === "poll" || text === "sse" ? text : "unknown";
}
function normalizeTargetState(value: unknown): string {
@@ -392,7 +392,7 @@ function outcomeLabel(value?: string): string {
}
function transportLabel(value?: string): string {
const labels: Record<string, string> = { sse: "实时流", rest_gap: "REST 补洞", poll: "轮询", none: "无", unknown: "未知" };
const labels: Record<string, string> = { sse: "实时流", sync_replay: "同步回放", detail_history: "明细/历史", poll: "轮询", none: "无", unknown: "未知" };
return labels[String(value ?? "unknown")] ?? String(value);
}