diff --git a/internal/cloud/kafka-event-bridge.test.ts b/internal/cloud/kafka-event-bridge.test.ts index 463072a5..907206a0 100644 --- a/internal/cloud/kafka-event-bridge.test.ts +++ b/internal/cloud/kafka-event-bridge.test.ts @@ -6,7 +6,7 @@ import { test } from "bun:test"; import { createCloudApiBunServer } from "./bun-server.ts"; import { buildCloudApiReadiness } from "./health-contract.ts"; -import { decodeCanonicalAgentRunKafkaMessage, kafkaEventBridgeConfig, projectAgentRunKafkaEventToHwlabEvent, relayHwlabKafkaOutboxOnce, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts"; +import { decodeCanonicalAgentRunKafkaMessage, kafkaEventBridgeConfig, projectAgentRunKafkaEventToHwlabEvent, publishAgentRunKafkaMessageLive, relayHwlabKafkaOutboxOnce, startHwlabKafkaEventBridge } from "./kafka-event-bridge.ts"; const PROJECTOR_ENV = Object.freeze({ HWLAB_KAFKA_DIRECT_PUBLISH_ENABLED: "false", @@ -396,6 +396,29 @@ test("live Kafka startup does not touch transactional v8 readiness", async () => await bridge.stop(); }); +test("live Kafka publishes Artificer runs without HWLAB trace lineage for the unscoped observer", async () => { + const published = []; + const message = canonicalMessage({ offset: "52", traceId: null, hwlabSessionId: null, sessionId: "sess_artificer_test" }); + const result = await publishAgentRunKafkaMessageLive({ + topic: "agentrun.event.v1", + partition: 0, + message + }, { + producer: { async send(input) { published.push(input); return [{ partition: 0, baseOffset: "7" }]; } }, + config: kafkaEventBridgeConfig(LIVE_ENV), + env: LIVE_ENV, + logger: null, + otelSpanEmitter() {} + }); + + assert.equal(result.published, true); + assert.equal(result.envelope.runId, "run_projector_test"); + assert.equal(result.envelope.sessionId, "sess_artificer_test"); + assert.equal(result.envelope.traceId, null); + assert.equal(published.length, 1); + assert.equal(published[0].topic, "hwlab.event.v1"); +}); + test("Cloud API starts with all six Kafka capabilities and without transactional v8 schema", async () => { let transactionalV8Reads = 0; const runtimeStore = new Proxy({}, { @@ -727,7 +750,7 @@ function kafkaProjectorHarness(overrides = {}, calls = []) { }; } -function canonicalMessage({ offset, traceId = "trc_projector_test", hwlabSessionId = "ses_projector_test" }) { +function canonicalMessage({ offset, traceId = "trc_projector_test", hwlabSessionId = "ses_projector_test", sessionId = hwlabSessionId }) { const eventId = `evt_projector_${offset}`; const runId = "run_projector_test"; return { @@ -740,10 +763,11 @@ function canonicalMessage({ offset, traceId = "trc_projector_test", hwlabSession eventId, sourceSeq: Number(offset) + 1, traceId, + sessionId, hwlabSessionId, runId, commandId: "cmd_projector_test", - run: { runId, status: "running", valuesPrinted: false }, + run: { runId, status: "running", sessionId, hwlabSessionId, valuesPrinted: false }, command: { commandId: "cmd_projector_test", runId, seq: 1, type: "prompt", state: "running", valuesPrinted: false }, event: { id: eventId, runId, seq: Number(offset) + 1, type: "assistant_message", payload: { traceId, hwlabSessionId, commandId: "cmd_projector_test", text: "projected answer" }, createdAt: "2026-07-10T10:00:00.000Z" }, valuesPrinted: false diff --git a/internal/cloud/kafka-event-bridge.ts b/internal/cloud/kafka-event-bridge.ts index 8689df09..56c7e60a 100644 --- a/internal/cloud/kafka-event-bridge.ts +++ b/internal/cloud/kafka-event-bridge.ts @@ -447,9 +447,6 @@ export async function publishAgentRunKafkaMessageLive(kafkaMessage = {}, { produ logWarn(logger, "hwlab-live-kafka-message-invalid", { topic: transport.sourceTopic, partition: transport.sourcePartition, offset: transport.sourceOffset, errorCode: decoded.error.code, valuesPrinted: false }); return { handled: true, invalid: true, published: false, valuesPrinted: false }; } - if (!stringValue(decoded.event.traceId) || !stringValue(decoded.event.hwlabSessionId)) { - return { handled: true, ignored: true, published: false, valuesPrinted: false }; - } const projected = projectAgentRunKafkaEventToHwlabEvent(decoded.event, { source: config.clientId, sourceTopic: transport.sourceTopic, diff --git a/scripts/gitops-render.test.ts b/scripts/gitops-render.test.ts index af3f13f6..93159497 100644 --- a/scripts/gitops-render.test.ts +++ b/scripts/gitops-render.test.ts @@ -702,6 +702,11 @@ test("v03 render keeps node identity as data instead of generated structure", as assert.equal(cloudApiEnv.get("HWLAB_KAFKA_AGENTRUN_EVENT_GROUP_ID"), "hwlab-v03-agentrun-event-direct-publish"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_PROJECTOR_GROUP_ID"), "hwlab-v03-agentrun-event-projector"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_HWLAB_EVENT_GROUP_ID"), "hwlab-v03-workbench-live-sse"); + assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_GROUP_PREFIX"), "hwlab-v03-agent-observer-replay"); + assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_TIMEOUT_MS"), "30000"); + assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_SCAN_LIMIT"), "1000000"); + assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_MATCHED_EVENT_LIMIT"), "10000"); + assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_LIVE_BUFFER_LIMIT"), "2000"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_HWLAB_DEBUG_EVENT_TOPIC"), "hwlab.event.debug.v1"); assert.equal(cloudApiEnv.get("HWLAB_KAFKA_HWLAB_DEBUG_GROUP_PREFIX"), "hwlab-v03-workbench-isolated-debug"); assert.equal(cloudApiEnv.get("HWLAB_WORKBENCH_KAFKA_DEBUG_REPLAY_LIMIT"), "200"); diff --git a/scripts/src/gitops-render/core.mjs b/scripts/src/gitops-render/core.mjs index 9a1a28eb..239b4cb5 100644 --- a/scripts/src/gitops-render/core.mjs +++ b/scripts/src/gitops-render/core.mjs @@ -1583,6 +1583,13 @@ function transformWorkloads({ workloads, deploy, catalog, source, sourceBranch = upsertEnv(container.env, "HWLAB_AGENT_OBSERVER_CHANGE_HIGHLIGHT_MS", String(agentObserverPolicy.changeHighlightMs)); upsertEnv(container.env, "HWLAB_AGENT_OBSERVER_STALE_AFTER_MS", String(agentObserverPolicy.staleAfterMs)); upsertEnv(container.env, "HWLAB_AGENT_OBSERVER_REDUCED_MOTION", String(agentObserverPolicy.reducedMotion)); + if (serviceId === "hwlab-cloud-api") { + upsertEnv(container.env, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_GROUP_PREFIX", String(agentObserverPolicy.retentionGroupPrefix)); + upsertEnv(container.env, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_TIMEOUT_MS", String(agentObserverPolicy.retentionTimeoutMs)); + upsertEnv(container.env, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_SCAN_LIMIT", String(agentObserverPolicy.retentionScanLimit)); + upsertEnv(container.env, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_MATCHED_EVENT_LIMIT", String(agentObserverPolicy.retentionEventLimit)); + upsertEnv(container.env, "HWLAB_WORKBENCH_KAFKA_REFRESH_REPLAY_LIVE_BUFFER_LIMIT", String(agentObserverPolicy.liveBufferLimit)); + } if (serviceId === "hwlab-cloud-web") { upsertEnv(container.env, "HWLAB_AGENT_OBSERVER_NODE_ID", runtimeNodeIdForProfile(deploy, profile, nodeId)); upsertEnv(container.env, "HWLAB_AGENT_OBSERVER_LANE", profile); diff --git a/tools/hwlab-cli/agent-observer.test.ts b/tools/hwlab-cli/agent-observer.test.ts new file mode 100644 index 00000000..22e6db14 --- /dev/null +++ b/tools/hwlab-cli/agent-observer.test.ts @@ -0,0 +1,39 @@ +import assert from "node:assert/strict"; +import { test } from "bun:test"; + +import { collectAgentObserverRuns } from "../src/hwlab-cli/agent-observer-client.ts"; + +test("agent observer CLI collects unscoped Artificer runs through retention-to-live SSE", async () => { + const envelope = { + schema: "hwlab.event.v1", + runId: "run_artificer_test", + sessionId: "sess_artificer_test", + producedAt: "2026-07-16T02:24:51.000Z", + context: { projectId: "pikainc/pikaoa", runId: "run_artificer_test" }, + event: { runId: "run_artificer_test", lane: "release", backend: "gpt-pika", runStatus: "running", createdAt: "2026-07-16T02:24:50.000Z" } + }; + const frames = [ + `event: hwlab.event.v1\ndata: ${JSON.stringify(envelope)}\n\n`, + `event: agent-observer.handoff\ndata: ${JSON.stringify({ replay: { scannedCount: 12, deliveredCount: 1 } })}\n\n`, + `event: agent-observer.connected\ndata: ${JSON.stringify({ replay: { scannedCount: 12, deliveredCount: 1 } })}\n\n` + ].join(""); + let authorization = ""; + const result = await collectAgentObserverRuns({ + fetchImpl: async (_url, init) => { + authorization = new Headers(init?.headers).get("authorization") ?? ""; + return new Response(frames, { status: 200, headers: { "content-type": "text/event-stream" } }); + }, + url: "https://lab-dev.hwpod.com/v1/agent-observer/events", + apiKey: "hwl_live_test_only", + timeoutMs: 1000, + limit: 20 + }); + + assert.equal(authorization, "Bearer hwl_live_test_only"); + assert.equal(result.phase, "live"); + assert.equal(result.eventCount, 1); + assert.equal(result.runCount, 1); + assert.equal(result.runs[0].runId, "run_artificer_test"); + assert.equal(result.runs[0].lane, "release"); + assert.equal(result.runs[0].backend, "gpt-pika"); +}); diff --git a/tools/src/hwlab-cli-lib.ts b/tools/src/hwlab-cli-lib.ts index 7a323ce8..a69d257f 100644 --- a/tools/src/hwlab-cli-lib.ts +++ b/tools/src/hwlab-cli-lib.ts @@ -5,6 +5,7 @@ import { createHash } from "node:crypto"; import { chmod, mkdir, readFile, readdir, rm, writeFile } from "node:fs/promises"; import path from "node:path"; import { caseCommand } from "./hwlab-caserun-lib.ts"; +import { collectAgentObserverRuns } from "./hwlab-cli/agent-observer-client.ts"; import { computeCodeAgentComposerState, isCodeAgentSessionUnusableStatus } from "./hwlab-cli/composer-policy.ts"; import { traceDisplayRows, traceNoiseEventCount } from "./hwlab-cli/trace-renderer.ts"; import { resolveRuntimeEndpoint, runtimeEndpointVisibility, sameRuntimeEndpointScope } from "./runtime-endpoint-resolver.ts"; @@ -86,6 +87,7 @@ async function clientCommand(context: any) { if (group === "runtime") return runtimeCommand(next); if (group === "gateway") return gatewayCommand(next); if (group === "agent") return agentCommand(next); + if (group === "agent-observer") return agentObserverCommand(next); if (group === "session" || group === "sessions") return sessionCommand(next); if (["harness", "harness-ops", "harness-opt"].includes(group)) return harnessCommand(next); if (group === "workbench") return workbenchCommand(next); @@ -105,6 +107,7 @@ function clientSubcommandHelp(group: string) { if (group === "runtime") return runtimeHelp(); if (group === "gateway") return gatewayHelp(); if (group === "agent") return agentHelp(); + if (group === "agent-observer") return agentObserverHelp(); if (group === "session" || group === "sessions") return sessionHelp(); if (["harness", "harness-ops", "harness-opt"].includes(group)) return harnessHelp(); if (group === "workbench") return workbenchHelp(); @@ -2280,6 +2283,43 @@ function agentSendWaitPolicy(traceId: string) { }; } +async function agentObserverCommand(context: any) { + const subcommand = context.rest[0] || "list"; + if (wantsHelp(context)) return agentObserverHelp(); + if (subcommand !== "list") throw cliError("unsupported_agent_observer_command", `unsupported agent-observer command: ${subcommand}`, { subcommand }); + const endpoint = runtimeEndpoint(context.parsed, context.env, "web"); + 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 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/events`, + apiKey, + timeoutMs, + limit + }); + return ok("client.agent-observer.list", { + baseUrl: endpoint.baseUrl, + runtimeEndpoint: runtimeEndpointVisibility(endpoint), + route: route("GET", "/v1/agent-observer/events"), + authority: "hwlab.event.v1", + ...observation, + valuesPrinted: false + }); +} + +function agentObserverHelp() { + return ok("client.agent-observer.help", { + serviceRuntime: false, + imagePublished: false, + commands: ["list [--limit 100] [--timeout-ms 30000]"], + auth: "HWLAB_API_KEY", + authority: "hwlab.event.v1" + }); +} + async function workbenchCommand(context: any) { const subcommand = context.rest[0] || "summary"; if (wantsHelp(context)) return workbenchHelp(); diff --git a/tools/src/hwlab-cli/agent-observer-client.ts b/tools/src/hwlab-cli/agent-observer-client.ts new file mode 100644 index 00000000..f80358a7 --- /dev/null +++ b/tools/src/hwlab-cli/agent-observer-client.ts @@ -0,0 +1,151 @@ +type FetchLike = typeof fetch; + +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 runs = new Map(); + let replay: unknown = null; + let phase = "replay"; + let eventCount = 0; + try { + const response = await fetchImpl(url, { + method: "GET", + headers: { accept: "text/event-stream", authorization: `Bearer ${apiKey}` }, + signal: controller.signal + }); + 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) }); + } + 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."); + } + } + } + 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 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 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 } }); +}