From 6dd02cabc7b4d50ea668e7810469d5148503c379 Mon Sep 17 00:00:00 2001 From: lyon Date: Sat, 20 Jun 2026 08:45:51 +0800 Subject: [PATCH] feat: add durable workbench facts store --- .../migrations/0001_cloud_core_skeleton.sql | 271 ++++++++++- internal/db/runtime-store.test.ts | 197 ++++++++ internal/db/runtime-store.ts | 425 +++++++++++++++++- internal/db/schema.test.ts | 17 + internal/db/schema.ts | 104 ++++- 5 files changed, 1008 insertions(+), 6 deletions(-) diff --git a/internal/db/migrations/0001_cloud_core_skeleton.sql b/internal/db/migrations/0001_cloud_core_skeleton.sql index dcc7cf0a..76c86391 100644 --- a/internal/db/migrations/0001_cloud_core_skeleton.sql +++ b/internal/db/migrations/0001_cloud_core_skeleton.sql @@ -195,6 +195,275 @@ CREATE INDEX IF NOT EXISTS idx_workbench_projection_state_status_retry_updated O CREATE INDEX IF NOT EXISTS idx_workbench_projection_state_run_command ON workbench_projection_state(run_id, command_id); CREATE INDEX IF NOT EXISTS idx_workbench_projection_state_session_updated ON workbench_projection_state(session_id, updated_at DESC); +CREATE TABLE IF NOT EXISTS workbench_sessions ( + session_id TEXT PRIMARY KEY, + owner_user_id TEXT REFERENCES users(id), + project_id TEXT, + conversation_id TEXT, + thread_id TEXT, + status TEXT NOT NULL DEFAULT 'unknown', + last_trace_id TEXT, + projected_seq INTEGER NOT NULL DEFAULT 0, + source_seq INTEGER NOT NULL DEFAULT 0, + source_event_id TEXT, + terminal BOOLEAN NOT NULL DEFAULT false, + sealed BOOLEAN NOT NULL DEFAULT false, + session_json TEXT NOT NULL DEFAULT '{}', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_workbench_sessions_owner_updated ON workbench_sessions(owner_user_id, updated_at DESC, session_id); +CREATE INDEX IF NOT EXISTS idx_workbench_sessions_status_updated ON workbench_sessions(status, updated_at DESC, session_id); + +CREATE TABLE IF NOT EXISTS workbench_messages ( + message_id TEXT PRIMARY KEY, + session_id TEXT NOT NULL, + turn_id TEXT, + trace_id TEXT, + role TEXT NOT NULL, + status TEXT NOT NULL DEFAULT 'unknown', + projected_seq INTEGER NOT NULL DEFAULT 0, + source_seq INTEGER NOT NULL DEFAULT 0, + source_event_id TEXT, + terminal BOOLEAN NOT NULL DEFAULT false, + sealed BOOLEAN NOT NULL DEFAULT false, + message_json TEXT NOT NULL DEFAULT '{}', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_workbench_messages_session_updated ON workbench_messages(session_id, updated_at ASC, message_id); +CREATE INDEX IF NOT EXISTS idx_workbench_messages_trace ON workbench_messages(trace_id, message_id); +CREATE INDEX IF NOT EXISTS idx_workbench_messages_turn ON workbench_messages(turn_id, message_id); + +CREATE TABLE IF NOT EXISTS workbench_parts ( + part_id TEXT PRIMARY KEY, + message_id TEXT NOT NULL, + session_id TEXT NOT NULL, + turn_id TEXT, + trace_id TEXT, + part_index INTEGER NOT NULL DEFAULT 0, + part_type TEXT NOT NULL DEFAULT 'text', + status TEXT NOT NULL DEFAULT 'unknown', + projected_seq INTEGER NOT NULL DEFAULT 0, + source_seq INTEGER NOT NULL DEFAULT 0, + source_event_id TEXT, + terminal BOOLEAN NOT NULL DEFAULT false, + sealed BOOLEAN NOT NULL DEFAULT false, + part_json TEXT NOT NULL DEFAULT '{}', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_workbench_parts_message_index ON workbench_parts(message_id, part_index, part_id); +CREATE INDEX IF NOT EXISTS idx_workbench_parts_session_trace ON workbench_parts(session_id, trace_id, part_index); + +CREATE TABLE IF NOT EXISTS workbench_turns ( + turn_id TEXT PRIMARY KEY, + session_id TEXT NOT NULL, + trace_id TEXT, + message_id TEXT, + status TEXT NOT NULL DEFAULT 'unknown', + projected_seq INTEGER NOT NULL DEFAULT 0, + source_seq INTEGER NOT NULL DEFAULT 0, + source_event_id TEXT, + terminal BOOLEAN NOT NULL DEFAULT false, + sealed BOOLEAN NOT NULL DEFAULT false, + final_response_json TEXT NOT NULL DEFAULT 'null', + diagnostic_json TEXT NOT NULL DEFAULT '{}', + turn_json TEXT NOT NULL DEFAULT '{}', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_workbench_turns_session_updated ON workbench_turns(session_id, updated_at DESC, turn_id); +CREATE INDEX IF NOT EXISTS idx_workbench_turns_trace ON workbench_turns(trace_id, turn_id); + +CREATE TABLE IF NOT EXISTS workbench_trace_events ( + id TEXT PRIMARY KEY, + trace_id TEXT NOT NULL, + session_id TEXT, + turn_id TEXT, + message_id TEXT, + source_seq INTEGER NOT NULL DEFAULT 0, + source_event_id TEXT, + projected_seq INTEGER NOT NULL DEFAULT 0, + event_type TEXT NOT NULL DEFAULT 'event', + terminal BOOLEAN NOT NULL DEFAULT false, + sealed BOOLEAN NOT NULL DEFAULT false, + event_json TEXT NOT NULL DEFAULT '{}', + occurred_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); +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 TABLE IF NOT EXISTS workbench_projection_checkpoints ( + trace_id TEXT PRIMARY KEY, + session_id TEXT, + turn_id TEXT, + run_id TEXT, + command_id TEXT, + projected_seq INTEGER NOT NULL DEFAULT 0, + source_seq INTEGER NOT NULL DEFAULT 0, + source_event_id TEXT, + projection_status TEXT NOT NULL DEFAULT 'projecting', + projection_health TEXT NOT NULL DEFAULT 'healthy', + terminal BOOLEAN NOT NULL DEFAULT false, + sealed BOOLEAN NOT NULL DEFAULT false, + diagnostic_json TEXT NOT NULL DEFAULT '{}', + checkpoint_json TEXT NOT NULL DEFAULT '{}', + created_at TEXT NOT NULL, + updated_at TEXT NOT NULL +); +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); + +INSERT INTO workbench_sessions ( + session_id, + owner_user_id, + project_id, + conversation_id, + thread_id, + status, + last_trace_id, + projected_seq, + source_seq, + source_event_id, + terminal, + sealed, + session_json, + created_at, + updated_at +) +SELECT + id, + owner_user_id, + project_id, + conversation_id, + thread_id, + status, + last_trace_id, + 0, + 0, + id, + status IN ('completed', 'failed', 'canceled', 'cancelled', 'archived'), + status IN ('completed', 'failed', 'canceled', 'cancelled', 'archived') AND ended_at IS NOT NULL, + COALESCE(NULLIF(session_json, ''), '{}'), + COALESCE(started_at, updated_at, CURRENT_TIMESTAMP::text), + COALESCE(updated_at, ended_at, started_at, CURRENT_TIMESTAMP::text) +FROM agent_sessions +WHERE id IS NOT NULL +ON CONFLICT (session_id) DO NOTHING; + +INSERT INTO workbench_trace_events ( + id, + trace_id, + session_id, + turn_id, + message_id, + source_seq, + source_event_id, + projected_seq, + event_type, + terminal, + sealed, + event_json, + occurred_at, + updated_at +) +SELECT + id, + trace_id, + agent_session_id, + trace_id, + NULL, + 0, + id, + 0, + 'agent_trace_event', + false, + false, + COALESCE(NULLIF(event_json, ''), '{}'), + occurred_at, + occurred_at +FROM agent_trace_events +WHERE id IS NOT NULL AND trace_id IS NOT NULL +ON CONFLICT (id) DO NOTHING; + +INSERT INTO workbench_turns ( + turn_id, + session_id, + trace_id, + message_id, + status, + projected_seq, + source_seq, + source_event_id, + terminal, + sealed, + final_response_json, + diagnostic_json, + turn_json, + created_at, + updated_at +) +SELECT + trace_id, + COALESCE(session_id, trace_id), + trace_id, + NULL, + projection_status, + last_projected_seq, + last_agentrun_seq, + trace_id || ':' || last_agentrun_seq::text, + projection_status IN ('caught_up', 'terminal'), + projection_status = 'terminal', + 'null', + COALESCE(NULLIF(projection_json, ''), '{}'), + COALESCE(NULLIF(projection_json, ''), '{}'), + COALESCE(created_at, updated_at, CURRENT_TIMESTAMP::text), + COALESCE(updated_at, created_at, CURRENT_TIMESTAMP::text) +FROM workbench_projection_state +WHERE trace_id IS NOT NULL +ON CONFLICT (turn_id) DO NOTHING; + +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 +) +SELECT + trace_id, + session_id, + trace_id, + run_id, + command_id, + last_projected_seq, + last_agentrun_seq, + trace_id || ':' || last_agentrun_seq::text, + projection_status, + projection_health, + projection_status IN ('caught_up', 'terminal'), + projection_status = 'terminal', + COALESCE(NULLIF(projection_json, ''), '{}'), + COALESCE(NULLIF(projection_json, ''), '{}'), + COALESCE(created_at, updated_at, CURRENT_TIMESTAMP::text), + COALESCE(updated_at, created_at, CURRENT_TIMESTAMP::text) +FROM workbench_projection_state +WHERE trace_id IS NOT NULL +ON CONFLICT (trace_id) DO NOTHING; + CREATE TABLE IF NOT EXISTS evidence_records ( id TEXT PRIMARY KEY, project_id TEXT, @@ -215,7 +484,7 @@ CREATE TABLE IF NOT EXISTS hwlab_schema_migrations ( INSERT INTO hwlab_schema_migrations (id, schema_version, applied_at, migration_json) VALUES ( '0001_cloud_core_skeleton', - 'runtime-durable-postgres-v3', + 'runtime-durable-postgres-v4', CURRENT_TIMESTAMP, '{"path":"internal/db/migrations/0001_cloud_core_skeleton.sql","runtime":"cloud-api"}' ) diff --git a/internal/db/runtime-store.test.ts b/internal/db/runtime-store.test.ts index 0eb60e0d..a339e432 100644 --- a/internal/db/runtime-store.test.ts +++ b/internal/db/runtime-store.test.ts @@ -694,6 +694,118 @@ test("configured postgres runtime persists and queries Workbench projection stat assert.ok(queryClient.calls.some((call) => call.sql.startsWith("INSERT INTO workbench_projection_state"))); }); +test("configured postgres runtime persists and queries Workbench durable facts", async () => { + const queryClient = createFakePostgresClient({ migrationReady: true }); + const store = createConfiguredCloudRuntimeStore({ + env: { + HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres", + HWLAB_CLOUD_RUNTIME_DURABLE: "true" + }, + dbUrl: "postgres://hwlab_redacted@db.example.invalid:5432/hwlab", + queryClient, + now: () => "2026-06-20T10:00:00.000Z" + }); + + const write = await store.writeWorkbenchFacts({ + facts: { + sessions: [{ + sessionId: "ses_fact_1", + ownerUserId: "usr_fact_1", + projectId: "prj_fact_1", + conversationId: "cnv_fact_1", + threadId: "thread-fact-1", + status: "running", + lastTraceId: "trc_fact_1", + projectedSeq: 11, + sourceSeq: 12, + sourceEventId: "evt_fact_12", + updatedAt: "2026-06-20T10:00:01.000Z" + }], + messages: [{ + messageId: "msg_fact_1", + sessionId: "ses_fact_1", + turnId: "turn_fact_1", + traceId: "trc_fact_1", + role: "assistant", + status: "streaming", + projectedSeq: 13, + sourceSeq: 14, + sourceEventId: "evt_fact_14", + updatedAt: "2026-06-20T10:00:02.000Z" + }], + parts: [{ + partId: "part_fact_1", + messageId: "msg_fact_1", + sessionId: "ses_fact_1", + turnId: "turn_fact_1", + traceId: "trc_fact_1", + partIndex: 0, + partType: "text", + status: "streaming", + updatedAt: "2026-06-20T10:00:03.000Z" + }], + turns: [{ + turnId: "turn_fact_1", + sessionId: "ses_fact_1", + traceId: "trc_fact_1", + messageId: "msg_fact_1", + status: "completed", + finalResponse: { text: "done" }, + terminal: true, + sealed: true, + updatedAt: "2026-06-20T10:00:04.000Z" + }], + traceEvents: [{ + id: "wte_fact_1", + traceId: "trc_fact_1", + sessionId: "ses_fact_1", + turnId: "turn_fact_1", + messageId: "msg_fact_1", + sourceSeq: 15, + projectedSeq: 16, + eventType: "assistant.delta", + occurredAt: "2026-06-20T10:00:05.000Z" + }], + checkpoints: [{ + traceId: "trc_fact_1", + sessionId: "ses_fact_1", + turnId: "turn_fact_1", + runId: "run_fact_1", + commandId: "cmd_fact_1", + projectedSeq: 16, + sourceSeq: 15, + projectionStatus: "caught_up", + projectionHealth: "healthy", + terminal: true, + sealed: true, + updatedAt: "2026-06-20T10:00:06.000Z" + }] + } + }); + const loaded = await store.queryWorkbenchFacts({ sessionId: "ses_fact_1", limit: 10 }); + const readiness = await store.readiness(); + + assert.equal(write.written, true); + assert.equal(write.facts.sessions[0].terminal, false); + assert.equal(loaded.count, 6); + assert.equal(loaded.facts.sessions[0].sessionId, "ses_fact_1"); + assert.equal(loaded.facts.messages[0].messageId, "msg_fact_1"); + assert.equal(loaded.facts.parts[0].partId, "part_fact_1"); + assert.equal(loaded.facts.turns[0].finalResponse.text, "done"); + assert.equal(loaded.facts.traceEvents[0].eventType, "assistant.delta"); + assert.equal(loaded.facts.checkpoints[0].projectionStatus, "caught_up"); + assert.equal(readiness.counts.workbenchSessions, 1); + assert.equal(readiness.counts.workbenchMessages, 1); + assert.equal(readiness.counts.workbenchParts, 1); + assert.equal(readiness.counts.workbenchTurns, 1); + assert.equal(readiness.counts.workbenchTraceEvents, 1); + assert.equal(readiness.counts.workbenchProjectionCheckpoints, 1); + assert.ok(queryClient.calls.some((call) => call.sql.startsWith("CREATE INDEX IF NOT EXISTS idx_workbench_sessions_owner_updated"))); + assert.ok(queryClient.calls.some((call) => call.sql.startsWith("INSERT INTO workbench_sessions"))); + assert.ok(queryClient.calls.some((call) => call.sql.startsWith("INSERT INTO workbench_projection_checkpoints"))); + assert.ok(queryClient.calls.some((call) => call.sql.startsWith("SELECT session_json FROM workbench_sessions WHERE session_id = $1"))); +}); + test("configured postgres runtime pushes trace event query filters into SQL", async () => { const queryClient = createFakePostgresClient({ migrationReady: true }); const store = createConfiguredCloudRuntimeStore({ @@ -734,6 +846,12 @@ function createFakePostgresClient({ worker_sessions: new Map(), agent_trace_events: new Map(), workbench_projection_state: new Map(), + workbench_sessions: new Map(), + workbench_messages: new Map(), + workbench_parts: new Map(), + workbench_turns: new Map(), + workbench_trace_events: new Map(), + workbench_projection_checkpoints: new Map(), hwlab_schema_migrations: new Map() }; if (migrationReady) { @@ -775,6 +893,9 @@ function createFakePostgresClient({ if (sql.startsWith("CREATE INDEX IF NOT EXISTS idx_workbench_projection_state_")) { return { rows: [] }; } + if (sql.startsWith("CREATE INDEX IF NOT EXISTS idx_workbench_")) { + return { rows: [] }; + } if (sql.startsWith("SELECT gateway_session_json FROM gateway_sessions")) { return jsonSelect(state.gateway_sessions, params[0], "gateway_session_json"); } @@ -830,6 +951,24 @@ function createFakePostgresClient({ rows.sort((left, right) => String(left.updated_at).localeCompare(String(right.updated_at)) || String(left.trace_id).localeCompare(String(right.trace_id))); return { rows: rows.map((record) => ({ projection_json: record.projection_json })) }; } + if (sql.startsWith("SELECT session_json FROM workbench_sessions")) { + return workbenchFactRows(state.workbench_sessions, "session_json", sql, params, readErrorCode); + } + if (sql.startsWith("SELECT message_json FROM workbench_messages")) { + return workbenchFactRows(state.workbench_messages, "message_json", sql, params, readErrorCode); + } + if (sql.startsWith("SELECT part_json FROM workbench_parts")) { + return workbenchFactRows(state.workbench_parts, "part_json", sql, params, readErrorCode); + } + if (sql.startsWith("SELECT turn_json FROM workbench_turns")) { + return workbenchFactRows(state.workbench_turns, "turn_json", sql, params, readErrorCode); + } + if (sql.startsWith("SELECT event_json FROM workbench_trace_events")) { + return workbenchFactRows(state.workbench_trace_events, "event_json", sql, params, readErrorCode); + } + if (sql.startsWith("SELECT checkpoint_json FROM workbench_projection_checkpoints")) { + return workbenchFactRows(state.workbench_projection_checkpoints, "checkpoint_json", sql, params, readErrorCode); + } if (sql.startsWith("SELECT metadata_json FROM evidence_records")) { if (readErrorCode) { const error = new Error("evidence read query failed"); @@ -896,6 +1035,42 @@ function createFakePostgresClient({ }); return { rows: [] }; } + if (sql.startsWith("INSERT INTO workbench_sessions")) { + state.workbench_sessions.set(params[0], { + session_id: params[0], owner_user_id: params[1], project_id: params[2], conversation_id: params[3], thread_id: params[4], status: params[5], last_trace_id: params[6], projected_seq: params[7], source_seq: params[8], source_event_id: params[9], terminal: params[10], sealed: params[11], session_json: params[12], created_at: params[13], updated_at: params[14] + }); + return { rows: [] }; + } + if (sql.startsWith("INSERT INTO workbench_messages")) { + state.workbench_messages.set(params[0], { + message_id: params[0], session_id: params[1], turn_id: params[2], trace_id: params[3], role: params[4], status: params[5], projected_seq: params[6], source_seq: params[7], source_event_id: params[8], terminal: params[9], sealed: params[10], message_json: params[11], created_at: params[12], updated_at: params[13] + }); + return { rows: [] }; + } + if (sql.startsWith("INSERT INTO workbench_parts")) { + state.workbench_parts.set(params[0], { + part_id: params[0], message_id: params[1], session_id: params[2], turn_id: params[3], trace_id: params[4], part_index: params[5], part_type: params[6], status: params[7], projected_seq: params[8], source_seq: params[9], source_event_id: params[10], terminal: params[11], sealed: params[12], part_json: params[13], created_at: params[14], updated_at: params[15] + }); + return { rows: [] }; + } + if (sql.startsWith("INSERT INTO workbench_turns")) { + state.workbench_turns.set(params[0], { + turn_id: params[0], session_id: params[1], trace_id: params[2], message_id: params[3], status: params[4], projected_seq: params[5], source_seq: params[6], source_event_id: params[7], terminal: params[8], sealed: params[9], final_response_json: params[10], diagnostic_json: params[11], turn_json: params[12], created_at: params[13], updated_at: params[14] + }); + return { rows: [] }; + } + if (sql.startsWith("INSERT INTO workbench_trace_events")) { + state.workbench_trace_events.set(params[0], { + id: params[0], trace_id: params[1], session_id: params[2], turn_id: params[3], message_id: params[4], source_seq: params[5], source_event_id: params[6], projected_seq: params[7], event_type: params[8], terminal: params[9], sealed: params[10], event_json: params[11], occurred_at: params[12], updated_at: params[13] + }); + return { rows: [] }; + } + if (sql.startsWith("INSERT INTO workbench_projection_checkpoints")) { + state.workbench_projection_checkpoints.set(params[0], { + trace_id: params[0], session_id: params[1], turn_id: params[2], run_id: params[3], command_id: params[4], projected_seq: params[5], source_seq: params[6], source_event_id: params[7], projection_status: params[8], projection_health: params[9], terminal: params[10], sealed: params[11], diagnostic_json: params[12], checkpoint_json: params[13], created_at: params[14], updated_at: params[15] + }); + return { rows: [] }; + } throw new Error(`unexpected sql: ${sql}`); } }; @@ -916,3 +1091,25 @@ function jsonSelect(map, id, column) { rows: row ? [{ [column]: row[column] }] : [] }; } + +function workbenchFactRows(map, jsonColumn, sql, params, readErrorCode) { + if (readErrorCode) { + const error = new Error("workbench fact read query failed"); + error.code = readErrorCode; + throw error; + } + const rows = [...map.values()] + .filter((record) => sqlRecordMatches(record, sql, params)) + .sort((left, right) => String(left.updated_at).localeCompare(String(right.updated_at)) || String(left.id ?? left.session_id ?? left.message_id ?? left.part_id ?? left.turn_id ?? left.trace_id).localeCompare(String(right.id ?? right.session_id ?? right.message_id ?? right.part_id ?? right.turn_id ?? right.trace_id))) + .map((record) => ({ [jsonColumn]: record[jsonColumn] })); + const limit = sql.includes(" LIMIT $") ? Number(params.at(-1)) : null; + return { rows: Number.isInteger(limit) && limit > 0 ? rows.slice(0, limit) : rows }; +} + +function sqlRecordMatches(record, sql, params) { + for (const column of Object.keys(record)) { + const match = sql.match(new RegExp(`(?:^|[^a-z_])${column} = \\$(\\d+)`, "u")); + if (match && record[column] !== params[Number(match[1]) - 1]) return false; + } + return true; +} diff --git a/internal/db/runtime-store.ts b/internal/db/runtime-store.ts index e25ce7e6..01bc92a0 100644 --- a/internal/db/runtime-store.ts +++ b/internal/db/runtime-store.ts @@ -53,7 +53,13 @@ const postgresCountTables = Object.freeze([ "account_workspaces", "worker_sessions", "agent_trace_events", - "workbench_projection_state" + "workbench_projection_state", + "workbench_sessions", + "workbench_messages", + "workbench_parts", + "workbench_turns", + "workbench_trace_events", + "workbench_projection_checkpoints" ]); const postgresRuntimeReadIndexes = Object.freeze([ @@ -62,7 +68,20 @@ const postgresRuntimeReadIndexes = Object.freeze([ "CREATE INDEX IF NOT EXISTS idx_agent_trace_events_worker_order ON agent_trace_events(worker_session_id, occurred_at, id)", "CREATE INDEX IF NOT EXISTS idx_workbench_projection_state_status_retry_updated ON workbench_projection_state(projection_status, next_retry_at, updated_at)", "CREATE INDEX IF NOT EXISTS idx_workbench_projection_state_run_command ON workbench_projection_state(run_id, command_id)", - "CREATE INDEX IF NOT EXISTS idx_workbench_projection_state_session_updated ON workbench_projection_state(session_id, updated_at DESC)" + "CREATE INDEX IF NOT EXISTS idx_workbench_projection_state_session_updated ON workbench_projection_state(session_id, updated_at DESC)", + "CREATE INDEX IF NOT EXISTS idx_workbench_sessions_owner_updated ON workbench_sessions(owner_user_id, updated_at DESC, session_id)", + "CREATE INDEX IF NOT EXISTS idx_workbench_sessions_status_updated ON workbench_sessions(status, updated_at DESC, session_id)", + "CREATE INDEX IF NOT EXISTS idx_workbench_messages_session_updated ON workbench_messages(session_id, updated_at ASC, message_id)", + "CREATE INDEX IF NOT EXISTS idx_workbench_messages_trace ON workbench_messages(trace_id, message_id)", + "CREATE INDEX IF NOT EXISTS idx_workbench_messages_turn ON workbench_messages(turn_id, message_id)", + "CREATE INDEX IF NOT EXISTS idx_workbench_parts_message_index ON workbench_parts(message_id, part_index, part_id)", + "CREATE INDEX IF NOT EXISTS idx_workbench_parts_session_trace ON workbench_parts(session_id, trace_id, part_index)", + "CREATE INDEX IF NOT EXISTS idx_workbench_turns_session_updated ON workbench_turns(session_id, updated_at DESC, turn_id)", + "CREATE INDEX IF NOT EXISTS idx_workbench_turns_trace ON workbench_turns(trace_id, turn_id)", + "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 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)" ]); export function createCloudRuntimeStore(options = {}) { @@ -111,6 +130,12 @@ export class CloudRuntimeStore { this.evidenceRecords = new Map(); this.agentTraceEvents = new Map(); this.workbenchProjectionStates = new Map(); + this.workbenchSessions = new Map(); + this.workbenchMessages = new Map(); + this.workbenchParts = new Map(); + this.workbenchTurns = new Map(); + this.workbenchTraceEvents = new Map(); + this.workbenchProjectionCheckpoints = new Map(); } summary() { @@ -133,7 +158,13 @@ export class CloudRuntimeStore { auditEvents: this.auditEvents.size, evidenceRecords: this.evidenceRecords.size, agentTraceEvents: this.agentTraceEvents.size, - workbenchProjectionStates: this.workbenchProjectionStates.size + workbenchProjectionStates: this.workbenchProjectionStates.size, + workbenchSessions: this.workbenchSessions.size, + workbenchMessages: this.workbenchMessages.size, + workbenchParts: this.workbenchParts.size, + workbenchTurns: this.workbenchTurns.size, + workbenchTraceEvents: this.workbenchTraceEvents.size, + workbenchProjectionCheckpoints: this.workbenchProjectionCheckpoints.size } }); } @@ -557,6 +588,43 @@ export class CloudRuntimeStore { }; } + writeWorkbenchFacts(params = {}, requestMeta = {}) { + const facts = normalizeWorkbenchFacts(params.facts ?? params, requestMeta, this.now()); + for (const fact of facts.sessions) this.workbenchSessions.set(fact.sessionId, fact); + for (const fact of facts.messages) this.workbenchMessages.set(fact.messageId, fact); + for (const fact of facts.parts) this.workbenchParts.set(fact.partId, fact); + for (const fact of facts.turns) this.workbenchTurns.set(fact.turnId, fact); + for (const fact of facts.traceEvents) this.workbenchTraceEvents.set(fact.id, fact); + for (const fact of facts.checkpoints) this.workbenchProjectionCheckpoints.set(fact.traceId, fact); + return { + written: true, + facts, + persistence: this.summary() + }; + } + + queryWorkbenchFacts(params = {}) { + const sessions = [...this.workbenchSessions.values()].filter((record) => matchesQuery(record, params, ["sessionId", "ownerUserId", "projectId", "conversationId", "threadId", "status"])); + const messages = [...this.workbenchMessages.values()].filter((record) => matchesQuery(record, params, ["messageId", "sessionId", "turnId", "traceId", "role", "status"])); + const parts = [...this.workbenchParts.values()].filter((record) => matchesQuery(record, params, ["partId", "messageId", "sessionId", "turnId", "traceId", "partType", "status"])); + const turns = [...this.workbenchTurns.values()].filter((record) => matchesQuery(record, params, ["turnId", "sessionId", "traceId", "messageId", "status"])); + const traceEvents = [...this.workbenchTraceEvents.values()].filter((record) => matchesQuery(record, params, ["id", "traceId", "sessionId", "turnId", "messageId", "eventType"])); + const checkpoints = [...this.workbenchProjectionCheckpoints.values()].filter((record) => matchesQuery(record, params, ["traceId", "sessionId", "turnId", "runId", "commandId", "projectionStatus", "projectionHealth"])); + const facts = { + sessions: limitResults(sessions.sort(compareWorkbenchFactRecords), params.limit), + messages: limitResults(messages.sort(compareWorkbenchFactRecords), params.limit), + parts: limitResults(parts.sort(compareWorkbenchPartFactRecords), params.limit), + turns: limitResults(turns.sort(compareWorkbenchFactRecords), params.limit), + traceEvents: limitResults(traceEvents.sort(compareWorkbenchTraceEventFacts), params.limit), + checkpoints: limitResults(checkpoints.sort(compareWorkbenchFactRecords), params.limit) + }; + return { + facts, + count: Object.values(facts).reduce((sum, rows) => sum + rows.length, 0), + persistence: this.summary() + }; + } + recordAuditEvent(input) { return this.writeAuditEvent(this.buildAuditEventFromInput(input, input.requestMeta), input.requestMeta).auditEvent; } @@ -1127,6 +1195,53 @@ export class PostgresCloudRuntimeStore { }; } + async writeWorkbenchFacts(params = {}, requestMeta = {}) { + const readiness = this.summary(); + if (readiness.ready !== true || readiness.durable !== true || readiness.liveRuntimeEvidence !== true) { + await this.assertReadyForWrites(); + } + const result = this.memory.writeWorkbenchFacts(params, requestMeta); + for (const fact of result.facts.sessions) await this.persistWorkbenchSessionFact(fact); + for (const fact of result.facts.messages) await this.persistWorkbenchMessageFact(fact); + for (const fact of result.facts.parts) await this.persistWorkbenchPartFact(fact); + 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); + return withPersistence(result, this.summary()); + } + + async queryWorkbenchFacts(params = {}) { + await this.assertReadyForDurableReads("workbench.facts.query"); + const facts = { + sessions: await this.queryWorkbenchFactRows("workbench.sessions.query", "workbench_sessions", "session_json", params, { sessionId: "session_id", ownerUserId: "owner_user_id", projectId: "project_id", conversationId: "conversation_id", threadId: "thread_id", status: "status" }), + messages: await this.queryWorkbenchFactRows("workbench.messages.query", "workbench_messages", "message_json", params, { messageId: "message_id", sessionId: "session_id", turnId: "turn_id", traceId: "trace_id", role: "role", status: "status" }), + parts: await this.queryWorkbenchFactRows("workbench.parts.query", "workbench_parts", "part_json", params, { partId: "part_id", messageId: "message_id", sessionId: "session_id", turnId: "turn_id", traceId: "trace_id", partType: "part_type", status: "status" }), + turns: await this.queryWorkbenchFactRows("workbench.turns.query", "workbench_turns", "turn_json", params, { turnId: "turn_id", sessionId: "session_id", traceId: "trace_id", messageId: "message_id", status: "status" }), + traceEvents: await this.queryWorkbenchFactRows("workbench.trace-events.query", "workbench_trace_events", "event_json", params, { id: "id", traceId: "trace_id", sessionId: "session_id", turnId: "turn_id", messageId: "message_id", eventType: "event_type" }), + checkpoints: await this.queryWorkbenchFactRows("workbench.projection-checkpoints.query", "workbench_projection_checkpoints", "checkpoint_json", params, { traceId: "trace_id", sessionId: "session_id", turnId: "turn_id", runId: "run_id", commandId: "command_id", projectionStatus: "projection_status", projectionHealth: "projection_health" }) + }; + return withPersistence({ facts, count: Object.values(facts).reduce((sum, rows) => sum + rows.length, 0) }, this.summary()); + } + + async queryWorkbenchFactRows(method, table, jsonColumn, params = {}, columnMap = {}) { + const queryParams = []; + const clauses = []; + for (const [field, column] of Object.entries(columnMap)) { + const value = textOr(params[field], ""); + if (!value) continue; + queryParams.push(value); + clauses.push(`${column} = $${queryParams.length}`); + } + const limit = positiveIntegerOrNull(params.limit); + let sql = `SELECT ${jsonColumn} FROM ${table}${clauses.length ? ` WHERE ${clauses.join(" AND ")}` : ""} ORDER BY updated_at ASC`; + if (limit) { + queryParams.push(limit); + sql += ` LIMIT $${queryParams.length}`; + } + const result = await this.queryDurableReadRows(method, sql, queryParams); + return result.rows.map((row) => parseJsonColumn(row[jsonColumn], null)).filter(Boolean); + } + async hydrateOperationRefs(params = {}) { await this.hydrateGatewaySession(params.gatewaySessionId); await this.hydrateBoxResource(params.resourceId); @@ -1191,6 +1306,24 @@ export class PostgresCloudRuntimeStore { for (const record of changedRecords(this.memory.workbenchProjectionStates, before.workbenchProjectionStates)) { await this.persistWorkbenchProjectionState(record); } + for (const record of changedRecords(this.memory.workbenchSessions, before.workbenchSessions)) { + await this.persistWorkbenchSessionFact(record); + } + for (const record of changedRecords(this.memory.workbenchMessages, before.workbenchMessages)) { + await this.persistWorkbenchMessageFact(record); + } + for (const record of changedRecords(this.memory.workbenchParts, before.workbenchParts)) { + await this.persistWorkbenchPartFact(record); + } + for (const record of changedRecords(this.memory.workbenchTurns, before.workbenchTurns)) { + await this.persistWorkbenchTurnFact(record); + } + for (const record of changedRecords(this.memory.workbenchTraceEvents, before.workbenchTraceEvents)) { + await this.persistWorkbenchTraceEventFact(record); + } + for (const record of changedRecords(this.memory.workbenchProjectionCheckpoints, before.workbenchProjectionCheckpoints)) { + await this.persistWorkbenchProjectionCheckpoint(record); + } } async assertReadyForWrites() { @@ -1366,6 +1499,48 @@ export class PostgresCloudRuntimeStore { ); } + async persistWorkbenchSessionFact(record) { + await this.query( + "INSERT INTO workbench_sessions (session_id, owner_user_id, project_id, conversation_id, thread_id, status, last_trace_id, projected_seq, source_seq, source_event_id, terminal, sealed, session_json, created_at, updated_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15) ON CONFLICT (session_id) DO UPDATE SET owner_user_id = EXCLUDED.owner_user_id, project_id = EXCLUDED.project_id, conversation_id = EXCLUDED.conversation_id, thread_id = EXCLUDED.thread_id, status = EXCLUDED.status, last_trace_id = EXCLUDED.last_trace_id, projected_seq = EXCLUDED.projected_seq, source_seq = EXCLUDED.source_seq, source_event_id = EXCLUDED.source_event_id, terminal = EXCLUDED.terminal, sealed = EXCLUDED.sealed, session_json = EXCLUDED.session_json, updated_at = EXCLUDED.updated_at", + [record.sessionId, record.ownerUserId, record.projectId, record.conversationId, record.threadId, record.status, record.lastTraceId, record.projectedSeq, record.sourceSeq, record.sourceEventId, record.terminal, record.sealed, stableJson(record), record.createdAt, record.updatedAt] + ); + } + + async persistWorkbenchMessageFact(record) { + await this.query( + "INSERT INTO workbench_messages (message_id, session_id, turn_id, trace_id, role, status, projected_seq, source_seq, source_event_id, terminal, sealed, message_json, created_at, updated_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14) ON CONFLICT (message_id) DO UPDATE SET session_id = EXCLUDED.session_id, turn_id = EXCLUDED.turn_id, trace_id = EXCLUDED.trace_id, role = EXCLUDED.role, status = EXCLUDED.status, projected_seq = EXCLUDED.projected_seq, source_seq = EXCLUDED.source_seq, source_event_id = EXCLUDED.source_event_id, terminal = EXCLUDED.terminal, sealed = EXCLUDED.sealed, message_json = EXCLUDED.message_json, updated_at = EXCLUDED.updated_at", + [record.messageId, record.sessionId, record.turnId, record.traceId, record.role, record.status, record.projectedSeq, record.sourceSeq, record.sourceEventId, record.terminal, record.sealed, stableJson(record), record.createdAt, record.updatedAt] + ); + } + + async persistWorkbenchPartFact(record) { + await this.query( + "INSERT INTO workbench_parts (part_id, message_id, session_id, turn_id, trace_id, part_index, part_type, status, projected_seq, source_seq, source_event_id, terminal, sealed, part_json, created_at, updated_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16) ON CONFLICT (part_id) DO UPDATE SET message_id = EXCLUDED.message_id, session_id = EXCLUDED.session_id, turn_id = EXCLUDED.turn_id, trace_id = EXCLUDED.trace_id, part_index = EXCLUDED.part_index, part_type = EXCLUDED.part_type, status = EXCLUDED.status, projected_seq = EXCLUDED.projected_seq, source_seq = EXCLUDED.source_seq, source_event_id = EXCLUDED.source_event_id, terminal = EXCLUDED.terminal, sealed = EXCLUDED.sealed, part_json = EXCLUDED.part_json, updated_at = EXCLUDED.updated_at", + [record.partId, record.messageId, record.sessionId, record.turnId, record.traceId, record.partIndex, record.partType, record.status, record.projectedSeq, record.sourceSeq, record.sourceEventId, record.terminal, record.sealed, stableJson(record), record.createdAt, record.updatedAt] + ); + } + + async persistWorkbenchTurnFact(record) { + await this.query( + "INSERT INTO workbench_turns (turn_id, session_id, trace_id, message_id, status, projected_seq, source_seq, source_event_id, terminal, sealed, final_response_json, diagnostic_json, turn_json, created_at, updated_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15) ON CONFLICT (turn_id) DO UPDATE SET session_id = EXCLUDED.session_id, trace_id = EXCLUDED.trace_id, message_id = EXCLUDED.message_id, status = EXCLUDED.status, projected_seq = EXCLUDED.projected_seq, source_seq = EXCLUDED.source_seq, source_event_id = EXCLUDED.source_event_id, terminal = EXCLUDED.terminal, sealed = EXCLUDED.sealed, final_response_json = EXCLUDED.final_response_json, diagnostic_json = EXCLUDED.diagnostic_json, turn_json = EXCLUDED.turn_json, updated_at = EXCLUDED.updated_at", + [record.turnId, record.sessionId, record.traceId, record.messageId, record.status, record.projectedSeq, record.sourceSeq, record.sourceEventId, record.terminal, record.sealed, stableJson(record.finalResponse), stableJson(record.diagnostic), stableJson(record), record.createdAt, record.updatedAt] + ); + } + + async persistWorkbenchTraceEventFact(record) { + await this.query( + "INSERT INTO workbench_trace_events (id, trace_id, session_id, turn_id, message_id, source_seq, source_event_id, projected_seq, event_type, terminal, sealed, event_json, occurred_at, updated_at) VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14) ON CONFLICT (id) DO UPDATE SET trace_id = EXCLUDED.trace_id, session_id = EXCLUDED.session_id, turn_id = EXCLUDED.turn_id, message_id = EXCLUDED.message_id, source_seq = EXCLUDED.source_seq, source_event_id = EXCLUDED.source_event_id, projected_seq = EXCLUDED.projected_seq, event_type = EXCLUDED.event_type, terminal = EXCLUDED.terminal, sealed = EXCLUDED.sealed, event_json = EXCLUDED.event_json, occurred_at = EXCLUDED.occurred_at, updated_at = EXCLUDED.updated_at", + [record.id, record.traceId, record.sessionId, record.turnId, record.messageId, record.sourceSeq, record.sourceEventId, record.projectedSeq, record.eventType, record.terminal, record.sealed, stableJson(record), record.occurredAt, record.updatedAt] + ); + } + + 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 = EXCLUDED.session_id, turn_id = EXCLUDED.turn_id, run_id = EXCLUDED.run_id, command_id = EXCLUDED.command_id, projected_seq = EXCLUDED.projected_seq, source_seq = EXCLUDED.source_seq, source_event_id = EXCLUDED.source_event_id, projection_status = EXCLUDED.projection_status, projection_health = EXCLUDED.projection_health, terminal = EXCLUDED.terminal, sealed = EXCLUDED.sealed, diagnostic_json = EXCLUDED.diagnostic_json, checkpoint_json = EXCLUDED.checkpoint_json, updated_at = EXCLUDED.updated_at", + [record.traceId, record.sessionId, record.turnId, record.runId, record.commandId, record.projectedSeq, record.sourceSeq, record.sourceEventId, record.projectionStatus, record.projectionHealth, record.terminal, record.sealed, stableJson(record.diagnostic), stableJson(record), record.createdAt, record.updatedAt] + ); + } + async readCounts() { const counts = {}; for (const table of postgresCountTables) { @@ -1630,6 +1805,218 @@ function normalizeWorkbenchProjectionState(value, requestMeta = {}, now) { }); } +function normalizeWorkbenchFacts(value, requestMeta = {}, now) { + const input = normalizeJsonObject(value); + return { + sessions: arrayInput(input.sessions ?? input.session).map((item) => normalizeWorkbenchSessionFact(item, requestMeta, now)), + messages: arrayInput(input.messages ?? input.message).map((item) => normalizeWorkbenchMessageFact(item, requestMeta, now)), + parts: arrayInput(input.parts ?? input.part).map((item) => normalizeWorkbenchPartFact(item, requestMeta, now)), + turns: arrayInput(input.turns ?? input.turn).map((item) => normalizeWorkbenchTurnFact(item, requestMeta, now)), + traceEvents: arrayInput(input.traceEvents ?? input.traceEvent).map((item) => normalizeWorkbenchTraceEventFact(item, requestMeta, now)), + checkpoints: arrayInput(input.checkpoints ?? input.checkpoint).map((item) => normalizeWorkbenchProjectionCheckpoint(item, requestMeta, now)) + }; +} + +function normalizeWorkbenchSessionFact(value, requestMeta = {}, now) { + const input = normalizeJsonObject(value); + const sessionId = textOr(input.sessionId ?? input.id ?? requestMeta.sessionId, ""); + if (!sessionId) throw requiredWorkbenchFactError("workbench session fact", ["sessionId"]); + const sourceSeq = nonNegativeInteger(input.sourceSeq); + const projectedSeq = nonNegativeInteger(input.projectedSeq ?? sourceSeq); + const timestamp = timestampOr(input.updatedAt ?? input.createdAt, now); + return pruneUndefined({ + ...input, + id: sessionId, + sessionId, + ownerUserId: textOr(input.ownerUserId ?? requestMeta.ownerUserId, "") || null, + projectId: textOr(input.projectId ?? requestMeta.projectId, "") || null, + conversationId: textOr(input.conversationId ?? requestMeta.conversationId, "") || null, + threadId: textOr(input.threadId ?? requestMeta.threadId, "") || null, + status: textOr(input.status, "unknown"), + lastTraceId: textOr(input.lastTraceId ?? input.traceId ?? requestMeta.traceId, "") || null, + projectedSeq, + sourceSeq, + sourceEventId: textOr(input.sourceEventId ?? input.eventId, "") || null, + terminal: Boolean(input.terminal), + sealed: Boolean(input.sealed), + createdAt: timestampOr(input.createdAt, timestamp), + updatedAt: timestamp, + valuesPrinted: false + }); +} + +function normalizeWorkbenchMessageFact(value, requestMeta = {}, now) { + const input = normalizeJsonObject(value); + const sessionId = textOr(input.sessionId ?? requestMeta.sessionId, ""); + if (!sessionId) throw requiredWorkbenchFactError("workbench message fact", ["sessionId"]); + const traceId = textOr(input.traceId ?? requestMeta.traceId, "") || null; + const turnId = textOr(input.turnId ?? requestMeta.turnId ?? traceId, "") || null; + const role = textOr(input.role, "agent"); + const sourceSeq = nonNegativeInteger(input.sourceSeq); + const projectedSeq = nonNegativeInteger(input.projectedSeq ?? sourceSeq); + const timestamp = timestampOr(input.updatedAt ?? input.createdAt, now); + const messageId = textOr(input.messageId ?? input.id, "") || stableWorkbenchFactId("msg", { sessionId, turnId, traceId, role, sourceSeq, projectedSeq }); + return pruneUndefined({ + ...input, + id: messageId, + messageId, + sessionId, + turnId, + traceId, + role, + status: textOr(input.status, "unknown"), + projectedSeq, + sourceSeq, + sourceEventId: textOr(input.sourceEventId ?? input.eventId, "") || null, + terminal: Boolean(input.terminal), + sealed: Boolean(input.sealed), + createdAt: timestampOr(input.createdAt, timestamp), + updatedAt: timestamp, + valuesPrinted: false + }); +} + +function normalizeWorkbenchPartFact(value, requestMeta = {}, now) { + const input = normalizeJsonObject(value); + const sessionId = textOr(input.sessionId ?? requestMeta.sessionId, ""); + const messageId = textOr(input.messageId ?? requestMeta.messageId, ""); + if (!sessionId || !messageId) throw requiredWorkbenchFactError("workbench part fact", ["sessionId", "messageId"]); + const traceId = textOr(input.traceId ?? requestMeta.traceId, "") || null; + const turnId = textOr(input.turnId ?? requestMeta.turnId ?? traceId, "") || null; + const partIndex = nonNegativeInteger(input.partIndex ?? input.index); + const sourceSeq = nonNegativeInteger(input.sourceSeq); + const projectedSeq = nonNegativeInteger(input.projectedSeq ?? sourceSeq); + const timestamp = timestampOr(input.updatedAt ?? input.createdAt, now); + const partType = textOr(input.partType ?? input.type, "text"); + const partId = textOr(input.partId ?? input.id, "") || stableWorkbenchFactId("wpt", { sessionId, messageId, turnId, traceId, partIndex, partType }); + return pruneUndefined({ + ...input, + id: partId, + partId, + messageId, + sessionId, + turnId, + traceId, + partIndex, + partType, + status: textOr(input.status, "unknown"), + projectedSeq, + sourceSeq, + sourceEventId: textOr(input.sourceEventId ?? input.eventId, "") || null, + terminal: Boolean(input.terminal), + sealed: Boolean(input.sealed), + createdAt: timestampOr(input.createdAt, timestamp), + updatedAt: timestamp, + valuesPrinted: false + }); +} + +function normalizeWorkbenchTurnFact(value, requestMeta = {}, now) { + const input = normalizeJsonObject(value); + const sessionId = textOr(input.sessionId ?? requestMeta.sessionId, ""); + const traceId = textOr(input.traceId ?? requestMeta.traceId, "") || null; + const turnId = textOr(input.turnId ?? input.id ?? requestMeta.turnId ?? traceId, ""); + if (!sessionId || !turnId) throw requiredWorkbenchFactError("workbench turn fact", ["sessionId", "turnId"]); + const sourceSeq = nonNegativeInteger(input.sourceSeq); + const projectedSeq = nonNegativeInteger(input.projectedSeq ?? sourceSeq); + const timestamp = timestampOr(input.updatedAt ?? input.createdAt, now); + return pruneUndefined({ + ...input, + id: turnId, + turnId, + sessionId, + traceId, + messageId: textOr(input.messageId ?? requestMeta.messageId, "") || null, + status: textOr(input.status, "unknown"), + projectedSeq, + sourceSeq, + sourceEventId: textOr(input.sourceEventId ?? input.eventId, "") || null, + terminal: Boolean(input.terminal), + sealed: Boolean(input.sealed), + finalResponse: input.finalResponse ?? null, + diagnostic: normalizeJsonObject(input.diagnostic ?? input.projection ?? input.blocker), + createdAt: timestampOr(input.createdAt, timestamp), + updatedAt: timestamp, + valuesPrinted: false + }); +} + +function normalizeWorkbenchTraceEventFact(value, requestMeta = {}, now) { + const input = normalizeJsonObject(value); + const traceId = textOr(input.traceId ?? requestMeta.traceId, ""); + if (!traceId) throw requiredWorkbenchFactError("workbench trace event fact", ["traceId"]); + const sessionId = textOr(input.sessionId ?? requestMeta.sessionId, "") || null; + const turnId = textOr(input.turnId ?? requestMeta.turnId ?? traceId, "") || null; + const sourceSeq = nonNegativeInteger(input.sourceSeq ?? input.seq); + const projectedSeq = nonNegativeInteger(input.projectedSeq ?? sourceSeq); + const occurredAt = timestampOr(input.occurredAt ?? input.createdAt ?? input.timestamp, now); + const eventType = textOr(input.eventType ?? input.type ?? input.kind, "event"); + const id = textOr(input.id ?? input.eventId, "") || stableWorkbenchFactId("wte", { traceId, sourceSeq, projectedSeq, eventType, sourceEventId: input.sourceEventId ?? null, occurredAt }); + return pruneUndefined({ + ...input, + id, + traceId, + sessionId, + turnId, + messageId: textOr(input.messageId ?? requestMeta.messageId, "") || null, + sourceSeq, + sourceEventId: textOr(input.sourceEventId ?? input.eventId, "") || null, + projectedSeq, + eventType, + terminal: Boolean(input.terminal), + sealed: Boolean(input.sealed), + occurredAt, + updatedAt: timestampOr(input.updatedAt, occurredAt), + valuesPrinted: false + }); +} + +function normalizeWorkbenchProjectionCheckpoint(value, requestMeta = {}, now) { + const input = normalizeJsonObject(value); + const traceId = textOr(input.traceId ?? requestMeta.traceId, ""); + if (!traceId) throw requiredWorkbenchFactError("workbench projection checkpoint", ["traceId"]); + const sourceSeq = nonNegativeInteger(input.sourceSeq ?? input.lastSourceSeq ?? input.lastAgentRunSeq); + const projectedSeq = nonNegativeInteger(input.projectedSeq ?? input.lastProjectedSeq ?? sourceSeq); + const timestamp = timestampOr(input.updatedAt ?? input.lastProjectedAt, now); + return pruneUndefined({ + ...input, + id: traceId, + traceId, + sessionId: textOr(input.sessionId ?? requestMeta.sessionId, "") || null, + turnId: textOr(input.turnId ?? requestMeta.turnId ?? traceId, "") || null, + runId: textOr(input.runId ?? input.sourceRunId, "") || null, + commandId: textOr(input.commandId ?? input.sourceCommandId, "") || null, + projectedSeq, + sourceSeq, + sourceEventId: textOr(input.sourceEventId ?? input.eventId, "") || null, + projectionStatus: projectionStatusOr(input.projectionStatus, "projecting"), + projectionHealth: projectionHealthOr(input.projectionHealth, "healthy"), + terminal: Boolean(input.terminal), + sealed: Boolean(input.sealed), + diagnostic: normalizeJsonObject(input.diagnostic ?? input.blocker), + createdAt: timestampOr(input.createdAt, timestamp), + updatedAt: timestamp, + valuesPrinted: false + }); +} + +function arrayInput(value) { + if (Array.isArray(value)) return value; + if (value && typeof value === "object") return [value]; + return []; +} + +function stableWorkbenchFactId(prefix, value) { + return `${prefix}_${sha256Hex(value).slice(0, 32)}`; +} + +function requiredWorkbenchFactError(label, fields) { + return new HwlabProtocolError(`${label} requires ${fields.join(", ")}`, { + code: ERROR_CODES.invalidParams, + data: { required: fields } + }); +} + function compareTraceEventRecords(left, right) { const leftSeq = positiveIntegerOrNull(left?.event?.seq); const rightSeq = positiveIntegerOrNull(right?.event?.seq); @@ -1643,6 +2030,24 @@ function compareProjectionStateRecords(left, right) { String(left?.traceId ?? "").localeCompare(String(right?.traceId ?? "")); } +function compareWorkbenchFactRecords(left, right) { + return String(left?.updatedAt ?? left?.createdAt ?? "").localeCompare(String(right?.updatedAt ?? right?.createdAt ?? "")) || + String(left?.id ?? left?.sessionId ?? left?.messageId ?? left?.turnId ?? left?.traceId ?? "").localeCompare(String(right?.id ?? right?.sessionId ?? right?.messageId ?? right?.turnId ?? right?.traceId ?? "")); +} + +function compareWorkbenchPartFactRecords(left, right) { + return String(left?.messageId ?? "").localeCompare(String(right?.messageId ?? "")) || + nonNegativeInteger(left?.partIndex) - nonNegativeInteger(right?.partIndex) || + compareWorkbenchFactRecords(left, right); +} + +function compareWorkbenchTraceEventFacts(left, right) { + return nonNegativeInteger(left?.sourceSeq) - nonNegativeInteger(right?.sourceSeq) || + nonNegativeInteger(left?.projectedSeq) - nonNegativeInteger(right?.projectedSeq) || + String(left?.occurredAt ?? "").localeCompare(String(right?.occurredAt ?? "")) || + String(left?.id ?? "").localeCompare(String(right?.id ?? "")); +} + function traceEventLevel(value) { const text = textOr(value, "info").toLowerCase(); if (["failed", "failure", "error", "timeout"].includes(text)) return "error"; @@ -2179,7 +2584,13 @@ function snapshotMemory(memory) { auditEvents: new Map(memory.auditEvents), evidenceRecords: new Map(memory.evidenceRecords), agentTraceEvents: new Map(memory.agentTraceEvents), - workbenchProjectionStates: new Map(memory.workbenchProjectionStates) + workbenchProjectionStates: new Map(memory.workbenchProjectionStates), + workbenchSessions: new Map(memory.workbenchSessions), + workbenchMessages: new Map(memory.workbenchMessages), + workbenchParts: new Map(memory.workbenchParts), + workbenchTurns: new Map(memory.workbenchTurns), + workbenchTraceEvents: new Map(memory.workbenchTraceEvents), + workbenchProjectionCheckpoints: new Map(memory.workbenchProjectionCheckpoints) }; } @@ -2226,6 +2637,12 @@ function toCountKey(table) { if (table === "worker_sessions") return "workerSessions"; if (table === "agent_trace_events") return "agentTraceEvents"; if (table === "workbench_projection_state") return "workbenchProjectionStates"; + if (table === "workbench_sessions") return "workbenchSessions"; + if (table === "workbench_messages") return "workbenchMessages"; + if (table === "workbench_parts") return "workbenchParts"; + if (table === "workbench_turns") return "workbenchTurns"; + if (table === "workbench_trace_events") return "workbenchTraceEvents"; + if (table === "workbench_projection_checkpoints") return "workbenchProjectionCheckpoints"; return table; } diff --git a/internal/db/schema.test.ts b/internal/db/schema.test.ts index 37ddb7e4..9cdcefd9 100644 --- a/internal/db/schema.test.ts +++ b/internal/db/schema.test.ts @@ -65,6 +65,23 @@ test("initial migration exposes columns required by the durable runtime adapter" assert.equal(CLOUD_CORE_MIGRATIONS[0].runtimeDurableMigrationTable, CLOUD_RUNTIME_DURABLE_MIGRATION_TABLE); }); +test("initial migration declares Workbench fact backfill sources", async () => { + const sql = await readFile(path.join(repoRoot, "internal/db/migrations/0001_cloud_core_skeleton.sql"), "utf8"); + + for (const [target, source] of [ + ["workbench_sessions", "agent_sessions"], + ["workbench_trace_events", "agent_trace_events"], + ["workbench_turns", "workbench_projection_state"], + ["workbench_projection_checkpoints", "workbench_projection_state"] + ]) { + assert.match(sql, new RegExp(`INSERT INTO ${target} \\(`, "u"), `missing ${target} backfill`); + assert.match(sql, new RegExp(`FROM ${source}\\b`, "u"), `missing ${source} source for ${target}`); + } + + assert.match(sql, /ON CONFLICT \(session_id\) DO NOTHING/u); + assert.match(sql, /ON CONFLICT \(trace_id\) DO NOTHING/u); +}); + test("protocol record guards catch schema drift before runtime writes", () => { assertProtocolRecord("gatewaySession", { gatewaySessionId: "gws_01J00000000000000000000000", diff --git a/internal/db/schema.ts b/internal/db/schema.ts index 5743ed79..794d2d43 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-v3"; +export const CLOUD_RUNTIME_DURABLE_ADAPTER_SCHEMA_VERSION = "runtime-durable-postgres-v4"; export const CLOUD_RUNTIME_DURABLE_MIGRATION_TABLE = "hwlab_schema_migrations"; export const FROZEN_CLOUD_CORE_TABLES = Object.freeze([...TABLES]); @@ -164,6 +164,108 @@ export const CLOUD_RUNTIME_DURABLE_TABLE_COLUMNS = Object.freeze({ "projection_json", "created_at", "updated_at" + ]), + workbench_sessions: Object.freeze([ + "session_id", + "owner_user_id", + "project_id", + "conversation_id", + "thread_id", + "status", + "last_trace_id", + "projected_seq", + "source_seq", + "source_event_id", + "terminal", + "sealed", + "session_json", + "created_at", + "updated_at" + ]), + workbench_messages: Object.freeze([ + "message_id", + "session_id", + "turn_id", + "trace_id", + "role", + "status", + "projected_seq", + "source_seq", + "source_event_id", + "terminal", + "sealed", + "message_json", + "created_at", + "updated_at" + ]), + workbench_parts: Object.freeze([ + "part_id", + "message_id", + "session_id", + "turn_id", + "trace_id", + "part_index", + "part_type", + "status", + "projected_seq", + "source_seq", + "source_event_id", + "terminal", + "sealed", + "part_json", + "created_at", + "updated_at" + ]), + workbench_turns: Object.freeze([ + "turn_id", + "session_id", + "trace_id", + "message_id", + "status", + "projected_seq", + "source_seq", + "source_event_id", + "terminal", + "sealed", + "final_response_json", + "diagnostic_json", + "turn_json", + "created_at", + "updated_at" + ]), + workbench_trace_events: Object.freeze([ + "id", + "trace_id", + "session_id", + "turn_id", + "message_id", + "source_seq", + "source_event_id", + "projected_seq", + "event_type", + "terminal", + "sealed", + "event_json", + "occurred_at", + "updated_at" + ]), + workbench_projection_checkpoints: Object.freeze([ + "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" ]) });