diff --git a/internal/cloud/code-agent-trace-store.ts b/internal/cloud/code-agent-trace-store.ts index 2f792cf4..2aa23472 100644 --- a/internal/cloud/code-agent-trace-store.ts +++ b/internal/cloud/code-agent-trace-store.ts @@ -1,5 +1,5 @@ /* - * SPEC: PJ2026-010403 API契约 draft-2026-06-17-r0; PJ2026-0102 Agent编排 draft-2026-06-17-r0; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0 + * SPEC: PJ2026-010403 API契约 draft-2026-06-17-r0; PJ2026-0102 Agent编排 draft-2026-06-17-r0; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0; PJ2026-0104010803 唯一投影 draft-2026-06-20-p1-zero-split-durable-realtime * 职责: Code Agent trace capture。实时内存 trace 可加速 SSE/详情面板,但 append 事件必须同步投影到 durable store。 */ import { randomUUID } from "node:crypto"; @@ -249,16 +249,22 @@ export function createDurableCodeAgentTraceStore({ traceStore = defaultCodeAgent function persistTraceEvent(traceEventStore, event, meta = {}, workbenchEventWriter = null) { try { - Promise.resolve(traceEventStore.writeAgentTraceEvent({ event }, { + const projectionContext = { traceId: event.traceId, agentSessionId: event.agentSessionId ?? event.sessionId ?? meta.agentSessionId ?? meta.sessionId, sessionId: event.sessionId ?? meta.sessionId - })).catch(() => {}); + }; + Promise.resolve(traceEventStore.writeAgentTraceEvent({ event }, projectionContext)).catch((error) => { + console.error(`[workbench-projection] durable trace event write failed: traceId=${event.traceId} sourceEventId=${event.sourceEventId ?? "none"} error=${error?.message ?? error}`); + }); if (typeof workbenchEventWriter === "function") { - Promise.resolve(workbenchEventWriter(event, meta)).catch(() => {}); + Promise.resolve(workbenchEventWriter(event, meta)).catch((error) => { + console.error(`[workbench-projection] workbench event writer failed: traceId=${event.traceId} error=${error?.message ?? error}`); + }); } - } catch { - // Trace capture must not fail the user turn when the durable projection is temporarily unavailable. + } catch (error) { + // Trace capture must not fail the user turn, but projection write errors must be visible for diagnostics. + console.error(`[workbench-projection] persistTraceEvent sync error: traceId=${event.traceId} error=${error?.message ?? error}`); } } diff --git a/internal/cloud/workbench-projection-writer.ts b/internal/cloud/workbench-projection-writer.ts index 13a2d959..32e368bc 100644 --- a/internal/cloud/workbench-projection-writer.ts +++ b/internal/cloud/workbench-projection-writer.ts @@ -1,5 +1,5 @@ /* - * SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-18-p0-unique-projection; PJ2026-010205 HWLAB接入 draft-2026-06-17-r0. + * SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-18-p0-unique-projection; draft-2026-06-20-p1-zero-split-durable-realtime; PJ2026-010205 HWLAB接入 draft-2026-06-17-r0. * 职责: WorkbenchProjectionWriter 组件入口。唯一封装 Code Agent/AgentRun facts 到 Workbench session projection facts 的持久化写入。 */ import { createHash } from "node:crypto"; diff --git a/internal/db/migrations/0001_cloud_core_skeleton.sql b/internal/db/migrations/0001_cloud_core_skeleton.sql index 76c86391..6b4efac5 100644 --- a/internal/db/migrations/0001_cloud_core_skeleton.sql +++ b/internal/db/migrations/0001_cloud_core_skeleton.sql @@ -294,6 +294,8 @@ CREATE TABLE IF NOT EXISTS workbench_trace_events ( ); CREATE INDEX IF NOT EXISTS idx_workbench_trace_events_trace_seq ON workbench_trace_events(trace_id, source_seq, projected_seq, id); CREATE INDEX IF NOT EXISTS idx_workbench_trace_events_session_seq ON workbench_trace_events(session_id, source_seq, projected_seq, id); +CREATE UNIQUE INDEX IF NOT EXISTS idx_workbench_trace_events_trace_source_event ON workbench_trace_events(trace_id, source_event_id) WHERE source_event_id IS NOT NULL; +CREATE UNIQUE INDEX IF NOT EXISTS idx_workbench_trace_events_trace_projected_seq ON workbench_trace_events(trace_id, projected_seq); CREATE TABLE IF NOT EXISTS workbench_projection_checkpoints ( trace_id TEXT PRIMARY KEY, @@ -316,6 +318,25 @@ CREATE TABLE IF NOT EXISTS workbench_projection_checkpoints ( CREATE INDEX IF NOT EXISTS idx_workbench_projection_checkpoints_status_updated ON workbench_projection_checkpoints(projection_status, updated_at DESC, trace_id); CREATE INDEX IF NOT EXISTS idx_workbench_projection_checkpoints_session_updated ON workbench_projection_checkpoints(session_id, updated_at DESC, trace_id); +CREATE TABLE IF NOT EXISTS workbench_projection_outbox ( + outbox_seq BIGSERIAL PRIMARY KEY, + trace_id TEXT NOT NULL, + session_id TEXT, + turn_id TEXT, + message_id TEXT, + projected_seq INTEGER NOT NULL DEFAULT 0, + source_seq INTEGER NOT NULL DEFAULT 0, + source_event_id TEXT, + commit_type TEXT NOT NULL DEFAULT 'event', + terminal BOOLEAN NOT NULL DEFAULT false, + sealed BOOLEAN NOT NULL DEFAULT false, + payload_json TEXT NOT NULL DEFAULT '{}', + created_at TEXT NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_workbench_projection_outbox_trace_seq ON workbench_projection_outbox(trace_id, outbox_seq); +CREATE INDEX IF NOT EXISTS idx_workbench_projection_outbox_session_seq ON workbench_projection_outbox(session_id, outbox_seq); +CREATE INDEX IF NOT EXISTS idx_workbench_projection_outbox_after_seq ON workbench_projection_outbox(outbox_seq); + INSERT INTO workbench_sessions ( session_id, owner_user_id, diff --git a/internal/db/runtime-store.ts b/internal/db/runtime-store.ts index b3e56921..3d32b60a 100644 --- a/internal/db/runtime-store.ts +++ b/internal/db/runtime-store.ts @@ -1,5 +1,5 @@ /* - * SPEC: PJ2026-010403 API契约 draft-2026-06-17-r0; PJ2026-0102 Agent编排 draft-2026-06-17-r0; PJ2026-0104010803 Workbench唯一投影 draft-2026-06-19-p1-agentrun-incremental-cursor. + * SPEC: PJ2026-010403 API契约 draft-2026-06-17-r0; PJ2026-0102 Agent编排 draft-2026-06-17-r0; PJ2026-0104010803 Workbench唯一投影 draft-2026-06-19-p1-agentrun-incremental-cursor; draft-2026-06-20-p1-zero-split-durable-realtime. * 职责: Cloud runtime durable store。Code Agent trace/session/projection cursor 投影事实必须从同一持久化适配器写入和恢复。 */ import { createHash, randomUUID } from "node:crypto"; @@ -1276,6 +1276,48 @@ export class PostgresCloudRuntimeStore { for (const fact of result.facts.turns) await this.persistWorkbenchTurnFact(fact); for (const fact of result.facts.traceEvents) await this.persistWorkbenchTraceEventFact(fact); for (const fact of result.facts.checkpoints) await this.persistWorkbenchProjectionCheckpoint(fact); + for (const fact of result.facts.traceEvents) { + try { + await this.persistWorkbenchProjectionOutbox({ + traceId: fact.traceId, + sessionId: fact.sessionId, + turnId: fact.turnId, + messageId: fact.messageId, + projectedSeq: fact.projectedSeq, + sourceSeq: fact.sourceSeq, + sourceEventId: fact.sourceEventId, + commitType: fact.terminal ? "terminal" : "event", + terminal: Boolean(fact.terminal), + sealed: Boolean(fact.sealed), + payload: { type: fact.eventType, projectedSeq: fact.projectedSeq, valuesRedacted: true }, + createdAt: fact.occurredAt ?? fact.updatedAt ?? this.now() + }); + } catch (outboxError) { + console.error(`[workbench-projection] outbox write failed: traceId=${fact.traceId} projectedSeq=${fact.projectedSeq} error=${outboxError?.message ?? outboxError}`); + } + } + for (const fact of result.facts.turns) { + if (fact.terminal || fact.sealed) { + try { + await this.persistWorkbenchProjectionOutbox({ + traceId: fact.traceId, + sessionId: fact.sessionId, + turnId: fact.turnId, + messageId: fact.messageId, + projectedSeq: fact.projectedSeq, + sourceSeq: fact.sourceSeq, + sourceEventId: fact.sourceEventId, + commitType: "terminal", + terminal: true, + sealed: true, + payload: { status: fact.status, failureKind: fact.failureKind, valuesRedacted: true }, + createdAt: fact.updatedAt ?? this.now() + }); + } catch (outboxError) { + console.error(`[workbench-projection] terminal outbox write failed: traceId=${fact.traceId} error=${outboxError?.message ?? outboxError}`); + } + } + } return withPersistence(result, this.summary()); } @@ -1632,6 +1674,51 @@ export class PostgresCloudRuntimeStore { ); } + async persistWorkbenchProjectionOutbox(record) { + await this.query( + "INSERT INTO workbench_projection_outbox (trace_id, session_id, turn_id, message_id, projected_seq, source_seq, source_event_id, commit_type, terminal, sealed, payload_json, created_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12)", + [record.traceId, record.sessionId, record.turnId, record.messageId, record.projectedSeq, record.sourceSeq, record.sourceEventId, record.commitType ?? "event", Boolean(record.terminal), Boolean(record.sealed), stableJson(record.payload ?? record), record.createdAt ?? this.now()] + ); + } + + async readWorkbenchProjectionOutbox({ afterSeq = 0, limit = 100, traceId = null, sessionId = null } = {}) { + const params = []; + let where = "outbox_seq > $1"; + params.push(Number(afterSeq) || 0); + let paramIdx = 2; + if (traceId) { + where += ` AND trace_id = $${paramIdx}`; + params.push(traceId); + paramIdx += 1; + } + if (sessionId) { + where += ` AND session_id = $${paramIdx}`; + params.push(sessionId); + paramIdx += 1; + } + params.push(Math.min(Math.max(Number(limit) || 100, 1), 500)); + const result = await this.query( + `SELECT outbox_seq, trace_id, session_id, turn_id, message_id, projected_seq, source_seq, source_event_id, commit_type, terminal, sealed, payload_json, created_at FROM workbench_projection_outbox WHERE ${where} ORDER BY outbox_seq ASC LIMIT $${paramIdx}`, + params + ); + return (result.rows ?? []).map((row) => ({ + outboxSeq: Number(row.outbox_seq), + traceId: row.trace_id, + sessionId: row.session_id, + turnId: row.turn_id, + messageId: row.message_id, + projectedSeq: Number(row.projected_seq), + sourceSeq: Number(row.source_seq), + sourceEventId: row.source_event_id, + commitType: row.commit_type, + terminal: Boolean(row.terminal), + sealed: Boolean(row.sealed), + payload: typeof row.payload_json === "string" ? JSON.parse(row.payload_json) : row.payload_json, + createdAt: row.created_at, + valuesPrinted: false + })); + } + async persistWorkbenchProjectionCheckpoint(record) { await this.query( "INSERT INTO workbench_projection_checkpoints (trace_id, session_id, turn_id, run_id, command_id, projected_seq, source_seq, source_event_id, projection_status, projection_health, terminal, sealed, diagnostic_json, checkpoint_json, created_at, updated_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16) ON CONFLICT (trace_id) DO UPDATE SET session_id = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.session_id ELSE workbench_projection_checkpoints.session_id END, turn_id = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.turn_id ELSE workbench_projection_checkpoints.turn_id END, run_id = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.run_id ELSE workbench_projection_checkpoints.run_id END, command_id = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.command_id ELSE workbench_projection_checkpoints.command_id END, projected_seq = GREATEST(workbench_projection_checkpoints.projected_seq, EXCLUDED.projected_seq), source_seq = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.source_seq ELSE workbench_projection_checkpoints.source_seq END, source_event_id = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.source_event_id ELSE workbench_projection_checkpoints.source_event_id END, projection_status = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.projection_status ELSE workbench_projection_checkpoints.projection_status END, projection_health = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.projection_health ELSE workbench_projection_checkpoints.projection_health END, terminal = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.terminal ELSE workbench_projection_checkpoints.terminal END, sealed = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.sealed ELSE workbench_projection_checkpoints.sealed END, diagnostic_json = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.diagnostic_json ELSE workbench_projection_checkpoints.diagnostic_json END, checkpoint_json = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.checkpoint_json ELSE workbench_projection_checkpoints.checkpoint_json END, updated_at = CASE WHEN EXCLUDED.projected_seq >= workbench_projection_checkpoints.projected_seq THEN EXCLUDED.updated_at ELSE workbench_projection_checkpoints.updated_at END", diff --git a/internal/db/schema.ts b/internal/db/schema.ts index 794d2d43..ac93ad39 100644 --- a/internal/db/schema.ts +++ b/internal/db/schema.ts @@ -5,7 +5,7 @@ import { TABLES } from "../protocol/index.mjs"; export const CLOUD_CORE_MIGRATION_ID = "0001_cloud_core_skeleton"; -export const CLOUD_RUNTIME_DURABLE_ADAPTER_SCHEMA_VERSION = "runtime-durable-postgres-v4"; +export const CLOUD_RUNTIME_DURABLE_ADAPTER_SCHEMA_VERSION = "runtime-durable-postgres-v5"; export const CLOUD_RUNTIME_DURABLE_MIGRATION_TABLE = "hwlab_schema_migrations"; export const FROZEN_CLOUD_CORE_TABLES = Object.freeze([...TABLES]); @@ -266,6 +266,21 @@ export const CLOUD_RUNTIME_DURABLE_TABLE_COLUMNS = Object.freeze({ "checkpoint_json", "created_at", "updated_at" + ]), + workbench_projection_outbox: Object.freeze([ + "outbox_seq", + "trace_id", + "session_id", + "turn_id", + "message_id", + "projected_seq", + "source_seq", + "source_event_id", + "commit_type", + "terminal", + "sealed", + "payload_json", + "created_at" ]) });