Files
pikasTech-HWLAB/internal/cloud/workbench-turn-projection.ts
T

511 lines
22 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; draft-2026-06-20-p1-zero-split-durable-realtime; 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", "retrying", "pending", "queued", "accepted", "dispatching", "streaming", "active", "processing", "busy", "creating"]);
const CANCEL_FINAL_RESPONSE_TEXT = "hwlab-user-cancel";
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 traceTerminal = traceTerminalEvidence(trace);
const terminalEvidence = terminalTurnEvidence({ result, traceTerminal });
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(terminalEvidence?.finalResponse) : null;
const agentRun = objectValue(result?.agentRun ?? trace?.agentRun);
const lastEvent = traceLastEvent(trace);
const timing = createWorkbenchTurnTimingProjection({ result, session, trace, status, terminal });
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,
timing,
startedAt: timing.startedAt,
lastEventAt: timing.lastEventAt,
finishedAt: timing.finishedAt,
durationMs: timing.durationMs,
valuesRedacted: true
};
}
export function createWorkbenchTurnTimingProjection({ result = null, session = null, trace = null, status = null, terminal = null } = {}) {
const events = Array.isArray(trace?.events) ? trace.events : [];
const firstEvent = events[0] ?? null;
const lastEvent = traceLastEvent(trace);
const terminalEvent = terminalTraceEventForTiming(events);
// Session lifecycle timing is not turn timing authority. Mixing session
// startedAt/createdAt with turn terminal duration makes the running timer
// jump backwards when the turn completes.
const directTiming = objectValue(result?.timing) ?? objectValue(trace?.timing) ?? null;
const agentRun = objectValue(result?.agentRun ?? trace?.agentRun);
const traceSummary = objectValue(result?.traceSummary ?? trace?.traceSummary);
const sessionSnapshot = objectValue(session?.session);
const normalizedStatus = normalizeWorkbenchStatus(status ?? result?.status ?? trace?.status ?? session?.status);
const isTerminal = typeof terminal === "boolean" ? terminal : TERMINAL_STATUSES.has(normalizedStatus) && !RUNNING_STATUSES.has(normalizedStatus);
const startedAt = firstTimestamp(
result?.startedAt,
directTiming?.startedAt,
result?.traceStartedAt,
trace?.startedAt,
traceSummary?.startedAt,
firstEvent?.startedAt,
firstEvent?.createdAt,
firstEvent?.occurredAt
);
const lastEventAt = latestTimestamp(
result?.lastEventAt,
directTiming?.lastEventAt,
trace?.lastEventAt,
traceSummary?.lastEventAt,
lastEvent?.lastEventAt,
lastEvent?.updatedAt,
lastEvent?.createdAt,
lastEvent?.occurredAt,
trace?.updatedAt,
result?.updatedAt
);
const finishedAt = isTerminal ? firstTimestamp(
result?.finishedAt,
directTiming?.finishedAt,
trace?.finishedAt,
traceSummary?.finishedAt,
trace?.terminalEvidence?.updatedAt,
terminalEvent?.updatedAt,
terminalEvent?.createdAt,
terminalEvent?.occurredAt,
result?.updatedAt,
trace?.updatedAt,
lastEventAt
) : null;
const durationMs = durationValue(
result?.durationMs,
directTiming?.durationMs,
trace?.durationMs,
traceSummary?.durationMs,
trace?.elapsedMs,
result?.elapsedMs,
elapsedBetween(startedAt, finishedAt ?? lastEventAt)
);
return {
startedAt,
lastEventAt,
finishedAt,
durationMs,
valuesRedacted: true
};
}
export function projectionDiagnostics({ traceId = null, projection = null, result = null, trace = null, refreshError = null } = {}) {
const turn = projection ?? createWorkbenchTurnProjection({ traceId, result, trace });
const sealedCompleted = turn.terminal === true && turn.status === "completed" && Boolean(turn.finalResponse?.text);
const source = sealedCompleted ? null : projectionDiagnosticSource(trace);
const rawBlocker = sealedCompleted ? null : refreshError ?? source?.blocker ?? trace?.blocker ?? result?.blocker ?? result?.error ?? null;
const retryingProviderInterruption = !sealedCompleted && turn.running === true ? retryableProviderInterruptionEvidence(rawBlocker) : null;
const blocker = retryingProviderInterruption ? null : rawBlocker;
const hasProjectionInput = hasTraceProjection(trace) || Boolean(result || result?.agentRun);
const sourceStatus = sealedCompleted ? null : normalizeProjectionStatus(source?.projectionStatus ?? trace?.projectionStatus);
const effectiveSourceStatus = retryingProviderInterruption && sourceStatus === "blocked" ? "projecting" : sourceStatus;
const status = sealedCompleted ? "caught-up" : effectiveSourceStatus ?? (blocker ? "blocked" : turn.terminal ? "caught-up" : hasProjectionInput ? "projecting" : "unknown");
const diagnostic = blocker ? diagnosticBlocker(blocker) : null;
const sourceHealth = sealedCompleted ? null : normalizeProjectionHealth(source?.projectionHealth ?? trace?.projectionHealth);
const effectiveSourceHealth = retryingProviderInterruption && (sourceHealth === "degraded" || sourceHealth === "unavailable" || sourceHealth === "stalled") ? "projecting" : sourceHealth;
const projectionHealth = sealedCompleted ? "healthy" : effectiveSourceHealth
?? projectionHealthFor({ status, turn, hasProjectionInput, blocker: diagnostic });
const staleMs = sealedCompleted ? null : 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,
retryingProviderInterruption,
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";
if (status !== "completed" && retryableProviderInterruptionEvidence(direct, trace)) return null;
return { source: "trace-terminal-evidence", status, evidence: direct, 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, traceTerminal = null } = {}) {
const resultStatus = terminalStatusFromValue(
result?.terminalStatus
?? result?.agentRun?.terminalStatus
?? result?.agentRun?.commandState
?? result?.agentRun?.status
?? result?.agentRun?.runStatus
);
if (resultStatus) {
if (resultStatus !== "completed" && retryableProviderInterruptionEvidence(result, result?.agentRun, result?.providerTrace, traceTerminal?.evidence)) return null;
return { source: "result", status: resultStatus, finalResponse: terminalFinalResponse(resultStatus, result, traceTerminal), valuesRedacted: true };
}
const statusOnly = terminalStatusFromValue(result?.status);
if (statusOnly && resultHasTerminalAuthority(result, traceTerminal)) {
if (statusOnly !== "completed" && retryableProviderInterruptionEvidence(result, result?.agentRun, result?.providerTrace, traceTerminal?.evidence)) return null;
return { source: "result", status: statusOnly, finalResponse: terminalFinalResponse(statusOnly, result, traceTerminal), valuesRedacted: true };
}
if (traceTerminal && retryableProviderInterruptionEvidence(traceTerminal.evidence)) return null;
return traceTerminal;
}
function terminalFinalResponse(status, result = null, traceTerminal = null) {
const direct = traceTerminal?.finalResponse ?? result?.finalResponse;
if (direct) return direct;
if (normalizeWorkbenchStatus(status) !== "canceled") return null;
return {
text: CANCEL_FINAL_RESPONSE_TEXT,
status: "canceled",
traceId: textValue(result?.traceId ?? traceTerminal?.finalResponse?.traceId ?? traceTerminal?.evidence?.traceId) || null,
valuesPrinted: false
};
}
function resultHasTerminalAuthority(result = null, traceTerminal = null) {
if (!result || typeof result !== "object") return false;
if (traceTerminal) return true;
if (result.terminal === true || result.sealed === true) return true;
if (result.error || result.blocker) return true;
if (firstTimestamp(result.finishedAt, result.completedAt, result.endedAt)) return true;
// A completed-looking status plus final text is not enough authority to seal
// a Workbench turn. For AgentRun-backed turns the final text can arrive before
// the runner terminal event; sealing here makes the visible completed card
// keep changing once later trace snapshots arrive.
return false;
}
function activeTurnEvidence({ result = null, session = null, trace = null } = {}) {
const retryEvidence = retryableProviderInterruptionEvidence(result, result?.agentRun, result?.providerTrace, trace);
if (retryEvidence) return { source: "provider-retry", status: "retrying", evidence: retryEvidence, valuesRedacted: true };
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 = []) {
const finalResponse = terminalAssistantEventFinalResponse(events);
for (let index = events.length - 1; index >= 0; index -= 1) {
const event = events[index];
if (!event || typeof event !== "object") continue;
if (retryableProviderInterruptionEvidence(event)) continue;
const terminal = event.terminal === 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, index),
eventType: textValue(event.type ?? event.label) || null,
finalResponse,
evidence: finalResponse ? { textPresent: true, source: "trace-terminal-assistant-event", valuesRedacted: true } : null,
valuesRedacted: true
};
}
return null;
}
function terminalAssistantEventFinalResponse(events = []) {
for (let index = events.length - 1; index >= 0; index -= 1) {
const event = events[index];
if (!event || typeof event !== "object") continue;
if (!isAssistantTraceEvent(event)) continue;
if (!(event.terminal === true || event.final === true || event.replyAuthority === true)) continue;
const text = projectionText(event.finalResponse, event.text, event.content, event.message, event.summary, event.payload?.text, event.payload?.content, event.payload?.message);
if (!text) continue;
const status = terminalStatusFromValue(event.status ?? event.terminalStatus ?? event.payload?.terminalStatus) ?? "completed";
return {
text,
status,
traceId: textValue(event.traceId) || null,
seq: eventSeq(event, index),
eventType: textValue(event.type ?? event.label) || null,
valuesPrinted: false
};
}
return null;
}
function isAssistantTraceEvent(event = {}) {
const type = String(event.type ?? event.eventType ?? "").trim().toLowerCase();
if (type === "assistant" || type === "assistant_message") return true;
return /assistant:message|assistant_message/u.test(String(event.label ?? "").toLowerCase());
}
function activeTraceEventStatus(events = []) {
for (const event of [...events].reverse()) {
if (retryableProviderInterruptionEvidence(event)) return "retrying";
const status = normalizeWorkbenchStatus(event?.status ?? event?.type);
if (RUNNING_STATUSES.has(status)) return normalizeActiveStatus(status);
}
return null;
}
function retryableProviderInterruptionEvidence(...values) {
for (const value of values) {
const direct = retryableProviderInterruptionRecord(value);
if (direct) return direct;
}
return null;
}
const RETRYABLE_AGENTRUN_TRANSPORT_FAILURE_KINDS = new Set([
"agentrun-connect-failed",
"agentrun-proxy-exec-failed",
"agentrun-manager-fetch-failed",
"agentrun-timeout"
]);
function isRetryableAgentRunTransportFailureKind(kind) {
return RETRYABLE_AGENTRUN_TRANSPORT_FAILURE_KINDS.has(normalizedProjectionFailureKind(kind));
}
function retryableProviderInterruptionRecord(value) {
const record = objectValue(value);
if (!record) return null;
const kind = normalizedProjectionFailureKind(record.failureKind ?? record.errorCode ?? record.code ?? record.name);
if ((kind === "provider-stream-disconnected" || kind === "provider-unavailable" || isRetryableAgentRunTransportFailureKind(kind)) && record.willRetry === true) {
return {
failureKind: kind,
willRetry: true,
retryAttempt: numberOrNull(record.retryAttempt),
retryMax: numberOrNull(record.retryMax),
retryDelayMs: numberOrNull(record.retryDelayMs),
nextRetryAt: textValue(record.nextRetryAt) || null,
runId: textValue(record.runId) || null,
commandId: textValue(record.commandId) || null,
valuesRedacted: true
};
}
for (const nested of [record.payload, record.providerTrace, record.agentRun, record.error, record.blocker, record.terminalEvidence]) {
const found = retryableProviderInterruptionRecord(nested);
if (found) return found;
}
return null;
}
function normalizedProjectionFailureKind(value) {
const kind = normalizeWorkbenchStatus(value);
if (kind === "agentrun-unreachable") return "agentrun-connect-failed";
return kind;
}
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 terminalTraceEventForTiming(events = []) {
for (let index = events.length - 1; index >= 0; index -= 1) {
const event = events[index];
if (!event || typeof event !== "object") continue;
if (event.terminal === true || event.final === true || event.replyAuthority === true || terminalStatusFromValue(event.status ?? event.terminalStatus)) return event;
}
return null;
}
function firstTimestamp(...values) {
for (const value of values) {
const ms = Date.parse(String(value ?? ""));
if (Number.isFinite(ms)) return new Date(ms).toISOString();
}
return null;
}
function latestTimestamp(...values) {
let latest = null;
for (const value of values) {
const ms = Date.parse(String(value ?? ""));
if (!Number.isFinite(ms)) continue;
latest = latest === null ? ms : Math.max(latest, ms);
}
return latest === null ? null : new Date(latest).toISOString();
}
function durationValue(...values) {
let max = null;
for (const value of values) {
const number = Number(value);
if (!Number.isFinite(number) || number < 0) continue;
const duration = Math.trunc(number);
max = max === null ? duration : Math.max(max, duration);
}
return max;
}
function elapsedBetween(startedAt, endedAt) {
const start = Date.parse(String(startedAt ?? ""));
const end = Date.parse(String(endedAt ?? ""));
if (!Number.isFinite(start) || !Number.isFinite(end) || end < start) return null;
return end - start;
}
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 numberOrNull(value) {
const number = Number(value);
return Number.isFinite(number) ? number : null;
}
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
};
}