diff --git a/internal/cloud/server-agent-observer-http.ts b/internal/cloud/server-agent-observer-http.ts index 78cc272d..065f31c5 100644 --- a/internal/cloud/server-agent-observer-http.ts +++ b/internal/cloud/server-agent-observer-http.ts @@ -25,7 +25,6 @@ export async function handleAgentObserverHttp(request, response, options = {}) { } const requestUrl = new URL(request.url ?? "/v1/agent-observer/events", "http://agent-observer.local"); const resumeAfterId = text(requestUrl.searchParams.get("afterId") ?? request.headers?.["last-event-id"]); - const liveOnly = requestUrl.searchParams.get("mode") === "live"; try { await (bridge.liveReady ?? bridge.ready); @@ -76,20 +75,7 @@ export async function handleAgentObserverHttp(request, response, options = {}) { if (!response.writableEnded) response.end(); }; - if (liveOnly) { - const unsubscribe = bridge.subscribeLiveHwlabEvents((envelope) => { void enqueue("hwlab.event.v1", envelope); }); - cleanup.push(() => unsubscribe?.()); - await enqueue("agent-observer.connected", { - type: "connected", - status: "live", - deliverySemantics: "agentrun-snapshot-then-kafka-live", - authority: "hwlab.event.v1", - replay: null, - serverSentAt: new Date().toISOString(), - valuesPrinted: false - }); - } else { - const handoff = createWorkbenchKafkaRefreshHandoff({ + const handoff = createWorkbenchKafkaRefreshHandoff({ allowUnscoped: true, resumeAfterId, liveBufferLimit: policy.liveBufferLimit, @@ -123,15 +109,14 @@ export async function handleAgentObserverHttp(request, response, options = {}) { valuesPrinted: false }), onFailure: fail - }); - cleanup.push(() => handoff.stop("connection-closed")); + }); + cleanup.push(() => handoff.stop("connection-closed")); - try { - await handoff.start(); - } catch (error) { - await fail(error); - return; - } + try { + await handoff.start(); + } catch (error) { + await fail(error); + return; } if (!active()) return; diff --git a/internal/cloud/server-agent-observer-runs-http.ts b/internal/cloud/server-agent-observer-runs-http.ts deleted file mode 100644 index 30dbf1f1..00000000 --- a/internal/cloud/server-agent-observer-runs-http.ts +++ /dev/null @@ -1,102 +0,0 @@ -// Responsibility: expose the AgentRun manager durable run collection through the authenticated HWLAB origin. - -const DEFAULT_AGENTRUN_MGR_URL = "http://agentrun-mgr.agentrun-v01.svc.cluster.local:8080"; -const AGENTRUN_MANAGER_HOST_PATTERN = /^agentrun-mgr\.[a-z0-9]([-a-z0-9]*[a-z0-9])?\.svc\.cluster\.local$/u; - -export async function handleAgentObserverRunsHttp(request, response, url, options = {}) { - if (request.method !== "GET") { - response.writeHead(405, { allow: "GET", "content-type": "application/json; charset=utf-8" }); - response.end(JSON.stringify({ ok: false, error: { code: "method_not_allowed", message: "GET required." } })); - return; - } - - const env = options.env ?? process.env; - const managerUrl = resolveManagerUrl(env); - const apiKey = text(env.AGENTRUN_API_KEY) || text(env.HWLAB_CODE_AGENT_AGENTRUN_API_KEY) || text(env.HWLAB_PROVIDER_PROFILE_AGENTRUN_API_KEY) || text(env.HWLAB_API_KEY); - const timeoutMs = positiveInteger(env.HWLAB_CODE_AGENT_AGENTRUN_HTTP_TIMEOUT_MS, 30_000); - const target = new URL("/api/v1/runs", managerUrl); - target.search = url.search; - const startedAt = Date.now(); - const controller = new AbortController(); - const timer = setTimeout(() => controller.abort(), timeoutMs); - - try { - const headers: Record = { - accept: "application/json", - "x-source-service-id": "hwlab-cloud-api", - "x-hwlab-delegation": "agent-observer-runs-v1", - }; - if (apiKey) headers.authorization = `Bearer ${apiKey}`; - const delegated = await (options.fetchImpl ?? fetch)(target, { method: "GET", headers, signal: controller.signal }); - const body = await delegated.text(); - const payload = parseJson(body); - const data = payload?.data ?? payload; - if (!delegated.ok || payload?.ok === false || !data || typeof data !== "object") { - sendJson(response, delegated.status || 502, { - ok: false, - error: { - code: text(payload?.error?.code) || text(payload?.failureKind) || "agentrun_run_collection_failed", - message: text(payload?.error?.message) || text(payload?.message) || `AgentRun run collection returned HTTP ${delegated.status}.`, - httpStatus: delegated.status, - elapsedMs: Date.now() - startedAt, - valuesPrinted: false, - }, - }); - return; - } - sendJson(response, 200, { - ok: true, - status: "ok", - authority: "agentrun-manager-durable-store", - ...data, - elapsedMs: Date.now() - startedAt, - valuesPrinted: false, - }); - } catch (error) { - const timedOut = controller.signal.aborted || error?.name === "AbortError"; - sendJson(response, timedOut ? 504 : 502, { - ok: false, - error: { - code: timedOut ? "agentrun_run_collection_timeout" : "agentrun_run_collection_unavailable", - message: timedOut ? `AgentRun run collection exceeded ${timeoutMs}ms.` : (error instanceof Error ? error.message : String(error)), - elapsedMs: Date.now() - startedAt, - valuesPrinted: false, - }, - }); - } finally { - clearTimeout(timer); - } -} - -function resolveManagerUrl(env) { - const raw = text(env.AGENTRUN_MGR_URL) || text(env.HWLAB_CODE_AGENT_AGENTRUN_MGR_URL) || text(env.AGENTRUN_MANAGER_URL) || DEFAULT_AGENTRUN_MGR_URL; - const parsed = new URL(raw); - const allowLocal = truthy(env.HWLAB_CODE_AGENT_AGENTRUN_ALLOW_NON_K3S_URL); - const localAllowed = allowLocal && ["127.0.0.1", "localhost"].includes(parsed.hostname); - if (!AGENTRUN_MANAGER_HOST_PATTERN.test(parsed.hostname) && !localAllowed) { - throw Object.assign(new Error("AgentRun manager URL must use the declared internal k3s service."), { code: "agentrun_manager_url_invalid" }); - } - return raw.replace(/\/+$/u, ""); -} - -function positiveInteger(value, fallback) { - const number = Number(value); - return Number.isInteger(number) && number > 0 ? number : fallback; -} - -function truthy(value) { - return ["1", "true", "yes", "on"].includes(String(value ?? "").trim().toLowerCase()); -} - -function text(value) { - return typeof value === "string" && value.trim() ? value.trim() : ""; -} - -function parseJson(value) { - try { return value ? JSON.parse(value) : null; } catch { return null; } -} - -function sendJson(response, status, payload) { - response.writeHead(status, { "content-type": "application/json; charset=utf-8", "cache-control": "no-store" }); - response.end(JSON.stringify(payload)); -} diff --git a/internal/cloud/server.ts b/internal/cloud/server.ts index 1a2080cb..9eafa5f5 100644 --- a/internal/cloud/server.ts +++ b/internal/cloud/server.ts @@ -73,7 +73,6 @@ import { handleWorkbenchDebugFakeSseHttp } from "./workbench-debug-fake-sse.ts"; import { handleWorkbenchKafkaSseDebugHttp } from "./workbench-kafka-sse-debug.ts"; import { handleWorkbenchLaunchHttp } from "./server-workbench-launch-http.ts"; import { handleAgentObserverHttp } from "./server-agent-observer-http.ts"; -import { handleAgentObserverRunsHttp } from "./server-agent-observer-runs-http.ts"; import { startWorkbenchEmptySessionGc } from "./workbench-empty-session-gc.ts"; import { startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts"; import { handleM3IoControlHttp } from "./server-m3-http.ts"; @@ -773,11 +772,6 @@ async function handleRestAdapter(request, response, url, options) { return; } - if (url.pathname === "/v1/agent-observer/runs") { - await handleAgentObserverRunsHttp(request, response, url, options); - return; - } - if (url.pathname === "/v1/workbench/launches") { await handleWorkbenchLaunchHttp(request, response, options); return; @@ -1063,7 +1057,7 @@ function navIdForRestPath(pathname, method = "GET") { if (pathname === "/v1/users/me/profile" || pathname === "/v1/users/me/password") return "system.settings"; if (pathname === "/v1/workbench/debug/fake-sse" || pathname.startsWith("/v1/workbench/debug/fake-sse/") || pathname === "/v1/workbench/debug/kafka-sse" || pathname.startsWith("/v1/workbench/debug/kafka-sse/")) return "workbench.debug"; if (pathname === "/v1/workbench/commands" || pathname === "/v1/workbench/events" || pathname === "/v1/workbench/projection-events" || pathname === "/v1/workbench/launches" || pathname === "/v1/workbench/sessions" || pathname.startsWith("/v1/workbench/sessions/") || pathname.startsWith("/v1/workbench/turns/") || pathname.startsWith("/v1/workbench/traces/")) return "workbench.code"; - if (pathname === "/v1/agent-observer/events" || pathname === "/v1/agent-observer/runs") return "workbench.code"; + if (pathname === "/v1/agent-observer/events") return "workbench.code"; if (pathname === "/v1/agent/chat" || pathname === "/v1/agent/sessions" || pathname.startsWith("/v1/agent/sessions/") || pathname === "/v1/agent/chat/inspect" || pathname.startsWith("/v1/agent/chat/result/") || pathname.startsWith("/v1/agent/turns/") || pathname.startsWith("/v1/agent/traces/") || pathname === "/v1/agent/chat/cancel" || pathname === "/v1/agent/chat/steer") return "workbench.code"; if (pathname === "/v1/admin/provider-profiles" || pathname.startsWith("/v1/admin/provider-profiles/")) return "admin.providerProfiles"; if (pathname === "/v1/admin/secrets" || pathname.startsWith("/v1/admin/secrets/")) return "admin.secrets"; diff --git a/tools/src/hwlab-cli-lib.ts b/tools/src/hwlab-cli-lib.ts index 411e4c7c..bf3c3d54 100644 --- a/tools/src/hwlab-cli-lib.ts +++ b/tools/src/hwlab-cli-lib.ts @@ -2293,26 +2293,20 @@ async function agentObserverCommand(context: any) { const apiKey = explicitApiKey(context.parsed, context.env); if (!apiKey) throw cliError("hwlab_api_key_required", "client agent-observer list requires HWLAB_API_KEY.", { credentialSource: "HWLAB_API_KEY" }); const timeoutMs = numberOption(context.parsed.timeoutMs) ?? DEFAULT_TIMEOUT_MS; - const page = numberOption(context.parsed.page) ?? 1; - const pageSize = numberOption(context.parsed.pageSize) ?? numberOption(context.parsed.limit) ?? 25; - if (!Number.isInteger(page) || page <= 0) throw cliError("invalid_page", "--page must be a positive integer.", { page }); - if (!Number.isInteger(pageSize) || pageSize <= 0 || pageSize > 100) throw cliError("invalid_page_size", "--page-size must be between 1 and 100.", { pageSize }); - const query = new URLSearchParams({ page: String(page), pageSize: String(pageSize) }); - for (const [option, parameter] of [["status", "status"], ["projectId", "projectId"], ["workspace", "workspace"], ["backendProfile", "backendProfile"], ["providerId", "providerId"], ["search", "search"], ["updatedAfter", "updatedAfter"], ["sort", "sort"]] as const) { - const value = text(context.parsed[option]); - if (value) query.set(parameter, value); - } + const limit = numberOption(context.parsed.limit) ?? 100; + if (!Number.isInteger(limit) || limit <= 0) throw cliError("invalid_limit", "--limit must be a positive integer.", { limit }); const observation = await collectAgentObserverRuns({ fetchImpl: context.fetchImpl, - url: `${endpoint.baseUrl}/v1/agent-observer/runs?${query.toString()}`, + url: `${endpoint.baseUrl}/v1/agent-observer/events`, apiKey, timeoutMs, + limit }); return ok("client.agent-observer.list", { baseUrl: endpoint.baseUrl, runtimeEndpoint: runtimeEndpointVisibility(endpoint), - route: route("GET", "/v1/agent-observer/runs"), - authority: "agentrun-manager-durable-store", + route: route("GET", "/v1/agent-observer/events"), + authority: "hwlab.event.v1", ...observation, valuesPrinted: false }); @@ -2322,9 +2316,9 @@ function agentObserverHelp() { return ok("client.agent-observer.help", { serviceRuntime: false, imagePublished: false, - commands: ["list [--page 1] [--page-size 25] [--status ] [--project-id ] [--workspace ] [--backend-profile ] [--provider-id ] [--search ] [--updated-after ] [--sort updated-desc|updated-asc|status-asc|id-asc] [--timeout-ms 30000]"], + commands: ["list [--limit 100] [--timeout-ms 30000]"], auth: "HWLAB_API_KEY", - authority: "agentrun-manager-durable-store" + authority: "hwlab.event.v1" }); } diff --git a/tools/src/hwlab-cli/agent-observer-client.ts b/tools/src/hwlab-cli/agent-observer-client.ts index 9495cbe8..f80358a7 100644 --- a/tools/src/hwlab-cli/agent-observer-client.ts +++ b/tools/src/hwlab-cli/agent-observer-client.ts @@ -1,54 +1,151 @@ type FetchLike = typeof fetch; -export async function collectAgentObserverRuns({ fetchImpl, url, apiKey, timeoutMs }: { +type ObserverRun = { + runId: string; + taskId: string | null; + commandId: string | null; + sessionId: string | null; + projectId: string | null; + workspace: string | null; + lane: string | null; + backend: string | null; + providerId: string | null; + status: string | null; + lastEventAt: string | null; +}; + +export async function collectAgentObserverRuns({ fetchImpl, url, apiKey, timeoutMs, limit }: { fetchImpl: FetchLike; url: string; apiKey: string; timeoutMs: number; + limit: number; }) { const controller = new AbortController(); const timer = setTimeout(() => controller.abort(), timeoutMs); - const startedAt = Date.now(); + const runs = new Map(); + let replay: unknown = null; + let phase = "replay"; + let eventCount = 0; try { const response = await fetchImpl(url, { method: "GET", - headers: { accept: "application/json", authorization: `Bearer ${apiKey}` }, - signal: controller.signal, + headers: { accept: "text/event-stream", authorization: `Bearer ${apiKey}` }, + signal: controller.signal }); - const body = await response.text(); - let payload: any = null; - try { payload = body ? JSON.parse(body) : null; } catch { /* handled below */ } - if (!response.ok || payload?.ok === false) { - throw observerError( - text(payload?.error?.code) ?? "agent_observer_request_failed", - text(payload?.error?.message) ?? `GET /v1/agent-observer/runs returned HTTP ${response.status}.`, - { httpStatus: response.status, bodyBytes: Buffer.byteLength(body) }, - ); + if (!response.ok) { + const body = await response.text(); + throw observerError("agent_observer_request_failed", `GET /v1/agent-observer/events returned HTTP ${response.status}.`, { httpStatus: response.status, bodyBytes: Buffer.byteLength(body) }); } - if (!payload || !Array.isArray(payload.items)) { - throw observerError("agent_observer_response_invalid", "Agent observer run collection did not return an items array."); + const reader = response.body?.getReader(); + if (!reader) throw observerError("agent_observer_stream_missing", "Agent observer response did not provide an SSE body."); + const decoder = new TextDecoder(); + let buffer = ""; + let connected = false; + while (!connected) { + const chunk = await reader.read(); + if (chunk.done) break; + buffer += decoder.decode(chunk.value, { stream: true }); + const frames = splitSseFrames(buffer); + buffer = frames.remainder; + for (const frame of frames.complete) { + const parsed = parseSseFrame(frame); + if (!parsed) continue; + if (parsed.event === "hwlab.event.v1") { + eventCount += 1; + mergeRun(runs, parsed.data); + } else if (parsed.event === "agent-observer.handoff") { + phase = "handoff"; + replay = objectValue(parsed.data)?.replay ?? replay; + } else if (parsed.event === "agent-observer.connected") { + phase = "live"; + replay = objectValue(parsed.data)?.replay ?? replay; + connected = true; + break; + } else if (parsed.event === "agent-observer.error") { + const error = objectValue(parsed.data); + throw observerError(text(error?.code) || "agent_observer_stream_failed", text(error?.message) || "Agent observer SSE reported an error."); + } + } } - return { - ...payload, - elapsedMs: Date.now() - startedAt, - returnedRunCount: payload.items.length, - truncated: payload.hasNextPage === true, - valuesPrinted: false, - }; + await reader.cancel().catch(() => undefined); + if (!connected) throw observerError("agent_observer_handoff_incomplete", "Agent observer SSE ended before retention-to-live handoff completed."); } catch (error: any) { - if (controller.signal.aborted && error?.code === undefined) { - throw observerError("agent_observer_timeout", `Agent observer run collection exceeded ${timeoutMs}ms.`, { timeoutMs }); - } + if (controller.signal.aborted && error?.code === undefined) throw observerError("agent_observer_timeout", `Agent observer did not complete replay within ${timeoutMs}ms.`, { timeoutMs }); throw error; } finally { clearTimeout(timer); + controller.abort(); } + + const allRuns = [...runs.values()].sort((left, right) => String(right.lastEventAt ?? "").localeCompare(String(left.lastEventAt ?? ""))); + return { + phase, + replay, + eventCount, + runCount: allRuns.length, + returnedRunCount: Math.min(allRuns.length, limit), + truncated: allRuns.length > limit, + runs: allRuns.slice(0, limit), + valuesPrinted: false + }; } -function text(value: unknown): string | null { +function mergeRun(runs: Map, value: unknown) { + const envelope = objectValue(value); + const event = objectValue(envelope?.event); + const context = objectValue(envelope?.context); + const runId = firstText(envelope?.runId, context?.runId, event?.runId); + if (!runId) return; + const previous = runs.get(runId); + runs.set(runId, { + runId, + taskId: firstText(event?.taskId, previous?.taskId), + commandId: firstText(envelope?.commandId, context?.commandId, event?.commandId, previous?.commandId), + sessionId: firstText(envelope?.sessionId, envelope?.hwlabSessionId, event?.sessionId, previous?.sessionId), + projectId: firstText(context?.projectId, event?.projectId, previous?.projectId), + workspace: firstText(event?.workspace, previous?.workspace), + lane: firstText(event?.lane, previous?.lane), + backend: firstText(event?.backend, previous?.backend), + providerId: firstText(event?.providerId, previous?.providerId), + status: firstText(event?.runStatus, event?.status, previous?.status), + lastEventAt: firstText(event?.createdAt, envelope?.producedAt, previous?.lastEventAt) + }); +} + +function splitSseFrames(value: string) { + const normalized = value.replace(/\r\n/gu, "\n"); + const parts = normalized.split("\n\n"); + return { complete: parts.slice(0, -1), remainder: parts.at(-1) ?? "" }; +} + +function parseSseFrame(frame: string) { + let event = "message"; + const data = []; + for (const line of frame.split("\n")) { + if (line.startsWith("event:")) event = line.slice(6).trim(); + if (line.startsWith("data:")) data.push(line.slice(5).trimStart()); + } + if (data.length === 0) return null; + try { return { event, data: JSON.parse(data.join("\n")) }; } catch { return null; } +} + +function objectValue(value: any): any { + return value && typeof value === "object" && !Array.isArray(value) ? value : null; +} + +function text(value: unknown) { return typeof value === "string" && value.trim() ? value.trim() : null; } +function firstText(...values: unknown[]) { + for (const value of values) { + const result = text(value); + if (result) return result; + } + return null; +} + function observerError(code: string, message: string, details: Record = {}) { return Object.assign(new Error(message), { code, details: { ...details, valuesPrinted: false } }); } diff --git a/web/hwlab-cloud-web/src/api/agentObserver.ts b/web/hwlab-cloud-web/src/api/agentObserver.ts index a54dff26..a8c17c2f 100644 --- a/web/hwlab-cloud-web/src/api/agentObserver.ts +++ b/web/hwlab-cloud-web/src/api/agentObserver.ts @@ -1,50 +1,5 @@ -import { fetchJson } from "@/api/client"; - -export interface AgentObserverRunPageQuery { - page: number; - pageSize: number; - status?: string; - projectId?: string; - workspace?: string; - backendProfile?: string; - providerId?: string; - search?: string; - updatedAfter?: string; - sort?: string; -} - -export interface AgentObserverRunPagePayload { - items: Array>; - page: number; - pageSize: number; - total: number; - totalPages: number; - hasNextPage: boolean; - facets: Record; - authority: string; - elapsedMs: number; - valuesPrinted: false; -} - -export async function fetchAgentObserverRuns(query: AgentObserverRunPageQuery) { - const search = new URLSearchParams({ page: String(query.page), pageSize: String(query.pageSize) }); - if (query.status) search.set("status", query.status); - if (query.projectId) search.set("projectId", query.projectId); - if (query.workspace) search.set("workspace", query.workspace); - if (query.backendProfile) search.set("backendProfile", query.backendProfile); - if (query.providerId) search.set("providerId", query.providerId); - if (query.search) search.set("search", query.search); - if (query.updatedAfter) search.set("updatedAfter", query.updatedAfter); - if (query.sort) search.set("sort", query.sort); - return await fetchJson(`/v1/agent-observer/runs?${search.toString()}`, { - timeoutMs: 10_000, - timeoutName: "AgentRun 分页快照", - }); -} - export interface AgentObserverStreamCallbacks { afterId?: string | null; - liveOnly?: boolean; onEnvelope: (payload: unknown, eventId: string | null) => void; onHandoff: (payload: unknown) => void; onConnected: (payload: unknown) => void; @@ -56,7 +11,6 @@ export interface AgentObserverStreamCallbacks { export function openAgentObserverStream(callbacks: AgentObserverStreamCallbacks): EventSource { const url = new URL("/v1/agent-observer/events", window.location.origin); if (callbacks.afterId) url.searchParams.set("afterId", callbacks.afterId); - if (callbacks.liveOnly) url.searchParams.set("mode", "live"); const source = new EventSource(`${url.pathname}${url.search}`, { withCredentials: true }); source.addEventListener("hwlab.event.v1", (event) => callbacks.onEnvelope(parseEvent(event), event instanceof MessageEvent ? event.lastEventId || null : null)); source.addEventListener("agent-observer.handoff", (event) => callbacks.onHandoff(parseEvent(event))); diff --git a/web/hwlab-cloud-web/src/stores/agent-observer.ts b/web/hwlab-cloud-web/src/stores/agent-observer.ts index bee8376f..2261168b 100644 --- a/web/hwlab-cloud-web/src/stores/agent-observer.ts +++ b/web/hwlab-cloud-web/src/stores/agent-observer.ts @@ -1,6 +1,6 @@ import { computed, ref } from "vue"; import { defineStore } from "pinia"; -import { fetchAgentObserverRuns, openAgentObserverStream } from "@/api/agentObserver"; +import { openAgentObserverStream } from "@/api/agentObserver"; import { agentObserverRuntimeConfig, type AgentObserverRuntimeConfig } from "@/config/agent-observer"; import { asRecord, nonEmptyString } from "@/utils"; @@ -114,16 +114,6 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { const lastConfirmedEventId = ref(null); const retentionWindowAdvanced = ref(false); const error = ref(null); - const page = ref(1); - const pageSize = ref(25); - const total = ref(0); - const totalPages = ref(0); - const facetProjects = ref([]); - const facetWorkspaces = ref([]); - const facetBackends = ref([]); - const facetProviders = ref([]); - const snapshotLoaded = ref(false); - const loadingPage = ref(false); const nowMs = ref(Date.now()); const config = ref(null); const expandedRelationIds = ref([]); @@ -134,23 +124,30 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { let source: EventSource | null = null; let reconnectTimer: number | null = null; let staleTimer: number | null = null; - let snapshotRefreshTimer: number | null = null; - let snapshotRequest = 0; let flushScheduled = false; let manuallyClosed = false; const runs = computed(() => orderedIds.value.map((id) => entities.value[id]).filter((item): item is AgentObserverRun => Boolean(item))); const selectedRun = computed(() => selectedRunId.value ? entities.value[selectedRunId.value] ?? null : null); - const visibleRuns = computed(() => runs.value); + const visibleRuns = computed(() => { + const needle = search.value.trim().toLowerCase(); + const result = runs.value.filter((run) => { + if (statusFilter.value !== "all" && statusGroup(run.runStatus) !== statusFilter.value) return false; + if (projectFilter.value !== "all" && !matchesScopeFilter(run, projectFilter.value)) return false; + if (agentFilter.value !== "all" && (run.backend ?? run.providerId) !== agentFilter.value) return false; + if (!insideTimeWindow(run.updatedAt ?? run.startedAt, timeWindow.value, nowMs.value)) return false; + if (!needle) return true; + return [run.id, run.taskId, run.commandId, run.sessionId, run.traceId, run.projectId, run.workspace, run.backend, run.providerId, run.errorCode, run.lastMessage] + .some((value) => value?.toLowerCase().includes(needle)); + }); + return result.sort((left, right) => compareRuns(left, right, sort.value)); + }); const projects = computed(() => unique(runs.value.map((run) => run.projectId))); const scopeOptions = computed(() => [ - ...facetProjects.value.map((value) => ({ value: `project:${value}`, label: `项目 · ${value}` })), - ...facetWorkspaces.value.map((value) => ({ value: `workspace:${value}`, label: `Workspace · ${value}` })) - ]); - const agents = computed(() => [ - ...facetBackends.value.map((value) => ({ value: `backend:${value}`, label: `Agent · ${value}` })), - ...facetProviders.value.map((value) => ({ value: `provider:${value}`, label: `Profile · ${value}` })), + ...unique(runs.value.map((run) => run.projectId)).map((value) => ({ value: `project:${value}`, label: `项目 · ${value}` })), + ...unique(runs.value.map((run) => run.workspace)).map((value) => ({ value: `workspace:${value}`, label: `Workspace · ${value}` })) ]); + const agents = computed(() => unique(runs.value.map((run) => run.backend ?? run.providerId))); const selectedEvents = computed(() => { const run = selectedRun.value; if (!run || !config.value) return []; @@ -160,7 +157,10 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { const selectedLogs = computed(() => selectedEvents.value.filter((event) => ["output", "tool", "error"].includes(event.type) || event.label.includes(":output:"))); const selectedInEventWindow = computed(() => { if (!selectedRunId.value) return true; - return Boolean(selectedRun.value); + const run = selectedRun.value; + if (!run) return false; + const retained = new Set(events.value.map((event) => event.id)); + return run.eventIds.some((id) => retained.has(id)); }); const stats = computed(() => { const values = visibleRuns.value; @@ -170,7 +170,7 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { capacityByRunner.set(run.runnerId, { used: run.runnerCapacityUsed, total: run.runnerCapacityTotal }); } return { - total: total.value, + total: values.length, running: values.filter((run) => statusGroup(run.runStatus) === "running").length, waiting: values.filter((run) => statusGroup(run.runStatus) === "waiting").length, failed: values.filter((run) => statusGroup(run.runStatus) === "failed").length, @@ -202,12 +202,10 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { return; } manuallyClosed = false; - phase.value = "reconnecting"; + phase.value = "replay"; error.value = null; - void loadPage(); source = openAgentObserverStream({ afterId: lastConfirmedEventId.value, - liveOnly: true, onEnvelope: enqueueEnvelope, onHandoff: () => { phase.value = "handoff"; @@ -247,68 +245,6 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { reconnectTimer = null; if (manual && staleTimer !== null) window.clearInterval(staleTimer); if (manual) staleTimer = null; - if (manual && snapshotRefreshTimer !== null) window.clearTimeout(snapshotRefreshTimer); - if (manual) snapshotRefreshTimer = null; - } - - async function loadPage(): Promise { - const request = ++snapshotRequest; - loadingPage.value = true; - if (statusFilter.value === "unknown") { - entities.value = {}; - orderedIds.value = []; - total.value = 0; - totalPages.value = 0; - snapshotLoaded.value = true; - loadingPage.value = false; - return; - } - const scope = splitFilter(projectFilter.value); - const agent = splitFilter(agentFilter.value); - const result = await fetchAgentObserverRuns({ - page: page.value, - pageSize: pageSize.value, - status: statusQuery(statusFilter.value), - projectId: scope.kind === "project" ? scope.value : undefined, - workspace: scope.kind === "workspace" ? scope.value : undefined, - backendProfile: agent.kind === "backend" ? agent.value : undefined, - providerId: agent.kind === "provider" ? agent.value : undefined, - search: search.value.trim() || undefined, - updatedAfter: updatedAfterQuery(timeWindow.value, nowMs.value), - sort: sort.value, - }); - if (request !== snapshotRequest) return; - loadingPage.value = false; - if (!result.ok || !result.data) { - error.value = result.error ?? "AgentRun 分页快照加载失败。"; - snapshotLoaded.value = true; - return; - } - const incoming = result.data.items - .map((item) => snapshotRun(item, nonEmptyString(item.id) ? entities.value[nonEmptyString(item.id) as string] : undefined)) - .filter((item): item is AgentObserverRun => Boolean(item)); - const next: Record = {}; - for (const item of incoming) next[item.id] = item; - entities.value = next; - orderedIds.value = incoming.map((item) => item.id); - total.value = result.data.total; - totalPages.value = result.data.totalPages; - const facets = asRecord(result.data.facets); - facetProjects.value = stringArray(facets?.projects); - facetWorkspaces.value = stringArray(facets?.workspaces); - facetBackends.value = stringArray(facets?.backends); - facetProviders.value = stringArray(facets?.providers); - snapshotLoaded.value = true; - lastEventAt.value = incoming.reduce((latest, item) => latestTimestamp(latest, item.updatedAt), lastEventAt.value); - error.value = null; - } - - function refreshSnapshotSoon(): void { - if (snapshotRefreshTimer !== null) return; - snapshotRefreshTimer = window.setTimeout(() => { - snapshotRefreshTimer = null; - void loadPage(); - }, 250); } function reconnect(): void { @@ -388,7 +324,6 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { } const previous = entities.value[runId]; - if (!previous && page.value !== 1) return true; const sequenceIsCurrent = !previous || (sourceSeq !== null && (previous.lastSourceSeq === null || sourceSeq > previous.lastSourceSeq)); lastEventAt.value = latestTimestamp(lastEventAt.value, createdAt) ?? new Date().toISOString(); @@ -434,7 +369,6 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { entities.value = { ...entities.value, [runId]: next }; orderedIds.value = [runId, ...orderedIds.value.filter((id) => id !== runId)]; trimEntities(); - refreshSnapshotSoon(); return true; } @@ -454,9 +388,8 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { } function trimEntities(): void { - const limit = Math.min(config.value?.reducerEntityLimit ?? pageSize.value, pageSize.value); - if (orderedIds.value.length <= limit) return; - const keep = orderedIds.value.slice(0, limit); + if (!config.value || orderedIds.value.length <= config.value.reducerEntityLimit) return; + const keep = orderedIds.value.slice(0, config.value.reducerEntityLimit); if (selectedRunId.value && !keep.includes(selectedRunId.value) && entities.value[selectedRunId.value]) keep[keep.length - 1] = selectedRunId.value; const next: Record = {}; for (const id of keep) if (entities.value[id]) next[id] = entities.value[id]; @@ -481,80 +414,13 @@ export const useAgentObserverStore = defineStore("agentObserver", () => { return { entities, orderedIds, events, runs, visibleRuns, selectedRun, selectedEvents, selectedLogs, selectedInEventWindow, selectedRunId, view, density, search, statusFilter, projectFilter, agentFilter, timeWindow, sort, phase, lastEventAt, - page, pageSize, total, totalPages, snapshotLoaded, loadingPage, connectedAt, lastConfirmedEventId, retentionWindowAdvanced, error, nowMs, config, projects, scopeOptions, agents, stats, dataAgeMs, relationRows, relationWindowLimited: computed(() => relationWindow.value.limited), relationNodeCount: computed(() => relationWindow.value.nodeCount), relationOmittedRunCount: computed(() => relationWindow.value.omittedRunCount), - connect, disconnect, reconnect, loadPage, selectRun, toggleRelation + connect, disconnect, reconnect, selectRun, toggleRelation }; }); -function snapshotRun(value: Record, previous?: AgentObserverRun): AgentObserverRun | null { - const id = nonEmptyString(value.id); - if (!id) return null; - const runStatus = nonEmptyString(value.runStatus) ?? "unknown"; - const updatedAt = nonEmptyString(value.updatedAt); - return { - id, - title: nonEmptyString(value.title), - taskId: nonEmptyString(value.taskId), - taskStatus: nonEmptyString(value.taskStatus) ?? runStatus, - commandId: nonEmptyString(value.commandId), - attemptId: nonEmptyString(value.attemptId), - attemptStatus: nonEmptyString(value.attemptStatus) ?? runStatus, - runnerId: nonEmptyString(value.runnerId), - runnerCapacityUsed: integer(value.runnerCapacityUsed), - runnerCapacityTotal: integer(value.runnerCapacityTotal), - nodeId: nonEmptyString(value.nodeId), - lane: nonEmptyString(value.lane), - sessionId: nonEmptyString(value.sessionId), - traceId: nonEmptyString(value.traceId), - projectId: nonEmptyString(value.projectId), - workspace: nonEmptyString(value.workspace), - backend: nonEmptyString(value.backend), - providerId: nonEmptyString(value.providerId), - runStatus, - commandState: nonEmptyString(value.commandState), - commandType: nonEmptyString(value.commandType), - phase: nonEmptyString(value.phase) ?? runStatus, - result: nonEmptyString(value.result) ?? nonEmptyString(value.terminalStatus), - replyAuthority: null, - final: ["completed", "failed", "blocked", "cancelled"].includes(runStatus), - errorCode: nonEmptyString(value.errorCode), - lastMessage: nonEmptyString(value.lastMessage), - startedAt: nonEmptyString(value.startedAt), - finishedAt: nonEmptyString(value.finishedAt), - updatedAt, - lastSourceSeq: previous?.lastSourceSeq ?? null, - eventIds: previous?.eventIds ?? [], - changedAtMs: previous?.updatedAt === updatedAt ? previous.changedAtMs : Date.now(), - }; -} - -function splitFilter(value: string): { kind: string | null; value: string | null } { - const separator = value.indexOf(":"); - if (value === "all" || separator < 1) return { kind: null, value: null }; - return { kind: value.slice(0, separator), value: value.slice(separator + 1) || null }; -} - -function statusQuery(value: string): string | undefined { - if (value === "running") return "claimed,running"; - if (value === "waiting") return "pending"; - if (value === "failed") return "failed,blocked"; - if (value === "terminal") return "completed,cancelled"; - if (value === "unknown") return "unknown"; - return undefined; -} - -function updatedAfterQuery(value: string, now: number): string | undefined { - const duration = value === "15m" ? 15 * 60_000 : value === "1h" ? 60 * 60_000 : value === "6h" ? 6 * 60 * 60_000 : 0; - return duration > 0 ? new Date(now - duration).toISOString() : undefined; -} - -function stringArray(value: unknown): string[] { - return Array.isArray(value) ? unique(value.map(nonEmptyString)) : []; -} - export function buildRelationWindow( runs: AgentObserverRun[], selectedRunId: string | null, @@ -767,10 +633,32 @@ function stableEnvelopeIdentity(payload: unknown): string | null { return nonEmptyString(envelope?.eventId) ?? nonEmptyString(envelope?.sourceEventId); } +function compareRuns(left: AgentObserverRun, right: AgentObserverRun, sort: string): number { + if (sort === "id-asc") return left.id.localeCompare(right.id); + if (sort === "status-asc") return left.runStatus.localeCompare(right.runStatus) || left.id.localeCompare(right.id); + const leftTime = Date.parse(left.updatedAt ?? left.startedAt ?? "") || 0; + const rightTime = Date.parse(right.updatedAt ?? right.startedAt ?? "") || 0; + return sort === "updated-asc" ? leftTime - rightTime : rightTime - leftTime; +} + function unique(values: Array): string[] { return [...new Set(values.filter((value): value is string => Boolean(value)))].sort(); } +function matchesScopeFilter(run: AgentObserverRun, filter: string): boolean { + if (filter.startsWith("project:")) return run.projectId === filter.slice("project:".length); + if (filter.startsWith("workspace:")) return run.workspace === filter.slice("workspace:".length); + return run.projectId === filter; +} + +function insideTimeWindow(value: string | null, window: string, nowMs: number): boolean { + if (window === "all") return true; + const duration = window === "15m" ? 15 * 60_000 : window === "1h" ? 60 * 60_000 : window === "6h" ? 6 * 60 * 60_000 : 0; + if (!duration || !value) return false; + const timestamp = Date.parse(value); + return Number.isFinite(timestamp) && timestamp >= nowMs - duration; +} + function explicitBoolean(value: unknown): boolean | null { return typeof value === "boolean" ? value : null; } diff --git a/web/hwlab-cloud-web/src/styles/agent-observer.css b/web/hwlab-cloud-web/src/styles/agent-observer.css index 9801c50d..ace2b642 100644 --- a/web/hwlab-cloud-web/src/styles/agent-observer.css +++ b/web/hwlab-cloud-web/src/styles/agent-observer.css @@ -47,8 +47,7 @@ .agent-observer-workspace { grid-row: -2 / -1; min-width: 0; min-height: 0; display: grid; grid-template-columns: minmax(0, 1fr); overflow: hidden; } .agent-observer-workspace.has-inspector { grid-template-columns: minmax(0, 1fr) minmax(340px, 420px); } -.agent-observer-content { min-width: 0; min-height: 0; overflow: hidden; background: var(--console-surface); display: grid; grid-template-rows: minmax(0, 1fr) auto; } -.agent-observer-content > .table-pagination { border-top: 1px solid var(--console-border); background: var(--console-surface); } +.agent-observer-content { min-width: 0; min-height: 0; overflow: hidden; background: var(--console-surface); } .agent-observer-skeleton { height: 100%; padding: 12px; display: grid; align-content: start; gap: 7px; } .agent-observer-skeleton span { height: 48px; border-radius: 3px; background: var(--console-surface-muted); } .agent-observer-empty { height: 100%; display: grid; place-content: center; justify-items: center; gap: 8px; padding: 24px; color: var(--console-muted); text-align: center; } diff --git a/web/hwlab-cloud-web/src/views/agents/AgentRunsView.vue b/web/hwlab-cloud-web/src/views/agents/AgentRunsView.vue index 3857f303..13b70c5e 100644 --- a/web/hwlab-cloud-web/src/views/agents/AgentRunsView.vue +++ b/web/hwlab-cloud-web/src/views/agents/AgentRunsView.vue @@ -18,7 +18,6 @@ import { } from "lucide-vue-next"; import AgentRunInspector from "@/components/agents/AgentRunInspector.vue"; import BaseDrawer from "@/components/common/BaseDrawer.vue"; -import Pagination from "@/components/common/Pagination.vue"; import StatusBadge from "@/components/common/StatusBadge.vue"; import VirtualGrid from "@/components/common/VirtualGrid.vue"; import VirtualList from "@/components/common/VirtualList.vue"; @@ -45,7 +44,7 @@ const treeRows = computed(() => observer.relationRows); const ageLabel = computed(() => formatAge(observer.dataAgeMs)); const streamLabel = computed(() => ({ idle: "未连接", replay: "回放", handoff: "交接", live: "Live", reconnecting: "重连中", stale: "陈旧", error: "传输错误" }[observer.phase])); const selectedMissing = computed(() => Boolean(observer.selectedRunId && !observer.selectedRun)); -const loading = computed(() => observer.loadingPage && !observer.snapshotLoaded); +const loading = computed(() => observer.runs.length === 0 && ["idle", "replay", "handoff"].includes(observer.phase)); const statusItems = computed(() => [ { id: "running", label: "运行中", value: observer.stats.running, tone: "info" }, { id: "waiting", label: "等待", value: observer.stats.waiting, tone: "warn" }, @@ -62,7 +61,6 @@ const statusItems = computed(() => [ { id: "total", label: "筛选结果", value: observer.stats.total }, ]); let inspectorMedia: MediaQueryList | null = null; -let filterTimer: number | null = null; onMounted(() => { applyRouteState(); @@ -74,28 +72,14 @@ onMounted(() => { onBeforeUnmount(() => { observer.disconnect(); - if (filterTimer !== null) window.clearTimeout(filterTimer); inspectorMedia?.removeEventListener("change", updateInspectorMode); }); watch(() => route.fullPath, applyRouteState); watch( - () => [observer.view, observer.density, observer.search, observer.statusFilter, observer.projectFilter, observer.agentFilter, observer.timeWindow, observer.sort, observer.page, observer.selectedRunId], + () => [observer.view, observer.density, observer.search, observer.statusFilter, observer.projectFilter, observer.agentFilter, observer.timeWindow, observer.sort, observer.selectedRunId], () => { if (!urlApplying.value) void syncUrl(); }, ); -watch( - () => [observer.search, observer.statusFilter, observer.projectFilter, observer.agentFilter, observer.timeWindow, observer.sort], - () => { - if (urlApplying.value) return; - const pageChanged = observer.page !== 1; - observer.page = 1; - if (pageChanged) return; - if (filterTimer !== null) window.clearTimeout(filterTimer); - filterTimer = window.setTimeout(() => { filterTimer = null; void observer.loadPage(); }, 250); - }, -); -watch(() => observer.page, () => { if (!urlApplying.value) void observer.loadPage(); }); - function updateInspectorMode(): void { compactInspector.value = inspectorMedia?.matches ?? false; } @@ -110,7 +94,6 @@ function applyRouteState(): void { observer.agentFilter = textQuery(route.query.agent) || "all"; observer.timeWindow = textQuery(route.query.window) || "all"; observer.sort = textQuery(route.query.sort) || "updated-desc"; - observer.page = positivePage(route.query.page); observer.selectRun(typeof route.params.runId === "string" ? route.params.runId : null); queueMicrotask(() => { urlApplying.value = false; }); } @@ -125,7 +108,6 @@ async function syncUrl(): Promise { agent: observer.agentFilter === "all" ? "" : observer.agentFilter, window: observer.timeWindow === "all" ? "" : observer.timeWindow, sort: observer.sort === "updated-desc" ? "" : observer.sort, - page: observer.page > 1 ? String(observer.page) : "", }); await router.replace(observer.selectedRunId ? { name: "AgentRunDetail", params: { runId: observer.selectedRunId }, query } @@ -165,12 +147,11 @@ function formatAge(value: number | null): string { if (value === null || !Number function parseView(value: unknown): AgentObserverView { return value === "cards" || value === "tree" ? value : "table"; } function textQuery(value: unknown): string { return typeof value === "string" ? value : ""; } function compactQuery(value: Record): Record { return Object.fromEntries(Object.entries(value).filter(([, item]) => item)); } -function positivePage(value: unknown): number { const page = Number(value); return Number.isInteger(page) && page > 0 ? page : 1; }