mirror of
https://github.com/MoonshotAI/kimi-code.git
synced 2026-08-23 07:37:18 +00:00
fix: reconcile MCP workspaces concurrently
This commit is contained in:
parent
69e1d79262
commit
6ed3fa0630
3 changed files with 76 additions and 25 deletions
|
|
@ -112,6 +112,24 @@ export type IWaitUntilData<T> = Omit<T, 'waitUntil' | 'signal'>;
|
|||
export class AsyncEmitter<T extends IWaitUntil> extends Emitter<T> {
|
||||
private _asyncDeliveryQueue?: LinkedList<[(event: T) => void, IWaitUntilData<T>]>;
|
||||
|
||||
async fireAsyncConcurrent(data: IWaitUntilData<T>, signal: AbortSignal): Promise<void> {
|
||||
if (this.isDisposed || this._listeners === undefined || signal.aborted) {
|
||||
return;
|
||||
}
|
||||
const snapshot = Array.from(this._listeners);
|
||||
await Promise.all(
|
||||
snapshot.map((entry) =>
|
||||
this.deliverAsync(
|
||||
(event) => {
|
||||
entry.listener.call(entry.thisArg, event);
|
||||
},
|
||||
data,
|
||||
signal,
|
||||
),
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
async fireAsync(data: IWaitUntilData<T>, signal: AbortSignal): Promise<void> {
|
||||
if (this.isDisposed || this._listeners === undefined) {
|
||||
return;
|
||||
|
|
@ -129,32 +147,37 @@ export class AsyncEmitter<T extends IWaitUntil> extends Emitter<T> {
|
|||
|
||||
while (this._asyncDeliveryQueue.size > 0 && !signal.aborted) {
|
||||
const [deliver, eventData] = this._asyncDeliveryQueue.shift()!;
|
||||
const thenables: Promise<unknown>[] = [];
|
||||
await this.deliverAsync(deliver, eventData, signal);
|
||||
}
|
||||
}
|
||||
|
||||
const event = {
|
||||
...eventData,
|
||||
signal,
|
||||
waitUntil: (p: Promise<unknown>): void => {
|
||||
if (Object.isFrozen(thenables)) {
|
||||
throw new Error('waitUntil can NOT be called asynchronously');
|
||||
}
|
||||
thenables.push(p);
|
||||
},
|
||||
} as T;
|
||||
|
||||
try {
|
||||
deliver(event);
|
||||
} catch (error) {
|
||||
onUnexpectedError(error);
|
||||
continue;
|
||||
}
|
||||
|
||||
void Object.freeze(thenables);
|
||||
const settled = await Promise.allSettled(thenables);
|
||||
for (const result of settled) {
|
||||
if (result.status === 'rejected') {
|
||||
onUnexpectedError(result.reason);
|
||||
private async deliverAsync(
|
||||
deliver: (event: T) => void,
|
||||
data: IWaitUntilData<T>,
|
||||
signal: AbortSignal,
|
||||
): Promise<void> {
|
||||
const thenables: Promise<unknown>[] = [];
|
||||
const event = {
|
||||
...data,
|
||||
signal,
|
||||
waitUntil: (p: Promise<unknown>): void => {
|
||||
if (Object.isFrozen(thenables)) {
|
||||
throw new Error('waitUntil can NOT be called asynchronously');
|
||||
}
|
||||
thenables.push(p);
|
||||
},
|
||||
} as T;
|
||||
try {
|
||||
deliver(event);
|
||||
} catch (error) {
|
||||
onUnexpectedError(error);
|
||||
return;
|
||||
}
|
||||
void Object.freeze(thenables);
|
||||
const settled = await Promise.allSettled(thenables);
|
||||
for (const result of settled) {
|
||||
if (result.status === 'rejected') {
|
||||
onUnexpectedError(result.reason);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -162,7 +162,7 @@ export class McpConfigStore extends Disposable implements IMcpConfigStore {
|
|||
await this.storage.write(CONFIG_SCOPE, MCP_CONFIG_KEY, textEncoder.encode(text), {
|
||||
atomic: true,
|
||||
});
|
||||
await this.writeEmitter.fireAsync({}, NO_ABORT);
|
||||
await this.writeEmitter.fireAsyncConcurrent({}, NO_ABORT);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -317,5 +317,33 @@ describe('McpConfigStore', () => {
|
|||
await mutation;
|
||||
expect(completed).toBe(true);
|
||||
});
|
||||
|
||||
it('starts asynchronous listeners concurrently before waiting for completion', async () => {
|
||||
let resolveStarted!: () => void;
|
||||
const started = new Promise<void>((resolve) => {
|
||||
resolveStarted = resolve;
|
||||
});
|
||||
let release!: () => void;
|
||||
const gate = new Promise<void>((resolve) => {
|
||||
release = resolve;
|
||||
});
|
||||
let secondStarted = false;
|
||||
store.onDidWrite((event) => {
|
||||
resolveStarted();
|
||||
event.waitUntil(gate);
|
||||
});
|
||||
store.onDidWrite(() => {
|
||||
secondStarted = true;
|
||||
});
|
||||
|
||||
const mutation = store.add(stdioServer('alpha'));
|
||||
await started;
|
||||
await Promise.resolve();
|
||||
const secondStartedBeforeRelease = secondStarted;
|
||||
release();
|
||||
await mutation;
|
||||
|
||||
expect(secondStartedBeforeRelease).toBe(true);
|
||||
});
|
||||
});
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue