diff --git a/config/max-lines-baseline.txt b/config/max-lines-baseline.txt index caae6aff3c82..07dc3eb0d85d 100644 --- a/config/max-lines-baseline.txt +++ b/config/max-lines-baseline.txt @@ -896,7 +896,6 @@ src/gateway/tools-invoke-http.test.ts src/gateway/watch-node-http.ts src/gateway/worker-environments/bootstrap.ts src/gateway/worker-environments/inference-runtime.ts -src/gateway/worker-environments/live-events.ts src/gateway/worker-environments/service.test.ts src/gateway/worker-environments/service.ts src/gateway/worker-environments/store.ts diff --git a/src/gateway/worker-environments/live-event-projection.ts b/src/gateway/worker-environments/live-event-projection.ts new file mode 100644 index 000000000000..381a0452250e --- /dev/null +++ b/src/gateway/worker-environments/live-event-projection.ts @@ -0,0 +1,123 @@ +import type { WorkerLiveEventParams } from "../../../packages/gateway-protocol/src/schema/worker-admission.js"; +import { + capLiveExecResult, + sanitizeToolArgs, + sanitizeToolResult, +} from "../../agents/embedded-agent-subscribe.tools.js"; +import { normalizeToolName } from "../../agents/tool-policy.js"; +import { formatSqliteSessionFileMarker } from "../../config/sessions/sqlite-marker.js"; +import { createTrajectoryRuntimeRecorder } from "../../trajectory/runtime.js"; + +export type WorkerLiveTrajectoryTarget = { + agentId?: string; + sessionId: string; + sessionKey: string; + storePath: string; +}; + +export type WorkerLiveTrajectoryRecorder = ReturnType; + +export function prepareWorkerLiveEventData( + event: WorkerLiveEventParams["event"], +): Record { + const payload = structuredClone(event.payload) as Record; + if (event.kind !== "tool") { + return payload; + } + const toolName = normalizeToolName(event.payload.name); + payload.name = toolName; + if (event.payload.phase === "start") { + payload.args = sanitizeToolArgs(event.payload.args); + } else if (event.payload.phase === "update") { + const partialResult = sanitizeToolResult(event.payload.partialResult); + payload.partialResult = toolName === "exec" ? capLiveExecResult(partialResult) : partialResult; + } else { + const result = sanitizeToolResult(event.payload.result); + payload.result = toolName === "exec" ? capLiveExecResult(result) : result; + } + return payload; +} + +export function isDefinitiveWorkerTerminalEvent(event: WorkerLiveEventParams["event"]): boolean { + return ( + event.kind === "lifecycle" && + (event.payload.phase === "end" || + (event.payload.phase === "error" && + (event.payload.aborted === true || event.payload.fallbackExhaustedFailure === true))) + ); +} + +export function createWorkerLiveTrajectoryRecorder(params: { + runId: string; + target: WorkerLiveTrajectoryTarget; +}): WorkerLiveTrajectoryRecorder { + const agentId = params.target.agentId ?? "main"; + return createTrajectoryRuntimeRecorder({ + runId: params.runId, + sessionId: params.target.sessionId, + sessionKey: params.target.sessionKey, + sessionFile: formatSqliteSessionFileMarker({ + agentId, + sessionId: params.target.sessionId, + storePath: params.target.storePath, + }), + }); +} + +export function recordWorkerLiveTrajectoryEvent( + recorder: WorkerLiveTrajectoryRecorder, + event: WorkerLiveEventParams["event"], +): void { + if (!recorder) { + return; + } + const data = prepareWorkerLiveEventData(event); + let recorded = false; + if (event.kind === "tool") { + if (event.payload.phase === "start") { + recorder.recordEvent("tool.call", data); + recorded = true; + } else if (event.payload.phase === "result") { + recorder.recordEvent("tool.result", { + ...data, + success: !event.payload.isError, + }); + recorded = true; + } + } else if (event.kind === "approval") { + recorder.recordEvent(`approval.${event.payload.phase}`, data); + recorded = true; + } else if (event.kind === "lifecycle") { + if (event.payload.phase === "start") { + recorder.recordEvent("session.started", { ...data, backend: "cloud-worker" }); + recorded = true; + } else if (event.payload.phase === "fallback_step") { + recorder.recordEvent("model.fallback_step", data); + recorded = true; + } else if (event.payload.phase === "finishing") { + recorder.recordEvent("model.finishing", data); + recorded = true; + } else if ( + (event.payload.phase === "end" || event.payload.phase === "error") && + isDefinitiveWorkerTerminalEvent(event) + ) { + const failed = event.payload.phase === "error"; + const interrupted = event.payload.aborted === true; + recorder.recordEvent("model.completed", { + ...data, + ...(failed ? { promptError: event.payload.error } : {}), + }); + recorder.recordEvent("session.ended", { + ...data, + status: interrupted ? "interrupted" : failed ? "error" : "success", + }); + recorded = true; + } + } + if (!recorded) { + return; + } + // Live delivery is authoritative; trajectory diagnostics must never reject a + // worker event. SQLite flushing begins synchronously and failures stay isolated. + void recorder.flush().catch(() => undefined); +} diff --git a/src/gateway/worker-environments/live-events.test.ts b/src/gateway/worker-environments/live-events.test.ts index 4fa9ad3a6e82..2b7fc0419c42 100644 --- a/src/gateway/worker-environments/live-events.test.ts +++ b/src/gateway/worker-environments/live-events.test.ts @@ -18,6 +18,8 @@ import { sweepStaleRunContexts, type AgentEventRuntimePayload as Event, } from "../../infra/agent-events.js"; +import { closeOpenClawAgentDatabasesForTest } from "../../state/openclaw-agent-db.js"; +import { loadSqliteTrajectoryRuntimeEventRowsSync } from "../../trajectory/runtime-store.sqlite.js"; import type { WorkerConnectionIdentity as Identity } from "./connection-identity.js"; import { createWorkerLiveEventReceiver, @@ -126,9 +128,86 @@ describe("worker live events", () => { afterEach(async () => { unsubscribe?.(); rx.clear(); + closeOpenClawAgentDatabasesForTest(); await fs.rm(root, { recursive: true, force: true }); }); + it("persists cloud-worker progress for sessions tail", async () => { + const credential = ["trajectory", "credential", "secret"].join("-"); + ack(live(1, lifecycle({ phase: "start", startedAt: 100 }))); + ack( + live( + 2, + tool({ + phase: "start", + name: "write", + toolCallId: "call-write", + args: { path: "proof.txt", credential }, + }), + ), + ); + ack( + live( + 3, + tool({ + phase: "result", + name: "write", + toolCallId: "call-write", + isError: false, + result: { status: "written", credential }, + }), + ), + ); + const terminal = live(4, lifecycle({ phase: "end", startedAt: 100, endedAt: 200 })); + ack(terminal); + ack(terminal); + await Promise.resolve(); + + const rows = loadSqliteTrajectoryRuntimeEventRowsSync({ + agentId: "main", + sessionId: SID, + storePath: store, + }); + expect(rows.map((row) => row.event.type)).toEqual([ + "session.started", + "tool.call", + "tool.result", + "model.completed", + "session.ended", + ]); + expect(rows[2]?.event.data).toMatchObject({ name: "write", success: true }); + expect(rows[4]?.event.data).toMatchObject({ status: "success" }); + expect(JSON.stringify(rows)).not.toContain(credential); + }); + + it("records aborted cloud-worker terminals as interrupted", async () => { + const credential = ["lifecycle", "credential", "value"].join("-"); + ack(live(1, lifecycle({ phase: "start", startedAt: 100 }))); + ack( + live( + 2, + lifecycle({ + phase: "error", + startedAt: 100, + endedAt: 200, + aborted: true, + error: `cancelled after Bearer ${credential}`, + }), + ), + ); + + const rows = loadSqliteTrajectoryRuntimeEventRowsSync({ + agentId: "main", + sessionId: SID, + storePath: store, + }); + expect(rows.at(-1)?.event).toMatchObject({ + type: "session.ended", + data: { status: "interrupted" }, + }); + expect(JSON.stringify(rows)).not.toContain(credential); + }); + it("maps and sanitizes kinds", () => { const credential = ["fixture", "credential", "value"].join("-"); const output = (char: string, status: string) => ({ diff --git a/src/gateway/worker-environments/live-events.ts b/src/gateway/worker-environments/live-events.ts index 0336d12991f6..bf1f6682297b 100644 --- a/src/gateway/worker-environments/live-events.ts +++ b/src/gateway/worker-environments/live-events.ts @@ -4,12 +4,6 @@ import type { WorkerLiveEventParams, WorkerLiveEventResult, } from "../../../packages/gateway-protocol/src/schema/worker-admission.js"; -import { - capLiveExecResult, - sanitizeToolArgs, - sanitizeToolResult, -} from "../../agents/embedded-agent-subscribe.tools.js"; -import { normalizeToolName } from "../../agents/tool-policy.js"; import { onSessionIdentityMutation, type SessionIdentityMutation, @@ -25,6 +19,14 @@ import { releaseAgentRunContext, } from "../../infra/agent-events.js"; import type { WorkerConnectionIdentity } from "./connection-identity.js"; +import { + createWorkerLiveTrajectoryRecorder, + isDefinitiveWorkerTerminalEvent, + prepareWorkerLiveEventData, + recordWorkerLiveTrajectoryEvent, + type WorkerLiveTrajectoryRecorder, + type WorkerLiveTrajectoryTarget, +} from "./live-event-projection.js"; import { resolveWorkerSessionTarget } from "./session-target.js"; const DEFAULT_WINDOW_SIZE = 128; @@ -42,13 +44,10 @@ type OwnedLiveRun = { claimId: string; controlUiVisible: boolean; lifecycleGeneration: string; + trajectoryRecorder: WorkerLiveTrajectoryRecorder; }; -type LiveEventTarget = { - agentId?: string; - sessionId: string; - sessionKey: string; -}; +type LiveEventTarget = WorkerLiveTrajectoryTarget; type WorkerLiveSessionBinding = Readonly<{ environmentId: string; @@ -115,6 +114,7 @@ function resolveLiveEventTarget( ...(target.agentId ? { agentId: target.agentId } : {}), sessionId: target.sessionId, sessionKey: target.sessionKey, + storePath: target.storePath, }; } @@ -163,34 +163,6 @@ function matchesSessionIdentityMutation( ); } -function prepareLiveEventData(event: WorkerLiveEventParams["event"]): Record { - const payload = structuredClone(event.payload) as Record; - if (event.kind !== "tool") { - return payload; - } - const toolName = normalizeToolName(event.payload.name); - payload.name = toolName; - if (event.payload.phase === "start") { - payload.args = sanitizeToolArgs(event.payload.args); - } else if (event.payload.phase === "update") { - const partialResult = sanitizeToolResult(event.payload.partialResult); - payload.partialResult = toolName === "exec" ? capLiveExecResult(partialResult) : partialResult; - } else { - const result = sanitizeToolResult(event.payload.result); - payload.result = toolName === "exec" ? capLiveExecResult(result) : result; - } - return payload; -} - -function isDefinitiveTerminal(event: WorkerLiveEventParams["event"]): boolean { - return ( - event.kind === "lifecycle" && - (event.payload.phase === "end" || - (event.payload.phase === "error" && - (event.payload.aborted === true || event.payload.fallbackExhaustedFailure === true))) - ); -} - export function createWorkerLiveEventReceiver(options: WorkerLiveEventReceiverOptions) { const boundSessions = new Map(); const sessionBindings = new Map(); @@ -503,7 +475,7 @@ export function createWorkerLiveEventReceiver(options: WorkerLiveEventReceiverOp } const pendingRunId = pending.request.runId; if (countedRunIds.has(pendingRunId)) { - if (isDefinitiveTerminal(pending.request.event)) { + if (isDefinitiveWorkerTerminalEvent(pending.request.event)) { return true; } continue; @@ -614,6 +586,7 @@ export function createWorkerLiveEventReceiver(options: WorkerLiveEventReceiverOp claimId, controlUiVisible, lifecycleGeneration, + trajectoryRecorder: createWorkerLiveTrajectoryRecorder({ runId, target: window.target }), }; window.activeRuns.set(runId, claimed); return claimed; @@ -628,7 +601,7 @@ export function createWorkerLiveEventReceiver(options: WorkerLiveEventReceiverOp if ("ok" in owned) { return owned; } - const definitiveTerminal = isDefinitiveTerminal(request.event); + const definitiveTerminal = isDefinitiveWorkerTerminalEvent(request.event); if (definitiveTerminal) { // Emission runs synchronous listeners that can clear this claim reentrantly. // Fence first so terminal delivery cannot reopen the run ID. @@ -638,10 +611,11 @@ export function createWorkerLiveEventReceiver(options: WorkerLiveEventReceiverOp { runId: request.runId, stream: request.event.kind, - data: prepareLiveEventData(request.event), + data: prepareWorkerLiveEventData(request.event), }, owned.claimId, ); + recordWorkerLiveTrajectoryEvent(owned.trajectoryRecorder, request.event); // Gateway handler owns cleanup so detach can revoke deferred terminal delivery. return undefined; }; @@ -788,4 +762,3 @@ export function createWorkerLiveEventReceiver(options: WorkerLiveEventReceiverOp } export type WorkerLiveEventReceiver = ReturnType; -/* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */ diff --git a/src/gateway/worker-environments/session-target.ts b/src/gateway/worker-environments/session-target.ts index a42b62285c00..f4466312afc3 100644 --- a/src/gateway/worker-environments/session-target.ts +++ b/src/gateway/worker-environments/session-target.ts @@ -9,7 +9,7 @@ import { export type ResolvedWorkerSessionTarget = Omit< SessionTranscriptWriteScope, - "sessionId" | "sessionKey" + "sessionId" | "sessionKey" | "storePath" > & { sessionEntry: NonNullable>; sessionId: string; @@ -18,6 +18,7 @@ export type ResolvedWorkerSessionTarget = Omit< string, NonNullable> >; + storePath: string; }; export function resolveWorkerSessionTarget(