fix: preserve workbench hydrate semantics
This commit is contained in:
@@ -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 });
|
||||
});
|
||||
@@ -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;
|
||||
}
|
||||
@@ -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<ApiResult<WorkbenchMessagePageResponse>> {
|
||||
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 {
|
||||
|
||||
Reference in New Issue
Block a user