Files
pikasTech-HWLAB/internal/cloud/workbench-read-model.ts
T

66 lines
3.2 KiB
TypeScript

/*
* SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-18-p0-unique-projection; PJ2026-010403 API契约 draft-2026-06-18-r1.
* 职责: WorkbenchReadModel 组件入口。只通过 WorkbenchFactsStore 读取 durable projection facts,不执行 AgentRun sync/finalize/repair。
*/
import { createWorkbenchFactsStore } from "./workbench-facts-store.ts";
export function createWorkbenchReadModel(options = {}, actor = null) {
const facts = createWorkbenchFactsStore(options, actor);
return {
facts,
listSessions: (input = {}) => facts.listSessions(input),
getSessionById: (sessionId) => facts.getSessionById(sessionId),
getSessionByRouteId: (routeId) => facts.getSessionByRouteId(routeId),
getSessionByTraceId: (traceId) => facts.getSessionByTraceId(traceId),
resultForTrace: (traceId) => facts.resultForTrace(traceId),
traceSnapshot: (traceId) => facts.traceSnapshot(traceId),
traceSnapshotSync: (traceId) => facts.traceSnapshotSync(traceId),
subscribeTrace: (traceId, listener) => facts.subscribeTrace(traceId, listener),
canReadOwner: (ownerUserId) => facts.canReadOwner(ownerUserId),
projectionDiagnostics: ({ traceId, result = null, trace = null } = {}) => projectionDiagnostics({ traceId, result, trace })
};
}
export function projectionDiagnostics({ traceId, result = null, trace = null } = {}) {
const agentRun = result?.agentRun && typeof result.agentRun === "object" ? result.agentRun : null;
const events = Array.isArray(trace?.events) ? trace.events : [];
const lastEvent = trace?.lastEvent ?? events.at(-1) ?? null;
const lastProjectedSeq = Number(lastEvent?.seq ?? trace?.lastProjectedSeq ?? 0);
const hasTraceProjection = Boolean(trace && trace.status !== "missing" && (events.length > 0 || lastEvent));
const sourceRunId = agentRun?.runId ?? lastEvent?.runId ?? null;
const sourceCommandId = agentRun?.commandId ?? lastEvent?.commandId ?? null;
const blocker = result?.blocker ?? result?.error ?? null;
const status = blocker
? "blocked"
: !hasTraceProjection && (result || agentRun)
? "projecting"
: isTerminalStatus(result?.status ?? trace?.status)
? "caught-up"
: result || trace?.status !== "missing"
? "projecting"
: "unknown";
return {
projectionStatus: status,
lastProjectedSeq: Number.isFinite(lastProjectedSeq) && lastProjectedSeq > 0 ? Math.trunc(lastProjectedSeq) : null,
sourceRunId,
sourceCommandId,
blocker: blocker ? diagnosticBlocker(blocker) : null,
updatedAt: trace?.updatedAt ?? result?.updatedAt ?? null,
valuesRedacted: true
};
}
function diagnosticBlocker(value = {}) {
if (!value || typeof value !== "object") return { code: "projection_blocked", summary: String(value), valuesRedacted: true };
return {
code: String(value.code ?? value.errorCode ?? "projection_blocked"),
summary: String(value.summary ?? value.userMessage ?? value.message ?? value.code ?? "Projection is blocked."),
retryable: typeof value.retryable === "boolean" ? value.retryable : null,
valuesRedacted: true
};
}
function isTerminalStatus(value) {
return ["completed", "failed", "blocked", "timeout", "cancelled", "canceled", "idle"].includes(String(value ?? "").trim().toLowerCase());
}