fix: coalesce Workbench realtime session refresh (#1902)
This commit is contained in:
@@ -23,6 +23,8 @@ const ACTIVE_TURN_GAP_INITIAL_DELAY_MS = 1_000;
|
||||
const ACTIVE_TURN_GAP_MAX_DELAY_MS = 5_000;
|
||||
const SESSION_LIST_PAGE_LIMIT = 20;
|
||||
const WORKBENCH_READ_HYDRATION_CONCURRENCY = 3;
|
||||
const SESSION_LIST_REALTIME_REFRESH_DELAY_MS = 1_200;
|
||||
const SESSION_LIST_TERMINAL_REFRESH_DELAY_MS = 400;
|
||||
|
||||
interface HydrateOptions {
|
||||
sessionId?: string | null;
|
||||
@@ -70,6 +72,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
let realtimeGapTimer: number | null = null;
|
||||
let activeTurnGapTimer: number | null = null;
|
||||
const sessionListRefreshInFlight = new Map<string, Promise<void>>();
|
||||
const sessionListRefreshTimers = new Map<string, number>();
|
||||
let workbenchReadHydrationActive = 0;
|
||||
const workbenchReadHydrationQueue: Array<() => void> = [];
|
||||
|
||||
@@ -267,7 +270,7 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
async function refreshSessions(includeSessionId: string | null = activeSessionId.value): Promise<void> {
|
||||
const requestIncludeSessionId = firstNonEmptyString(includeSessionId);
|
||||
const requestLimit = currentSessionListLimit();
|
||||
const requestKey = `${requestIncludeSessionId ?? ""}|${requestLimit}`;
|
||||
const requestKey = sessionListRefreshKey(requestIncludeSessionId, requestLimit);
|
||||
const existing = sessionListRefreshInFlight.get(requestKey);
|
||||
if (existing) return existing;
|
||||
const refresh = refreshSessionsNow(requestIncludeSessionId, requestLimit).finally(() => {
|
||||
@@ -294,6 +297,27 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
if (sessions.value.length === 0) error.value = response.error ?? "session list unavailable";
|
||||
}
|
||||
|
||||
function scheduleSessionListRefresh(includeSessionId: string | null | undefined = activeSessionId.value, delayMs = SESSION_LIST_REALTIME_REFRESH_DELAY_MS): void {
|
||||
const requestIncludeSessionId = firstNonEmptyString(includeSessionId);
|
||||
const requestLimit = currentSessionListLimit();
|
||||
const requestKey = sessionListRefreshKey(requestIncludeSessionId, requestLimit);
|
||||
if (typeof window === "undefined") {
|
||||
void refreshSessions(requestIncludeSessionId);
|
||||
return;
|
||||
}
|
||||
const existing = sessionListRefreshTimers.get(requestKey);
|
||||
if (existing) window.clearTimeout(existing);
|
||||
const timer = window.setTimeout(() => {
|
||||
sessionListRefreshTimers.delete(requestKey);
|
||||
void refreshSessions(requestIncludeSessionId);
|
||||
}, Math.max(0, Math.trunc(delayMs)));
|
||||
sessionListRefreshTimers.set(requestKey, timer);
|
||||
}
|
||||
|
||||
function sessionListRefreshKey(includeSessionId: string | null | undefined, limit: number): string {
|
||||
return `${firstNonEmptyString(includeSessionId) ?? ""}|${limit}`;
|
||||
}
|
||||
|
||||
async function loadMoreSessions(): Promise<void> {
|
||||
if (sessionListLoadingMore.value || !sessionListHasMore.value) return;
|
||||
const cursor = sessionListNextCursor.value;
|
||||
@@ -624,7 +648,7 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
||||
status: response.data.status ?? "running"
|
||||
};
|
||||
applyTurnStatusSnapshot(canonicalTraceId, response.data);
|
||||
void refreshSessions(sessionId);
|
||||
scheduleSessionListRefresh(sessionId, SESSION_LIST_TERMINAL_REFRESH_DELAY_MS);
|
||||
if ((response.data as AgentChatResultResponse).terminal === true || isTerminalMessageStatus(response.data.status)) {
|
||||
completeTrace(canonicalTraceId, response.data as AgentChatResultResponse);
|
||||
return true;
|
||||
@@ -653,7 +677,7 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
||||
}
|
||||
if (message.status === "running") chatPending.value = false;
|
||||
currentRequest.value = null;
|
||||
void refreshSessions(message.sessionId ?? selectedSessionId.value);
|
||||
scheduleSessionListRefresh(message.sessionId ?? selectedSessionId.value, SESSION_LIST_TERMINAL_REFRESH_DELAY_MS);
|
||||
}
|
||||
|
||||
async function cancelRunningTrace(): Promise<void> {
|
||||
@@ -1038,8 +1062,8 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
||||
|
||||
async function hydrateRealtimeGap(reason: string): Promise<void> {
|
||||
recordActivity(`realtime-gap:${reason}`);
|
||||
scheduleSessionListRefresh(activeSessionId.value, SESSION_LIST_REALTIME_REFRESH_DELAY_MS);
|
||||
await Promise.all([
|
||||
refreshSessions(),
|
||||
hydrateTurnStatusAuthority(messages.value),
|
||||
hydrateTraceEvents(messages.value)
|
||||
]);
|
||||
@@ -1114,7 +1138,7 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
||||
return { ...message, ...messageTimingPatchForMerge(message, trace), runnerTrace, error: clearCompletedDiagnostics ? null : error ?? message.error ?? null, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, updatedAt: new Date().toISOString() };
|
||||
}));
|
||||
markWorkbenchTraceProjected(traceId);
|
||||
void refreshSessions(ownerSessionId);
|
||||
scheduleSessionListRefresh(ownerSessionId, SESSION_LIST_REALTIME_REFRESH_DELAY_MS);
|
||||
}
|
||||
|
||||
function completeTrace(traceId: string, result: AgentChatResultResponse): void {
|
||||
@@ -1148,7 +1172,7 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
||||
void clearActiveTrace(traceId, "trace-terminal");
|
||||
restartRealtime("trace-terminal");
|
||||
}
|
||||
void refreshSessions(ownerSessionId);
|
||||
scheduleSessionListRefresh(ownerSessionId, SESSION_LIST_TERMINAL_REFRESH_DELAY_MS);
|
||||
}
|
||||
|
||||
async function hydrateTerminalMessageDiagnostics(): Promise<void> {
|
||||
@@ -1183,7 +1207,7 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
||||
return { ...message, ...messageTimingPatchForMerge(message, result), title: normalizeWorkbenchMessageTitle(message.role, message.title), runnerTrace, error, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, blocker: projection?.blocker ?? null, agentRun: agentRun ?? undefined, updatedAt: new Date().toISOString() };
|
||||
}));
|
||||
void hydrateTraceEvents(serverState.value.messagesBySessionId[ownerSessionId] ?? []);
|
||||
void refreshSessions(ownerSessionId);
|
||||
scheduleSessionListRefresh(ownerSessionId, SESSION_LIST_TERMINAL_REFRESH_DELAY_MS);
|
||||
}
|
||||
|
||||
function applyRealtimeProjectionError(event: WorkbenchRealtimeEvent): void {
|
||||
@@ -1225,7 +1249,7 @@ function nonBlockingProjection(projection: ProjectionDiagnostic | null): Project
|
||||
chatPending.value = false;
|
||||
currentRequest.value = null;
|
||||
void clearActiveTrace(traceId, "trace-infrastructure-error");
|
||||
void refreshSessions(selectedSessionId.value);
|
||||
scheduleSessionListRefresh(selectedSessionId.value, SESSION_LIST_TERMINAL_REFRESH_DELAY_MS);
|
||||
}
|
||||
|
||||
function projectOptimisticRunningTurn(input: { sessionId: string; threadId: string | null; traceId: string; userText: string }): void {
|
||||
@@ -1365,12 +1389,14 @@ function sessionFromWorkbenchSession(value: unknown): WorkbenchSessionRecord | n
|
||||
async function loadWorkbenchSession(sessionId: string, seed: WorkbenchSessionRecord | null = null): Promise<WorkbenchSessionRecord | null> {
|
||||
const requestId = normalizeWorkbenchSessionRouteId(sessionId);
|
||||
if (!requestId) return null;
|
||||
const normalizedRequestId = normalizeWorkbenchSessionId(requestId);
|
||||
const eagerMessages = normalizedRequestId ? api.workbench.sessionMessages(normalizedRequestId, { limit: 100 }) : null;
|
||||
const detail = await api.workbench.session(requestId);
|
||||
if (!detail.ok) return null;
|
||||
const detailSession = sessionFromWorkbenchSession(detail.data?.session);
|
||||
const id = detailSession?.sessionId ?? normalizeWorkbenchSessionId(requestId) ?? seed?.sessionId;
|
||||
const id = detailSession?.sessionId ?? normalizedRequestId ?? seed?.sessionId;
|
||||
if (!id) return null;
|
||||
const messages = await api.workbench.sessionMessages(id, { limit: 100 });
|
||||
const messages = eagerMessages && id === normalizedRequestId ? await eagerMessages : await api.workbench.sessionMessages(id, { limit: 100 });
|
||||
const base = detailSession ?? seed;
|
||||
if (!base) return null;
|
||||
const page = messages.ok ? messages.data : null;
|
||||
|
||||
Reference in New Issue
Block a user