fix(models): drop the worker thread from the catalog sync
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 <noreply@anthropic.com>
This commit is contained in:
@@ -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;
|
||||
|
||||
@@ -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) })
|
||||
);
|
||||
Reference in New Issue
Block a user