import type { BaseEvent } from "../types.js"; import type { AgentUtility } from "./attribution.js"; import { computeRoutingRegret, type RoutingRegretMetrics } from "./routing-regret.js"; export interface WeeklyMetrics { tasks: number; terminalTasks: number; successfulTasks: number; verifiedSuccessfulTasks: number; taskCompletionRate: number; successfulCompletionRate: number; verifiedSuccessRate: number; humanReworkRate: number; defectEscapeRate: number | null; defectEscape: { eligibleTasks: number; escapedTasks: number }; totalCredits: number; totalUsdReported: number; totalTokens: number; inputTokens: number; cachedInputTokens: number; outputTokens: number; cacheEfficiency: number | null; modelCalls: number; passedVerifications: number; feedback: { tasks: number; good: number; fixed: number; bad: number; implicit: number; reworkTasks: number; repairTasks: number; coverageRate: number }; profileDistribution: Record; topologyDistribution: Record; modelDistribution: Record; providerDistribution: Record; costToGreen?: { creditSamples: number; usdSamples: number; timeSamples: number; averageCredits?: number; averageUsdReported?: number; averageTimeToFirstGreenMs?: number; averageTimeToFinalGreenMs?: number }; contextEfficiency: { workerInputTokens: number; uniqueEstimatedTaskContextTokens: number | null; contextDuplicationCoefficient: number | null; usedEnvelopeFacts: number | null; transmittedEnvelopeFacts: number | null; envelopeEfficiency: number | null }; coordination: { handoffs: number; blackboardEvents: number; coordinationTokens: number; coordinationTax: number | null; modelHopCount: number; escalations: number; creditsBeforeEscalation: number; warRoomConvergence: number | null }; review: { completed: number; approved: number; disagreements: number; confirmedFindings: number; failures: number }; attribution: { spawnedAgents: number; attributedAgents: number; attributionCoverageRate: number | null; agentUtilizationRate: number | null; zeroImpactAgentRate: number | null; uniqueContributionRate: number | null; duplicateFindingRate: number | null; creditsPerUsedFact: number | null; creditsPerAcceptedDecision: number | null; outputDensity: number | null }; routingRegret: RoutingRegretMetrics; dataQualityLimitations: string[]; } export interface WeeklyCohort { profile: string; topology: string; configVersions: string[]; policyVersions: string[]; providers: string[]; models: string[]; tasks: number; verifiedSuccessfulTasks: number; verifiedSuccessRate: number; humanReworkRate: number; totalCredits: number; totalTokens: number; } function number(event: BaseEvent, key: string): number | undefined { const value = event[key]; return typeof value === "number" && Number.isFinite(value) && value >= 0 ? value : undefined; } function ratio(numerator: number, denominator: number): number { return denominator ? numerator / denominator : 0; } function nullableRatio(numerator: number, denominator: number): number | null { return denominator ? numerator / denominator : null; } function tokens(event: BaseEvent): number { return (number(event, "inputTokens") ?? 0) + (number(event, "cachedInputTokens") ?? 0) + (number(event, "outputTokens") ?? 0); } function time(event: BaseEvent): number | undefined { const value = Date.parse(event.timestamp); return Number.isFinite(value) ? value : undefined; } function text(event: BaseEvent, key: string, fallback = "unknown"): string { return typeof event[key] === "string" && event[key] ? event[key] as string : fallback; } function sorted(events: readonly BaseEvent[]): BaseEvent[] { return events.map((event, index) => ({ event, index })).sort((left, right) => (time(left.event) ?? 0) - (time(right.event) ?? 0) || left.index - right.index).map(({ event }) => event); } function distribution(values: readonly string[]): Record { const counts: Record = {}; for (const value of values) counts[value] = (counts[value] ?? 0) + 1; return Object.fromEntries(Object.entries(counts).sort(([left], [right]) => left.localeCompare(right))); } function byTask(events: readonly BaseEvent[]): Map { const taskIds = new Set(events.filter((event) => event.eventType === "request.received" || event.eventType === "request.classified").map((event) => event.taskId)); const tasks = new Map(); for (const event of events) if (taskIds.has(event.taskId)) tasks.set(event.taskId, [...(tasks.get(event.taskId) ?? []), event]); for (const [taskId, taskEvents] of tasks) tasks.set(taskId, sorted(taskEvents)); return tasks; } function finalResult(events: readonly BaseEvent[]): BaseEvent | undefined { for (let index = events.length - 1; index >= 0; index -= 1) if (events[index]!.eventType === "result.completed") return events[index]; return undefined; } function taskTopology(events: readonly BaseEvent[]): string { for (let index = events.length - 1; index >= 0; index -= 1) { const event = events[index]!; if (event.eventType === "route.validated") return text(event, "finalTopology"); if (event.eventType === "request.classified") return text(event, "topology"); } return "unknown"; } export function aggregateWeeklyMetrics(events: readonly BaseEvent[]): WeeklyMetrics { const ordered = sorted(events); const tasks = byTask(ordered); const calls = ordered.filter((event) => event.eventType === "model.call.completed" || event.eventType === "model.call.failed"); const results = [...tasks.values()].map(finalResult).filter((event): event is BaseEvent => Boolean(event)); const successful = results.filter((event) => event.success === true); const verified = results.filter((event) => event.success === true && event.verified === true); const inputTokens = calls.reduce((sum, event) => sum + (number(event, "inputTokens") ?? 0), 0); const cachedInputTokens = calls.reduce((sum, event) => sum + (number(event, "cachedInputTokens") ?? 0), 0); const outputTokens = calls.reduce((sum, event) => sum + (number(event, "outputTokens") ?? 0), 0); const totalTokens = inputTokens + cachedInputTokens + outputTokens; const totalCredits = calls.reduce((sum, event) => sum + (number(event, "creditsEstimated") ?? 0), 0); const totalUsdReported = calls.reduce((sum, event) => sum + (number(event, "usdReported") ?? 0), 0); const feedbackEvents = ordered.filter((event) => event.eventType === "feedback.explicit" || event.eventType === "feedback.implicit" || event.eventType === "user.feedback"); const feedbackTasks = new Set(feedbackEvents.map((event) => event.taskId)); const fixedTasks = new Set(feedbackEvents.filter((event) => event.kind === "fixed").map((event) => event.taskId)); const implicitTasks = new Set(feedbackEvents.filter((event) => event.eventType === "feedback.implicit" || event.explicit === false).map((event) => event.taskId)); const reworkTasks = new Set([...fixedTasks, ...implicitTasks]); const repairTasks = new Set(ordered.filter((event) => event.eventType === "repair.started" || event.eventType === "repair.attempted").map((event) => event.taskId)); const defectEscapeEligibleTasks = [...tasks.values()].filter((taskEvents) => { const result = finalResult(taskEvents); return result?.success === true && result.verified === true; }); const escapedDefectTasks = defectEscapeEligibleTasks.filter((taskEvents) => { const result = finalResult(taskEvents)!; return taskEvents.slice(taskEvents.indexOf(result) + 1).some((event) => (event.eventType === "feedback.explicit" || event.eventType === "feedback.implicit" || event.eventType === "user.feedback") && (event.kind === "fixed" || event.eventType === "feedback.implicit" || event.explicit === false)); }).length; const workerCalls = calls.filter((event) => ["scout", "writer", "reviewer", "deep", "reducer"].includes(text(event, "role", ""))); const workerInputTokens = workerCalls.reduce((sum, event) => sum + (number(event, "inputTokens") ?? 0) + (number(event, "cachedInputTokens") ?? 0), 0); const workerNodes = new Set(); let completeContextCoverage = workerCalls.length > 0; for (const call of workerCalls) { const nodeId = typeof call.nodeId === "string" && call.nodeId ? call.nodeId : undefined; if (nodeId) workerNodes.add(`${call.taskId}:${nodeId}`); else completeContextCoverage = false; } const measuredNodes = new Set(); const uniqueContexts = new Map(); for (const measurement of ordered.filter((event) => event.eventType === "task.context.measured")) { const nodeId = typeof measurement.nodeId === "string" && measurement.nodeId ? measurement.nodeId : undefined; if (!nodeId) continue; const nodeKey = `${measurement.taskId}:${nodeId}`; if (!workerNodes.has(nodeKey)) continue; const hash = typeof measurement.contextHash === "string" && measurement.contextHash ? measurement.contextHash : undefined; const estimatedTokens = number(measurement, "contextEstimatedTokens"); if (!hash || estimatedTokens === undefined) { completeContextCoverage = false; continue; } measuredNodes.add(nodeKey); const contextKey = `${measurement.taskId}:${hash}`; const previous = uniqueContexts.get(contextKey); if (previous !== undefined && previous !== estimatedTokens) completeContextCoverage = false; else uniqueContexts.set(contextKey, estimatedTokens); } for (const key of workerNodes) if (!measuredNodes.has(key)) completeContextCoverage = false; const uniqueEstimatedTaskContextTokens = completeContextCoverage ? [...uniqueContexts.values()].reduce((sum, value) => sum + value, 0) : null; const contextDuplicationCoefficient = uniqueEstimatedTaskContextTokens === null ? null : nullableRatio(workerInputTokens, uniqueEstimatedTaskContextTokens); const envelopeTasks = [...tasks.values()].filter((taskEvents) => finalResult(taskEvents) && taskEvents.some((event) => event.eventType === "handoff.sent" && event.kind === "task-envelope")); let completeEnvelopeCoverage = envelopeTasks.length > 0; let measuredUsedEnvelopeFacts = 0; let measuredTransmittedEnvelopeFacts = 0; for (const taskEvents of envelopeTasks) { const measurement = taskEvents.findLast((event) => event.eventType === "envelope.measured"); const used = measurement && number(measurement, "usedEnvelopeFacts"); const transmitted = measurement && number(measurement, "transmittedEnvelopeFacts"); if (measurement?.complete !== true || used === undefined || transmitted === undefined || used > transmitted) { completeEnvelopeCoverage = false; continue; } measuredUsedEnvelopeFacts += used; measuredTransmittedEnvelopeFacts += transmitted; } const usedEnvelopeFacts = completeEnvelopeCoverage ? measuredUsedEnvelopeFacts : null; const transmittedEnvelopeFacts = completeEnvelopeCoverage ? measuredTransmittedEnvelopeFacts : null; const envelopeEfficiency = usedEnvelopeFacts === null || transmittedEnvelopeFacts === null ? null : nullableRatio(usedEnvelopeFacts, transmittedEnvelopeFacts); const warRoomTasks = [...tasks.values()].filter((taskEvents) => taskTopology(taskEvents) === "warroom"); let completeWarRoomCoverage = warRoomTasks.length > 0; let openHypothesesReduction = 0; let communicationRounds = 0; for (const taskEvents of warRoomTasks) { const measurement = taskEvents.findLast((event) => event.eventType === "warroom.converged" || (event.eventType === "escalation.triggered" && event.reason === "warroom-no-convergence")); const reduction = measurement?.openHypothesesReduction; const round = measurement && number(measurement, "round"); if (typeof reduction !== "number" || !Number.isFinite(reduction) || round === undefined || !Number.isInteger(round) || round < 1) { completeWarRoomCoverage = false; continue; } openHypothesesReduction += reduction; communicationRounds += round; } const warRoomConvergence = completeWarRoomCoverage ? nullableRatio(openHypothesesReduction, communicationRounds) : null; let creditSamples = 0; let usdSamples = 0; let timeSamples = 0; let greenCredits = 0; let greenUsd = 0; let firstGreenMs = 0; let finalGreenMs = 0; for (const [taskId, taskEvents] of tasks) { const result = finalResult(taskEvents); if (result?.success !== true || result.verified !== true) continue; const taskCalls = taskEvents.filter((event) => event.eventType === "model.call.completed" || event.eventType === "model.call.failed"); if (taskCalls.length && taskCalls.every((event) => number(event, "creditsEstimated") !== undefined)) { creditSamples += 1; greenCredits += taskCalls.reduce((sum, event) => sum + number(event, "creditsEstimated")!, 0); } if (taskCalls.length && taskCalls.every((event) => number(event, "usdReported") !== undefined)) { usdSamples += 1; greenUsd += taskCalls.reduce((sum, event) => sum + number(event, "usdReported")!, 0); } const startedAt = taskEvents.find((event) => event.eventType === "request.received" || event.eventType === "request.classified"); const firstGreen = taskEvents.find((event) => event.eventType === "verification.completed" && event.passed === true); const start = startedAt && time(startedAt); const first = firstGreen && time(firstGreen); const last = time(result); if (start !== undefined && first !== undefined && last !== undefined && first >= start && last >= first) { timeSamples += 1; firstGreenMs += first - start; finalGreenMs += last - start; } void taskId; } let interactionTokens = 0; const handoffKeys = new Set(); for (const event of ordered.filter((entry) => entry.eventType === "handoff.sent")) { if (typeof event.contentHash === "string") handoffKeys.add(`${event.taskId}:${text(event, "fromNodeId", "")}:${text(event, "toNodeId", "")}:${event.contentHash}`); interactionTokens += number(event, "payloadEstimatedTokens") ?? Math.ceil((number(event, "payloadChars") ?? 0) / 4); } for (const event of ordered.filter((entry) => entry.eventType === "blackboard.event")) { const key = typeof event.contentHash === "string" ? `${event.taskId}:${text(event, "fromNodeId", "")}:${text(event, "toNodeId", "")}:${event.contentHash}` : undefined; if (!key || !handoffKeys.has(key)) interactionTokens += number(event, "payloadEstimatedTokens") ?? Math.ceil((number(event, "payloadChars") ?? 0) / 4); } const reducerTokens = calls.filter((event) => event.role === "reducer").reduce((sum, event) => sum + tokens(event), 0); const coordinationTokens = interactionTokens + reducerTokens; let modelHopCount = 0; for (const taskEvents of tasks.values()) { let current: string | undefined; for (const event of taskEvents) { if (event.eventType === "route.validated") current = text(event, "modelFinal", current ?? "") || current; else if (event.eventType === "root.selection.completed") current = text(event, "actualModel", current ?? "") || current; else if (event.eventType === "model.changed") { const previous = text(event, "previousModel", current ?? ""); const next = text(event, "model"); if (previous && next && previous !== next) modelHopCount += 1; if (next) current = next; } else if (event.eventType === "escalation.triggered") { const next = text(event, "model"); if (current && next && current !== next) modelHopCount += 1; if (next) current = next; } } } let creditsBeforeEscalation = 0; for (const taskEvents of tasks.values()) { const escalation = taskEvents.find((event) => event.eventType === "escalation.triggered"); const at = escalation && time(escalation); if (at === undefined) continue; creditsBeforeEscalation += taskEvents.filter((event) => (event.eventType === "model.call.completed" || event.eventType === "model.call.failed") && (time(event) ?? Infinity) < at).reduce((sum, event) => sum + (number(event, "creditsEstimated") ?? 0), 0); } const spawned = new Set(ordered.filter((event) => event.eventType === "agent.spawned" && text(event, "role", "") !== "reviewer").map((event) => `${event.taskId}:${text(event, "nodeId", event.eventId)}`)); const attributions = new Map(); for (const event of ordered.filter((entry) => entry.eventType === "attribution.created")) attributions.set(`${event.taskId}:${text(event, "nodeId", event.eventId)}`, event); const attributedSpawned = [...spawned].flatMap((key) => attributions.has(key) ? [attributions.get(key)!] : []); const factsTotal = attributedSpawned.reduce((sum, event) => sum + (number(event, "factsTotal") ?? 0), 0); const factsUsed = attributedSpawned.reduce((sum, event) => sum + (number(event, "factsUsed") ?? 0), 0); const uniqueFacts = attributedSpawned.reduce((sum, event) => sum + (number(event, "uniqueFactsUsed") ?? 0), 0); const duplicateFacts = attributedSpawned.reduce((sum, event) => sum + (number(event, "duplicateFacts") ?? 0), 0); const decisions = attributedSpawned.reduce((sum, event) => sum + (number(event, "decisionCount") ?? 0), 0); const usedAgents = attributedSpawned.filter((event) => (number(event, "factsUsed") ?? 0) > 0).length; const zeroImpactAgents = attributedSpawned.filter((event) => number(event, "factsUsed") === 0).length; const completeAttribution = spawned.size > 0 && attributedSpawned.length === spawned.size; const usedFactCredits = attributedSpawned.reduce((sum, event) => sum + (number(event, "creditsPerUsedFinding") ?? 0) * (number(event, "factsUsed") ?? 0), 0); const acceptedDecisionCredits = attributedSpawned.reduce((sum, event) => sum + (number(event, "creditsPerAcceptedDecision") ?? 0) * (number(event, "decisionCount") ?? 0), 0); const routingRegret = computeRoutingRegret(tasks.values()); const taskCount = tasks.size; const limitations = [ ...(calls.some((event) => number(event, "creditsEstimated") === undefined) ? ["Some model calls lack estimated credits."] : []), ...(calls.some((event) => number(event, "usdReported") === undefined) ? ["Pi-estimated USD is incomplete; credits remain estimates."] : []), ...(spawned.size > attributedSpawned.length ? ["Some spawned agents lack attribution events."] : []), ...(defectEscapeEligibleTasks.length ? ["Defect Escape Rate is an observed lower bound based on fixed or implicit rework feedback after verified success."] : ["Defect Escape Rate unavailable: no verified successful tasks in the period."]), ...(contextDuplicationCoefficient === null ? ["Context Duplication Coefficient unavailable: worker context measurements are absent, incomplete, conflicting, or have no positive unique-context denominator."] : []), ...(envelopeEfficiency === null ? ["Envelope Efficiency unavailable: eligible terminal task-envelope measurements are absent, incomplete, invalid, or have no transmitted-fact denominator."] : []), ...(warRoomConvergence === null ? ["War Room Convergence unavailable: final War Room measurements are absent, incomplete, or invalid."] : []), ...routingRegret.limitations, ]; const costToGreen = creditSamples || usdSamples || timeSamples ? { creditSamples, usdSamples, timeSamples, ...(creditSamples ? { averageCredits: greenCredits / creditSamples } : {}), ...(usdSamples ? { averageUsdReported: greenUsd / usdSamples } : {}), ...(timeSamples ? { averageTimeToFirstGreenMs: firstGreenMs / timeSamples, averageTimeToFinalGreenMs: finalGreenMs / timeSamples } : {}), } : undefined; return { tasks: taskCount, terminalTasks: results.length, successfulTasks: successful.length, verifiedSuccessfulTasks: verified.length, taskCompletionRate: ratio(results.length, taskCount), successfulCompletionRate: ratio(successful.length, taskCount), verifiedSuccessRate: ratio(verified.length, taskCount), humanReworkRate: ratio(reworkTasks.size, taskCount), defectEscapeRate: nullableRatio(escapedDefectTasks, defectEscapeEligibleTasks.length), defectEscape: { eligibleTasks: defectEscapeEligibleTasks.length, escapedTasks: escapedDefectTasks }, totalCredits, totalUsdReported, totalTokens, inputTokens, cachedInputTokens, outputTokens, cacheEfficiency: nullableRatio(cachedInputTokens, inputTokens + cachedInputTokens), modelCalls: calls.length, passedVerifications: ordered.filter((event) => event.eventType === "verification.completed" && event.passed === true).length, feedback: { tasks: feedbackTasks.size, good: feedbackEvents.filter((event) => event.kind === "good").length, fixed: feedbackEvents.filter((event) => event.kind === "fixed").length, bad: feedbackEvents.filter((event) => event.kind === "bad").length, implicit: feedbackEvents.filter((event) => event.eventType === "feedback.implicit" || event.explicit === false).length, reworkTasks: reworkTasks.size, repairTasks: repairTasks.size, coverageRate: ratio(feedbackTasks.size, taskCount), }, profileDistribution: distribution([...tasks.values()].map((taskEvents) => taskEvents[0]?.profile ?? "unknown")), topologyDistribution: distribution([...tasks.values()].map(taskTopology)), modelDistribution: distribution(calls.map((event) => text(event, "model"))), providerDistribution: distribution(calls.map((event) => text(event, "provider"))), ...(costToGreen ? { costToGreen } : {}), contextEfficiency: { workerInputTokens, uniqueEstimatedTaskContextTokens, contextDuplicationCoefficient, usedEnvelopeFacts, transmittedEnvelopeFacts, envelopeEfficiency, }, coordination: { handoffs: ordered.filter((event) => event.eventType === "handoff.sent").length, blackboardEvents: ordered.filter((event) => event.eventType === "blackboard.event").length, coordinationTokens, coordinationTax: nullableRatio(coordinationTokens, totalTokens), modelHopCount, escalations: ordered.filter((event) => event.eventType === "escalation.triggered").length, creditsBeforeEscalation, warRoomConvergence, }, review: { completed: ordered.filter((event) => event.eventType === "review.completed" && (number(event, "reviewers") ?? number(event, "required") ?? 0) > 0).length, approved: ordered.filter((event) => event.eventType === "review.completed" && (number(event, "reviewers") ?? number(event, "required") ?? 0) > 0 && event.approved === true).length, disagreements: ordered.filter((event) => event.eventType === "review.completed" && (number(event, "reviewers") ?? number(event, "required") ?? 0) > 0 && event.disagreement === true).length, confirmedFindings: ordered.filter((event) => event.eventType === "review.completed" && (number(event, "reviewers") ?? number(event, "required") ?? 0) > 0).reduce((sum, event) => sum + (number(event, "confirmedFindings") ?? 0), 0), failures: ordered.filter((event) => event.eventType === "review.failed").length, }, attribution: { spawnedAgents: spawned.size, attributedAgents: attributedSpawned.length, attributionCoverageRate: nullableRatio(attributedSpawned.length, spawned.size), agentUtilizationRate: completeAttribution ? usedAgents / spawned.size : null, zeroImpactAgentRate: completeAttribution ? zeroImpactAgents / spawned.size : null, uniqueContributionRate: completeAttribution ? nullableRatio(uniqueFacts, factsTotal) : null, duplicateFindingRate: completeAttribution ? nullableRatio(duplicateFacts, factsTotal) : null, creditsPerUsedFact: completeAttribution ? nullableRatio(usedFactCredits, factsUsed) : null, creditsPerAcceptedDecision: completeAttribution ? nullableRatio(acceptedDecisionCredits, decisions) : null, outputDensity: completeAttribution ? nullableRatio(factsUsed + decisions, outputTokens) : null, }, routingRegret: routingRegret.metrics, dataQualityLimitations: limitations, }; } export function buildWeeklyCohorts(events: readonly BaseEvent[]): WeeklyCohort[] { const grouped = new Map; events: BaseEvent[] }>(); for (const taskEvents of byTask(events).values()) { const calls = taskEvents.filter((event) => event.eventType === "model.call.completed" || event.eventType === "model.call.failed"); const dimensions = { profile: taskEvents[0]?.profile ?? "unknown", topology: taskTopology(taskEvents), configVersions: [...new Set(taskEvents.map((event) => event.configVersion))].sort(), policyVersions: [...new Set(taskEvents.map((event) => event.policyVersion))].sort(), providers: [...new Set(calls.map((event) => text(event, "provider")))].sort(), models: [...new Set(calls.map((event) => text(event, "model")))].sort(), }; if (!dimensions.providers.length) dimensions.providers.push("unknown"); if (!dimensions.models.length) dimensions.models.push("unknown"); const key = JSON.stringify(dimensions); const cohort = grouped.get(key) ?? { dimensions, events: [] }; cohort.events.push(...taskEvents); grouped.set(key, cohort); } return [...grouped.values()].sort((left, right) => JSON.stringify(left.dimensions).localeCompare(JSON.stringify(right.dimensions))).map(({ dimensions, events: cohortEvents }) => { const metrics = aggregateWeeklyMetrics(cohortEvents); return { ...dimensions, tasks: metrics.tasks, verifiedSuccessfulTasks: metrics.verifiedSuccessfulTasks, verifiedSuccessRate: metrics.verifiedSuccessRate, humanReworkRate: metrics.humanReworkRate, totalCredits: metrics.totalCredits, totalTokens: metrics.totalTokens }; }); } export function aggregateUtilities(utilities: AgentUtility[]) { if (!utilities.length) return { agentUtilizationRate: null, zeroImpactAgentRate: null, uniqueContributionRate: null, duplicateFindingRate: null }; const count = utilities.length; return { agentUtilizationRate: utilities.filter((utility) => !utility.zeroImpact).length / count, zeroImpactAgentRate: utilities.filter((utility) => utility.zeroImpact).length / count, uniqueContributionRate: utilities.reduce((sum, utility) => sum + utility.uniqueContributionRatio, 0) / count, duplicateFindingRate: utilities.reduce((sum, utility) => sum + utility.duplicateRatio, 0) / count }; }