feat: add durable workbench facts store
This commit is contained in:
@@ -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"}'
|
||||
)
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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",
|
||||
|
||||
+103
-1
@@ -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"
|
||||
])
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user