Merge pull request #2365 from pikasTech/fix/2356-colada-trace-hydrate
fix: bound Workbench trace hydrate merge
This commit is contained in:
@@ -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;
|
||||
}
|
||||
@@ -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 };
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user