refactor: move workbench realtime planning out of store

This commit is contained in:
UniDesk Codex
2026-06-30 23:58:59 +08:00
parent 586a740010
commit 6c09342f0c
4 changed files with 210 additions and 28 deletions
+9 -2
View File
@@ -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");
@@ -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>): 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 }
};
}
@@ -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<WorkbenchRealtimeEvent["turn"]> }
| { 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 [];
}
}
+48 -26
View File
@@ -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;
}
}