@@ -53,11 +53,13 @@ type LaneState = {
5353 */
5454const COMMAND_QUEUE_STATE_KEY = Symbol . for ( "openclaw.commandQueueState" ) ;
5555
56- const queueState = resolveGlobalSingleton ( COMMAND_QUEUE_STATE_KEY , ( ) => ( {
57- gatewayDraining : false ,
58- lanes : new Map < string , LaneState > ( ) ,
59- nextTaskId : 1 ,
60- } ) ) ;
56+ function getQueueState ( ) {
57+ return resolveGlobalSingleton ( COMMAND_QUEUE_STATE_KEY , ( ) => ( {
58+ gatewayDraining : false ,
59+ lanes : new Map < string , LaneState > ( ) ,
60+ nextTaskId : 1 ,
61+ } ) ) ;
62+ }
6163
6264function normalizeLane ( lane : string ) : string {
6365 return lane . trim ( ) || CommandLane . Main ;
@@ -68,6 +70,7 @@ function getLaneDepth(state: LaneState): number {
6870}
6971
7072function getLaneState ( lane : string ) : LaneState {
73+ const queueState = getQueueState ( ) ;
7174 const existing = queueState . lanes . get ( lane ) ;
7275 if ( existing ) {
7376 return existing ;
@@ -120,7 +123,7 @@ function drainLane(lane: string) {
120123 ) ;
121124 }
122125 logLaneDequeue ( lane , waitedMs , state . queue . length ) ;
123- const taskId = queueState . nextTaskId ++ ;
126+ const taskId = getQueueState ( ) . nextTaskId ++ ;
124127 const taskGeneration = state . generation ;
125128 state . activeTaskIds . add ( taskId ) ;
126129 void ( async ( ) => {
@@ -163,7 +166,7 @@ function drainLane(lane: string) {
163166 * `GatewayDrainingError` instead of being silently killed on shutdown.
164167 */
165168export function markGatewayDraining ( ) : void {
166- queueState . gatewayDraining = true ;
169+ getQueueState ( ) . gatewayDraining = true ;
167170}
168171
169172export function setCommandLaneConcurrency ( lane : string , maxConcurrent : number ) {
@@ -181,6 +184,7 @@ export function enqueueCommandInLane<T>(
181184 onWait ?: ( waitMs : number , queuedAhead : number ) => void ;
182185 } ,
183186) : Promise < T > {
187+ const queueState = getQueueState ( ) ;
184188 if ( queueState . gatewayDraining ) {
185189 return Promise . reject ( new GatewayDrainingError ( ) ) ;
186190 }
@@ -213,7 +217,7 @@ export function enqueueCommand<T>(
213217
214218export function getQueueSize ( lane : string = CommandLane . Main ) {
215219 const resolved = normalizeLane ( lane ) ;
216- const state = queueState . lanes . get ( resolved ) ;
220+ const state = getQueueState ( ) . lanes . get ( resolved ) ;
217221 if ( ! state ) {
218222 return 0 ;
219223 }
@@ -222,15 +226,15 @@ export function getQueueSize(lane: string = CommandLane.Main) {
222226
223227export function getTotalQueueSize ( ) {
224228 let total = 0 ;
225- for ( const s of queueState . lanes . values ( ) ) {
229+ for ( const s of getQueueState ( ) . lanes . values ( ) ) {
226230 total += getLaneDepth ( s ) ;
227231 }
228232 return total ;
229233}
230234
231235export function clearCommandLane ( lane : string = CommandLane . Main ) {
232236 const cleaned = normalizeLane ( lane ) ;
233- const state = queueState . lanes . get ( cleaned ) ;
237+ const state = getQueueState ( ) . lanes . get ( cleaned ) ;
234238 if ( ! state ) {
235239 return 0 ;
236240 }
@@ -257,6 +261,7 @@ export function clearCommandLane(lane: string = CommandLane.Main) {
257261 * `enqueueCommandInLane()` call (which may never come).
258262 */
259263export function resetAllLanes ( ) : void {
264+ const queueState = getQueueState ( ) ;
260265 queueState . gatewayDraining = false ;
261266 const lanesToDrain : string [ ] = [ ] ;
262267 for ( const state of queueState . lanes . values ( ) ) {
@@ -278,6 +283,7 @@ export function resetAllLanes(): void {
278283 * (excludes queued-but-not-started entries).
279284 */
280285export function getActiveTaskCount ( ) : number {
286+ const queueState = getQueueState ( ) ;
281287 let total = 0 ;
282288 for ( const s of queueState . lanes . values ( ) ) {
283289 total += s . activeTaskIds . size ;
@@ -297,6 +303,7 @@ export function waitForActiveTasks(timeoutMs: number): Promise<{ drained: boolea
297303 // Keep shutdown/drain checks responsive without busy looping.
298304 const POLL_INTERVAL_MS = 50 ;
299305 const deadline = Date . now ( ) + timeoutMs ;
306+ const queueState = getQueueState ( ) ;
300307 const activeAtStart = new Set < number > ( ) ;
301308 for ( const state of queueState . lanes . values ( ) ) {
302309 for ( const taskId of state . activeTaskIds ) {
0 commit comments