Files
pikasTech-HWLAB/internal/db/runtime-store.test.ts
T

1803 lines
79 KiB
TypeScript

import assert from "node:assert/strict";
import { test } from "bun:test";
import { handleJsonRpcRequest } from "../cloud/json-rpc.ts";
import { ERROR_CODES, validateResponse } from "../protocol/index.mjs";
import {
RUNTIME_DURABLE_ADAPTER_AUTH_BLOCKED,
RUNTIME_DURABLE_ADAPTER_MIGRATION_BLOCKED,
RUNTIME_DURABLE_ADAPTER_QUERY_BLOCKED,
RUNTIME_DURABILITY_REQUIRED_EVIDENCE,
RUNTIME_DURABLE_ADAPTER_SCHEMA_BLOCKED,
RUNTIME_DURABLE_ADAPTER_SSL_BLOCKED,
RUNTIME_STORE_KIND_POSTGRES,
buildPostgresPoolConfig,
createCloudRuntimeStore,
createConfiguredCloudRuntimeStore
} from "./runtime-store.ts";
import {
CLOUD_CORE_MIGRATION_ID,
CLOUD_RUNTIME_DURABLE_ADAPTER_SCHEMA_VERSION,
CLOUD_RUNTIME_DURABLE_TABLE_COLUMNS
} from "./schema.ts";
test("memory runtime store is never marked durable", () => {
const store = createCloudRuntimeStore();
const summary = store.summary();
assert.equal(summary.adapter, "memory");
assert.equal(summary.durable, false);
assert.equal(summary.status, "degraded");
assert.equal(summary.ready, undefined);
assert.equal(summary.adapterContract.secretMaterialRequiredInSource, false);
assert.equal(summary.durabilityContract.ready, false);
assert.equal(summary.durabilityContract.status, "blocked");
assert.equal(summary.durabilityContract.blockedLayer, "adapter");
assert.equal(summary.durabilityContract.adapterQueryRequired, true);
assert.equal(summary.durabilityContract.dbLiveEvidenceIsDurabilityEvidence, false);
assert.equal(summary.durabilityContract.secretMaterialRead, false);
});
test("workbench session input facts are queryable and idempotent in session aggregate order", () => {
const store = createCloudRuntimeStore({ now: () => "2026-06-25T20:00:00.000Z" });
const sessionId = "ses_input_schema";
const first = store.writeWorkbenchFacts({
facts: {
inputs: [{
sessionId,
turnId: "turn_input_schema_1",
traceId: "trc_input_schema_1",
messageId: "msg_input_schema_user_1",
delivery: "queue",
status: "admitted",
sourceEventId: "input:queue:1"
}]
}
});
const duplicate = store.writeWorkbenchFacts({ facts: { inputs: [first.facts.inputs[0]] } });
const second = store.writeWorkbenchFacts({
facts: {
inputs: [{
sessionId,
turnId: "turn_input_schema_1",
traceId: "trc_input_schema_1",
messageId: "msg_input_schema_user_1",
commandId: "cmd_input_schema_1",
delivery: "steer",
status: "promoted",
sourceEventId: "input:steer:1"
}]
}
});
const third = store.writeWorkbenchFacts({
facts: {
inputs: [{
sessionId,
turnId: "turn_input_schema_1",
traceId: "trc_input_schema_1",
messageId: "msg_input_schema_user_1",
commandId: "cmd_input_schema_1",
delivery: "cancel",
status: "promoted",
sourceEventId: "input:cancel:1"
}]
}
});
assert.equal(first.facts.inputs[0].admittedSeq, 1);
assert.equal(duplicate.facts.inputs[0].admittedSeq, 1);
assert.equal(second.facts.inputs[0].admittedSeq, 2);
assert.equal(second.facts.inputs[0].promotedSeq, 2);
assert.equal(third.facts.inputs[0].admittedSeq, 3);
assert.equal(third.facts.inputs[0].promotedSeq, 3);
const queried = store.queryWorkbenchFacts({ sessionId, families: ["inputs"] });
assert.equal(queried.facts.inputs.length, 3);
assert.deepEqual(queried.facts.inputs.map((input) => [input.delivery, input.status, input.admittedSeq, input.promotedSeq]), [
["queue", "admitted", 1, null],
["steer", "promoted", 2, 2],
["cancel", "promoted", 3, 3]
]);
assert.equal(queried.facts.inputs[1].commandId, "cmd_input_schema_1");
});
test("configured postgres runtime classifies auth blocker before schema and migration", async () => {
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
dbUrl: "postgres://hwlab_redacted@db.example.invalid:5432/hwlab",
queryClient: {
async query(sql) {
assert.match(sql, /information_schema\.columns/u);
const error = new Error("auth failed");
error.code = "28P01";
throw error;
}
}
});
const readiness = await store.readiness();
assert.equal(readiness.adapter, RUNTIME_STORE_KIND_POSTGRES);
assert.equal(readiness.durable, false);
assert.equal(readiness.ready, false);
assert.equal(readiness.status, "blocked");
assert.equal(readiness.blocker, RUNTIME_DURABLE_ADAPTER_AUTH_BLOCKED);
assert.equal(readiness.connection.queryAttempted, true);
assert.equal(readiness.connection.queryResult, "auth_blocked");
assert.equal(readiness.gates.auth.status, "blocked");
assert.equal(readiness.gates.schema.status, "not_checked");
assert.equal(readiness.gates.migration.status, "not_checked");
assert.equal(readiness.durabilityContract.ready, false);
assert.equal(readiness.durabilityContract.blockedLayer, "auth");
assert.equal(readiness.durabilityContract.dbLiveEvidenceIsDurabilityEvidence, false);
});
test("configured postgres runtime classifies schema driver errors separately from query failures", async () => {
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
queryClient: {
async query(sql) {
assert.match(sql, /information_schema\.columns/u);
const error = new Error("schema missing");
error.code = "3F000";
throw error;
}
}
});
const readiness = await store.readiness();
assert.equal(readiness.blocker, RUNTIME_DURABLE_ADAPTER_SCHEMA_BLOCKED);
assert.equal(readiness.connection.queryAttempted, true);
assert.equal(readiness.connection.queryResult, "schema_blocked");
assert.equal(readiness.gates.ssl.status, "ready");
assert.equal(readiness.gates.auth.status, "ready");
assert.equal(readiness.gates.schema.status, "blocked");
assert.equal(readiness.gates.migration.status, "not_checked");
assert.equal(readiness.durabilityContract.blockedLayer, "schema");
});
test("configured postgres runtime classifies SSL negotiation blocker separately from auth", async () => {
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
queryClient: {
async query(sql) {
assert.match(sql, /information_schema\.columns/u);
const error = new Error("The server does not support SSL connections");
error.code = "08P01";
throw error;
}
}
});
const readiness = await store.readiness();
assert.equal(readiness.adapter, RUNTIME_STORE_KIND_POSTGRES);
assert.equal(readiness.durable, false);
assert.equal(readiness.ready, false);
assert.equal(readiness.status, "blocked");
assert.equal(readiness.blocker, RUNTIME_DURABLE_ADAPTER_SSL_BLOCKED);
assert.equal(readiness.connection.queryAttempted, true);
assert.equal(readiness.connection.queryResult, "ssl_negotiation_blocked");
assert.equal(readiness.connection.errorCode, "08P01");
assert.equal(readiness.gates.ssl.status, "blocked");
assert.equal(readiness.gates.auth.status, "not_checked");
assert.equal(readiness.gates.schema.status, "not_checked");
assert.equal(readiness.gates.migration.status, "not_checked");
assert.equal(readiness.durabilityContract.ready, false);
assert.equal(readiness.durabilityContract.blockedLayer, "ssl");
assert.equal(readiness.durabilityContract.secretMaterialRead, false);
assert.equal(JSON.stringify(readiness).includes("The server does not support SSL connections"), false);
});
test("postgres pool config makes HWLAB_CLOUD_DB_SSL_MODE authoritative over URL sslmode", () => {
const disabled = buildPostgresPoolConfig({
dbUrl: "postgres://hwlab_user:fixture-pass@db.example.test:5432/hwlab?sslmode=require&application_name=hwlab",
sslMode: "disable",
timeoutMs: "2500",
poolMax: "20"
});
assert.equal(disabled.ssl, false);
assert.equal(disabled.connectionTimeoutMillis, 2500);
assert.equal(disabled.max, 20);
const disabledUrl = new URL(disabled.connectionString);
assert.equal(disabledUrl.searchParams.get("sslmode"), null);
assert.equal(disabledUrl.searchParams.get("ssl"), null);
assert.equal(disabledUrl.searchParams.get("application_name"), "hwlab");
const required = buildPostgresPoolConfig({
dbUrl: "postgres://hwlab_user:fixture-pass@db.example.test:5432/hwlab?sslmode=disable",
sslMode: "require"
});
assert.deepEqual(required.ssl, { rejectUnauthorized: false });
assert.equal(new URL(required.connectionString).searchParams.get("sslmode"), "require");
assert.equal(new URL(required.connectionString).searchParams.get("uselibpqcompat"), "true");
});
test("postgres pool config inherits URL sslmode when env override is absent", () => {
const disabled = buildPostgresPoolConfig({
dbUrl: "postgres://hwlab_user:fixture-pass@db.example.test:5432/hwlab?sslmode=disable&application_name=hwlab"
});
assert.equal(disabled.ssl, false);
const disabledUrl = new URL(disabled.connectionString);
assert.equal(disabledUrl.searchParams.get("sslmode"), null);
assert.equal(disabledUrl.searchParams.get("application_name"), "hwlab");
const required = buildPostgresPoolConfig({
dbUrl: "postgres://hwlab_user:fixture-pass@db.example.test:5432/hwlab?sslmode=require"
});
assert.deepEqual(required.ssl, { rejectUnauthorized: false });
assert.equal(new URL(required.connectionString).searchParams.get("sslmode"), "require");
});
test("configured postgres runtime passes normalized pool config to pg", async () => {
const pools = [];
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true",
HWLAB_CLOUD_DB_PROBE_TIMEOUT_MS: "1800",
HWLAB_CLOUD_DB_POOL_MAX: "20"
},
dbUrl: "postgres://hwlab_user:fixture-pass@db.example.test:5432/hwlab?sslmode=require",
sslMode: "disable",
pgModuleLoader: async () => ({
Pool: class FakePool {
constructor(config) {
this.config = config;
pools.push(this);
}
async query() {
return { rows: [] };
}
}
})
});
await store.readiness();
assert.equal(pools.length, 1);
assert.equal(pools[0].config.ssl, false);
assert.equal(pools[0].config.connectionTimeoutMillis, 1800);
assert.equal(pools[0].config.max, 20);
assert.equal(new URL(pools[0].config.connectionString).searchParams.get("sslmode"), null);
assert.equal(store.summary().blocker, RUNTIME_DURABLE_ADAPTER_SCHEMA_BLOCKED);
assert.equal(store.summary().durabilityContract.blockedLayer, "schema");
});
test("configured postgres runtime consumes pg pool idle errors as degraded readiness", async () => {
const pools = [];
const warnings = [];
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true",
HWLAB_CLOUD_DB_PROBE_TIMEOUT_MS: "1800"
},
dbUrl: "postgres://hwlab_user:fixture-pass@db.example.test:5432/hwlab?sslmode=require",
logger: { warn(entry) { warnings.push(entry); } },
pgModuleLoader: async () => ({
Pool: class FakePool {
constructor(config) {
this.config = config;
this.handlers = new Map();
pools.push(this);
}
on(event, handler) {
this.handlers.set(event, handler);
return this;
}
async query() {
return { rows: [] };
}
emit(event, error) {
this.handlers.get(event)?.(error);
}
}
})
});
await store.readiness();
const error = new Error("connection timed out after fixture secret host");
error.code = "ETIMEDOUT";
pools[0].emit("error", error);
const summary = store.summary();
assert.equal(summary.blocker, RUNTIME_DURABLE_ADAPTER_QUERY_BLOCKED);
assert.equal(summary.connection.queryResult, "query_blocked");
assert.equal(summary.connection.errorCode, "ETIMEDOUT");
assert.equal(summary.durabilityContract.blockedLayer, "durability_query");
assert.equal(warnings.length, 1);
assert.equal(warnings[0].event, "postgres_runtime_pool_error");
assert.equal(warnings[0].valuesRedacted, true);
assert.equal(JSON.stringify(warnings).includes("fixture secret host"), false);
});
test("configured postgres runtime classifies pg_hba no-encryption rejection as SSL before auth", async () => {
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
sslMode: "disable",
queryClient: {
async query(sql) {
assert.match(sql, /information_schema\.columns/u);
const error = new Error('no pg_hba.conf entry for host "[redacted]", user "[redacted]", database "[redacted]", no encryption');
error.code = "28000";
throw error;
}
}
});
const readiness = await store.readiness();
assert.equal(readiness.blocker, RUNTIME_DURABLE_ADAPTER_SSL_BLOCKED);
assert.equal(readiness.connection.queryAttempted, true);
assert.equal(readiness.connection.queryResult, "ssl_negotiation_blocked");
assert.equal(readiness.connection.errorCode, "28000");
assert.equal(readiness.gates.ssl.status, "blocked");
assert.equal(readiness.gates.auth.status, "not_checked");
assert.equal(readiness.durabilityContract.blockedLayer, "ssl");
assert.equal(readiness.durabilityContract.dbLiveEvidenceIsDurabilityEvidence, false);
assert.equal(JSON.stringify(readiness).includes("pg_hba.conf"), false);
});
test("configured postgres runtime reports schema blocker without green readiness", async () => {
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
dbUrl: "postgres://hwlab_redacted@db.example.invalid:5432/hwlab",
queryClient: {
async query(sql) {
assert.match(sql, /information_schema\.columns/u);
return { rows: [] };
}
}
});
const readiness = await store.readiness();
assert.equal(readiness.adapter, RUNTIME_STORE_KIND_POSTGRES);
assert.equal(readiness.durable, false);
assert.equal(readiness.durableRequested, true);
assert.equal(readiness.ready, false);
assert.equal(readiness.status, "blocked");
assert.equal(readiness.blocker, RUNTIME_DURABLE_ADAPTER_SCHEMA_BLOCKED);
assert.equal(readiness.liveRuntimeEvidence, false);
assert.equal(readiness.fixtureEvidence, false);
assert.equal(readiness.schema.checked, true);
assert.ok(readiness.schema.missingTables.includes("gateway_sessions"));
assert.equal(readiness.migration.checked, false);
assert.equal(readiness.gates.schema.status, "blocked");
assert.equal(readiness.gates.migration.status, "not_checked");
assert.equal(readiness.durabilityContract.ready, false);
assert.equal(readiness.durabilityContract.blockedLayer, "schema");
assert.equal(readiness.durabilityContract.dbLiveEvidenceIsDurabilityEvidence, false);
assert.equal(JSON.stringify(readiness).includes("redacted@db.example"), false);
});
test("configured postgres runtime keeps durable false when migration ledger is missing", async () => {
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
dbUrl: "postgres://hwlab_redacted@db.example.invalid:5432/hwlab",
queryClient: createFakePostgresClient({ migrationReady: false })
});
const readiness = await store.readiness();
assert.equal(readiness.adapter, RUNTIME_STORE_KIND_POSTGRES);
assert.equal(readiness.durable, false);
assert.equal(readiness.durableRequested, true);
assert.equal(readiness.ready, false);
assert.equal(readiness.status, "blocked");
assert.equal(readiness.blocker, RUNTIME_DURABLE_ADAPTER_MIGRATION_BLOCKED);
assert.equal(readiness.schema.ready, true);
assert.equal(readiness.migration.checked, true);
assert.equal(readiness.migration.ready, false);
assert.equal(readiness.migration.missing, true);
assert.equal(readiness.gates.schema.status, "ready");
assert.equal(readiness.gates.migration.status, "blocked");
assert.equal(readiness.liveRuntimeEvidence, false);
assert.equal(readiness.durabilityContract.ready, false);
assert.equal(readiness.durabilityContract.blockedLayer, "migration");
assert.equal(readiness.durabilityContract.dbLiveEvidenceIsDurabilityEvidence, false);
});
test("configured postgres runtime classifies migration query failures separately from generic readiness queries", async () => {
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
queryClient: createFakePostgresClient({ migrationErrorCode: "57014" })
});
const readiness = await store.readiness();
assert.equal(readiness.blocker, RUNTIME_DURABLE_ADAPTER_MIGRATION_BLOCKED);
assert.equal(readiness.schema.ready, true);
assert.equal(readiness.migration.checked, true);
assert.equal(readiness.migration.errorCode, "57014");
assert.equal(readiness.connection.queryAttempted, true);
assert.equal(readiness.connection.queryResult, "migration_blocked");
assert.equal(readiness.gates.ssl.status, "ready");
assert.equal(readiness.gates.auth.status, "ready");
assert.equal(readiness.gates.schema.status, "ready");
assert.equal(readiness.gates.migration.status, "blocked");
assert.equal(readiness.gates.durability.status, "not_checked");
assert.equal(readiness.durabilityContract.blockedLayer, "migration");
});
test("configured postgres runtime classifies durable read query blocker after schema and migration", async () => {
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
dbUrl: "postgres://hwlab_redacted@db.example.invalid:5432/hwlab",
queryClient: createFakePostgresClient({ migrationReady: true, countErrorCode: "57014" })
});
const readiness = await store.readiness();
assert.equal(readiness.adapter, RUNTIME_STORE_KIND_POSTGRES);
assert.equal(readiness.durable, false);
assert.equal(readiness.ready, false);
assert.equal(readiness.status, "blocked");
assert.equal(readiness.blocker, RUNTIME_DURABLE_ADAPTER_QUERY_BLOCKED);
assert.equal(readiness.schema.ready, true);
assert.equal(readiness.migration.ready, true);
assert.equal(readiness.connection.queryAttempted, true);
assert.equal(readiness.connection.queryResult, "query_blocked");
assert.equal(readiness.gates.auth.status, "ready");
assert.equal(readiness.gates.schema.status, "ready");
assert.equal(readiness.gates.migration.status, "ready");
assert.equal(readiness.gates.durability.status, "blocked");
assert.equal(readiness.durabilityContract.ready, false);
assert.equal(readiness.durabilityContract.blockedLayer, "durability_query");
assert.equal(readiness.durabilityContract.dbLiveEvidenceIsDurabilityEvidence, false);
});
test("JSON-RPC audit/evidence durable queries return blocked errors instead of empty success", async () => {
const secretDbUrl = "postgres://hwlab_user:super-secret-password@db.example.invalid:5432/hwlab";
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
dbUrl: secretDbUrl,
queryClient: createFakePostgresClient({ migrationReady: true, countErrorCode: "57014" })
});
const context = { runtimeStore: store, dbProbe: { probe: false } };
const meta = {
traceId: "trc_01J00000000000000000002000",
actorId: "usr_01J00000000000000000002000",
serviceId: "hwlab-cloud-web",
environment: "dev"
};
for (const [id, method] of [
["req_01J00000000000000000002001", "audit.event.query"],
["req_01J00000000000000000002002", "evidence.record.query"]
]) {
const response = await handleJsonRpcRequest(
{
jsonrpc: "2.0",
id,
method,
params: {
projectId: "prj_01J00000000000000000002000"
},
meta
},
context
);
validateResponse(response);
assert.equal(Object.hasOwn(response, "result"), false);
assert.equal(response.error.code, ERROR_CODES.internalError);
assert.equal(response.error.data.method, method);
assert.equal(response.error.data.adapter, "postgres");
assert.equal(response.error.data.status, "blocked");
assert.equal(response.error.data.blocker, RUNTIME_DURABLE_ADAPTER_QUERY_BLOCKED);
assert.equal(response.error.data.durable, false);
assert.equal(response.error.data.durableRequested, true);
assert.equal(response.error.data.liveRuntimeEvidence, false);
assert.equal(response.error.data.queryResult, "query_blocked");
assert.equal(response.error.data.blockedLayer, "durability_query");
assert.equal(response.error.data.requiredEvidence, RUNTIME_DURABILITY_REQUIRED_EVIDENCE);
assert.equal(response.error.data.dbLiveEvidenceIsDurabilityEvidence, false);
assert.equal(response.error.data.secretMaterialRead, false);
assert.equal(response.error.data.valueRedacted, true);
assert.equal(response.error.data.endpointRedacted, true);
assert.equal(Object.hasOwn(response.error.data, "events"), false);
assert.equal(Object.hasOwn(response.error.data, "records"), false);
assert.equal(Object.hasOwn(response.error.data, "count"), false);
const serialized = JSON.stringify(response);
assert.equal(serialized.includes("super-secret-password"), false);
assert.equal(serialized.includes(secretDbUrl), false);
}
});
test("JSON-RPC audit/evidence table read failures return blocked errors", async () => {
const queryClient = createFakePostgresClient({ migrationReady: true, readErrorCode: "57014" });
const store = createConfiguredCloudRuntimeStore({
env: {
HWLAB_CLOUD_RUNTIME_ADAPTER: "postgres",
HWLAB_CLOUD_RUNTIME_DURABLE: "true"
},
dbUrl: "postgres://hwlab_user:another-secret@db.example.invalid:5432/hwlab",
queryClient
});
const context = { runtimeStore: store, dbProbe: { probe: false } };
const meta = {
traceId: "trc_01J00000000000000000002200",
actorId: "usr_01J00000000000000000002200",
serviceId: "hwlab-cloud-web",
environment: "dev"
};
for (const [id, method] of [
["req_01J00000000000000000002201", "audit.event.query"],
["req_01J00000000000000000002202", "evidence.record.query"]
]) {
const response = await handleJsonRpcRequest(
{
jsonrpc: "2.0",
id,
method,
params: {
projectId: "prj_01J00000000000000000002200"
},
meta
},
context
);
validateResponse(response);
assert.equal(Object.hasOwn(response, "result"), false);
assert.equal(response.error.code, ERROR_CODES.internalError);
assert.equal(response.error.data.method, method);
assert.equal(response.error.data.blocker, RUNTIME_DURABLE_ADAPTER_QUERY_BLOCKED);
assert.equal(response.error.data.queryResult, "query_blocked");
assert.equal(response.error.data.blockedLayer, "durability_query");
const serialized = JSON.stringify(response);
assert.equal(serialized.includes("another-secret"), false);
assert.equal(serialized.includes("events"), false);
assert.equal(serialized.includes("records"), false);
assert.equal(serialized.includes("\"count\":0"), false);
}
});
test("JSON-RPC durable queries return successful empty result sets for true empty tables", 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
});
const context = { runtimeStore: store, dbProbe: { probe: false } };
const meta = {
traceId: "trc_01J00000000000000000002100",
actorId: "usr_01J00000000000000000002100",
serviceId: "hwlab-cloud-web",
environment: "dev"
};
async function rpc(id, method, params = {}) {
const response = await handleJsonRpcRequest(
{
jsonrpc: "2.0",
id,
method,
params,
meta
},
context
);
validateResponse(response);
assert.equal(Object.hasOwn(response, "error"), false, response.error?.message);
return response.result;
}
const audit = await rpc("req_01J00000000000000000002101", "audit.event.query", {
projectId: "prj_01J00000000000000000002100"
});
assert.deepEqual(audit.events, []);
assert.equal(audit.count, 0);
assert.equal(audit.persistence.adapter, "postgres");
assert.equal(audit.persistence.durable, true);
assert.equal(audit.persistence.ready, true);
const evidence = await rpc("req_01J00000000000000000002102", "evidence.record.query", {
projectId: "prj_01J00000000000000000002100"
});
assert.deepEqual(evidence.records, []);
assert.equal(evidence.count, 0);
assert.equal(evidence.persistence.adapter, "postgres");
assert.equal(evidence.persistence.durable, true);
assert.equal(evidence.persistence.ready, true);
});
test("configured postgres runtime persists and queries records through query client", 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-05-22T00:00:00.000Z"
});
const context = { runtimeStore: store, dbProbe: { probe: false } };
const meta = {
traceId: "trc_01J00000000000000000001000",
actorId: "usr_01J00000000000000000001000",
serviceId: "hwlab-cloud-web",
environment: "dev"
};
async function rpc(id, method, params = {}) {
const response = await handleJsonRpcRequest(
{
jsonrpc: "2.0",
id,
method,
params,
meta
},
context
);
validateResponse(response);
assert.equal(Object.hasOwn(response, "error"), false, response.error?.message);
return response.result;
}
const health = await rpc("req_01J00000000000000000001000", "system.health");
assert.equal(health.runtime.adapter, "postgres");
assert.equal(health.runtime.durable, true);
assert.equal(health.runtime.durableRequested, true);
assert.equal(health.runtime.durableCapable, true);
assert.equal(health.runtime.ready, true);
assert.equal(health.runtime.liveRuntimeEvidence, true);
assert.equal(health.runtime.fixtureEvidence, false);
assert.equal(health.runtime.schema.ready, true);
assert.equal(health.runtime.migration.ready, true);
assert.equal(health.runtime.gates.migration.status, "ready");
assert.equal(health.runtime.gates.durability.status, "ready");
assert.equal(health.runtime.durabilityContract.ready, true);
assert.equal(health.runtime.durabilityContract.blockedLayer, null);
assert.equal(health.runtime.durabilityContract.dbLiveEvidenceIsDurabilityEvidence, false);
assert.equal(health.readiness.components.runtime, "ready");
assert.equal(health.readiness.durability.status, "ready");
assert.equal(health.readiness.durability.dbLiveEvidenceIsDurabilityEvidence, false);
await rpc("req_01J00000000000000000001001", "gateway.session.register", {
projectId: "prj_01J00000000000000000001000",
gatewaySessionId: "gws_01J00000000000000000001000",
gatewayId: "gtw_01J00000000000000000001000",
serviceId: "hwlab-gateway",
endpoint: "http://127.0.0.1:7101"
});
await rpc("req_01J00000000000000000001002", "box.resource.register", {
projectId: "prj_01J00000000000000000001000",
gatewaySessionId: "gws_01J00000000000000000001000",
resourceId: "res_01J00000000000000000001000",
boxId: "box_01J00000000000000000001000",
resourceType: "board",
state: "available"
});
await rpc("req_01J00000000000000000001003", "box.capability.report", {
capabilityId: "cap_01J00000000000000000001000",
resourceId: "res_01J00000000000000000001000",
projectId: "prj_01J00000000000000000001000",
name: "shell.exec",
direction: "bidirectional",
valueType: "object",
mutatesState: true
});
const invoke = await rpc("req_01J00000000000000000001004", "hardware.invoke.shell", {
projectId: "prj_01J00000000000000000001000",
gatewaySessionId: "gws_01J00000000000000000001000",
resourceId: "res_01J00000000000000000001000",
capabilityId: "cap_01J00000000000000000001000",
input: {
command: "echo durable"
}
});
const audit = await rpc("req_01J00000000000000000001005", "audit.event.query", {
projectId: "prj_01J00000000000000000001000"
});
const evidence = await rpc("req_01J00000000000000000001006", "evidence.record.query", {
projectId: "prj_01J00000000000000000001000"
});
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("INSERT INTO gateway_sessions")));
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("INSERT INTO evidence_records")));
assert.ok(audit.events.some((event) => event.action === "hardware.invoke.shell"));
assert.equal(evidence.count, 1);
assert.equal(evidence.records[0].operationId, invoke.operationId);
});
test("configured postgres runtime persists and queries Code Agent trace events", 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-17T03:50:00.000Z"
});
const traceId = "trc_01J00000000000000000002000";
const write = await store.writeAgentTraceEvent({
event: {
traceId,
seq: 1,
type: "runner",
status: "completed",
label: "runner:completed",
sessionId: "ses_01J00000000000000000002000",
message: "turn completed",
terminal: true,
createdAt: "2026-06-17T03:50:01.000Z",
valuesPrinted: false
}
});
const events = await store.queryAgentTraceEvents({ traceId });
const readiness = await store.readiness();
assert.equal(write.written, true);
assert.equal(write.traceEvent.traceId, traceId);
assert.equal(write.traceEvent.agentSessionId, "ses_01J00000000000000000002000");
assert.equal(events.count, 1);
assert.equal(events.events[0].traceId, traceId);
assert.equal(events.events[0].terminal, true);
assert.equal(readiness.counts.agentTraceEvents, 1);
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("CREATE INDEX IF NOT EXISTS idx_agent_trace_events_trace_order")));
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("INSERT INTO agent_trace_events")));
});
test("configured postgres runtime persists and queries Workbench projection state", 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-19T13:05:00.000Z"
});
const write = await store.writeWorkbenchProjectionState({
projectionState: {
traceId: "trc_projection_state",
sessionId: "ses_projection_state",
conversationId: "cnv_projection_state",
threadId: "thread-projection-state",
runId: "run_projection_state",
commandId: "cmd_projection_state",
lastSourceSeq: 3700,
lastProjectedSeq: 129,
sourceLatestSeq: 3701,
projectionStatus: "projecting",
resultSyncState: "pending",
updatedAt: "2026-06-19T13:05:00.000Z"
}
});
const loaded = await store.getWorkbenchProjectionState({ traceId: "trc_projection_state" });
const due = await store.queryWorkbenchProjectionStates({ projectionStatuses: ["projecting"], dueAt: "2026-06-19T13:06:00.000Z", limit: 10 });
const readiness = await store.readiness();
assert.equal(write.written, true);
assert.equal(write.projectionState.lastSourceSeq, 3700);
assert.equal(write.projectionState.lastAgentRunSeq, 3700);
assert.equal(loaded.found, true);
assert.equal(loaded.projectionState.commandId, "cmd_projection_state");
assert.equal(due.count, 1);
assert.equal(due.states[0].resultSyncState, "pending");
assert.equal(readiness.counts.workbenchProjectionStates, 1);
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("CREATE INDEX IF NOT EXISTS idx_workbench_projection_state_status_retry_updated")));
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: {
inputs: [{
sessionId: "ses_fact_1",
turnId: "turn_fact_1",
traceId: "trc_fact_1",
messageId: "msg_fact_1_user",
commandId: "cmd_fact_1",
delivery: "steer",
status: "promoted",
sourceEventId: "evt_input_fact_1",
updatedAt: "2026-06-20T10:00:00.500Z"
}],
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.inputs[0].delivery, "steer");
assert.equal(write.facts.inputs[0].status, "promoted");
assert.equal(write.facts.inputs[0].promotedSeq, write.facts.inputs[0].admittedSeq);
assert.equal(write.facts.sessions[0].terminal, false);
assert.equal(loaded.count, 7);
assert.equal(loaded.facts.inputs[0].delivery, "steer");
assert.equal(loaded.facts.inputs[0].promotedSeq, loaded.facts.inputs[0].admittedSeq);
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_session_inputs")));
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 appends Workbench aggregate events and outbox in one transaction", 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-24T12:00:00.000Z"
});
const write = await store.writeWorkbenchFacts({ facts: {
traceEvents: [{ id: "wte_event_stream_1", traceId: "trc_event_stream", sessionId: "ses_event_stream", turnId: "turn_event_stream", sourceSeq: 1, sourceEventId: "source-event-1", projectedSeq: 1, eventType: "assistant", occurredAt: "2026-06-24T12:00:01.000Z" }],
turns: [{ turnId: "turn_event_stream", sessionId: "ses_event_stream", traceId: "trc_event_stream", status: "completed", sourceSeq: 2, sourceEventId: "source-terminal-2", projectedSeq: 2, terminal: true, sealed: true, updatedAt: "2026-06-24T12:00:02.000Z" }]
} });
const outbox = await store.readWorkbenchProjectionOutbox({ afterSeq: 0, traceId: "trc_event_stream", limit: 10 });
assert.equal(write.events.length, 2);
assert.deepEqual(write.events.map((event) => event.aggregateId), ["session:ses_event_stream", "session:ses_event_stream"]);
assert.deepEqual(write.events.map((event) => event.aggregateSeq), [1, 2]);
assert.equal(outbox.length, 2);
assert.equal(outbox[0].eventSeq, write.events[0].eventSeq);
assert.equal(outbox[0].aggregateId, "session:ses_event_stream");
assert.equal(outbox[1].commitType, "terminal");
assert.equal(outbox[1].aggregateSeq, 2);
assert.ok(queryClient.calls.some((call) => call.sql === "BEGIN"));
assert.ok(queryClient.calls.some((call) => call.sql === "COMMIT"));
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("INSERT INTO workbench_events")));
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("INSERT INTO workbench_projection_outbox")));
});
test("configured postgres runtime reuses aggregate event sequence for duplicate source event", 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-24T12:10:00.000Z"
});
const facts = { traceEvents: [{ id: "wte_duplicate_source", traceId: "trc_duplicate_source", sessionId: "ses_duplicate_source", turnId: "turn_duplicate_source", sourceSeq: 1, sourceEventId: "source-dup", projectedSeq: 1, eventType: "assistant", occurredAt: "2026-06-24T12:10:01.000Z" }] };
const first = await store.writeWorkbenchFacts({ facts });
const repeated = await store.writeWorkbenchFacts({ facts });
assert.equal(first.events[0].eventSeq, repeated.events[0].eventSeq);
assert.equal(first.events[0].aggregateSeq, repeated.events[0].aggregateSeq);
assert.equal(first.events[0].projectionRevision, repeated.events[0].projectionRevision);
});
test("configured postgres runtime rolls back Workbench facts when outbox append fails", async () => {
const queryClient = createFakePostgresClient({ migrationReady: true, outboxErrorCode: "XX000" });
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-24T12:20:00.000Z"
});
await assert.rejects(
() => store.writeWorkbenchFacts({ facts: { traceEvents: [{ id: "wte_outbox_fail", traceId: "trc_outbox_fail", sessionId: "ses_outbox_fail", turnId: "turn_outbox_fail", sourceSeq: 1, sourceEventId: "source-outbox-fail", projectedSeq: 1, eventType: "assistant", occurredAt: "2026-06-24T12:20:01.000Z" }] } }),
/projection outbox insert failed/u
);
assert.ok(queryClient.calls.some((call) => call.sql === "BEGIN"));
assert.ok(queryClient.calls.some((call) => call.sql === "ROLLBACK"));
assert.equal(queryClient.calls.some((call) => call.sql === "COMMIT"), false);
});
test("configured postgres runtime allocates Workbench projectedSeq idempotently by source event", 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-20T12:30:00.000Z"
});
const traceId = "trc_pg_projected_seq_idempotent";
const first = await store.allocateWorkbenchProjectedSeq({ traceId, sessionId: "ses_pg_projected_seq_idempotent", sourceSeq: 1, sourceEventId: "source-a" });
await store.writeWorkbenchFacts({ facts: {
traceEvents: [{ id: "wte_pg_source_a", traceId, sessionId: "ses_pg_projected_seq_idempotent", sourceSeq: 1, sourceEventId: "source-a", projectedSeq: first.projectedSeq, eventType: "backend", occurredAt: "2026-06-20T12:30:01.000Z" }],
checkpoints: [{ traceId, sessionId: "ses_pg_projected_seq_idempotent", projectedSeq: first.projectedSeq, sourceSeq: 1, sourceEventId: "source-a", projectionStatus: "projecting", projectionHealth: "healthy", updatedAt: "2026-06-20T12:30:01.000Z" }]
} });
const second = await store.allocateWorkbenchProjectedSeq({ traceId, sessionId: "ses_pg_projected_seq_idempotent", sourceSeq: 1, sourceEventId: "source-b" });
await store.writeWorkbenchFacts({ facts: {
traceEvents: [{ id: "wte_pg_source_b", traceId, sessionId: "ses_pg_projected_seq_idempotent", sourceSeq: 1, sourceEventId: "source-b", projectedSeq: second.projectedSeq, eventType: "backend", occurredAt: "2026-06-20T12:30:02.000Z" }],
checkpoints: [{ traceId, sessionId: "ses_pg_projected_seq_idempotent", projectedSeq: second.projectedSeq, sourceSeq: 1, sourceEventId: "source-b", projectionStatus: "projecting", projectionHealth: "healthy", updatedAt: "2026-06-20T12:30:02.000Z" }]
} });
const repeated = await store.allocateWorkbenchProjectedSeq({ traceId, sessionId: "ses_pg_projected_seq_idempotent", sourceSeq: 1, sourceEventId: "source-a" });
await store.writeWorkbenchFacts({ facts: { checkpoints: [{ traceId, sessionId: "ses_pg_projected_seq_idempotent", projectedSeq: repeated.projectedSeq, sourceSeq: 1, sourceEventId: "source-a", projectionStatus: "projecting", projectionHealth: "healthy", updatedAt: "2026-06-20T12:30:03.000Z" }] } });
const loaded = await store.queryWorkbenchFacts({ traceId, families: ["traceEvents", "checkpoints"], limit: 10 });
assert.equal(first.projectedSeq, 1);
assert.equal(second.projectedSeq, 2);
assert.equal(repeated.projectedSeq, 1);
assert.deepEqual(loaded.facts.traceEvents.map((event) => event.projectedSeq), [1, 2]);
assert.equal(loaded.facts.checkpoints[0].projectedSeq, 2);
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("SELECT pg_advisory_xact_lock")));
});
test("configured postgres runtime writes Workbench facts without re-running blocked read readiness", async () => {
const queryClient = createFakePostgresClient({ migrationReady: true, countErrorCode: "53300" });
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-25T02:50:00.000Z"
});
const readiness = await store.readiness();
assert.equal(readiness.blocker, RUNTIME_DURABLE_ADAPTER_QUERY_BLOCKED);
assert.equal(readiness.connection.queryResult, "query_blocked");
queryClient.calls.length = 0;
const write = await store.writeWorkbenchFacts({
facts: {
sessions: [{ sessionId: "ses_write_after_read_blocked", ownerUserId: "usr_write_after_read_blocked", status: "running", updatedAt: "2026-06-25T02:50:01.000Z" }]
}
});
assert.equal(write.written, true);
assert.equal(queryClient.calls.filter((call) => isReadinessQuery(call.sql)).length, 0);
assert.equal(queryClient.calls.some((call) => call.sql.startsWith("CREATE INDEX IF NOT EXISTS")), false);
assert.ok(queryClient.calls.some((call) => call.sql === "BEGIN"));
assert.ok(queryClient.calls.some((call) => call.sql === "COMMIT"));
assert.ok(queryClient.calls.some((call) => call.sql.startsWith("INSERT INTO workbench_sessions")));
});
test("configured postgres runtime reuses proven readiness for Workbench durable reads", 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"
});
await store.writeWorkbenchFacts({
facts: {
sessions: [{ sessionId: "ses_ready_cached", ownerUserId: "usr_ready", status: "running", updatedAt: "2026-06-20T10:00:01.000Z" }]
}
});
await store.readiness();
queryClient.calls.length = 0;
const first = await store.queryWorkbenchFacts({ families: ["sessions"], limit: 10 });
const second = await store.queryWorkbenchFacts({ families: ["sessions"], limit: 10 });
assert.equal(first.count, 1);
assert.equal(second.count, 1);
assert.equal(queryClient.calls.filter((call) => isReadinessQuery(call.sql)).length, 0);
assert.equal(queryClient.calls.filter((call) => call.sql.startsWith("SELECT session_json FROM workbench_sessions")).length, 2);
});
test("configured postgres runtime can query thin Workbench session summaries", 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"
});
await store.writeWorkbenchFacts({
facts: {
sessions: [{
sessionId: "ses_summary_projection",
ownerUserId: "usr_summary_projection",
projectId: "prj_summary_projection",
conversationId: "cnv_summary_projection",
threadId: "thr_summary_projection",
status: "running",
lastTraceId: "trc_summary_projection",
projectedSeq: 7,
sourceSeq: 7,
sourceEventId: "evt_summary_projection",
terminal: false,
sealed: false,
sessionJson: {
providerProfile: "codex-fast",
messages: [{ role: "user", text: "large payload".repeat(500) }],
prompt: "should stay out of list projection"
},
createdAt: "2026-06-20T09:59:00.000Z",
updatedAt: "2026-06-20T10:00:00.000Z"
}]
}
});
await store.readiness();
queryClient.calls.length = 0;
const loaded = await store.queryWorkbenchFacts({ families: ["sessions"], sessionProjection: "summary", limit: 10 });
const session = loaded.facts.sessions[0];
const readCall = queryClient.calls.find((call) => call.sql.includes("FROM workbench_sessions"));
assert.equal(session.sessionId, "ses_summary_projection");
assert.equal(session.providerProfile, "codex-fast");
assert.deepEqual(session.sessionJson, { providerProfile: "codex-fast", valuesRedacted: true, secretMaterialStored: false });
assert.equal(session.sessionJson.messages, undefined);
assert.ok(readCall?.sql.startsWith("SELECT session_id, owner_user_id"));
assert.equal(readCall.sql.includes("SELECT session_json FROM workbench_sessions"), false);
assert.ok(readCall.sql.includes("providerProfile"));
});
test("configured postgres runtime orders Workbench sessions by sessionsOrder updated_desc", 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"
});
await store.writeWorkbenchFacts({
facts: {
sessions: [
{ sessionId: "ses_order_old", ownerUserId: "usr_order", status: "running", updatedAt: "2026-06-15T10:00:00.000Z" },
{ sessionId: "ses_order_new", ownerUserId: "usr_order", status: "running", updatedAt: "2026-06-20T10:00:00.000Z" },
{ sessionId: "ses_order_middle", ownerUserId: "usr_order", status: "running", updatedAt: "2026-06-18T10:00:00.000Z" }
]
}
});
queryClient.calls.length = 0;
const loaded = await store.queryWorkbenchFacts({ families: ["sessions"], sessionsOrder: "updated_desc", limit: 10 });
const readCall = queryClient.calls.find((call) => call.sql.startsWith("SELECT session_json FROM workbench_sessions"));
assert.deepEqual(loaded.facts.sessions.map((session) => session.sessionId), ["ses_order_new", "ses_order_middle", "ses_order_old"]);
assert.match(readCall?.sql ?? "", /ORDER BY updated_at DESC/u);
});
test("configured postgres runtime queries requested workbench fact families sequentially", async () => {
const expectedReads = 4;
let started = 0;
let releaseReads;
let resolveAllStarted;
const readGate = new Promise((resolve) => {
releaseReads = resolve;
});
const allStarted = new Promise((resolve) => {
resolveAllStarted = resolve;
});
const queryClient = createFakePostgresClient({
migrationReady: true,
beforeWorkbenchFactRead: async (sql) => {
if (!/^SELECT (message_json|part_json|turn_json|checkpoint_json) FROM/u.test(sql)) return;
started += 1;
if (started === expectedReads) resolveAllStarted();
await readGate;
}
});
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"
});
await store.writeWorkbenchFacts({
facts: {
sessions: [],
messages: [{
messageId: "msg_parallel",
sessionId: "ses_parallel",
turnId: "turn_parallel",
traceId: "trc_parallel",
role: "assistant",
status: "completed",
text: "done",
updatedAt: "2026-06-20T10:00:01.000Z"
}],
parts: [{
partId: "part_parallel",
messageId: "msg_parallel",
sessionId: "ses_parallel",
turnId: "turn_parallel",
traceId: "trc_parallel",
partIndex: 0,
partType: "text",
status: "completed",
text: "done",
updatedAt: "2026-06-20T10:00:02.000Z"
}],
turns: [{
turnId: "turn_parallel",
sessionId: "ses_parallel",
traceId: "trc_parallel",
messageId: "msg_parallel",
status: "completed",
finalResponse: { text: "done" },
updatedAt: "2026-06-20T10:00:03.000Z"
}],
traceEvents: [],
checkpoints: [{
traceId: "trc_parallel",
sessionId: "ses_parallel",
turnId: "turn_parallel",
runId: "run_parallel",
commandId: "cmd_parallel",
projectionStatus: "caught_up",
projectionHealth: "healthy",
updatedAt: "2026-06-20T10:00:04.000Z"
}]
}
});
const pending = store.queryWorkbenchFacts({
sessionId: "ses_parallel",
traceId: "trc_parallel",
families: ["messages", "parts", "turns", "checkpoints"],
limit: 10
});
const startedConcurrently = await Promise.race([
allStarted.then(() => true),
new Promise((resolve) => setTimeout(() => resolve(false), 150))
]);
releaseReads();
const loaded = await pending;
assert.equal(startedConcurrently, false);
assert.equal(started, expectedReads);
assert.equal(loaded.count, 4);
assert.equal(loaded.facts.messages[0].messageId, "msg_parallel");
assert.equal(loaded.facts.parts[0].partId, "part_parallel");
assert.equal(loaded.facts.turns[0].turnId, "turn_parallel");
assert.equal(loaded.facts.checkpoints[0].traceId, "trc_parallel");
});
test("configured postgres runtime pushes trace event query filters into SQL", 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-17T03:50:00.000Z"
});
await store.writeAgentTraceEvent({ event: { traceId: "trc_sql_filter_a", sessionId: "ses_filter_a", seq: 1, level: "info", label: "a", createdAt: "2026-06-17T03:50:01.000Z" } });
await store.writeAgentTraceEvent({ event: { traceId: "trc_sql_filter_b", sessionId: "ses_filter_b", seq: 1, level: "info", label: "b", createdAt: "2026-06-17T03:50:02.000Z" } });
const events = await store.queryAgentTraceEvents({ sessionId: "ses_filter_b" });
assert.equal(events.count, 1);
assert.equal(events.events[0].traceId, "trc_sql_filter_b");
const select = queryClient.calls.findLast((call) => call.sql.startsWith("SELECT event_json FROM agent_trace_events"));
assert.match(select.sql, /WHERE agent_session_id = \$1/u);
assert.deepEqual(select.params, ["ses_filter_b"]);
});
function isReadinessQuery(sql) {
return sql.includes("information_schema.columns")
|| sql.startsWith("SELECT id, schema_version FROM hwlab_schema_migrations")
|| sql.startsWith("SELECT COUNT(*)::int AS count FROM ");
}
function createFakePostgresClient({
migrationReady = true,
migrationErrorCode = null,
countErrorCode = null,
readErrorCode = null,
outboxErrorCode = null,
beforeWorkbenchFactRead = null
} = {}) {
const state = {
gateway_sessions: new Map(),
box_resources: new Map(),
box_capabilities: new Map(),
hardware_operations: new Map(),
audit_events: new Map(),
evidence_records: new Map(),
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_event_sequences: new Map(),
workbench_events: new Map(),
workbench_session_inputs: new Map(),
workbench_projection_checkpoints: new Map(),
workbench_projection_outbox: new Map(),
hwlab_schema_migrations: new Map()
};
if (migrationReady) {
state.hwlab_schema_migrations.set(CLOUD_CORE_MIGRATION_ID, {
id: CLOUD_CORE_MIGRATION_ID,
schema_version: CLOUD_RUNTIME_DURABLE_ADAPTER_SCHEMA_VERSION
});
}
const calls = [];
return {
calls,
async query(sql, params = []) {
calls.push({ sql, params });
if (["BEGIN", "COMMIT", "ROLLBACK"].includes(sql)) {
return { rows: [] };
}
if (sql.includes("information_schema.columns")) {
return { rows: schemaRows() };
}
if (sql.startsWith("SELECT COUNT(*)::int AS count FROM ")) {
if (countErrorCode) {
const error = new Error("count query failed");
error.code = countErrorCode;
throw error;
}
const table = sql.match(/FROM ([a-z_]+)/u)?.[1];
return { rows: [{ count: state[table]?.size ?? 0 }] };
}
if (sql.startsWith("SELECT id, schema_version FROM hwlab_schema_migrations")) {
if (migrationErrorCode) {
const error = new Error("migration query failed");
error.code = migrationErrorCode;
throw error;
}
const row = state.hwlab_schema_migrations.get(params[0]);
return { rows: row ? [row] : [] };
}
if (sql.startsWith("CREATE INDEX IF NOT EXISTS idx_agent_trace_events_")) {
return { rows: [] };
}
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 pg_advisory_xact_lock")) {
return { rows: [{ pg_advisory_xact_lock: null }] };
}
if (sql.startsWith("SELECT projected_seq FROM workbench_trace_events WHERE trace_id = $1 AND source_event_id = $2")) {
const rows = [...state.workbench_trace_events.values()]
.filter((record) => record.trace_id === params[0] && record.source_event_id === params[1])
.sort((left, right) => Number(left.projected_seq) - Number(right.projected_seq));
return { rows: rows.slice(0, 1).map((record) => ({ projected_seq: record.projected_seq })) };
}
if (sql.startsWith("SELECT GREATEST(COALESCE((SELECT MAX(projected_seq) FROM workbench_trace_events")) {
const traceId = params[0];
const eventMax = [...state.workbench_trace_events.values()]
.filter((record) => record.trace_id === traceId)
.reduce((max, record) => Math.max(max, Number(record.projected_seq ?? 0)), 0);
const checkpointMax = Number(state.workbench_projection_checkpoints.get(traceId)?.projected_seq ?? 0);
return { rows: [{ max_projected_seq: Math.max(eventMax, checkpointMax) }] };
}
if (sql.startsWith("SELECT event_seq, aggregate_seq, projection_revision FROM workbench_events")) {
const row = [...state.workbench_events.values()]
.filter((record) => record.event_id === params[0] || (params[1] !== null && params[1] !== undefined && record.aggregate_id === params[2] && record.source_event_id === params[1]))
.sort((left, right) => Number(left.event_seq) - Number(right.event_seq))[0] ?? null;
return { rows: row ? [{ event_seq: row.event_seq, aggregate_seq: row.aggregate_seq, projection_revision: row.projection_revision }] : [] };
}
if (sql.startsWith("SELECT gateway_session_json FROM gateway_sessions")) {
return jsonSelect(state.gateway_sessions, params[0], "gateway_session_json");
}
if (sql.startsWith("SELECT resource_json FROM box_resources")) {
return jsonSelect(state.box_resources, params[0], "resource_json");
}
if (sql.startsWith("SELECT capability_json FROM box_capabilities")) {
return jsonSelect(state.box_capabilities, params[0], "capability_json");
}
if (sql.startsWith("SELECT event_json FROM audit_events")) {
if (readErrorCode) {
const error = new Error("audit read query failed");
error.code = readErrorCode;
throw error;
}
return { rows: [...state.audit_events.values()].map((record) => ({ event_json: record.event_json })) };
}
if (sql.startsWith("SELECT event_json FROM agent_trace_events")) {
if (readErrorCode) {
const error = new Error("agent trace event read query failed");
error.code = readErrorCode;
throw error;
}
const traceId = sql.includes("trace_id = $1") ? params[0] : null;
const sessionId = sql.includes("agent_session_id = $1") ? params[0] : null;
const workerSessionId = sql.includes("worker_session_id = $1") ? params[0] : null;
const level = sql.includes("level = $1") ? params[0] : null;
const rows = [...state.agent_trace_events.values()]
.filter((record) => !traceId || record.trace_id === traceId)
.filter((record) => !sessionId || record.agent_session_id === sessionId)
.filter((record) => !workerSessionId || record.worker_session_id === workerSessionId)
.filter((record) => !level || record.level === level)
.sort((left, right) => String(left.occurred_at).localeCompare(String(right.occurred_at)) || String(left.id).localeCompare(String(right.id)))
.map((record) => ({ event_json: record.event_json }));
return { rows };
}
if (sql.startsWith("SELECT projection_json FROM workbench_projection_state")) {
if (readErrorCode) {
const error = new Error("workbench projection state read query failed");
error.code = readErrorCode;
throw error;
}
let rows = [...state.workbench_projection_state.values()];
if (sql.includes("trace_id = $1")) rows = rows.filter((record) => record.trace_id === params[0]);
if (sql.includes("projection_status = ANY")) {
const statuses = params.find((param) => Array.isArray(param)) ?? [];
rows = rows.filter((record) => statuses.includes(record.projection_status));
}
if (sql.includes("next_retry_at IS NULL")) {
const dueAt = params.find((param) => typeof param === "string" && /^\d{4}-\d{2}-\d{2}T/u.test(param));
rows = rows.filter((record) => !record.next_retry_at || !dueAt || record.next_retry_at <= dueAt);
}
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_id, owner_user_id") && sql.includes("FROM workbench_sessions")) {
if (beforeWorkbenchFactRead) await beforeWorkbenchFactRead(sql, params);
return workbenchSessionSummaryRows(state.workbench_sessions, sql, params, readErrorCode);
}
if (sql.startsWith("SELECT session_json FROM workbench_sessions")) {
if (beforeWorkbenchFactRead) await beforeWorkbenchFactRead(sql, params);
return workbenchFactRows(state.workbench_sessions, "session_json", sql, params, readErrorCode);
}
if (sql.startsWith("SELECT input_json FROM workbench_session_inputs")) {
if (beforeWorkbenchFactRead) await beforeWorkbenchFactRead(sql, params);
return workbenchFactRows(state.workbench_session_inputs, "input_json", sql, params, readErrorCode);
}
if (sql.startsWith("SELECT message_json FROM workbench_messages")) {
if (beforeWorkbenchFactRead) await beforeWorkbenchFactRead(sql, params);
return workbenchFactRows(state.workbench_messages, "message_json", sql, params, readErrorCode);
}
if (sql.startsWith("SELECT part_json FROM workbench_parts")) {
if (beforeWorkbenchFactRead) await beforeWorkbenchFactRead(sql, params);
return workbenchFactRows(state.workbench_parts, "part_json", sql, params, readErrorCode);
}
if (sql.startsWith("SELECT turn_json FROM workbench_turns")) {
if (beforeWorkbenchFactRead) await beforeWorkbenchFactRead(sql, params);
return workbenchFactRows(state.workbench_turns, "turn_json", sql, params, readErrorCode);
}
if (sql.startsWith("SELECT event_json FROM workbench_trace_events")) {
if (beforeWorkbenchFactRead) await beforeWorkbenchFactRead(sql, params);
return workbenchFactRows(state.workbench_trace_events, "event_json", sql, params, readErrorCode);
}
if (sql.startsWith("SELECT checkpoint_json FROM workbench_projection_checkpoints")) {
if (beforeWorkbenchFactRead) await beforeWorkbenchFactRead(sql, params);
return workbenchFactRows(state.workbench_projection_checkpoints, "checkpoint_json", sql, params, readErrorCode);
}
if (sql.startsWith("SELECT outbox_seq, event_seq, aggregate_id")) {
const afterSeq = Number(params[0] ?? 0);
const traceFilter = sql.includes("trace_id = $2") ? params[1] : null;
const sessionFilter = sql.includes("session_id = $2") ? params[1] : sql.includes("session_id = $3") ? params[2] : null;
const limit = Number(params[params.length - 1] ?? 100);
const rows = [...state.workbench_projection_outbox.values()]
.filter((record) => Number(record.outbox_seq) > afterSeq)
.filter((record) => !traceFilter || record.trace_id === traceFilter)
.filter((record) => !sessionFilter || record.session_id === sessionFilter)
.sort((left, right) => Number(left.outbox_seq) - Number(right.outbox_seq))
.slice(0, limit);
return { rows };
}
if (sql.startsWith("SELECT metadata_json FROM evidence_records")) {
if (readErrorCode) {
const error = new Error("evidence read query failed");
error.code = readErrorCode;
throw error;
}
return { rows: [...state.evidence_records.values()].map((record) => ({ metadata_json: record.metadata_json })) };
}
if (sql.startsWith("INSERT INTO gateway_sessions")) {
state.gateway_sessions.set(params[0], { gateway_session_json: params[6] });
return { rows: [] };
}
if (sql.startsWith("INSERT INTO box_resources")) {
state.box_resources.set(params[0], { resource_json: params[5] });
return { rows: [] };
}
if (sql.startsWith("INSERT INTO box_capabilities")) {
state.box_capabilities.set(params[0], { capability_json: params[3] });
return { rows: [] };
}
if (sql.startsWith("INSERT INTO hardware_operations")) {
state.hardware_operations.set(params[0], { operation_json: params[4] });
return { rows: [] };
}
if (sql.startsWith("INSERT INTO audit_events")) {
state.audit_events.set(params[0], { event_json: params[8] });
return { rows: [] };
}
if (sql.startsWith("INSERT INTO evidence_records")) {
state.evidence_records.set(params[0], { metadata_json: params[5] });
return { rows: [] };
}
if (sql.startsWith("INSERT INTO agent_trace_events")) {
state.agent_trace_events.set(params[0], {
id: params[0],
trace_id: params[1],
agent_session_id: params[2],
worker_session_id: params[3],
level: params[4],
message: params[5],
event_json: params[6],
occurred_at: params[7]
});
return { rows: [] };
}
if (sql.startsWith("INSERT INTO workbench_projection_state")) {
state.workbench_projection_state.set(params[0], {
trace_id: params[0],
session_id: params[1],
conversation_id: params[2],
thread_id: params[3],
run_id: params[4],
command_id: params[5],
last_agentrun_seq: params[6],
last_projected_seq: params[7],
upstream_latest_seq: params[8],
projection_status: params[9],
projection_health: params[10],
result_sync_state: params[11],
next_retry_at: params[17],
projection_json: params[18],
created_at: params[19],
updated_at: params[20]
});
return { rows: [] };
}
if (sql.startsWith("INSERT INTO workbench_event_sequences")) {
if (!state.workbench_event_sequences.has(params[0])) {
state.workbench_event_sequences.set(params[0], {
aggregate_id: params[0], aggregate_type: params[1], session_id: params[2], turn_id: params[3], trace_id: params[4], last_seq: 0, created_at: params[5], updated_at: params[6]
});
}
return { rows: [] };
}
if (sql.startsWith("UPDATE workbench_event_sequences SET last_seq")) {
const record = state.workbench_event_sequences.get(params[0]);
const nextSeq = Number(record?.last_seq ?? 0) + 1;
state.workbench_event_sequences.set(params[0], { ...(record ?? { aggregate_id: params[0] }), last_seq: nextSeq, updated_at: params[1] });
return { rows: [{ last_seq: nextSeq }] };
}
if (sql.startsWith("INSERT INTO workbench_events")) {
const eventSeq = state.workbench_events.size + 1;
const record = {
event_seq: eventSeq,
event_id: params[0],
aggregate_id: params[1],
aggregate_type: params[2],
aggregate_seq: params[3],
session_id: params[4],
turn_id: params[5],
trace_id: params[6],
message_id: params[7],
source_run_id: params[8],
source_command_id: params[9],
source_seq: params[10],
source_event_id: params[11],
event_type: params[12],
projection_revision: params[13],
terminal: params[14],
sealed: params[15],
payload_json: params[16],
occurred_at: params[17],
committed_at: params[18]
};
const existing = [...state.workbench_events.values()].find((item) => item.event_id === record.event_id) ?? null;
state.workbench_events.set(existing?.event_seq ?? eventSeq, existing ? { ...existing, projection_revision: Math.max(Number(existing.projection_revision ?? 0), Number(record.projection_revision ?? 0)), terminal: Boolean(existing.terminal) || Boolean(record.terminal), sealed: Boolean(existing.sealed) || Boolean(record.sealed), payload_json: record.payload_json, committed_at: record.committed_at } : record);
const saved = existing ? state.workbench_events.get(existing.event_seq) : record;
return { rows: [{ event_seq: saved.event_seq, aggregate_seq: saved.aggregate_seq, projection_revision: saved.projection_revision }] };
}
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_session_inputs")) {
const existing = state.workbench_session_inputs.get(params[0]) ?? null;
state.workbench_session_inputs.set(params[0], {
input_id: params[0], session_id: params[1], turn_id: params[2], trace_id: params[3], message_id: params[4], command_id: params[5] ?? existing?.command_id ?? null, delivery: params[6], admitted_seq: existing?.admitted_seq > 0 ? existing.admitted_seq : params[7], promoted_seq: params[8] ?? existing?.promoted_seq ?? null, status: params[9], error_code: params[10], source_event_id: params[11] ?? existing?.source_event_id ?? null, input_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_outbox")) {
if (outboxErrorCode) {
const error = new Error("projection outbox insert failed");
error.code = outboxErrorCode;
throw error;
}
const outboxSeq = state.workbench_projection_outbox.size + 1;
state.workbench_projection_outbox.set(outboxSeq, {
outbox_seq: outboxSeq,
event_seq: params[0],
aggregate_id: params[1],
aggregate_seq: params[2],
projection_revision: params[3],
trace_id: params[4],
session_id: params[5],
turn_id: params[6],
message_id: params[7],
projected_seq: params[8],
source_seq: params[9],
source_event_id: params[10],
commit_type: params[11],
terminal: params[12],
sealed: params[13],
payload_json: params[14],
created_at: params[15]
});
return { rows: [] };
}
if (sql.startsWith("INSERT INTO workbench_projection_checkpoints")) {
const previous = state.workbench_projection_checkpoints.get(params[0]) ?? null;
const projectedSeq = sql.includes("RETURNING projected_seq")
? Math.max(Number(previous?.projected_seq ?? 0) + (previous ? 1 : 0), Number(params[5] ?? 0))
: previous && Number(params[5] ?? 0) < Number(previous.projected_seq ?? 0)
? Number(previous.projected_seq)
: Number(params[5] ?? 0);
if (previous && !sql.includes("RETURNING projected_seq") && Number(params[5] ?? 0) < Number(previous.projected_seq ?? 0)) {
return { rows: [] };
}
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: projectedSeq, 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: sql.includes("RETURNING projected_seq") ? [{ projected_seq: projectedSeq }] : [] };
}
throw new Error(`unexpected sql: ${sql}`);
}
};
}
function schemaRows() {
return Object.entries(CLOUD_RUNTIME_DURABLE_TABLE_COLUMNS).flatMap(([table, columns]) =>
columns.map((column) => ({
table_name: table,
column_name: column
}))
);
}
function jsonSelect(map, id, column) {
const row = map.get(id);
return {
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 direction = sql.includes("ORDER BY updated_at DESC") ? -1 : 1;
const rows = [...map.values()]
.filter((record) => sqlRecordMatches(record, sql, params))
.sort((left, right) => direction * (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 workbenchSessionSummaryRows(map, sql, params, readErrorCode) {
if (readErrorCode) {
const error = new Error("workbench fact read query failed");
error.code = readErrorCode;
throw error;
}
const direction = sql.includes("ORDER BY updated_at DESC") ? -1 : 1;
const rows = [...map.values()]
.filter((record) => sqlRecordMatches(record, sql, params))
.sort((left, right) => direction * (String(left.updated_at).localeCompare(String(right.updated_at)) || String(left.session_id).localeCompare(String(right.session_id))))
.map((record) => {
const session = JSON.parse(record.session_json || "{}");
const providerProfile = session.providerProfile ?? session.sessionJson?.providerProfile ?? null;
return {
session_id: record.session_id,
owner_user_id: record.owner_user_id,
project_id: record.project_id,
conversation_id: record.conversation_id,
thread_id: record.thread_id,
status: record.status,
last_trace_id: record.last_trace_id,
projected_seq: record.projected_seq,
source_seq: record.source_seq,
source_event_id: record.source_event_id,
terminal: record.terminal,
sealed: record.sealed,
created_at: record.created_at,
updated_at: record.updated_at,
provider_profile: providerProfile
};
});
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;
}