From ce6234da3fbbeead758748e285f624344b2577ed Mon Sep 17 00:00:00 2001 From: weishu Date: Tue, 30 Dec 2025 13:26:33 +0800 Subject: [PATCH] fix: prevent race condition in MessageQueue2 wait logic --- cli/src/utils/MessageQueue2.ts | 54 ++++++++++++++++++---------------- 1 file changed, 29 insertions(+), 25 deletions(-) diff --git a/cli/src/utils/MessageQueue2.ts b/cli/src/utils/MessageQueue2.ts index 45f254ed..6ba5fcdd 100644 --- a/cli/src/utils/MessageQueue2.ts +++ b/cli/src/utils/MessageQueue2.ts @@ -288,22 +288,18 @@ export class MessageQueue2 { */ private waitForMessages(abortSignal?: AbortSignal): Promise { return new Promise((resolve) => { + let settled = false; let abortHandler: (() => void) | null = null; + let waiterFunc: (hasMessages: boolean) => void; - // Set up abort handler - if (abortSignal) { - abortHandler = () => { - logger.debug('[MessageQueue2] Wait aborted'); - // Clear waiter if it's still set - if (this.waiter === waiterFunc) { - this.waiter = null; - } - resolve(false); - }; - abortSignal.addEventListener('abort', abortHandler); - } - - const waiterFunc = (hasMessages: boolean) => { + const finish = (hasMessages: boolean) => { + if (settled) { + return; + } + settled = true; + if (this.waiter === waiterFunc) { + this.waiter = null; + } // Clean up abort handler if (abortHandler && abortSignal) { abortSignal.removeEventListener('abort', abortHandler); @@ -311,26 +307,34 @@ export class MessageQueue2 { resolve(hasMessages); }; + waiterFunc = (hasMessages: boolean) => { + finish(hasMessages); + }; + + // Set up abort handler + if (abortSignal) { + abortHandler = () => { + logger.debug('[MessageQueue2] Wait aborted'); + finish(false); + }; + abortSignal.addEventListener('abort', abortHandler); + } + + // Set the waiter before checking the queue to avoid missed notifications + this.waiter = waiterFunc; + // Check again in case messages arrived or queue closed while setting up if (this.queue.length > 0) { - if (abortHandler && abortSignal) { - abortSignal.removeEventListener('abort', abortHandler); - } - resolve(true); + finish(true); return; } if (this.closed || abortSignal?.aborted) { - if (abortHandler && abortSignal) { - abortSignal.removeEventListener('abort', abortHandler); - } - resolve(false); + finish(false); return; } - // Set the waiter - this.waiter = waiterFunc; logger.debug('[MessageQueue2] Waiting for messages...'); }); } -} \ No newline at end of file +}