fix: prevent race conditions in sticky round-robin
Adds a mutex to serialize account selection and updates in the proxy engine. This ensures that concurrent requests respect the sticky limit and don't distribute to the same account simultaneously. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Sonnet 4.5 <noreply@anthropic.com>
This commit is contained in:
parent
4f292aae63
commit
3ad2f8dc58
1 changed files with 81 additions and 67 deletions
|
|
@ -2,6 +2,9 @@ import { getProviderConnections, validateApiKey, updateProviderConnection, getSe
|
||||||
import { isAccountUnavailable, getUnavailableUntil } from "open-sse/services/accountFallback.js";
|
import { isAccountUnavailable, getUnavailableUntil } from "open-sse/services/accountFallback.js";
|
||||||
import * as log from "../utils/logger.js";
|
import * as log from "../utils/logger.js";
|
||||||
|
|
||||||
|
// Mutex to prevent race conditions during account selection
|
||||||
|
let selectionMutex = Promise.resolve();
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Get provider credentials from localDb
|
* Get provider credentials from localDb
|
||||||
* Filters out unavailable accounts and returns the selected account based on strategy
|
* Filters out unavailable accounts and returns the selected account based on strategy
|
||||||
|
|
@ -9,86 +12,97 @@ import * as log from "../utils/logger.js";
|
||||||
* @param {string|null} excludeConnectionId - Connection ID to exclude (for retry with next account)
|
* @param {string|null} excludeConnectionId - Connection ID to exclude (for retry with next account)
|
||||||
*/
|
*/
|
||||||
export async function getProviderCredentials(provider, excludeConnectionId = null) {
|
export async function getProviderCredentials(provider, excludeConnectionId = null) {
|
||||||
const connections = await getProviderConnections({ provider, isActive: true });
|
// Acquire mutex to prevent race conditions
|
||||||
|
const currentMutex = selectionMutex;
|
||||||
|
let resolveMutex;
|
||||||
|
selectionMutex = new Promise(resolve => { resolveMutex = resolve; });
|
||||||
|
|
||||||
if (connections.length === 0) {
|
try {
|
||||||
log.warn("AUTH", `No credentials for ${provider}`);
|
await currentMutex;
|
||||||
return null;
|
|
||||||
}
|
|
||||||
|
|
||||||
// Filter out unavailable accounts and excluded connection
|
const connections = await getProviderConnections({ provider, isActive: true });
|
||||||
const availableConnections = connections.filter(c => {
|
|
||||||
if (excludeConnectionId && c.id === excludeConnectionId) return false;
|
|
||||||
if (isAccountUnavailable(c.rateLimitedUntil)) return false;
|
|
||||||
return true;
|
|
||||||
});
|
|
||||||
|
|
||||||
if (availableConnections.length === 0) {
|
if (connections.length === 0) {
|
||||||
log.warn("AUTH", `All ${connections.length} accounts for ${provider} unavailable`);
|
log.warn("AUTH", `No credentials for ${provider}`);
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
const settings = await getSettings();
|
// Filter out unavailable accounts and excluded connection
|
||||||
const strategy = settings.fallbackStrategy || "fill-first";
|
const availableConnections = connections.filter(c => {
|
||||||
|
if (excludeConnectionId && c.id === excludeConnectionId) return false;
|
||||||
let connection;
|
if (isAccountUnavailable(c.rateLimitedUntil)) return false;
|
||||||
if (strategy === "round-robin") {
|
return true;
|
||||||
const stickyLimit = settings.stickyRoundRobinLimit || 3;
|
|
||||||
|
|
||||||
// Sort by lastUsed (most recent first) to find current candidate
|
|
||||||
const byRecency = [...availableConnections].sort((a, b) => {
|
|
||||||
if (!a.lastUsedAt && !b.lastUsedAt) return (a.priority || 999) - (b.priority || 999);
|
|
||||||
if (!a.lastUsedAt) return 1;
|
|
||||||
if (!b.lastUsedAt) return -1;
|
|
||||||
return new Date(b.lastUsedAt) - new Date(a.lastUsedAt);
|
|
||||||
});
|
});
|
||||||
|
|
||||||
const current = byRecency[0];
|
if (availableConnections.length === 0) {
|
||||||
const currentCount = current?.consecutiveUseCount || 0;
|
log.warn("AUTH", `All ${connections.length} accounts for ${provider} unavailable`);
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
if (current && current.lastUsedAt && currentCount < stickyLimit) {
|
const settings = await getSettings();
|
||||||
// Stay with current account
|
const strategy = settings.fallbackStrategy || "fill-first";
|
||||||
connection = current;
|
|
||||||
// Update lastUsedAt and increment count
|
let connection;
|
||||||
updateProviderConnection(connection.id, {
|
if (strategy === "round-robin") {
|
||||||
lastUsedAt: new Date().toISOString(),
|
const stickyLimit = settings.stickyRoundRobinLimit || 3;
|
||||||
consecutiveUseCount: (connection.consecutiveUseCount || 0) + 1
|
|
||||||
}).catch(() => {});
|
// Sort by lastUsed (most recent first) to find current candidate
|
||||||
} else {
|
const byRecency = [...availableConnections].sort((a, b) => {
|
||||||
// Pick the least recently used (excluding current if possible)
|
|
||||||
const sortedByOldest = [...availableConnections].sort((a, b) => {
|
|
||||||
if (!a.lastUsedAt && !b.lastUsedAt) return (a.priority || 999) - (b.priority || 999);
|
if (!a.lastUsedAt && !b.lastUsedAt) return (a.priority || 999) - (b.priority || 999);
|
||||||
if (!a.lastUsedAt) return -1;
|
if (!a.lastUsedAt) return 1;
|
||||||
if (!b.lastUsedAt) return 1;
|
if (!b.lastUsedAt) return -1;
|
||||||
return new Date(a.lastUsedAt) - new Date(b.lastUsedAt);
|
return new Date(b.lastUsedAt) - new Date(a.lastUsedAt);
|
||||||
});
|
});
|
||||||
|
|
||||||
connection = sortedByOldest[0];
|
const current = byRecency[0];
|
||||||
|
const currentCount = current?.consecutiveUseCount || 0;
|
||||||
|
|
||||||
// Update lastUsedAt and reset count to 1
|
if (current && current.lastUsedAt && currentCount < stickyLimit) {
|
||||||
updateProviderConnection(connection.id, {
|
// Stay with current account
|
||||||
lastUsedAt: new Date().toISOString(),
|
connection = current;
|
||||||
consecutiveUseCount: 1
|
// Update lastUsedAt and increment count (await to ensure persistence)
|
||||||
}).catch(() => {});
|
await updateProviderConnection(connection.id, {
|
||||||
|
lastUsedAt: new Date().toISOString(),
|
||||||
|
consecutiveUseCount: (connection.consecutiveUseCount || 0) + 1
|
||||||
|
});
|
||||||
|
} else {
|
||||||
|
// Pick the least recently used (excluding current if possible)
|
||||||
|
const sortedByOldest = [...availableConnections].sort((a, b) => {
|
||||||
|
if (!a.lastUsedAt && !b.lastUsedAt) return (a.priority || 999) - (b.priority || 999);
|
||||||
|
if (!a.lastUsedAt) return -1;
|
||||||
|
if (!b.lastUsedAt) return 1;
|
||||||
|
return new Date(a.lastUsedAt) - new Date(b.lastUsedAt);
|
||||||
|
});
|
||||||
|
|
||||||
|
connection = sortedByOldest[0];
|
||||||
|
|
||||||
|
// Update lastUsedAt and reset count to 1 (await to ensure persistence)
|
||||||
|
await updateProviderConnection(connection.id, {
|
||||||
|
lastUsedAt: new Date().toISOString(),
|
||||||
|
consecutiveUseCount: 1
|
||||||
|
});
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
// Default: fill-first (already sorted by priority in getProviderConnections)
|
||||||
|
connection = availableConnections[0];
|
||||||
}
|
}
|
||||||
} else {
|
|
||||||
// Default: fill-first (already sorted by priority in getProviderConnections)
|
|
||||||
connection = availableConnections[0];
|
|
||||||
}
|
|
||||||
|
|
||||||
return {
|
return {
|
||||||
apiKey: connection.apiKey,
|
apiKey: connection.apiKey,
|
||||||
accessToken: connection.accessToken,
|
accessToken: connection.accessToken,
|
||||||
refreshToken: connection.refreshToken,
|
refreshToken: connection.refreshToken,
|
||||||
projectId: connection.projectId,
|
projectId: connection.projectId,
|
||||||
copilotToken: connection.providerSpecificData?.copilotToken,
|
copilotToken: connection.providerSpecificData?.copilotToken,
|
||||||
providerSpecificData: connection.providerSpecificData,
|
providerSpecificData: connection.providerSpecificData,
|
||||||
connectionId: connection.id,
|
connectionId: connection.id,
|
||||||
// Include current status for optimization check
|
// Include current status for optimization check
|
||||||
testStatus: connection.testStatus,
|
testStatus: connection.testStatus,
|
||||||
lastError: connection.lastError,
|
lastError: connection.lastError,
|
||||||
rateLimitedUntil: connection.rateLimitedUntil
|
rateLimitedUntil: connection.rateLimitedUntil
|
||||||
};
|
};
|
||||||
|
} finally {
|
||||||
|
if (resolveMutex) resolveMutex();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue