import { createHash } from "node:crypto"; import { workerLevelConfigSignature } from "./config-core.ts"; import { projectStorageKey } from "./project-key.ts"; import { LEVEL_ORDER } from "./types.ts"; import type { ClusterConfig, ClusterLevel, ClusterState, ClusterTaskInput, LearningConfig, LearningEvidence, TaskAttempt, TaskLevelSelection, TaskRuntime, } from "./types.ts"; export const LEARNING_RETENTION_MS = 30 * 24 * 60 * 60 * 1_000; export function learningOptions(config: ClusterConfig): LearningConfig { return config.learning; } export function learningRetentionMs(config: ClusterConfig): number { return learningOptions(config).historyRetentionDays * 24 * 60 * 60 * 1_000; } function levelIndex(level: ClusterLevel): number { return LEVEL_ORDER.indexOf(level); } function higherLevel(left: ClusterLevel, right: ClusterLevel): ClusterLevel { return levelIndex(left) >= levelIndex(right) ? left : right; } function highestLevel(levels: ClusterLevel[]): ClusterLevel | undefined { return levels.reduce((highest, level) => highest === undefined ? level : higherLevel(highest, level), undefined); } function isClusterLevel(value: unknown): value is ClusterLevel { return value === "low" || value === "medium" || value === "high"; } function isNonEmptyString(value: unknown): value is string { return typeof value === "string" && value.trim() !== ""; } export function isLearningEvidence(value: unknown): value is LearningEvidence { if (typeof value !== "object" || value === null || Array.isArray(value)) return false; const record = value as Record; return ( isNonEmptyString(record.projectKey) && isNonEmptyString(record.taskFingerprint) && isNonEmptyString(record.taskType) && isClusterLevel(record.fromLevel) && isClusterLevel(record.toLevel) && levelIndex(record.toLevel) === levelIndex(record.fromLevel) + 1 && isNonEmptyString(record.runId) && typeof record.completedAt === "number" && Number.isFinite(record.completedAt) && isNonEmptyString(record.configSignature) ); } function normalizeFingerprintText(value: string): string { return value.trim().replace(/\s+/g, " "); } export function normalizeTaskType(value: string): string { return normalizeFingerprintText(value).toLowerCase(); } export function taskFingerprint(task: Pick): string { const payload = { taskType: normalizeTaskType(task.taskType), title: normalizeFingerprintText(task.title), task: normalizeFingerprintText(task.task), acceptanceCriteria: task.acceptanceCriteria.map(normalizeFingerprintText), }; return createHash("sha256").update(JSON.stringify(payload)).digest("hex"); } export function unexpiredEvidence(evidence: LearningEvidence[], now: number, retentionMs = LEARNING_RETENTION_MS): LearningEvidence[] { const expiresAt = now - retentionMs; return evidence.filter((item) => item.completedAt >= expiresAt && item.completedAt <= now); } function normalWorkerCompletion(attempt: TaskAttempt): boolean { const worker = attempt.worker; return worker?.exitCode === 0 && !worker.errorMessage && worker.stopReason !== "error" && worker.stopReason !== "aborted"; } function qualityUpgrades(task: TaskRuntime, projectKey: string, configSignature: string, fallbackCompletedAt: number): LearningEvidence[] { if (task.status !== "completed" || task.attempts.length < 2) return []; if (task.attempts.some((attempt) => !normalWorkerCompletion(attempt) || !attempt.review || !["pass", "retry", "escalate"].includes(attempt.review.decision))) return []; const finalAttempt = task.attempts.at(-1); if (!finalAttempt || finalAttempt.review?.decision !== "pass") return []; const failedLevels = new Set( task.attempts .slice(0, -1) .filter((attempt) => (attempt.review?.decision === "retry" || attempt.review?.decision === "escalate") && levelIndex(attempt.level) < levelIndex(finalAttempt.level), ) .map((attempt) => attempt.level), ); return LEVEL_ORDER.flatMap((fromLevel, index) => { const toLevel = LEVEL_ORDER[index + 1]; if (!toLevel || !failedLevels.has(fromLevel) || levelIndex(toLevel) > levelIndex(finalAttempt.level)) return []; return [{ projectKey, taskFingerprint: taskFingerprint(task), taskType: normalizeTaskType(task.taskType), fromLevel, toLevel, runId: "", completedAt: task.finishedAt ?? fallbackCompletedAt, configSignature, }]; }); } export function collectQualityUpgrades(state: ClusterState, config: ClusterConfig, now = Date.now()): LearningEvidence[] { if (state.status === "cancelled" || !learningOptions(config).enabled) return []; const configSignature = workerLevelConfigSignature(config); const completedAt = state.finishedAt ?? now; return state.tasks.flatMap((task) => qualityUpgrades(task, projectStorageKey(task.cwd ?? state.cwd), configSignature, completedAt) .map((evidence) => ({ ...evidence, runId: state.runId })), ); } export function selectInitialLevel( task: ClusterTaskInput, evidence: LearningEvidence[], config: ClusterConfig, projectKey: string, now = Date.now(), ): TaskLevelSelection { const learning = learningOptions(config); if (!learning.enabled) { return { initialLevel: task.level, levelSelectionReason: `主 agent 请求 ${task.level};学习机制已关闭,保持请求等级。`, }; } const signature = workerLevelConfigSignature(config); const matchingEvidence = unexpiredEvidence(evidence, now, learningRetentionMs(config)).filter((item) => item.configSignature === signature); const fingerprintEvidence = matchingEvidence.filter((item) => item.projectKey === projectKey && item.taskFingerprint === taskFingerprint(task)); const fingerprintLevel = highestLevel(fingerprintEvidence.map((item) => item.toLevel)); const normalizedTaskType = normalizeTaskType(task.taskType); const taskTypeEvidence = matchingEvidence.filter((item) => normalizeTaskType(item.taskType) === normalizedTaskType); const taskTypeRunIdsByLevel = new Map>(); for (const item of taskTypeEvidence) { const runIds = taskTypeRunIdsByLevel.get(item.toLevel) ?? new Set(); runIds.add(`${item.projectKey}\u0000${item.runId}`); taskTypeRunIdsByLevel.set(item.toLevel, runIds); } const taskTypeLevel = highestLevel( LEVEL_ORDER.filter((level) => (taskTypeRunIdsByLevel.get(level)?.size ?? 0) >= learning.taskTypeUpgradeThreshold), ); const learnedLevel = highestLevel([fingerprintLevel, taskTypeLevel].filter((level): level is ClusterLevel => level !== undefined)); const sources: string[] = []; if (fingerprintLevel) sources.push(`任务指纹有 ${fingerprintEvidence.length} 条有效升级证据,最低起始等级为 ${fingerprintLevel}`); if (taskTypeLevel) sources.push(`任务类型 ${normalizedTaskType} 有 ${taskTypeRunIdsByLevel.get(taskTypeLevel)?.size ?? 0} 个不同运行支持最低起始等级 ${taskTypeLevel}`); if (!learnedLevel) { const taskTypeRunCount = Math.max(0, ...Array.from(taskTypeRunIdsByLevel.values(), (runIds) => runIds.size)); const taskTypeProgress = taskTypeRunCount > 0 ? `任务类型 ${normalizedTaskType} 最多仅有 ${taskTypeRunCount} 个不同运行支持同一目标等级,尚未达到 ${learning.taskTypeUpgradeThreshold} 个运行的门槛;` : ""; return { initialLevel: task.level, levelSelectionReason: `主 agent 请求 ${task.level};${taskTypeProgress}无匹配当前 worker 配置签名且未过期的可用学习最低等级。`, }; } const initialLevel = higherLevel(task.level, learnedLevel); if (initialLevel === task.level) { return { initialLevel, levelSelectionReason: `主 agent 请求 ${task.level};${sources.join(";")};学习机制不自动降级,保持 ${initialLevel}。`, }; } return { initialLevel, levelSelectionReason: `主 agent 请求 ${task.level};${sources.join(";")};将实际初始等级提升到 ${initialLevel}。`, }; }