feat: read workbench get from durable facts
This commit is contained in:
@@ -138,12 +138,7 @@ test("workbench trace event page exposes monotonic cursor range for restored mix
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const runtimeStore = {
|
||||
async queryAgentTraceEvents(input = {}) {
|
||||
assert.equal(input.traceId, traceId);
|
||||
return { events: mixedEvents, count: mixedEvents.length };
|
||||
}
|
||||
};
|
||||
const runtimeStore = createDurableFactsRuntimeStore({ sessions: [{ session, events: mixedEvents, status: "completed" }] });
|
||||
const server = createCloudApiServer({ accessController, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
@@ -207,7 +202,7 @@ test("workbench read model exposes session, messages, turn, and trace without wr
|
||||
secretMaterialStored: false
|
||||
}
|
||||
};
|
||||
const listInputs = [];
|
||||
const factQueries = [];
|
||||
traceStore.append(traceId, { type: "request", status: "accepted", label: "request:accepted" });
|
||||
traceStore.append(traceId, { type: "result", status: "completed", label: "result:completed", terminal: true });
|
||||
results.set(traceId, {
|
||||
@@ -224,7 +219,6 @@ test("workbench read model exposes session, messages, turn, and trace without wr
|
||||
const accessController = {
|
||||
store: {
|
||||
async listAgentSessionsForUser(input = {}) {
|
||||
listInputs.push(input);
|
||||
return [session];
|
||||
},
|
||||
async getAgentSession(sessionId) { return sessionId === session.id ? session : null; },
|
||||
@@ -235,7 +229,22 @@ test("workbench read model exposes session, messages, turn, and trace without wr
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: results });
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{
|
||||
session,
|
||||
events: [
|
||||
{ seq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-17T00:00:00.000Z" },
|
||||
{ seq: 2, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-17T00:00:03.000Z" }
|
||||
],
|
||||
status: "completed",
|
||||
finalText: "pong",
|
||||
runId: "run_workbench_read_model",
|
||||
commandId: "cmd_workbench_read_model",
|
||||
lastProjectedSeq: 2
|
||||
}],
|
||||
queries: factQueries
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
try {
|
||||
@@ -245,13 +254,13 @@ test("workbench read model exposes session, messages, turn, and trace without wr
|
||||
assert.equal(sessions.body.contractVersion, "workbench-sessions-v1");
|
||||
assert.equal(sessions.body.sessions[0].sessionId, session.id);
|
||||
assert.equal(sessions.body.sessions[0].turnSummary.traceId, traceId);
|
||||
assert.equal(listInputs[0].projectId, undefined);
|
||||
assert.equal(listInputs[0].limit, 21);
|
||||
assert.equal(factQueries[0].projectId, undefined);
|
||||
assert.equal(factQueries[0].limit, 21);
|
||||
|
||||
const staleProject = await getJson(port, `/v1/workbench/sessions?projectId=prj_stale_filter&includeSessionId=${encodeURIComponent(session.id)}`);
|
||||
assert.equal(staleProject.status, 400);
|
||||
assert.equal(staleProject.body.error.code, "workbench_authority_removed");
|
||||
assert.equal(listInputs.length, 1);
|
||||
assert.equal(factQueries.length, 1);
|
||||
|
||||
const detail = await getJson(port, `/v1/workbench/sessions/${encodeURIComponent(session.id)}`);
|
||||
assert.equal(detail.status, 200);
|
||||
@@ -333,14 +342,11 @@ test("workbench read model recovers trace events from durable projection without
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
let traceQueryCount = 0;
|
||||
const runtimeStore = {
|
||||
async queryAgentTraceEvents(params = {}) {
|
||||
traceQueryCount += 1;
|
||||
assert.equal(params.traceId, traceId);
|
||||
return { events: durableEvents, count: durableEvents.length };
|
||||
}
|
||||
};
|
||||
const factQueries = [];
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{ session, events: durableEvents, status: "completed", finalText: "durable pong", lastProjectedSeq: 2 }],
|
||||
queries: factQueries
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
@@ -362,7 +368,7 @@ test("workbench read model recovers trace events from durable projection without
|
||||
assert.equal(trace.body.events.length, 2);
|
||||
assert.equal(trace.body.traceStatus, "completed");
|
||||
assert.equal(trace.body.hasMore, false);
|
||||
assert.equal(traceQueryCount, 1);
|
||||
assert.ok(factQueries.length >= 3);
|
||||
} finally {
|
||||
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
|
||||
}
|
||||
@@ -415,7 +421,21 @@ test("workbench read model projects terminal result atomically across session, m
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: results });
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{
|
||||
session,
|
||||
events: [
|
||||
{ seq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-17T02:09:58.000Z" },
|
||||
{ seq: 2, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-17T02:10:00.000Z" }
|
||||
],
|
||||
status: "completed",
|
||||
finalText,
|
||||
runId: "run_workbench_terminal_atomic",
|
||||
commandId: "cmd_workbench_terminal_atomic",
|
||||
lastProjectedSeq: 2
|
||||
}]
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
try {
|
||||
@@ -498,7 +518,15 @@ test("workbench read model projects current turn running state from trace projec
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: createCodeAgentChatResultStore() });
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{
|
||||
session,
|
||||
events: [{ seq: 1, type: "backend", status: "running", label: "runner:created", terminal: false, createdAt: session.updatedAt }],
|
||||
status: "running",
|
||||
lastProjectedSeq: 1
|
||||
}]
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
try {
|
||||
@@ -536,7 +564,15 @@ test("workbench read model projects current turn running state for idle session
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: createCodeAgentChatResultStore() });
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{
|
||||
session,
|
||||
events: [{ seq: 1, type: "backend", status: "running", label: "runner:created", terminal: false, createdAt: session.updatedAt }],
|
||||
status: "running",
|
||||
lastProjectedSeq: 1
|
||||
}]
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
try {
|
||||
@@ -598,7 +634,21 @@ test("workbench read model projects completed current turn for idle session summ
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const backendPerformanceStore = createBackendPerformanceStore({ env: { POD_NAMESPACE: "hwlab-v03", HWLAB_GITOPS_TARGET: "v03", HWLAB_RUNTIME_LANE: "v03", HWLAB_CODE_AGENT_AGENTRUN_PROVIDER_ID: "D601" } });
|
||||
const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: results, backendPerformanceStore });
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{
|
||||
session,
|
||||
events: [
|
||||
{ seq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-17T00:02:58.000Z" },
|
||||
{ seq: 2, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-17T00:03:00.000Z" }
|
||||
],
|
||||
status: "completed",
|
||||
finalText,
|
||||
runId: "run_workbench_idle_completed",
|
||||
commandId: "cmd_workbench_idle_completed",
|
||||
lastProjectedSeq: 2
|
||||
}]
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results, backendPerformanceStore });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
try {
|
||||
@@ -685,7 +735,21 @@ test("workbench read model lets terminal result override stale running session s
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const server = createCloudApiServer({ accessController, traceStore, codeAgentChatResults: results });
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{
|
||||
session,
|
||||
events: [
|
||||
{ seq: 1, type: "backend", status: "running", label: "runner:created", terminal: false, createdAt: "2026-06-18T16:30:00.000Z" },
|
||||
{ seq: 2, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-18T16:31:04.000Z" }
|
||||
],
|
||||
status: "completed",
|
||||
finalText,
|
||||
runId: "run_workbench_stale_running_completed",
|
||||
commandId: "cmd_workbench_stale_running_completed",
|
||||
lastProjectedSeq: 2
|
||||
}]
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
try {
|
||||
@@ -720,7 +784,7 @@ test("workbench read model lets terminal result override stale running session s
|
||||
}
|
||||
});
|
||||
|
||||
test("workbench session list uses compact trace result without durable trace hydration", async () => {
|
||||
test("workbench session list uses durable projection checkpoint without trace/result hydration", async () => {
|
||||
const traceStore = createCodeAgentTraceStore();
|
||||
const traceId = "trc_workbench_compact_list_summary";
|
||||
const finalText = "compact list summary OK";
|
||||
@@ -771,13 +835,22 @@ test("workbench session list uses compact trace result without durable trace hyd
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
let durableQueryCount = 0;
|
||||
const runtimeStore = {
|
||||
async queryAgentTraceEvents() {
|
||||
durableQueryCount += 1;
|
||||
throw new Error("session list must not hydrate durable trace events");
|
||||
}
|
||||
};
|
||||
const factQueries = [];
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{
|
||||
session,
|
||||
events: [
|
||||
{ seq: 1, sourceSeq: 1, type: "request", status: "accepted", label: "request:accepted", createdAt: "2026-06-19T15:39:10.000Z" },
|
||||
{ seq: 42, sourceSeq: 42, type: "result", status: "completed", label: "result:completed", terminal: true, createdAt: "2026-06-19T15:40:00.000Z" }
|
||||
],
|
||||
status: "completed",
|
||||
finalText,
|
||||
runId: "run_workbench_compact_list_summary",
|
||||
commandId: "cmd_workbench_compact_list_summary",
|
||||
lastProjectedSeq: 42
|
||||
}],
|
||||
queries: factQueries
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
@@ -791,13 +864,13 @@ test("workbench session list uses compact trace result without durable trace hyd
|
||||
assert.equal(sessions.body.sessions[0].turnSummary.status, "completed");
|
||||
assert.equal(sessions.body.sessions[0].turnSummary.eventCount, 42);
|
||||
assert.equal(sessions.body.sessions[0].projectionStatus, "caught-up");
|
||||
assert.equal(durableQueryCount, 0);
|
||||
assert.ok(factQueries.length >= 1);
|
||||
} finally {
|
||||
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
|
||||
}
|
||||
});
|
||||
|
||||
test("workbench read model recovers terminal session status from durable trace", async () => {
|
||||
test("workbench read model reads terminal session status from durable facts", async () => {
|
||||
const traceStore = createCodeAgentTraceStore();
|
||||
const traceId = "trc_workbench_durable_terminal_after_memory_running";
|
||||
const session = {
|
||||
@@ -833,14 +906,9 @@ test("workbench read model recovers terminal session status from durable trace",
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
let durableQueryCount = 0;
|
||||
const runtimeStore = {
|
||||
async queryAgentTraceEvents(params = {}) {
|
||||
durableQueryCount += 1;
|
||||
assert.equal(params.traceId, traceId);
|
||||
return { events: durableEvents, count: durableEvents.length };
|
||||
}
|
||||
};
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{ session, events: durableEvents, status: "completed", finalText: "OK", lastProjectedSeq: 2 }]
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: createCodeAgentChatResultStore() });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
@@ -848,17 +916,15 @@ test("workbench read model recovers terminal session status from durable trace",
|
||||
const { port } = server.address();
|
||||
const sessions = await getJson(port, `/v1/workbench/sessions?includeSessionId=${encodeURIComponent(session.id)}`);
|
||||
assert.equal(sessions.status, 200);
|
||||
assert.equal(sessions.body.sessions[0].status, "running");
|
||||
assert.equal(sessions.body.sessions[0].running, true);
|
||||
assert.equal(sessions.body.sessions[0].terminal, false);
|
||||
assert.equal(durableQueryCount, 0);
|
||||
assert.equal(sessions.body.sessions[0].status, "completed");
|
||||
assert.equal(sessions.body.sessions[0].running, false);
|
||||
assert.equal(sessions.body.sessions[0].terminal, true);
|
||||
|
||||
const detail = await getJson(port, `/v1/workbench/sessions/${encodeURIComponent(session.id)}`);
|
||||
assert.equal(detail.status, 200);
|
||||
assert.equal(detail.body.session.status, "completed");
|
||||
assert.equal(detail.body.session.running, false);
|
||||
assert.equal(detail.body.session.terminal, true);
|
||||
assert.equal(durableQueryCount, 1);
|
||||
|
||||
const messages = await getJson(port, `/v1/workbench/sessions/${encodeURIComponent(session.id)}/messages?limit=10`);
|
||||
assert.equal(messages.status, 200);
|
||||
@@ -927,12 +993,17 @@ test("workbench read model keeps non-terminal durable tool completed events out
|
||||
async ensureBootstrap() {},
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const runtimeStore = {
|
||||
async queryAgentTraceEvents(params = {}) {
|
||||
assert.equal(params.traceId, traceId);
|
||||
return { events: durableEvents, count: durableEvents.length };
|
||||
}
|
||||
};
|
||||
const runtimeStore = createDurableFactsRuntimeStore({
|
||||
sessions: [{
|
||||
session,
|
||||
events: durableEvents,
|
||||
status: "running",
|
||||
runId: "run_workbench_nonterminal_tool_completed",
|
||||
commandId: "cmd_workbench_nonterminal_tool_completed",
|
||||
projectionStatus: "projecting",
|
||||
lastProjectedSeq: 2
|
||||
}]
|
||||
});
|
||||
const server = createCloudApiServer({ accessController, traceStore, runtimeStore, codeAgentChatResults: results });
|
||||
await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve));
|
||||
|
||||
@@ -1053,7 +1124,7 @@ test("workbench read model exposes runtime trace projection query failures as pr
|
||||
async authenticate() { return { ok: true, actor: ACTOR, session: { id: "uss_workbench_reader" } }; }
|
||||
};
|
||||
const runtimeStore = {
|
||||
async queryAgentTraceEvents() {
|
||||
async queryWorkbenchFacts() {
|
||||
const error = new Error("runtime query failed");
|
||||
error.code = "TEST_RUNTIME_QUERY_FAILED";
|
||||
throw error;
|
||||
@@ -1065,18 +1136,12 @@ test("workbench read model exposes runtime trace projection query failures as pr
|
||||
try {
|
||||
const { port } = server.address();
|
||||
const turn = await getJson(port, `/v1/workbench/turns/${encodeURIComponent(traceId)}`);
|
||||
assert.equal(turn.status, 200);
|
||||
assert.equal(turn.body.turn.status, "running");
|
||||
assert.equal(turn.body.projectionHealth, "unavailable");
|
||||
assert.equal(turn.body.projection.blocker.code, "projection_store_unavailable");
|
||||
assert.equal(turn.body.blocker.code, "projection_store_unavailable");
|
||||
assert.equal(turn.status, 503);
|
||||
assert.equal(turn.body.error.code, "projection_store_unavailable");
|
||||
|
||||
const trace = await getJson(port, `/v1/workbench/traces/${encodeURIComponent(traceId)}/events?limit=10`);
|
||||
assert.equal(trace.status, 200);
|
||||
assert.equal(trace.body.events.length, 1);
|
||||
assert.equal(trace.body.traceStatus, "running");
|
||||
assert.equal(trace.body.projectionHealth, "unavailable");
|
||||
assert.equal(trace.body.projection.blocker.code, "projection_store_unavailable");
|
||||
assert.equal(trace.status, 503);
|
||||
assert.equal(trace.body.error.code, "projection_store_unavailable");
|
||||
} finally {
|
||||
await new Promise((resolve, reject) => server.close((error) => error ? reject(error) : resolve()));
|
||||
}
|
||||
@@ -1090,6 +1155,198 @@ async function getJson(port, path) {
|
||||
};
|
||||
}
|
||||
|
||||
function createDurableFactsRuntimeStore({ sessions = [], queryError = null, queries = [] } = {}) {
|
||||
const facts = emptyFacts();
|
||||
for (const item of sessions) mergeFacts(facts, buildDurableFactsForSession(item));
|
||||
return {
|
||||
queries,
|
||||
async queryWorkbenchFacts(params = {}) {
|
||||
queries.push(params);
|
||||
if (queryError) throw queryError;
|
||||
const filtered = filterFacts(facts, params);
|
||||
return {
|
||||
facts: filtered,
|
||||
count: Object.values(filtered).reduce((sum, rows) => sum + rows.length, 0),
|
||||
persistence: { adapter: "test-durable-workbench-facts", durable: true }
|
||||
};
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
function buildDurableFactsForSession({ session, events = [], status = null, finalText = null, runId = null, commandId = null, projectionStatus = null, projectionHealth = "healthy", lastProjectedSeq = null } = {}) {
|
||||
const traceId = session.lastTraceId ?? session.session?.lastTraceId;
|
||||
const projectedStatus = normalizeTestStatus(status ?? session.status);
|
||||
const terminal = testTerminalStatuses.has(projectedStatus);
|
||||
const normalizedMessages = normalizeTestMessages(session, { traceId, status: projectedStatus, finalText, terminal });
|
||||
const normalizedEvents = normalizeTestEvents(events, session, { traceId });
|
||||
const projectedSeq = Number.isFinite(Number(lastProjectedSeq)) ? Number(lastProjectedSeq) : Math.max(normalizedEvents.length, normalizedMessages.length, 0);
|
||||
const checkpointStatus = projectionStatus ?? (terminal ? "caught_up" : "projecting");
|
||||
return {
|
||||
sessions: [{
|
||||
sessionId: session.id,
|
||||
ownerUserId: session.ownerUserId,
|
||||
ownerRole: session.ownerRole ?? null,
|
||||
projectId: session.projectId ?? null,
|
||||
conversationId: session.conversationId ?? null,
|
||||
threadId: session.threadId ?? null,
|
||||
agentId: session.agentId ?? "hwlab-code-agent",
|
||||
status: projectedStatus,
|
||||
lastTraceId: traceId,
|
||||
providerProfile: session.session?.providerProfile ?? null,
|
||||
projectedSeq,
|
||||
sourceSeq: projectedSeq,
|
||||
sourceEventId: traceId,
|
||||
terminal,
|
||||
sealed: terminal,
|
||||
createdAt: session.startedAt ?? session.createdAt ?? session.updatedAt ?? null,
|
||||
updatedAt: session.updatedAt ?? session.endedAt ?? session.startedAt ?? null,
|
||||
valuesRedacted: true
|
||||
}],
|
||||
messages: normalizedMessages,
|
||||
parts: normalizedMessages.filter((message) => message.text).map((message, index) => ({
|
||||
partId: `prt_${session.id}_${index}`,
|
||||
messageId: message.messageId,
|
||||
sessionId: session.id,
|
||||
turnId: message.turnId,
|
||||
traceId: message.traceId,
|
||||
partIndex: 0,
|
||||
partType: "text",
|
||||
status: message.status,
|
||||
text: message.text,
|
||||
projectedSeq: message.projectedSeq,
|
||||
sourceSeq: message.sourceSeq,
|
||||
sourceEventId: message.sourceEventId,
|
||||
terminal: message.terminal,
|
||||
sealed: message.sealed,
|
||||
updatedAt: message.updatedAt,
|
||||
valuesRedacted: true
|
||||
})),
|
||||
turns: traceId ? [{
|
||||
turnId: traceId,
|
||||
sessionId: session.id,
|
||||
traceId,
|
||||
messageId: normalizedMessages.find((message) => message.role !== "user")?.messageId ?? null,
|
||||
status: projectedStatus,
|
||||
projectedSeq,
|
||||
sourceSeq: projectedSeq,
|
||||
sourceEventId: traceId,
|
||||
terminal,
|
||||
sealed: terminal,
|
||||
finalResponse: finalText ? { text: finalText, status: projectedStatus, traceId, valuesPrinted: false } : null,
|
||||
diagnostic: { projectionStatus: checkpointStatus, projectionHealth, valuesRedacted: true },
|
||||
createdAt: session.startedAt ?? session.updatedAt ?? null,
|
||||
updatedAt: session.updatedAt ?? null,
|
||||
valuesRedacted: true
|
||||
}] : [],
|
||||
traceEvents: normalizedEvents,
|
||||
checkpoints: traceId ? [{
|
||||
traceId,
|
||||
sessionId: session.id,
|
||||
turnId: traceId,
|
||||
runId,
|
||||
commandId,
|
||||
projectedSeq,
|
||||
sourceSeq: projectedSeq,
|
||||
sourceEventId: traceId,
|
||||
projectionStatus: checkpointStatus,
|
||||
projectionHealth,
|
||||
terminal,
|
||||
sealed: terminal,
|
||||
diagnostic: { projectionStatus: checkpointStatus, projectionHealth, valuesRedacted: true },
|
||||
createdAt: session.startedAt ?? session.updatedAt ?? null,
|
||||
updatedAt: session.updatedAt ?? null,
|
||||
valuesRedacted: true
|
||||
}] : []
|
||||
};
|
||||
}
|
||||
|
||||
function normalizeTestMessages(session, { traceId, status, finalText, terminal }) {
|
||||
const raw = Array.isArray(session.session?.messages) ? session.session.messages : [];
|
||||
return raw.map((message, index) => {
|
||||
const role = message.role ?? (index % 2 === 0 ? "user" : "agent");
|
||||
const isAssistant = role === "agent" || role === "assistant";
|
||||
const text = isAssistant && finalText ? finalText : String(message.text ?? message.content ?? "");
|
||||
const messageStatus = isAssistant && finalText ? status : normalizeTestStatus(message.status ?? (role === "user" ? "sent" : status));
|
||||
return {
|
||||
messageId: message.messageId ?? `msg_${session.id}_${index}`,
|
||||
sessionId: session.id,
|
||||
turnId: message.turnId ?? traceId,
|
||||
traceId: message.traceId ?? traceId,
|
||||
role,
|
||||
status: messageStatus,
|
||||
projectedSeq: index + 1,
|
||||
sourceSeq: index + 1,
|
||||
sourceEventId: `${session.id}:message:${index}`,
|
||||
terminal: isAssistant && terminal,
|
||||
sealed: isAssistant && terminal,
|
||||
text,
|
||||
createdAt: message.createdAt ?? session.startedAt ?? session.updatedAt ?? null,
|
||||
updatedAt: message.updatedAt ?? session.updatedAt ?? null,
|
||||
valuesRedacted: true
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
function normalizeTestEvents(events, session, { traceId }) {
|
||||
return events.map((event, index) => {
|
||||
const seq = Number.isFinite(Number(event.seq)) ? Number(event.seq) : index + 1;
|
||||
const sourceSeq = Number.isFinite(Number(event.sourceSeq)) ? Number(event.sourceSeq) : seq;
|
||||
return {
|
||||
...event,
|
||||
id: event.id ?? `wte_${session.id}_${index}`,
|
||||
traceId: event.traceId ?? traceId,
|
||||
sessionId: event.sessionId ?? session.id,
|
||||
turnId: event.turnId ?? traceId,
|
||||
messageId: event.messageId ?? null,
|
||||
sourceSeq,
|
||||
sourceEventId: event.sourceEventId ?? `${traceId}:${sourceSeq}:${index}`,
|
||||
projectedSeq: seq,
|
||||
eventType: event.eventType ?? event.type ?? event.label ?? "event",
|
||||
terminal: event.terminal === true,
|
||||
sealed: event.terminal === true,
|
||||
occurredAt: event.occurredAt ?? event.createdAt ?? session.updatedAt ?? null,
|
||||
updatedAt: event.updatedAt ?? event.createdAt ?? session.updatedAt ?? null,
|
||||
valuesRedacted: true
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
function filterFacts(facts, params = {}) {
|
||||
const filtered = {
|
||||
sessions: facts.sessions.filter((record) => matchesFact(record, params, ["sessionId", "ownerUserId", "projectId", "conversationId", "threadId", "status"]) && (!params.traceId || record.lastTraceId === params.traceId)),
|
||||
messages: facts.messages.filter((record) => matchesFact(record, params, ["messageId", "sessionId", "turnId", "traceId", "role", "status"])),
|
||||
parts: facts.parts.filter((record) => matchesFact(record, params, ["partId", "messageId", "sessionId", "turnId", "traceId", "partType", "status"])),
|
||||
turns: facts.turns.filter((record) => matchesFact(record, params, ["turnId", "sessionId", "traceId", "messageId", "status"])),
|
||||
traceEvents: facts.traceEvents.filter((record) => matchesFact(record, params, ["id", "traceId", "sessionId", "turnId", "messageId", "eventType"])),
|
||||
checkpoints: facts.checkpoints.filter((record) => matchesFact(record, params, ["traceId", "sessionId", "turnId", "runId", "commandId", "projectionStatus", "projectionHealth"]))
|
||||
};
|
||||
const limit = Number.parseInt(params.limit ?? "", 10);
|
||||
if (!Number.isInteger(limit) || limit <= 0) return filtered;
|
||||
return Object.fromEntries(Object.entries(filtered).map(([key, value]) => [key, value.slice(0, limit)]));
|
||||
}
|
||||
|
||||
function matchesFact(record, params, fields) {
|
||||
return fields.every((field) => params[field] === undefined || params[field] === null || record[field] === params[field]);
|
||||
}
|
||||
|
||||
function mergeFacts(target, source) {
|
||||
for (const key of Object.keys(target)) target[key].push(...source[key]);
|
||||
return target;
|
||||
}
|
||||
|
||||
function emptyFacts() {
|
||||
return { sessions: [], messages: [], parts: [], turns: [], traceEvents: [], checkpoints: [] };
|
||||
}
|
||||
|
||||
function normalizeTestStatus(value) {
|
||||
const status = String(value ?? "").trim().toLowerCase().replace(/_/gu, "-");
|
||||
if (status === "active" || status === "busy" || status === "processing") return "running";
|
||||
if (status === "cancelled") return "canceled";
|
||||
return status || "unknown";
|
||||
}
|
||||
|
||||
const testTerminalStatuses = new Set(["completed", "failed", "blocked", "timeout", "canceled", "idle"]);
|
||||
|
||||
async function getSseEvents(port, path, count) {
|
||||
const controller = new AbortController();
|
||||
const response = await fetch(`http://127.0.0.1:${port}${path}`, { signal: controller.signal });
|
||||
|
||||
@@ -14,7 +14,7 @@ import {
|
||||
sendJson
|
||||
} from "./server-http-utils.ts";
|
||||
import { createWorkbenchReadModel } from "./workbench-read-model.ts";
|
||||
import { createWorkbenchTurnProjection, RUNNING_STATUSES, TERMINAL_STATUSES } from "./workbench-turn-projection.ts";
|
||||
import { createWorkbenchTurnProjection, durableTraceStatus, RUNNING_STATUSES, TERMINAL_STATUSES } from "./workbench-turn-projection.ts";
|
||||
|
||||
const DEFAULT_PAGE_LIMIT = 50;
|
||||
const DEFAULT_SESSION_LIST_LIMIT = 20;
|
||||
@@ -318,14 +318,21 @@ async function handleWorkbenchSessionList(response, url, options, actor) {
|
||||
const includeSessionId = safeSessionId(url.searchParams.get("includeSessionId"));
|
||||
const includeRouteId = includeSessionId ?? safeConversationId(url.searchParams.get("includeSessionId"));
|
||||
const readModel = createWorkbenchReadModel(options, actor);
|
||||
const naturalPage = await (perf ? perf.measure("workbench_session_page_query", () => readModel.listSessions({ limit: limit + 1, offset })) : readModel.listSessions({ limit: limit + 1, offset }));
|
||||
const pageSessions = naturalPage.slice(0, limit);
|
||||
let responseSessions = pageSessions;
|
||||
if (includeRouteId && !pageSessions.some((session) => sessionMatchesRouteId(session, includeRouteId))) {
|
||||
const included = await (perf ? perf.measure("workbench_session_include_query", () => readModel.getSessionByRouteId(includeRouteId)) : readModel.getSessionByRouteId(includeRouteId));
|
||||
if (included) responseSessions = [included, ...pageSessions];
|
||||
const pageResult = await (perf ? perf.measure("workbench_session_page_query", () => readModel.queryFacts({ ownerUserId: actor.role === "admin" ? undefined : actor.id, limit: limit + offset + 1 })) : readModel.queryFacts({ ownerUserId: actor.role === "admin" ? undefined : actor.id, limit: limit + offset + 1 }));
|
||||
if (pageResult.error) return sendJson(response, 503, workbenchProjectionStoreError(pageResult.error));
|
||||
let pageFacts = pageResult.facts;
|
||||
let naturalPage = visibleFactSessions(pageFacts, actor).slice(offset, offset + limit + 1);
|
||||
if (includeRouteId && !naturalPage.some((session) => factSessionMatchesRouteId(session, includeRouteId))) {
|
||||
const includeResult = await (perf ? perf.measure("workbench_session_include_query", () => queryFactsByRouteId(readModel, includeRouteId)) : queryFactsByRouteId(readModel, includeRouteId));
|
||||
if (includeResult.error) return sendJson(response, 503, workbenchProjectionStoreError(includeResult.error));
|
||||
const included = visibleFactSessions(includeResult.facts, actor)[0] ?? null;
|
||||
if (included) {
|
||||
pageFacts = mergeFactSets(pageFacts, includeResult.facts);
|
||||
naturalPage = [included, ...naturalPage];
|
||||
}
|
||||
}
|
||||
const summaries = await (perf ? perf.measure("workbench_session_summary_projection", () => sessionListSummaries(readModel, responseSessions, options)) : sessionListSummaries(readModel, responseSessions, options));
|
||||
const pageSessions = naturalPage.slice(0, limit);
|
||||
const summaries = pageSessions.map((session) => factSessionSummary(session, pageFacts)).filter(Boolean);
|
||||
const hasMore = naturalPage.length > limit;
|
||||
sendJson(response, 200, {
|
||||
ok: true,
|
||||
@@ -348,16 +355,428 @@ function sessionMatchesRouteId(session, routeId) {
|
||||
return Boolean(conversationId && session?.conversationId === conversationId);
|
||||
}
|
||||
|
||||
async function queryFactsByRouteId(readModel, routeId) {
|
||||
const sessionId = safeSessionId(routeId);
|
||||
if (sessionId) return readModel.queryFacts({ sessionId, limit: MAX_PAGE_LIMIT });
|
||||
const conversationId = safeConversationId(routeId);
|
||||
if (conversationId) return readModel.queryFacts({ conversationId, limit: MAX_PAGE_LIMIT });
|
||||
return { facts: emptyFactSet(), count: 0, durable: false, persistence: null };
|
||||
}
|
||||
|
||||
function visibleFactSessions(facts, actor) {
|
||||
return factArray(facts?.sessions)
|
||||
.filter((session) => canActorReadFactSession(session, actor))
|
||||
.sort(compareFactSessionsDesc);
|
||||
}
|
||||
|
||||
function canActorReadFactSession(session, actor) {
|
||||
if (!session || !actor) return false;
|
||||
if (normalizeStatus(session.status) === "archived") return false;
|
||||
if (actor.role === "admin") return true;
|
||||
return session.ownerUserId === actor.id;
|
||||
}
|
||||
|
||||
function factSessionMatchesRouteId(session, routeId) {
|
||||
const sessionId = safeSessionId(routeId);
|
||||
if (sessionId && factSessionId(session) === sessionId) return true;
|
||||
const conversationId = safeConversationId(routeId);
|
||||
return Boolean(conversationId && session?.conversationId === conversationId);
|
||||
}
|
||||
|
||||
function factSessionSummary(session, facts = {}) {
|
||||
const sessionId = factSessionId(session);
|
||||
if (!sessionId) return null;
|
||||
const traceId = factLastTraceId(session) ?? factLatestTraceIdForSession(facts, sessionId);
|
||||
const projection = traceId ? factProjectionForTrace(facts, traceId) : null;
|
||||
const turn = traceId ? factTurnForTrace(facts, traceId) : null;
|
||||
const messages = factMessagesForSession(session, facts);
|
||||
const status = normalizeStatus(turn?.status ?? session?.status);
|
||||
return {
|
||||
sessionId,
|
||||
threadId: safeOpaqueId(session?.threadId) ?? (textValue(session?.threadId) || null),
|
||||
agentId: session?.agentId ?? "hwlab-code-agent",
|
||||
status,
|
||||
running: RUNNING_STATUSES.has(status),
|
||||
terminal: TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status),
|
||||
lastTraceId: traceId,
|
||||
projection,
|
||||
projectionStatus: projection?.projectionStatus ?? null,
|
||||
projectionHealth: projection?.projectionHealth ?? null,
|
||||
staleMs: projection?.staleMs ?? null,
|
||||
blocker: projection?.blocker ?? null,
|
||||
providerProfile: textValue(session?.providerProfile ?? session?.sessionJson?.providerProfile) || null,
|
||||
messageCount: messages.length,
|
||||
firstUserMessagePreview: firstUserPreview(messages),
|
||||
updatedAt: factUpdatedAt(session),
|
||||
turnSummary: turn ? {
|
||||
turnId: factTurnId(turn, traceId),
|
||||
traceId,
|
||||
status: normalizeStatus(turn.status ?? status),
|
||||
running: RUNNING_STATUSES.has(normalizeStatus(turn.status ?? status)),
|
||||
terminal: TERMINAL_STATUSES.has(normalizeStatus(turn.status ?? status)) && !RUNNING_STATUSES.has(normalizeStatus(turn.status ?? status)),
|
||||
eventCount: projection?.lastProjectedSeq ?? factTraceSnapshot(facts, traceId).eventCount,
|
||||
projection,
|
||||
projectionStatus: projection?.projectionStatus ?? null,
|
||||
projectionHealth: projection?.projectionHealth ?? null,
|
||||
staleMs: projection?.staleMs ?? null,
|
||||
blocker: projection?.blocker ?? null,
|
||||
updatedAt: factUpdatedAt(turn)
|
||||
} : null,
|
||||
valuesRedacted: true
|
||||
};
|
||||
}
|
||||
|
||||
function factSessionDetail(session, facts = {}) {
|
||||
const sessionId = factSessionId(session);
|
||||
return {
|
||||
...factSessionSummary(session, facts),
|
||||
metadata: {
|
||||
startedAt: session?.startedAt ?? session?.createdAt ?? null,
|
||||
endedAt: session?.endedAt ?? null,
|
||||
ownerUserId: session?.ownerUserId ?? null
|
||||
},
|
||||
messagePageUrl: sessionId ? `/v1/workbench/sessions/${encodeURIComponent(sessionId)}/messages` : null,
|
||||
valuesRedacted: true,
|
||||
secretMaterialStored: false
|
||||
};
|
||||
}
|
||||
|
||||
function factMessagesForSession(session, facts = {}) {
|
||||
const sessionId = factSessionId(session);
|
||||
if (!sessionId) return [];
|
||||
const partsByMessageId = new Map();
|
||||
for (const part of factArray(facts.parts)) {
|
||||
if (part.sessionId !== sessionId) continue;
|
||||
const messageId = textValue(part.messageId);
|
||||
if (!messageId) continue;
|
||||
const existing = partsByMessageId.get(messageId) ?? [];
|
||||
existing.push(part);
|
||||
partsByMessageId.set(messageId, existing);
|
||||
}
|
||||
return factArray(facts.messages)
|
||||
.filter((message) => message.sessionId === sessionId)
|
||||
.sort(compareFactMessagesAsc)
|
||||
.map((message) => factMessageDto(message, partsByMessageId.get(message.messageId) ?? []))
|
||||
.filter(Boolean);
|
||||
}
|
||||
|
||||
function factMessageDto(message, parts = []) {
|
||||
const messageId = safeMessageId(message?.messageId) || textValue(message?.messageId);
|
||||
if (!messageId) return null;
|
||||
const traceId = safeTraceId(message?.traceId) ?? null;
|
||||
const text = projectionText(message?.text, message?.content, message?.message, message?.finalResponse);
|
||||
const normalizedParts = parts.length > 0
|
||||
? [...parts].sort(compareFactPartsAsc).map((part) => factPartDto(part, messageId, traceId)).filter(Boolean)
|
||||
: text ? [partFact({ type: "text", text, status: message?.status }, 0, messageId, traceId)] : [];
|
||||
return {
|
||||
messageId,
|
||||
role: textValue(message?.role) || "agent",
|
||||
sessionId: (safeSessionId(message?.sessionId) ?? textValue(message?.sessionId)) || null,
|
||||
traceId,
|
||||
turnId: safeTurnId(message?.turnId) || traceId,
|
||||
status: normalizeStatus(message?.status),
|
||||
parts: normalizedParts,
|
||||
text: text || "",
|
||||
textPreview: text ? text.slice(0, 240) : null,
|
||||
createdAt: textValue(message?.createdAt) || null,
|
||||
updatedAt: factUpdatedAt(message),
|
||||
projectionStatus: null,
|
||||
projectionHealth: null,
|
||||
valuesRedacted: message?.valuesRedacted !== false
|
||||
};
|
||||
}
|
||||
|
||||
function factPartDto(part, messageId, traceId) {
|
||||
const partId = safePartId(part?.partId) || textValue(part?.partId) || `prt_${hash(`${messageId}:${part?.partIndex ?? 0}:${part?.partType ?? part?.type ?? "text"}`).slice(0, 24)}`;
|
||||
const text = projectionText(part?.text, part?.content, part?.message);
|
||||
return {
|
||||
partId,
|
||||
messageId,
|
||||
traceId: safeTraceId(part?.traceId) ?? traceId,
|
||||
type: textValue(part?.partType ?? part?.type) || "text",
|
||||
text: text || null,
|
||||
status: normalizeStatus(part?.status),
|
||||
toolName: textValue(part?.toolName ?? part?.name) || null,
|
||||
createdAt: textValue(part?.createdAt ?? part?.occurredAt) || null,
|
||||
valuesRedacted: part?.valuesRedacted !== false
|
||||
};
|
||||
}
|
||||
|
||||
function factTurnForTrace(facts = {}, traceId, turnId = null) {
|
||||
const safeTrace = safeTraceId(traceId);
|
||||
if (!safeTrace) return null;
|
||||
const requestedTurn = safeTurnId(turnId);
|
||||
const matches = factArray(facts.turns)
|
||||
.filter((turn) => turn.traceId === safeTrace && (!requestedTurn || turn.turnId === requestedTurn || turn.turnId === safeTrace))
|
||||
.sort(compareFactRecordsDesc);
|
||||
return matches[0] ?? null;
|
||||
}
|
||||
|
||||
function factTurnSnapshot({ turn = null, session = null, facts = {}, traceId, turnId = null } = {}) {
|
||||
const safeTrace = safeTraceId(traceId);
|
||||
const resolvedTurnId = factTurnId(turn, safeTrace) ?? safeTurnId(turnId) ?? safeTrace;
|
||||
const status = normalizeStatus(turn?.status ?? session?.status);
|
||||
const messages = session ? factMessagesForSession(session, facts) : [];
|
||||
const userMessage = messages.find((message) => message.role === "user") ?? null;
|
||||
const assistantMessage = [...messages].reverse().find((message) => isAssistantLikeRole(message.role)) ?? null;
|
||||
const checkpoint = safeTrace ? factCheckpointForTrace(facts, safeTrace) : null;
|
||||
const finalText = projectionText(turn?.finalResponse, turn?.assistantText, turn?.text, assistantMessage?.text);
|
||||
const trace = factTraceSnapshot(facts, safeTrace);
|
||||
return {
|
||||
turnId: resolvedTurnId,
|
||||
traceId: safeTrace,
|
||||
status,
|
||||
running: RUNNING_STATUSES.has(status),
|
||||
terminal: TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status),
|
||||
sessionId: factSessionId(session) ?? turn?.sessionId ?? checkpoint?.sessionId ?? null,
|
||||
threadId: safeOpaqueId(session?.threadId) ?? (textValue(session?.threadId) || null),
|
||||
userMessageId: userMessage?.messageId ?? null,
|
||||
assistantMessageId: assistantMessage?.messageId ?? turn?.messageId ?? null,
|
||||
assistantText: finalText ?? null,
|
||||
finalResponse: turn?.finalResponse ?? (finalText ? { text: finalText, status, traceId: safeTrace, valuesPrinted: false } : null),
|
||||
agentRun: checkpoint ? {
|
||||
runId: textValue(checkpoint.runId) || null,
|
||||
commandId: textValue(checkpoint.commandId) || null,
|
||||
status,
|
||||
lastSeq: factSeq(checkpoint),
|
||||
valuesRedacted: true
|
||||
} : null,
|
||||
trace: {
|
||||
traceId: safeTrace,
|
||||
status: trace.status,
|
||||
eventCount: trace.eventCount,
|
||||
updatedAt: trace.updatedAt
|
||||
},
|
||||
urls: {
|
||||
self: resolvedTurnId ? `/v1/workbench/turns/${encodeURIComponent(resolvedTurnId)}` : null,
|
||||
traceEvents: safeTrace ? `/v1/workbench/traces/${encodeURIComponent(safeTrace)}/events` : null
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
function factTraceSnapshot(facts = {}, traceId) {
|
||||
const safeTrace = safeTraceId(traceId);
|
||||
const events = factArray(facts.traceEvents)
|
||||
.filter((event) => !safeTrace || event.traceId === safeTrace)
|
||||
.sort(compareFactTraceEventsAsc)
|
||||
.map(factTraceEventDto);
|
||||
const status = durableTraceStatus(events);
|
||||
const lastEvent = events.at(-1) ?? null;
|
||||
return {
|
||||
traceId: safeTrace,
|
||||
status,
|
||||
eventCount: events.length,
|
||||
events,
|
||||
lastEvent,
|
||||
updatedAt: lastEvent?.updatedAt ?? lastEvent?.createdAt ?? factCheckpointForTrace(facts, safeTrace)?.updatedAt ?? null,
|
||||
valuesRedacted: true,
|
||||
secretMaterialStored: false
|
||||
};
|
||||
}
|
||||
|
||||
function factTraceEventDto(event, index) {
|
||||
const seq = Number.isFinite(Number(event?.seq)) && Number(event.seq) > 0 ? Math.trunc(Number(event.seq)) : factSeq(event) || index + 1;
|
||||
return {
|
||||
...event,
|
||||
id: textValue(event?.id) || textValue(event?.sourceEventId) || null,
|
||||
seq,
|
||||
sourceSeq: factSeq(event) || seq,
|
||||
type: textValue(event?.type ?? event?.eventType) || "event",
|
||||
label: textValue(event?.label ?? event?.eventType ?? event?.type) || null,
|
||||
status: normalizeStatus(event?.status),
|
||||
createdAt: textValue(event?.createdAt ?? event?.occurredAt) || null,
|
||||
updatedAt: factUpdatedAt(event),
|
||||
valuesRedacted: event?.valuesRedacted !== false
|
||||
};
|
||||
}
|
||||
|
||||
function factProjectionForTrace(facts = {}, traceId) {
|
||||
const checkpoint = factCheckpointForTrace(facts, traceId);
|
||||
if (!checkpoint) {
|
||||
return {
|
||||
projectionStatus: "unknown",
|
||||
projectionHealth: "unknown",
|
||||
lastProjectedSeq: null,
|
||||
sourceRunId: null,
|
||||
sourceCommandId: null,
|
||||
staleMs: null,
|
||||
blocker: null,
|
||||
updatedAt: null,
|
||||
valuesRedacted: true
|
||||
};
|
||||
}
|
||||
const projectionStatus = normalizeProjectionStatus(checkpoint.projectionStatus);
|
||||
const projectionHealth = normalizeProjectionHealth(checkpoint.projectionHealth, projectionStatus);
|
||||
const diagnostic = objectValue(checkpoint.diagnostic);
|
||||
return {
|
||||
projectionStatus,
|
||||
projectionHealth,
|
||||
lastProjectedSeq: factSeq(checkpoint),
|
||||
sourceRunId: textValue(checkpoint.runId ?? checkpoint.sourceRunId) || null,
|
||||
sourceCommandId: textValue(checkpoint.commandId ?? checkpoint.sourceCommandId) || null,
|
||||
staleMs: null,
|
||||
blocker: diagnostic.blocker ?? checkpoint.blocker ?? null,
|
||||
updatedAt: factUpdatedAt(checkpoint),
|
||||
valuesRedacted: true
|
||||
};
|
||||
}
|
||||
|
||||
function factCheckpointForTrace(facts = {}, traceId) {
|
||||
const safeTrace = safeTraceId(traceId);
|
||||
if (!safeTrace) return null;
|
||||
return factArray(facts.checkpoints)
|
||||
.filter((checkpoint) => checkpoint.traceId === safeTrace)
|
||||
.sort(compareFactRecordsDesc)[0] ?? null;
|
||||
}
|
||||
|
||||
function workbenchProjectionStoreError(error) {
|
||||
return workbenchError("projection_store_unavailable", "Workbench durable projection store is unavailable.", { causeCode: error?.code ?? null });
|
||||
}
|
||||
|
||||
function factSessionId(session) {
|
||||
return safeSessionId(session?.sessionId ?? session?.id) ?? null;
|
||||
}
|
||||
|
||||
function factLastTraceId(session) {
|
||||
return safeTraceId(session?.lastTraceId ?? session?.traceId) ?? null;
|
||||
}
|
||||
|
||||
function factLatestTraceIdForSession(facts = {}, sessionId) {
|
||||
const turn = factArray(facts.turns).filter((item) => item.sessionId === sessionId).sort(compareFactRecordsDesc)[0] ?? null;
|
||||
return safeTraceId(turn?.traceId) ?? null;
|
||||
}
|
||||
|
||||
function factTurnId(turn, fallbackTraceId = null) {
|
||||
return safeTurnId(turn?.turnId) ?? safeTraceId(turn?.turnId) ?? safeTraceId(fallbackTraceId) ?? null;
|
||||
}
|
||||
|
||||
function factUpdatedAt(record) {
|
||||
return textValue(record?.updatedAt ?? record?.occurredAt ?? record?.createdAt) || null;
|
||||
}
|
||||
|
||||
function factSeq(record) {
|
||||
for (const value of [record?.projectedSeq, record?.sourceSeq, record?.seq]) {
|
||||
const parsed = Number(value);
|
||||
if (Number.isFinite(parsed) && parsed >= 0) return Math.trunc(parsed);
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
function normalizeProjectionStatus(value) {
|
||||
const status = normalizeStatus(value);
|
||||
if (status === "caught-up") return "caught-up";
|
||||
if (status === "caughtup") return "caught-up";
|
||||
return ["projecting", "blocked", "stalled", "unknown"].includes(status) ? status : "unknown";
|
||||
}
|
||||
|
||||
function normalizeProjectionHealth(value, projectionStatus = "unknown") {
|
||||
const health = normalizeStatus(value);
|
||||
if (health === "healthy" && projectionStatus === "caught-up") return "caught-up";
|
||||
if (health === "healthy" && projectionStatus === "projecting") return "projecting";
|
||||
if (["caught-up", "projecting", "degraded", "stalled", "unavailable", "unknown"].includes(health)) return health;
|
||||
return projectionStatus === "caught-up" || projectionStatus === "projecting" ? projectionStatus : "unknown";
|
||||
}
|
||||
|
||||
function compareFactSessionsDesc(left, right) {
|
||||
return compareTimestampDesc(factUpdatedAt(left), factUpdatedAt(right)) || compareText(factSessionId(left), factSessionId(right));
|
||||
}
|
||||
|
||||
function compareFactMessagesAsc(left, right) {
|
||||
return compareNumberAsc(factSeq(left), factSeq(right)) || compareTimestampAsc(left?.createdAt ?? left?.updatedAt, right?.createdAt ?? right?.updatedAt) || compareText(left?.messageId, right?.messageId);
|
||||
}
|
||||
|
||||
function compareFactPartsAsc(left, right) {
|
||||
return compareNumberAsc(left?.partIndex, right?.partIndex) || compareNumberAsc(factSeq(left), factSeq(right)) || compareText(left?.partId, right?.partId);
|
||||
}
|
||||
|
||||
function compareFactTraceEventsAsc(left, right) {
|
||||
return compareNumberAsc(factSeq(left), factSeq(right)) || compareTimestampAsc(left?.occurredAt ?? left?.createdAt, right?.occurredAt ?? right?.createdAt) || compareText(left?.id, right?.id);
|
||||
}
|
||||
|
||||
function compareFactRecordsDesc(left, right) {
|
||||
return compareNumberDesc(factSeq(left), factSeq(right)) || compareTimestampDesc(factUpdatedAt(left), factUpdatedAt(right));
|
||||
}
|
||||
|
||||
function compareTimestampAsc(left, right) {
|
||||
return timestampMs(left) - timestampMs(right);
|
||||
}
|
||||
|
||||
function compareTimestampDesc(left, right) {
|
||||
return timestampMs(right) - timestampMs(left);
|
||||
}
|
||||
|
||||
function compareNumberAsc(left, right) {
|
||||
return numericValue(left, Number.MAX_SAFE_INTEGER) - numericValue(right, Number.MAX_SAFE_INTEGER);
|
||||
}
|
||||
|
||||
function compareNumberDesc(left, right) {
|
||||
return numericValue(right, -1) - numericValue(left, -1);
|
||||
}
|
||||
|
||||
function compareText(left, right) {
|
||||
return textValue(left).localeCompare(textValue(right));
|
||||
}
|
||||
|
||||
function timestampMs(value) {
|
||||
const parsed = Date.parse(String(value ?? ""));
|
||||
return Number.isFinite(parsed) ? parsed : 0;
|
||||
}
|
||||
|
||||
function numericValue(value, fallback) {
|
||||
const parsed = Number(value);
|
||||
return Number.isFinite(parsed) ? parsed : fallback;
|
||||
}
|
||||
|
||||
function mergeFactSets(...sets) {
|
||||
const merged = emptyFactSet();
|
||||
const keys = {
|
||||
sessions: "sessionId",
|
||||
messages: "messageId",
|
||||
parts: "partId",
|
||||
turns: "turnId",
|
||||
traceEvents: "id",
|
||||
checkpoints: "traceId"
|
||||
};
|
||||
for (const set of sets) {
|
||||
for (const [key, idKey] of Object.entries(keys)) {
|
||||
const seen = new Set(merged[key].map((item) => textValue(item?.[idKey])));
|
||||
for (const item of factArray(set?.[key])) {
|
||||
const id = textValue(item?.[idKey]);
|
||||
if (id && seen.has(id)) continue;
|
||||
merged[key].push(item);
|
||||
if (id) seen.add(id);
|
||||
}
|
||||
}
|
||||
}
|
||||
return merged;
|
||||
}
|
||||
|
||||
function emptyFactSet() {
|
||||
return {
|
||||
sessions: [],
|
||||
messages: [],
|
||||
parts: [],
|
||||
turns: [],
|
||||
traceEvents: [],
|
||||
checkpoints: []
|
||||
};
|
||||
}
|
||||
|
||||
function factArray(value) {
|
||||
return Array.isArray(value) ? value : [];
|
||||
}
|
||||
|
||||
async function handleWorkbenchSessionDetail(response, options, actor, sessionId) {
|
||||
const readModel = createWorkbenchReadModel(options, actor);
|
||||
const session = await readModel.getSessionByRouteId(sessionId);
|
||||
const result = await queryFactsByRouteId(readModel, sessionId);
|
||||
if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error));
|
||||
const session = visibleFactSessions(result.facts, actor)[0] ?? null;
|
||||
if (!session) return sendJson(response, 404, workbenchError("workbench_session_not_found", "Workbench session is not visible to the current actor.", { sessionId }));
|
||||
const projectionOptions = await sessionProjectionOptions(readModel, session, options);
|
||||
sendJson(response, 200, {
|
||||
ok: true,
|
||||
status: "found",
|
||||
contractVersion: "workbench-session-detail-v1",
|
||||
session: sessionDetail(session, projectionOptions),
|
||||
session: factSessionDetail(session, result.facts),
|
||||
valuesRedacted: true,
|
||||
secretMaterialStored: false
|
||||
});
|
||||
@@ -365,19 +784,21 @@ async function handleWorkbenchSessionDetail(response, options, actor, sessionId)
|
||||
|
||||
async function handleWorkbenchMessagePage(response, url, options, actor, sessionId) {
|
||||
const readModel = createWorkbenchReadModel(options, actor);
|
||||
const session = await readModel.getSessionByRouteId(sessionId);
|
||||
const result = await queryFactsByRouteId(readModel, sessionId);
|
||||
if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error));
|
||||
const session = visibleFactSessions(result.facts, actor)[0] ?? null;
|
||||
if (!session) return sendJson(response, 404, workbenchError("workbench_session_not_found", "Workbench session is not visible to the current actor.", { sessionId }));
|
||||
const projectionOptions = await sessionProjectionOptions(readModel, session, options);
|
||||
const messages = sessionMessages(session, projectionOptions);
|
||||
const messages = factMessagesForSession(session, result.facts);
|
||||
const limit = boundedLimit(url.searchParams.get("limit"));
|
||||
const offset = cursorOffset(url.searchParams.get("cursor") ?? url.searchParams.get("after"));
|
||||
const page = messages.slice(offset, offset + limit);
|
||||
const nextOffset = offset + page.length;
|
||||
const resolvedSessionId = factSessionId(session);
|
||||
sendJson(response, 200, {
|
||||
ok: true,
|
||||
status: "succeeded",
|
||||
contractVersion: "workbench-message-page-v1",
|
||||
sessionId: session.id,
|
||||
sessionId: resolvedSessionId,
|
||||
messages: page,
|
||||
count: page.length,
|
||||
total: messages.length,
|
||||
@@ -400,30 +821,24 @@ async function handleWorkbenchTurnSnapshot(response, url, options, actor, rawTur
|
||||
return sendJson(response, 400, body);
|
||||
}
|
||||
const readModel = createWorkbenchReadModel(options, actor);
|
||||
const result = readModel.resultForTrace(traceId);
|
||||
const session = await readModel.getSessionByTraceId(traceId);
|
||||
if (result?.ownerUserId && !readModel.canReadOwner(result.ownerUserId)) {
|
||||
const body = workbenchError("agent_session_owner_required", "Only the session owner or admin can read this Workbench turn.", { traceId });
|
||||
recordWorkbenchTurnReadMetric(options, url, 403, body, startedAt);
|
||||
return sendJson(response, 403, body);
|
||||
}
|
||||
const resultVisible = result && (session || actor.role === "admin" || Boolean(result.ownerUserId));
|
||||
const trace = await readModel.traceSnapshot(traceId);
|
||||
const turnProjection = createWorkbenchTurnProjection({ turnId, traceId, result, session, trace });
|
||||
const status = turnProjection.status;
|
||||
const found = Boolean(resultVisible || session);
|
||||
const result = await readModel.queryFacts({ traceId, limit: MAX_PAGE_LIMIT });
|
||||
if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error));
|
||||
const session = visibleFactSessions(result.facts, actor).find((item) => item.lastTraceId === traceId) ?? visibleFactSessions(result.facts, actor)[0] ?? null;
|
||||
const turn = factTurnForTrace(result.facts, traceId, turnId);
|
||||
const found = Boolean(session || turn);
|
||||
if (!found) {
|
||||
const body = workbenchError("workbench_turn_not_found", "Workbench turn is not visible to the current actor.", { turnId, traceId });
|
||||
recordWorkbenchTurnReadMetric(options, url, 404, body, startedAt);
|
||||
return sendJson(response, 404, body);
|
||||
}
|
||||
const projection = readModel.projectionDiagnostics({ traceId, result, trace, projection: turnProjection });
|
||||
recordWorkbenchProjectionMetric(options, { result, trace, projection: turnProjection, diagnostic: projection });
|
||||
const projection = factProjectionForTrace(result.facts, traceId);
|
||||
const snapshot = factTurnSnapshot({ turn, session, facts: result.facts, traceId, turnId });
|
||||
const status = snapshot.status;
|
||||
const body = {
|
||||
ok: true,
|
||||
status,
|
||||
contractVersion: "workbench-turn-snapshot-v1",
|
||||
turn: turnSnapshot({ projection: turnProjection, result, session, trace }),
|
||||
turn: snapshot,
|
||||
projection,
|
||||
projectionStatus: projection.projectionStatus,
|
||||
projectionHealth: projection.projectionHealth,
|
||||
@@ -435,6 +850,10 @@ async function handleWorkbenchTurnSnapshot(response, url, options, actor, rawTur
|
||||
valuesRedacted: true,
|
||||
secretMaterialStored: false
|
||||
};
|
||||
recordWorkbenchProjectionMetric(options, {
|
||||
projection: { ...projection, status, eventCount: snapshot.trace?.eventCount ?? projection.lastProjectedSeq ?? 0 },
|
||||
diagnostic: projection
|
||||
});
|
||||
recordWorkbenchTurnReadMetric(options, url, 200, body, startedAt);
|
||||
sendJson(response, 200, body);
|
||||
}
|
||||
@@ -443,27 +862,21 @@ async function handleWorkbenchTraceEventPage(response, url, options, actor, rawT
|
||||
const traceId = safeTraceId(rawTraceId);
|
||||
if (!traceId) return sendJson(response, 400, workbenchError("invalid_trace_id", "traceId must start with trc_.", { traceId: rawTraceId }));
|
||||
const readModel = createWorkbenchReadModel(options, actor);
|
||||
const result = readModel.resultForTrace(traceId);
|
||||
const session = await readModel.getSessionByTraceId(traceId);
|
||||
if (result?.ownerUserId && !readModel.canReadOwner(result.ownerUserId)) {
|
||||
return sendJson(response, 403, workbenchError("agent_session_owner_required", "Only the session owner or admin can read this Workbench trace.", { traceId }));
|
||||
}
|
||||
const resultVisible = result && (session || actor.role === "admin" || Boolean(result.ownerUserId));
|
||||
if (!resultVisible && !session) {
|
||||
const result = await readModel.queryFacts({ traceId, limit: MAX_PAGE_LIMIT });
|
||||
if (result.error) return sendJson(response, 503, workbenchProjectionStoreError(result.error));
|
||||
const session = visibleFactSessions(result.facts, actor).find((item) => item.lastTraceId === traceId) ?? visibleFactSessions(result.facts, actor)[0] ?? null;
|
||||
if (!session && result.facts.traceEvents.length === 0) {
|
||||
return sendJson(response, 404, workbenchError("workbench_trace_not_found", "Workbench trace is not visible to the current actor.", { traceId }));
|
||||
}
|
||||
const trace = await readModel.traceSnapshot(traceId);
|
||||
const page = traceEventPage(trace, tracePageOptions(url));
|
||||
const turnProjection = createWorkbenchTurnProjection({ traceId, result, session, trace });
|
||||
const projection = readModel.projectionDiagnostics({ traceId, result, trace, projection: turnProjection });
|
||||
recordWorkbenchProjectionMetric(options, { result, trace, projection: turnProjection, diagnostic: projection });
|
||||
const page = traceEventPage(factTraceSnapshot(result.facts, traceId), tracePageOptions(url));
|
||||
const projection = factProjectionForTrace(result.facts, traceId);
|
||||
sendJson(response, 200, {
|
||||
ok: true,
|
||||
status: "succeeded",
|
||||
contractVersion: "workbench-trace-events-v1",
|
||||
traceId,
|
||||
sessionId: safeSessionId(result?.sessionId ?? result?.session?.sessionId ?? session?.id) ?? null,
|
||||
threadId: safeOpaqueId(result?.threadId ?? result?.session?.threadId ?? session?.threadId) ?? (textValue(session?.threadId) || null),
|
||||
sessionId: factSessionId(session),
|
||||
threadId: safeOpaqueId(session?.threadId) ?? (textValue(session?.threadId) || null),
|
||||
...page,
|
||||
projection,
|
||||
projectionStatus: projection.projectionStatus,
|
||||
|
||||
@@ -16,6 +16,34 @@ export function createWorkbenchFactsStore(options = {}, actor = null) {
|
||||
const runtimeStore = options.runtimeStore ?? null;
|
||||
const results = options.codeAgentChatResults ?? null;
|
||||
|
||||
async function queryFacts(params = {}) {
|
||||
if (typeof runtimeStore?.queryWorkbenchFacts === "function") {
|
||||
try {
|
||||
const result = await runtimeStore.queryWorkbenchFacts(params);
|
||||
return {
|
||||
facts: normalizeWorkbenchFactsResult(result?.facts),
|
||||
count: Number.isFinite(Number(result?.count)) ? Number(result.count) : 0,
|
||||
durable: true,
|
||||
persistence: result?.persistence ?? null
|
||||
};
|
||||
} catch (error) {
|
||||
return {
|
||||
facts: emptyWorkbenchFacts(),
|
||||
count: 0,
|
||||
durable: true,
|
||||
error,
|
||||
projection: projectionStoreUnavailableTrace(safeTraceId(params.traceId) ?? "trc_unassigned", error)
|
||||
};
|
||||
}
|
||||
}
|
||||
return {
|
||||
facts: emptyWorkbenchFacts(),
|
||||
count: 0,
|
||||
durable: false,
|
||||
persistence: null
|
||||
};
|
||||
}
|
||||
|
||||
async function getSessionById(sessionId) {
|
||||
const safeId = safeSessionId(sessionId);
|
||||
if (!safeId) return null;
|
||||
@@ -96,6 +124,7 @@ export function createWorkbenchFactsStore(options = {}, actor = null) {
|
||||
}
|
||||
|
||||
return {
|
||||
queryFacts,
|
||||
getSessionById,
|
||||
getSessionByRouteId,
|
||||
getSessionByTraceId,
|
||||
@@ -250,3 +279,22 @@ function normalizeStatus(value) {
|
||||
if (text === "cancelled") return "canceled";
|
||||
return text || "unknown";
|
||||
}
|
||||
|
||||
function normalizeWorkbenchFactsResult(facts = {}) {
|
||||
const normalized = emptyWorkbenchFacts();
|
||||
for (const key of Object.keys(normalized)) {
|
||||
normalized[key] = Array.isArray(facts?.[key]) ? facts[key].map((item) => ({ ...item })) : [];
|
||||
}
|
||||
return normalized;
|
||||
}
|
||||
|
||||
function emptyWorkbenchFacts() {
|
||||
return {
|
||||
sessions: [],
|
||||
messages: [],
|
||||
parts: [],
|
||||
turns: [],
|
||||
traceEvents: [],
|
||||
checkpoints: []
|
||||
};
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ export function createWorkbenchReadModel(options = {}, actor = null) {
|
||||
const facts = createWorkbenchFactsStore(options, actor);
|
||||
return {
|
||||
facts,
|
||||
queryFacts: (input = {}) => facts.queryFacts(input),
|
||||
listSessions: (input = {}) => facts.listSessions(input),
|
||||
getSessionById: (sessionId) => facts.getSessionById(sessionId),
|
||||
getSessionByRouteId: (routeId) => facts.getSessionByRouteId(routeId),
|
||||
|
||||
@@ -604,7 +604,7 @@ export class CloudRuntimeStore {
|
||||
}
|
||||
|
||||
queryWorkbenchFacts(params = {}) {
|
||||
const sessions = [...this.workbenchSessions.values()].filter((record) => matchesQuery(record, params, ["sessionId", "ownerUserId", "projectId", "conversationId", "threadId", "status"]));
|
||||
const sessions = [...this.workbenchSessions.values()].filter((record) => matchesWorkbenchSessionFactQuery(record, params));
|
||||
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"]));
|
||||
@@ -1213,7 +1213,7 @@ export class PostgresCloudRuntimeStore {
|
||||
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" }),
|
||||
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", traceId: "last_trace_id", lastTraceId: "last_trace_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" }),
|
||||
@@ -2166,6 +2166,13 @@ function matchesQuery(record, params, fields) {
|
||||
return true;
|
||||
}
|
||||
|
||||
function matchesWorkbenchSessionFactQuery(record, params) {
|
||||
if (!matchesQuery(record, params, ["sessionId", "ownerUserId", "projectId", "conversationId", "threadId", "status"])) return false;
|
||||
const traceId = textOr(params.traceId ?? params.lastTraceId, "");
|
||||
if (traceId && record.lastTraceId !== traceId && record.traceId !== traceId) return false;
|
||||
return true;
|
||||
}
|
||||
|
||||
function limitResults(records, limit) {
|
||||
const parsed = Number.parseInt(limit ?? "", 10);
|
||||
if (!Number.isInteger(parsed) || parsed <= 0) {
|
||||
|
||||
Reference in New Issue
Block a user