/* * SPEC: PJ2026-0104010803 Workbench唯一投影 draft-2026-06-20-p2-terminal-outbox-recovery; draft-2026-06-27-p0-read-model-timeline-contract; draft-2026-06-28-p0-d518-session-timeline-consistency; PJ2026-010403 API契约 draft-2026-06-20-p2-terminal-outbox-recovery; draft-2026-06-27-p0-workbench-read-api-contract; PJ2026-010401 Web工作台 draft-2026-06-20-p2-terminal-outbox-recovery; draft-2026-06-27-p0-workbench-read-model-contract; PJ2026-01060505 Workbench Performance draft-2026-06-17-p0 * 职责: Workbench durable fact normalization and DTO projection helpers; no HTTP or realtime transport ownership. */ import { createHash } from "node:crypto"; import { defaultCodeAgentTraceStore } from "./code-agent-trace-store.ts"; import { parsePositiveInteger, safeConversationId, safeOpaqueId, safeSessionId, safeTraceId, sendJson } from "./server-http-utils.ts"; import { createWorkbenchReadModel } from "./workbench-read-model.ts"; import { createWorkbenchRuntimeClient } from "./workbench-runtime-client.ts"; import { buildWorkbenchSessionDetail, compactLaunchContext, includeMessagesForSessionDetail } from "./workbench-session-detail-response.ts"; import { durableTraceStatus, RUNNING_STATUSES, terminalFinalResponse, TERMINAL_STATUSES } from "./workbench-turn-projection.ts"; import { emitCodeAgentOtelSpan, emitHttpServerRequestSpan } from "./otel-trace.ts"; 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 workbenchError(code, message, extra = {}) { return { ok: false, status: "failed", error: { code, message, ...extra }, valuesRedacted: true, secretMaterialStored: false }; } function normalizeRole(role) { const value = textValue(role).toLowerCase(); return value === "agent" ? "assistant" : value; } function partFact(part, index, messageId, traceId) { const type = textValue(part?.type) || "text"; const text = messageAuthorityTextValue(part?.text ?? part?.content ?? part?.message); return { partId: safePartId(part?.partId ?? part?.id) || `prt_${hash(`${messageId}:${index}:${type}`).slice(0, 24)}`, messageId, traceId, type, text: text || null, status: normalizeStatus(part?.status ?? "completed"), toolName: textValue(part?.toolName ?? part?.name) || null, createdAt: textValue(part?.createdAt ?? part?.timestamp) || null, valuesRedacted: part?.valuesRedacted !== false }; } export function factSessionSummary(session, facts = {}) { const sessionId = factSessionId(session); if (!sessionId) return null; const currentTrace = factCurrentTraceForSession(facts, sessionId, session); const traceId = currentTrace?.traceId ?? factLastTraceId(session) ?? factLatestTraceIdForSession(facts, sessionId); const projection = traceId ? factProjectionForTrace(facts, traceId) : null; const turn = traceId ? factTurnForTrace(facts, traceId) : null; const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null; const trace = traceId ? factTraceSnapshot(facts, traceId) : null; const messages = factMessagesForSession(session, facts); const terminalMessage = traceId ? factTerminalMessageForTrace(messages, traceId) : null; const traceStatus = normalizeTerminalStatus(trace?.status); const checkpointStatus = normalizeTerminalStatus(checkpoint?.status); const messageTerminalStatus = factTerminalStatusFromMessageProjection(terminalMessage); const turnStatus = normalizeStatus(turn?.status); const launchContext = compactLaunchContext(session?.sessionJson?.launchContext); const status = normalizeStatus(checkpointStatus ?? traceStatus ?? messageTerminalStatus ?? nonUnknownStatus(turnStatus) ?? nonUnknownStatus(currentTrace?.status) ?? session?.status); const timing = traceId ? factCombinedTimingProjection(status, checkpoint, messageTerminalStatus ? terminalMessage : null, turn, trace) : factTimingProjection(session, status); const title = sessionTitleFromMessages(messages); const preview = sessionPreviewFromMessages(messages); return { sessionId, threadId: safeOpaqueId(session?.threadId) ?? (textValue(session?.threadId) || null), agentId: session?.agentId ?? "hwlab-code-agent", title, preview, titleSource: title ? "message-projection" : null, previewSource: preview ? "message-projection" : null, 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, launchContext, messageCount: messages.length, firstUserMessagePreview: firstUserPreview(messages), updatedAt: factUpdatedAt(session), turnSummary: turn ? { turnId: factTurnId(turn, traceId), traceId, status, running: RUNNING_STATUSES.has(status), terminal: TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status), eventCount: projection?.lastProjectedSeq ?? trace?.eventCount, projection, projectionStatus: projection?.projectionStatus ?? null, projectionHealth: projection?.projectionHealth ?? null, staleMs: projection?.staleMs ?? null, blocker: projection?.blocker ?? null, timing, startedAt: timing.startedAt, lastEventAt: timing.lastEventAt, finishedAt: timing.finishedAt, durationMs: timing.durationMs, updatedAt: factUpdatedAt(turn) } : null, valuesRedacted: true }; } export function factSessionDetail(session, facts = {}, options = {}) { const sessionId = factSessionId(session); return buildWorkbenchSessionDetail({ summary: factSessionSummary(session, facts), session, sessionId, includeMessages: options.includeMessages === true, messages: factMessagesForSession(session, facts) }); } export 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); } const messageGroups = canonicalFactMessageGroupsForSession( factArray(facts.messages) .filter((message) => message.sessionId === sessionId) .sort((left, right) => compareFactMessagesForSessionAsc(left, right, facts)), partsByMessageId, facts ); return messageGroups .sort((left, right) => compareFactMessagesForSessionAsc(left.message, right.message, facts)) .map((group) => factMessageDto(group.message, group.parts, facts)) .filter(Boolean); } export function canonicalFactMessageGroupsForSession(messages = [], partsByMessageId = new Map(), facts = {}) { const groups = new Map(); for (const message of messages) { const key = factMessageCanonicalGroupKey(message); const existing = groups.get(key) ?? []; existing.push(message); groups.set(key, existing); } return [...groups.values()].map((groupMessages) => { const message = selectCanonicalFactMessage(groupMessages, partsByMessageId, facts); const aliasIds = uniqueText(groupMessages.map((item) => item?.messageId)); const parts = aliasIds.flatMap((messageId) => partsByMessageId.get(messageId) ?? []); return { message, parts }; }); } export function factMessageCanonicalGroupKey(message) { const messageId = textValue(message?.messageId) || textValue(message?.id); const role = textValue(message?.role) || "agent"; const traceId = safeTraceId(message?.traceId); const sessionId = factSessionId(message) ?? textValue(message?.sessionId) ?? "session"; if (traceId && isAssistantLikeRole(role)) return `assistant:${sessionId}:${traceId}`; return `message:${messageId || sessionId}:${role}:${traceId || "none"}`; } export function selectCanonicalFactMessage(messages = [], partsByMessageId = new Map(), facts = {}) { if (messages.length <= 1) return messages[0] ?? null; const traceId = safeTraceId(messages[0]?.traceId); const preferredMessageId = traceId ? workbenchLifecycleMessageId(traceId, "agent") : null; const turnMessageId = traceId ? textValue(factTurnForTrace(facts, traceId)?.messageId) || null : null; return [...messages].sort((left, right) => { const score = factMessageCanonicalScore(right, partsByMessageId, preferredMessageId, turnMessageId) - factMessageCanonicalScore(left, partsByMessageId, preferredMessageId, turnMessageId); return score || compareFactRecordsDesc(left, right); })[0] ?? messages[0] ?? null; } export function factMessageCanonicalScore(message, partsByMessageId = new Map(), preferredMessageId = null, turnMessageId = null) { const messageId = textValue(message?.messageId) || textValue(message?.id); let score = 0; if (preferredMessageId && messageId === preferredMessageId) score += 10000; if (turnMessageId && messageId === turnMessageId) score += 1000; if (factMessageHasFinalResponsePart(message, partsByMessageId)) score += 100; if (message?.sealed === true || message?.terminal === true) score += 50; if (isTerminalProjectionStatus(normalizeStatus(message?.status))) score += 25; score += Math.min(10, factSeq(message) ?? 0); return score; } export function factMessageHasFinalResponsePart(message, partsByMessageId = new Map()) { const messageId = textValue(message?.messageId) || textValue(message?.id); if (!messageId) return false; return (partsByMessageId.get(messageId) ?? []).some((part) => (textValue(part?.partType ?? part?.type) === "final_response") && Boolean(projectionText(part?.text, part?.content, part?.message))); } export function workbenchLifecycleMessageId(traceId, role = "agent") { const traceSuffix = (safeTraceId(traceId) || String(traceId || "trace")) .replace(/^trc_/u, "") .replace(/[^A-Za-z0-9_.:-]/gu, "_") .slice(0, 48) || "trace"; return `msg_${traceSuffix}_${role === "user" ? "user" : "agent"}`; } export function factMessageDto(message, parts = [], facts = {}) { const messageId = safeMessageId(message?.messageId) || textValue(message?.messageId); if (!messageId) return null; const traceId = safeTraceId(message?.traceId) ?? null; const role = textValue(message?.role) || "agent"; const turn = traceId ? factTurnForTrace(facts, traceId) : null; const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null; const messageStatus = normalizeStatus(message?.status); const assistantLike = isAssistantLikeRole(role); const legacyText = projectionText(message?.text, message?.content, message?.message, message?.finalResponse); const baseParts = parts.length > 0 ? [...parts].sort(compareFactPartsAsc).map((part) => factPartDto(part, messageId, traceId)).filter(Boolean) : !assistantLike && legacyText ? [partFact({ type: "text", text: legacyText, status: message?.status }, 0, messageId, traceId)] : []; const checkpointStatus = assistantLike ? normalizeTerminalStatus(checkpoint?.status) : null; const preliminaryStatus = normalizeStatus(checkpointStatus ?? (assistantLike ? turn?.status : null) ?? messageStatus); const syntheticTerminalFinalResponse = assistantLike && !factPartsHaveFinalResponse(baseParts) ? factSyntheticTerminalFinalResponse(preliminaryStatus, traceId, message, checkpoint, turn) : null; const normalizedParts = syntheticTerminalFinalResponse ? [...baseParts, partFact({ type: "final_response", text: syntheticTerminalFinalResponse.text, status: syntheticTerminalFinalResponse.status }, baseParts.length, messageId, traceId)] : baseParts; const messageTerminalStatus = assistantLike ? factTerminalStatusFromMessageProjection(message, normalizedParts) : null; const status = normalizeStatus(checkpointStatus ?? messageTerminalStatus ?? (assistantLike ? turn?.status : null) ?? messageStatus); const timing = assistantLike ? factCombinedTimingProjection(status, checkpoint, messageTerminalStatus ? { ...message, parts: normalizedParts } : null, turn, message) : factTimingProjection(message, status); const text = factMessageAuthorityText({ role, status, parts: normalizedParts, legacyText }); return { messageId, role, sessionId: (safeSessionId(message?.sessionId) ?? textValue(message?.sessionId)) || null, traceId, turnId: safeTurnId(message?.turnId) || traceId, status, parts: normalizedParts, text: text || "", textPreview: text ? text.slice(0, 240) : null, createdAt: textValue(message?.createdAt) || null, updatedAt: factUpdatedAt(message), timing, startedAt: timing.startedAt, lastEventAt: timing.lastEventAt, finishedAt: timing.finishedAt, durationMs: timing.durationMs, projectionStatus: null, projectionHealth: null, valuesRedacted: message?.valuesRedacted !== false }; } export function factMessageAuthorityText({ role, status, parts = [], legacyText = null } = {}) { if (!isAssistantLikeRole(role)) return firstFactPartText(parts) ?? legacyText ?? ""; if (!isTerminalProjectionStatus(status)) return ""; return firstFactPartText(parts, "final_response") ?? ""; } export function factPartsHaveFinalResponse(parts = []) { return parts.some((part) => part?.type === "final_response" && Boolean(projectionText(part?.text))); } export function factSyntheticTerminalFinalResponse(status, traceId = null, ...records) { const terminalStatus = normalizeTerminalStatus(status); if (!terminalStatus || terminalStatus === "completed") return null; return terminalFinalResponse(terminalStatus, { traceId, status: terminalStatus, records }, { evidence: records }); } export function isTerminalProjectionStatus(status) { return TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status); } export function firstFactPartText(parts = [], type = null) { for (const part of parts) { if (type && part?.type !== type) continue; const text = projectionText(part?.text); if (text) return text; } return null; } export function factMessageFinalResponseText(message) { if (!message || !isAssistantLikeRole(message.role) || !isTerminalProjectionStatus(message.status)) return null; return firstFactPartText(message.parts, "final_response"); } export function factTerminalMessageForTrace(messages = [], traceId = null) { const safeTrace = safeTraceId(traceId); return [...factArray(messages)].reverse().find((message) => { if (!message || !isAssistantLikeRole(message.role)) return false; if (safeTrace && message.traceId !== safeTrace) return false; return Boolean(factTerminalStatusFromMessageProjection(message)); }) ?? null; } export function factTerminalStatusFromMessageProjection(message = null, parts = null) { if (!message || !isAssistantLikeRole(message.role)) return null; const sourceParts = Array.isArray(parts) ? parts : factArray(message.parts); if (!firstFactPartText(sourceParts, "final_response")) return null; const messageStatus = normalizeTerminalStatus(message.status); if (messageStatus) return messageStatus; for (const part of sourceParts) { if (textValue(part?.partType ?? part?.type) !== "final_response") continue; if (!projectionText(part?.text, part?.content, part?.message)) continue; const partStatus = normalizeTerminalStatus(part?.status); if (partStatus) return partStatus; if (part?.terminal === true || part?.sealed === true) return "completed"; } return message.terminal === true || message.sealed === true ? "completed" : null; } export 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 }; } export 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; } export function factTurnSnapshot({ turn = null, session = null, facts = {}, traceId, turnId = null } = {}) { const safeTrace = safeTraceId(traceId); const resolvedTurnId = factTurnId(turn, safeTrace) ?? safeTurnId(turnId) ?? safeTrace; 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 terminalMessage = safeTrace ? factTerminalMessageForTrace(messages, safeTrace) : null; const checkpoint = safeTrace ? factCheckpointForTrace(facts, safeTrace) : null; const trace = factTraceSnapshot(facts, safeTrace); const checkpointStatus = normalizeTerminalStatus(checkpoint?.status); const traceStatus = normalizeTerminalStatus(trace.status); const messageTerminalStatus = factTerminalStatusFromMessageProjection(terminalMessage); const status = normalizeStatus(checkpointStatus ?? traceStatus ?? messageTerminalStatus ?? turn?.status ?? session?.status); const terminal = isTerminalProjectionStatus(status); const assistantText = terminal ? factMessageFinalResponseText(assistantMessage) : null; const timing = factCombinedTimingProjection(status, checkpoint, messageTerminalStatus ? terminalMessage : null, turn, trace); return { turnId: resolvedTurnId, traceId: safeTrace, status, running: RUNNING_STATUSES.has(status), terminal, 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: assistantText ?? null, finalResponse: assistantText ? { text: assistantText, sealed: true, source: "message-part" } : null, timing, startedAt: timing.startedAt, lastEventAt: timing.lastEventAt, finishedAt: timing.finishedAt, durationMs: timing.durationMs, 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, timing: trace.timing, startedAt: trace.startedAt, lastEventAt: trace.lastEventAt, finishedAt: trace.finishedAt, durationMs: trace.durationMs, updatedAt: trace.updatedAt }, urls: { self: resolvedTurnId ? `/v1/workbench/turns/${encodeURIComponent(resolvedTurnId)}` : null, traceEvents: safeTrace ? `/v1/workbench/traces/${encodeURIComponent(safeTrace)}/events` : null } }; } export function factTraceSnapshot(facts = {}, traceId) { const safeTrace = safeTraceId(traceId); const events = factArray(facts.traceEvents) .filter((event) => !safeTrace || event.traceId === safeTrace) .sort(compareFactTraceEventsAsc) .map(factTraceEventDto) .filter(Boolean); const status = durableTraceStatus(events); const lastEvent = events.at(-1) ?? null; const checkpoint = factCheckpointForTrace(facts, safeTrace); const timing = factTraceTimingProjection(events, checkpoint, status); return { traceId: safeTrace, status, eventCount: events.length, events, lastEvent, timing, startedAt: timing.startedAt, lastEventAt: timing.lastEventAt, finishedAt: timing.finishedAt, durationMs: timing.durationMs, updatedAt: lastEvent?.updatedAt ?? lastEvent?.createdAt ?? checkpoint?.updatedAt ?? null, valuesRedacted: true, secretMaterialStored: false }; } export function factTraceTimingProjection(events = [], checkpoint = null, status = null) { const checkpointTiming = factTimingProjection(checkpoint, status); const normalizedStatus = normalizeStatus(status ?? checkpoint?.status); const terminal = TERMINAL_STATUSES.has(normalizedStatus) && !RUNNING_STATUSES.has(normalizedStatus); const observedAt = new Date().toISOString(); const eventStartTimes = factArray(events).flatMap((event) => [event?.createdAt, event?.occurredAt]); const eventActivityTimes = factArray(events).flatMap((event) => [event?.createdAt, event?.occurredAt, event?.updatedAt]); const terminalEventTimes = factArray(events) .filter(factTraceEventIsTerminalAuthority) .flatMap(factTraceEventTerminalTimes); const startedAt = firstTimestampIso(checkpointTiming.startedAt, ...eventStartTimes); const finishedAt = terminal ? latestTimestampIso(checkpointTiming.finishedAt, ...terminalEventTimes) : null; const lastEventAt = terminal && finishedAt ? finishedAt : latestTimestampIso(checkpointTiming.lastEventAt, ...eventActivityTimes); const durationMs = elapsedFactMs(startedAt, terminal ? finishedAt : observedAt); const lastEventAgeMs = terminal ? null : elapsedFactMs(lastEventAt, observedAt); return { ...checkpointTiming, startedAt, lastEventAt, finishedAt, durationMs, observedAt: terminal ? null : observedAt, lastEventAgeMs, valuesRedacted: true }; } export function factTraceEventIsTerminalAuthority(event = null) { const eventType = textValue(event?.eventType ?? event?.type); return event?.terminal === true || event?.sealed === true || eventType === "terminal"; } export function factTraceEventTerminalTimes(event = null) { const source = objectValue(event?.timing); return [source?.finishedAt, event?.finishedAt, source?.lastEventAt, event?.lastEventAt, event?.occurredAt, event?.createdAt]; } export function factTraceEventDto(event, index) { const seq = factProjectedSeq(event); if (!seq) return null; const sourceSeq = Number.isFinite(Number(event?.sourceSeq)) && Number(event.sourceSeq) >= 0 ? Math.trunc(Number(event.sourceSeq)) : null; return { ...event, id: textValue(event?.id) || textValue(event?.sourceEventId) || null, seq, projectedSeq: seq, sourceSeq, 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 }; } export 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); const timing = factTimingProjection(checkpoint, normalizeStatus(checkpoint.status)); 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, timing, startedAt: timing.startedAt, lastEventAt: timing.lastEventAt, finishedAt: timing.finishedAt, durationMs: timing.durationMs, updatedAt: factUpdatedAt(checkpoint), valuesRedacted: true }; } export function factTimingProjection(record = null, status = null) { const source = objectValue(record?.timing); const observedAt = new Date().toISOString(); const startedAt = timestampIso(source?.startedAt ?? record?.startedAt ?? record?.createdAt); const lastEventAt = timestampIso(source?.lastEventAt ?? record?.lastEventAt ?? record?.updatedAt ?? record?.createdAt); const normalizedStatus = normalizeStatus(status ?? record?.status); const terminal = record?.terminal === true || TERMINAL_STATUSES.has(normalizedStatus) && !RUNNING_STATUSES.has(normalizedStatus); const finishedAt = terminal ? timestampIso(source?.finishedAt ?? record?.finishedAt ?? record?.completedAt) : null; const durationMs = elapsedFactMs(startedAt, terminal ? finishedAt : observedAt); const lastEventAgeMs = terminal ? null : elapsedFactMs(lastEventAt, observedAt); return { startedAt, lastEventAt, finishedAt, durationMs, observedAt: terminal ? null : observedAt, lastEventAgeMs, valuesRedacted: source?.valuesRedacted !== false }; } export function factCombinedTimingProjection(status = null, ...records) { const normalizedStatus = normalizeStatus(status ?? records.find((record) => record?.status)?.status); const terminal = TERMINAL_STATUSES.has(normalizedStatus) && !RUNNING_STATUSES.has(normalizedStatus); const observedAt = new Date().toISOString(); const timings = records.map((record) => factTimingSource(record)).filter(Boolean); const startedAt = firstTimestampIso(...timings.map((timing) => timing.startedAt)); const finishedAt = terminal ? latestTimestampIso(...timings.map((timing) => timing.finishedAt)) : null; const lastEventAt = terminal && finishedAt ? finishedAt : latestTimestampIso(...timings.map((timing) => timing.lastEventAt)); const durationMs = elapsedFactMs(startedAt, terminal ? finishedAt : observedAt); const lastEventAgeMs = terminal ? null : elapsedFactMs(lastEventAt, observedAt); return { startedAt, lastEventAt, finishedAt, durationMs, observedAt: terminal ? null : observedAt, lastEventAgeMs, valuesRedacted: true }; } export function factTimingSource(record = null) { if (!record) return null; const source = objectValue(record?.timing); return { startedAt: timestampIso(source?.startedAt ?? record?.startedAt ?? record?.createdAt), lastEventAt: timestampIso(source?.lastEventAt ?? record?.lastEventAt ?? record?.updatedAt ?? record?.createdAt), finishedAt: timestampIso(source?.finishedAt ?? record?.finishedAt ?? record?.completedAt) }; } export function firstTimestampIso(...values) { for (const value of values) { const timestamp = timestampIso(value); if (timestamp) return timestamp; } return null; } export function latestTimestampIso(...values) { let latest = null; let latestMs = Number.NEGATIVE_INFINITY; for (const value of values) { const timestamp = timestampIso(value); if (!timestamp) continue; const ms = Date.parse(timestamp); if (!Number.isFinite(ms) || ms < latestMs) continue; latest = timestamp; latestMs = ms; } return latest; } export function elapsedFactMs(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 Math.trunc(end - start); } export function timestampIso(value) { const ms = Date.parse(String(value ?? "")); return Number.isFinite(ms) ? new Date(ms).toISOString() : null; } export 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; } export function workbenchProjectionStoreError(error) { const data = objectValue(error?.data) ?? {}; const retryAfterNumber = Number(data.retryAfterMs); const retryAfterMs = Number.isFinite(retryAfterNumber) && retryAfterNumber >= 0 ? Math.trunc(retryAfterNumber) : null; const retryAttempt = nonNegativeIntegerOrNull(data.retryAttempt ?? data.retryAttempts); const retryMax = nonNegativeIntegerOrNull(data.retryMax); const retryDelaysMs = Array.isArray(data.retryDelaysMs) ? data.retryDelaysMs.map(nonNegativeIntegerOrNull).filter((value) => value !== null) : []; return workbenchError("projection_store_unavailable", "Workbench durable projection store is unavailable.", { causeCode: error?.code ?? null, queryResult: textValue(data.queryResult) || null, blockedLayer: textValue(data.blockedLayer) || null, retryable: data.retryable !== false, transient: data.transient === true || data.retryable === true, retryAfterMs: retryAfterMs ?? null, retryAttempt, retryMax, retryAttempts: retryAttempt, retryExhausted: data.retryExhausted === true, retryDelaysMs: retryDelaysMs.length > 0 ? retryDelaysMs : null }); } export function nonNegativeIntegerOrNull(value) { const parsed = Number(value); return Number.isFinite(parsed) && parsed >= 0 ? Math.trunc(parsed) : null; } export function factSessionId(session) { return safeSessionId(session?.sessionId ?? session?.id) ?? null; } export function factLastTraceId(session) { return safeTraceId(session?.lastTraceId ?? session?.traceId) ?? null; } export function factLatestTraceIdForSession(facts = {}, sessionId) { const turn = factArray(facts.turns).filter((item) => item.sessionId === sessionId).sort(compareFactRecordsDesc)[0] ?? null; return safeTraceId(turn?.traceId) ?? null; } export function factCurrentTraceForSession(facts = {}, sessionId, session = null) { const candidates = []; const sessionTraceId = factLastTraceId(session); if (sessionTraceId) candidates.push(factTraceCandidate(session, sessionTraceId, normalizeStatus(session?.status), 0)); for (const message of factArray(facts.messages)) { if (message?.sessionId !== sessionId || !safeTraceId(message?.traceId)) continue; candidates.push(factMessageTraceCandidate(message)); } for (const turn of factArray(facts.turns)) { if (turn?.sessionId !== sessionId || !safeTraceId(turn?.traceId)) continue; candidates.push(factTraceCandidate(turn, safeTraceId(turn.traceId), normalizeStatus(turn?.status), 3)); } return candidates.sort(compareTraceCandidatesDesc)[0] ?? null; } export function factMessageTraceCandidate(message) { const role = textValue(message?.role); let status = normalizeStatus(message?.status); let priority = 1; if (isAssistantLikeRole(role)) { priority = 2; if (status === "unknown") status = "running"; } else if (safeTraceId(message?.traceId)) { status = "running"; } return factTraceCandidate(message, safeTraceId(message?.traceId), status, priority); } export function factTraceCandidate(record, traceId, status, priority) { return { traceId: safeTraceId(traceId) ?? null, status: normalizeStatus(status), seq: factSeq(record), updatedAt: factUpdatedAt(record), priority }; } export function compareTraceCandidatesDesc(left, right) { return compareNumberDesc(left?.seq, right?.seq) || compareTimestampDesc(left?.updatedAt, right?.updatedAt) || compareNumberDesc(left?.priority, right?.priority); } export function nonUnknownStatus(status) { return normalizeStatus(status) === "unknown" ? null : status; } export function factTurnId(turn, fallbackTraceId = null) { return safeTurnId(turn?.turnId) ?? safeTraceId(turn?.turnId) ?? safeTraceId(fallbackTraceId) ?? null; } export function factUpdatedAt(record) { return textValue(record?.updatedAt ?? record?.occurredAt ?? record?.createdAt) || null; } export 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; } export function factProjectedSeq(record) { const parsed = Number(record?.projectedSeq); return Number.isFinite(parsed) && parsed > 0 ? Math.trunc(parsed) : null; } export 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"; } export 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"; } export function compareFactSessionsDesc(left, right) { return compareTimestampDesc(factUpdatedAt(left), factUpdatedAt(right)) || compareText(factSessionId(left), factSessionId(right)); } export function compareFactMessagesAsc(left, right) { return compareNumberAsc(factSeq(left), factSeq(right)) || compareTimestampAsc(left?.createdAt ?? left?.updatedAt, right?.createdAt ?? right?.updatedAt) || compareText(left?.messageId, right?.messageId); } export function compareFactMessagesForSessionAsc(left, right, facts = {}) { const leftKey = factMessageTimelineKey(left, facts); const rightKey = factMessageTimelineKey(right, facts); return compareOptionalTimestampAsc(leftKey.timelineAnchorAt, rightKey.timelineAnchorAt) || compareOptionalNumberAsc(leftKey.turnSeq, rightKey.turnSeq) || compareTimestampAsc(leftKey.turnStartedAt, rightKey.turnStartedAt) || compareText(leftKey.traceId, rightKey.traceId) || compareOptionalNumberAsc(leftKey.roleRank, rightKey.roleRank) || compareTimestampAsc(leftKey.messageCreatedAt, rightKey.messageCreatedAt) || compareOptionalNumberAsc(leftKey.messageSeq, rightKey.messageSeq) || compareText(left?.messageId, right?.messageId); } export function factMessageTimelineKey(message, facts = {}) { const traceId = safeTraceId(message?.traceId) ?? null; const turn = traceId ? factTurnForTrace(facts, traceId) : null; const checkpoint = traceId ? factCheckpointForTrace(facts, traceId) : null; const anchorMessage = traceId ? factTimelineAnchorMessageForTrace(facts, traceId) : null; return { traceId: traceId ?? "", timelineAnchorAt: firstTimestampIso(anchorMessage?.createdAt, anchorMessage?.updatedAt), turnSeq: firstFiniteNumber(factSeq(turn), factSeq(checkpoint)), turnStartedAt: firstTimestampIso(turn?.startedAt, checkpoint?.startedAt, message?.createdAt, message?.updatedAt), roleRank: factMessageTimelineRoleRank(message?.role), messageCreatedAt: timestampIso(message?.createdAt ?? message?.updatedAt), messageSeq: factSeq(message) }; } export function factTimelineAnchorMessageForTrace(facts = {}, traceId) { const safeTrace = safeTraceId(traceId); if (!safeTrace) return null; const messages = factArray(facts.messages).filter((message) => message.traceId === safeTrace); const userMessages = messages.filter((message) => normalizeRole(message?.role) === "user"); const candidates = userMessages.length > 0 ? userMessages : messages; return [...candidates].sort((left, right) => compareOptionalTimestampAsc(left?.createdAt ?? left?.updatedAt, right?.createdAt ?? right?.updatedAt) || compareOptionalNumberAsc(factSeq(left), factSeq(right)) || compareText(left?.messageId, right?.messageId))[0] ?? null; } export function factMessageTimelineRoleRank(role) { const normalized = normalizeRole(role); if (normalized === "user") return 0; if (isAssistantLikeRole(normalized)) return 1; return 2; } export function firstFiniteNumber(...values) { for (const value of values) { const parsed = Number(value); if (Number.isFinite(parsed)) return Math.trunc(parsed); } return null; } export function compareOptionalNumberAsc(left, right) { const leftNumber = Number(left); const rightNumber = Number(right); const leftComparable = Number.isFinite(leftNumber) ? leftNumber : Number.MAX_SAFE_INTEGER; const rightComparable = Number.isFinite(rightNumber) ? rightNumber : Number.MAX_SAFE_INTEGER; return leftComparable - rightComparable; } export function compareFactPartsAsc(left, right) { return compareNumberAsc(left?.partIndex, right?.partIndex) || compareNumberAsc(factSeq(left), factSeq(right)) || compareText(left?.partId, right?.partId); } export function compareFactTraceEventsAsc(left, right) { return compareNumberAsc(factProjectedSeq(left), factProjectedSeq(right)); } export function compareFactRecordsDesc(left, right) { return compareNumberDesc(factSeq(left), factSeq(right)) || compareTimestampDesc(factUpdatedAt(left), factUpdatedAt(right)); } export function compareTimestampAsc(left, right) { return timestampMs(left) - timestampMs(right); } export function compareTimestampDesc(left, right) { return timestampMs(right) - timestampMs(left); } export function compareOptionalTimestampAsc(left, right) { const leftMs = optionalTimestampMs(left); const rightMs = optionalTimestampMs(right); return (leftMs ?? Number.MAX_SAFE_INTEGER) - (rightMs ?? Number.MAX_SAFE_INTEGER); } export function compareNumberAsc(left, right) { return numericValue(left, Number.MAX_SAFE_INTEGER) - numericValue(right, Number.MAX_SAFE_INTEGER); } export function compareNumberDesc(left, right) { return numericValue(right, -1) - numericValue(left, -1); } export function compareText(left, right) { return textValue(left).localeCompare(textValue(right)); } export function timestampMs(value) { const parsed = Date.parse(String(value ?? "")); return Number.isFinite(parsed) ? parsed : 0; } export function optionalTimestampMs(value) { const parsed = Date.parse(String(value ?? "")); return Number.isFinite(parsed) ? parsed : null; } export function numericValue(value, fallback) { const parsed = Number(value); return Number.isFinite(parsed) ? parsed : fallback; } export 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; } export function emptyFactSet() { return { sessions: [], messages: [], parts: [], turns: [], traceEvents: [], checkpoints: [] }; } export function factArray(value) { return Array.isArray(value) ? value : []; } export function uniqueText(values = []) { return [...new Set(values.map((value) => textValue(value)).filter(Boolean))]; } export function eventSeq(event, index) { const seq = Number(event?.seq); return Number.isFinite(seq) && seq > 0 ? Math.trunc(seq) : index + 1; } export function normalizeStatus(value) { const text = textValue(value).toLowerCase().replace(/_/gu, "-"); if (text === "cancelled") return "canceled"; return text || "unknown"; } export function normalizeTerminalStatus(value) { const status = normalizeStatus(value); return TERMINAL_STATUSES.has(status) && !RUNNING_STATUSES.has(status) ? status : null; } export function firstUserPreview(messages) { return messages.find((message) => message.role === "user")?.textPreview ?? null; } export function latestUserPreview(messages) { return [...messages].reverse().find((message) => message.role === "user")?.textPreview ?? null; } export function latestValidMessagePreview(messages) { return [...messages].reverse().find((message) => message.role === "user" || isAssistantLikeRole(message.role))?.textPreview ?? null; } export function latestAssistantPreview(messages) { return [...messages].reverse().find((message) => isAssistantLikeRole(message.role))?.textPreview ?? null; } export function sessionTitleFromMessages(messages = []) { return boundedPreviewText(latestUserPreview(messages) ?? firstUserPreview(messages) ?? latestAssistantPreview(messages)); } export function sessionPreviewFromMessages(messages = []) { return boundedPreviewText(latestValidMessagePreview(messages) ?? firstUserPreview(messages) ?? latestAssistantPreview(messages)); } export function boundedPreviewText(value, maxLength = 160) { const text = textValue(value).replace(/\s+/gu, " ").trim(); if (!text) return null; return text.length > maxLength ? `${text.slice(0, maxLength - 3)}...` : text; } export function isAssistantLikeRole(role) { return role === "assistant" || role === "agent"; } export function objectValue(value) { return value && typeof value === "object" && !Array.isArray(value) ? value : {}; } export function textValue(value) { return String(value ?? "").trim(); } export function messageAuthorityTextValue(value) { const text = String(value ?? "").replace(/\r\n?/gu, "\n"); if (!text.trim() || text.trim() === "[object Object]") return ""; return text; } export function safeTurnId(value) { const text = textValue(value); return /^turn_[A-Za-z0-9_.:-]+$/u.test(text) ? text : null; } export function safeMessageId(value) { const text = textValue(value); return /^msg_[A-Za-z0-9_.:-]+$/u.test(text) ? text : null; } export function safePartId(value) { const text = textValue(value); return /^prt_[A-Za-z0-9_.:-]+$/u.test(text) ? text : null; } export function hash(value) { return createHash("sha256").update(String(value)).digest("hex"); }