From 216c69d9f605abef2972248aa09c03beadccaac3 Mon Sep 17 00:00:00 2001 From: UniDesk Codex Date: Fri, 3 Jul 2026 02:33:28 +0800 Subject: [PATCH] fix: preserve workbench hydrate semantics --- ...bench-session-messages-read-budget.test.ts | 36 +++++++++++++++++ .../workbench-session-messages-read-budget.ts | 40 +++++++++++++++++++ web/hwlab-cloud-web/src/stores/workbench.ts | 26 ++++-------- 3 files changed, 84 insertions(+), 18 deletions(-) create mode 100644 web/hwlab-cloud-web/src/stores/workbench-session-messages-read-budget.test.ts create mode 100644 web/hwlab-cloud-web/src/stores/workbench-session-messages-read-budget.ts diff --git a/web/hwlab-cloud-web/src/stores/workbench-session-messages-read-budget.test.ts b/web/hwlab-cloud-web/src/stores/workbench-session-messages-read-budget.test.ts new file mode 100644 index 00000000..e48147a2 --- /dev/null +++ b/web/hwlab-cloud-web/src/stores/workbench-session-messages-read-budget.test.ts @@ -0,0 +1,36 @@ +import assert from "node:assert/strict"; +import { test } from "bun:test"; + +import type { ChatMessage } from "@/types"; +import { projectionMergeCommitSummary, workbenchSessionMessagesReadKey } from "./workbench-session-messages-read-budget"; +import { selectActiveTurnStatusRefreshTraceIds } from "./workbench-session"; + +function message(id: string): ChatMessage { + return { id, role: "agent", text: "", status: "running", createdAt: "2026-07-02T00:00:00.000Z" } as ChatMessage; +} + +test("session messages singleflight key separates limit and force class", () => { + const limit8 = workbenchSessionMessagesReadKey({ sessionId: "ses_key", limit: 8, force: false }); + const limit20 = workbenchSessionMessagesReadKey({ sessionId: "ses_key", limit: 20, force: false }); + const force20 = workbenchSessionMessagesReadKey({ sessionId: "ses_key", limit: 20, force: true }); + + assert.notEqual(limit8, limit20); + assert.notEqual(limit20, force20); + assert.equal(limit20, workbenchSessionMessagesReadKey({ sessionId: "ses_key", limit: 20, force: false })); +}); + +test("projection merge summary reports no write without blocking follow-up hydrate", () => { + const existing = [{ ...message("msg_same"), traceId: "trc_current", sessionId: "ses_same" }]; + const summary = projectionMergeCommitSummary(existing, existing); + const hydrateTraceIds = selectActiveTurnStatusRefreshTraceIds({ messages: existing, currentRequestTraceId: "trc_current", limit: 1 }); + + assert.deepEqual(summary, { changed: false, changedCount: 0 }); + assert.deepEqual(hydrateTraceIds, ["trc_current"]); +}); + +test("projection merge summary reports changed message references", () => { + const existing = [message("msg_old")]; + const merged = [message("msg_new")]; + + assert.deepEqual(projectionMergeCommitSummary(existing, merged), { changed: true, changedCount: 1 }); +}); diff --git a/web/hwlab-cloud-web/src/stores/workbench-session-messages-read-budget.ts b/web/hwlab-cloud-web/src/stores/workbench-session-messages-read-budget.ts new file mode 100644 index 00000000..01f50d92 --- /dev/null +++ b/web/hwlab-cloud-web/src/stores/workbench-session-messages-read-budget.ts @@ -0,0 +1,40 @@ +// SPEC: pikasTech/HWLAB#2356 Workbench bounded request storm follow-up. +// Responsibility: Pure budgeting helpers for session messages read de-duplication and projection merge writes. + +import type { ChatMessage } from "@/types"; +import { firstNonEmptyString } from "@/utils"; +import { composeWorkbenchScopedKey } from "@/utils/workbench-key"; + +export interface SessionMessagesReadKeyInput { + sessionId: string | null | undefined; + limit: number | null | undefined; + force?: boolean | null; +} + +export interface ProjectionMergeCommitSummary { + changed: boolean; + changedCount: number; +} + +export function workbenchSessionMessagesReadKey(input: SessionMessagesReadKeyInput): string { + return composeWorkbenchScopedKey( + "workbench.session-messages.read", + firstNonEmptyString(input.sessionId) ?? "unknown-session", + sessionMessagesReadLimitPart(input.limit), + input.force === true ? "force" : "fresh" + ); +} + +export function projectionMergeCommitSummary(existing: ChatMessage[], merged: ChatMessage[]): ProjectionMergeCommitSummary { + let changedCount = Math.max(0, merged.length - existing.length); + const sharedLength = Math.min(existing.length, merged.length); + for (let index = 0; index < sharedLength; index += 1) { + if (existing[index] !== merged[index]) changedCount += 1; + } + return { changed: changedCount > 0, changedCount }; +} + +function sessionMessagesReadLimitPart(value: number | null | undefined): number { + const number = Number(value); + return Number.isFinite(number) && number > 0 ? Math.trunc(number) : 1; +} diff --git a/web/hwlab-cloud-web/src/stores/workbench.ts b/web/hwlab-cloud-web/src/stores/workbench.ts index b3ba60a1..51f29349 100644 --- a/web/hwlab-cloud-web/src/stores/workbench.ts +++ b/web/hwlab-cloud-web/src/stores/workbench.ts @@ -70,6 +70,7 @@ import { planWorkbenchRealtimeApply, planWorkbenchRealtimeRecovery, type Workben import { useWorkbenchColadaMutations } from "./workbench-colada-mutations"; import { useWorkbenchColadaQueries } from "./workbench-colada-queries"; import { useWorkbenchColadaReducer } from "./workbench-colada-reducer"; +import { projectionMergeCommitSummary, workbenchSessionMessagesReadKey } from "./workbench-session-messages-read-budget"; const WORKBENCH_SESSION_PROJECTION_SIGNAL_CHANNEL = "hwlab.workbench.sessionProjection.v1"; const WORKBENCH_SESSION_PROJECTION_SIGNAL_KEY = "hwlab.workbench.sessionProjectionSignal.v1"; @@ -516,7 +517,7 @@ export const useWorkbenchStore = defineStore("workbench", () => { if (!response.ok || !response.data) return; const pageMessages = Array.isArray(response.data.messages) ? response.data.messages.map((message) => normalizeChatMessage(message as ChatMessage)) : []; const merged = mergeMessageProjectionPage(id, pageMessages, { limit: sessionMessageProjectionWindowLimit() }); - if (!commitMergedSessionMessages(id, pageMessages, merged, { limit: sessionMessageProjectionWindowLimit(), reason: "session-message-page" })) return; + commitMergedSessionMessages(id, pageMessages, merged, { limit: sessionMessageProjectionWindowLimit(), reason: "session-message-page" }); await hydrateTurnStatusAuthority(merged, { limit: 1, reason: "session-message-page" }); hydrateTerminalTraceGaps(merged, "session-message-page"); } @@ -542,13 +543,13 @@ export const useWorkbenchStore = defineStore("workbench", () => { } const pageMessages = Array.isArray(response.data.messages) ? response.data.messages.map((message) => normalizeChatMessage(message as ChatMessage)) : []; const merged = mergeMessageProjectionPage(id, pageMessages, { traceId, limit: traceMessageProjectionWindowLimit() }); - if (!commitMergedSessionMessages(id, pageMessages, merged, { traceId, limit: traceMessageProjectionWindowLimit(), reason: "trace-message-page" })) return; + commitMergedSessionMessages(id, pageMessages, merged, { traceId, limit: traceMessageProjectionWindowLimit(), reason: "trace-message-page" }); await hydrateTurnStatusAuthority(merged, { traceId, limit: 1, reason: "trace-message-page" }); if (!traceProjectionIsTerminalSealed(traceId, merged)) hydrateTerminalTraceGaps(merged, `trace-message-page:${traceId}`); } function fetchSessionMessagesPage(sessionId: string, options: SessionMessagesReadOptions): Promise> { - const key = composeWorkbenchScopedKey("workbench.session-messages.read", sessionId); + const key = workbenchSessionMessagesReadKey({ sessionId, limit: options.limit, force: options.force }); return sessionMessagesReadSingleflight.run(key, () => workbenchColadaQueries.fetchSessionMessages(sessionId, { limit: options.limit, minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force }), { reason: options.reason }); } @@ -567,9 +568,9 @@ export const useWorkbenchStore = defineStore("workbench", () => { return mergeBoundedProjectionMessages(existing, incoming); } - function commitMergedSessionMessages(sessionId: string, incoming: ChatMessage[], merged: ChatMessage[], options: { traceId?: string | null; limit: number; reason: string }): boolean { + function commitMergedSessionMessages(sessionId: string, incoming: ChatMessage[], merged: ChatMessage[], options: { traceId?: string | null; limit: number; reason: string }): void { const existing = serverState.value.messagesBySessionId[sessionId] ?? []; - const changedCount = projectionChangedMessageCount(existing, merged); + const summary = projectionMergeCommitSummary(existing, merged); recordWorkbenchRuntimeDiagnostic({ module: "workbench-message-projection", sessionId, @@ -582,23 +583,12 @@ export const useWorkbenchStore = defineStore("workbench", () => { inputCount: incoming.length, existingCount: existing.length, mergedCount: merged.length, - changedCount, + changedCount: summary.changedCount, limit: options.limit, valuesRedacted: true } }); - if (changedCount === 0) return false; - rememberSessionMessages(sessionId, merged); - return true; - } - - function projectionChangedMessageCount(existing: ChatMessage[], merged: ChatMessage[]): number { - let changed = Math.max(0, merged.length - existing.length); - const sharedLength = Math.min(existing.length, merged.length); - for (let index = 0; index < sharedLength; index += 1) { - if (existing[index] !== merged[index]) changed += 1; - } - return changed; + if (summary.changed) rememberSessionMessages(sessionId, merged); } function mergeMessageProjectionMessage(message: ChatMessage, existing: ChatMessage[]): ChatMessage {