diff --git a/gui/src/pages/Logs.tsx b/gui/src/pages/Logs.tsx index 07aec2b12a..d211823664 100644 --- a/gui/src/pages/Logs.tsx +++ b/gui/src/pages/Logs.tsx @@ -26,6 +26,7 @@ import { sanitizeLogEntryRouteDecision, validCachedRouteDecision, } from "./log-route-decision"; +import { mergeLogDelta, parseLogPollResponse } from "./log-poll"; function logsCacheKey(apiBase: string): string { return `ocx.logs.list.v1:${apiBase}`; @@ -379,6 +380,11 @@ export default function Logs({ apiBase }: { apiBase: string }) { const logRetryRef = useRef<{ key: string; failures: number; nextAttemptAt: number; error: unknown }>( { key: resourceKey, failures: 0, nextAttemptAt: 0, error: null }, ); + const logPollRef = useRef<{ + key: string; + cursor: string | null; + rows: LogEntry[]; + }>({ key: resourceKey, cursor: null, rows: cachedLogs ?? [] }); const localeTag = LOCALES.find(l => l.code === locale)?.htmlLang; // The proxy's own zone, so timestamps read the same as the server's logs rather than being // silently shifted into the viewer's zone (#725). Fetched once: it cannot change while the @@ -433,15 +439,40 @@ export default function Logs({ apiBase }: { apiBase: string }) { logRetryRef.current = retry; } if (retry.failures > 0 && Date.now() < retry.nextAttemptAt) throw retry.error; + + let pollState = logPollRef.current; + if (pollState.key !== resourceKey) { + pollState = { key: resourceKey, cursor: null, rows: readSessionListCache(resourceKey) ?? [] }; + logPollRef.current = pollState; + } + + const currentCursor = pollState.cursor; + const url = currentCursor + ? `${apiBase}/api/logs?limit=2000&cursor=${encodeURIComponent(currentCursor)}` + : `${apiBase}/api/logs?limit=2000`; + try { - const res = await fetch(`${apiBase}/api/logs?limit=2000`, { signal }); + const res = await fetch(url, { signal }); if (!res.ok) throw new Error(`${res.status} ${res.statusText}`.trim()); - const body = await res.json() as LogEntry[] | { logs?: LogEntry[] }; - const raw = Array.isArray(body) ? body : (body.logs ?? []); - const next = raw.map(sanitizeLogEntryRouteDecision); + const body = await res.json(); + const parsed = parseLogPollResponse(body); + const sanitizedIncoming = parsed.rows.map(sanitizeLogEntryRouteDecision); + + let nextRows: LogEntry[]; + if (currentCursor && parsed.cursorCapable && !parsed.reset) { + nextRows = mergeLogDelta(pollState.rows, sanitizedIncoming, 2000); + } else { + nextRows = sanitizedIncoming; + } + + logPollRef.current = { + key: resourceKey, + cursor: parsed.cursor, + rows: nextRows, + }; logRetryRef.current = { key: resourceKey, failures: 0, nextAttemptAt: 0, error: null }; - writeSessionListCache(resourceKey, next); - return next; + writeSessionListCache(resourceKey, nextRows); + return nextRows; } catch (error) { if (signal.aborted) throw error; const normalized = error ?? new Error("log request failed"); @@ -473,6 +504,7 @@ export default function Logs({ apiBase }: { apiBase: string }) { const fetchLogs = logsResource.refresh; const retryLogs = useCallback(() => { logRetryRef.current = { key: resourceKey, failures: 0, nextAttemptAt: 0, error: null }; + logPollRef.current = { key: resourceKey, cursor: null, rows: logPollRef.current.rows }; fetchLogs({ forceLoading: true }); }, [fetchLogs, resourceKey]); diff --git a/gui/src/pages/log-poll.ts b/gui/src/pages/log-poll.ts new file mode 100644 index 0000000000..1f22be362b --- /dev/null +++ b/gui/src/pages/log-poll.ts @@ -0,0 +1,54 @@ +export interface ParsedLogPollResponse { + rows: T[]; + cursor: string | null; + reset: boolean; + cursorCapable: boolean; +} + +export function parseLogPollResponse(body: unknown): ParsedLogPollResponse { + if (Array.isArray(body)) { + return { + rows: body as T[], + cursor: null, + reset: false, + cursorCapable: false, + }; + } + if (typeof body === "object" && body !== null) { + const candidate = body as { logs?: unknown; cursor?: unknown; reset?: unknown }; + const rows = Array.isArray(candidate.logs) ? (candidate.logs as T[]) : []; + const cursor = typeof candidate.cursor === "string" && candidate.cursor.trim() ? candidate.cursor.trim() : null; + const reset = candidate.reset === true; + const cursorCapable = cursor !== null || reset || Object.hasOwn(candidate, "cursor") || Object.hasOwn(candidate, "reset"); + return { + rows, + cursor, + reset, + cursorCapable, + }; + } + return { + rows: [], + cursor: null, + reset: false, + cursorCapable: false, + }; +} + +export function mergeLogDelta( + previous: readonly T[], + incoming: readonly T[], + cap = 2000, +): T[] { + if (incoming.length === 0) { + return previous.length > cap ? previous.slice(previous.length - cap) : [...previous]; + } + const incomingIds = new Set( + incoming + .map(row => row.requestId) + .filter((id): id is string => typeof id === "string" && id.length > 0), + ); + const filteredPrev = previous.filter(row => !row.requestId || !incomingIds.has(row.requestId)); + const merged = [...filteredPrev, ...incoming]; + return merged.length > cap ? merged.slice(merged.length - cap) : merged; +} diff --git a/gui/tests/log-poll.test.ts b/gui/tests/log-poll.test.ts new file mode 100644 index 0000000000..b2995547ee --- /dev/null +++ b/gui/tests/log-poll.test.ts @@ -0,0 +1,73 @@ +import { describe, expect, test } from "bun:test"; +import { mergeLogDelta, parseLogPollResponse } from "../src/pages/log-poll"; + +describe("log-poll helpers", () => { + test("merge appends delta rows, replaces duplicate ids, and keeps newest cap", () => { + expect(mergeLogDelta( + [{ requestId: "a", value: 1 }, { requestId: "b", value: 1 }], + [{ requestId: "b", value: 2 }, { requestId: "c", value: 1 }], + 2, + )).toEqual([{ requestId: "b", value: 2 }, { requestId: "c", value: 1 }]); + }); + + test("merge preserves all rows when under cap and no duplicates", () => { + expect(mergeLogDelta( + [{ requestId: "a", value: 1 }], + [{ requestId: "b", value: 2 }], + 10, + )).toEqual([{ requestId: "a", value: 1 }, { requestId: "b", value: 2 }]); + }); + + test("merge returns previous capped rows when incoming is empty", () => { + expect(mergeLogDelta( + [{ requestId: "a", value: 1 }, { requestId: "b", value: 2 }], + [], + 2, + )).toEqual([{ requestId: "a", value: 1 }, { requestId: "b", value: 2 }]); + }); + + test("legacy array and object bodies are full snapshots", () => { + expect(parseLogPollResponse([{ requestId: "a" }])).toEqual({ + rows: [{ requestId: "a" }], + cursor: null, + reset: false, + cursorCapable: false, + }); + expect(parseLogPollResponse({ logs: [{ requestId: "a" }] })).toEqual({ + rows: [{ requestId: "a" }], + cursor: null, + reset: false, + cursorCapable: false, + }); + }); + + test("cursor response preserves reset metadata and cursorCapable flag", () => { + expect(parseLogPollResponse({ logs: [], cursor: "opaque", reset: true })).toEqual({ + rows: [], + cursor: "opaque", + reset: true, + cursorCapable: true, + }); + expect(parseLogPollResponse({ logs: [{ requestId: "x" }], cursor: "c-1", reset: false })).toEqual({ + rows: [{ requestId: "x" }], + cursor: "c-1", + reset: false, + cursorCapable: true, + }); + }); + + test("malformed response returns empty rows and cursorCapable false", () => { + expect(parseLogPollResponse(null)).toEqual({ + rows: [], + cursor: null, + reset: false, + cursorCapable: false, + }); + expect(parseLogPollResponse("string")).toEqual({ + rows: [], + cursor: null, + reset: false, + cursorCapable: false, + }); + }); +}); diff --git a/gui/tests/logs-auto-refresh.test.tsx b/gui/tests/logs-auto-refresh.test.tsx index 758eeb14fa..a9eec10aa7 100644 --- a/gui/tests/logs-auto-refresh.test.tsx +++ b/gui/tests/logs-auto-refresh.test.tsx @@ -540,3 +540,80 @@ test("Logs: an intercepted helper row is badged and filterable", async () => { await act(async () => { root.unmount(); }); }); + +test("Logs: incremental cursor delta polling, delta append, empty delta, reset, and legacy fallback", async () => { + const calls: string[] = []; + let pollStep = 0; + + globalThis.fetch = (async (input) => { + const url = String(input); + calls.push(url); + if (!url.includes("/api/logs")) return new Response(null, { status: 404 }); + + if (pollStep === 0) { + // Step 0: Initial fetch -> returns initial log and cursor + return jsonResponse({ logs: [sampleLog], cursor: "cursor-0", reset: false }); + } + if (pollStep === 1) { + // Step 1: Empty delta -> cursor unchanged, keep existing rows + return jsonResponse({ logs: [], cursor: "cursor-0", reset: false }); + } + if (pollStep === 2) { + // Step 2: Delta with new row -> append updatedLog + return jsonResponse({ logs: [updatedLog], cursor: "cursor-1", reset: false }); + } + if (pollStep === 3) { + // Step 3: Eviction / reset: true -> replaces list + const resetLog = { ...sampleLog, requestId: "req-reset", model: "gpt-reset" }; + return jsonResponse({ logs: [resetLog], cursor: "cursor-reset", reset: true }); + } + if (pollStep === 4) { + // Step 4: Legacy server returns array without cursor + const legacyLog = { ...sampleLog, requestId: "req-legacy", model: "gpt-legacy" }; + return jsonResponse([legacyLog]); + } + // Step 5: Next poll after legacy response must be full request without cursor + return jsonResponse([sampleLog]); + }) as typeof fetch; + + const { root, container } = await mountLogs(); + + await flushMicrotasks(); + expectTableLoaded(container, "gpt-test"); + const logCalls = () => calls.filter(u => u.includes("/api/logs")); + expect(logCalls()[0]).toBe("http://localhost/api/logs?limit=2000"); + + // Advance to Step 1 (Empty delta) + pollStep = 1; + await advanceSilentRefresh(2000); + expect(logCalls().at(-1)).toContain("cursor=cursor-0"); + expectTableLoaded(container, "gpt-test"); + + // Advance to Step 2 (Delta with updatedLog) + pollStep = 2; + await advanceSilentRefresh(2000); + expect(logCalls().at(-1)).toContain("cursor=cursor-0"); + expectTableLoaded(container, "gpt-test"); + expect(container.textContent).toContain("gpt-updated"); + + // Advance to Step 3 (Reset replaces whole list) + pollStep = 3; + await advanceSilentRefresh(2000); + expect(logCalls().at(-1)).toContain("cursor=cursor-1"); + expectTableLoaded(container, "gpt-reset"); + expect(container.textContent).not.toContain("gpt-updated"); + + // Advance to Step 4 (Legacy server response) + pollStep = 4; + await advanceSilentRefresh(2000); + expect(logCalls().at(-1)).toContain("cursor=cursor-reset"); + expectTableLoaded(container, "gpt-legacy"); + + // Advance to Step 5 (Next poll after legacy response has no cursor) + pollStep = 5; + await advanceSilentRefresh(2000); + expect(logCalls().at(-1)).toBe("http://localhost/api/logs?limit=2000"); + expectTableLoaded(container, "gpt-test"); + + await act(async () => { root.unmount(); }); +}); diff --git a/src/server/management/logs-usage-routes.ts b/src/server/management/logs-usage-routes.ts index 03d8f96f59..fb03127e01 100644 --- a/src/server/management/logs-usage-routes.ts +++ b/src/server/management/logs-usage-routes.ts @@ -68,6 +68,7 @@ import { import type { OcxClaudeCodeConfig, OcxConfig, OcxCustomModel, OcxProviderConfig } from "../../types"; import { drainAndShutdown } from "../lifecycle"; import { filterRequestLogs, filteredRequestLogCount, getRequestLogEntries, type RequestLogEntry } from "../request-log"; +import { decodeRequestLogCursor, encodeRequestLogCursor, sliceRequestLogsAfterCursor } from "../request-log-cursor"; import { estimateComboCost, estimateRequestCost, normalizeCostTokens, tokensPerSecond } from "../../usage/cost"; import { userCostOverlayVersion } from "../../usage/user-cost-overlays"; import type { PersistedUsageAttempt } from "../../usage/log"; @@ -133,12 +134,23 @@ export async function handleLogsUsageRoutes(ctx: ManagementContext): Promise, +): string { + const payload: SerializedCursorV1 = { + v: 1, + t: entry.timestamp, + id: entry.requestId, + }; + return Buffer.from(JSON.stringify(payload)).toString("base64url"); +} + +export function decodeRequestLogCursor(raw: string): RequestLogCursor | null { + if (typeof raw !== "string" || !raw || raw.length > 512) return null; + try { + const json = Buffer.from(raw, "base64url").toString("utf8"); + const parsed = JSON.parse(json); + if (typeof parsed !== "object" || parsed === null || Array.isArray(parsed)) return null; + if (parsed.v !== 1) return null; + if (typeof parsed.t !== "number" || !Number.isFinite(parsed.t) || parsed.t < 0) return null; + if (typeof parsed.id !== "string" || parsed.id.length === 0 || parsed.id.length > 256) return null; + return { timestamp: parsed.t, requestId: parsed.id }; + } catch { + return null; + } +} + +export function sliceRequestLogsAfterCursor( + logs: readonly RequestLogEntry[], + cursor: RequestLogCursor, +): { entries: RequestLogEntry[]; reset: boolean } { + if (logs.length === 0) return { entries: [], reset: true }; + // Search from the end towards the beginning since recent cursors are near the end + let matchIndex = -1; + for (let i = logs.length - 1; i >= 0; i--) { + const entry = logs[i]; + if (entry && entry.timestamp === cursor.timestamp && entry.requestId === cursor.requestId) { + matchIndex = i; + break; + } + } + if (matchIndex === -1) { + return { entries: [...logs], reset: true }; + } + return { entries: logs.slice(matchIndex + 1), reset: false }; +} diff --git a/tests/management-api-logs-metrics.test.ts b/tests/management-api-logs-metrics.test.ts index 0255005e2f..f3bec90e5a 100644 --- a/tests/management-api-logs-metrics.test.ts +++ b/tests/management-api-logs-metrics.test.ts @@ -57,6 +57,68 @@ function baseEntry(overrides: Partial): RequestLogEntry { } describe("GET /api/logs display metrics", () => { + test("supports cursor delta polling, empty delta, eviction reset, and invalid cursor rejection", async () => { + addRequestLog(baseEntry({ requestId: "req-1", timestamp: 1000, provider: "anthropic" })); + addRequestLog(baseEntry({ requestId: "req-2", timestamp: 2000, provider: "anthropic" })); + + // Initial request returns the full window, an opaque cursor, and reset: false + const initUrl = new URL("http://localhost/api/logs"); + const initRes = await handleManagementAPI(new Request(initUrl), initUrl, config); + expect(initRes?.status).toBe(200); + const initBody = await initRes!.json() as { logs: Array<{ requestId: string }>; cursor?: string; reset?: boolean; total: number }; + expect(initBody.logs.map(r => r.requestId)).toEqual(["req-1", "req-2"]); + expect(typeof initBody.cursor).toBe("string"); + expect(initBody.reset).toBe(false); + expect(initBody.total).toBe(2); + + const cursor1 = initBody.cursor!; + + // Empty delta when no new logs have been added + const pollUrl = new URL(`http://localhost/api/logs?cursor=${encodeURIComponent(cursor1)}`); + const pollRes = await handleManagementAPI(new Request(pollUrl), pollUrl, config); + expect(pollRes?.status).toBe(200); + const pollBody = await pollRes!.json() as { logs: Array<{ requestId: string }>; cursor?: string; reset?: boolean; total: number }; + expect(pollBody.logs).toEqual([]); + expect(pollBody.cursor).toBe(cursor1); + expect(pollBody.reset).toBe(false); + expect(pollBody.total).toBe(2); + + // Live cursor returns only new entries + addRequestLog(baseEntry({ requestId: "req-3", timestamp: 3000, provider: "anthropic" })); + const pollRes2 = await handleManagementAPI(new Request(pollUrl), pollUrl, config); + expect(pollRes2?.status).toBe(200); + const pollBody2 = await pollRes2!.json() as { logs: Array<{ requestId: string }>; cursor?: string; reset?: boolean; total: number }; + expect(pollBody2.logs.map(r => r.requestId)).toEqual(["req-3"]); + expect(pollBody2.cursor).not.toBe(cursor1); + expect(pollBody2.reset).toBe(false); + expect(pollBody2.total).toBe(3); + + // Filtered delta advances the raw cursor even if 0 matching rows in delta + const filterUrl = new URL(`http://localhost/api/logs?provider=openai&cursor=${encodeURIComponent(cursor1)}`); + const filterRes = await handleManagementAPI(new Request(filterUrl), filterUrl, config); + expect(filterRes?.status).toBe(200); + const filterBody = await filterRes!.json() as { logs: Array<{ requestId: string }>; cursor?: string; reset?: boolean; total: number }; + expect(filterBody.logs).toEqual([]); + expect(filterBody.cursor).toBe(pollBody2.cursor); + expect(filterBody.total).toBe(0); + + // Stale/evicted cursor returns full window with reset: true + const fakeCursor = Buffer.from(JSON.stringify({ v: 1, t: 999, id: "evicted-req" })).toString("base64url"); + const staleUrl = new URL(`http://localhost/api/logs?cursor=${fakeCursor}`); + const staleRes = await handleManagementAPI(new Request(staleUrl), staleUrl, config); + expect(staleRes?.status).toBe(200); + const staleBody = await staleRes!.json() as { logs: Array<{ requestId: string }>; cursor?: string; reset?: boolean; total: number }; + expect(staleBody.logs.map(r => r.requestId)).toEqual(["req-1", "req-2", "req-3"]); + expect(staleBody.reset).toBe(true); + + // Malformed cursor fails closed with 400 + const badUrl = new URL("http://localhost/api/logs?cursor=not-valid-cursor"); + const badRes = await handleManagementAPI(new Request(badUrl), badUrl, config); + expect(badRes?.status).toBe(400); + const badBody = await badRes!.json(); + expect(badBody).toEqual({ error: { code: "invalid_cursor", message: "invalid cursor" } }); + }); + test("reports filtered total before limit pagination", async () => { addRequestLog(baseEntry({ requestId: "ok-a", provider: "anthropic", status: 200 })); addRequestLog(baseEntry({ requestId: "ok-b", provider: "anthropic", status: 200 })); diff --git a/tests/request-log-cursor.test.ts b/tests/request-log-cursor.test.ts new file mode 100644 index 0000000000..f033462977 --- /dev/null +++ b/tests/request-log-cursor.test.ts @@ -0,0 +1,71 @@ +import { describe, expect, test } from "bun:test"; +import { + decodeRequestLogCursor, + encodeRequestLogCursor, + sliceRequestLogsAfterCursor, +} from "../src/server/request-log-cursor"; +import type { RequestLogEntry } from "../src/server/request-log"; + +function makeEntry(requestId: string, timestamp: number): RequestLogEntry { + return { + requestId, + timestamp, + model: "test-model", + provider: "test-provider", + status: 200, + durationMs: 100, + usageStatus: "reported", + }; +} + +describe("request-log cursor", () => { + test("round-trips timestamp and request id", () => { + const encoded = encodeRequestLogCursor({ timestamp: 1700000000000, requestId: "req-2" }); + expect(decodeRequestLogCursor(encoded)).toEqual({ timestamp: 1700000000000, requestId: "req-2" }); + }); + + test("rejects malformed, empty, and type-confused payloads", () => { + expect(decodeRequestLogCursor("")).toBeNull(); + expect(decodeRequestLogCursor("not-base64-json")).toBeNull(); + expect(decodeRequestLogCursor("!!!")).toBeNull(); + expect(decodeRequestLogCursor(Buffer.from(JSON.stringify({ t: "1", id: 2 })).toString("base64url"))).toBeNull(); + expect(decodeRequestLogCursor(Buffer.from(JSON.stringify({ v: 2, t: 1, id: "req" })).toString("base64url"))).toBeNull(); + expect(decodeRequestLogCursor(Buffer.from(JSON.stringify({ v: 1, t: -1, id: "req" })).toString("base64url"))).toBeNull(); + expect(decodeRequestLogCursor(Buffer.from(JSON.stringify({ v: 1, t: Number.NaN, id: "req" })).toString("base64url"))).toBeNull(); + expect(decodeRequestLogCursor(Buffer.from(JSON.stringify({ v: 1, t: 1, id: "" })).toString("base64url"))).toBeNull(); + expect(decodeRequestLogCursor(Buffer.from(JSON.stringify({ v: 1, t: 1, id: "a".repeat(300) })).toString("base64url"))).toBeNull(); + expect(decodeRequestLogCursor(Buffer.from(JSON.stringify([1, 2, 3])).toString("base64url"))).toBeNull(); + expect(decodeRequestLogCursor(Buffer.from(JSON.stringify(null)).toString("base64url"))).toBeNull(); + }); + + test("ring slice returns only rows after a live cursor", () => { + const rows = [makeEntry("req-1", 1), makeEntry("req-2", 2), makeEntry("req-3", 3)]; + expect(sliceRequestLogsAfterCursor(rows, { timestamp: 2, requestId: "req-2" })).toEqual({ + entries: [rows[2]], + reset: false, + }); + }); + + test("ring slice returns empty entries when cursor matches the newest item", () => { + const rows = [makeEntry("req-1", 1), makeEntry("req-2", 2), makeEntry("req-3", 3)]; + expect(sliceRequestLogsAfterCursor(rows, { timestamp: 3, requestId: "req-3" })).toEqual({ + entries: [], + reset: false, + }); + }); + + test("ring slice requests a reset when the cursor was evicted or not found", () => { + const rows = [makeEntry("req-3", 3), makeEntry("req-4", 4)]; + expect(sliceRequestLogsAfterCursor(rows, { timestamp: 2, requestId: "req-2" })).toEqual({ + entries: rows, + reset: true, + }); + }); + + test("ring slice handles empty logs array with reset", () => { + expect(sliceRequestLogsAfterCursor([], { timestamp: 1, requestId: "req-1" })).toEqual({ + entries: [], + reset: true, + }); + }); +});