mirror of
https://github.com/diegosouzapw/OmniRoute.git
synced 2026-08-22 07:02:16 +03:00
fix(embeddings): cool down account on hard errors (402/401/5xx) (#10529)
Merged — locally validated (focused embedding-cooldown tests green, gates green). Good parity with the chat path's existing pattern. Thanks!
This commit is contained in:
@@ -1073,11 +1073,6 @@
|
||||
"count": 1
|
||||
}
|
||||
},
|
||||
"src/lib/embeddings/service.ts": {
|
||||
"no-restricted-imports": {
|
||||
"count": 1
|
||||
}
|
||||
},
|
||||
"src/lib/evals/runtime.ts": {
|
||||
"no-restricted-imports": {
|
||||
"count": 1
|
||||
|
||||
@@ -10,13 +10,14 @@ import { errorResponse, unavailableResponse } from "@omniroute/open-sse/utils/er
|
||||
import { HTTP_STATUS } from "@omniroute/open-sse/config/constants.ts";
|
||||
import * as log from "@/sse/utils/logger";
|
||||
import { toJsonErrorPayload } from "@/shared/utils/upstreamError";
|
||||
import { getProviderCredentials, clearRecoveredProviderState } from "@/sse/services/auth";
|
||||
import {
|
||||
getCachedProviderNodes,
|
||||
getComboByName,
|
||||
getCombos,
|
||||
getDatabaseSettings,
|
||||
} from "@/lib/localDb";
|
||||
getProviderCredentials,
|
||||
clearRecoveredProviderState,
|
||||
markAccountUnavailable,
|
||||
} from "@/sse/services/auth";
|
||||
import { getCachedProviderNodes } from "@/lib/db/readCache";
|
||||
import { getComboByName, getCombos } from "@/lib/db/combos";
|
||||
import { getDatabaseSettings } from "@/lib/db/databaseSettings";
|
||||
import { resolveProxyForConnection } from "@/lib/db/settings";
|
||||
import { runWithProxyContext } from "@omniroute/open-sse/utils/proxyFetch.ts";
|
||||
import { handleComboChat } from "@omniroute/open-sse/services/combo.ts";
|
||||
@@ -309,7 +310,7 @@ export async function createEmbeddingResponse(
|
||||
// #10347 — thread the selected connection id so handleEmbedding can cool the
|
||||
// account on a hard upstream failure (previously always null on /v1/embeddings).
|
||||
connectionId:
|
||||
((credentials as { connectionId?: string } | null)?.connectionId) ||
|
||||
(credentials as { connectionId?: string } | null)?.connectionId ||
|
||||
options.connectionId ||
|
||||
connectionIdForProxy ||
|
||||
null,
|
||||
@@ -340,6 +341,33 @@ export async function createEmbeddingResponse(
|
||||
});
|
||||
}
|
||||
|
||||
// #10347: cool down the account on hard errors (402 subscription expired,
|
||||
// 401 revoked, 403 forbidden, 404 model gone, 429 rate limit, 5xx server
|
||||
// errors) so the next embedding request skips this account. Mirrors chat.ts
|
||||
// behavior.
|
||||
// Skip for 400 (bad request) — the account is fine, the request was wrong.
|
||||
// Best-effort: don't block the error response on the DB write.
|
||||
const HARD_ERROR_STATUSES = new Set([401, 402, 403, 404, 429, 500, 502, 503, 504]);
|
||||
if (
|
||||
credentials &&
|
||||
"connectionId" in credentials &&
|
||||
typeof credentials.connectionId === "string" &&
|
||||
HARD_ERROR_STATUSES.has(result.status)
|
||||
) {
|
||||
markAccountUnavailable(
|
||||
credentials.connectionId,
|
||||
result.status,
|
||||
result.error || "Embedding provider error",
|
||||
provider,
|
||||
resolvedModel || null
|
||||
).catch((err) => {
|
||||
log.debug(
|
||||
"EMBED",
|
||||
`Cooldown write failed for ${provider}/${credentials.connectionId?.slice(0, 8)}: ${err}`
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
responseHeaders.set("Content-Type", "application/json");
|
||||
const errorPayload = toJsonErrorPayload(result.error, "Embedding provider error");
|
||||
return new Response(JSON.stringify(errorPayload), {
|
||||
|
||||
95
tests/unit/embedding-account-cooldown-10347.test.ts
Normal file
95
tests/unit/embedding-account-cooldown-10347.test.ts
Normal file
@@ -0,0 +1,95 @@
|
||||
import test from "node:test";
|
||||
import assert from "node:assert/strict";
|
||||
import fs from "node:fs";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
|
||||
const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-embed-cooldown-"));
|
||||
process.env.DATA_DIR = TEST_DATA_DIR;
|
||||
process.env.API_KEY_SECRET = process.env.API_KEY_SECRET || "embed-cooldown-test-secret";
|
||||
|
||||
const core = await import("../../src/lib/db/core.ts");
|
||||
const providersDb = await import("../../src/lib/db/providers.ts");
|
||||
const auth = await import("../../src/sse/services/auth.ts");
|
||||
|
||||
async function resetStorage() {
|
||||
core.resetDbInstance();
|
||||
fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true });
|
||||
fs.mkdirSync(TEST_DATA_DIR, { recursive: true });
|
||||
}
|
||||
|
||||
async function seedConnection(provider: string): Promise<string> {
|
||||
const conn = await providersDb.createProviderConnection({
|
||||
provider,
|
||||
authType: "apikey",
|
||||
apiKey: `${provider}-key`,
|
||||
isActive: true,
|
||||
testStatus: "active",
|
||||
});
|
||||
return (conn as Record<string, unknown>).id as string;
|
||||
}
|
||||
|
||||
test.after(() => {
|
||||
core.resetDbInstance();
|
||||
fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
test("#10347: markAccountUnavailable triggers cooldown on embedding 402", async () => {
|
||||
await resetStorage();
|
||||
const connId = await seedConnection("mistral");
|
||||
|
||||
const result = await auth.markAccountUnavailable(
|
||||
connId,
|
||||
402,
|
||||
"Check your subscription on https://admin.mistral.ai/subscription",
|
||||
"mistral",
|
||||
"mistral-embed"
|
||||
);
|
||||
|
||||
assert.strictEqual(result.shouldFallback, true, "402 must trigger account cooldown");
|
||||
|
||||
// Verify the connection was marked — 402 is terminal (credits_exhausted),
|
||||
// which sets testStatus but not rateLimitedUntil
|
||||
const conn = await providersDb.getProviderConnectionById(connId);
|
||||
assert.strictEqual(
|
||||
conn.testStatus,
|
||||
"credits_exhausted",
|
||||
"402 must mark connection credits_exhausted"
|
||||
);
|
||||
});
|
||||
|
||||
test("#10347: markAccountUnavailable triggers cooldown on embedding 500", async () => {
|
||||
await resetStorage();
|
||||
const connId = await seedConnection("mistral");
|
||||
|
||||
const result = await auth.markAccountUnavailable(
|
||||
connId,
|
||||
500,
|
||||
"Internal server error",
|
||||
"mistral",
|
||||
"mistral-embed"
|
||||
);
|
||||
|
||||
assert.strictEqual(result.shouldFallback, true, "500 must trigger account cooldown");
|
||||
const conn = await providersDb.getProviderConnectionById(connId);
|
||||
assert.strictEqual(conn.testStatus, "unavailable", "500 must mark connection unavailable");
|
||||
});
|
||||
|
||||
test("#10347: embedding 400 (bad request) does NOT trigger account cooldown", async () => {
|
||||
await resetStorage();
|
||||
const connId = await seedConnection("mistral");
|
||||
|
||||
// 400 bad request is a client error, not an account issue — should not cool down
|
||||
const result = await auth.markAccountUnavailable(
|
||||
connId,
|
||||
400,
|
||||
"Invalid embedding input format",
|
||||
"mistral",
|
||||
"mistral-embed"
|
||||
);
|
||||
|
||||
// Generic 400 returns shouldFallback:false (not account-fallback-worthy)
|
||||
assert.strictEqual(result.shouldFallback, false, "400 bad request must not trigger cooldown");
|
||||
const conn = await providersDb.getProviderConnectionById(connId);
|
||||
assert.strictEqual(conn.testStatus, "active", "400 must keep connection active");
|
||||
});
|
||||
184
tests/unit/embedding-cooldown-integration-10347.test.ts
Normal file
184
tests/unit/embedding-cooldown-integration-10347.test.ts
Normal file
@@ -0,0 +1,184 @@
|
||||
import test from "node:test";
|
||||
import assert from "node:assert/strict";
|
||||
import fs from "node:fs";
|
||||
import os from "node:os";
|
||||
import path from "node:path";
|
||||
|
||||
// #10347 integration: exercise createEmbeddingResponse end-to-end with a
|
||||
// mocked upstream that returns 402, then verify the connection gets cooled
|
||||
// down. This proves the production code path actually calls
|
||||
// markAccountUnavailable — the direct-call tests in
|
||||
// embedding-account-cooldown-10347.test.ts would pass even if the
|
||||
// production block were removed.
|
||||
|
||||
const TEST_DATA_DIR = fs.mkdtempSync(path.join(os.tmpdir(), "omniroute-embed-int-"));
|
||||
process.env.DATA_DIR = TEST_DATA_DIR;
|
||||
process.env.API_KEY_SECRET = process.env.API_KEY_SECRET || "embed-int-test-secret";
|
||||
|
||||
const core = await import("../../src/lib/db/core.ts");
|
||||
const providersDb = await import("../../src/lib/db/providers.ts");
|
||||
const auth = await import("../../src/sse/services/auth.ts");
|
||||
|
||||
function resetStorage() {
|
||||
core.resetDbInstance();
|
||||
fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true });
|
||||
fs.mkdirSync(TEST_DATA_DIR, { recursive: true });
|
||||
}
|
||||
|
||||
async function seedConnection(
|
||||
provider: string,
|
||||
overrides: Record<string, unknown> = {}
|
||||
): Promise<string> {
|
||||
const conn = await providersDb.createProviderConnection({
|
||||
provider,
|
||||
authType: "apikey",
|
||||
apiKey: `${provider}-key`,
|
||||
isActive: true,
|
||||
testStatus: "active",
|
||||
...overrides,
|
||||
});
|
||||
return (conn as Record<string, unknown>).id as string;
|
||||
}
|
||||
|
||||
test.after(() => {
|
||||
core.resetDbInstance();
|
||||
fs.rmSync(TEST_DATA_DIR, { recursive: true, force: true });
|
||||
});
|
||||
|
||||
test("createEmbeddingResponse marks connection on upstream 402", async () => {
|
||||
resetStorage();
|
||||
const connId = await seedConnection("mistral");
|
||||
|
||||
// Mock upstream to return 402 (subscription expired).
|
||||
const originalFetch = globalThis.fetch;
|
||||
globalThis.fetch = (async () =>
|
||||
new Response(
|
||||
JSON.stringify({ error: "Check your subscription on https://admin.mistral.ai/subscription" }),
|
||||
{ status: 402, headers: { "Content-Type": "application/json" } }
|
||||
)) as typeof globalThis.fetch;
|
||||
|
||||
try {
|
||||
// Import the service AFTER seeding the DB so its module-level caches
|
||||
// see the seeded connection.
|
||||
const { createEmbeddingResponse } = await import("../../src/lib/embeddings/service.ts");
|
||||
|
||||
// Call the production path — this must exercise the markAccountUnavailable
|
||||
// block we added in #10347.
|
||||
const res = await createEmbeddingResponse(
|
||||
{ model: "mistral-embed", input: "hello" },
|
||||
{ connectionId: connId }
|
||||
);
|
||||
|
||||
assert.equal(res.status, 402, "must return upstream status");
|
||||
|
||||
// Give the fire-and-forget markAccountUnavailable call time to settle.
|
||||
await new Promise((r) => setImmediate(r));
|
||||
await new Promise((r) => setImmediate(r));
|
||||
|
||||
// Verify the connection was actually marked — this is the assertion
|
||||
// that would FAIL if the production block were removed.
|
||||
const conn = await providersDb.getProviderConnectionById(connId);
|
||||
assert.equal(
|
||||
conn.testStatus,
|
||||
"credits_exhausted",
|
||||
"402 must mark connection credits_exhausted via production code path"
|
||||
);
|
||||
} finally {
|
||||
globalThis.fetch = originalFetch;
|
||||
}
|
||||
});
|
||||
|
||||
test("cooled account is skipped on next request — second connection selected", async () => {
|
||||
resetStorage();
|
||||
const conn1 = await seedConnection("mistral", { apiKey: "mistral-key-1" });
|
||||
const conn2 = await seedConnection("mistral", { apiKey: "mistral-key-2" });
|
||||
|
||||
const originalFetch = globalThis.fetch;
|
||||
let fetchCallCount = 0;
|
||||
|
||||
try {
|
||||
const { createEmbeddingResponse } = await import("../../src/lib/embeddings/service.ts");
|
||||
|
||||
// First request: upstream returns 402 → conn1 gets cooled.
|
||||
globalThis.fetch = (async () => {
|
||||
fetchCallCount++;
|
||||
return new Response(JSON.stringify({ error: "subscription expired" }), {
|
||||
status: 402,
|
||||
headers: { "Content-Type": "application/json" },
|
||||
});
|
||||
}) as typeof globalThis.fetch;
|
||||
|
||||
const res1 = await createEmbeddingResponse(
|
||||
{ model: "mistral-embed", input: "hello" },
|
||||
{ connectionId: conn1 }
|
||||
);
|
||||
assert.equal(res1.status, 402);
|
||||
|
||||
// Wait for fire-and-forget cooldown write.
|
||||
await new Promise((r) => setImmediate(r));
|
||||
await new Promise((r) => setImmediate(r));
|
||||
|
||||
// Verify conn1 is cooled.
|
||||
const conn1After = await providersDb.getProviderConnectionById(conn1);
|
||||
assert.equal(conn1After.testStatus, "credits_exhausted", "conn1 must be cooled");
|
||||
|
||||
// Second request: upstream returns 200.
|
||||
globalThis.fetch = (async () => {
|
||||
fetchCallCount++;
|
||||
return new Response(
|
||||
JSON.stringify({
|
||||
data: [{ embedding: [0.1, 0.2], index: 0 }],
|
||||
model: "mistral-embed",
|
||||
usage: { prompt_tokens: 1, total_tokens: 1 },
|
||||
}),
|
||||
{ status: 200, headers: { "Content-Type": "application/json" } }
|
||||
);
|
||||
}) as typeof globalThis.fetch;
|
||||
|
||||
// Call without specifying connectionId — credential selection should
|
||||
// skip conn1 (credits_exhausted) and pick conn2.
|
||||
const res2 = await createEmbeddingResponse({ model: "mistral-embed", input: "world" }, {});
|
||||
assert.equal(res2.status, 200, "second request must succeed via conn2");
|
||||
|
||||
// Verify conn2 is still healthy.
|
||||
const conn2After = await providersDb.getProviderConnectionById(conn2);
|
||||
assert.equal(conn2After.testStatus, "active", "conn2 must remain active");
|
||||
} finally {
|
||||
globalThis.fetch = originalFetch;
|
||||
}
|
||||
});
|
||||
|
||||
test("createEmbeddingResponse skips 400 (bad request) — no cooldown", async () => {
|
||||
resetStorage();
|
||||
const connId = await seedConnection("mistral");
|
||||
|
||||
const originalFetch = globalThis.fetch;
|
||||
globalThis.fetch = (async () =>
|
||||
new Response(JSON.stringify({ error: "Invalid embedding input format" }), {
|
||||
status: 400,
|
||||
headers: { "Content-Type": "application/json" },
|
||||
})) as typeof globalThis.fetch;
|
||||
|
||||
try {
|
||||
const { createEmbeddingResponse } = await import("../../src/lib/embeddings/service.ts");
|
||||
|
||||
const res = await createEmbeddingResponse(
|
||||
{ model: "mistral-embed", input: "hello" },
|
||||
{ connectionId: connId }
|
||||
);
|
||||
|
||||
assert.equal(res.status, 400, "must return upstream status");
|
||||
|
||||
await new Promise((r) => setImmediate(r));
|
||||
await new Promise((r) => setImmediate(r));
|
||||
|
||||
const conn = await providersDb.getProviderConnectionById(connId);
|
||||
assert.equal(
|
||||
conn.testStatus,
|
||||
"active",
|
||||
"400 must NOT mark connection — account is fine, request was wrong"
|
||||
);
|
||||
} finally {
|
||||
globalThis.fetch = originalFetch;
|
||||
}
|
||||
});
|
||||
Reference in New Issue
Block a user