File size: 6,093 Bytes
88c4c60 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 | // Concurrency stress test β simulate many parallel saveRequestUsage / saveRequestDetail
// to verify atomic counter, no data loss, no race conditions.
import fs from "node:fs";
import os from "node:os";
import path from "node:path";
import { describe, it, expect, beforeAll, afterAll, vi } from "vitest";
const originalDataDir = process.env.DATA_DIR;
let tempDir;
let db;
beforeAll(async () => {
tempDir = fs.mkdtempSync(path.join(os.tmpdir(), "9router-concurrent-"));
process.env.DATA_DIR = tempDir;
vi.resetModules();
db = await import("@/lib/db/index.js");
await db.initDb();
});
afterAll(() => {
if (tempDir) fs.rmSync(tempDir, { recursive: true, force: true });
if (originalDataDir === undefined) delete process.env.DATA_DIR;
else process.env.DATA_DIR = originalDataDir;
});
describe("DB Concurrency β atomic safety", () => {
it("100 parallel saveRequestUsage β no count loss", async () => {
const N = 100;
const promises = [];
for (let i = 0; i < N; i++) {
promises.push(db.saveRequestUsage({
provider: "openai", model: "gpt-4", connectionId: "c1",
tokens: { prompt_tokens: 10, completion_tokens: 5 },
endpoint: "/v1/chat", status: "ok",
}));
}
await Promise.all(promises);
const stats = await db.getUsageStats("24h");
expect(stats.totalRequests).toBe(N);
expect(stats.byProvider.openai.requests).toBe(N);
expect(stats.byProvider.openai.promptTokens).toBe(N * 10);
const hist = await db.getUsageHistory({ provider: "openai" });
expect(hist.length).toBe(N);
});
it("200 parallel saveRequestDetail β all flushed", async () => {
await db.updateSettings({ enableObservability: true, observabilityBatchSize: 10 });
const N = 200;
const promises = [];
for (let i = 0; i < N; i++) {
promises.push(db.saveRequestDetail({
id: `det-${i}`, provider: "openai", model: "gpt-4",
connectionId: "c1", status: "ok",
tokens: { prompt_tokens: 1 }, request: { i }, response: { ok: true },
}));
}
await Promise.all(promises);
// Wait for any timer-based flush
await new Promise((r) => setTimeout(r, 6000));
const list = await db.getRequestDetails({ provider: "openai", pageSize: 500 });
expect(list.pagination.totalItems).toBeGreaterThanOrEqual(N);
}, 15000);
it("mixed concurrent: usage + details + connections + aliases", async () => {
const ops = [];
for (let i = 0; i < 50; i++) {
ops.push(db.saveRequestUsage({
provider: "anthropic", model: `m-${i % 3}`, connectionId: "c2",
tokens: { prompt_tokens: 20 }, status: "ok",
}));
ops.push(db.setModelAlias(`a-${i}`, `target-${i}`));
ops.push(db.disableModels("openai", [`d-${i}`]));
}
await Promise.all(ops);
const aliases = await db.getModelAliases();
expect(Object.keys(aliases).filter((k) => k.startsWith("a-")).length).toBe(50);
const disabled = await db.getDisabledByProvider("openai");
expect(disabled.length).toBeGreaterThanOrEqual(50);
const stats = await db.getUsageStats("24h");
expect(stats.byProvider.anthropic.requests).toBe(50);
}, 30000);
it("updateSettings parallel β no merge loss", async () => {
const N = 50;
await db.updateSettings({ counter: 0 });
const promises = [];
for (let i = 0; i < N; i++) {
promises.push(db.updateSettings({ [`field${i}`]: `v${i}` }));
}
await Promise.all(promises);
const s = await db.getSettings();
for (let i = 0; i < N; i++) {
expect(s[`field${i}`]).toBe(`v${i}`); // all updates preserved
}
});
it("OAuth refresh race: parallel updateProviderConnection on same id", async () => {
const conn = await db.createProviderConnection({
provider: "oauth-test", authType: "oauth", email: "x@y.com",
accessToken: "initial", refreshToken: "rt-initial",
});
// 20 parallel updates each with a unique field
const N = 20;
const promises = [];
for (let i = 0; i < N; i++) {
promises.push(db.updateProviderConnection(conn.id, { [`marker${i}`]: i }));
}
await Promise.all(promises);
const after = await db.getProviderConnectionById(conn.id);
for (let i = 0; i < N; i++) {
expect(after[`marker${i}`]).toBe(i); // no field lost
}
expect(after.refreshToken).toBe("rt-initial"); // base preserved
});
it("addCustomModel race: parallel duplicate adds β only 1 inserted", async () => {
const N = 30;
const promises = [];
for (let i = 0; i < N; i++) {
promises.push(db.addCustomModel({ providerAlias: "racep", id: "racemodel", type: "llm", name: "r" }));
}
const results = await Promise.all(promises);
const trueCount = results.filter((r) => r === true).length;
expect(trueCount).toBe(1); // exactly one wins
const all = await db.getCustomModels();
expect(all.filter((m) => m.providerAlias === "racep" && m.id === "racemodel").length).toBe(1);
});
it("updatePricing race: parallel adds different models β all merged", async () => {
const N = 30;
const promises = [];
for (let i = 0; i < N; i++) {
promises.push(db.updatePricing({ "race-prov": { [`m${i}`]: { input: i, output: i * 2 } } }));
}
await Promise.all(promises);
const p = await db.getPricing();
for (let i = 0; i < N; i++) {
expect(p["race-prov"][`m${i}`]).toEqual({ input: i, output: i * 2 });
}
});
it("daily summary aggregates correctly under parallel writes", async () => {
const N = 50;
const promises = [];
for (let i = 0; i < N; i++) {
promises.push(db.saveRequestUsage({
provider: "google", model: "gemini-pro", connectionId: "cG",
tokens: { prompt_tokens: 100, completion_tokens: 50 },
status: "ok",
}));
}
await Promise.all(promises);
const stats = await db.getUsageStats("7d");
const g = stats.byProvider.google;
expect(g).toBeDefined();
expect(g.requests).toBe(N);
expect(g.promptTokens).toBe(N * 100);
expect(g.completionTokens).toBe(N * 50);
});
});
|