feat: PJ2026-0104010803 P1 durable projection outbox + DB constraints + error visibility (#1754)

DB migration (0001_cloud_core_skeleton.sql):
- Add (trace_id, source_event_id) unique index on workbench_trace_events
- Add (trace_id, projected_seq) unique index on workbench_trace_events
- Add workbench_projection_outbox table with BIGSERIAL outbox_seq

Schema (schema.ts):
- Bump CLOUD_RUNTIME_DURABLE_ADAPTER_SCHEMA_VERSION to v5
- Add workbench_projection_outbox to CLOUD_RUNTIME_DURABLE_TABLE_COLUMNS

Runtime store (runtime-store.ts):
- Add persistWorkbenchProjectionOutbox method
- Add readWorkbenchProjectionOutbox method for SSE cursor replay
- Write outbox rows in writeWorkbenchFacts for trace events and terminal turns
- Outbox write errors are logged, not silently swallowed

Trace store (code-agent-trace-store.ts):
- persistTraceEvent no longer uses .catch(() => {}) to swallow errors
- Projection write errors are logged with traceId/sourceEventId for diagnostics
- Sync errors also logged instead of empty catch block

SPEC headers updated to draft-2026-06-20-p1-zero-split-durable-realtime

Closes #1747
Refs #1742
This commit is contained in:
Lyon
2026-06-20 19:29:44 +08:00
committed by GitHub
parent e54a0001c0
commit b254c1875e
5 changed files with 138 additions and 9 deletions
+12 -6
View File
@@ -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}`);
}
}
@@ -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";
@@ -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,
+88 -1
View File
@@ -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",
+16 -1
View File
@@ -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"
])
});