Merge pull request #2357 from pikasTech/fix/2356-workbench-sse-flush-budget
fix(workbench): bound SSE flush and sealed projection work
This commit is contained in:
@@ -99,6 +99,49 @@ test("Workbench runtime diagnostics expose root cause labels without raw identif
|
||||
}
|
||||
});
|
||||
|
||||
test("Workbench SSE flush diagnostics expose bounded queue cost fields", () => {
|
||||
resetWorkbenchPerformanceForTest();
|
||||
const restoreBrowserRuntime = installBrowserProbeRuntime();
|
||||
try {
|
||||
recordWorkbenchRuntimeDiagnostic({
|
||||
module: "workbench-events",
|
||||
sessionId: "ses_secret",
|
||||
traceId: "trc_secret",
|
||||
outcome: "ok",
|
||||
diagnostic: {
|
||||
code: "workbench_sse_flush",
|
||||
reason: "yield",
|
||||
eventCount: 5,
|
||||
deliveredCount: 4,
|
||||
dropped: 1,
|
||||
chunkCount: 2,
|
||||
replacedByKey: 1,
|
||||
flushDurationMs: 17,
|
||||
maxItemsPerChunk: 4,
|
||||
maxChunkMs: 8
|
||||
}
|
||||
});
|
||||
|
||||
const events = drainWorkbenchPerformanceEventsForTest();
|
||||
const event = events.find((item) => item.kind === "workbench_ui_event" && item.eventType === "runtime_diagnostic");
|
||||
assert.equal(event?.state, "sample");
|
||||
assert.equal(event?.reason, "sse_flush");
|
||||
assert.equal(event?.valueMs, 17);
|
||||
assert.equal(event?.eventCount, 5);
|
||||
assert.equal(event?.deliveredCount, 4);
|
||||
assert.equal(event?.droppedCount, 1);
|
||||
assert.equal(event?.chunkCount, 2);
|
||||
assert.equal(event?.replacedByKey, 1);
|
||||
assert.equal(event?.flushDurationMs, 17);
|
||||
assert.equal(event?.maxItemsPerChunk, 4);
|
||||
assert.equal(event?.maxChunkMs, 8);
|
||||
assert.doesNotMatch(JSON.stringify(events), /ses_secret|trc_secret/u);
|
||||
} finally {
|
||||
resetWorkbenchPerformanceForTest();
|
||||
restoreBrowserRuntime();
|
||||
}
|
||||
});
|
||||
|
||||
test("Workbench API timing enriches from ResourceTiming and clears the browser buffer", () => {
|
||||
resetWorkbenchPerformanceForTest();
|
||||
let clearCount = 0;
|
||||
|
||||
@@ -37,6 +37,9 @@ test("Workbench runtime policy reads injected config while preserving defaults",
|
||||
workbenchSessionMessagesWindowLimit: 9,
|
||||
workbenchTraceMessagesWindowLimit: 4,
|
||||
workbenchRealtimeErrorGapFillMinMs: 0,
|
||||
workbenchRealtimeFlushMaxItemsPerChunk: 2,
|
||||
workbenchRealtimeFlushMaxChunkMs: 6,
|
||||
workbenchRealtimeFlushYieldMs: 5,
|
||||
defaultGatewayTimeoutMs: "bad"
|
||||
});
|
||||
|
||||
@@ -46,6 +49,9 @@ test("Workbench runtime policy reads injected config while preserving defaults",
|
||||
assert.equal(policy.workbenchSessionMessagesWindowLimit, 9);
|
||||
assert.equal(policy.workbenchTraceMessagesWindowLimit, 4);
|
||||
assert.equal(policy.workbenchRealtimeErrorGapFillMinMs, 0);
|
||||
assert.equal(policy.workbenchRealtimeFlushMaxItemsPerChunk, 2);
|
||||
assert.equal(policy.workbenchRealtimeFlushMaxChunkMs, 6);
|
||||
assert.equal(policy.workbenchRealtimeFlushYieldMs, 5);
|
||||
assert.equal(policy.defaultGatewayTimeoutMs, 120_000);
|
||||
});
|
||||
|
||||
@@ -127,6 +133,50 @@ test("coalesced event queue replaces keyed snapshots and preserves unkeyed event
|
||||
assert.deepEqual(flushed, [["new", "event-1", "event-2"]]);
|
||||
});
|
||||
|
||||
test("coalesced event queue chunks scheduled flushes and yields between chunks", () => {
|
||||
const flushed: string[] = [];
|
||||
const drain: unknown[] = [];
|
||||
const scheduled: (() => void)[] = [];
|
||||
const yielded: (() => void)[] = [];
|
||||
let clock = 0;
|
||||
const queue = createCoalescedEventQueue<{ key?: string; value: string }>({
|
||||
keyOf: (item) => item.key,
|
||||
maxItemsPerChunk: 10,
|
||||
maxChunkMs: 4,
|
||||
now: () => clock,
|
||||
schedule: (flush) => {
|
||||
scheduled.push(flush);
|
||||
return () => undefined;
|
||||
},
|
||||
yieldSchedule: (flush) => {
|
||||
yielded.push(flush);
|
||||
return () => undefined;
|
||||
},
|
||||
onFlush: (items) => {
|
||||
clock += 3;
|
||||
flushed.push(...items.map((item) => item.value));
|
||||
},
|
||||
onDrain: (info) => drain.push(info)
|
||||
});
|
||||
|
||||
queue.push({ key: "snapshot:trc", value: "old" });
|
||||
queue.push({ value: "event-1" });
|
||||
queue.push({ key: "snapshot:trc", value: "new" });
|
||||
queue.push({ value: "event-2" });
|
||||
queue.push({ value: "event-3" });
|
||||
|
||||
assert.equal(queue.size, 4);
|
||||
assert.equal(scheduled.length, 1);
|
||||
scheduled.shift()?.();
|
||||
assert.deepEqual(flushed, ["new", "event-1"]);
|
||||
assert.equal(queue.size, 2);
|
||||
assert.equal(yielded.length, 1);
|
||||
yielded.shift()?.();
|
||||
assert.deepEqual(flushed, ["new", "event-1", "event-2", "event-3"]);
|
||||
assert.equal(queue.size, 0);
|
||||
assert.deepEqual(drain, [{ eventCount: 5, deliveredCount: 4, droppedCount: 1, chunkCount: 2, replacedByKey: 1, flushDurationMs: 12, maxItemsPerChunk: 10, maxChunkMs: 4, reason: "yield" }]);
|
||||
});
|
||||
|
||||
test("keyed single-flight coalesces duplicate work and supports replacement", async () => {
|
||||
const flight = createKeyedSingleflight<number>();
|
||||
let calls = 0;
|
||||
@@ -366,6 +416,9 @@ test("realtime recovery planner gates refresh steps by transport actions and aut
|
||||
|
||||
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"]);
|
||||
|
||||
const terminalSealed = planWorkbenchRealtimeRecovery(recovery, { selectedSessionId: "ses_1", activeSessionId: "ses_1", fallbackTraceId: "trc_1", activeTraceAuthorized: true, terminalTraceSealed: true });
|
||||
assert.deepEqual(terminalSealed.steps.map((step) => step.type), []);
|
||||
});
|
||||
|
||||
test("health probe cache records ok and unavailable states", async () => {
|
||||
|
||||
@@ -4,7 +4,7 @@
|
||||
import type { ChatMessage, ProjectionDiagnostic, TraceEvent } from "@/types";
|
||||
import { createCoalescedEventQueue } from "@/utils/scheduler/coalesced-event-queue";
|
||||
import { composeWorkbenchScopedKey, firstScopePart } from "@/utils/workbench-key";
|
||||
import { recordWorkbenchSseLifecycle } from "@/utils/workbench-performance";
|
||||
import { recordWorkbenchRuntimeDiagnostic, recordWorkbenchSseLifecycle } from "@/utils/workbench-performance";
|
||||
|
||||
export interface WorkbenchRealtimeTraceSnapshot {
|
||||
traceId?: string | null;
|
||||
@@ -45,6 +45,9 @@ export interface WorkbenchEventStreamOptions {
|
||||
sessionId?: string | null;
|
||||
traceId?: string | null;
|
||||
afterSeq?: number | null;
|
||||
flushMaxItemsPerChunk?: number | null;
|
||||
flushMaxChunkMs?: number | null;
|
||||
flushYieldMs?: number | null;
|
||||
onEvent: (event: WorkbenchRealtimeEvent, eventName: string) => void;
|
||||
onOpen?: () => void;
|
||||
onError?: (event: Event) => void;
|
||||
@@ -82,8 +85,34 @@ export function connectWorkbenchEvents(options: WorkbenchEventStreamOptions): Wo
|
||||
const source = new EventSource(eventRoute, { withCredentials: true });
|
||||
const queue = createCoalescedEventQueue<QueuedRealtimeEvent>({
|
||||
keyOf: (item) => realtimeCoalesceKey(item.payload, item.eventName),
|
||||
maxItemsPerChunk: options.flushMaxItemsPerChunk,
|
||||
maxChunkMs: options.flushMaxChunkMs,
|
||||
yieldSchedule: (flush) => scheduleRealtimeFlushYield(flush, options.flushYieldMs),
|
||||
onFlush: (items) => {
|
||||
for (const item of items) options.onEvent(item.payload, item.eventName);
|
||||
},
|
||||
onDrain: (info) => {
|
||||
if (info.eventCount === 0 && info.droppedCount === 0) return;
|
||||
recordWorkbenchRuntimeDiagnostic({
|
||||
module: "workbench-events",
|
||||
sessionId: options.sessionId,
|
||||
traceId: options.traceId,
|
||||
outcome: "ok",
|
||||
diagnostic: {
|
||||
code: "workbench_sse_flush",
|
||||
reason: info.reason,
|
||||
source: "eventsource-coalesced-queue",
|
||||
eventCount: info.eventCount,
|
||||
deliveredCount: info.deliveredCount,
|
||||
dropped: info.droppedCount,
|
||||
chunkCount: info.chunkCount,
|
||||
replacedByKey: info.replacedByKey,
|
||||
flushDurationMs: Math.round(info.flushDurationMs),
|
||||
maxItemsPerChunk: info.maxItemsPerChunk,
|
||||
maxChunkMs: Math.round(info.maxChunkMs),
|
||||
valuesRedacted: true
|
||||
}
|
||||
});
|
||||
}
|
||||
});
|
||||
source.onopen = () => {
|
||||
@@ -109,7 +138,7 @@ export function connectWorkbenchEvents(options: WorkbenchEventStreamOptions): Wo
|
||||
};
|
||||
return {
|
||||
close() {
|
||||
queue.flush();
|
||||
queue.clear("stream-close");
|
||||
for (const { name, listener } of listeners) source.removeEventListener(name, listener);
|
||||
recordWorkbenchSseLifecycle({ state: "close", route: eventRoute, sessionId: options.sessionId, traceId: options.traceId });
|
||||
source.close();
|
||||
@@ -137,6 +166,16 @@ function parseRealtimeEvent(raw: string): WorkbenchRealtimeEvent | null {
|
||||
}
|
||||
}
|
||||
|
||||
function scheduleRealtimeFlushYield(flush: () => void, yieldMs: number | null | undefined): () => void {
|
||||
if (typeof window !== "undefined" && typeof window.requestAnimationFrame === "function" && (!Number.isFinite(Number(yieldMs)) || Number(yieldMs) <= 0)) {
|
||||
const id = window.requestAnimationFrame(() => flush());
|
||||
return () => window.cancelAnimationFrame(id);
|
||||
}
|
||||
const delay = Math.max(0, Math.trunc(Number(yieldMs ?? 0)));
|
||||
const id = setTimeout(() => flush(), delay);
|
||||
return () => clearTimeout(id);
|
||||
}
|
||||
|
||||
function realtimeCoalesceKey(event: WorkbenchRealtimeEvent, eventName: string): string | null {
|
||||
const sessionId = firstScopePart(event.sessionId, event.message?.sessionId, event.snapshot?.sessionId, event.event?.sessionId);
|
||||
const traceId = firstScopePart(event.traceId, event.message?.traceId, event.snapshot?.traceId, event.event?.traceId);
|
||||
|
||||
@@ -26,6 +26,9 @@ export interface WorkbenchRuntimePolicy {
|
||||
sessionListTerminalRefreshDelayMs: number;
|
||||
sessionListMinRefreshIntervalMs: number;
|
||||
workbenchRealtimeErrorGapFillMinMs: number;
|
||||
workbenchRealtimeFlushMaxItemsPerChunk: number;
|
||||
workbenchRealtimeFlushMaxChunkMs: number;
|
||||
workbenchRealtimeFlushYieldMs: number;
|
||||
workbenchActiveTraceRestGapFillInitialMs: number;
|
||||
workbenchActiveTraceRestGapFillRepeatMs: number;
|
||||
}
|
||||
@@ -55,6 +58,9 @@ const DEFAULT_WORKBENCH_RUNTIME_POLICY: WorkbenchRuntimePolicy = Object.freeze({
|
||||
sessionListTerminalRefreshDelayMs: 1_500,
|
||||
sessionListMinRefreshIntervalMs: 15_000,
|
||||
workbenchRealtimeErrorGapFillMinMs: 2_000,
|
||||
workbenchRealtimeFlushMaxItemsPerChunk: 4,
|
||||
workbenchRealtimeFlushMaxChunkMs: 8,
|
||||
workbenchRealtimeFlushYieldMs: 0,
|
||||
workbenchActiveTraceRestGapFillInitialMs: 2_500,
|
||||
workbenchActiveTraceRestGapFillRepeatMs: 5_000
|
||||
});
|
||||
@@ -86,6 +92,9 @@ export function workbenchRuntimePolicy(input: unknown = runtimePolicyConfig()):
|
||||
sessionListTerminalRefreshDelayMs: nonNegativeNumber(source.sessionListTerminalRefreshDelayMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.sessionListTerminalRefreshDelayMs),
|
||||
sessionListMinRefreshIntervalMs: nonNegativeNumber(source.sessionListMinRefreshIntervalMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.sessionListMinRefreshIntervalMs),
|
||||
workbenchRealtimeErrorGapFillMinMs: nonNegativeNumber(source.workbenchRealtimeErrorGapFillMinMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeErrorGapFillMinMs),
|
||||
workbenchRealtimeFlushMaxItemsPerChunk: positiveInteger(source.workbenchRealtimeFlushMaxItemsPerChunk, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeFlushMaxItemsPerChunk),
|
||||
workbenchRealtimeFlushMaxChunkMs: nonNegativeNumber(source.workbenchRealtimeFlushMaxChunkMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeFlushMaxChunkMs),
|
||||
workbenchRealtimeFlushYieldMs: nonNegativeNumber(source.workbenchRealtimeFlushYieldMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchRealtimeFlushYieldMs),
|
||||
workbenchActiveTraceRestGapFillInitialMs: nonNegativeNumber(source.workbenchActiveTraceRestGapFillInitialMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchActiveTraceRestGapFillInitialMs),
|
||||
workbenchActiveTraceRestGapFillRepeatMs: nonNegativeNumber(source.workbenchActiveTraceRestGapFillRepeatMs, DEFAULT_WORKBENCH_RUNTIME_POLICY.workbenchActiveTraceRestGapFillRepeatMs)
|
||||
};
|
||||
|
||||
@@ -32,6 +32,7 @@ export interface WorkbenchRealtimeRecoveryContext {
|
||||
activeSessionId?: string | null;
|
||||
fallbackTraceId?: string | null;
|
||||
activeTraceAuthorized?: boolean;
|
||||
terminalTraceSealed?: boolean;
|
||||
}
|
||||
|
||||
export type WorkbenchRealtimeRecoveryAuthority = "automatic-recovery";
|
||||
@@ -76,9 +77,9 @@ export function planWorkbenchRealtimeRecovery(recovery: WorkbenchStreamTransport
|
||||
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(automaticRecoveryStep({ type: "refresh-session-messages", sessionId, reason: "realtime-error:messages" }));
|
||||
if (actions.has("schedule-session-list") && sessionId) steps.push(automaticRecoveryStep({ type: "schedule-session-list", sessionId, reason: "realtime-error:session-list" }));
|
||||
if (traceId && context.activeTraceAuthorized === true) {
|
||||
if (context.terminalTraceSealed !== true && actions.has("refresh-session-messages") && sessionId && sessionId === context.activeSessionId) steps.push(automaticRecoveryStep({ type: "refresh-session-messages", sessionId, reason: "realtime-error:messages" }));
|
||||
if (context.terminalTraceSealed !== true && actions.has("schedule-session-list") && sessionId) steps.push(automaticRecoveryStep({ type: "schedule-session-list", sessionId, reason: "realtime-error:session-list" }));
|
||||
if (context.terminalTraceSealed !== true && traceId && context.activeTraceAuthorized === true) {
|
||||
if (actions.has("refresh-turn-status")) steps.push(automaticRecoveryStep({ type: "refresh-turn-status", traceId, reason: "realtime-error:turn-status" }));
|
||||
if (actions.has("hydrate-trace-events")) steps.push(automaticRecoveryStep({ type: "hydrate-trace-events", traceId, reason: "realtime-error:trace-events" }));
|
||||
}
|
||||
|
||||
@@ -1001,6 +1001,9 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
sessionId,
|
||||
traceId,
|
||||
errorRecoveryMinMs: runtimePolicy.workbenchRealtimeErrorGapFillMinMs,
|
||||
flushMaxItemsPerChunk: runtimePolicy.workbenchRealtimeFlushMaxItemsPerChunk,
|
||||
flushMaxChunkMs: runtimePolicy.workbenchRealtimeFlushMaxChunkMs,
|
||||
flushYieldMs: runtimePolicy.workbenchRealtimeFlushYieldMs,
|
||||
onOpen: () => undefined,
|
||||
onRecovery: (recovery) => handleRealtimeRecovery(recovery),
|
||||
onEvent: (event, eventName) => applyRealtimeEvent(event, eventName)
|
||||
@@ -1014,7 +1017,8 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
selectedSessionId: selectedSessionId.value,
|
||||
activeSessionId: activeSessionId.value,
|
||||
fallbackTraceId: traceId,
|
||||
activeTraceAuthorized: Boolean(traceId && shouldApplyActiveTraceAuthority(traceId, sessionId))
|
||||
activeTraceAuthorized: Boolean(traceId && shouldApplyActiveTraceAuthority(traceId, sessionId)),
|
||||
terminalTraceSealed: Boolean(traceId && traceTerminalBodyIsVisible(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);
|
||||
@@ -1209,7 +1213,12 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
function applyRealtimeTraceSnapshot(traceId: string | null | undefined, snapshot: WorkbenchRealtimeEvent["snapshot"]): void {
|
||||
const id = firstNonEmptyString(traceId, snapshot?.traceId);
|
||||
if (!id || !snapshot) return;
|
||||
if (!shouldApplyActiveTraceAuthority(id, traceResultSessionId(snapshot))) return;
|
||||
const sessionId = traceResultSessionId(snapshot);
|
||||
if (!shouldApplyActiveTraceAuthority(id, sessionId)) return;
|
||||
if (traceTerminalBodyIsVisible(id, sessionId)) {
|
||||
recordWorkbenchRuntimeDiagnostic({ module: "workbench-terminal-priority", sessionId, traceId: id, outcome: "ok", diagnostic: { code: "terminal_low_priority_sse_trace_skip", source: "realtime-trace-snapshot", valuesRedacted: true } });
|
||||
return;
|
||||
}
|
||||
applyTraceSnapshot(id, realtimeSnapshotToTraceSnapshot(id, snapshot));
|
||||
}
|
||||
|
||||
@@ -1218,6 +1227,10 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
if (!id) return;
|
||||
const sessionId = realtimeEvent ? realtimeEventSessionId(realtimeEvent) : traceResultSessionId(snapshot ?? event ?? null);
|
||||
if (!shouldApplyActiveTraceAuthority(id, sessionId)) return;
|
||||
if (traceTerminalBodyIsVisible(id, sessionId)) {
|
||||
recordWorkbenchRuntimeDiagnostic({ module: "workbench-terminal-priority", sessionId, traceId: id, outcome: "ok", diagnostic: { code: "terminal_low_priority_sse_trace_skip", source: "realtime-trace-event", valuesRedacted: true } });
|
||||
return;
|
||||
}
|
||||
const events = event ? [event] : Array.isArray(snapshot?.events) ? snapshot.events : [];
|
||||
markWorkbenchTraceEventsReceived({ traceId: id, events, transport: "sse", serverSentAt: realtimeEvent?.serverSentAt, eventCreatedAt: realtimeEvent?.eventCreatedAt, traceSeq: realtimeEvent?.traceSeq ?? realtimeEvent?.cursor?.traceSeq });
|
||||
applyTraceSnapshot(id, realtimeSnapshotToTraceSnapshot(id, snapshot ?? { traceId: id, status: event?.status, events }, events));
|
||||
|
||||
@@ -5,24 +5,69 @@ export type CoalescedQueueKey = string | null | undefined;
|
||||
export type CoalescedQueueCancel = () => void;
|
||||
export type CoalescedQueueScheduler = (flush: () => void) => CoalescedQueueCancel | void;
|
||||
|
||||
export interface CoalescedEventQueueChunkInfo {
|
||||
chunkIndex: number;
|
||||
eventCount: number;
|
||||
remainingCount: number;
|
||||
reason: string;
|
||||
}
|
||||
|
||||
export interface CoalescedEventQueueDrainInfo {
|
||||
eventCount: number;
|
||||
deliveredCount: number;
|
||||
droppedCount: number;
|
||||
chunkCount: number;
|
||||
replacedByKey: number;
|
||||
flushDurationMs: number;
|
||||
maxItemsPerChunk: number;
|
||||
maxChunkMs: number;
|
||||
reason: string;
|
||||
}
|
||||
|
||||
export interface CoalescedEventQueueOptions<T> {
|
||||
keyOf: (item: T) => CoalescedQueueKey;
|
||||
onFlush: (items: T[]) => void;
|
||||
onFlush: (items: T[], info: CoalescedEventQueueChunkInfo) => void;
|
||||
onDrain?: (info: CoalescedEventQueueDrainInfo) => void;
|
||||
schedule?: CoalescedQueueScheduler;
|
||||
yieldSchedule?: CoalescedQueueScheduler;
|
||||
maxItemsPerChunk?: number | null;
|
||||
maxChunkMs?: number | null;
|
||||
now?: () => number;
|
||||
}
|
||||
|
||||
interface CoalescedQueueCycle {
|
||||
startedAt: number;
|
||||
eventCount: number;
|
||||
deliveredCount: number;
|
||||
replacedByKey: number;
|
||||
chunkCount: number;
|
||||
reason: string;
|
||||
}
|
||||
|
||||
export class CoalescedEventQueue<T> {
|
||||
private readonly keyOf: (item: T) => CoalescedQueueKey;
|
||||
private readonly onFlush: (items: T[]) => void;
|
||||
private readonly onFlush: (items: T[], info: CoalescedEventQueueChunkInfo) => void;
|
||||
private readonly onDrain: ((info: CoalescedEventQueueDrainInfo) => void) | null;
|
||||
private readonly scheduleFlush: CoalescedQueueScheduler;
|
||||
private readonly scheduleYield: CoalescedQueueScheduler;
|
||||
private readonly maxItemsPerChunk: number;
|
||||
private readonly maxChunkMs: number;
|
||||
private readonly now: () => number;
|
||||
private items: T[] = [];
|
||||
private keyIndex = new Map<string, number>();
|
||||
private cancelScheduled: CoalescedQueueCancel | null = null;
|
||||
private scheduled = false;
|
||||
private cycle: CoalescedQueueCycle | null = null;
|
||||
|
||||
constructor(options: CoalescedEventQueueOptions<T>) {
|
||||
this.keyOf = options.keyOf;
|
||||
this.onFlush = options.onFlush;
|
||||
this.onDrain = options.onDrain ?? null;
|
||||
this.scheduleFlush = options.schedule ?? defaultScheduler;
|
||||
this.scheduleYield = options.yieldSchedule ?? defaultYieldScheduler;
|
||||
this.maxItemsPerChunk = positiveInteger(options.maxItemsPerChunk, 16);
|
||||
this.maxChunkMs = nonNegativeNumber(options.maxChunkMs, 8);
|
||||
this.now = options.now ?? monotonicNow;
|
||||
}
|
||||
|
||||
get size(): number {
|
||||
@@ -30,9 +75,12 @@ export class CoalescedEventQueue<T> {
|
||||
}
|
||||
|
||||
push(item: T): void {
|
||||
this.ensureCycle("scheduled");
|
||||
if (this.cycle) this.cycle.eventCount += 1;
|
||||
const key = normalizedKey(this.keyOf(item));
|
||||
if (key && this.keyIndex.has(key)) {
|
||||
this.items[this.keyIndex.get(key) ?? 0] = item;
|
||||
if (this.cycle) this.cycle.replacedByKey += 1;
|
||||
} else {
|
||||
if (key) this.keyIndex.set(key, this.items.length);
|
||||
this.items.push(item);
|
||||
@@ -41,32 +89,114 @@ export class CoalescedEventQueue<T> {
|
||||
}
|
||||
|
||||
flush(): T[] {
|
||||
if (this.cancelScheduled) {
|
||||
this.cancelScheduled();
|
||||
this.cancelScheduled = null;
|
||||
}
|
||||
this.cancelPendingSchedule();
|
||||
if (this.items.length === 0) return [];
|
||||
const items = this.items;
|
||||
this.items = [];
|
||||
this.keyIndex = new Map();
|
||||
this.onFlush(items);
|
||||
this.ensureCycle("manual");
|
||||
if (this.cycle) {
|
||||
this.cycle.deliveredCount += items.length;
|
||||
this.cycle.chunkCount += 1;
|
||||
}
|
||||
this.onFlush(items, { chunkIndex: this.cycle?.chunkCount ?? 1, eventCount: items.length, remainingCount: 0, reason: "manual" });
|
||||
this.finishCycle("manual", 0);
|
||||
return items;
|
||||
}
|
||||
|
||||
clear(): void {
|
||||
if (this.cancelScheduled) {
|
||||
this.cancelScheduled();
|
||||
this.cancelScheduled = null;
|
||||
}
|
||||
clear(reason = "clear"): void {
|
||||
this.cancelPendingSchedule();
|
||||
const droppedCount = this.items.length;
|
||||
this.items = [];
|
||||
this.keyIndex = new Map();
|
||||
this.finishCycle(reason, droppedCount);
|
||||
}
|
||||
|
||||
private ensureScheduled(): void {
|
||||
if (this.cancelScheduled) return;
|
||||
const cancel = this.scheduleFlush(() => this.flush());
|
||||
if (this.scheduled) return;
|
||||
this.scheduled = true;
|
||||
const cancel = this.scheduleFlush(() => this.flushScheduledChunk("scheduled"));
|
||||
this.cancelScheduled = typeof cancel === "function" ? cancel : null;
|
||||
}
|
||||
|
||||
private ensureYieldScheduled(): void {
|
||||
if (this.scheduled) return;
|
||||
this.scheduled = true;
|
||||
const cancel = this.scheduleYield(() => this.flushScheduledChunk("yield"));
|
||||
this.cancelScheduled = typeof cancel === "function" ? cancel : null;
|
||||
}
|
||||
|
||||
private flushScheduledChunk(reason: string): void {
|
||||
this.scheduled = false;
|
||||
this.cancelScheduled = null;
|
||||
if (this.items.length === 0) {
|
||||
this.finishCycle(reason, 0);
|
||||
return;
|
||||
}
|
||||
this.ensureCycle(reason);
|
||||
const chunkIndex = (this.cycle?.chunkCount ?? 0) + 1;
|
||||
if (this.cycle) this.cycle.chunkCount += 1;
|
||||
const chunkStartedAt = this.now();
|
||||
let deliveredInChunk = 0;
|
||||
while (this.items.length > 0 && deliveredInChunk < this.maxItemsPerChunk) {
|
||||
const item = this.takeNextItem();
|
||||
this.onFlush([item], { chunkIndex, eventCount: 1, remainingCount: this.items.length, reason });
|
||||
if (this.cycle) this.cycle.deliveredCount += 1;
|
||||
deliveredInChunk += 1;
|
||||
if (this.items.length > 0 && this.now() - chunkStartedAt >= this.maxChunkMs) break;
|
||||
}
|
||||
if (this.items.length > 0) {
|
||||
this.ensureYieldScheduled();
|
||||
return;
|
||||
}
|
||||
this.finishCycle(reason, 0);
|
||||
}
|
||||
|
||||
private takeNextItem(): T {
|
||||
const item = this.items[0];
|
||||
this.items = this.items.slice(1);
|
||||
this.rebuildKeyIndex();
|
||||
if (item === undefined) throw new Error("coalesced queue item missing");
|
||||
return item;
|
||||
}
|
||||
|
||||
private rebuildKeyIndex(): void {
|
||||
this.keyIndex = new Map();
|
||||
for (let index = 0; index < this.items.length; index += 1) {
|
||||
const item = this.items[index];
|
||||
const key = item === undefined ? null : normalizedKey(this.keyOf(item));
|
||||
if (key) this.keyIndex.set(key, index);
|
||||
}
|
||||
}
|
||||
|
||||
private ensureCycle(reason: string): void {
|
||||
if (this.cycle) return;
|
||||
this.cycle = { startedAt: this.now(), eventCount: 0, deliveredCount: 0, replacedByKey: 0, chunkCount: 0, reason };
|
||||
}
|
||||
|
||||
private finishCycle(reason: string, droppedCount: number): void {
|
||||
const cycle = this.cycle;
|
||||
if (!cycle) return;
|
||||
this.cycle = null;
|
||||
const dropped = Math.max(droppedCount, cycle.eventCount - cycle.deliveredCount);
|
||||
this.onDrain?.({
|
||||
eventCount: cycle.eventCount,
|
||||
deliveredCount: cycle.deliveredCount,
|
||||
droppedCount: dropped,
|
||||
chunkCount: cycle.chunkCount,
|
||||
replacedByKey: cycle.replacedByKey,
|
||||
flushDurationMs: Math.max(0, this.now() - cycle.startedAt),
|
||||
maxItemsPerChunk: this.maxItemsPerChunk,
|
||||
maxChunkMs: this.maxChunkMs,
|
||||
reason: reason || cycle.reason
|
||||
});
|
||||
}
|
||||
|
||||
private cancelPendingSchedule(): void {
|
||||
if (this.cancelScheduled) this.cancelScheduled();
|
||||
this.cancelScheduled = null;
|
||||
this.scheduled = false;
|
||||
}
|
||||
}
|
||||
|
||||
export function createCoalescedEventQueue<T>(options: CoalescedEventQueueOptions<T>): CoalescedEventQueue<T> {
|
||||
@@ -92,9 +222,34 @@ function defaultScheduler(flush: () => void): CoalescedQueueCancel {
|
||||
};
|
||||
}
|
||||
|
||||
function defaultYieldScheduler(flush: () => void): CoalescedQueueCancel {
|
||||
if (typeof requestAnimationFrame === "function" && typeof cancelAnimationFrame === "function") {
|
||||
const id = requestAnimationFrame(() => flush());
|
||||
return () => cancelAnimationFrame(id);
|
||||
}
|
||||
if (typeof setTimeout === "function" && typeof clearTimeout === "function") {
|
||||
const id = setTimeout(() => flush(), 0);
|
||||
return () => clearTimeout(id);
|
||||
}
|
||||
return defaultScheduler(flush);
|
||||
}
|
||||
|
||||
function normalizedKey(value: CoalescedQueueKey): string | null {
|
||||
if (typeof value !== "string") return null;
|
||||
const text = value.trim();
|
||||
return text ? text : null;
|
||||
}
|
||||
|
||||
function positiveInteger(value: number | null | undefined, fallback: number): number {
|
||||
const parsed = Number(value);
|
||||
return Number.isFinite(parsed) && parsed > 0 ? Math.max(1, Math.trunc(parsed)) : fallback;
|
||||
}
|
||||
|
||||
function nonNegativeNumber(value: number | null | undefined, fallback: number): number {
|
||||
const parsed = Number(value);
|
||||
return Number.isFinite(parsed) && parsed >= 0 ? parsed : fallback;
|
||||
}
|
||||
|
||||
function monotonicNow(): number {
|
||||
return typeof performance !== "undefined" && typeof performance.now === "function" ? performance.now() : Date.now();
|
||||
}
|
||||
|
||||
@@ -7,7 +7,7 @@ type WorkbenchEventKind = "workbench_journey" | "workbench_event_phase" | "workb
|
||||
type WorkbenchOutcome = "ok" | "timeout" | "error" | "dropped" | "stale" | "partial" | "empty" | "network" | "denied" | "unknown";
|
||||
type WorkbenchLoadingScope = "workbench" | "session_list" | "session_detail" | "api" | "sse" | "page";
|
||||
type WorkbenchUiState = "enter" | "exit" | "request" | "connect" | "open" | "message" | "close" | "error" | "sample";
|
||||
type WorkbenchUiReason = "hydrate" | "select_session" | "create_session" | "api_request" | "sse_connect" | "sse_open" | "sse_error" | "sse_close" | "pagehide" | "visibility_hidden" | "unknown";
|
||||
type WorkbenchUiReason = "hydrate" | "select_session" | "create_session" | "api_request" | "sse_connect" | "sse_open" | "sse_error" | "sse_close" | "sse_flush" | "pagehide" | "visibility_hidden" | "unknown";
|
||||
|
||||
interface WorkbenchPerformanceEvent {
|
||||
kind: WorkbenchEventKind;
|
||||
@@ -60,6 +60,14 @@ interface WorkbenchPerformanceEvent {
|
||||
recoveryAction?: string;
|
||||
transportState?: string;
|
||||
scopedKey?: string;
|
||||
eventCount?: number;
|
||||
deliveredCount?: number;
|
||||
droppedCount?: number;
|
||||
chunkCount?: number;
|
||||
replacedByKey?: number;
|
||||
flushDurationMs?: number;
|
||||
maxItemsPerChunk?: number;
|
||||
maxChunkMs?: number;
|
||||
}
|
||||
|
||||
interface TraceEventTimingInput {
|
||||
@@ -350,14 +358,16 @@ export function recordWorkbenchRuntimeDiagnostic(input: { module?: string | null
|
||||
const diagnostic = recordValue(input.diagnostic);
|
||||
const endedAt = wallNow();
|
||||
const code = diagnosticText(diagnostic, "code") ?? "runtime_diagnostic";
|
||||
const flushDurationMs = diagnosticNumber(diagnostic, "flushDurationMs");
|
||||
const isSseFlush = code === "workbench_sse_flush";
|
||||
enqueueWorkbenchUiEvent({
|
||||
eventType: "runtime_diagnostic",
|
||||
loadingScope: "sse",
|
||||
state: "error",
|
||||
reason: "sse_error",
|
||||
state: isSseFlush ? "sample" : "error",
|
||||
reason: isSseFlush ? "sse_flush" : "sse_error",
|
||||
route: pageRoute(),
|
||||
valueMs: 0,
|
||||
startedAtEpochMs: endedAt,
|
||||
valueMs: flushDurationMs ?? 0,
|
||||
startedAtEpochMs: Math.max(0, endedAt - (flushDurationMs ?? 0)),
|
||||
endedAtEpochMs: endedAt,
|
||||
sessionHash: hashIdentifier(input.sessionId ?? diagnosticText(diagnostic, "sessionId"), "ses"),
|
||||
traceHash: hashIdentifier(input.traceId ?? diagnosticText(diagnostic, "traceId"), "trc"),
|
||||
@@ -368,7 +378,15 @@ export function recordWorkbenchRuntimeDiagnostic(input: { module?: string | null
|
||||
rootCause: diagnosticText(diagnostic, "rootCause", "reason")?.slice(0, 160),
|
||||
recoveryAction: diagnosticText(diagnostic, "recoveryAction")?.slice(0, 160),
|
||||
transportState: diagnosticLabel(diagnosticText(diagnostic, "transportState")),
|
||||
scopedKey: hashIdentifier(diagnosticText(diagnostic, "scopedKey"), "scope")
|
||||
scopedKey: hashIdentifier(diagnosticText(diagnostic, "scopedKey"), "scope"),
|
||||
eventCount: diagnosticNumber(diagnostic, "eventCount"),
|
||||
deliveredCount: diagnosticNumber(diagnostic, "deliveredCount"),
|
||||
droppedCount: diagnosticNumber(diagnostic, "dropped", "droppedCount"),
|
||||
chunkCount: diagnosticNumber(diagnostic, "chunkCount"),
|
||||
replacedByKey: diagnosticNumber(diagnostic, "replacedByKey"),
|
||||
flushDurationMs,
|
||||
maxItemsPerChunk: diagnosticNumber(diagnostic, "maxItemsPerChunk"),
|
||||
maxChunkMs: diagnosticNumber(diagnostic, "maxChunkMs")
|
||||
});
|
||||
}
|
||||
|
||||
@@ -814,6 +832,16 @@ function diagnosticText(record: Record<string, unknown> | null, ...keys: string[
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function diagnosticNumber(record: Record<string, unknown> | null, ...keys: string[]): number | undefined {
|
||||
if (!record) return undefined;
|
||||
for (const key of keys) {
|
||||
const value = record[key];
|
||||
const number = Number(value);
|
||||
if (Number.isFinite(number) && number >= 0) return Math.round(number);
|
||||
}
|
||||
return undefined;
|
||||
}
|
||||
|
||||
function diagnosticLabel(value: unknown): string | undefined {
|
||||
const text = safeText(value).replace(/[^A-Za-z0-9_.:-]/gu, "_").slice(0, 80);
|
||||
return text || undefined;
|
||||
|
||||
@@ -17,6 +17,9 @@ export interface WorkbenchStreamTransportRestartInput {
|
||||
traceId?: string | null;
|
||||
afterSeq?: number | null;
|
||||
errorRecoveryMinMs?: number | null;
|
||||
flushMaxItemsPerChunk?: number | null;
|
||||
flushMaxChunkMs?: number | null;
|
||||
flushYieldMs?: number | null;
|
||||
onOpen?: () => void;
|
||||
onError?: (event: Event) => void;
|
||||
onState?: (state: WorkbenchStreamTransportState) => void;
|
||||
@@ -89,6 +92,9 @@ export class WorkbenchStreamTransportRuntime {
|
||||
sessionId: input.sessionId ?? null,
|
||||
traceId: input.traceId ?? null,
|
||||
afterSeq: input.afterSeq ?? this.cursorByKey.get(key)?.outboxSeq ?? null,
|
||||
flushMaxItemsPerChunk: input.flushMaxItemsPerChunk,
|
||||
flushMaxChunkMs: input.flushMaxChunkMs,
|
||||
flushYieldMs: input.flushYieldMs,
|
||||
onOpen: () => {
|
||||
this.markLive(key);
|
||||
this.emitState(input, "open", null);
|
||||
|
||||
Reference in New Issue
Block a user