/** * model-roundrobin.ts — 多虚拟模型轮询引擎。 * * 设计要点: * - 每个预设(含 config.json 作为 "default" 组)注册成独立的虚拟模型, * model id = 预设文件名。/model 里同时出现所有虚拟模型,选哪个用哪组候选。 * - 健康统计和 currentIndex 按预设组隔离。 * - 配置由可视化面板写入 config.json 和 presets/。面板保存后通过 * pi.events.emit("roundrobin:config-changed") 通知本模块 reloadRuntime()。 * - 失效候选(models.json 里被删/改名)容错跳过,不让一个坏引用炸掉整组; * 全部候选失效则该组 enabled=false(不注册虚拟模型),不影响其他组。 * - sticky 策略:成功停留;失败切下一个;当前候选过冷却期后下次请求回首选(index 0)。 */ import { appendFileSync, existsSync, statSync, rmSync, renameSync } from "node:fs"; import { join } from "node:path"; import { createAssistantMessageEventStream, streamSimple, type Api, type AssistantMessage, type AssistantMessageEventStream, type Context, type Model, type SimpleStreamOptions, } from "@earendil-works/pi-ai"; import { getAgentDir } from "@earendil-works/pi-coding-agent"; import type { ExtensionAPI, ExtensionContext, ExtensionUIContext } from "@earendil-works/pi-coding-agent"; import { loadAllGroupsConfig, DEFAULT_GROUP_NAME, type GroupConfig, type SpeedTestConfig, normalizeSpeedTest } from "./rr-config.js"; import { appendHealthEvent, getProviderHealthSnapshot as loadProviderHealthSnapshot, resetProviderHealth as resetProviderHealthFile, pruneHealthFile, providerHealthKey, loadEvents, aggregateEvents, } from "./health-store.js"; // ============================================================ // Types // ============================================================ type CandidateConfig = { provider: string; model: string; }; type Config = { virtualModel?: { id?: string; name?: string; reasoning?: boolean; input?: string[]; contextWindow?: number; maxTokens?: number; thinkingLevelMap?: Record; compat?: Record; }; candidates?: CandidateConfig[]; log?: boolean; timeoutMs?: number; cooldownMs?: number; strategy?: "sticky" | "round-robin" | "primary"; /** 单候选失败后冷却前的原地重试次数(不含首次)。0=不重试。 */ maxRetriesPerCandidate?: number; /** 可选的测速排序配置。详见 SpeedTestConfig。 */ speedTest?: SpeedTestConfig; }; type ResolvedCandidate = Model & { roundrobinLabel: string; }; type HealthStats = { success: number; fail: number; totalLatencyMs: number; lastFailAt: number | null; lastSuccessAt: number | null; }; type ResolvedVirtualModel = { id: string; name: string; reasoning: boolean; input: string[]; contextWindow: number; maxTokens?: number; thinkingLevelMap?: Record; compat?: Record; }; /** 单次测速结果(进程内, 不落盘; health.jsonl 只存聚合统计)。 */ type SpeedTestResult = { ok: boolean; ttftMs: number | null; latencyMs: number | null; error?: string; }; /** 一组虚拟模型(对应一个预设或 config.json 的 default)。 */ type Group = { name: string; // = 预设文件名;config.json 固定 "default" enabled: boolean; // ≥1 有效候选才 true virtualModel: Model; // id = name candidates: ResolvedCandidate[]; config: { virtualModel: ResolvedVirtualModel; candidates: CandidateConfig[]; timeoutMs: number; cooldownMs: number; strategy: NonNullable; log: boolean; maxRetriesPerCandidate: number; speedTest?: SpeedTestConfig; }; currentIndex: number; healthStats: Map; /** 最近一次测速的时刻(0=从未测过)。 */ lastSpeedTestAt: number; /** 最近一次测速的逐候选结果, key = roundrobinLabel。 */ lastSpeedTestResults: Map; /** 测速进行中标记, 防并发重测同组。 */ speedTestRunning: boolean; }; // ============================================================ // Constants & module-level state // ============================================================ const LOG_PATH = join(getAgentDir(), "roundrobin", "roundrobin.log"); const PROVIDER = "roundrobin"; const API = "model-roundrobin-api" as Api; const CONFIG_CHANGED_EVENT = "roundrobin:config-changed"; const groups = new Map(); // 模块级测速锁(按组名): reload 重建 group 对象后锁仍有效——旧测速闭包持旧对象, // 若锁存在 group 上, finally 清的是旧对象的锁, 新对象上锁永不释放(或硬编码 false // 导致并发双测)。按名存模块级, 旧测速释放锁能正确解锁, 新请求按名查锁正确等待。 const speedTestLocks = new Set(); function lockSpeedTest(group: Group): void { speedTestLocks.add(group.name); group.speedTestRunning = true; } function unlockSpeedTest(group: Group): void { speedTestLocks.delete(group.name); group.speedTestRunning = false; } function isSpeedTestLocked(group: Group): boolean { return speedTestLocks.has(group.name); } let statePi: ExtensionAPI | undefined; let stateModelRegistry: ExtensionContext["modelRegistry"] | undefined; let stateUi: ExtensionUIContext | undefined; // /rr-speedtest 详细结果 widget 的自动清除 timer (防重复触发时旧 timer 清掉新结果) let speedtestWidgetTimer: ReturnType | null = null; // ============================================================ // Logging & UI notify // ============================================================ const LOG_ROTATE_BYTES = 5 * 1024 * 1024; // 5MB 轮转阈值 let logLineCount = 0; function logLine(group: Group | null, message: string): void { const enabled = group ? group.config.log : true; if (!enabled) return; try { // 每 512 次检查一次大小(避免高频 append 的 stat 开销), 超 5MB 轮转保留一份 if ((logLineCount++ & 511) === 0) { try { const s = statSync(LOG_PATH); if (s.size > LOG_ROTATE_BYTES) { try { rmSync(`${LOG_PATH}.1`); } catch { /* ignore */ } try { renameSync(LOG_PATH, `${LOG_PATH}.1`); } catch { /* ignore */ } } } catch { /* 文件不存在等 */ } } const prefix = group ? `[${group.name}] ` : ""; appendFileSync(LOG_PATH, `${new Date().toISOString()} ${prefix}${message}\n`, "utf-8"); } catch { // ignore } } function notifyUi(text: string, level: "info" | "warning" | "error" = "info"): void { try { stateUi?.notify?.(text, level); } catch { // ignore } } // ============================================================ // Config → Group (per-name resolution) // ============================================================ function resolveVirtualModel(name: string, raw: Config["virtualModel"]): ResolvedVirtualModel { const v = raw ?? {}; return { id: name, // 强制: model id = 预设名(覆盖预设内 virtualModel.id) name: v.name || v.id || name, reasoning: v.reasoning ?? true, input: v.input || ["text", "image"], contextWindow: v.contextWindow ?? 200000, maxTokens: v.maxTokens, thinkingLevelMap: v.thinkingLevelMap, compat: v.compat, }; } function groupConfigFromParsed(parsed: Config): Group["config"] { const vm = resolveVirtualModel("", parsed.virtualModel); return { virtualModel: vm, candidates: Array.isArray(parsed.candidates) ? parsed.candidates : [], timeoutMs: typeof parsed.timeoutMs === "number" ? parsed.timeoutMs : 30000, cooldownMs: typeof parsed.cooldownMs === "number" ? parsed.cooldownMs : 60000, strategy: parsed.strategy === "round-robin" || parsed.strategy === "primary" ? parsed.strategy : "sticky", log: parsed.log !== false, maxRetriesPerCandidate: typeof parsed.maxRetriesPerCandidate === "number" ? parsed.maxRetriesPerCandidate : 2, speedTest: normalizeSpeedTest(parsed.speedTest), }; } function virtualModelFromGroup(group: Group): Model { const vm: Model = { id: group.config.virtualModel.id, // = group.name name: group.config.virtualModel.name, api: API, provider: PROVIDER, baseUrl: "model-roundrobin://local", reasoning: group.config.virtualModel.reasoning, input: group.config.virtualModel.input as Model["input"], cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, contextWindow: group.config.virtualModel.contextWindow, maxTokens: group.config.virtualModel.maxTokens ?? 16384, }; if (group.config.virtualModel.maxTokens !== undefined) { vm.maxTokens = group.config.virtualModel.maxTokens; } if (group.config.virtualModel.thinkingLevelMap) { vm.thinkingLevelMap = group.config.virtualModel.thinkingLevelMap; } if (group.config.virtualModel.compat) { vm.compat = group.config.virtualModel.compat as Model["compat"]; } return vm; } function resolveCandidates( candidates: CandidateConfig[], registry: ExtensionContext["modelRegistry"], group: Group | null, ): ResolvedCandidate[] { const resolved: ResolvedCandidate[] = []; const skipped: string[] = []; for (const candidate of candidates) { const model = registry.find(candidate.provider, candidate.model); if (!model) { logLine(group, `skip unknown model ${candidate.provider}/${candidate.model}`); skipped.push(`${candidate.provider}/${candidate.model}`); continue; } resolved.push({ ...model, roundrobinLabel: `${candidate.provider}/${candidate.model}` }); } // 批量通知, 避免逐个 notifyUi 刷屏 if (skipped.length) { const gname = group ? `[${group.name}] ` : ""; notifyUi(`roundrobin ${gname}跳过 ${skipped.length} 个失效候选: ${skipped.join(", ")}`, "warning"); } return resolved; } // ============================================================ // Health stats (per-group) // ============================================================ function healthKey(candidate: ResolvedCandidate): string { return candidate.roundrobinLabel; } function getHealth(group: Group, candidate: ResolvedCandidate): HealthStats { const key = healthKey(candidate); let h = group.healthStats.get(key); if (!h) { h = { success: 0, fail: 0, totalLatencyMs: 0, lastFailAt: null, lastSuccessAt: null }; group.healthStats.set(key, h); } return h; } function recordSuccess( group: Group, candidate: ResolvedCandidate, latencyMs: number, ttftMs?: number | null, probe?: boolean, ): void { const h = getHealth(group, candidate); h.success += 1; h.totalLatencyMs += latencyMs; h.lastSuccessAt = Date.now(); // 落盘到全局 7 天 health.jsonl(真实候选维度,跨 CLI 共享) appendHealthEvent({ ts: Date.now(), provider: candidate.provider, model: candidate.id, ok: true, latencyMs, ttftMs: ttftMs ?? null, ...(probe === true ? { probe: true } : {}), }); } function recordFailure(group: Group, candidate: ResolvedCandidate, probe?: boolean): void { const h = getHealth(group, candidate); h.fail += 1; h.lastFailAt = Date.now(); appendHealthEvent({ ts: Date.now(), provider: candidate.provider, model: candidate.id, ok: false, ...(probe === true ? { probe: true } : {}), }); } function isCoolingDown(group: Group, candidate: ResolvedCandidate): boolean { const h = group.healthStats.get(healthKey(candidate)); if (!h || h.lastFailAt === null) return false; return Date.now() - h.lastFailAt < group.config.cooldownMs; } /** * 轮询面板健康快照。 * - success/fail/rate/TTFT/延迟: 从 health.jsonl 聚合(跨 CLI 共享,重启不丢) * - currentIndex/coolingDown: 进程内 failover 运行时状态(重启清空,本就该进程内) */ function getHealthSnapshot(): unknown { // 从 health.jsonl 读 7 天聚合(与单 provider 健康同一份数据源) const agg = loadProviderHealthSnapshot() as { providers: Record; }; return { groups: [...groups.values()].map((g) => ({ name: g.name, enabled: g.enabled, strategy: g.config.strategy, timeoutMs: g.config.timeoutMs, cooldownMs: g.config.cooldownMs, maxRetriesPerCandidate: g.config.maxRetriesPerCandidate, currentIndex: g.currentIndex, // 进程内运行时 speedTest: { enabled: g.config.speedTest?.enabled === true, sortKey: g.config.speedTest?.sortKey ?? "ttft", minIntervalMs: g.config.speedTest?.minIntervalMs ?? 60000, timeoutMs: g.config.speedTest?.timeoutMs ?? 60000, concurrency: g.config.speedTest?.concurrency ?? 5, lastSpeedTestAt: g.lastSpeedTestAt || null, speedTestRunning: isSpeedTestLocked(g), }, candidates: g.candidates.map((c) => { // 从文件聚合查这个真实候选的 7 天统计 const h = agg.providers[providerHealthKey(c.provider, c.id)]; // 最近一次测速的逐候选结果 (进程内) const sr = g.lastSpeedTestResults.get(c.roundrobinLabel); return { provider: c.provider, model: c.id, label: c.roundrobinLabel, success: h?.success ?? 0, fail: h?.fail ?? 0, rate: h?.rate ?? null, avgTtftMs: h?.avgTtftMs ?? null, avgLatencyMs: h?.avgLatencyMs ?? null, coolingDown: isCoolingDown(g, c), // 进程内运行时 lastFailAt: h?.lastFailAt ?? null, lastSuccessAt: h?.lastSuccessAt ?? null, lastSpeedOk: sr?.ok ?? null, lastTtftMs: sr?.ttftMs ?? null, lastLatencyMs: sr?.latencyMs ?? null, lastSpeedError: sr?.error ?? null, // 动态超时: 直接复用 dynamicTimeoutMs(单一公式源, 避免双份维护漂移) dynamicTimeoutMs: dynamicTimeoutMs(g, c), }; }), })), }; } // ============================================================ // Provider registration (single call, all models) // ============================================================ function registerAllVirtualProviders(pi: ExtensionAPI): void { const allModels = [...groups.values()].filter((g) => g.enabled).map((g) => g.virtualModel); pi.registerProvider(PROVIDER, { baseUrl: "model-roundrobin://local", apiKey: "model-roundrobin", api: API, models: allModels, streamSimple: streamRoundRobin, }); } // ============================================================ // Reload: build groups from all preset configs // ============================================================ function buildGroups(allConfigs: GroupConfig[]): void { if (!stateModelRegistry) return; // 暂存旧 healthStats / 测速状态按 name 回填,避免 reload 丢历史统计 // lastSpeedTestAt 也回填:否则面板每次保存配置就把它重置为 0, // 触发“首次请求”判定重新生效 → 保存一次就重复测一次(用户痛点)。 const oldHealth = new Map>(); const oldSpeedAt = new Map(); const oldSpeedResults = new Map>(); for (const [name, g] of groups) { oldHealth.set(name, g.healthStats); oldSpeedAt.set(name, g.lastSpeedTestAt); oldSpeedResults.set(name, g.lastSpeedTestResults); } groups.clear(); for (const { name, config: parsed } of allConfigs) { const config = groupConfigFromParsed(parsed); config.virtualModel.id = name; // 强制 model id = 预设名(覆盖 groupConfigFromParsed 的空值) const candidates = resolveCandidates(config.candidates, stateModelRegistry, null); const group: Group = { name, enabled: candidates.length > 0, virtualModel: {} as Model, // 先占位,下面 virtualModelFromGroup 填 candidates, config, currentIndex: 0, healthStats: oldHealth.get(name) ?? new Map(), lastSpeedTestAt: oldSpeedAt.get(name) ?? 0, lastSpeedTestResults: ((): Map => { // 回填后按新候选 label 过滤: 删掉的候选 label 残留会污染 // getHealthSnapshot 的 lastSpeedOk/dynamicTimeout 等读数 const old = oldSpeedResults.get(name); if (!old) return new Map(); const live = new Set(candidates.map((c) => c.roundrobinLabel)); const filtered = new Map(); for (const [k, v] of old) if (live.has(k)) filtered.set(k, v); return filtered; })(), speedTestRunning: speedTestLocks.has(name), }; group.virtualModel = virtualModelFromGroup(group); groups.set(name, group); logLine( group, `load enabled=${group.enabled} candidates=${candidates.length}/${config.candidates.length}`, ); } } function reloadRuntime(): void { if (!stateModelRegistry || !statePi) { // 还没 session_start,无法 resolve;session_start 会读最新配置 return; } try { // 刷新 registry 以读到面板新保存的模型(否则新加候选 resolve 不到) stateModelRegistry.refresh(); const allConfigs = loadAllGroupsConfig(); buildGroups(allConfigs); registerAllVirtualProviders(statePi); const enabledNames = [...groups.values()].filter((g) => g.enabled).map((g) => g.name); logLine(null, `reload groups=${groups.size} enabled=${enabledNames.join(",") || "(none)"}`); // 初始测速改为懒触发:首次请求某组时才测(见 streamRoundRobin),不在 reload 时狂测所有组 } catch (e) { const msg = e instanceof Error ? e.message : String(e); logLine(null, `reload failed: ${msg}`); notifyUi(`roundrobin 热重载失败: ${msg}`, "error"); } } // ============================================================ // streamSimple — failover entry (routes by model.id = group name) // ============================================================ function makeErrorMessage( model: Model, message: string, stopReason: "error" | "aborted" = "error", ): AssistantMessage { return { role: "assistant", content: [], api: (model?.api ?? API) as Api, provider: model?.provider ?? PROVIDER, model: model?.id ?? "roundrobin", usage: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, totalTokens: 0, cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, }, stopReason, errorMessage: message, timestamp: Date.now(), } as AssistantMessage; } function streamRoundRobin( model: Model, context: Context, options?: SimpleStreamOptions, ): AssistantMessageEventStream { const outer = createAssistantMessageEventStream(); (async () => { const group = groups.get(model.id); if (!group || !group.enabled) { const message = makeErrorMessage(model, `roundrobin preset "${model.id}" not found or disabled`); outer.push({ type: "error", reason: "error", error: message }); outer.end(); return; } // 首次请求该组且测速已开启:同步等测速排完再发(用户要首次就享受排序)。 // 不同组独立判定(lastSpeedTestAt===0 = 本会话从未测过),reload 不重置该状态。 // 并发:若另一请求正在测同组,等它测完;若对方被 abort 中断没测成 // (lastSpeedTestAt 仍为 0),由本请求接手补测 — 而非未排序裸发。 if (group.config.speedTest?.enabled === true && group.lastSpeedTestAt === 0) { if (isSpeedTestLocked(group)) { logLine(group, "first-request: 另一请求正在测速, 等待完成后再发"); while (isSpeedTestLocked(group)) { if (options?.signal?.aborted) { outer.push({ type: "error", reason: "aborted", error: makeErrorMessage(model, "Request was aborted", "aborted") }); outer.end(); return; } const ok = await waitForRetry(250, options?.signal); if (!ok) { outer.push({ type: "error", reason: "aborted", error: makeErrorMessage(model, "Request was aborted", "aborted") }); outer.end(); return; } } } // 锁已释放: 正常完成(对方已测,lastSpeedTestAt>0 跳过) 或 对方被 abort 没测成(=0 则自己补测) if (group.lastSpeedTestAt === 0 && !isSpeedTestLocked(group)) { logLine(group, "first-request: 首次请求该组, 同步触发测速排序"); notifyUi(`[${group.name}] 首次请求触发测速排序…`, "info"); lockSpeedTest(group); try { await runSpeedTest(group, options?.signal); } catch (e) { logLine(group, `first-request speedtest failed: ${e instanceof Error ? e.message : String(e)}`); } finally { unlockSpeedTest(group); // 若首测被用户 abort 中断: 排序结果不可信(全候选被 abort 判败+进冷却), // 重置回“从未测过”, 下次请求重新首测; 否则 abort 会永久消耗首次资格。 if (options?.signal?.aborted) { group.lastSpeedTestAt = 0; logLine(group, "first-request speedtest aborted; reset lastSpeedTestAt=0"); } } } } const lastFailure = await tryCandidates(group, model, context, options, outer); if (lastFailure) { const reason = lastFailure.stopReason === "aborted" ? "aborted" : "error"; outer.push({ type: "error", reason, error: lastFailure }); outer.end(); } })().catch((error) => { const message = makeErrorMessage(model, error instanceof Error ? error.message : String(error)); outer.push({ type: "error", reason: "error", error: message }); outer.end(); }); return outer; } async function tryCandidates( group: Group, model: Model, context: Context, options: SimpleStreamOptions | undefined, outer: AssistantMessageEventStream, ): Promise { // 快照本请求期间的候选数组与策略: 防止测速(runSpeedTest) 重排序在并发请求中途 // 用 group.candidates = sorted 整组换数组, 导致本请求的 index 取到错位候选、且成功后 // group.currentIndex 被写回错位索引(黏错人)。本请求始终迭代此快照, 成功后按 label // 定位实时 group.candidates 中的真实位置再更新 currentIndex。 const candidates = group.candidates; const total = candidates.length; if (total === 0) { return makeErrorMessage(model, `roundrobin/${group.name}: no valid candidates`); } const strategy = group.config.strategy; let lastFailure: AssistantMessage | null = null; let firstTriedLabel: string | null = null; // 请求内首个轮到的候选, 用于区分“真 failover”与“回首选” let round = 0; let allCooldownClears = 0; // 连续全冷却清空次数, 防止候选持续失败时忙循环 const MAX_ALL_COOLDOWN_CLEARS = 3; // 超过则不再清空, 走正常等待/失败 const MAX_ROUNDS = 5; // 单请求最大轮数(防无限 failover 烧 token; 日志曾见 127 候选病态链) // 全冷却且清空额度耗尽后, 不再原地阻塞等满整个 cooldownMs(用户常配 600-900s, 会把 // 单个请求挂 5-15 分钟); 短退避后即返回错误, 让用户自己决定重试。 const MAX_COOLDOWN_FALLBACK_WAIT_MS = 5000; const currentIndex = Number.isInteger(group.currentIndex) && group.currentIndex >= 0 ? group.currentIndex % total : 0; // 回首选(index 0) 是 sticky / primary 的设计: 成功黏当前, 过冷却期后下次回首选。 // round-robin 的语义是“成功后指向下一候选, 均分流量” — 若也对它回首选, 只要首选活着 // 流量就永远全给 index 0, round-robin 名存实亡。故回首选仅对 sticky/primary 生效。 const firstRoundStartIndex = strategy !== "round-robin" && currentIndex !== 0 && !isCoolingDown(group, candidates[currentIndex]) ? 0 : currentIndex; while (true) { if (options?.signal?.aborted) { return makeErrorMessage(model, "Request was aborted", "aborted"); } if (round >= MAX_ROUNDS) { const reason = lastFailure?.errorMessage || "all candidates exhausted after max rounds"; logLine(group, `giving up after ${MAX_ROUNDS} rounds; last error: ${reason}`); return makeErrorMessage(model, reason, "error"); } // The retry cursor is request-local so concurrent requests cannot reset the // group's sticky state. Only a successful candidate updates currentIndex. const startIndex = round === 0 ? firstRoundStartIndex : 0; let attempted = 0; for (let offset = 0; offset < total; offset += 1) { if (options?.signal?.aborted) { return makeErrorMessage(model, "Request was aborted", "aborted"); } const index = (startIndex + offset) % total; const candidate = candidates[index]; if (isCoolingDown(group, candidate)) { logLine(group, `skip cooling-down ${candidate.roundrobinLabel}`); continue; } attempted += 1; if (firstTriedLabel === null) firstTriedLabel = candidate.roundrobinLabel; // 单候选原地重试: 失败后按指数退避重试, 最多 maxRetriesPerCandidate 次(不含首试)。 // 只有全部重试耗尽才 recordFailure 进冷却, 然后切下一个候选。 const maxAttempts = 1 + Math.max(0, group.config.maxRetriesPerCandidate); let attemptNo = 0; while (attemptNo < maxAttempts) { attemptNo += 1; if (options?.signal?.aborted) { return makeErrorMessage(model, "Request was aborted", "aborted"); } const isRetry = attemptNo > 1; if (isRetry) { const backoffMs = Math.min(30000, 2000 * 2 ** (attemptNo - 2)); logLine(group, `retry ${candidate.roundrobinLabel} #${attemptNo - 1}/${maxAttempts - 1} after ${backoffMs}ms`); if (!(await waitForRetry(backoffMs, options?.signal))) { return makeErrorMessage(model, "Request was aborted", "aborted"); } } const outcome = await attemptOneCandidate(group, model, candidate, context, options, outer, index); if (outcome.kind === "success") { // TUI 提示 failover: 仅在成功候选与首个轮到的不同时闪一下 (不进对话历史) if (firstTriedLabel !== null && candidate.roundrobinLabel !== firstTriedLabel) { notifyUi(`↔ 轮询故障转移到 ${candidate.roundrobinLabel}`); } return null; // outer 已 end } if (outcome.kind === "aborted") { return makeErrorMessage(model, outcome.reason || "Request was aborted", "aborted"); } if (outcome.kind === "terminal") { // 流中途出错(forwarded 后): 内容已吐出, 不能换 provider 重放, 直接终止。 // 但仍记一次失败 + 进冷却: 该渠道吐了一半就断属质量问题, health/smart 排序 // 必须能看到(否则 7 天统计里查不到“经常中途断”的渠道, smart 也不会压低它)。 recordFailure(group, candidate); return outcome.failure; } // outcome.kind === "fail": 流前失败 lastFailure = outcome.failure; // 仍有重试次数则继续 while; 耗尽则跳出进冷却。 // allCooldownClears 是请求级计数器, 成功分支已 return 无需重置。 } // 重试耗尽, 正式记一次失败 + 进冷却 recordFailure(group, candidate); logLine(group, `fail ${candidate.roundrobinLabel} after ${maxAttempts} attempts: ${lastFailure?.errorMessage || "unknown"}`); continue; } // === 测速模式分支: 整轮全炸后触发重测+排序 (speedTest.enabled=true) === if (group.config.speedTest?.enabled === true) { const now = Date.now(); const minInterval = group.config.speedTest?.minIntervalMs ?? 60000; const sinceLast = now - (group.lastSpeedTestAt || 0); if (sinceLast < minInterval) { // 节流期内: 不重测, 走原退避 (剩余节流时间与冷却键较短者) const throttleRemain = minInterval - sinceLast; const waitMs = Math.min(Math.max(25, throttleRemain), getRetryWaitMs(group)); logLine(group, `speedTest: throttle ${Math.round(throttleRemain / 1000)}s remaining; wait ${waitMs}ms`); if (!(await waitForRetry(waitMs, options?.signal))) { return makeErrorMessage(model, "Request was aborted", "aborted"); } round += 1; continue; } // 若已在测速中(init 并发), 退避让出一小段再轮 if (isSpeedTestLocked(group)) { const waitMs = Math.min(2000, getRetryWaitMs(group)); logLine(group, `speedTest: already running; wait ${waitMs}ms`); if (!(await waitForRetry(waitMs, options?.signal))) { return makeErrorMessage(model, "Request was aborted", "aborted"); } round += 1; continue; } lockSpeedTest(group); try { notifyUi(`[${group.name}] 重新测速排序中…`, "info"); const start = Date.now(); const { okCount, badCount } = await runSpeedTest(group, options?.signal); const elapsed = Date.now() - start; notifyUi( okCount > 0 ? `[${group.name}] 测速完成 ${elapsed}ms: 可用 ${okCount}/${total}, 按序重试` : `[${group.name}] 测速完成 ${elapsed}ms: 全部不可用`, okCount > 0 ? "info" : "error", ); logLine(group, `speedtest round triggered ok=${okCount} bad=${badCount}`); } catch (e) { logLine(group, `speedtest round failed: ${e instanceof Error ? e.message : String(e)}`); } finally { unlockSpeedTest(group); } // 测速已完成排序 + 清成功的候选冷却; 下一轮从头试 round += 1; continue; } // === 以下仅 speedTest.enabled=false 时执行原冷却逻辑 === // 整轮耗尽。若所有候选都在冷却, 立即清空冷却从头轮询(不等任何 CD), 并在 TUI 提示一次。 const allCooling = candidates.every((c) => isCoolingDown(group, c)); if (allCooling && allCooldownClears < MAX_ALL_COOLDOWN_CLEARS) { allCooldownClears += 1; notifyUi( `[${group.name}] 全部候选冷却中, 已清空冷却重试 (${allCooldownClears}/${MAX_ALL_COOLDOWN_CLEARS})`, "warning", ); logLine(group, `all candidates cooling down (#${allCooldownClears}/${MAX_ALL_COOLDOWN_CLEARS}); clearing cooldowns and retrying immediately`); for (const c of candidates) { const h = group.healthStats.get(healthKey(c)); if (h) h.lastFailAt = null; } // 不等, 直接进入下一轮 (startIndex = 0) round += 1; continue; } // 全冷却且清空额度已耗尽: 不再原地阻塞等满整个 cooldownMs(用户常配 600-900s, 会把单个 // 请求挂 5-15 分钟), 短退避后即返回错误, 让用户自己决定重试(TUI 提示实际等待时长)。 if (allCooling) { const waitMs = Math.min(MAX_COOLDOWN_FALLBACK_WAIT_MS, getRetryWaitMs(group)); notifyUi( `[${group.name}] 全部候选持续冷却, 放弃请求 (退避 ${Math.round(waitMs / 1000)}s 后失败); 请稍后重试`, "error", ); logLine(group, `all candidates still cooling down after ${MAX_ALL_COOLDOWN_CLEARS} clears; final backoff ${waitMs}ms then fail`); if (!(await waitForRetry(waitMs, options?.signal))) { return makeErrorMessage(model, "Request was aborted", "aborted"); } const reason = lastFailure?.errorMessage || `全部候选冷却中 (已重置 ${MAX_ALL_COOLDOWN_CLEARS} 次), 请稍后重试`; return makeErrorMessage(model, reason, "error"); } round += 1; const waitMs = getRetryWaitMs(group); const reason = lastFailure?.errorMessage || "unknown error"; logLine( group, `round ${round} exhausted (${attempted}/${total} attempted); retry from preferred in ${waitMs}ms; last error: ${reason}`, ); if (!(await waitForRetry(waitMs, options?.signal))) { return makeErrorMessage(model, "Request was aborted", "aborted"); } } } /** 单次尝试一个候选的返回结果。 */ type AttemptOutcome = | { kind: "success" } // 已成功并 outer.end() | { kind: "fail"; failure: AssistantMessage; reason: string } // 流前失败, 可重试 | { kind: "terminal"; failure: AssistantMessage; reason: string } // 流中途失败, 不可重试 | { kind: "aborted"; reason: string }; // 用户取消 /** 尝试单个候选一次。建流→首事件→流中→结束。 * 不在此处 recordFailure: 重试计数由调用方控制, 耗尽后统一记失败+冷却。 */ async function attemptOneCandidate( group: Group, model: Model, candidate: ResolvedCandidate, context: Context, options: SimpleStreamOptions | undefined, outer: AssistantMessageEventStream, _index: number, // 保留用于扩展(如排序前后的 offset 监控); currentIndex 已改由 label-find 计算 ): Promise { logLine(group, `use ${candidate.roundrobinLabel}`); // 动态超时: 测速实测 ttft×2.0(下限=timeoutMs 面板可调), 无测速结果回退 timeoutMs。 // 同一公式用于首字与流中 chunk 间空闲, 使慢卡候选不再无限栏 (仍走原重试/重试逻辑)。 const startTime = Date.now(); let forwarded = false; let iteratorEnded = false; let firstTtftMs: number | null = null; let iter: AsyncIterator<{ type: string; [k: string]: unknown }> | undefined; const attemptController = new AbortController(); const abortAttempt = (): void => attemptController.abort(); options?.signal?.addEventListener("abort", abortAttempt, { once: true }); if (options?.signal?.aborted) abortAttempt(); const attemptOptions: SimpleStreamOptions = { ...options, signal: attemptController.signal, }; // 首字与流中空闲共用动态超时; 流中用同一值避免 chunk 间体感担面 (实测 ttft 是整体响应节奏的好代理) const idleTimeoutMs = dynamicTimeoutMs(group, candidate); try { const auth = await waitWithAbort(stateAuth(candidate, attemptOptions), attemptController.signal); const stream = streamSimple(candidate, context, auth); iter = stream[Symbol.asyncIterator]() as AsyncIterator<{ type: string; [k: string]: unknown }>; const firstResult = await nextWithTimeout(iter, idleTimeoutMs, attemptController.signal); if (firstResult.done) { iteratorEnded = true; throw new Error("stream ended before start event"); } const firstEvent = firstResult.value; if (firstEvent.type === "error") { const errEvent = firstEvent as unknown as { error?: AssistantMessage }; const upstreamError = errEvent.error; const reason = upstreamError?.errorMessage || upstreamError?.stopReason || "unknown upstream error"; if (options?.signal?.aborted || upstreamError?.stopReason === "aborted") { iteratorEnded = true; return { kind: "aborted", reason }; } logLine(group, `fail ${candidate.roundrobinLabel}: ${reason}`); const failure = upstreamError ? { ...upstreamError, provider: PROVIDER, model: model.id, api: model.api } : makeErrorMessage(model, reason); iteratorEnded = true; return { kind: "fail", failure, reason }; } if (firstEvent.type !== "start") { throw new Error(`unexpected first event: ${firstEvent.type}`); } outer.push(firstResult.value as never); forwarded = true; firstTtftMs = Date.now() - startTime; while (true) { let nextResult: IteratorResult<{ type: string; [k: string]: unknown }>; try { nextResult = await nextWithAbort(iter, attemptController.signal, idleTimeoutMs); } catch (e) { // 流中 chunk 间空闲超时: 日志标注, 抬出为 fail (走原重试逻辑)。 if (options?.signal?.aborted) { return { kind: "aborted", reason: "Request was aborted" }; } const reason = e instanceof Error ? e.message : String(e); logLine(group, `fail ${candidate.roundrobinLabel}: stream idle timeout after ${idleTimeoutMs}ms`); const failure = makeErrorMessage(model, `stream idle timeout: ${reason}`); // 已转发过内容 → 不能换 provider 重放, 终止 if (forwarded) { iteratorEnded = true; return { kind: "terminal", failure, reason }; } iteratorEnded = true; return { kind: "fail", failure, reason }; } if (nextResult.done) { iteratorEnded = true; break; } const event = nextResult.value; if (event.type === "error") { const errEvent = event as unknown as { error: AssistantMessage }; const reason = errEvent.error.errorMessage || errEvent.error.stopReason || "unknown error"; if (options?.signal?.aborted || errEvent.error.stopReason === "aborted") { logLine(group, `aborted ${candidate.roundrobinLabel}: ${reason}`); } else { logLine(group, `fail ${candidate.roundrobinLabel}: ${reason}`); } const failure: AssistantMessage = { ...errEvent.error, provider: PROVIDER, model: model.id, api: model.api, }; iteratorEnded = true; return { kind: "terminal", failure, reason }; } const doneEvent = event as unknown as { type: string; message: AssistantMessage }; if (doneEvent.type === "done") { const stopReason = doneEvent.message.stopReason; // done 带 error/aborted: 很多渠道把流中错误以 done+error 形式送达(而非独立 error 事件)。 // 不能当成功 — 需记失败 + 按是否已转发分流(fail 可重试 / terminal 不重放), 否则用户拿到错误终止的回复、 // 候选不进冷却也不 failover, 与 terminal 设计矛盾。 if (stopReason === "error" || stopReason === "aborted") { // 用户主动取消(Esc)不应记候选失败/进冷却 — 原样终止 if (stopReason === "aborted" && options?.signal?.aborted) { iteratorEnded = true; return { kind: "aborted", reason: "Request was aborted" }; } const reason = doneEvent.message.errorMessage || stopReason || "unknown error"; const failure = { ...doneEvent.message, provider: PROVIDER, model: model.id, api: model.api } as AssistantMessage; if (forwarded) { // 已转发内容(吐了一半) — 不重放, 终止 + 记失败进冷却 logLine(group, `terminal ${candidate.roundrobinLabel} (done+${stopReason}): ${reason}`); iteratorEnded = true; return { kind: "terminal", failure, reason }; } // 未转发 — 可重试的流前失败 logLine(group, `fail ${candidate.roundrobinLabel} (done+${stopReason}): ${reason}`); iteratorEnded = true; return { kind: "fail", failure, reason }; } // strategy 语义: sticky=成功后黏住当前候选(冷却后回首选); round-robin=成功后 // 指向下一下候选(均分流); primary=恒回首选(0) — 只有首选冷却时才用备选。 // 按 label 在实时 group.candidates 上定位候选: 若测速已在并发请求中途重排 // (group.candidates = sorted), 快照 index 已错位; label-find 保证 currentIndex // 始终指向真实候选在新数组中的正确位置, 不黏错人。 const liveArr = group.candidates; const liveIdx = liveArr.findIndex((c) => c.roundrobinLabel === candidate.roundrobinLabel); group.currentIndex = liveIdx < 0 ? group.currentIndex // 候选已不在实时数组(理论不发生: 测速只重排不删候选) : group.config.strategy === "round-robin" ? (liveIdx + 1) % liveArr.length : group.config.strategy === "primary" ? 0 : liveIdx; const successLatencyMs = Date.now() - startTime; recordSuccess(group, candidate, successLatencyMs, firstTtftMs); logLine(group, `success ${candidate.roundrobinLabel} (${successLatencyMs}ms)`); // 改写 done message 的 provider/model 为虚拟模型 — 否则转发后 pi 的 // message_end 钱子见真实候选 provider(≠PROVIDER) 会再调 recordProviderSuccess, // 与上面的 recordSuccess 双写 health.jsonl(总数×2, 影响面板显示与计数)。 doneEvent.message = { ...doneEvent.message, provider: PROVIDER, model: model.id, api: model.api }; outer.push(event as never); iteratorEnded = true; outer.end(); return { kind: "success" }; } outer.push(event as never); } if (forwarded) { const reason = "stream ended without done or error event"; logLine(group, `fail ${candidate.roundrobinLabel}: ${reason}`); return { kind: "terminal", failure: makeErrorMessage(model, reason), reason }; } const reason = "stream ended before start event"; return { kind: "fail", failure: makeErrorMessage(model, reason), reason }; } catch (error) { if (options?.signal?.aborted) { return { kind: "aborted", reason: "Request was aborted" }; } const reason = error instanceof Error ? error.message : String(error); logLine(group, `fail ${candidate.roundrobinLabel}: ${reason}`); const failure = makeErrorMessage(model, reason); if (forwarded) { // The outer stream already emitted start. Replaying on another provider // could duplicate content or tool calls, so this error is terminal. return { kind: "terminal", failure, reason }; } return { kind: "fail", failure, reason }; } finally { options?.signal?.removeEventListener("abort", abortAttempt); attemptController.abort(); if (iter && !iteratorEnded) { try { // Cleanup is best-effort: a broken iterator.return() must not block // user cancellation or failover to the next candidate. void Promise.resolve(iter.return?.()).catch(() => {}); } catch { // ignore synchronous cleanup error } } } } function waitWithAbort(promise: Promise, signal?: AbortSignal): Promise { if (!signal) return promise; if (signal.aborted) return Promise.reject(new Error("Request was aborted")); return new Promise((resolve, reject) => { let settled = false; const finish = (callback: () => void): void => { if (settled) return; settled = true; signal.removeEventListener("abort", onAbort); callback(); }; const onAbort = (): void => finish(() => reject(new Error("Request was aborted"))); signal.addEventListener("abort", onAbort, { once: true }); if (signal.aborted) { onAbort(); return; } promise.then( (value) => finish(() => resolve(value)), (error) => finish(() => reject(error)), ); }); } /** 取候选动态超时: 优先用测速实测 ttft×2.0(下限=timeoutMs, 上限 120s); 无测速/ttft=null 回退 group.config.timeoutMs。 * 对首字与流中 chunk 间空闲都用同一公式, 避免慢卡候选无限桩。 * 系数/下限取舍: 测速 ttft 是 maxTokens=16 极简请求的首字, 真实请求首字普遍更慢 * (prompt 更长、思考更多), 故用 ×2.0 放大给抖动留余量; 下限取你配的 timeoutMs 因中转站首字常见 * 6~19s, 旧 10s 下限会误杀正常慢响应(实测 runway ttft 5.8~18.9s 被判 timeout)。 * 上限 120s: 病态样本(实测曾测出 ttft 327s)不该把超时上限跟看到 10 分钟, 失去保护。 */ function dynamicTimeoutMs(group: Group, candidate: ResolvedCandidate): number { const measured = group.lastSpeedTestResults.get(candidate.roundrobinLabel)?.ttftMs; if (typeof measured === "number" && measured > 0 && Number.isFinite(measured)) { // 下限 = 用户配的 timeoutMs (面板可调), 不再硬编码 20s — // 中转站波动大, 让用户自己决定多宽的容忍窗口; 没测速时也回退 timeoutMs, 语义一致。 // 上限 120s 仍硬编码(防病态样本把超时抬到无穷)。 return Math.max(group.config.timeoutMs, Math.min(120000, Math.round(measured * 2.0))); } return group.config.timeoutMs; } function nextWithAbort( iter: AsyncIterator<{ type: string; [k: string]: unknown }>, signal?: AbortSignal, idleTimeoutMs?: number, ): Promise> { if (signal?.aborted) return Promise.reject(new Error("Request was aborted")); if (typeof idleTimeoutMs !== "number" || !Number.isFinite(idleTimeoutMs) || idleTimeoutMs <= 0) { return waitWithAbort(iter.next(), signal); } return nextWithTimeout(iter, idleTimeoutMs, signal); } async function nextWithTimeout( iter: AsyncIterator<{ type: string; [k: string]: unknown }>, timeoutMs: number, signal?: AbortSignal, ): Promise> { if (signal?.aborted) throw new Error("Request was aborted"); return new Promise((resolve, reject) => { let settled = false; let timer: ReturnType; const finish = (callback: () => void): void => { if (settled) return; settled = true; clearTimeout(timer); signal?.removeEventListener("abort", onAbort); callback(); }; const onAbort = (): void => finish(() => reject(new Error("Request was aborted"))); timer = setTimeout( () => finish(() => reject(new Error(`timeout after ${timeoutMs}ms`))), timeoutMs, ); signal?.addEventListener("abort", onAbort, { once: true }); if (signal?.aborted) { onAbort(); return; } iter.next().then( (result) => finish(() => resolve(result)), (error) => finish(() => reject(error)), ); }); } function getRetryWaitMs(group: Group): number { const now = Date.now(); const waits = group.candidates .map((candidate) => { const lastFailAt = group.healthStats.get(healthKey(candidate))?.lastFailAt; return lastFailAt === null || lastFailAt === undefined ? 0 : Math.max(0, lastFailAt + group.config.cooldownMs - now); }) .filter((wait) => wait > 0); const nextCooldown = waits.length > 0 ? Math.min(...waits) : group.config.cooldownMs; return Math.max(25, nextCooldown); } function waitForRetry(delayMs: number, signal?: AbortSignal): Promise { if (signal?.aborted) return Promise.resolve(false); return new Promise((resolve) => { let settled = false; let timer: ReturnType; const finish = (completed: boolean): void => { if (settled) return; settled = true; clearTimeout(timer); signal?.removeEventListener("abort", onAbort); resolve(completed); }; const onAbort = (): void => finish(false); timer = setTimeout(() => finish(true), delayMs); signal?.addEventListener("abort", onAbort, { once: true }); if (signal?.aborted) onAbort(); }); } async function stateAuth( candidate: ResolvedCandidate, options: SimpleStreamOptions | undefined, ): Promise { const auth = await stateModelRegistry?.getApiKeyAndHeaders(candidate); if (!auth || !auth.ok) { return options; } return { ...options, apiKey: auth.apiKey, headers: auth.headers, }; } // ============================================================ // Provider health stats (全量: 所有 provider/model 的 TTFT + 成功率 + 总延迟) // 采集靠 message_update(首个 text_delta=TTFT) + message_end(latency=msg.timestamp); // 不再用 before_provider_request+pending 队列配对(中断/取消时 end 缺失会泄漏, // FIFO 错配会把 latency 算成"现在-2小时前")。message_end 自带 message.timestamp, // 直接算 latency 无需配对, 彻底避免泄漏与跨请求错配。持久化为 7 天 JSONL(跨 CLI 共享)。 // ============================================================ /** 每个 provider/model 的首个 text_delta 到达时刻(用于算 TTFT)。 * value.msgTs 识别是否属于当前请求(避免上次请求残留错配)。 * latency 用 message_end 的 message.timestamp 计算, 不依赖任何队列配对。 */ const firstTokenAt = new Map(); function recordProviderSuccess(provider: string, model: string, latencyMs: number, ttftMs: number | null): void { appendHealthEvent({ ts: Date.now(), provider, model, ok: true, latencyMs, ttftMs: ttftMs ?? null, }); } function recordProviderFailure(provider: string, model: string): void { appendHealthEvent({ ts: Date.now(), provider, model, ok: false, }); } function getProviderHealthSnapshot(): unknown { // 每次从磁盘聚合,多 CLI 共享同一份 7 天窗口 return loadProviderHealthSnapshot(); } function resetProviderHealth(): void { firstTokenAt.clear(); resetProviderHealthFile(); } // ============================================================ // Model liveness test — 发一个极小真对话请求验证模型可用性 // 复用 streamSimple + modelRegistry 鉴权, 读流首字/完成/错误。省钱: maxTokens=16。 // 每次刷新 registry 读最新 models.json; 未保存的改动仍测不到(只在面板内存)。 // ============================================================ type TestResult = { ok: boolean; error?: string; ttftMs?: number | null; latencyMs?: number; reply?: string; }; /** 测速入口(带重试): 单次失败且非 abort/鉴权问题 → 退避 1.5s 再试, 默认 retries=2(共 3 试)。 * 为什么重试: 测速失败会 recordFailure 进冷却+排序沉底, 一次网络抖动就把 86% 成功率 * 的可用渠道判死太苛刻; 两次连败才判失败, 抖动基本能救回(86%²≈98%能过), * 真挂的渠道两次都败判死正确, 多花的只是一次注定失败的请求(失败请求不耗 token)。 * 鉴权/配置类失败不重试(重试无意义); abort(Esc/超时联动)不重试直接返回。 */ async function testCandidateSpeed( candidate: ResolvedCandidate, opts: { prompt?: string; timeoutMs?: number; signal?: AbortSignal; retries?: number }, logTag = "testModel", ): Promise { // 默认重试 2 次(共 3 次尝试): 两次退避重试把暂时抖动的存活率拉满 // (86% 渠道 3 连败概率 ≈0.3%), 真挂渠道多花的只是注定失败的请求(不耗 token)。 const maxAttempts = 1 + Math.max(0, opts.retries ?? 2); let last: TestResult = { ok: false, error: "no attempt" }; for (let attempt = 1; attempt <= maxAttempts; attempt += 1) { if (opts.signal?.aborted) return { ...last, error: last.error || "aborted" }; const r = await testCandidateSpeedOnce(candidate, opts, attempt > 1 ? `${logTag}#retry${attempt}` : logTag); if (r.ok) return r; last = r; // 鉴权/配置类失败不重试 if (r.error === "扩展未初始化" || r.error === "未配置 API Key(保存后重试)") return r; if (attempt < maxAttempts && !opts.signal?.aborted) { logLine(null, `${logTag} attempt ${attempt} failed (${r.error || "unknown"}); retry in 1500ms`); await new Promise((resolve) => { const t = setTimeout(resolve, 1500); opts.signal?.addEventListener("abort", () => { clearTimeout(t); resolve(); }, { once: true }); }); } } return last; } /** 测速核心(单次): 对一个已 resolve 的候选发极简对话(maxTokens=16), 采 TTFT/latency/reply。 * 面板 testModel 和引擎 runSpeedTest 共用此函数 → 行为统一, 不会被当探活封号。 * signal 联动: 若外部 signal(用户 Esc) abort, 内部 AbortController 随之 abort, 立即终止。 */ async function testCandidateSpeedOnce( candidate: ResolvedCandidate, opts: { prompt?: string; timeoutMs?: number; signal?: AbortSignal }, logTag = "testModel", ): Promise { const prompt = opts.prompt || "欧拉函数的意义?"; const timeoutMs = opts.timeoutMs ?? 60000; if (!stateModelRegistry) { return { ok: false, error: "扩展未初始化" }; } const auth = await stateModelRegistry.getApiKeyAndHeaders(candidate); if (!auth || !auth.ok || !auth.apiKey) { logLine(null, `${logTag} FAIL: 鉴权失败 ok=${auth?.ok}`); return { ok: false, error: "未配置 API Key(保存后重试)" }; } // 测速 maxTokens: 普通候选 16(省钱); reasoning 候选抬到 2048 — // Anthropic 类 API 开 thinking 时要求 max_tokens > thinking budget(最小1024), // 16 直接 400 → reasoning 候选测速永远失败被沉底+冷却(误杀)。 // 2048 仍只测首字时延, 多出的额度只在模型真用思考时才消耗。 const probeMaxTokens = (candidate as { reasoning?: string }).reasoning ? 2048 : 16; logLine(null, `${logTag} auth ok, 调 streamSimple (maxTokens=${probeMaxTokens}, timeout=${timeoutMs}ms)`); const startTs = Date.now(); // UserMessage.content 本就是 string | block[] — 裸字符串合法, 无需 as never 绕类型 // (此前 as never 屏蔽了编译期信号, 若 pi-ai 改类型会静默运行时炸) const context: Context = { messages: [{ role: "user", content: prompt, timestamp: startTs }], }; const ctrl = new AbortController(); const timer = setTimeout(() => ctrl.abort(), timeoutMs); // 联动外部 signal (用户 Esc 打断测速); finally 必须 removeEventListener, // 否则每次尝试泄漏一个闭包挂在用户请求 signal 上(请求存活期间无法回收) const onOuterAbort = (): void => ctrl.abort(); if (opts.signal) { if (opts.signal.aborted) { ctrl.abort(); } else opts.signal.addEventListener("abort", onOuterAbort, { once: true }); } let ttftMs: number | null = null; let reply = ""; try { const stream: AssistantMessageEventStream = streamSimple(candidate, context, { apiKey: auth.apiKey, headers: auth.headers, maxTokens: probeMaxTokens, signal: ctrl.signal, } as SimpleStreamOptions); for await (const ev of stream as AsyncIterable<{ type: string; [k: string]: unknown }>) { if (ev.type === "start") { logLine(null, `${logTag} stream start`); } else if (ev.type === "text_delta" || ev.type === "thinking_delta") { // TTFT = 首个 text_delta(真正首字); reasoning 模型 maxTokens=16 可能全花在 // thinking 上、永不产 text_delta → 退而认 thinking_delta(模型开始产出)兑底。 // start 仅标志流建立, 不算 TTFT。 if (ttftMs == null) { ttftMs = Date.now() - startTs; logLine(null, `${logTag} first token (${ev.type === "thinking_delta" ? "thinking" : "text"}, ttft=${ttftMs}ms)`); } reply += (ev as { delta?: string }).delta || ""; } else if (ev.type === "done") { const msg = (ev as unknown as { message: AssistantMessage }).message; logLine(null, `${logTag} done: stopReason=${msg.stopReason} latency=${Date.now() - startTs}ms`); if (msg.stopReason === "error" || msg.stopReason === "aborted") { return { ok: false, error: msg.errorMessage || msg.stopReason || "aborted", ttftMs }; } return { ok: true, ttftMs, latencyMs: Date.now() - startTs, reply: reply.slice(0, 200) }; } else if (ev.type === "error") { const msg = (ev as unknown as { error: AssistantMessage }).error; logLine(null, `${logTag} ERROR event: ${msg?.errorMessage}`); return { ok: false, error: msg?.errorMessage || "stream error", ttftMs }; } } return { ok: false, error: "流结束但无 done 事件", ttftMs }; } catch (e) { const msg = e instanceof Error ? e.message : String(e); // 超时判定用 ctrl.signal.aborted 独立标志(timer 触发 abort), 不匹配文案 —— // 旧版 /aborted/i 会把上游 "request aborted by upstream" 错改写为超时, 污染 error。 // 用户 Esc(opts.signal.aborted) 联动 abort 也非超时, 保留原语义。 const isTimeout = ctrl.signal.aborted && !opts.signal?.aborted; const isUserAbort = opts.signal?.aborted === true; const shown = isTimeout ? `超时(${timeoutMs / 1000}s)` : isUserAbort ? "Request aborted" : msg; logLine(null, `${logTag} CATCH: ${shown}`); return { ok: false, error: shown, ttftMs }; } finally { clearTimeout(timer); if (opts.signal) opts.signal.removeEventListener("abort", onOuterAbort); } } /** 面板测试模型: find → 构 ResolvedCandidate → 复用 testCandidateSpeed(30s 超时)。 */ async function testModel(provider: string, model: string, prompt: string): Promise { logLine(null, `testModel START ${provider}/${model} prompt=${JSON.stringify(prompt).slice(0, 40)}`); if (!stateModelRegistry) { logLine(null, "testModel FAIL: registry 未初始化"); return { ok: false, error: "扩展未初始化" }; } // 刷新 registry 以读到面板新保存的模型(否则 session_start 后新增的模型 find 不到) stateModelRegistry.refresh(); const m = stateModelRegistry.find(provider, model); if (!m) { logLine(null, `testModel FAIL: ${provider}/${model} 未在 registry 找到`); return { ok: false, error: `模型 ${provider}/${model} 未在已加载配置中找到(若刚改配置请先保存)` }; } logLine(null, `testModel found: api=${(m as { api?: string }).api} baseUrl=${(m as { baseUrl?: string }).baseUrl}`); return testCandidateSpeed( { ...m, roundrobinLabel: `${provider}/${model}` }, { prompt, timeoutMs: 30000 }, ); } // ============================================================ // SpeedTest: 对组内所有候选并发测速 + 排序 + 落 health // - 成功: recordSuccess → 清 lastFailAt(出冷却), 按 sortKey 升序前排 // - 失败: recordFailure → 进冷却, 排末尾 // - 排序后写回 group.candidates, currentIndex=0 // - 通过 notify/log/lastSpeedTestResults 暴露状态给 TUI 和面板 // ============================================================ /** 并发限流: chunked Promise.allSettled, 每批 concurrency 个。 */ async function runSpeedTest( group: Group, signal?: AbortSignal, ): Promise<{ okCount: number; badCount: number }> { const st = group.config.speedTest; const prompt = st?.prompt || "欧拉函数的意义?"; const timeoutMs = st?.timeoutMs ?? 60000; const concurrency = Math.max(1, st?.concurrency ?? 5); const sortKey = st?.sortKey ?? "ttft"; const total = group.candidates.length; const results: { candidate: ResolvedCandidate; ok: boolean; ttftMs: number | null; latencyMs: number | null; error?: string }[] = []; let nextIdx = 0; const workers = Array.from({ length: Math.min(concurrency, total) }, async () => { while (true) { if (signal?.aborted) break; const i = nextIdx++; if (i >= total) break; const candidate = group.candidates[i]; logLine(group, `speedtest ${candidate.roundrobinLabel} (#${i + 1}/${total})`); const r = await testCandidateSpeed(candidate, { prompt, timeoutMs, signal, retries: st?.retries }, `speedtest[${group.name}]`); results.push({ candidate, ok: r.ok, ttftMs: r.ttftMs ?? null, latencyMs: r.latencyMs ?? null, error: r.error, }); // 落 health.jsonl(与真实请求路径同一份统计源, 打 probe 标区分—— // smart 的可靠性维度只看真实请求, 避免测速探活自打分污染) const h = getHealth(group, candidate); if (r.ok) { recordSuccess(group, candidate, r.latencyMs ?? 0, r.ttftMs ?? null, true); h.lastFailAt = null; // 测速成功 = 出冷却, 立即可用 } else { recordFailure(group, candidate, true); } } }); await Promise.all(workers); // 排序: 成功的按 sortKey 升序在前, 失败的排后(稳定顺序按 label) const ok = results.filter((r) => r.ok); const bad = results.filter((r) => !r.ok); // 未测候选(worker abort 时 break 跳过): 保留在原数组, 按原序排末尾, 不丢失。 // 否则 sorted 从 results 重建会永久删掉未测候选, 后续轮询再也轮不到(数据破坏)。 const testedLabels = new Set(results.map((r) => r.candidate.roundrobinLabel)); const untested = group.candidates.filter((c) => !testedLabels.has(c.roundrobinLabel)); // 中位数兑底: ttft/latency 缺失(null)的 ok 候选, 用组内实测中位数参与排序 // (语义: 该维度未知 → 假设中等水平; 比直接垫底 Infinity 合理——reasoning 模型 // 测速可能拿不到 text 首字, 但它整体可用, 不该被当作最慢)。 // 采样取全组(results, 含 bad 段): bad 候选只要首字拿到了, 其 ttft 同样是 // “组内典型首字水平”的有效样本——ok 段全 null 而 bad 段有值时, 旧版 // median=[] 兜底失败退 Infinity, 唯一能跑完的 ok 候选反被排尾。 // 只影响排序, 不改写持久化 ttftMs/latencyMs(health.jsonl 仍存真实 null)。 const median = (arr: number[]): number | null => { if (!arr.length) return null; const s = [...arr].sort((a, b) => a - b); const mid = Math.floor(s.length / 2); return s.length % 2 ? s[mid] : (s[mid - 1] + s[mid]) / 2; }; const ttftMedian = median(results.map((r) => r.ttftMs).filter((v): v is number => typeof v === "number")); const latMedian = median(results.map((r) => r.latencyMs).filter((v): v is number => typeof v === "number")); const tOf = (r: (typeof results)[number]): number => r.ttftMs ?? ttftMedian ?? Infinity; const lOf = (r: (typeof results)[number]): number => r.latencyMs ?? latMedian ?? Infinity; // smart 三维加权: 0.5×ttft_norm + 0.3×(1-reliability) + 0.2×latency_norm // ttft/latency 用 log 尺度归一化(组内, 0=最快): 纯 min-max 对离群值过敏—— // 一个 5s 拖尾候选会把 200ms/250ms 的差距压扁成 0.01, ttft 维度几乎失效; // log 尺度下同一组差距恢复到可分辨量级(长尾是指数型的, log 是共轭尺度)。 // reliability 用贝叶斯平滑 (success + k×0.5)/(total + k), k=5 先验强度: // 无历史→0.5中性; 样本不足时往中性拉回; 已知全败(0/10)→0.167 监底。 // 数据源: health.jsonl 7天 realOnly 聚合(排除 probe 探活事件, 跨进程共享, // 重启不丢) — 不再用进程内 healthStats(那本账混入了测速探活, 且重启清零)。 // 加权后分数越小越好(0=理想), 升序排序。失败候选仍 sink 到末尾。 let smartVal: ((r: (typeof results)[number]) => number) | null = null; if (sortKey === "smart") { const logNorm = (vals: number[]): ((v: number) => number) => { const finite = vals.filter((v) => Number.isFinite(v) && v > 0); if (!finite.length) return () => 0; // 全未知/非正 → 维度不区分, 全 0 const logs = finite.map((v) => Math.log(v)); const lo = Math.min(...logs); const hi = Math.max(...logs); // 近等值(hi/lo 差距 <25%)当不区分: 速度相仿的候选组(常见于同模型多渠道) // 不让微小抖动被 log 放大压过 reliability 等大跨度维度。log 只对指数级长尾有意义。 if (!(hi > lo) || hi - lo < Math.abs(lo) * 0.25 + 0.001) return () => 0; return (v: number): number => (Number.isFinite(v) && v > 0 ? (Math.log(v) - lo) / (hi - lo) : 1); }; const tNorm = logNorm(ok.map(tOf)); const lNorm = logNorm(ok.map(lOf)); // 可靠性: 一次读盘, 组内按 provider/model 查(7天 realOnly) const realAgg = aggregateEvents(loadEvents(), { realOnly: true }); const BETA_K = 5; // 贝叶斯先验强度: 相当于预设 5 次 50% 成功率 smartVal = (r): number => { const t = tNorm(tOf(r)); const l = lNorm(lOf(r)); const agg = realAgg[providerHealthKey(r.candidate.provider, r.candidate.id)]; const succ = agg?.success ?? 0; const fail = agg?.fail ?? 0; const rel = (succ + BETA_K * 0.5) / (succ + fail + BETA_K); // 贝叶斯平滑 return 0.5 * t + 0.3 * (1 - rel) + 0.2 * l; }; } const sortVal = (r: (typeof results)[number]): number => sortKey === "latency" ? lOf(r) : sortKey === "hybrid" ? (0.7 * tOf(r) + 0.3 * lOf(r)) : sortKey === "smart" ? (smartVal ? smartVal(r) : tOf(r)) : tOf(r); ok.sort((a, b) => sortVal(a) - sortVal(b)); // bad 段二级排序: 先按“离成功多近”排——拿到首字的(b/c类: 首字快但中途断/吐完无done) // 排在连 start 都没拿到的(a类: 鉴权失败/全程超时)之前, 再按首字快慢, 最后 label 兑底。 // 已采集的 ttftMs 不再丢给纯字母序(与 median 扩采样同一动机: 数据采了就用)。 bad.sort((a, b) => { const at = a.ttftMs ?? null; const bt = b.ttftMs ?? null; if (at === null && bt === null) return a.candidate.roundrobinLabel.localeCompare(b.candidate.roundrobinLabel); if (at === null) return 1; // a 无首字垫后 if (bt === null) return -1; // b 无首字垫后 if (at !== bt) return at - bt; // 都有首字: 快者在前 return a.candidate.roundrobinLabel.localeCompare(b.candidate.roundrobinLabel); }); const sorted = [...ok.map((r) => r.candidate), ...bad.map((r) => r.candidate), ...untested]; group.candidates = sorted; group.currentIndex = 0; group.lastSpeedTestAt = Date.now(); group.lastSpeedTestResults = new Map( results.map((r) => [r.candidate.roundrobinLabel, { ok: r.ok, ttftMs: r.ttftMs, latencyMs: r.latencyMs, error: r.error }]), ); // 构造排序摘要日志 (smart 值为 0~1 加权分, 非 ms) const valFmt = (r: (typeof results)[number]): string => sortKey === "smart" ? (smartVal ? smartVal(r).toFixed(3) : "?") : `${sortVal(r)}ms`; const summary = [...ok.map((r) => `${r.candidate.roundrobinLabel}:${valFmt(r)}`), ...bad.map((r) => `${r.candidate.roundrobinLabel}:✗( ${r.error || "fail"})`)].join(" "); logLine(group, `speedtest done ok=${ok.length}/${total} order=[${summary}]`); return { okCount: ok.length, badCount: bad.length }; } /** 主动测速: 对所有 enabled 且 speedTest.enabled=true 的组, 绕过节流立即重测。 * 供 /rr-speedtest 命令和面板按钮调用。返回逐组结果摘要。 */ async function manualSpeedTest(groupName?: string): Promise< Array<{ name: string; ok: boolean; okCount: number; total: number; error?: string }> > { const results: Array<{ name: string; ok: boolean; okCount: number; total: number; error?: string }> = []; for (const g of groups.values()) { if (groupName && g.name !== groupName) continue; if (!g.enabled) continue; if (g.config.speedTest?.enabled !== true) { results.push({ name: g.name, ok: false, okCount: 0, total: g.candidates.length, error: "测速未开启" }); continue; } if (isSpeedTestLocked(g)) { results.push({ name: g.name, ok: false, okCount: 0, total: g.candidates.length, error: "测速进行中" }); continue; } lockSpeedTest(g); try { notifyUi(`[${g.name}] 手动测速排序中…`, "info"); const start = Date.now(); const { okCount } = await runSpeedTest(g, undefined); const total = g.candidates.length; notifyUi( `[${g.name}] 手动测速完成 ${Date.now() - start}ms: 可用 ${okCount}/${total}`, okCount > 0 ? "info" : "error", ); results.push({ name: g.name, ok: okCount > 0, okCount, total }); } catch (e) { const msg = e instanceof Error ? e.message : String(e); logLine(g, `manual speedtest failed: ${msg}`); results.push({ name: g.name, ok: false, okCount: 0, total: g.candidates.length, error: msg }); } finally { unlockSpeedTest(g); } } return results; } // ============================================================ // Extension entry // ============================================================ export default function (pi: ExtensionAPI) { statePi = pi; // 启动期先注册(此时 modelRegistry 还没就绪,candidates 为空,会注册空 provider; // session_start 会真正 build groups 并重新注册) registerAllVirtualProviders(pi); pi.on("session_start", async (_event, ctx) => { stateModelRegistry = ctx.modelRegistry; stateUi = ctx.ui; try { // 剪掉 7 天外旧事件,控制 health.jsonl 体积(跨 CLI 共享) try { pruneHealthFile(); } catch { // ignore prune errors } const allConfigs = loadAllGroupsConfig(); buildGroups(allConfigs); registerAllVirtualProviders(pi); const enabledNames = [...groups.values()].filter((g) => g.enabled).map((g) => g.name); logLine(null, `session_start groups=${groups.size} enabled=${enabledNames.join(",") || "(none)"}`); // 初始测速改为懒触发:首次请求某组时才测(见 streamRoundRobin),session_start 不测 } catch (e) { const msg = e instanceof Error ? e.message : String(e); logLine(null, `session_start failed: ${msg}`); ctx.ui.notify(`roundrobin 启用失败: ${msg}`, "error"); } }); // 面板保存 rr 配置 / 预设后触发热重载 pi.events.on(CONFIG_CHANGED_EVENT, () => reloadRuntime()); // ---- Provider 全量健康统计: 采集 TTFT + 成功率 + 总延迟 ---- // TTFT: 首个 text_delta 到达时刻(真正首字), 按 provider/model 记首字。 // 注: pi 层 message_start 是 AgentMessage 创建时刻(≈请求发出), 用它算 TTFT 恒为 0; // 真正首字是 message_update 的 assistantMessageEvent.type === "text_delta"。 pi.on("message_update", (event, _ctx) => { const msg = event.message; if (msg.role !== "assistant") return; if (msg.provider === PROVIDER) return; // 跳过 roundrobin 虚拟模型(候选健康由 tryCandidates 直接落盘) const ame = event.assistantMessageEvent as { type?: string } | undefined; if (!ame || ame.type !== "text_delta") return; const key = providerHealthKey(msg.provider, msg.model); const ts = typeof msg.timestamp === "number" ? msg.timestamp : 0; const cur = firstTokenAt.get(key); // 首个 text_delta; 若上次请求残留(msgTs 不同)则覆盖 if (!cur || cur.msgTs !== ts) { firstTokenAt.set(key, { ts: Date.now(), msgTs: ts }); } }); // latency = 结束时刻 - 请求开始(message.timestamp)。不依赖队列配对, 无泄漏/错配。 pi.on("message_end", (event, _ctx) => { const msg = event.message; if (msg.role !== "assistant") return; if (msg.provider === PROVIDER) return; const key = providerHealthKey(msg.provider, msg.model); const isFail = msg.stopReason === "error" || msg.stopReason === "aborted"; const startTs = typeof msg.timestamp === "number" ? msg.timestamp : null; if (isFail) { recordProviderFailure(msg.provider, msg.model); } else { const latency = startTs != null ? Math.max(0, Date.now() - startTs) : null; const cur = firstTokenAt.get(key); const ttft = cur && startTs != null && cur.msgTs === startTs ? Math.max(0, cur.ts - startTs) : null; recordProviderSuccess(msg.provider, msg.model, latency ?? 0, ttft); } firstTokenAt.delete(key); }); // 暴露运行时统计给面板(只读) (globalThis as unknown as Record).__roundrobinGetHealth = (): unknown => getHealthSnapshot(); (globalThis as unknown as Record).__providerHealthGet = (): unknown => getProviderHealthSnapshot(); (globalThis as unknown as Record).__providerHealthReset = (): void => { resetProviderHealth(); }; // 暴露模型测试给面板 (globalThis as unknown as Record).__testModel = (provider: string, model: string, prompt: string): Promise => testModel(provider, model, prompt); // 暴露手动测速给面板 (绕过节流立即重测) (globalThis as unknown as Record).__rrManualSpeedTest = (groupName?: string): Promise => manualSpeedTest(groupName); // 主动测速命令: /rr-speedtest [组名] — 随时重新测速排序 // 测试环境可能未提供 registerCommand, 用可选链防护 pi.registerCommand?.("rr-speedtest", { description: "手动触发轮询组测速排序 (绕过节流, 立即重测)", getArgumentCompletions: (argumentPrefix: string) => [...groups.values()] .filter((g) => g.enabled && g.config.speedTest?.enabled === true && g.name.startsWith(argumentPrefix)) .map((g) => ({ value: g.name, label: g.name, description: `${g.candidates.length} 个候选`, })), handler: async (args, ctx) => { const groupName = args.trim() || undefined; const targets = [...groups.values()].filter( (g) => (!groupName || g.name === groupName) && g.enabled && g.config.speedTest?.enabled === true, ); if (targets.length === 0) { ctx.ui.notify( groupName ? `组 "${groupName}" 不存在或未开启测速 (在面板勾选 speedTest.enabled)` : "没有开启测速的组 (在面板勾选 speedTest.enabled)", "error", ); return; } ctx.ui.notify(`测速开始: ${targets.map((g) => g.name).join(", ")}`, "info"); const results = await manualSpeedTest(groupName); // 详细结果: 每个候选的 TTFT/延迟/状态, 按排序顺序, 显示在编辑器上方 widget。 // manualSpeedTest 结束后 group.lastSpeedTestResults 已更新, group.candidates 已重排。 const WIDGET_KEY = "rr-speedtest-result"; const lines: string[] = []; for (const r of results) { const g = groups.get(r.name); const resMap = g?.lastSpeedTestResults; const sk = g?.config.speedTest?.sortKey; lines.push(`[${r.name}] ${r.ok ? `可用 ${r.okCount}/${r.total}` : `失败: ${r.error || `${r.okCount}/${r.total}`}`}${sk ? ` sortKey=${sk}` : ""}`); if (g && resMap && resMap.size > 0) { g.candidates.forEach((c, i) => { const sr = resMap.get(c.roundrobinLabel); if (!sr) return; const ttft = sr.ttftMs != null ? `${(sr.ttftMs / 1000).toFixed(1)}s` : "-"; const lat = sr.latencyMs != null ? `${(sr.latencyMs / 1000).toFixed(1)}s` : "-"; const rank = sr.ok ? `#${i + 1}` : ` ✗`; const tail = sr.ok ? `TTFT ${ttft} 延迟 ${lat}` : `失败: ${sr.error || "unknown"}`; lines.push(` ${rank} ${c.roundrobinLabel} ${tail}`); }); } lines.push(""); } lines.push("(按排序顺序; #可用 / ✗不可用; 30s 后自动消失)"); try { ctx.ui.setWidget(WIDGET_KEY, lines, { placement: "aboveEditor" }); } catch { // 非 TUI 模式(rpc/json/print) 可能无 widget, 忽略 } if (speedtestWidgetTimer) clearTimeout(speedtestWidgetTimer); speedtestWidgetTimer = setTimeout(() => { try { ctx.ui.setWidget(WIDGET_KEY, undefined); } catch { /* ignore */ } speedtestWidgetTimer = null; }, 30000); for (const r of results) { ctx.ui.notify( `[${r.name}] ${r.ok ? `可用 ${r.okCount}/${r.total}(详细见上方)` : `失败: ${r.error || `${r.okCount}/${r.total}`}`}`, r.ok ? "info" : "error", ); } }, }); }