diff --git a/web/hwlab-cloud-web/scripts/check.ts b/web/hwlab-cloud-web/scripts/check.ts index 157916f1..7303736f 100644 --- a/web/hwlab-cloud-web/scripts/check.ts +++ b/web/hwlab-cloud-web/scripts/check.ts @@ -40,6 +40,7 @@ const requiredFiles = Object.freeze([ "src/stores/auth.ts", "src/stores/workbench.ts", "src/stores/workbench-event-reducer.ts", + "src/stores/workbench-realtime-plan.ts", "src/stores/workbench-timeline-model.ts", "src/stores/workbench-session-cache.ts", "src/composables/useWorkbenchScrollRuntime.ts", @@ -84,6 +85,7 @@ const workbenchRealtimeRuntimeSource = `${readWeb("src/utils/workbench-realtime- const workbenchRefreshRuntimeSource = readWeb("src/utils/workbench-refresh-runtime.ts"); const workbenchPerformanceSource = readWeb("src/utils/workbench-performance.ts"); const workbenchEventReducerSource = readWeb("src/stores/workbench-event-reducer.ts"); +const workbenchRealtimePlanSource = readWeb("src/stores/workbench-realtime-plan.ts"); const workbenchTimelineRuntimeSource = readWeb("src/stores/workbench-timeline-model.ts"); const workbenchScrollRuntimeSource = readWeb("src/composables/useWorkbenchScrollRuntime.ts"); const workbenchErrorRuntimeSource = readWeb("src/utils/workbench-error-runtime.ts"); @@ -140,6 +142,11 @@ assertIncludes(workbenchStoreSource, "createWorkbenchScheduledTaskRuntime", "Wor assertIncludes(workbenchEventReducerSource, "reduceWorkbenchRealtimeEvent", "Realtime event reducer must own SSE event classification"); assertIncludes(workbenchEventReducerSource, "workbench-event-reducer", "Realtime event reducer must emit module diagnostics for monitor root cause"); assertIncludes(workbenchStoreSource, "reduceWorkbenchRealtimeEvent", "Workbench store must consume migrated realtime reducer actions"); +assertIncludes(workbenchRealtimePlanSource, "planWorkbenchRealtimeApply", "Realtime planner must own store apply-step planning"); +assertIncludes(workbenchRealtimePlanSource, "planWorkbenchRealtimeRecovery", "Realtime planner must own recovery action planning"); +assertIncludes(workbenchStoreSource, "planWorkbenchRealtimeApply", "Workbench store must execute realtime apply plans instead of branching on reducer actions"); +assertIncludes(workbenchStoreSource, "planWorkbenchRealtimeRecovery", "Workbench store must execute realtime recovery plans instead of branching on recovery actions"); +assert.doesNotMatch(workbenchStoreSource, /function applyWorkbenchRealtimeAction[\s\S]{0,500}switch \(action\.type\)/u, "Workbench store must not switch directly on WorkbenchRealtimeAction"); assert.doesNotMatch(workbenchStoreSource, /function applyRealtimeEvent[\s\S]{0,900}event\.type ===/u, "Workbench store applyRealtimeEvent must not branch directly on raw SSE event.type"); assertIncludes(workbenchTimelineRuntimeSource, "TurnGap", "Timeline runtime must preserve OpenCode TurnGap row parity"); assertIncludes(workbenchTimelineRuntimeSource, "AssistantPart", "Timeline runtime must preserve OpenCode AssistantPart row parity"); @@ -163,8 +170,8 @@ assertIncludes(workbenchRefreshRuntimeSource, "if (existing)", "Scheduled refres assertIncludes(workbenchRefreshRuntimeSource, "replaceTimer", "Scheduled refresh runtime must make replacement explicit instead of ad hoc timer resets"); assertIncludes(workbenchPerformanceSource, "recordWorkbenchRuntimeDiagnostic", "Workbench performance probe must record runtime diagnostics for monitor root cause visibility"); assertIncludes(workbenchStoreSource, "recordWorkbenchRuntimeDiagnostic", "Workbench store must surface SSE recovery diagnostics to the performance probe"); -assertIncludes(workbenchStoreSource, "new Set(recovery.actions)", "Realtime recovery must consume transport-owned actions explicitly"); -assertIncludes(workbenchStoreSource, "actions.has(\"schedule-session-list\")", "Realtime stream errors must schedule bounded session list refreshes only when transport requests that action"); +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(workbenchStoreSource, "runtimePolicy.sessionListRealtimeRefreshDelayMs", "Realtime recovery delay must come from runtime policy instead of store constants"); 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"); diff --git a/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts b/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts index ae501925..04b0c3e6 100644 --- a/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts +++ b/web/hwlab-cloud-web/scripts/workbench-realtime-runtime.test.ts @@ -5,6 +5,7 @@ import assert from "node:assert/strict"; import test from "node:test"; import type { ChatMessage } from "../src/types/index.ts"; +import type { WorkbenchStreamTransportRecovery } from "../src/utils/workbench-realtime-runtime.ts"; import { workbenchRuntimePolicy } from "../src/config/workbench-runtime-policy.ts"; import { AsyncQueue, work } from "../src/utils/scheduler/async-queue.ts"; import { createCoalescedEventQueue } from "../src/utils/scheduler/coalesced-event-queue.ts"; @@ -16,6 +17,7 @@ import { checkWorkbenchHealth, createWorkbenchHealthProbeCache } from "../src/ut import { messageDiagnosticView } from "../src/utils/workbench-error-runtime.ts"; import { buildWorkbenchTimelineRows, normalizeWorkbenchTimelineMessages, workbenchTimelineSignature } from "../src/stores/workbench-timeline-model.ts"; import { reduceWorkbenchRealtimeEvent } from "../src/stores/workbench-event-reducer.ts"; +import { planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery } from "../src/stores/workbench-realtime-plan.ts"; import { cleanupWorkbenchServerStateDroppedSessions, cleanupWorkbenchServerStateSessions, createWorkbenchServerState, reduceWorkbenchServerState } from "../src/stores/workbench-server-state.ts"; import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "../src/stores/workbench-session-cache.ts"; @@ -279,6 +281,33 @@ test("realtime event reducer classifies SSE payloads before store side effects", assert.equal(error.diagnostic.traceId, "trc_2"); }); +test("realtime apply planner turns reducer actions into store steps", () => { + const reduced = reduceWorkbenchRealtimeEvent({ type: "trace.event", traceId: "trc_1", event: { traceId: "trc_1", label: "delta" }, snapshot: { traceId: "trc_1", status: "running" } }, "workbench.trace.event"); + const tracePlan = planWorkbenchRealtimeApply(reduced.action); + assert.deepEqual(tracePlan.steps.map((step) => step.type), ["apply-trace-event"]); + assert.equal(tracePlan.diagnostic.module, "workbench-realtime-plan"); + + const unavailable = planWorkbenchRealtimeApply({ type: "trace.unavailable", traceId: "trc_2", reason: "lost" }); + assert.deepEqual(unavailable.steps, [{ type: "clear-active-trace", traceId: "trc_2", reason: "lost" }]); + + const ignored = planWorkbenchRealtimeApply({ type: "ignore", reason: "unsupported" }); + assert.deepEqual(ignored.steps, []); +}); + +test("realtime recovery planner gates refresh steps by transport actions and authority", () => { + const recovery = recoveryEvent(["refresh-session-messages", "schedule-session-list", "refresh-turn-status", "hydrate-trace-events"]); + const authorized = planWorkbenchRealtimeRecovery(recovery, { selectedSessionId: "ses_1", activeSessionId: "ses_1", fallbackTraceId: "trc_1", activeTraceAuthorized: true }); + assert.deepEqual(authorized.steps.map((step) => step.type), ["refresh-session-messages", "schedule-session-list", "refresh-turn-status", "hydrate-trace-events"]); + assert.equal(authorized.sessionId, "ses_1"); + assert.equal(authorized.traceId, "trc_1"); + + const unauthorized = planWorkbenchRealtimeRecovery(recovery, { selectedSessionId: "ses_1", activeSessionId: "ses_1", fallbackTraceId: "trc_1", activeTraceAuthorized: false }); + assert.deepEqual(unauthorized.steps.map((step) => step.type), ["refresh-session-messages", "schedule-session-list"]); + + const inactive = planWorkbenchRealtimeRecovery(recovery, { selectedSessionId: "ses_1", activeSessionId: "ses_other", fallbackTraceId: "trc_1", activeTraceAuthorized: true }); + assert.deepEqual(inactive.steps.map((step) => step.type), ["schedule-session-list", "refresh-turn-status", "hydrate-trace-events"]); +}); + test("health probe cache records ok and unavailable states", async () => { const cache = createWorkbenchHealthProbeCache({ cacheMs: 100 }); const ok = await cache.probe({ key: "workbench", fetcher: async () => ({ ready: true }), classify: (value) => value.ready ? "ok" : "degraded" }); @@ -354,3 +383,15 @@ function agentMessage(overrides: Partial): ChatMessage { } as ChatMessage; } +function recoveryEvent(actions: WorkbenchStreamTransportRecovery["actions"]): WorkbenchStreamTransportRecovery { + return { + key: "workbench.realtime|ses_1|trc_1", + sessionId: "ses_1", + traceId: "trc_1", + tick: 1, + reason: "eventsource-error", + actions, + diagnostic: { code: "workbench_sse_recovery", valuesRedacted: true } + }; +} + diff --git a/web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts b/web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts new file mode 100644 index 00000000..b7615744 --- /dev/null +++ b/web/hwlab-cloud-web/src/stores/workbench-realtime-plan.ts @@ -0,0 +1,112 @@ +// SPEC: PJ2026-0106050514 Workbench实时运行面 draft-2026-06-30-p0-1297-spec-first. +// Responsibility: OpenCode-style realtime apply/recovery planning; the store executes plan steps without owning branching policy. +// Mechanical source: /root/opencode/packages/app/src/context/global-sync/event-reducer.ts:93-270 session/message reducer action boundary. +// Mechanical source: /root/opencode/packages/opencode/src/cli/cmd/run/session-data.ts:1-17 reducer side-effect boundary and replay notes. + +import type { WorkbenchRealtimeEvent } from "@/api/workbench-events"; +import type { WorkbenchStreamTransportRecovery } from "@/utils/workbench-realtime-runtime"; +import { firstNonEmptyString, normalizeWorkbenchSessionId } from "@/utils"; +import type { WorkbenchRealtimeAction } from "./workbench-event-reducer"; + +export type WorkbenchRealtimeApplyStep = + | { type: "apply-trace-snapshot"; traceId: string | null; snapshot: WorkbenchRealtimeEvent["snapshot"] } + | { type: "apply-trace-event"; traceId: string | null; event: WorkbenchRealtimeEvent["event"]; snapshot: WorkbenchRealtimeEvent["snapshot"]; realtimeEvent: WorkbenchRealtimeEvent } + | { type: "apply-message-snapshot"; realtimeEvent: WorkbenchRealtimeEvent } + | { type: "apply-turn-snapshot"; turn: NonNullable } + | { type: "apply-projection-error"; realtimeEvent: WorkbenchRealtimeEvent } + | { type: "clear-active-trace"; traceId: string; reason: string }; + +export interface WorkbenchRealtimeApplyPlan { + sourceActionType: WorkbenchRealtimeAction["type"]; + steps: WorkbenchRealtimeApplyStep[]; + diagnostic: { + module: "workbench-realtime-plan"; + sourceActionType: WorkbenchRealtimeAction["type"]; + stepTypes: string[]; + valuesRedacted: true; + }; +} + +export interface WorkbenchRealtimeRecoveryContext { + selectedSessionId?: string | null; + activeSessionId?: string | null; + fallbackTraceId?: string | null; + activeTraceAuthorized?: boolean; +} + +export type WorkbenchRealtimeRecoveryStep = + | { type: "refresh-session-messages"; sessionId: string; reason: string; force: true } + | { type: "schedule-session-list"; sessionId: string } + | { type: "refresh-turn-status"; traceId: string; force: true } + | { type: "hydrate-trace-events"; traceId: string; force: true }; + +export interface WorkbenchRealtimeRecoveryPlan { + sessionId: string | null; + traceId: string | null; + steps: WorkbenchRealtimeRecoveryStep[]; + diagnostic: { + module: "workbench-realtime-plan"; + reason: string; + actionTypes: string[]; + stepTypes: string[]; + valuesRedacted: true; + }; +} + +export function planWorkbenchRealtimeApply(action: WorkbenchRealtimeAction): WorkbenchRealtimeApplyPlan { + const steps = applySteps(action); + return { + sourceActionType: action.type, + steps, + diagnostic: { + module: "workbench-realtime-plan", + sourceActionType: action.type, + stepTypes: steps.map((step) => step.type), + valuesRedacted: true + } + }; +} + +export function planWorkbenchRealtimeRecovery(recovery: WorkbenchStreamTransportRecovery, context: WorkbenchRealtimeRecoveryContext): WorkbenchRealtimeRecoveryPlan { + const sessionId = normalizeWorkbenchSessionId(recovery.sessionId ?? context.selectedSessionId); + const traceId = firstNonEmptyString(recovery.traceId, context.fallbackTraceId); + const actions = new Set(recovery.actions); + const steps: WorkbenchRealtimeRecoveryStep[] = []; + if (actions.has("refresh-session-messages") && sessionId && sessionId === context.activeSessionId) steps.push({ type: "refresh-session-messages", sessionId, reason: "realtime-error:messages", force: true }); + if (actions.has("schedule-session-list") && sessionId) steps.push({ type: "schedule-session-list", sessionId }); + if (traceId && context.activeTraceAuthorized === true) { + if (actions.has("refresh-turn-status")) steps.push({ type: "refresh-turn-status", traceId, force: true }); + if (actions.has("hydrate-trace-events")) steps.push({ type: "hydrate-trace-events", traceId, force: true }); + } + return { + sessionId: sessionId ?? null, + traceId: traceId ?? null, + steps, + diagnostic: { + module: "workbench-realtime-plan", + reason: recovery.reason, + actionTypes: recovery.actions, + stepTypes: steps.map((step) => step.type), + valuesRedacted: true + } + }; +} + +function applySteps(action: WorkbenchRealtimeAction): WorkbenchRealtimeApplyStep[] { + switch (action.type) { + case "trace.snapshot": + return [{ type: "apply-trace-snapshot", traceId: action.traceId, snapshot: action.snapshot }]; + case "trace.event": + return [{ type: "apply-trace-event", traceId: action.traceId, event: action.event, snapshot: action.snapshot, realtimeEvent: action.realtimeEvent }]; + case "message.snapshot": + return [{ type: "apply-message-snapshot", realtimeEvent: action.realtimeEvent }]; + case "turn.snapshot": + return [{ type: "apply-turn-snapshot", turn: action.turn }]; + case "projection.error": + return [{ type: "apply-projection-error", realtimeEvent: action.realtimeEvent }]; + case "trace.unavailable": + return action.traceId ? [{ type: "clear-active-trace", traceId: action.traceId, reason: action.reason }] : []; + case "ignore": + return []; + } +} diff --git a/web/hwlab-cloud-web/src/stores/workbench.ts b/web/hwlab-cloud-web/src/stores/workbench.ts index 358f7d4c..d18cbf07 100644 --- a/web/hwlab-cloud-web/src/stores/workbench.ts +++ b/web/hwlab-cloud-web/src/stores/workbench.ts @@ -20,6 +20,7 @@ import { initialWorkbenchSessionIdFromLocation } from "./workbench-projection"; import { cleanupWorkbenchServerStateSessions, createWorkbenchServerState, reduceWorkbenchServerState, selectActiveMessages, selectActiveSession, selectSessionList, selectSessionStatusAuthority, selectTraceAuthorityById, selectTurnStatusAuthority, type WorkbenchServerAction } from "./workbench-server-state"; import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "./workbench-session-cache"; import { reduceWorkbenchRealtimeEvent, type WorkbenchRealtimeAction } from "./workbench-event-reducer"; +import { planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery, type WorkbenchRealtimeApplyStep, type WorkbenchRealtimeRecoveryStep } from "./workbench-realtime-plan"; const WORKBENCH_SESSION_PROJECTION_SIGNAL_CHANNEL = "hwlab.workbench.sessionProjection.v1"; const WORKBENCH_SESSION_PROJECTION_SIGNAL_KEY = "hwlab.workbench.sessionProjectionSignal.v1"; @@ -1021,17 +1022,35 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project } function handleRealtimeRecovery(recovery: WorkbenchStreamTransportRecovery): void { - const activeId = normalizeWorkbenchSessionId(recovery.sessionId ?? selectedSessionId.value); - const activeTraceId = firstNonEmptyString(recovery.traceId, realtimeTraceId()); - const actions = new Set(recovery.actions); - recordWorkbenchRuntimeDiagnostic({ module: "workbench-stream-transport", diagnostic: recovery.diagnostic, sessionId: activeId, traceId: activeTraceId, outcome: "network" }); - if (!activeId && !activeTraceId) return; - if (actions.has("refresh-session-messages") && activeId === activeSessionId.value) void refreshRealtimeSessionMessages(activeId, "realtime-error:messages", { force: true }); - if (actions.has("schedule-session-list") && activeId) scheduleSessionListRefresh(activeId, runtimePolicy.sessionListRealtimeRefreshDelayMs); - if (!activeTraceId || !shouldApplyActiveTraceAuthority(activeTraceId, activeId)) return; - if (actions.has("refresh-turn-status")) void refreshTurnStatusByTraceId(activeTraceId, { force: true }); - const message = latestMessageForTrace(activeTraceId); - if (actions.has("hydrate-trace-events") && message) void hydrateTraceEventsForMessage(message, { force: true }); + const sessionId = normalizeWorkbenchSessionId(recovery.sessionId ?? selectedSessionId.value); + const traceId = firstNonEmptyString(recovery.traceId, realtimeTraceId()); + const plan = planWorkbenchRealtimeRecovery(recovery, { + selectedSessionId: selectedSessionId.value, + activeSessionId: activeSessionId.value, + fallbackTraceId: traceId, + activeTraceAuthorized: Boolean(traceId && shouldApplyActiveTraceAuthority(traceId, sessionId)) + }); + recordWorkbenchRuntimeDiagnostic({ module: "workbench-stream-transport", diagnostic: recovery.diagnostic, sessionId: plan.sessionId, traceId: plan.traceId, outcome: "network" }); + for (const step of plan.steps) executeRealtimeRecoveryStep(step); + } + + function executeRealtimeRecoveryStep(step: WorkbenchRealtimeRecoveryStep): void { + switch (step.type) { + case "refresh-session-messages": + void refreshRealtimeSessionMessages(step.sessionId, step.reason, { force: step.force }); + return; + case "schedule-session-list": + scheduleSessionListRefresh(step.sessionId, runtimePolicy.sessionListRealtimeRefreshDelayMs); + return; + case "refresh-turn-status": + void refreshTurnStatusByTraceId(step.traceId, { force: step.force }); + return; + case "hydrate-trace-events": { + const message = latestMessageForTrace(step.traceId); + if (message) void hydrateTraceEventsForMessage(message, { force: step.force }); + return; + } + } } function scheduleActiveTraceRestGapFill(traceId: string | null | undefined, reason: string, delayMs = runtimePolicy.workbenchActiveTraceRestGapFillInitialMs): void { @@ -1153,26 +1172,29 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project } function applyWorkbenchRealtimeAction(action: WorkbenchRealtimeAction): void { - switch (action.type) { - case "trace.snapshot": - applyRealtimeTraceSnapshot(action.traceId, action.snapshot); + const plan = planWorkbenchRealtimeApply(action); + for (const step of plan.steps) applyWorkbenchRealtimePlanStep(step); + } + + function applyWorkbenchRealtimePlanStep(step: WorkbenchRealtimeApplyStep): void { + switch (step.type) { + case "apply-trace-snapshot": + applyRealtimeTraceSnapshot(step.traceId, step.snapshot); return; - case "trace.event": - applyRealtimeTraceEvent(action.traceId, action.event, action.snapshot, action.realtimeEvent); + case "apply-trace-event": + applyRealtimeTraceEvent(step.traceId, step.event, step.snapshot, step.realtimeEvent); return; - case "message.snapshot": - applyRealtimeMessageSnapshot(action.realtimeEvent); + case "apply-message-snapshot": + applyRealtimeMessageSnapshot(step.realtimeEvent); return; - case "turn.snapshot": - applyRealtimeTurnSnapshot(action.turn); + case "apply-turn-snapshot": + applyRealtimeTurnSnapshot(step.turn); return; - case "projection.error": - applyRealtimeProjectionError(action.realtimeEvent); + case "apply-projection-error": + applyRealtimeProjectionError(step.realtimeEvent); return; - case "trace.unavailable": - if (action.traceId) void clearActiveTrace(action.traceId, action.reason); - return; - case "ignore": + case "clear-active-trace": + void clearActiveTrace(step.traceId, step.reason); return; } }