From d01724556a0e5c9a3a8f562384d5bb0ac463359c Mon Sep 17 00:00:00 2001 From: decolua Date: Thu, 27 Aug 2026 18:39:44 +0700 Subject: [PATCH] fix(models): drop the worker thread from the catalog sync MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The worker resolved its own path through import.meta.url, which the bundler rewrites — so the running server looked for the file at a path that does not exist there: [modelCatalog] sync failed: Cannot find module '/Users/Working/router4/9router/src/lib/modelCatalog/worker.js' It was guarding against a 23ms JSON.parse that runs once a day, 60s after boot. Inlining it into sync.js costs that 23ms on an otherwise idle tick and removes both the failure mode and a whole file. Co-Authored-By: Claude Fable 5 --- src/lib/modelCatalog/sync.js | 191 +++++++++++++++++++++++---------- src/lib/modelCatalog/worker.js | 132 ----------------------- 2 files changed, 137 insertions(+), 186 deletions(-) delete mode 100644 src/lib/modelCatalog/worker.js diff --git a/src/lib/modelCatalog/sync.js b/src/lib/modelCatalog/sync.js index d30adc71..ff70d092 100644 --- a/src/lib/modelCatalog/sync.js +++ b/src/lib/modelCatalog/sync.js @@ -1,19 +1,30 @@ -// Background refresh of model capabilities from models.dev. +// Daily refresh of model capabilities from models.dev. // -// Failures are swallowed on purpose: a stale or missing catalog just means the -// hand-written capability tables keep deciding on their own. +// Downloads the catalog, keeps only what differs from the hand-written tables, +// and writes it next to the database. Failures are swallowed on purpose: a +// stale or missing file just means those tables keep deciding on their own. +import fs from "node:fs"; import path from "node:path"; -import { Worker } from "node:worker_threads"; -import { fileURLToPath } from "node:url"; import { CATALOG_FILE, CATALOG_RAW_FILE, invalidateCatalog, installCatalogSource } from "open-sse/providers/catalogOverride.js"; const CATALOG_URL = "https://models.dev/api.json"; -const WORKER_FILE = path.join(path.dirname(fileURLToPath(import.meta.url)), "worker.js"); +const FETCH_TIMEOUT_MS = 60000; + +export const SYNC_INTERVAL_MS = 24 * 60 * 60 * 1000; +const STARTUP_DELAY_MS = 60 * 1000; // let the server boot and serve first requests +const RETRY_DELAY_MS = 30 * 60 * 1000; + +const MODALITY_BY_INPUT = { image: "vision", pdf: "pdf", audio: "audioInput", video: "videoInput" }; +// Gateways disagree about the same model, so a modality needs a majority of +// them to declare it — one reseller mislabelling a text model must not win. +const MIN_SHARE = 0.5; +// Ignore limit differences below this: gateways round 200000 vs 202752. +const LIMIT_TOLERANCE = 0.1; // 9router provider id -> models.dev provider id, for context/maxOutput only. -// Providers absent here keep whatever the local pattern table resolves; the -// names that already match are resolved automatically. +// Providers absent here keep whatever the local pattern table resolves; names +// that already match are resolved automatically. const PROVIDER_ALIASES = { "glm": "zai", "glm-cn": "zhipuai", @@ -29,11 +40,6 @@ const PROVIDER_ALIASES = { "cloudflare-ai": "cloudflare-workers-ai", }; -export const SYNC_INTERVAL_MS = 24 * 60 * 60 * 1000; -const STARTUP_DELAY_MS = 60 * 1000; // let the server boot and serve first requests -const RETRY_DELAY_MS = 30 * 60 * 1000; -const WORKER_TIMEOUT_MS = 120000; - let state = { running: false, lastSync: null, lastError: null, lastResult: null, etag: null }; let timer = null; @@ -41,9 +47,95 @@ export function getSyncState() { return { ...state, file: CATALOG_FILE, url: CATALOG_URL, intervalMs: SYNC_INTERVAL_MS }; } +// "zai-org/GLM-4.6V:free" -> "glm-4.6v" +function baseId(modelId) { + const withoutVendor = modelId.includes("/") ? modelId.split("/").pop() : modelId; + return withoutVendor.toLowerCase().split(":")[0]; +} + +function writeAtomic(file, contents) { + fs.mkdirSync(path.dirname(file), { recursive: true }); + fs.writeFileSync(`${file}.tmp`, contents, "utf8"); + fs.renameSync(`${file}.tmp`, file); +} + +// Trimmed copy of the upstream catalog, kept for the add-models skill: same +// models, ~470KB instead of 4.3MB. +function slim(catalog) { + const out = {}; + for (const [providerId, provider] of Object.entries(catalog)) { + const models = {}; + for (const [modelId, model] of Object.entries(provider?.models || {})) { + models[modelId] = { + i: (model?.modalities?.input || []).filter((x) => x !== "text"), + c: model?.limit?.context, + o: model?.limit?.output, + r: model?.reasoning || undefined, + }; + } + out[providerId] = models; + } + return out; +} + +function build(catalog, entries) { + // Index once: per provider for limits, and tallied across all of them for + // modalities. + const byProvider = {}; + const tally = {}; + for (const [providerId, provider] of Object.entries(catalog)) { + const models = {}; + for (const [modelId, model] of Object.entries(provider?.models || {})) { + const id = baseId(modelId); + models[id] = model; + const counts = tally[id] || (tally[id] = { total: 0 }); + counts.total++; + for (const input of model?.modalities?.input || []) { + const key = MODALITY_BY_INPUT[input]; + if (key) counts[key] = (counts[key] || 0) + 1; + } + } + byProvider[providerId] = models; + } + + // Modalities belong to the model — every gateway serving it has the same + // weights — so they are keyed by model id and shared across providers. + const models = {}; + for (const [id, counts] of Object.entries(tally)) { + const declared = {}; + for (const key of Object.values(MODALITY_BY_INPUT)) { + if ((counts[key] || 0) / counts.total >= MIN_SHARE) declared[key] = true; + } + if (Object.keys(declared).length) models[id] = declared; + } + + // Limits belong to the gateway — each truncates differently — so only the + // matching provider's own numbers are used, keyed by provider + model. + const providers = {}; + for (const { provider, model, contextLength, current } of entries) { + const alias = PROVIDER_ALIASES[provider]; + const upstream = catalog[provider] ? provider : (alias && catalog[alias] ? alias : null); + const entry = upstream && byProvider[upstream]?.[baseId(model)]; + if (!entry) continue; + + const delta = {}; + const { context, output } = entry.limit || {}; + if (context > 0 && !contextLength + && Math.abs(context - current.contextWindow) / current.contextWindow > LIMIT_TOLERANCE) { + delta.contextWindow = context; + } + if (output > 0 + && Math.abs(output - current.maxOutput) / current.maxOutput > LIMIT_TOLERANCE) { + delta.maxOutput = output; + } + if (Object.keys(delta).length) (providers[provider] || (providers[provider] = {}))[model] = delta; + } + + return { models, providers }; +} + // Snapshot every registered model with its currently resolved capabilities, so -// the worker can compute a delta without importing app modules (it cannot -// resolve the bundler-only "open-sse/*" alias). +// build() can tell which upstream values are actually a change. async function collectEntries() { const [{ default: registry }, { getCapabilitiesForModel }] = await Promise.all([ import("open-sse/providers/registry/index.js"), @@ -65,52 +157,43 @@ async function collectEntries() { return entries; } -function runWorker(entries) { - return new Promise((resolve, reject) => { - const worker = new Worker(WORKER_FILE, { - workerData: { - url: CATALOG_URL, - etag: state.etag, - outFile: CATALOG_FILE, - rawFile: CATALOG_RAW_FILE, - entries, - providerAliases: PROVIDER_ALIASES, - }, - resourceLimits: { maxOldGenerationSizeMb: 512 }, - }); - - let settled = false; - const finish = (fn, value) => { - if (settled) return; - settled = true; - clearTimeout(timeout); - fn(value); - }; - const timeout = setTimeout(() => { - worker.terminate(); - finish(reject, new Error("sync timed out")); - }, WORKER_TIMEOUT_MS); - - worker.on("message", (msg) => { - if (msg?.ok) finish(resolve, msg.result); - else finish(reject, new Error(msg?.error || "sync failed")); - }); - worker.on("error", (err) => finish(reject, err)); - worker.on("exit", (code) => finish(reject, new Error(`worker exited with ${code}`))); - }); -} - -// Run one sync. Returns the worker summary, or null when it could not complete. +// Run one sync. Returns a summary, or null when it could not complete. export async function syncModelCatalog() { if (state.running) return null; state.running = true; try { - const result = await runWorker(await collectEntries()); - if (result.status === "updated") { - state.etag = result.etag; + const headers = { accept: "application/json" }; + if (state.etag) headers["if-none-match"] = state.etag; + const response = await fetch(CATALOG_URL, { headers, signal: AbortSignal.timeout(FETCH_TIMEOUT_MS) }); + + let result; + if (response.status === 304) { + result = { status: "unchanged" }; + } else if (!response.ok) { + throw new Error(`HTTP ${response.status}`); + } else { + // ~23ms to parse, once a day, on a server that is otherwise idle at this + // point — not worth a worker thread. + const catalog = await response.json(); + const etag = response.headers.get("etag") || null; + const { models, providers } = build(catalog, await collectEntries()); + const serialized = JSON.stringify({ v: 1, etag, syncedAt: Date.now(), models, providers }); + + writeAtomic(CATALOG_FILE, serialized); + writeAtomic(CATALOG_RAW_FILE, JSON.stringify(slim(catalog))); + + state.etag = etag; invalidateCatalog(); + result = { + status: "updated", + etag, + bytes: Buffer.byteLength(serialized), + models: Object.keys(models).length, + providers: Object.keys(providers).length, + }; console.log(`[modelCatalog] ${result.models} models, ${result.providers} providers, ${(result.bytes / 1024).toFixed(1)}KB`); } + state.lastSync = Date.now(); state.lastError = null; state.lastResult = result; diff --git a/src/lib/modelCatalog/worker.js b/src/lib/modelCatalog/worker.js deleted file mode 100644 index 39629ea0..00000000 --- a/src/lib/modelCatalog/worker.js +++ /dev/null @@ -1,132 +0,0 @@ -// Downloads models.dev and writes the capability deltas 9router reads. -// Runs in a worker thread: the 4MB parse would otherwise block requests for ~20ms. - -import { parentPort, workerData } from "node:worker_threads"; -import fs from "node:fs"; -import path from "node:path"; - -const FETCH_TIMEOUT_MS = 60000; -const MODALITY_BY_INPUT = { image: "vision", pdf: "pdf", audio: "audioInput", video: "videoInput" }; -// Gateways disagree about the same model, so a modality needs a majority of -// them to declare it — one reseller mislabelling a text model must not win. -const MIN_SHARE = 0.5; -// Ignore limit differences below this: gateways round 200000 vs 202752. -const LIMIT_TOLERANCE = 0.1; - -// "zai-org/GLM-4.6V:free" -> "glm-4.6v" -function baseId(modelId) { - const withoutVendor = modelId.includes("/") ? modelId.split("/").pop() : modelId; - return withoutVendor.toLowerCase().split(":")[0]; -} - -function build(catalog, entries, providerAliases) { - // Index the catalog once: per provider for limits, and tallied for modalities. - const byProvider = {}; - const tally = {}; - for (const [providerId, provider] of Object.entries(catalog)) { - const models = {}; - for (const [modelId, model] of Object.entries(provider?.models || {})) { - const id = baseId(modelId); - models[id] = model; - const counts = tally[id] || (tally[id] = { total: 0 }); - counts.total++; - for (const input of model?.modalities?.input || []) { - const key = MODALITY_BY_INPUT[input]; - if (key) counts[key] = (counts[key] || 0) + 1; - } - } - byProvider[providerId] = models; - } - - // Modalities belong to the model — every gateway serving it has the same - // weights — so they are keyed by model id and shared across providers. - const models = {}; - for (const [id, counts] of Object.entries(tally)) { - const declared = {}; - for (const key of Object.values(MODALITY_BY_INPUT)) { - if ((counts[key] || 0) / counts.total >= MIN_SHARE) declared[key] = true; - } - if (Object.keys(declared).length) models[id] = declared; - } - - // Limits belong to the gateway — each truncates differently — so only the - // matching provider's own numbers are used, keyed by provider + model. - const providers = {}; - for (const { provider, model, contextLength, current } of entries) { - const alias = providerAliases[provider]; - const upstream = catalog[provider] ? provider : (alias && catalog[alias] ? alias : null); - const entry = upstream && byProvider[upstream]?.[baseId(model)]; - if (!entry) continue; - - const delta = {}; - const { context, output } = entry.limit || {}; - if (context > 0 && !contextLength - && Math.abs(context - current.contextWindow) / current.contextWindow > LIMIT_TOLERANCE) { - delta.contextWindow = context; - } - if (output > 0 - && Math.abs(output - current.maxOutput) / current.maxOutput > LIMIT_TOLERANCE) { - delta.maxOutput = output; - } - if (Object.keys(delta).length) (providers[provider] || (providers[provider] = {}))[model] = delta; - } - - return { models, providers }; -} - -// Trimmed copy of the upstream catalog, kept for the add-models skill: same -// 7348 models, 470KB instead of 4.3MB, so a scan reads it in ~5ms. -function slim(catalog) { - const out = {}; - for (const [providerId, provider] of Object.entries(catalog)) { - const models = {}; - for (const [modelId, model] of Object.entries(provider?.models || {})) { - models[modelId] = { - i: (model?.modalities?.input || []).filter((x) => x !== "text"), - c: model?.limit?.context, - o: model?.limit?.output, - r: model?.reasoning || undefined, - }; - } - out[providerId] = models; - } - return out; -} - -function writeAtomic(file, contents) { - fs.mkdirSync(path.dirname(file), { recursive: true }); - fs.writeFileSync(`${file}.tmp`, contents, "utf8"); - fs.renameSync(`${file}.tmp`, file); -} - -async function run() { - const { url, etag, outFile, rawFile, entries, providerAliases } = workerData; - - const headers = { accept: "application/json" }; - if (etag) headers["if-none-match"] = etag; - const response = await fetch(url, { headers, signal: AbortSignal.timeout(FETCH_TIMEOUT_MS) }); - - if (response.status === 304) return { status: "unchanged" }; - if (!response.ok) throw new Error(`HTTP ${response.status}`); - - const catalog = await response.json(); - const nextEtag = response.headers.get("etag") || null; - const { models, providers } = build(catalog, entries, providerAliases); - const serialized = JSON.stringify({ v: 1, etag: nextEtag, syncedAt: Date.now(), models, providers }); - - writeAtomic(outFile, serialized); - if (rawFile) writeAtomic(rawFile, JSON.stringify(slim(catalog))); - - return { - status: "updated", - etag: nextEtag, - bytes: Buffer.byteLength(serialized), - models: Object.keys(models).length, - providers: Object.keys(providers).length, - }; -} - -run().then( - (result) => parentPort?.postMessage({ ok: true, result }), - (error) => parentPort?.postMessage({ ok: false, error: error?.message || String(error) }) -);