Skip to content

Commit 84cf6e5

Browse files
committed
perf(worker): eliminate DDL locks, batch user existence check, and add concurrent user fetching
1 parent 68aeaa2 commit 84cf6e5

4 files changed

Lines changed: 330 additions & 17 deletions

File tree

‎src/features/leaderboard/services/calculate-leaderboard.ts‎

Lines changed: 44 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,29 @@ function getRefreshLimit(): number {
5454
return getEnvInt("LEADERBOARD_REFRESH_LIMIT", 500);
5555
}
5656

57+
function getConcurrency(): number {
58+
return getEnvInt("LEADERBOARD_CONCURRENCY", 3);
59+
}
60+
61+
async function runConcurrent<T>(
62+
items: T[],
63+
concurrency: number,
64+
task: (item: T) => Promise<void>,
65+
): Promise<void> {
66+
if (items.length === 0) return;
67+
const limit = Math.max(1, Math.min(concurrency, items.length));
68+
let currentIndex = 0;
69+
70+
const workers = Array.from({ length: limit }, async () => {
71+
while (currentIndex < items.length) {
72+
const index = currentIndex++;
73+
await task(items[index]);
74+
}
75+
});
76+
77+
await Promise.all(workers);
78+
}
79+
5780
function getSourceUrl(country: string): string {
5881
const template = process.env.LEADERBOARD_SOURCE_URL_TEMPLATE?.trim();
5982
if (!template) {
@@ -92,6 +115,7 @@ export async function seedNewUsers(
92115
users: LeaderboardSourceEntry[],
93116
seedLimit: number,
94117
staleDays: number,
118+
concurrency: number = getConcurrency(),
95119
): Promise<{
96120
newUsersCount: number;
97121
skippedExistingCount: number;
@@ -104,15 +128,19 @@ export async function seedNewUsers(
104128
const fetchMetrics: { duration: number; errors: { part: string; reason: string }[] }[] = [];
105129

106130
const usersToSeed = users.slice(0, seedLimit);
131+
const existingSet = await db.getExistingUsernames(usersToSeed.map((u) => u.login));
107132

133+
const usersToFetch: LeaderboardSourceEntry[] = [];
108134
for (const user of usersToSeed) {
109-
try {
110-
const exists = await db.userExists(user.login);
111-
if (exists) {
112-
skippedExistingCount += 1;
113-
continue;
114-
}
135+
if (existingSet.has(user.login.toLowerCase())) {
136+
skippedExistingCount += 1;
137+
} else {
138+
usersToFetch.push(user);
139+
}
140+
}
115141

142+
await runConcurrent(usersToFetch, concurrency, async (user) => {
143+
try {
116144
const { data, metrics } = await getUserData(user.login, {
117145
cacheInRedis: false,
118146
withMetrics: true,
@@ -130,7 +158,7 @@ export async function seedNewUsers(
130158
} catch (e) {
131159
errors.push({ username: user.login, reason: e instanceof Error ? e.message : String(e) });
132160
}
133-
}
161+
});
134162

135163
return { newUsersCount, skippedExistingCount, errors, fetchMetrics };
136164
}
@@ -140,6 +168,7 @@ export async function refreshStaleUsers(
140168
country: string,
141169
refreshLimit: number,
142170
staleDays: number,
171+
concurrency: number = getConcurrency(),
143172
): Promise<{
144173
refreshedCount: number;
145174
errors: { username: string; reason: string }[];
@@ -150,12 +179,10 @@ export async function refreshStaleUsers(
150179
const fetchMetrics: { duration: number; errors: { part: string; reason: string }[] }[] = [];
151180

152181
const topUsers = await db.getTopUsers(country, refreshLimit);
182+
const now = new Date();
183+
const staleUsers = topUsers.filter((row) => row.stale_after < now);
153184

154-
for (const row of topUsers) {
155-
if (row.stale_after >= new Date()) {
156-
continue;
157-
}
158-
185+
await runConcurrent(staleUsers, concurrency, async (row) => {
159186
try {
160187
const { data, metrics } = await getUserData(row.username, {
161188
cacheInRedis: false,
@@ -174,7 +201,7 @@ export async function refreshStaleUsers(
174201
} catch (e) {
175202
errors.push({ username: row.username, reason: e instanceof Error ? e.message : String(e) });
176203
}
177-
}
204+
});
178205

179206
return { refreshedCount, errors, fetchMetrics };
180207
}
@@ -242,23 +269,24 @@ export async function calculateLeaderboard(
242269
refreshLimit?: number;
243270
staleDays?: number;
244271
displayLimit?: number;
272+
concurrency?: number;
245273
},
246274
): Promise<CalculateLeaderboardResponse> {
247275
const db = getDatabaseStore();
248-
await db.initializeSchema();
249276

250277
const staleDays = overrides?.staleDays ?? getStaleDays();
251278
const seedLimit = overrides?.seedLimit ?? getSeedLimit();
252279
const refreshLimit = overrides?.refreshLimit ?? getRefreshLimit();
280+
const concurrency = overrides?.concurrency ?? getConcurrency();
253281

254282
// 1. Fetch source users
255283
const sourceData = await fetchCommittersFromTop(country);
256284

257285
// 2a. Seed new users
258-
const seedResult = await seedNewUsers(db, sourceData.users, seedLimit, staleDays);
286+
const seedResult = await seedNewUsers(db, sourceData.users, seedLimit, staleDays, concurrency);
259287

260288
// 2b. Refresh stale users from DB top N
261-
const refreshResult = await refreshStaleUsers(db, country, refreshLimit, staleDays);
289+
const refreshResult = await refreshStaleUsers(db, country, refreshLimit, staleDays, concurrency);
262290

263291
// 3. Build leaderboard result
264292
const allErrors = [
Lines changed: 203 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,203 @@
1+
import { beforeEach, describe, expect, test, vi } from "vitest";
2+
3+
const mocks = vi.hoisted(() => ({
4+
getUserData: vi.fn(),
5+
calculateUserScore: vi.fn(),
6+
persistUserScores: vi.fn().mockResolvedValue(undefined),
7+
initializeSchema: vi.fn().mockResolvedValue(undefined),
8+
getExistingUsernames: vi.fn().mockResolvedValue(new Set<string>()),
9+
getTopUsers: vi.fn().mockResolvedValue([]),
10+
getLeaderboard: vi.fn().mockResolvedValue([]),
11+
getLeaderboardCount: vi.fn().mockResolvedValue(0),
12+
}));
13+
14+
vi.mock("@/lib/github", () => ({
15+
getUserData: mocks.getUserData,
16+
}));
17+
18+
vi.mock("@/features/scoring", () => ({
19+
calculateUserScore: mocks.calculateUserScore,
20+
}));
21+
22+
vi.mock("@/features/developer/services", () => ({
23+
persistUserScores: mocks.persistUserScores,
24+
}));
25+
26+
vi.mock("@/lib/db", () => ({
27+
getDatabaseStore: () => ({
28+
initializeSchema: mocks.initializeSchema,
29+
getExistingUsernames: mocks.getExistingUsernames,
30+
getTopUsers: mocks.getTopUsers,
31+
getLeaderboard: mocks.getLeaderboard,
32+
getLeaderboardCount: mocks.getLeaderboardCount,
33+
}),
34+
}));
35+
36+
vi.mock("@/lib/cache", () => ({
37+
getCacheConfigFromEnv: () => ({ enabled: false, namespace: "test" }),
38+
createCacheStore: () => ({ enabled: false, del: vi.fn() }),
39+
}));
40+
41+
import type { DatabaseStore } from "@/lib/db";
42+
import {
43+
seedNewUsers,
44+
refreshStaleUsers,
45+
calculateLeaderboard,
46+
} from "../services/calculate-leaderboard";
47+
48+
describe("calculate-leaderboard service optimizations", () => {
49+
beforeEach(() => {
50+
vi.clearAllMocks();
51+
});
52+
53+
describe("seedNewUsers", () => {
54+
test("batches username existence checks and skips existing users", async () => {
55+
const mockDb = {
56+
getExistingUsernames: vi.fn().mockResolvedValue(new Set(["alice"])),
57+
} as unknown as DatabaseStore;
58+
59+
mocks.getUserData.mockResolvedValue({
60+
data: { login: "bob" },
61+
metrics: { duration: 12, errors: [] },
62+
});
63+
mocks.calculateUserScore.mockReturnValue({ finalScore: 85 });
64+
65+
const users = [
66+
{
67+
rank: 1,
68+
name: "Alice",
69+
login: "alice",
70+
avatarUrl: "https://example.com/1.png",
71+
contributions: 100,
72+
},
73+
{
74+
rank: 2,
75+
name: "Bob",
76+
login: "bob",
77+
avatarUrl: "https://example.com/2.png",
78+
contributions: 50,
79+
},
80+
];
81+
82+
const result = await seedNewUsers(mockDb, users, 10, 30, 2);
83+
84+
// getExistingUsernames was called once with both logins
85+
expect(mockDb.getExistingUsernames).toHaveBeenCalledTimes(1);
86+
expect(mockDb.getExistingUsernames).toHaveBeenCalledWith(["alice", "bob"]);
87+
88+
// Alice was skipped, Bob was fetched
89+
expect(result.skippedExistingCount).toBe(1);
90+
expect(result.newUsersCount).toBe(1);
91+
expect(mocks.getUserData).toHaveBeenCalledTimes(1);
92+
expect(mocks.getUserData).toHaveBeenCalledWith("bob", {
93+
cacheInRedis: false,
94+
withMetrics: true,
95+
});
96+
expect(mocks.persistUserScores).toHaveBeenCalledTimes(1);
97+
});
98+
99+
test("handles errors for individual users without failing the batch", async () => {
100+
const mockDb = {
101+
getExistingUsernames: vi.fn().mockResolvedValue(new Set()),
102+
} as unknown as DatabaseStore;
103+
104+
mocks.getUserData
105+
.mockRejectedValueOnce(new Error("Rate limit exceeded"))
106+
.mockResolvedValueOnce({
107+
data: { login: "user2" },
108+
metrics: { duration: 10, errors: [] },
109+
});
110+
mocks.calculateUserScore.mockReturnValue({ finalScore: 90 });
111+
112+
const users = [
113+
{
114+
rank: 1,
115+
name: "User 1",
116+
login: "user1",
117+
avatarUrl: "https://example.com/1.png",
118+
contributions: 10,
119+
},
120+
{
121+
rank: 2,
122+
name: "User 2",
123+
login: "user2",
124+
avatarUrl: "https://example.com/2.png",
125+
contributions: 20,
126+
},
127+
];
128+
129+
const result = await seedNewUsers(mockDb, users, 10, 30, 2);
130+
131+
expect(result.newUsersCount).toBe(1);
132+
expect(result.errors).toHaveLength(1);
133+
expect(result.errors[0]).toEqual({
134+
username: "user1",
135+
reason: "Rate limit exceeded",
136+
});
137+
});
138+
});
139+
140+
describe("refreshStaleUsers", () => {
141+
test("only refreshes users whose stale_after is in the past", async () => {
142+
const mockDb = {
143+
getTopUsers: vi.fn().mockResolvedValue([
144+
{ username: "staleUser", stale_after: new Date(Date.now() - 10000), final_score: 80 },
145+
{ username: "freshUser", stale_after: new Date(Date.now() + 100000), final_score: 90 },
146+
]),
147+
} as unknown as DatabaseStore;
148+
149+
mocks.getUserData.mockResolvedValue({
150+
data: { login: "staleUser" },
151+
metrics: { duration: 15, errors: [] },
152+
});
153+
mocks.calculateUserScore.mockReturnValue({ finalScore: 85 });
154+
155+
const result = await refreshStaleUsers(mockDb, "sweden", 10, 30, 2);
156+
157+
expect(result.refreshedCount).toBe(1);
158+
expect(mocks.getUserData).toHaveBeenCalledTimes(1);
159+
expect(mocks.getUserData).toHaveBeenCalledWith("staleUser", {
160+
cacheInRedis: false,
161+
withMetrics: true,
162+
});
163+
expect(mocks.persistUserScores).toHaveBeenCalledTimes(1);
164+
});
165+
});
166+
167+
describe("calculateLeaderboard", () => {
168+
test("does NOT call db.initializeSchema during execution", async () => {
169+
process.env.LEADERBOARD_SOURCE_URL_TEMPLATE = "https://example.com/{country}.yml";
170+
171+
const sampleYaml = `
172+
title: Sweden
173+
total_user_count: 1
174+
users:
175+
- rank: 1
176+
name: Tester
177+
login: tester
178+
avatarUrl: https://example.com/tester.png
179+
contributions: 10
180+
`;
181+
vi.stubGlobal(
182+
"fetch",
183+
vi.fn().mockResolvedValue({
184+
ok: true,
185+
text: vi.fn().mockResolvedValue(sampleYaml),
186+
}),
187+
);
188+
189+
mocks.getExistingUsernames.mockResolvedValue(new Set(["tester"]));
190+
mocks.getTopUsers.mockResolvedValue([]);
191+
mocks.getLeaderboard.mockResolvedValue([]);
192+
mocks.getLeaderboardCount.mockResolvedValue(1);
193+
194+
await calculateLeaderboard("sweden");
195+
196+
// Verify lock-eliminating fix: initializeSchema must NOT be called
197+
expect(mocks.initializeSchema).not.toHaveBeenCalled();
198+
expect(mocks.getExistingUsernames).toHaveBeenCalledWith(["tester"]);
199+
200+
vi.unstubAllGlobals();
201+
});
202+
});
203+
});

‎src/lib/db/db-store.ts‎

Lines changed: 23 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -73,7 +73,7 @@ function getPoolConfig(): PoolConfig {
7373
connectionString,
7474
max: isServerless ? 2 : 10,
7575
idleTimeoutMillis: isServerless ? 10_000 : 30_000,
76-
connectionTimeoutMillis: 5_000,
76+
connectionTimeoutMillis: 10_000,
7777
};
7878
}
7979

@@ -279,6 +279,28 @@ export class DatabaseStore {
279279
return (result.rowCount ?? 0) > 0;
280280
}
281281

282+
async getExistingUsernames(usernames: string[]): Promise<Set<string>> {
283+
if (!usernames.length) {
284+
return new Set();
285+
}
286+
287+
const lowerUsernames = Array.from(
288+
new Set(usernames.map((u) => u.trim().toLowerCase()).filter(Boolean)),
289+
);
290+
291+
if (lowerUsernames.length === 0) {
292+
return new Set();
293+
}
294+
295+
const client = getPool();
296+
const result = await client.query<{ username: string }>(
297+
"SELECT LOWER(username) AS username FROM github_users WHERE LOWER(username) = ANY($1::text[])",
298+
[lowerUsernames],
299+
);
300+
301+
return new Set(result.rows.map((r) => r.username.toLowerCase()));
302+
}
303+
282304
// ── Leaderboard operations ──────────────────────────────────────────
283305

284306
async getLeaderboard(country: string, limit: number = 500): Promise<LeaderboardUserRow[]> {

0 commit comments

Comments
 (0)