fix: restore Workbench session titles from Kafka replay
This commit is contained in:
@@ -2,8 +2,9 @@ import assert from "node:assert/strict";
|
||||
import { test } from "bun:test";
|
||||
|
||||
import type { ChatMessage, TraceEvent } from "../types";
|
||||
import { projectWorkbenchLiveKafkaMessage, projectWorkbenchLiveKafkaUserMessage, workbenchAgentMessageIdForTrace, workbenchLiveKafkaEnvelope, workbenchLiveKafkaEventIsUserMessage, workbenchLiveKafkaProjectionTarget, workbenchUserMessageIdForTrace } from "./workbench-live-kafka-event";
|
||||
import { projectWorkbenchLiveKafkaMessage, projectWorkbenchLiveKafkaUserMessage, projectWorkbenchLiveKafkaUserSession, workbenchAgentMessageIdForTrace, workbenchLiveKafkaEnvelope, workbenchLiveKafkaEventIsUserMessage, workbenchLiveKafkaProjectionTarget, workbenchUserMessageIdForTrace } from "./workbench-live-kafka-event";
|
||||
import { createWorkbenchServerState, reduceWorkbenchServerState, selectActiveMessages, type WorkbenchServerState } from "./workbench-server-state";
|
||||
import { sessionToSessionTab } from "./workbench-session";
|
||||
|
||||
test("only canonical HWLAB Kafka envelopes bypass terminal priority", () => {
|
||||
assert.equal(workbenchLiveKafkaEnvelope("hwlab.event.v1"), true);
|
||||
@@ -41,6 +42,54 @@ test("live Kafka user event restores and idempotently updates the user bubble",
|
||||
assert.equal(second?.updatedAt, event.createdAt);
|
||||
});
|
||||
|
||||
test("Kafka user replay rebuilds the session rail title without optimistic state", () => {
|
||||
const sessionId = "ses_live_user_title";
|
||||
const traceId = "trc_live_user_title";
|
||||
const firstEvent = traceEvent(1, "user", {
|
||||
userMessageId: "msg_live_user_title_user",
|
||||
messageId: "msg_live_user_title_user",
|
||||
text: "persisted session title",
|
||||
createdAt: "2026-07-10T10:00:00.000Z"
|
||||
});
|
||||
const firstMessage = projectWorkbenchLiveKafkaUserMessage({
|
||||
previous: null,
|
||||
traceId,
|
||||
sessionId,
|
||||
event: firstEvent,
|
||||
receivedAt: "2026-07-10T10:01:00.000Z"
|
||||
});
|
||||
assert.ok(firstMessage);
|
||||
const replayedSession = projectWorkbenchLiveKafkaUserSession({ previous: null, message: firstMessage, sessionId });
|
||||
assert.ok(replayedSession);
|
||||
assert.equal(replayedSession.firstUserMessagePreview, "persisted session title");
|
||||
assert.equal(replayedSession.updatedAt, firstEvent.createdAt);
|
||||
assert.equal(sessionToSessionTab(replayedSession, sessionId).label, "persisted session title");
|
||||
|
||||
const repeatedSession = projectWorkbenchLiveKafkaUserSession({ previous: replayedSession, message: firstMessage, sessionId });
|
||||
assert.deepEqual(repeatedSession, replayedSession);
|
||||
|
||||
const secondEvent = traceEvent(2, "user", {
|
||||
traceId: "trc_live_user_title_second",
|
||||
userMessageId: "msg_live_user_title_second_user",
|
||||
messageId: "msg_live_user_title_second_user",
|
||||
text: "later user input",
|
||||
createdAt: "2026-07-10T10:02:00.000Z"
|
||||
});
|
||||
const secondMessage = projectWorkbenchLiveKafkaUserMessage({
|
||||
previous: null,
|
||||
traceId: "trc_live_user_title_second",
|
||||
sessionId,
|
||||
event: secondEvent,
|
||||
receivedAt: "2026-07-10T10:03:00.000Z"
|
||||
});
|
||||
assert.ok(secondMessage);
|
||||
const updatedSession = projectWorkbenchLiveKafkaUserSession({ previous: replayedSession, message: secondMessage, sessionId });
|
||||
assert.ok(updatedSession);
|
||||
assert.equal(updatedSession.firstUserMessagePreview, "persisted session title");
|
||||
assert.equal(updatedSession.updatedAt, secondEvent.createdAt);
|
||||
assert.equal(sessionToSessionTab(updatedSession, sessionId).label, "persisted session title");
|
||||
});
|
||||
|
||||
test("HTTP-first and Kafka-first admission converge on one stable optimistic user bubble", () => {
|
||||
const sessionId = "ses_live_user_race";
|
||||
const traceId = "trc_live_user_race";
|
||||
|
||||
@@ -2,7 +2,7 @@
|
||||
// Responsibility: project one transparent hwlab.event.v1 business event into visible turn state without replay/finalizer/polling.
|
||||
|
||||
import { mergeRunnerTrace } from "../composables/workbench-trace-snapshot";
|
||||
import type { ChatMessage, TraceEvent, WorkbenchTurnTimingProjection } from "../types";
|
||||
import type { ChatMessage, TraceEvent, WorkbenchSessionRecord, WorkbenchTurnTimingProjection } from "../types";
|
||||
import { firstNonEmptyString } from "../utils";
|
||||
import { traceAssistantLifecycleEventIsSuperseded } from "../../../../tools/src/hwlab-cli/trace-renderer.ts";
|
||||
|
||||
@@ -32,6 +32,13 @@ export interface WorkbenchLiveKafkaUserProjectionInput {
|
||||
threadId?: string | null;
|
||||
}
|
||||
|
||||
export interface WorkbenchLiveKafkaUserSessionProjectionInput {
|
||||
previous: WorkbenchSessionRecord | null;
|
||||
message: ChatMessage;
|
||||
sessionId: string;
|
||||
threadId?: string | null;
|
||||
}
|
||||
|
||||
export function reduceWorkbenchLiveKafkaMessageState(previous: WorkbenchLiveMessageState, event: TraceEvent, previousEvents: TraceEvent[] = []): WorkbenchLiveMessageState {
|
||||
if (previous.terminal) return previous;
|
||||
const terminal = workbenchLiveKafkaEventIsTerminal(event);
|
||||
@@ -101,6 +108,20 @@ export function projectWorkbenchLiveKafkaUserMessage(input: WorkbenchLiveKafkaUs
|
||||
} as ChatMessage;
|
||||
}
|
||||
|
||||
export function projectWorkbenchLiveKafkaUserSession(input: WorkbenchLiveKafkaUserSessionProjectionInput): WorkbenchSessionRecord | null {
|
||||
const text = input.message.role === "user" ? firstNonEmptyString(input.message.text, input.message.content) : null;
|
||||
const lastUserMessageAt = timestamp(input.message.createdAt) ?? timestamp(input.message.updatedAt);
|
||||
if (!text || !lastUserMessageAt) return null;
|
||||
return {
|
||||
...(input.previous ?? {}),
|
||||
sessionId: input.sessionId,
|
||||
threadId: firstNonEmptyString(input.threadId, input.message.threadId, input.previous?.threadId),
|
||||
firstUserMessagePreview: firstNonEmptyString(input.previous?.firstUserMessagePreview, text),
|
||||
lastUserMessageAt,
|
||||
updatedAt: lastUserMessageAt
|
||||
};
|
||||
}
|
||||
|
||||
export function workbenchLiveKafkaEventIsTerminal(event: TraceEvent | null | undefined): boolean {
|
||||
return Boolean(event?.terminal === true || firstNonEmptyString(event?.eventType) === "terminal" || firstNonEmptyString(event?.type) === "result");
|
||||
}
|
||||
|
||||
@@ -21,7 +21,7 @@ import { initialWorkbenchSessionIdFromLocation } from "./workbench-projection";
|
||||
import { cleanupWorkbenchServerStateSessions, selectActiveMessages, selectActiveSession, selectSessionList, selectSessionStatusAuthority, selectTraceAuthorityById, selectTurnStatusAuthority, type WorkbenchServerAction } from "./workbench-server-state";
|
||||
import { cleanupDroppedWorkbenchSessionCaches, trimWorkbenchSessionCache } from "./workbench-session-cache";
|
||||
import { reduceWorkbenchRealtimeEvent, workbenchRealtimeEventIsBusinessActivity, type WorkbenchRealtimeAction } from "./workbench-event-reducer";
|
||||
import { projectWorkbenchLiveKafkaMessage, projectWorkbenchLiveKafkaUserMessage, workbenchLiveKafkaEnvelope, workbenchLiveKafkaProjectionTarget } from "./workbench-live-kafka-event";
|
||||
import { projectWorkbenchLiveKafkaMessage, projectWorkbenchLiveKafkaUserMessage, projectWorkbenchLiveKafkaUserSession, workbenchLiveKafkaEnvelope, workbenchLiveKafkaProjectionTarget } from "./workbench-live-kafka-event";
|
||||
import { projectRejectedWorkbenchAdmission } from "./workbench-admission-failure";
|
||||
import { messageHasSealedTerminalResult, messageIsSealedTerminal, traceAuthorityIsSealed } from "./workbench-terminal-authority";
|
||||
import {
|
||||
@@ -1075,14 +1075,8 @@ export const useWorkbenchStore = defineStore("workbench", () => {
|
||||
if (userMessage) {
|
||||
reduceServerState({ type: "message.upsert", sessionId: ownerSessionId, message: userMessage });
|
||||
const existing = sessions.value.find((session) => session.sessionId === ownerSessionId) ?? null;
|
||||
if (existing) {
|
||||
const lastUserMessageAt = firstNonEmptyString(userMessage.createdAt, userMessage.updatedAt);
|
||||
rememberSessionList(mergeSessionIntoList(sessions.value, {
|
||||
...existing,
|
||||
lastUserMessageAt,
|
||||
updatedAt: lastUserMessageAt ?? existing.updatedAt
|
||||
}));
|
||||
}
|
||||
const session = projectWorkbenchLiveKafkaUserSession({ previous: existing, message: userMessage, sessionId: ownerSessionId, threadId });
|
||||
if (session) rememberSessionList(mergeSessionIntoList(sessions.value, session));
|
||||
}
|
||||
else error.value = "workbench_live_user_message_invalid";
|
||||
return;
|
||||
|
||||
Reference in New Issue
Block a user