Merge pull request #2365 from pikasTech/fix/2356-colada-trace-hydrate

fix: bound Workbench trace hydrate merge
This commit is contained in:
Lyon
2026-07-03 05:55:44 +08:00
committed by GitHub
5 changed files with 221 additions and 32 deletions
@@ -0,0 +1,49 @@
import assert from "node:assert/strict";
import { test } from "bun:test";
import type { ChatMessage, TraceEvent } from "@/types";
import { mergeRunnerTrace } from "./useTraceSubscription";
function event(projectedSeq: number, label = `event-${projectedSeq}`): TraceEvent {
return { projectedSeq, label, type: "event" };
}
function trace(input: NonNullable<ChatMessage["runnerTrace"]>): NonNullable<ChatMessage["runnerTrace"]> {
return input;
}
test("mergeRunnerTrace preserves event identity for replayed trace-api range", () => {
const previousEvents = [event(1), event(2)];
const previous = trace({ traceId: "trc_replay", eventSource: "trace-api", events: previousEvents, nextProjectedSeq: 2, range: { toProjectedSeq: 2 }, hasMore: true });
const next = trace({ traceId: "trc_replay", eventSource: "trace-api", events: [event(2)], range: { afterProjectedSeq: 1, fromProjectedSeq: 2, toProjectedSeq: 2, returned: 1 }, hasMore: true });
const merged = mergeRunnerTrace(previous, next);
assert.equal(merged.events, previousEvents);
assert.equal(merged.events?.length, 2);
});
test("mergeRunnerTrace appends only delta events after the canonical cursor", () => {
const previousEvents = [event(1), event(2)];
const delta = [event(3), event(4)];
const previous = trace({ traceId: "trc_delta", eventSource: "trace-api", events: previousEvents, nextProjectedSeq: 2, range: { toProjectedSeq: 2 }, hasMore: true });
const next = trace({ traceId: "trc_delta", eventSource: "trace-api", events: delta, range: { afterProjectedSeq: 2, fromProjectedSeq: 3, toProjectedSeq: 4, returned: 2 }, hasMore: false, fullTraceLoaded: true });
const merged = mergeRunnerTrace(previous, next);
assert.notEqual(merged.events, previousEvents);
assert.deepEqual(merged.events?.map((item) => item.projectedSeq), [1, 2, 3, 4]);
assert.equal(merged.events?.[0], previousEvents[0]);
assert.equal(merged.events?.[2], delta[0]);
});
test("mergeRunnerTrace can seal terminal metadata without requiring trace rows", () => {
const previousEvents = [event(1), event(2)];
const previous = trace({ traceId: "trc_terminal", eventSource: "trace-api", events: previousEvents, nextProjectedSeq: 2, range: { toProjectedSeq: 2 }, hasMore: true, status: "running" });
const terminal = trace({ traceId: "trc_terminal", eventSource: "turn-authority", events: [], status: "completed", traceStatus: "completed", fullTraceLoaded: false });
const merged = mergeRunnerTrace(previous, terminal);
assert.equal(merged.status, "completed");
assert.equal(merged.events, previousEvents);
});
@@ -109,14 +109,14 @@ export function mergeRunnerTrace(previous: ChatMessage["runnerTrace"], next: Non
if (!previous) return next;
const previousEvents = Array.isArray(previous.events) ? previous.events : [];
const nextEvents = Array.isArray(next.events) ? next.events : [];
const previousCursor = traceProjectedSeqCursor(previous, previousEvents);
const nextAuthoritative = next.eventSource === "trace-api";
const previousAuthoritative = previous.eventSource === "trace-api";
const keepPreviousEvents = previousEvents.length > 0 && !nextAuthoritative && (previousAuthoritative || next.eventsCompacted === true || nextEvents.length <= previousEvents.length);
const mergeAuthoritativeEvents = nextAuthoritative && previousEvents.length > 0 && (next.hasMore === true || next.fullTraceLoaded !== true || previous.fullTraceLoaded !== true || nextEvents.length < previousEvents.length);
const events = nextAuthoritative
? mergeAuthoritativeEvents ? mergeTraceEvents(previousEvents, nextEvents) : nextEvents
? previousEvents.length > 0 ? mergeTraceEvents(previousEvents, nextEvents, { previousCursor, nextRange: next.range }) : nextEvents
: keepPreviousEvents && previousAuthoritative && next.eventsCompacted === true ? previousEvents
: keepPreviousEvents ? mergeTraceEvents(previousEvents, nextEvents) : nextEvents.length > previousEvents.length ? mergeTraceEvents(previousEvents, nextEvents) : nextEvents;
: keepPreviousEvents ? mergeTraceEvents(previousEvents, nextEvents, { previousCursor, nextRange: next.range }) : nextEvents.length > previousEvents.length ? mergeTraceEvents(previousEvents, nextEvents, { previousCursor, nextRange: next.range }) : nextEvents;
const eventCount = keepPreviousEvents ? previous.eventCount ?? events.length : next.eventCount ?? previous.eventCount ?? events.length;
const timing = mergeTraceTimingProjection(previous, next);
return {
@@ -209,16 +209,79 @@ function traceTimingIsTerminal(value: unknown): boolean {
return ["completed", "failed", "cancelled", "canceled", "error"].includes(status) || traceTimingCandidate(value)?.finishedAt != null;
}
function mergeTraceEvents(previousEvents: TraceEvent[], nextEvents: TraceEvent[]): TraceEvent[] {
const merged: TraceEvent[] = [];
interface TraceEventMergeOptions {
previousCursor?: number | null;
nextRange?: TraceSnapshot["range"] | null;
}
function mergeTraceEvents(previousEvents: TraceEvent[], nextEvents: TraceEvent[], options: TraceEventMergeOptions = {}): TraceEvent[] {
if (previousEvents.length === 0) return nextEvents;
if (nextEvents.length === 0 || previousEvents === nextEvents) return previousEvents;
const previousCursor = finiteTraceCursor(options.previousCursor) ?? traceEventsProjectedSeqCursor(previousEvents);
const nextToProjectedSeq = traceRangeToProjectedSeq(options.nextRange);
if (previousCursor !== null && nextToProjectedSeq !== null && nextToProjectedSeq <= previousCursor) return previousEvents;
const appended = previousCursor !== null ? appendTraceEventDelta(previousEvents, nextEvents, previousCursor) : null;
if (appended) return appended;
const merged = [...previousEvents];
const keys = new Set<string>();
for (const event of [...previousEvents, ...nextEvents]) {
for (const event of previousEvents) {
const key = traceEventIdentity(event);
if (key) keys.add(key);
}
let changed = false;
for (const event of nextEvents) {
const key = traceEventIdentity(event);
if (key && keys.has(key)) continue;
if (key) keys.add(key);
merged.push(event);
changed = true;
}
return merged.sort((left, right) => traceEventSortSeq(left) - traceEventSortSeq(right));
if (!changed) return previousEvents;
return traceEventsAreSorted(merged) ? merged : merged.sort((left, right) => traceEventSortSeq(left) - traceEventSortSeq(right));
}
function appendTraceEventDelta(previousEvents: TraceEvent[], nextEvents: TraceEvent[], previousCursor: number): TraceEvent[] | null {
const delta: TraceEvent[] = [];
for (const event of nextEvents) {
const seq = traceEventProjectedSeq(event);
if (!Number.isFinite(seq)) return null;
if (seq > previousCursor) delta.push(event);
}
if (delta.length === 0) return previousEvents;
if (!traceEventsAreSorted(delta)) return null;
return [...previousEvents, ...delta];
}
function traceProjectedSeqCursor(trace: ChatMessage["runnerTrace"] | TraceSnapshot | null | undefined, events: TraceEvent[] = Array.isArray(trace?.events) ? trace.events : []): number | null {
return finiteTraceCursor(trace?.nextProjectedSeq) ?? traceRangeToProjectedSeq(trace?.range) ?? traceEventsProjectedSeqCursor(events);
}
function traceRangeToProjectedSeq(range: TraceSnapshot["range"] | null | undefined): number | null {
return finiteTraceCursor(range?.toProjectedSeq);
}
function traceEventsProjectedSeqCursor(events: TraceEvent[]): number | null {
let cursor: number | null = null;
for (const event of events) {
const seq = traceEventProjectedSeq(event);
if (Number.isFinite(seq) && (cursor === null || seq > cursor)) cursor = Math.trunc(seq);
}
return cursor;
}
function traceEventsAreSorted(events: TraceEvent[]): boolean {
let previous = Number.NEGATIVE_INFINITY;
for (const event of events) {
const seq = traceEventSortSeq(event);
if (seq < previous) return false;
previous = seq;
}
return true;
}
function finiteTraceCursor(value: unknown): number | null {
const number = Number(value);
return Number.isFinite(number) && number >= 0 ? Math.trunc(number) : null;
}
function traceEventSortSeq(event: TraceEvent): number {
@@ -0,0 +1,53 @@
import assert from "node:assert/strict";
import { test } from "bun:test";
import type { AgentChatResultResponse, ChatMessage } from "../types";
import { terminalSealResultWithoutTraceEvents, traceHydrationProjectedSeq, traceNextProjectedSeq } from "./workbench-trace-hydration";
test("trace hydration cursor prefers range metadata over scanning events", () => {
const trace: ChatMessage["runnerTrace"] = {
traceId: "trc_cursor_metadata",
nextProjectedSeq: 42,
range: { toProjectedSeq: 42 },
events: [{ projectedSeq: 1, type: "event" }]
};
assert.equal(traceHydrationProjectedSeq(trace), 42);
});
test("trace result next cursor uses metadata before event fallback", () => {
const result: AgentChatResultResponse = {
traceId: "trc_result_cursor",
status: "running",
range: { toProjectedSeq: 12 },
events: [{ projectedSeq: 3, type: "event" }]
};
assert.equal(traceNextProjectedSeq(result, 8), 12);
});
test("terminal seal strips heavy trace rows while preserving final authority", () => {
const result: AgentChatResultResponse = {
traceId: "trc_terminal_seal",
status: "completed",
terminal: true,
finalResponse: { text: "final answer" },
events: [{ projectedSeq: 1, type: "event" }],
traceEvents: [{ projectedSeq: 2, type: "event" }],
runnerTrace: {
traceId: "trc_terminal_seal",
status: "completed",
traceStatus: "completed",
events: [{ projectedSeq: 3, type: "event" }]
}
};
const sealed = terminalSealResultWithoutTraceEvents(result);
assert.equal(sealed.status, "completed");
assert.equal(sealed.finalResponse?.text, "final answer");
assert.equal(sealed.events, undefined);
assert.equal(sealed.traceEvents, undefined);
assert.equal(sealed.runnerTrace?.events, undefined);
assert.equal(sealed.runnerTrace?.traceStatus, "completed");
});
@@ -0,0 +1,43 @@
import type { AgentChatResultResponse, ChatMessage, TraceEvent } from "../types";
export function traceHydrationProjectedSeq(trace: ChatMessage["runnerTrace"]): number {
const cursor = traceHydrationCursor(trace);
if (cursor !== null) return cursor;
const events = Array.isArray(trace?.events) ? trace.events : [];
return eventsProjectedSeq(events, 0);
}
export function traceNextProjectedSeq(result: AgentChatResultResponse, fallback: number): number {
const cursor = traceHydrationCursor(result);
if (cursor !== null) return cursor;
const events = Array.isArray(result.events) ? result.events : Array.isArray(result.traceEvents) ? result.traceEvents : [];
return eventsProjectedSeq(events, fallback);
}
export function terminalSealResultWithoutTraceEvents(result: AgentChatResultResponse): AgentChatResultResponse {
const { events: _events, traceEvents: _traceEvents, runnerTrace, ...rest } = result;
if (!runnerTrace) return rest;
const { events: _runnerEvents, ...traceRest } = runnerTrace;
return { ...rest, runnerTrace: traceRest };
}
function traceHydrationCursor(source: { nextProjectedSeq?: unknown; range?: { toProjectedSeq?: unknown } | null } | null | undefined): number | null {
return finiteProjectedSeq(source?.nextProjectedSeq) ?? finiteProjectedSeq(source?.range?.toProjectedSeq);
}
function finiteProjectedSeq(value: unknown): number | null {
const number = Number(value);
return Number.isFinite(number) && number >= 0 ? Math.trunc(number) : null;
}
function eventsProjectedSeq(events: TraceEvent[], fallback: number): number {
return events.reduce((max, event) => {
const seq = traceEventProjectedSeq(event);
return Number.isFinite(seq) && seq > max ? seq : max;
}, fallback);
}
function traceEventProjectedSeq(event: TraceEvent): number {
const seq = Number(event.projectedSeq);
return Number.isFinite(seq) ? Math.trunc(seq) : Number.NaN;
}
+6 -25
View File
@@ -71,6 +71,7 @@ import { useWorkbenchColadaMutations } from "./workbench-colada-mutations";
import { useWorkbenchColadaQueries } from "./workbench-colada-queries";
import { useWorkbenchColadaReducer } from "./workbench-colada-reducer";
import { projectionMergeCommitSummary, workbenchSessionDetailReadKey, workbenchSessionMessagesReadKey } from "./workbench-session-messages-read-budget";
import { terminalSealResultWithoutTraceEvents, traceHydrationProjectedSeq, traceNextProjectedSeq } from "./workbench-trace-hydration";
const WORKBENCH_SESSION_PROJECTION_SIGNAL_CHANNEL = "hwlab.workbench.sessionProjection.v1";
const WORKBENCH_SESSION_PROJECTION_SIGNAL_KEY = "hwlab.workbench.sessionProjectionSignal.v1";
@@ -951,6 +952,11 @@ export const useWorkbenchStore = defineStore("workbench", () => {
const activityLabel = firstNonEmptyString(result.lastEventLabel, result.status);
if (ownerSessionId === activeSessionId.value && (events.length > 0 || turnResultIsTerminalForMerge(result))) recordActivity(`trace:${activityLabel ?? "hydrated"}`);
markWorkbenchTraceEventsReceived({ traceId, events, transport: "rest_gap" });
if (turnResultIsTerminalForMerge(result)) {
const terminalSeal = terminalSealResultWithoutTraceEvents(result);
rememberTurnStatus(traceId, terminalSeal);
projectTurnAuthorityToMessages(traceId, terminalSeal, "trace-hydration-terminal-seal");
}
updateSessionMessages(ownerSessionId, (source) => source.map((message) => {
if (!messageMatchesTraceAuthority(message, traceId, authoritySessionId, ownerSessionId)) return message;
const traceDetailStatus = firstNonEmptyString(result.traceStatus, result.runnerTrace?.traceStatus, result.runnerTrace?.status);
@@ -2049,31 +2055,6 @@ export async function liveCall<T>(label: string, call: () => Promise<ApiResult<T
}
}
function traceHydrationProjectedSeq(trace: ChatMessage["runnerTrace"]): number {
const rangeNext = Number(trace?.nextProjectedSeq ?? trace?.range?.toProjectedSeq);
if (Number.isFinite(rangeNext) && rangeNext > 0 && trace?.hasMore === true) return Math.trunc(rangeNext);
const events = Array.isArray(trace?.events) ? trace.events : [];
return events.reduce((max, event) => {
const seq = traceEventProjectedSeq(event);
return Number.isFinite(seq) && seq > max ? Math.trunc(seq) : max;
}, 0);
}
function traceNextProjectedSeq(result: AgentChatResultResponse, fallback: number): number {
const direct = Number(result.nextProjectedSeq ?? result.range?.toProjectedSeq);
if (Number.isFinite(direct) && direct >= 0) return Math.trunc(direct);
const events = Array.isArray(result.events) ? result.events : Array.isArray(result.traceEvents) ? result.traceEvents : [];
return events.reduce((max, event) => {
const seq = traceEventProjectedSeq(event);
return Number.isFinite(seq) && seq > max ? Math.trunc(seq) : max;
}, fallback);
}
function traceEventProjectedSeq(event: TraceEvent): number {
const seq = Number(event.projectedSeq);
return Number.isFinite(seq) ? Math.trunc(seq) : Number.NaN;
}
function makeMessage(role: ChatMessage["role"], text: string, status: ChatMessage["status"], extra: Partial<ChatMessage> = {}): ChatMessage {
return { id: nextProtocolId("msg"), role, title: extra.title ?? (role === "user" ? "用户" : "Code Agent"), text, status, createdAt: new Date().toISOString(), ...extra };
}