Merge pull request #2352 from pikasTech/feat/2351-workbench-merge-cost
fix: 降低 Workbench 消息投影自动刷新成本
This commit is contained in:
@@ -0,0 +1,79 @@
|
|||||||
|
// SPEC: pikasTech/HWLAB#2351 Workbench JSON parse/projection merge cost reduction.
|
||||||
|
// Responsibility: Regression tests for bounded Workbench message projection refresh planning.
|
||||||
|
|
||||||
|
import assert from "node:assert/strict";
|
||||||
|
import test from "node:test";
|
||||||
|
|
||||||
|
import type { ChatMessage } from "../src/types/index.ts";
|
||||||
|
import { boundedProjectionMessageLimit, mergeBoundedProjectionMessages, selectProjectionMessageWindow, traceProjectionIsTerminalSealed } from "../src/stores/workbench-message-projection-budget.ts";
|
||||||
|
|
||||||
|
test("projection window keeps only bounded tail for automatic session refresh", () => {
|
||||||
|
const messages = Array.from({ length: 8 }, (_, index) => message(`msg_${index}`, `trc_${index}`));
|
||||||
|
const selected = selectProjectionMessageWindow(messages, { limit: 3 });
|
||||||
|
|
||||||
|
assert.deepEqual(selected.map((item) => item.id), ["msg_5", "msg_6", "msg_7"]);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("projection window keeps active trace outside the bounded tail", () => {
|
||||||
|
const messages = Array.from({ length: 8 }, (_, index) => message(`msg_${index}`, `trc_${index}`));
|
||||||
|
const selected = selectProjectionMessageWindow(messages, { traceId: "trc_2", limit: 3 });
|
||||||
|
|
||||||
|
assert.deepEqual(selected.map((item) => item.id), ["msg_2", "msg_5", "msg_6", "msg_7"]);
|
||||||
|
});
|
||||||
|
|
||||||
|
test("bounded projection merge updates window messages without dropping history", () => {
|
||||||
|
const existing = [
|
||||||
|
message("msg_1", "trc_1", { text: "old one" }),
|
||||||
|
message("msg_2", "trc_2", { text: "old two" }),
|
||||||
|
message("msg_3", "trc_3", { text: "old three" })
|
||||||
|
];
|
||||||
|
const incoming = [
|
||||||
|
message("msg_2", "trc_2", { text: "new two" }),
|
||||||
|
message("msg_4", "trc_4", { text: "new four" })
|
||||||
|
];
|
||||||
|
|
||||||
|
const merged = mergeBoundedProjectionMessages(existing, incoming);
|
||||||
|
|
||||||
|
assert.deepEqual(merged.map((item) => item.id), ["msg_1", "msg_2", "msg_3", "msg_4"]);
|
||||||
|
assert.equal(merged[0]?.text, "old one");
|
||||||
|
assert.equal(merged[1]?.text, "new two");
|
||||||
|
assert.equal(merged[3]?.text, "new four");
|
||||||
|
});
|
||||||
|
|
||||||
|
test("bounded projection merge replaces optimistic agent row by trace when message ids differ", () => {
|
||||||
|
const existing = [message("msg_optimistic", "trc_same", { text: "optimistic" })];
|
||||||
|
const incoming = [message("msg_canonical", "trc_same", { text: "canonical" })];
|
||||||
|
|
||||||
|
const merged = mergeBoundedProjectionMessages(existing, incoming);
|
||||||
|
|
||||||
|
assert.equal(merged.length, 1);
|
||||||
|
assert.equal(merged[0]?.id, "msg_canonical");
|
||||||
|
assert.equal(merged[0]?.text, "canonical");
|
||||||
|
});
|
||||||
|
|
||||||
|
test("terminal sealed trace skips redundant message projection refresh", () => {
|
||||||
|
const sealed = message("msg_terminal", "trc_terminal", {
|
||||||
|
status: "completed",
|
||||||
|
text: "final answer",
|
||||||
|
finalResponse: { text: "final answer" }
|
||||||
|
});
|
||||||
|
const unsealed = message("msg_unsealed", "trc_unsealed", { status: "completed", text: "" });
|
||||||
|
|
||||||
|
assert.equal(traceProjectionIsTerminalSealed("trc_terminal", [sealed, unsealed]), true);
|
||||||
|
assert.equal(traceProjectionIsTerminalSealed("trc_unsealed", [sealed, unsealed]), false);
|
||||||
|
assert.equal(boundedProjectionMessageLimit(0, 5), 5);
|
||||||
|
});
|
||||||
|
|
||||||
|
function message(id: string, traceId: string, extra: Partial<ChatMessage> = {}): ChatMessage {
|
||||||
|
return {
|
||||||
|
id,
|
||||||
|
messageId: id,
|
||||||
|
role: "agent",
|
||||||
|
title: "Code Agent",
|
||||||
|
text: "",
|
||||||
|
status: "running",
|
||||||
|
traceId,
|
||||||
|
createdAt: "2026-07-02T00:00:00.000Z",
|
||||||
|
...extra
|
||||||
|
};
|
||||||
|
}
|
||||||
@@ -34,6 +34,8 @@ test("Workbench runtime policy reads injected config while preserving defaults",
|
|||||||
sessionListPageLimit: 7,
|
sessionListPageLimit: 7,
|
||||||
traceHydrationMaxPages: 2,
|
traceHydrationMaxPages: 2,
|
||||||
workbenchSessionDetailMinRefreshMs: 1234,
|
workbenchSessionDetailMinRefreshMs: 1234,
|
||||||
|
workbenchSessionMessagesWindowLimit: 9,
|
||||||
|
workbenchTraceMessagesWindowLimit: 4,
|
||||||
workbenchRealtimeErrorGapFillMinMs: 0,
|
workbenchRealtimeErrorGapFillMinMs: 0,
|
||||||
defaultGatewayTimeoutMs: "bad"
|
defaultGatewayTimeoutMs: "bad"
|
||||||
});
|
});
|
||||||
@@ -41,6 +43,8 @@ test("Workbench runtime policy reads injected config while preserving defaults",
|
|||||||
assert.equal(policy.sessionListPageLimit, 7);
|
assert.equal(policy.sessionListPageLimit, 7);
|
||||||
assert.equal(policy.traceHydrationMaxPages, 2);
|
assert.equal(policy.traceHydrationMaxPages, 2);
|
||||||
assert.equal(policy.workbenchSessionDetailMinRefreshMs, 1234);
|
assert.equal(policy.workbenchSessionDetailMinRefreshMs, 1234);
|
||||||
|
assert.equal(policy.workbenchSessionMessagesWindowLimit, 9);
|
||||||
|
assert.equal(policy.workbenchTraceMessagesWindowLimit, 4);
|
||||||
assert.equal(policy.workbenchRealtimeErrorGapFillMinMs, 0);
|
assert.equal(policy.workbenchRealtimeErrorGapFillMinMs, 0);
|
||||||
assert.equal(policy.defaultGatewayTimeoutMs, 120_000);
|
assert.equal(policy.defaultGatewayTimeoutMs, 120_000);
|
||||||
});
|
});
|
||||||
@@ -450,4 +454,3 @@ function recoveryEvent(actions: WorkbenchStreamTransportRecovery["actions"]): Wo
|
|||||||
diagnostic: { code: "workbench_sse_recovery", valuesRedacted: true }
|
diagnostic: { code: "workbench_sse_recovery", valuesRedacted: true }
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -18,6 +18,8 @@ export interface WorkbenchRuntimePolicy {
|
|||||||
workbenchTraceEventsMinRefreshMs: number;
|
workbenchTraceEventsMinRefreshMs: number;
|
||||||
workbenchSessionDetailMinRefreshMs: number;
|
workbenchSessionDetailMinRefreshMs: number;
|
||||||
workbenchSessionMessagesMinRefreshMs: number;
|
workbenchSessionMessagesMinRefreshMs: number;
|
||||||
|
workbenchSessionMessagesWindowLimit: number;
|
||||||
|
workbenchTraceMessagesWindowLimit: number;
|
||||||
workbenchRealtimeSessionMessagesMinRefreshMs: number;
|
workbenchRealtimeSessionMessagesMinRefreshMs: number;
|
||||||
workbenchTraceEventsTimeoutMs: number;
|
workbenchTraceEventsTimeoutMs: number;
|
||||||
sessionListRealtimeRefreshDelayMs: number;
|
sessionListRealtimeRefreshDelayMs: number;
|
||||||
@@ -45,6 +47,8 @@ const DEFAULT_WORKBENCH_RUNTIME_POLICY: WorkbenchRuntimePolicy = Object.freeze({
|
|||||||
workbenchTraceEventsMinRefreshMs: 4_000,
|
workbenchTraceEventsMinRefreshMs: 4_000,
|
||||||
workbenchSessionDetailMinRefreshMs: 5_000,
|
workbenchSessionDetailMinRefreshMs: 5_000,
|
||||||
workbenchSessionMessagesMinRefreshMs: 5_000,
|
workbenchSessionMessagesMinRefreshMs: 5_000,
|
||||||
|
workbenchSessionMessagesWindowLimit: 20,
|
||||||
|
workbenchTraceMessagesWindowLimit: 8,
|
||||||
workbenchRealtimeSessionMessagesMinRefreshMs: 1_000,
|
workbenchRealtimeSessionMessagesMinRefreshMs: 1_000,
|
||||||
workbenchTraceEventsTimeoutMs: 5_000,
|
workbenchTraceEventsTimeoutMs: 5_000,
|
||||||
sessionListRealtimeRefreshDelayMs: 5_000,
|
sessionListRealtimeRefreshDelayMs: 5_000,
|
||||||
@@ -74,6 +78,8 @@ export function workbenchRuntimePolicy(input: unknown = runtimePolicyConfig()):
|
|||||||
workbenchTraceEventsMinRefreshMs: nonNegativeNumber(source.workbenchTraceEventsMinRefreshMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchTraceEventsMinRefreshMs),
|
workbenchTraceEventsMinRefreshMs: nonNegativeNumber(source.workbenchTraceEventsMinRefreshMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchTraceEventsMinRefreshMs),
|
||||||
workbenchSessionDetailMinRefreshMs: nonNegativeNumber(source.workbenchSessionDetailMinRefreshMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchSessionDetailMinRefreshMs),
|
workbenchSessionDetailMinRefreshMs: nonNegativeNumber(source.workbenchSessionDetailMinRefreshMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchSessionDetailMinRefreshMs),
|
||||||
workbenchSessionMessagesMinRefreshMs: nonNegativeNumber(source.workbenchSessionMessagesMinRefreshMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchSessionMessagesMinRefreshMs),
|
workbenchSessionMessagesMinRefreshMs: nonNegativeNumber(source.workbenchSessionMessagesMinRefreshMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchSessionMessagesMinRefreshMs),
|
||||||
|
workbenchSessionMessagesWindowLimit: positiveInteger(source.workbenchSessionMessagesWindowLimit, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchSessionMessagesWindowLimit),
|
||||||
|
workbenchTraceMessagesWindowLimit: positiveInteger(source.workbenchTraceMessagesWindowLimit, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchTraceMessagesWindowLimit),
|
||||||
workbenchRealtimeSessionMessagesMinRefreshMs: nonNegativeNumber(source.workbenchRealtimeSessionMessagesMinRefreshMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeSessionMessagesMinRefreshMs),
|
workbenchRealtimeSessionMessagesMinRefreshMs: nonNegativeNumber(source.workbenchRealtimeSessionMessagesMinRefreshMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeSessionMessagesMinRefreshMs),
|
||||||
workbenchTraceEventsTimeoutMs: positiveNumber(source.workbenchTraceEventsTimeoutMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchTraceEventsTimeoutMs),
|
workbenchTraceEventsTimeoutMs: positiveNumber(source.workbenchTraceEventsTimeoutMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchTraceEventsTimeoutMs),
|
||||||
sessionListRealtimeRefreshDelayMs: nonNegativeNumber(source.sessionListRealtimeRefreshDelayMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.sessionListRealtimeRefreshDelayMs),
|
sessionListRealtimeRefreshDelayMs: nonNegativeNumber(source.sessionListRealtimeRefreshDelayMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.sessionListRealtimeRefreshDelayMs),
|
||||||
|
|||||||
@@ -0,0 +1,72 @@
|
|||||||
|
// SPEC: pikasTech/HWLAB#2351 Workbench JSON parse/projection merge cost reduction.
|
||||||
|
// Responsibility: Pure bounded-window planning for Workbench session message projection refreshes.
|
||||||
|
|
||||||
|
import type { ChatMessage } from "@/types";
|
||||||
|
import { firstNonEmptyString } from "@/utils";
|
||||||
|
import { isTerminalMessageStatus, messageHasTerminalResponse } from "./workbench-message-projection-runtime";
|
||||||
|
|
||||||
|
export interface ProjectionWindowOptions {
|
||||||
|
traceId?: string | null;
|
||||||
|
limit: number;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function boundedProjectionMessageLimit(value: number | null | undefined, fallback: number): number {
|
||||||
|
const number = Number(value);
|
||||||
|
if (Number.isFinite(number) && number > 0) return Math.max(1, Math.trunc(number));
|
||||||
|
return Math.max(1, Math.trunc(fallback));
|
||||||
|
}
|
||||||
|
|
||||||
|
export function selectProjectionMessageWindow(source: ChatMessage[], options: ProjectionWindowOptions): ChatMessage[] {
|
||||||
|
const limit = boundedProjectionMessageLimit(options.limit, source.length || 1);
|
||||||
|
const tail = source.slice(-limit);
|
||||||
|
const traceId = firstNonEmptyString(options.traceId);
|
||||||
|
if (!traceId) return tail;
|
||||||
|
const selected = new Set(tail);
|
||||||
|
for (const message of source) {
|
||||||
|
if (messageTraceId(message) === traceId) selected.add(message);
|
||||||
|
}
|
||||||
|
return source.filter((message) => selected.has(message));
|
||||||
|
}
|
||||||
|
|
||||||
|
export function mergeBoundedProjectionMessages(existing: ChatMessage[], incoming: ChatMessage[]): ChatMessage[] {
|
||||||
|
if (existing.length === 0) return incoming;
|
||||||
|
if (incoming.length === 0) return existing;
|
||||||
|
const incomingByKey = new Map<string, ChatMessage>();
|
||||||
|
for (const message of incoming) {
|
||||||
|
const key = messageProjectionKey(message);
|
||||||
|
if (key) incomingByKey.set(key, message);
|
||||||
|
}
|
||||||
|
if (incomingByKey.size === 0) return existing;
|
||||||
|
const seen = new Set<string>();
|
||||||
|
const merged = existing.map((message) => {
|
||||||
|
const key = messageProjectionKey(message);
|
||||||
|
if (!key) return message;
|
||||||
|
seen.add(key);
|
||||||
|
return incomingByKey.get(key) ?? message;
|
||||||
|
});
|
||||||
|
for (const message of incoming) {
|
||||||
|
const key = messageProjectionKey(message);
|
||||||
|
if (!key || seen.has(key)) continue;
|
||||||
|
seen.add(key);
|
||||||
|
merged.push(message);
|
||||||
|
}
|
||||||
|
return merged;
|
||||||
|
}
|
||||||
|
|
||||||
|
export function traceProjectionIsTerminalSealed(traceId: string | null | undefined, source: ChatMessage[]): boolean {
|
||||||
|
const id = firstNonEmptyString(traceId);
|
||||||
|
if (!id) return false;
|
||||||
|
const message = [...source].reverse().find((item) => item.role === "agent" && messageTraceId(item) === id) ?? null;
|
||||||
|
return Boolean(message && isTerminalMessageStatus(message.status) && messageHasTerminalResponse(message));
|
||||||
|
}
|
||||||
|
|
||||||
|
function messageProjectionKey(message: ChatMessage): string | null {
|
||||||
|
const traceId = messageTraceId(message);
|
||||||
|
if (traceId) return `trace:${message.role}:${traceId}`;
|
||||||
|
const messageId = firstNonEmptyString(message.messageId, message.id);
|
||||||
|
return messageId ? `message:${messageId}` : null;
|
||||||
|
}
|
||||||
|
|
||||||
|
function messageTraceId(message: ChatMessage): string | null {
|
||||||
|
return firstNonEmptyString(message.traceId, message.runnerTrace?.traceId) ?? null;
|
||||||
|
}
|
||||||
@@ -20,6 +20,7 @@ import { cleanupWorkbenchServerStateSessions, selectActiveMessages, selectActive
|
|||||||
import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "./workbench-session-cache";
|
import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "./workbench-session-cache";
|
||||||
import { reduceWorkbenchRealtimeEvent, type WorkbenchRealtimeAction } from "./workbench-event-reducer";
|
import { reduceWorkbenchRealtimeEvent, type WorkbenchRealtimeAction } from "./workbench-event-reducer";
|
||||||
import { messageHasSealedTerminalResult, messageIsSealedTerminal, traceAuthorityIsSealed } from "./workbench-terminal-authority";
|
import { messageHasSealedTerminalResult, messageIsSealedTerminal, traceAuthorityIsSealed } from "./workbench-terminal-authority";
|
||||||
|
import { boundedProjectionMessageLimit, mergeBoundedProjectionMessages, selectProjectionMessageWindow, traceProjectionIsTerminalSealed } from "./workbench-message-projection-budget";
|
||||||
import {
|
import {
|
||||||
agentErrorDisplayText,
|
agentErrorDisplayText,
|
||||||
agentErrorFromApiFailure,
|
agentErrorFromApiFailure,
|
||||||
@@ -491,10 +492,10 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
|||||||
async function refreshSessionMessageProjectionPage(sessionId: string | null | undefined, options: { force?: boolean } = {}): Promise<void> {
|
async function refreshSessionMessageProjectionPage(sessionId: string | null | undefined, options: { force?: boolean } = {}): Promise<void> {
|
||||||
const id = normalizeWorkbenchSessionId(sessionId);
|
const id = normalizeWorkbenchSessionId(sessionId);
|
||||||
if (!id) return;
|
if (!id) return;
|
||||||
const response = await workbenchColadaQueries.fetchSessionMessages(id, { limit: 100, minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force });
|
const response = await workbenchColadaQueries.fetchSessionMessages(id, { limit: sessionMessageProjectionWindowLimit(), minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force });
|
||||||
if (!response.ok || !response.data) return;
|
if (!response.ok || !response.data) return;
|
||||||
const pageMessages = Array.isArray(response.data.messages) ? response.data.messages.map((message) => normalizeChatMessage(message as ChatMessage)) : [];
|
const pageMessages = Array.isArray(response.data.messages) ? response.data.messages.map((message) => normalizeChatMessage(message as ChatMessage)) : [];
|
||||||
const merged = mergeMessageProjectionPage(id, pageMessages);
|
const merged = mergeMessageProjectionPage(id, pageMessages, { limit: sessionMessageProjectionWindowLimit() });
|
||||||
rememberSessionMessages(id, merged);
|
rememberSessionMessages(id, merged);
|
||||||
await hydrateTurnStatusAuthority(merged);
|
await hydrateTurnStatusAuthority(merged);
|
||||||
hydrateTerminalTraceGaps(merged, "session-message-page");
|
hydrateTerminalTraceGaps(merged, "session-message-page");
|
||||||
@@ -510,7 +511,9 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
|||||||
async function refreshMessageProjectionForTrace(sessionId: string | null | undefined, traceId: string, options: { force?: boolean } = {}): Promise<void> {
|
async function refreshMessageProjectionForTrace(sessionId: string | null | undefined, traceId: string, options: { force?: boolean } = {}): Promise<void> {
|
||||||
const id = normalizeWorkbenchSessionId(sessionId);
|
const id = normalizeWorkbenchSessionId(sessionId);
|
||||||
if (!id) return;
|
if (!id) return;
|
||||||
const response = await workbenchColadaQueries.fetchSessionMessages(id, { limit: 100, minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force });
|
const existing = serverState.value.messagesBySessionId[id] ?? [];
|
||||||
|
if (traceProjectionIsTerminalSealed(traceId, existing)) return;
|
||||||
|
const response = await workbenchColadaQueries.fetchSessionMessages(id, { limit: traceMessageProjectionWindowLimit(), minIntervalMs: runtimePolicy.workbenchSessionMessagesMinRefreshMs, force: options.force });
|
||||||
if (!response.ok || !response.data) {
|
if (!response.ok || !response.data) {
|
||||||
if (shouldSuppressTransientWorkbenchReadFailure(response)) return;
|
if (shouldSuppressTransientWorkbenchReadFailure(response)) return;
|
||||||
if (traceHasTerminalResponse(traceId, messages.value)) return;
|
if (traceHasTerminalResponse(traceId, messages.value)) return;
|
||||||
@@ -518,15 +521,25 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
const pageMessages = Array.isArray(response.data.messages) ? response.data.messages.map((message) => normalizeChatMessage(message as ChatMessage)) : [];
|
const pageMessages = Array.isArray(response.data.messages) ? response.data.messages.map((message) => normalizeChatMessage(message as ChatMessage)) : [];
|
||||||
const merged = mergeMessageProjectionPage(id, pageMessages);
|
const merged = mergeMessageProjectionPage(id, pageMessages, { traceId, limit: traceMessageProjectionWindowLimit() });
|
||||||
rememberSessionMessages(id, merged);
|
rememberSessionMessages(id, merged);
|
||||||
await hydrateTurnStatusAuthority(merged);
|
await hydrateTurnStatusAuthority(merged);
|
||||||
hydrateTerminalTraceGaps(merged, `trace-message-page:${traceId}`);
|
if (!traceProjectionIsTerminalSealed(traceId, merged)) hydrateTerminalTraceGaps(merged, `trace-message-page:${traceId}`);
|
||||||
}
|
}
|
||||||
|
|
||||||
function mergeMessageProjectionPage(sessionId: string, pageMessages: ChatMessage[]): ChatMessage[] {
|
function sessionMessageProjectionWindowLimit(): number {
|
||||||
|
return boundedProjectionMessageLimit(runtimePolicy.workbenchSessionMessagesWindowLimit, runtimePolicy.sessionListPageLimit);
|
||||||
|
}
|
||||||
|
|
||||||
|
function traceMessageProjectionWindowLimit(): number {
|
||||||
|
return boundedProjectionMessageLimit(runtimePolicy.workbenchTraceMessagesWindowLimit, runtimePolicy.traceHydrationAutoQueueLimit);
|
||||||
|
}
|
||||||
|
|
||||||
|
function mergeMessageProjectionPage(sessionId: string, pageMessages: ChatMessage[], options: { traceId?: string | null; limit: number }): ChatMessage[] {
|
||||||
const existing = serverState.value.messagesBySessionId[sessionId] ?? [];
|
const existing = serverState.value.messagesBySessionId[sessionId] ?? [];
|
||||||
return pageMessages.map((message) => mergeMessageProjectionMessage(message, existing));
|
const windowMessages = selectProjectionMessageWindow(pageMessages, options);
|
||||||
|
const incoming = windowMessages.map((message) => mergeMessageProjectionMessage(message, existing));
|
||||||
|
return mergeBoundedProjectionMessages(existing, incoming);
|
||||||
}
|
}
|
||||||
|
|
||||||
function mergeMessageProjectionMessage(message: ChatMessage, existing: ChatMessage[]): ChatMessage {
|
function mergeMessageProjectionMessage(message: ChatMessage, existing: ChatMessage[]): ChatMessage {
|
||||||
@@ -999,6 +1012,11 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
|||||||
if (ownerBefore === activeSessionId.value) recordActivity(reason);
|
if (ownerBefore === activeSessionId.value) recordActivity(reason);
|
||||||
await refreshTurnStatusByTraceId(id);
|
await refreshTurnStatusByTraceId(id);
|
||||||
const ownerSessionId = traceOwnerSessionId(id, turnStatusAuthority.value[id]?.sessionId ?? null) ?? ownerBefore;
|
const ownerSessionId = traceOwnerSessionId(id, turnStatusAuthority.value[id]?.sessionId ?? null) ?? ownerBefore;
|
||||||
|
const ownerMessagesBeforeProjection = ownerSessionId ? serverState.value.messagesBySessionId[ownerSessionId] ?? [] : messages.value;
|
||||||
|
if (traceProjectionIsTerminalSealed(id, ownerMessagesBeforeProjection)) {
|
||||||
|
clearActiveTraceRestGapFill(id);
|
||||||
|
return;
|
||||||
|
}
|
||||||
if (ownerSessionId) {
|
if (ownerSessionId) {
|
||||||
scheduleSessionListRefresh(ownerSessionId, runtimePolicy.sessionListRealtimeRefreshDelayMs);
|
scheduleSessionListRefresh(ownerSessionId, runtimePolicy.sessionListRealtimeRefreshDelayMs);
|
||||||
await refreshMessageProjectionForTrace(ownerSessionId, id);
|
await refreshMessageProjectionForTrace(ownerSessionId, id);
|
||||||
@@ -1288,10 +1306,10 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
|||||||
rememberTurnStatus(traceId, result);
|
rememberTurnStatus(traceId, result);
|
||||||
markWorkbenchTraceProjected(traceId);
|
markWorkbenchTraceProjected(traceId);
|
||||||
clearActiveTraceRestGapFill(traceId);
|
clearActiveTraceRestGapFill(traceId);
|
||||||
void refreshMessageProjectionForTrace(ownerSessionId, traceId, { force: options.forceRead });
|
|
||||||
const ownerMessages = serverState.value.messagesBySessionId[ownerSessionId] ?? [];
|
const ownerMessages = serverState.value.messagesBySessionId[ownerSessionId] ?? [];
|
||||||
const terminalMessage = latestMessageForTrace(traceId, ownerMessages);
|
const terminalMessage = latestMessageForTrace(traceId, ownerMessages);
|
||||||
if (options.forceRead && terminalMessage) void hydrateTraceEventsForMessage(terminalMessage, { force: true });
|
if (!traceProjectionIsTerminalSealed(traceId, ownerMessages)) void refreshMessageProjectionForTrace(ownerSessionId, traceId, { force: options.forceRead });
|
||||||
|
if (options.forceRead && terminalMessage && !traceProjectionIsTerminalSealed(traceId, ownerMessages)) void hydrateTraceEventsForMessage(terminalMessage, { force: true });
|
||||||
if (ownerSessionId === activeSessionId.value) {
|
if (ownerSessionId === activeSessionId.value) {
|
||||||
chatPending.value = false;
|
chatPending.value = false;
|
||||||
currentRequest.value = null;
|
currentRequest.value = null;
|
||||||
|
|||||||
+2
@@ -30,6 +30,8 @@ declare global {
|
|||||||
workbenchTraceEventsMinRefreshMs?: number;
|
workbenchTraceEventsMinRefreshMs?: number;
|
||||||
workbenchSessionDetailMinRefreshMs?: number;
|
workbenchSessionDetailMinRefreshMs?: number;
|
||||||
workbenchSessionMessagesMinRefreshMs?: number;
|
workbenchSessionMessagesMinRefreshMs?: number;
|
||||||
|
workbenchSessionMessagesWindowLimit?: number;
|
||||||
|
workbenchTraceMessagesWindowLimit?: number;
|
||||||
workbenchRealtimeSessionMessagesMinRefreshMs?: number;
|
workbenchRealtimeSessionMessagesMinRefreshMs?: number;
|
||||||
workbenchTraceEventsTimeoutMs?: number;
|
workbenchTraceEventsTimeoutMs?: number;
|
||||||
sessionListRealtimeRefreshDelayMs?: number;
|
sessionListRealtimeRefreshDelayMs?: number;
|
||||||
|
|||||||
Reference in New Issue
Block a user