Files
pikasTech-HWLAB/internal/cloud/workbench-turn-projection.ts
T
2026-06-19 15:40:12 +08:00

259 lines
12 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
/*
* SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-19-p0-projector-resume; PJ2026-010403 API契约 draft-2026-06-18-r1.
* 职责: 生成唯一 Workbench turn projectiontrace event 状态只作为输入证据,不直接成为 turn lifecycle。
*/
export const TERMINAL_STATUSES = new Set(["completed", "failed", "blocked", "timeout", "cancelled", "canceled", "idle"]);
export const RUNNING_STATUSES = new Set(["running", "pending", "queued", "accepted", "dispatching", "streaming", "active", "processing", "busy", "creating"]);
export function createWorkbenchTurnProjection({ turnId = null, traceId = null, result = null, session = null, trace = null } = {}) {
const projectionTraceId = textValue(traceId ?? trace?.traceId ?? result?.traceId ?? session?.lastTraceId) || null;
const projectionTurnId = textValue(turnId) || projectionTraceId;
const terminalEvidence = terminalTurnEvidence({ result, trace });
const activeEvidence = activeTurnEvidence({ result, session, trace });
const status = terminalEvidence?.status ?? activeEvidence?.status ?? "unknown";
const running = RUNNING_STATUSES.has(status);
const terminal = Boolean(terminalEvidence && TERMINAL_STATUSES.has(status) && !running);
const finalText = terminal ? projectionText(result?.finalResponse, result?.assistantText, result?.reply, result?.text, result?.summary, trace?.finalResponse, trace?.terminalEvidence?.finalResponse) : null;
const agentRun = objectValue(result?.agentRun ?? trace?.agentRun);
const lastEvent = traceLastEvent(trace);
return {
turnId: projectionTurnId,
traceId: projectionTraceId,
status,
running,
terminal,
source: terminalEvidence?.source ?? activeEvidence?.source ?? null,
terminalEvidence,
finalResponse: finalText ? { text: finalText, status, traceId: projectionTraceId, valuesPrinted: false } : null,
assistantText: finalText,
lastProjectedSeq: traceLastSeq(trace),
sourceRunId: agentRun?.runId ?? lastEvent?.runId ?? lastEvent?.payload?.runId ?? null,
sourceCommandId: agentRun?.commandId ?? lastEvent?.commandId ?? lastEvent?.payload?.commandId ?? null,
eventCount: normalizedEventCount(trace),
updatedAt: trace?.updatedAt ?? result?.updatedAt ?? session?.updatedAt ?? null,
valuesRedacted: true
};
}
export function projectionDiagnostics({ traceId = null, projection = null, result = null, trace = null, refreshError = null } = {}) {
const turn = projection ?? createWorkbenchTurnProjection({ traceId, result, trace });
const source = projectionDiagnosticSource(trace);
const blocker = refreshError ?? source?.blocker ?? trace?.blocker ?? result?.blocker ?? result?.error ?? null;
const hasProjectionInput = hasTraceProjection(trace) || Boolean(result || result?.agentRun);
const sourceStatus = normalizeProjectionStatus(source?.projectionStatus ?? trace?.projectionStatus);
const status = sourceStatus ?? (blocker ? "blocked" : turn.terminal ? "caught-up" : hasProjectionInput ? "projecting" : "unknown");
const diagnostic = blocker ? diagnosticBlocker(blocker) : null;
const projectionHealth = normalizeProjectionHealth(source?.projectionHealth ?? trace?.projectionHealth)
?? projectionHealthFor({ status, turn, hasProjectionInput, blocker: diagnostic });
const staleMs = projectionStaleMs(source?.staleMs ?? trace?.staleMs, turn.updatedAt ?? trace?.updatedAt ?? result?.updatedAt);
return {
projectionStatus: status,
projectionHealth,
lastProjectedSeq: turn.lastProjectedSeq ?? null,
sourceRunId: turn.sourceRunId ?? null,
sourceCommandId: turn.sourceCommandId ?? null,
staleMs,
blocker: diagnostic,
updatedAt: turn.updatedAt ?? null,
valuesRedacted: true
};
}
export function durableTraceStatus(events = []) {
const terminal = terminalTraceEventEvidence(events);
if (terminal) return terminal.status;
const activeStatus = activeTraceEventStatus(events);
return activeStatus ?? (events.length > 0 ? "running" : "missing");
}
export function traceTerminalEvidence(trace = null) {
const direct = objectValue(trace?.terminalEvidence);
if (direct) {
const status = terminalStatusFromValue(direct.status ?? direct.terminalStatus ?? trace?.status) ?? "completed";
return { source: "trace-terminal-evidence", status, evidence: direct, valuesRedacted: true };
}
const finalResponse = objectValue(trace?.finalResponse);
if (finalResponse) {
const status = terminalStatusFromValue(finalResponse.status ?? trace?.status) ?? "completed";
return { source: "trace-final-response", status, evidence: { textPresent: Boolean(projectionText(finalResponse)), valuesRedacted: true }, valuesRedacted: true };
}
const events = Array.isArray(trace?.events) ? trace.events : [];
return terminalTraceEventEvidence(events);
}
export function normalizeWorkbenchStatus(value) {
const text = String(value ?? "").trim().toLowerCase().replace(/_/gu, "-");
if (text === "cancelled") return "canceled";
return text || "unknown";
}
function terminalTurnEvidence({ result = null, trace = null } = {}) {
const resultStatus = terminalStatusFromValue(result?.status ?? result?.terminalStatus ?? result?.agentRun?.terminalStatus);
if (resultStatus) return { source: "result", status: resultStatus, valuesRedacted: true };
return traceTerminalEvidence(trace);
}
function activeTurnEvidence({ result = null, session = null, trace = null } = {}) {
const resultStatus = firstStatus([
result?.status,
result?.agentRun?.status,
result?.agentRun?.commandState,
result?.agentRun?.runStatus
], RUNNING_STATUSES);
if (resultStatus) return { source: "result", status: normalizeActiveStatus(resultStatus), valuesRedacted: true };
const traceStatus = firstStatus([trace?.status, trace?.traceStatus], RUNNING_STATUSES);
if (traceStatus) return { source: "trace", status: normalizeActiveStatus(traceStatus), valuesRedacted: true };
if (hasTraceProjection(trace)) return { source: "trace-events", status: "running", valuesRedacted: true };
const sessionStatus = firstStatus([session?.status, session?.session?.sessionStatus], RUNNING_STATUSES);
if (sessionStatus) return { source: "session", status: normalizeActiveStatus(sessionStatus), valuesRedacted: true };
return null;
}
function firstStatus(values, accepted) {
for (const value of values) {
const status = normalizeWorkbenchStatus(value);
if (accepted.has(status)) return status;
}
return null;
}
function terminalStatusFromValue(value) {
const status = normalizeWorkbenchStatus(value);
return TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status) ? status : null;
}
function normalizeActiveStatus(status) {
return status === "active" || status === "busy" || status === "processing" || status === "creating" ? "running" : status;
}
function terminalTraceEventEvidence(events = []) {
for (const event of [...events].reverse()) {
if (!event || typeof event !== "object") continue;
const terminal = event.terminal === true || event.final === true || event.replyAuthority === true;
if (!terminal) continue;
const status = terminalStatusFromValue(event.status ?? event.terminalStatus) ?? (event.error || event.errorCode ? "failed" : "completed");
return {
source: "trace-terminal-event",
status,
seq: eventSeq(event, events.indexOf(event)),
eventType: textValue(event.type ?? event.label) || null,
valuesRedacted: true
};
}
return null;
}
function activeTraceEventStatus(events = []) {
for (const event of [...events].reverse()) {
const status = normalizeWorkbenchStatus(event?.status ?? event?.type);
if (RUNNING_STATUSES.has(status)) return normalizeActiveStatus(status);
}
return null;
}
function hasTraceProjection(trace) {
return Boolean(trace && trace.status !== "missing" && (normalizedEventCount(trace) > 0 || traceLastEvent(trace)));
}
function normalizedEventCount(trace) {
const events = Array.isArray(trace?.events) ? trace.events : [];
const count = Number(trace?.eventCount ?? events.length);
return Number.isFinite(count) && count >= 0 ? Math.trunc(count) : events.length;
}
function traceLastSeq(trace) {
const events = Array.isArray(trace?.events) ? trace.events : [];
const lastEvent = traceLastEvent(trace);
const indexedMax = events.reduce((max, event, index) => Math.max(max, eventSeq(event, index)), 0);
return Math.max(indexedMax, lastEvent ? eventSeq(lastEvent, Math.max(0, events.length - 1)) : 0) || null;
}
function traceLastEvent(trace) {
const events = Array.isArray(trace?.events) ? trace.events : [];
return trace?.lastEvent ?? events.at(-1) ?? null;
}
function projectionDiagnosticSource(trace) {
const direct = objectValue(trace?.projection);
if (direct) return direct;
if (trace?.projectionHealth || trace?.projectionStatus || trace?.blocker) return trace;
return null;
}
function normalizeProjectionStatus(value) {
const text = String(value ?? "").trim().toLowerCase().replace(/_/gu, "-");
return ["caught-up", "projecting", "blocked", "stalled", "unknown"].includes(text) ? text : null;
}
function normalizeProjectionHealth(value) {
const text = String(value ?? "").trim().toLowerCase().replace(/_/gu, "-");
return ["caught-up", "projecting", "degraded", "stalled", "unavailable", "unknown"].includes(text) ? text : null;
}
function projectionHealthFor({ status, turn, hasProjectionInput, blocker }) {
if (blocker?.code === "projection_store_unavailable") return "unavailable";
if (blocker) return status === "stalled" ? "stalled" : "degraded";
if (turn.terminal || status === "caught-up") return "caught-up";
if (status === "stalled") return "stalled";
if (status === "projecting" || hasProjectionInput) return "projecting";
return "unknown";
}
function projectionStaleMs(value, updatedAt) {
const direct = Number(value);
if (Number.isFinite(direct) && direct >= 0) return Math.trunc(direct);
const updatedAtMs = Date.parse(String(updatedAt ?? ""));
if (!Number.isFinite(updatedAtMs)) return null;
return Math.max(0, Date.now() - updatedAtMs);
}
function eventSeq(event, index) {
const seq = Number(event?.seq);
return Number.isFinite(seq) && seq > 0 ? Math.trunc(seq) : index + 1;
}
function projectionText(...values) {
for (const value of values) {
if (value && typeof value === "object") {
const nested = messageAuthorityTextValue(value.text ?? value.content ?? value.message ?? value.summary ?? value.preview ?? value.title);
if (nested) return nested;
continue;
}
const text = messageAuthorityTextValue(value);
if (text) return text;
}
return null;
}
function objectValue(value) {
return value && typeof value === "object" && !Array.isArray(value) ? value : null;
}
function textValue(value) {
return String(value ?? "").trim();
}
function messageAuthorityTextValue(value) {
const text = String(value ?? "").replace(/\r\n?/gu, "\n");
if (!text.trim() || text.trim() === "[object Object]") return "";
return text;
}
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"),
layer: textValue(value.layer ?? value.source) || null,
category: textValue(value.category ?? value.type) || null,
summary: String(value.summary ?? value.userMessage ?? value.message ?? value.code ?? "Projection is blocked."),
message: textValue(value.message ?? value.summary ?? value.userMessage) || null,
userMessage: textValue(value.userMessage ?? value.zh) || null,
retryable: typeof value.retryable === "boolean" ? value.retryable : null,
runId: textValue(value.runId) || null,
commandId: textValue(value.commandId) || null,
timeoutMs: Number.isFinite(Number(value.timeoutMs)) ? Math.trunc(Number(value.timeoutMs)) : null,
valuesRedacted: true
};
}