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

308 lines
14 KiB
TypeScript

import assert from "node:assert/strict";
import { randomUUID } from "node:crypto";
import { test } from "bun:test";
import pg from "pg";
import { buildWorkbenchProjectionEventFacts } from "../cloud/workbench-projection-writer.ts";
import { PostgresCloudRuntimeStore } from "./runtime-store-postgres.ts";
import { commitAgentRunKafkaProjection as commitProjection } from "./runtime-store-postgres-kafka.ts";
const dbUrl = String(process.env.HWLAB_KAFKA_PROJECTOR_INTEGRATION_DB_URL ?? "").trim();
const enabled = Boolean(dbUrl) && process.env.HWLAB_KAFKA_PROJECTOR_INTEGRATION_CONFIRM_NON_PRODUCTION === "1";
if (!enabled) {
test.skip("Kafka projector real PostgreSQL transaction and fault-injection contract", () => {});
} else {
test("Kafka projector real PostgreSQL transaction and fault-injection contract", async () => {
const { Pool } = pg;
const pool = new Pool({ connectionString: dbUrl, ssl: false, max: 4 });
const suffix = randomUUID().replaceAll("-", "");
const traceId = `trc_it_${suffix}`;
const sessionId = `ses_it_${suffix}`;
const ownerUserId = `usr_it_${suffix}`;
const now = "2026-07-10T12:00:00.000Z";
const store = new PostgresCloudRuntimeStore({
dbUrl,
sslMode: "disable",
env: { ...process.env, HWLAB_CLOUD_DB_SSL_MODE: "disable" },
queryClient: pool,
now: () => now,
logger: { error() {}, warn() {}, info() {} }
});
try {
const readiness = await store.readiness();
assert.equal(readiness.ready, true);
assert.equal(readiness.durable, true);
assert.equal(readiness.migration?.ready, true);
await pool.query(
"INSERT INTO users (id, username, display_name, role, status, created_at, updated_at) VALUES ($1,$2,$2,'user','active',$3,$3)",
[ownerUserId, `integration-${suffix}`, now]
);
await store.writeWorkbenchSessionAdmissionFact({
fact: {
sessionId,
ownerUserId,
projectId: `prj_it_${suffix}`,
conversationId: `cnv_it_${suffix}`,
threadId: `thread_it_${suffix}`,
status: "queued",
lastTraceId: traceId,
projectedSeq: 0,
sourceSeq: 0,
terminal: false,
sealed: false,
sessionJson: { sessionId, launchContext: { source: "integration-test" }, valuesRedacted: true },
createdAt: now,
updatedAt: now
}
});
const running = await project(store, {
traceId,
sessionId,
sourceEventId: `evt_running_${suffix}`,
sourceSeq: 1,
sourceOffset: "100",
inputSha256: "1".repeat(64),
event: { eventType: "assistant_message", status: "running", assistantText: "partial", createdAt: now }
});
assert.equal(running.duplicate, false);
assert.equal(running.projectedSeq, 1);
assert.equal(await count(pool, "workbench_kafka_inbox", "trace_id = $1", [traceId]), 1);
assert.equal(await count(pool, "workbench_trace_events", "trace_id = $1", [traceId]), 1);
assert.equal(await count(pool, "workbench_projection_outbox", "trace_id = $1", [traceId]), 2);
assert.equal(await count(pool, "hwlab_kafka_outbox", "partition_key = $1", [traceId]), 1);
const sameTransport = await project(store, {
traceId,
sessionId,
sourceEventId: `evt_running_${suffix}`,
sourceSeq: 1,
sourceOffset: "100",
inputSha256: "1".repeat(64),
event: { eventType: "assistant_message", status: "running", assistantText: "partial", createdAt: now }
});
const newOffsetSameEvent = await project(store, {
traceId,
sessionId,
sourceEventId: `evt_running_${suffix}`,
sourceSeq: 1,
sourceOffset: "101",
inputSha256: "1".repeat(64),
event: { eventType: "assistant_message", status: "running", assistantText: "partial", createdAt: now }
});
assert.equal(sameTransport.duplicate, true);
assert.equal(newOffsetSameEvent.duplicate, true);
assert.equal(await count(pool, "workbench_kafka_inbox", "trace_id = $1", [traceId]), 1);
assert.equal(await count(pool, "workbench_projection_outbox", "trace_id = $1", [traceId]), 2);
const conflict = await project(store, {
traceId,
sessionId,
sourceEventId: `evt_running_${suffix}`,
sourceSeq: 1,
sourceOffset: "102",
inputSha256: "2".repeat(64),
event: { eventType: "assistant_message", status: "running", assistantText: "mutated", createdAt: now }
});
assert.equal(conflict.conflict, true);
assert.equal(await count(pool, "workbench_kafka_dlq", "source_event_id = $1", [`evt_running_${suffix}`]), 1);
const projectedInbox = await pool.query(
"SELECT status, projected_seq FROM workbench_kafka_inbox WHERE source_topic = $1 AND source_event_id = $2",
["agentrun.event.v1", `evt_running_${suffix}`]
);
assert.deepEqual(projectedInbox.rows[0], { status: "projected", projected_seq: 1 });
await pool.query(`
CREATE OR REPLACE FUNCTION issue_2464_delay_inbox_insert() RETURNS trigger AS $$
BEGIN
IF NEW.source_topic = 'agentrun.event.v1' AND NEW.source_partition = 0 AND NEW.source_offset = 300 THEN
PERFORM pg_sleep(0.25);
END IF;
RETURN NEW;
END;
$$ LANGUAGE plpgsql;
DROP TRIGGER IF EXISTS issue_2464_delay_inbox_insert ON workbench_kafka_inbox;
CREATE TRIGGER issue_2464_delay_inbox_insert BEFORE INSERT ON workbench_kafka_inbox
FOR EACH ROW EXECUTE FUNCTION issue_2464_delay_inbox_insert();
`);
try {
const concurrent = await Promise.allSettled([
project(store, {
traceId: `trc_transport_a_${suffix}`,
sessionId: `ses_transport_a_${suffix}`,
sourceEventId: `evt_transport_a_${suffix}`,
sourceSeq: 1,
sourceOffset: "300",
inputSha256: "6".repeat(64),
event: { eventType: "backend_status", status: "running", createdAt: now }
}),
project(store, {
traceId: `trc_transport_b_${suffix}`,
sessionId: `ses_transport_b_${suffix}`,
sourceEventId: `evt_transport_b_${suffix}`,
sourceSeq: 1,
sourceOffset: "300",
inputSha256: "7".repeat(64),
event: { eventType: "backend_status", status: "running", createdAt: now }
})
]);
assert.deepEqual(concurrent.map((result) => result.status), ["fulfilled", "fulfilled"]);
const values = concurrent.map((result) => {
if (result.status !== "fulfilled") throw result.reason;
return result.value;
});
assert.equal(values.filter((result) => result.conflict === true).length, 1);
assert.equal(values.filter((result) => result.conflict !== true && result.duplicate !== true).length, 1);
assert.equal(await count(pool, "workbench_kafka_inbox", "source_topic = $1 AND source_partition = 0 AND source_offset = 300", ["agentrun.event.v1"]), 1);
assert.equal(await count(pool, "workbench_kafka_dlq", "source_topic = $1 AND source_partition = 0 AND source_offset = 300", ["agentrun.event.v1"]), 1);
} finally {
await pool.query("DROP TRIGGER IF EXISTS issue_2464_delay_inbox_insert ON workbench_kafka_inbox; DROP FUNCTION IF EXISTS issue_2464_delay_inbox_insert()");
}
const faultTraceId = `trc_fault_${suffix}`;
const faultSessionId = `ses_fault_${suffix}`;
const faultStore = new PostgresCloudRuntimeStore({
dbUrl,
sslMode: "disable",
env: { ...process.env, HWLAB_CLOUD_DB_SSL_MODE: "disable" },
queryClient: faultInjectingPool(pool, /^INSERT INTO workbench_projection_outbox/u),
now: () => now,
logger: { error() {}, warn() {}, info() {} }
});
await assert.rejects(project(faultStore, {
traceId: faultTraceId,
sessionId: faultSessionId,
sourceEventId: `evt_fault_${suffix}`,
sourceSeq: 1,
sourceOffset: "200",
inputSha256: "3".repeat(64),
event: { eventType: "assistant_message", status: "running", assistantText: "must roll back", createdAt: now }
}), /injected projection outbox failure/u);
assert.equal(await count(pool, "workbench_kafka_inbox", "trace_id = $1", [faultTraceId]), 0);
assert.equal(await count(pool, "workbench_trace_events", "trace_id = $1", [faultTraceId]), 0);
assert.equal(await count(pool, "workbench_projection_outbox", "trace_id = $1", [faultTraceId]), 0);
assert.equal(await count(pool, "hwlab_kafka_outbox", "partition_key = $1", [faultTraceId]), 0);
const terminal = await project(store, {
traceId,
sessionId,
sourceEventId: `evt_terminal_${suffix}`,
sourceSeq: 2,
sourceOffset: "103",
inputSha256: "4".repeat(64),
event: {
eventType: "terminal_status",
status: "completed",
terminal: true,
assistantText: "authoritative final",
finalResponse: { text: "authoritative final", status: "completed" },
startedAt: now,
finishedAt: "2026-07-10T12:00:01.000Z",
durationMs: 1000,
createdAt: "2026-07-10T12:00:01.000Z"
}
});
assert.equal(terminal.projectedSeq, 2);
const sealed = await pool.query(
"SELECT projected_seq, terminal, sealed, checkpoint_json FROM workbench_projection_checkpoints WHERE trace_id = $1",
[traceId]
);
assert.equal(sealed.rows[0].projected_seq, 2);
assert.equal(sealed.rows[0].terminal, true);
assert.equal(sealed.rows[0].sealed, true);
assert.equal(JSON.parse(sealed.rows[0].checkpoint_json).finalResponse.text, "authoritative final");
assert.equal(await count(pool, "workbench_projection_outbox", "trace_id = $1", [traceId]), 5);
const late = await project(store, {
traceId,
sessionId,
sourceEventId: `evt_late_${suffix}`,
sourceSeq: 3,
sourceOffset: "104",
inputSha256: "5".repeat(64),
event: { eventType: "backend_status", status: "running", terminal: false, createdAt: "2026-07-10T12:00:02.000Z" }
});
assert.equal(late.suppressedAfterSeal, true);
assert.equal(late.projectedSeq, 2);
assert.equal(await count(pool, "workbench_trace_events", "trace_id = $1", [traceId]), 2);
assert.equal(await count(pool, "workbench_projection_outbox", "trace_id = $1", [traceId]), 5);
const lateInbox = await pool.query(
"SELECT status, projected_seq FROM workbench_kafka_inbox WHERE source_topic = $1 AND source_event_id = $2",
["agentrun.event.v1", `evt_late_${suffix}`]
);
assert.deepEqual(lateInbox.rows[0], { status: "projected", projected_seq: 2 });
const firstSync = await store.readAtomicWorkbenchProjectionSync({ traceId, afterOutboxSeq: 0, limit: 100 });
assert.equal(firstSync.events.length, 5);
assert.equal(firstSync.hasMore, false);
assert.equal(firstSync.cursorOutboxSeq, firstSync.cutoffOutboxSeq);
assert.equal(new Set(firstSync.events.map((event) => event.outboxSeq)).size, firstSync.events.length);
assert.equal(firstSync.facts.checkpoints[0].projectedSeq, 2);
assert.equal(firstSync.facts.turns[0].finalResponse.text, "authoritative final");
const caughtUp = await store.readAtomicWorkbenchProjectionSync({ traceId, afterOutboxSeq: firstSync.cursorOutboxSeq, limit: 100, deltaOnly: true });
assert.equal(caughtUp.events.length, 0);
assert.equal(caughtUp.cursorOutboxSeq, firstSync.cursorOutboxSeq);
} finally {
await pool.end();
}
}, 30_000);
}
async function project(store, { traceId, sessionId, sourceEventId, sourceSeq, sourceOffset, inputSha256, event }) {
const canonicalEvent = {
schema: "agentrun.event.v1",
eventId: sourceEventId,
sourceSeq,
traceId,
hwlabSessionId: sessionId,
runId: `run_${traceId}`,
commandId: `cmd_${traceId}`,
event
};
const projectedEvent = { ...event, traceId, sessionId, sourceEventId, sourceSeq, runId: canonicalEvent.runId, commandId: canonicalEvent.commandId };
return commitProjection(store, {
transport: { sourceTopic: "agentrun.event.v1", sourcePartition: 0, sourceOffset, sourceKey: canonicalEvent.runId, inputSha256 },
canonicalEvent,
projectedEvent,
requestMeta: { traceId, sessionId, valuesPrinted: false },
factsFactory: ({ projectedSeq, previousCheckpoint, projectedAt }) => buildWorkbenchProjectionEventFacts({ projectedSeq, previousCheckpoint, projectedAt, event: projectedEvent }),
hwlabEvent: { schema: "hwlab.event.v1", eventId: `hwlab:${sourceEventId}`, traceId, sessionId, sourceEventId, sourceSeq, valuesPrinted: false },
hwlabTopic: "hwlab.event.v1",
partitionKey: traceId,
headers: { sourceEventId }
});
}
function faultInjectingPool(pool, pattern) {
return {
async connect() {
const client = await pool.connect();
let injected = false;
return {
async query(sql, params) {
if (!injected && pattern.test(String(sql))) {
injected = true;
const error = new Error("injected projection outbox failure");
error.code = "XX2464";
throw error;
}
return client.query(sql, params);
},
release() {
client.release();
}
};
}
};
}
async function count(pool, table, where, params) {
if (!/^[a-z_]+$/u.test(table)) throw new Error("unsafe test table");
const result = await pool.query(`SELECT COUNT(*)::int AS count FROM ${table} WHERE ${where}`, params);
return Number(result.rows[0].count);
}