mirror of
https://github.com/danny-avila/LibreChat.git
synced 2026-09-06 14:39:10 +00:00
fix: Make MCP aggregate writes atomic
This commit is contained in:
parent
bf7e77872c
commit
eb117d1b49
4 changed files with 152 additions and 14 deletions
|
|
@ -45,6 +45,15 @@ const runtimePlaceholderYamlEntry: t.ParsedServerConfig = {
|
|||
headers: { 'X-User-Id': '{{LIBRECHAT_USER_ID}}' },
|
||||
};
|
||||
|
||||
const runtimeApiKeyPlaceholderYamlEntry: t.ParsedServerConfig = {
|
||||
...startupDeferredYamlEntry,
|
||||
apiKey: {
|
||||
source: 'admin',
|
||||
authorization_type: 'bearer',
|
||||
key: '{{LIBRECHAT_OPENID_ACCESS_TOKEN}}',
|
||||
},
|
||||
};
|
||||
|
||||
describe('MCPServersRegistry.setResolvedInstructions', () => {
|
||||
let registry: MCPServersRegistry;
|
||||
|
||||
|
|
@ -139,6 +148,7 @@ describe('MCPServersRegistry.setResolvedInstructions', () => {
|
|||
it.each([
|
||||
['OAuth', oauthDeferredYamlEntry],
|
||||
['runtime placeholders', runtimePlaceholderYamlEntry],
|
||||
['runtime placeholders in an admin API key', runtimeApiKeyPlaceholderYamlEntry],
|
||||
[
|
||||
'configured-oauth-block',
|
||||
{
|
||||
|
|
@ -390,6 +400,7 @@ describe('UserConnectionManager.backfillResolvedInstructions', () => {
|
|||
},
|
||||
],
|
||||
['runtime placeholders', runtimePlaceholderYamlEntry],
|
||||
['runtime placeholders in an admin API key', runtimeApiKeyPlaceholderYamlEntry],
|
||||
[
|
||||
'a configured oauth block with requiresOAuth stamped false',
|
||||
{
|
||||
|
|
|
|||
|
|
@ -13,21 +13,46 @@ import { BaseRegistryCache } from './BaseRegistryCache';
|
|||
* caused by SCAN under concurrent load in large deployments (see GitHub #11624, #12408).
|
||||
*
|
||||
* Trade-offs:
|
||||
* - `add/update/remove` use a serialized read-modify-write on the aggregate key via a
|
||||
* promise-based mutex. This prevents concurrent writes from racing within a single
|
||||
* process (e.g., during `Promise.allSettled` initialization of multiple servers).
|
||||
* - In-memory writes use a serialized read-modify-write via a promise-based mutex.
|
||||
* Redis writes use single-key Lua mutations, so replicas cannot overwrite one another.
|
||||
* - The entire config map is serialized/deserialized on every operation. With typical MCP
|
||||
* deployments (~5-50 servers), the JSON payload is small (10-50KB).
|
||||
* - Cross-instance visibility is preserved: all instances read/write the same Redis key,
|
||||
* so reinspection results propagate automatically after readThroughCache TTL expiry.
|
||||
*
|
||||
* `patch` uses a Redis-side Lua compare-and-set when Redis backs this cache.
|
||||
* This keeps simultaneous replicas from losing another server's patch and makes
|
||||
* `resolvedInstructions` first-write-wins. Initialization writes remain leader
|
||||
* coordinated, and manual reinspection is intentionally rare.
|
||||
* All mutations use Redis-side Lua when Redis backs this cache. This keeps
|
||||
* simultaneous replicas from losing one another's changes and makes
|
||||
* `resolvedInstructions` first-write-wins.
|
||||
*/
|
||||
const AGGREGATE_KEY = '__all__';
|
||||
|
||||
const MUTATE_AGGREGATE_ENTRY = `
|
||||
local operation = ARGV[1]
|
||||
local serverName = ARGV[2]
|
||||
local encoded = redis.call('GET', KEYS[1])
|
||||
local envelope
|
||||
if encoded then
|
||||
envelope = cjson.decode(encoded)
|
||||
else
|
||||
envelope = { value = {} }
|
||||
end
|
||||
if not envelope.value then envelope.value = {} end
|
||||
local existing = envelope.value[serverName]
|
||||
if operation == 'add' and existing then return -1 end
|
||||
if (operation == 'update' or operation == 'remove') and not existing then return 0 end
|
||||
local ttl = redis.call('PTTL', KEYS[1])
|
||||
if operation == 'remove' then
|
||||
envelope.value[serverName] = nil
|
||||
else
|
||||
local config = cjson.decode(ARGV[3])
|
||||
config.updatedAt = tonumber(ARGV[4])
|
||||
envelope.value[serverName] = config
|
||||
end
|
||||
redis.call('SET', KEYS[1], cjson.encode(envelope))
|
||||
if ttl > 0 then redis.call('PEXPIRE', KEYS[1], ttl) end
|
||||
return 1
|
||||
`;
|
||||
|
||||
/** Keyv stores its serialized value in an envelope with a `value` member. This script
|
||||
* updates that envelope atomically and preserves a configured Redis expiration. */
|
||||
const PATCH_AGGREGATE_ENTRY = `
|
||||
|
|
@ -100,6 +125,24 @@ export class ServerConfigsCacheRedisAggregateKey
|
|||
return `${prefix}${this.cache.namespace}:${AGGREGATE_KEY}`;
|
||||
}
|
||||
|
||||
private async mutateRedisEntry(
|
||||
operation: 'add' | 'update' | 'upsert' | 'remove',
|
||||
serverName: string,
|
||||
config?: ParsedServerConfig,
|
||||
updatedAt?: number,
|
||||
): Promise<number> {
|
||||
const result = await evalKeyvRedisScript(MUTATE_AGGREGATE_ENTRY, {
|
||||
keys: [this.aggregateRedisKey()],
|
||||
arguments: [
|
||||
operation,
|
||||
serverName,
|
||||
config ? JSON.stringify(config) : '',
|
||||
updatedAt != null ? String(updatedAt) : '',
|
||||
],
|
||||
});
|
||||
return typeof result === 'number' ? result : 0;
|
||||
}
|
||||
|
||||
/**
|
||||
* Serializes write operations to prevent concurrent read-modify-write races.
|
||||
* Reads (`get`, `getAll`) are not serialized — they can run concurrently.
|
||||
|
|
@ -149,6 +192,22 @@ export class ServerConfigsCacheRedisAggregateKey
|
|||
public async add(serverName: string, config: ParsedServerConfig): Promise<AddServerResult> {
|
||||
if (this.leaderOnly) await this.leaderCheck('add MCP servers');
|
||||
return this.withWriteLock(async () => {
|
||||
const storedConfig = { ...config, updatedAt: Date.now() };
|
||||
if (this.usesRedisStore()) {
|
||||
const result = await this.mutateRedisEntry(
|
||||
'add',
|
||||
serverName,
|
||||
storedConfig,
|
||||
storedConfig.updatedAt,
|
||||
);
|
||||
if (result === -1) {
|
||||
throw new Error(
|
||||
`Server "${serverName}" already exists in cache. Use update() to modify existing configs.`,
|
||||
);
|
||||
}
|
||||
this.successCheck(`add ${this.namespace} server "${serverName}"`, result === 1);
|
||||
return { serverName, config: storedConfig };
|
||||
}
|
||||
// Force fresh Redis read so the read-modify-write uses current data,
|
||||
// not a snapshot that may predate this write. Distinct from the finally-block
|
||||
// invalidation which cleans up after the write completes or throws.
|
||||
|
|
@ -159,7 +218,6 @@ export class ServerConfigsCacheRedisAggregateKey
|
|||
`Server "${serverName}" already exists in cache. Use update() to modify existing configs.`,
|
||||
);
|
||||
}
|
||||
const storedConfig = { ...config, updatedAt: Date.now() };
|
||||
const newAll = { ...all, [serverName]: storedConfig };
|
||||
const success = await this.cache.set(AGGREGATE_KEY, newAll);
|
||||
this.successCheck(`add ${this.namespace} server "${serverName}"`, success);
|
||||
|
|
@ -170,6 +228,17 @@ export class ServerConfigsCacheRedisAggregateKey
|
|||
public async update(serverName: string, config: ParsedServerConfig): Promise<void> {
|
||||
if (this.leaderOnly) await this.leaderCheck('update MCP servers');
|
||||
return this.withWriteLock(async () => {
|
||||
const updatedAt = Date.now();
|
||||
if (this.usesRedisStore()) {
|
||||
const result = await this.mutateRedisEntry('update', serverName, config, updatedAt);
|
||||
if (result === 0) {
|
||||
throw new Error(
|
||||
`Server "${serverName}" does not exist in cache. Use add() to create new configs.`,
|
||||
);
|
||||
}
|
||||
this.successCheck(`update ${this.namespace} server "${serverName}"`, result === 1);
|
||||
return;
|
||||
}
|
||||
this.invalidateLocalSnapshot(); // Force fresh Redis read (see add() comment)
|
||||
const all = await this.getAll();
|
||||
if (!all[serverName]) {
|
||||
|
|
@ -177,7 +246,7 @@ export class ServerConfigsCacheRedisAggregateKey
|
|||
`Server "${serverName}" does not exist in cache. Use add() to create new configs.`,
|
||||
);
|
||||
}
|
||||
const newAll = { ...all, [serverName]: { ...config, updatedAt: Date.now() } };
|
||||
const newAll = { ...all, [serverName]: { ...config, updatedAt } };
|
||||
const success = await this.cache.set(AGGREGATE_KEY, newAll);
|
||||
this.successCheck(`update ${this.namespace} server "${serverName}"`, success);
|
||||
});
|
||||
|
|
@ -186,9 +255,15 @@ export class ServerConfigsCacheRedisAggregateKey
|
|||
public async upsert(serverName: string, config: ParsedServerConfig): Promise<void> {
|
||||
if (this.leaderOnly) await this.leaderCheck('upsert MCP servers');
|
||||
return this.withWriteLock(async () => {
|
||||
const updatedAt = Date.now();
|
||||
if (this.usesRedisStore()) {
|
||||
const result = await this.mutateRedisEntry('upsert', serverName, config, updatedAt);
|
||||
this.successCheck(`upsert ${this.namespace} server "${serverName}"`, result === 1);
|
||||
return;
|
||||
}
|
||||
this.invalidateLocalSnapshot();
|
||||
const all = await this.getAll();
|
||||
const newAll = { ...all, [serverName]: { ...config, updatedAt: Date.now() } };
|
||||
const newAll = { ...all, [serverName]: { ...config, updatedAt } };
|
||||
const success = await this.cache.set(AGGREGATE_KEY, newAll);
|
||||
this.successCheck(`upsert ${this.namespace} server "${serverName}"`, success);
|
||||
});
|
||||
|
|
@ -236,6 +311,14 @@ export class ServerConfigsCacheRedisAggregateKey
|
|||
public async remove(serverName: string): Promise<void> {
|
||||
if (this.leaderOnly) await this.leaderCheck('remove MCP servers');
|
||||
return this.withWriteLock(async () => {
|
||||
if (this.usesRedisStore()) {
|
||||
const result = await this.mutateRedisEntry('remove', serverName);
|
||||
if (result === 0) {
|
||||
throw new Error(`Failed to remove server "${serverName}" in cache.`);
|
||||
}
|
||||
this.successCheck(`remove ${this.namespace} server "${serverName}"`, result === 1);
|
||||
return;
|
||||
}
|
||||
this.invalidateLocalSnapshot(); // Force fresh Redis read (see add() comment)
|
||||
const all = await this.getAll();
|
||||
if (!all[serverName]) {
|
||||
|
|
|
|||
|
|
@ -247,6 +247,42 @@ describe('ServerConfigsCacheRedisAggregateKey Integration Tests', () => {
|
|||
expect(result.server1.resolvedInstructions).toBe('server one instructions');
|
||||
expect(result.server2.resolvedInstructions).toBe('server two instructions');
|
||||
});
|
||||
|
||||
it('routes every aggregate mutation through Redis-side atomic updates', async () => {
|
||||
const replica = new ServerConfigsCacheRedisAggregateKey('agg-test', false);
|
||||
const cacheSetSpy = jest.spyOn(replica['cache'], 'set');
|
||||
|
||||
await replica.add('atomic-server', mockConfig1);
|
||||
await replica.update('atomic-server', mockConfig2);
|
||||
await replica.upsert('atomic-server', mockConfig3);
|
||||
await replica.remove('atomic-server');
|
||||
|
||||
expect(cacheSetSpy).not.toHaveBeenCalled();
|
||||
cacheSetSpy.mockRestore();
|
||||
});
|
||||
|
||||
it('preserves patches concurrent with whole-entry mutations on other replicas', async () => {
|
||||
const patchReplica = new ServerConfigsCacheRedisAggregateKey('agg-test', false);
|
||||
const writerReplica = new ServerConfigsCacheRedisAggregateKey('agg-test', false);
|
||||
|
||||
for (let i = 0; i < 20; i++) {
|
||||
const patchedName = `patched-${i}`;
|
||||
const updatedName = `updated-${i}`;
|
||||
await cache.add(patchedName, mockConfig1);
|
||||
await cache.add(updatedName, mockConfig2);
|
||||
|
||||
await expect(
|
||||
Promise.all([
|
||||
patchReplica.patch(patchedName, { resolvedInstructions: `instructions-${i}` }),
|
||||
writerReplica.update(updatedName, { ...mockConfig3, description: `updated-${i}` }),
|
||||
]),
|
||||
).resolves.toEqual([true, undefined]);
|
||||
|
||||
const result = await cache.getAll();
|
||||
expect(result[patchedName].resolvedInstructions).toBe(`instructions-${i}`);
|
||||
expect(result[updatedName].description).toBe(`updated-${i}`);
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
describe('reset operation', () => {
|
||||
|
|
|
|||
|
|
@ -211,9 +211,9 @@ type UserScopedConnectionConfig = Pick<
|
|||
ParsedServerConfig,
|
||||
'requiresOAuth' | 'source' | 'dbId' | 'startup'
|
||||
> & {
|
||||
/** Loosened like the fields below: raw (pre-inspection) configs carry an
|
||||
* optional `source`, and the gating predicates only read `apiKey?.source`. */
|
||||
apiKey?: { source?: 'user' | 'admin' } | null;
|
||||
/** Loosened like the fields below: raw (pre-inspection) configs carry
|
||||
* optional API-key fields, and the gating predicates only inspect them. */
|
||||
apiKey?: { key?: string; source?: 'user' | 'admin' } | null;
|
||||
args?: string[];
|
||||
/** Loosened from the parsed shapes so raw (pre-inspection) configs qualify;
|
||||
* scoping predicates only check key presence */
|
||||
|
|
@ -230,7 +230,15 @@ type UserScopedConnectionConfig = Pick<
|
|||
};
|
||||
|
||||
function placeholderBearingFields(config: UserScopedConnectionConfig): PlaceholderValue[] {
|
||||
return [config.args, config.env, config.headers, config.oauth, config.oauth_headers, config.url];
|
||||
return [
|
||||
config.apiKey?.key,
|
||||
config.args,
|
||||
config.env,
|
||||
config.headers,
|
||||
config.oauth,
|
||||
config.oauth_headers,
|
||||
config.url,
|
||||
];
|
||||
}
|
||||
|
||||
/** Whether a server should use MCP OAuth handling. */
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue