From ba86dccc8d343021aa88bf8330668a90b9814830 Mon Sep 17 00:00:00 2001 From: chengyongru <61816729+chengyongru@users.noreply.github.com> Date: Wed, 15 Jul 2026 16:46:05 +0800 Subject: [PATCH] fix(webui): correct activity timer duration (#4649) --- nanobot/webui/transcript.py | 70 +++++++++++++++---- tests/channels/test_websocket_channel.py | 14 +++- tests/utils/test_webui_transcript.py | 36 ++++++++++ .../thread/AgentActivityCluster.tsx | 16 ++++- .../src/components/thread/ThreadMessages.tsx | 1 + webui/src/lib/activity-timeline.ts | 36 ++++++++-- webui/src/tests/thread-messages.test.tsx | 45 +++++++++++- 7 files changed, 195 insertions(+), 23 deletions(-) diff --git a/nanobot/webui/transcript.py b/nanobot/webui/transcript.py index 28ac7c9f..63a9aac9 100644 --- a/nanobot/webui/transcript.py +++ b/nanobot/webui/transcript.py @@ -584,9 +584,30 @@ def _append_to_active_transcript(session_key: str, obj: dict[str, Any]) -> None: os.fsync(f.fileno()) +def _now_ms() -> int: + return int(time.time() * 1000) + + +def _valid_created_at_ms(value: Any) -> int | None: + if isinstance(value, bool): + return None + if isinstance(value, (int, float)) and value >= 0 and value < 10_000_000_000_000_000: + return int(value) + return None + + +def _record_for_append(obj: dict[str, Any]) -> dict[str, Any]: + if _valid_created_at_ms(obj.get("created_at_ms")) is not None: + return obj + record = dict(obj) + record["created_at_ms"] = _now_ms() + return record + + def append_transcript_object(session_key: str, obj: dict[str, Any]) -> None: - _append_to_active_transcript(session_key, obj) - if obj.get("event") == "turn_end": + record = _record_for_append(obj) + _append_to_active_transcript(session_key, record) + if record.get("event") == "turn_end": _rotate_active_transcript_if_needed(session_key) @@ -1257,13 +1278,19 @@ def replay_transcript_to_ui_messages( active_activity_segment_id: str | None = None active_file_edit_segment_id: str | None = None activity_segment_counter = 0 - _ts_base = int(time.time() * 1000) + _ts_base = _now_ms() closed_turn_ids: set[str] = set() replay_turn_aliases: dict[str, str] = {} def _new_id(prefix: str, idx: int) -> str: return f"{prefix}-{idx}-{uuid.uuid4().hex[:8]}" + def _created_at_ms(rec: dict[str, Any], idx: int) -> int: + created_at_ms = _valid_created_at_ms(rec.get("created_at_ms")) + if created_at_ms is not None: + return created_at_ms + return _ts_base + idx + def _new_activity_segment(*, activate: bool = True) -> str: nonlocal active_activity_segment_id, activity_segment_counter activity_segment_counter += 1 @@ -1330,6 +1357,7 @@ def replay_transcript_to_ui_messages( chunk: str, idx: int, turn_fields: dict[str, Any] | None = None, + created_at_ms: int | None = None, ) -> None: turn_fields = turn_fields or {} for i in range(len(prev) - 1, -1, -1): @@ -1379,7 +1407,7 @@ def replay_transcript_to_ui_messages( "reasoningStreaming": True, "activitySegmentId": segment, **turn_fields, - "createdAt": _ts_base + idx, + "createdAt": created_at_ms if created_at_ms is not None else _ts_base + idx, }, ) @@ -1474,7 +1502,7 @@ def replay_transcript_to_ui_messages( } return - def absorb_complete(extra: dict[str, Any], idx: int) -> None: + def absorb_complete(extra: dict[str, Any], idx: int, created_at_ms: int) -> None: nonlocal active_activity_segment_id, active_file_edit_segment_id last = messages[-1] if messages else None if last and is_reasoning_only_placeholder(last) and _same_turn(last, extra): @@ -1489,7 +1517,7 @@ def replay_transcript_to_ui_messages( { "id": _new_id("as", idx), "role": "assistant", - "createdAt": _ts_base + idx, + "createdAt": created_at_ms, **extra, }, ) @@ -1576,6 +1604,7 @@ def replay_transcript_to_ui_messages( edits: list[dict[str, Any]], idx: int, turn_fields: dict[str, Any] | None = None, + created_at_ms: int | None = None, ) -> None: nonlocal active_file_edit_segment_id turn_fields = turn_fields or {} @@ -1606,7 +1635,7 @@ def replay_transcript_to_ui_messages( "fileEdits": [], "activitySegmentId": segment, **turn_fields, - "createdAt": _ts_base + idx, + "createdAt": created_at_ms if created_at_ms is not None else _ts_base + idx, }, ) target_index = len(messages) - 1 @@ -1671,7 +1700,7 @@ def replay_transcript_to_ui_messages( "role": "user", "content": text_s, **_turn_fields(rec, "user"), - "createdAt": _ts_base + idx, + "createdAt": _created_at_ms(rec, idx), } if media_att: row["media"] = media_att @@ -1695,6 +1724,7 @@ def replay_transcript_to_ui_messages( [e for e in raw_edits if isinstance(e, dict)], idx, _turn_fields(rec, "activity"), + _created_at_ms(rec, idx), ) continue @@ -1719,7 +1749,7 @@ def replay_transcript_to_ui_messages( "content": "", "isStreaming": True, **_turn_fields(rec, "answer"), - "createdAt": _ts_base + idx, + "createdAt": _created_at_ms(rec, idx), }, ) buffer_parts.append(chunk) @@ -1751,7 +1781,7 @@ def replay_transcript_to_ui_messages( "content": final_text, "isStreaming": True, **_turn_fields(rec, "answer"), - "createdAt": _ts_base + idx, + "createdAt": _created_at_ms(rec, idx), }, ) else: @@ -1775,7 +1805,13 @@ def replay_transcript_to_ui_messages( if not isinstance(chunk, str) or not chunk: continue close_file_edit_phase_before_activity() - attach_reasoning_chunk(messages, chunk, idx, _turn_fields(rec, "reasoning")) + attach_reasoning_chunk( + messages, + chunk, + idx, + _turn_fields(rec, "reasoning"), + _created_at_ms(rec, idx), + ) continue if ev == "reasoning_end": @@ -1797,7 +1833,13 @@ def replay_transcript_to_ui_messages( if not isinstance(line, str) or not line: continue close_file_edit_phase_before_activity() - attach_reasoning_chunk(messages, line, idx, _turn_fields(rec, "reasoning")) + attach_reasoning_chunk( + messages, + line, + idx, + _turn_fields(rec, "reasoning"), + _created_at_ms(rec, idx), + ) close_reasoning(messages) continue if kind in ("tool_hint", "progress"): @@ -1853,7 +1895,7 @@ def replay_transcript_to_ui_messages( **({"toolEvents": visible_structured_events} if visible_structured_events else {}), "activitySegmentId": segment, **_turn_fields(rec, "activity"), - "createdAt": _ts_base + idx, + "createdAt": _created_at_ms(rec, idx), }, ) continue @@ -1878,7 +1920,7 @@ def replay_transcript_to_ui_messages( extra["latencyMs"] = int(lat) extra.update(_turn_fields(rec, "answer")) extra.update(_source_fields(rec)) - absorb_complete(extra, idx) + absorb_complete(extra, idx, _created_at_ms(rec, idx)) if media: suppress_until_turn_end = True continue diff --git a/tests/channels/test_websocket_channel.py b/tests/channels/test_websocket_channel.py index 42410956..7ae7cb03 100644 --- a/tests/channels/test_websocket_channel.py +++ b/tests/channels/test_websocket_channel.py @@ -393,14 +393,23 @@ async def test_webui_message_envelope_marks_inbound_metadata(bus: MagicMock) -> assert msg.metadata["webui_turn_id"] == "turn-1" assert msg.metadata["_wants_stream"] is True lines = read_transcript_lines("websocket:chat-1") - assert lines == [{ + assert len(lines) == 1 + assert {key: lines[0].get(key) for key in ( + "event", + "chat_id", + "text", + "turn_id", + "turn_phase", + "turn_seq", + )} == { "event": "user", "chat_id": "chat-1", "text": "hello", "turn_id": "turn-1", "turn_phase": "user", "turn_seq": 1, - }] + } + assert isinstance(lines[0].get("created_at_ms"), int) @pytest.mark.asyncio @@ -2702,6 +2711,7 @@ async def test_webui_message_envelope_appends_user_transcript( assert isinstance(line.get("turn_id"), str) assert line.get("turn_phase") == "user" assert line.get("turn_seq") == 1 + assert isinstance(line.get("created_at_ms"), int) inbound = bus.publish_inbound.await_args.args[0] assert inbound.chat_id == "source" assert inbound.content == "round1" diff --git a/tests/utils/test_webui_transcript.py b/tests/utils/test_webui_transcript.py index 55a7d405..067a45fa 100644 --- a/tests/utils/test_webui_transcript.py +++ b/tests/utils/test_webui_transcript.py @@ -25,6 +25,17 @@ def test_append_and_read_roundtrip(tmp_path, monkeypatch) -> None: assert lines[0]["text"] == "hello" +def test_append_stamps_created_at_ms(tmp_path, monkeypatch) -> None: + monkeypatch.setattr("nanobot.config.paths.get_data_dir", lambda: tmp_path) + monkeypatch.setattr("nanobot.webui.transcript.time.time", lambda: 1_700_000_000.0) + key = "websocket:t-created-at" + + append_transcript_object(key, {"event": "user", "chat_id": "t-created-at", "text": "hello"}) + + lines = read_transcript_lines(key) + assert lines[0]["created_at_ms"] == 1_700_000_000_000 + + def _force_small_transcript_budget(monkeypatch, *, limit: int = 520, target: int = 260) -> None: monkeypatch.setattr("nanobot.webui.transcript._MAX_TRANSCRIPT_FILE_BYTES", limit) monkeypatch.setattr("nanobot.webui.transcript._TARGET_ACTIVE_TRANSCRIPT_BYTES", target) @@ -291,6 +302,31 @@ def test_replay_delta_and_turn_end(tmp_path, monkeypatch) -> None: assert msgs[1]["latencyMs"] == 42 +def test_replay_uses_persisted_created_at_ms() -> None: + msgs = replay_transcript_to_ui_messages( + [ + { + "event": "user", + "chat_id": "t-created-at", + "text": "q", + "created_at_ms": 1_700_000_000_000, + }, + { + "event": "message", + "chat_id": "t-created-at", + "kind": "tool_hint", + "text": "exec()", + "created_at_ms": 1_700_000_230_000, + }, + ], + ) + + assert [message["createdAt"] for message in msgs] == [ + 1_700_000_000_000, + 1_700_000_230_000, + ] + + def test_thread_response_does_not_mark_completed_message_tool_tail_pending( tmp_path, monkeypatch, diff --git a/webui/src/components/thread/AgentActivityCluster.tsx b/webui/src/components/thread/AgentActivityCluster.tsx index 6728cd07..d0e4f757 100644 --- a/webui/src/components/thread/AgentActivityCluster.tsx +++ b/webui/src/components/thread/AgentActivityCluster.tsx @@ -175,6 +175,8 @@ interface AgentActivityClusterProps { hasBodyBelow: boolean; /** Persisted end-to-end turn latency from the assistant answer, used for history replay. */ turnLatencyMs?: number; + /** User turn start timestamp for live activity before the first trace/reasoning row. */ + startedAtMs?: number; cliApps?: CliAppInfo[]; mcpPresets?: McpPresetInfo[]; onOpenFilePreview?: (path: string) => void; @@ -189,6 +191,7 @@ export function AgentActivityCluster({ isTurnStreaming, hasBodyBelow, turnLatencyMs, + startedAtMs, cliApps = [], mcpPresets = [], onOpenFilePreview, @@ -251,7 +254,13 @@ export function AgentActivityCluster({ const hasVisibleActivity = reasoningSteps > 0 || toolCalls > 0 || cliCount > 0 || mcpCount > 0 || fileCount > 0; const hasOnlyFileActivity = fileCount > 0 && messages.every(messageHasOnlyFileActivity); const hasNonReasoningActivity = toolCalls > 0 || cliCount > 0 || mcpCount > 0 || fileCount > 0; - const durationMs = activityDurationMs(messages, isTurnStreaming, now, turnLatencyMs); + const durationMs = activityDurationMs( + messages, + isTurnStreaming, + now, + turnLatencyMs, + startedAtMs, + ); const activityDuration = formatActivityDuration(durationMs); const thoughtLabel = hasNonReasoningActivity ? isTurnStreaming @@ -654,6 +663,7 @@ function activityDurationMs( active: boolean, now: number, completedLatencyMs?: number, + activeStartedAtMs?: number, ): number { if (!active && Number.isFinite(completedLatencyMs) && completedLatencyMs! >= 0) { return Math.round(completedLatencyMs!); @@ -662,7 +672,9 @@ function activityDurationMs( .map((message) => message.createdAt) .filter((value) => Number.isFinite(value)); if (!timestamps.length) return 0; - const first = Math.min(...timestamps); + const first = active && Number.isFinite(activeStartedAtMs) + ? activeStartedAtMs! + : Math.min(...timestamps); const last = active && first > 1_000_000_000_000 ? now : Math.max(...timestamps); diff --git a/webui/src/components/thread/ThreadMessages.tsx b/webui/src/components/thread/ThreadMessages.tsx index 29a68505..2aabb25f 100644 --- a/webui/src/components/thread/ThreadMessages.tsx +++ b/webui/src/components/thread/ThreadMessages.tsx @@ -104,6 +104,7 @@ export function ThreadMessages({ isTurnStreaming={liveActivityClusterIndices.has(index)} hasBodyBelow={hasBodyBelow} turnLatencyMs={unit.turnLatencyMs} + startedAtMs={unit.startedAtMs} cliApps={cliApps} mcpPresets={mcpPresets} onOpenFilePreview={onOpenFilePreview} diff --git a/webui/src/lib/activity-timeline.ts b/webui/src/lib/activity-timeline.ts index ce1ebbcf..3f012ba8 100644 --- a/webui/src/lib/activity-timeline.ts +++ b/webui/src/lib/activity-timeline.ts @@ -35,7 +35,13 @@ export interface ActivityGroup { } export type TurnUnit = - | { type: "activity"; messages: UIMessage[]; items: ActivityItem[]; turnLatencyMs?: number } + | { + type: "activity"; + messages: UIMessage[]; + items: ActivityItem[]; + turnLatencyMs?: number; + startedAtMs?: number; + } | { type: "message"; message: UIMessage }; interface NormalizeActivityTimelineOptions { @@ -91,11 +97,16 @@ export function normalizeActivityTimeline( const units: TurnUnit[] = []; let turnMessages: UIMessage[] = []; let activeTurnId: string | undefined; + let activeTurnStartedAtMs: number | undefined; const flushTurn = (flushOptions: NormalizeActivityTimelineOptions = {}) => { - if (turnMessages.length === 0) return; + if (turnMessages.length === 0) { + activeTurnId = undefined; + return; + } const turnUnits: TurnUnit[] = []; + const turnStartedAtMs = activeTurnStartedAtMs; const orderedTurnMessages = orderMessagesByTurnSeq(turnMessages); const visibleMessages = visibleMessagesForTurn(orderedTurnMessages); let visibleIndex = 0; @@ -103,7 +114,12 @@ export function normalizeActivityTimeline( const flushActivityMessages = () => { if (!activityMessages.length) return; - pushActivityUnits(turnUnits, activityMessages, visibleMessages.slice(visibleIndex)); + pushActivityUnits( + turnUnits, + activityMessages, + visibleMessages.slice(visibleIndex), + turnStartedAtMs, + ); activityMessages = []; }; @@ -130,6 +146,7 @@ export function normalizeActivityTimeline( units.push(...normalizeCompletedTurnUnits(turnUnits, flushOptions)); turnMessages = []; activeTurnId = undefined; + activeTurnStartedAtMs = undefined; }; for (const message of messages) { @@ -137,6 +154,7 @@ export function normalizeActivityTimeline( flushTurn(); units.push({ type: "message", message }); activeTurnId = message.turnId; + activeTurnStartedAtMs = validCreatedAtMs(message.createdAt); continue; } @@ -206,7 +224,16 @@ function visibleMessagesForTurn(messages: UIMessage[]): UIMessage[] { return visibleMessages; } -function pushActivityUnits(units: TurnUnit[], activityMessages: UIMessage[], visibleMessages: UIMessage[]) { +function validCreatedAtMs(value: unknown): number | undefined { + return typeof value === "number" && Number.isFinite(value) ? value : undefined; +} + +function pushActivityUnits( + units: TurnUnit[], + activityMessages: UIMessage[], + visibleMessages: UIMessage[], + startedAtMs?: number, +) { let runMessages: UIMessage[] = []; let runBucket: "file" | "other" | undefined; let runSegmentId: string | undefined; @@ -218,6 +245,7 @@ function pushActivityUnits(units: TurnUnit[], activityMessages: UIMessage[], vis messages: runMessages, items: runMessages.flatMap(activityItemsForMessage), turnLatencyMs: activityTurnLatencyMs(runMessages, visibleMessages), + startedAtMs, }); runMessages = []; runBucket = undefined; diff --git a/webui/src/tests/thread-messages.test.tsx b/webui/src/tests/thread-messages.test.tsx index c84dc6bf..76152941 100644 --- a/webui/src/tests/thread-messages.test.tsx +++ b/webui/src/tests/thread-messages.test.tsx @@ -1,5 +1,5 @@ import { render, screen } from "@testing-library/react"; -import { describe, expect, it } from "vitest"; +import { afterEach, describe, expect, it, vi } from "vitest"; import { assistantCopyFlags, @@ -9,6 +9,10 @@ import { } from "@/components/thread/ThreadMessages"; import type { UIMessage } from "@/lib/types"; +afterEach(() => { + vi.useRealTimers(); +}); + describe("ThreadMessages", () => { it("groups consecutive reasoning and tool rows into one timeline before the answer", () => { const messages: UIMessage[] = [ @@ -293,6 +297,45 @@ describe("ThreadMessages", () => { expect(screen.queryByLabelText(/editing foo\.txt/i)).not.toBeInTheDocument(); }); + it("times live activity from the user turn start", () => { + vi.useFakeTimers(); + const startedAt = 1_700_000_000_000; + vi.setSystemTime(startedAt + 230_000); + const messages: UIMessage[] = [ + { + id: "u1", + role: "user", + content: "run it", + turnId: "turn-1", + turnPhase: "user", + turnSeq: 1, + createdAt: startedAt, + }, + { + id: "t1", + role: "tool", + kind: "trace", + content: "exec()", + traces: ["exec()"], + turnId: "turn-1", + turnPhase: "activity", + turnSeq: 2, + createdAt: startedAt + 220_000, + }, + ]; + + const units = buildDisplayUnits(messages, true); + + expect( + units[1].type === "activity" ? units[1].startedAtMs : undefined, + ).toBe(startedAt); + + render(); + + expect(screen.getByText("Working for 3m 50s")).toBeInTheDocument(); + expect(screen.queryByText("Working for 10s")).not.toBeInTheDocument(); + }); + it("folds final answer reasoning into the preceding activity timeline", () => { const messages: UIMessage[] = [ {