Skip to content

Commit a9da52d

Browse files
committed
refactor(core): make event and queue state lazy
1 parent f6b3377 commit a9da52d

2 files changed

Lines changed: 31 additions & 18 deletions

File tree

src/infra/agent-events.ts

Lines changed: 14 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -29,16 +29,19 @@ type AgentEventState = {
2929

3030
const AGENT_EVENT_STATE_KEY = Symbol.for("openclaw.agentEvents.state");
3131

32-
const state = resolveGlobalSingleton<AgentEventState>(AGENT_EVENT_STATE_KEY, () => ({
33-
seqByRun: new Map<string, number>(),
34-
listeners: new Set<(evt: AgentEventPayload) => void>(),
35-
runContextById: new Map<string, AgentRunContext>(),
36-
}));
32+
function getAgentEventState(): AgentEventState {
33+
return resolveGlobalSingleton<AgentEventState>(AGENT_EVENT_STATE_KEY, () => ({
34+
seqByRun: new Map<string, number>(),
35+
listeners: new Set<(evt: AgentEventPayload) => void>(),
36+
runContextById: new Map<string, AgentRunContext>(),
37+
}));
38+
}
3739

3840
export function registerAgentRunContext(runId: string, context: AgentRunContext) {
3941
if (!runId) {
4042
return;
4143
}
44+
const state = getAgentEventState();
4245
const existing = state.runContextById.get(runId);
4346
if (!existing) {
4447
state.runContextById.set(runId, { ...context });
@@ -59,18 +62,19 @@ export function registerAgentRunContext(runId: string, context: AgentRunContext)
5962
}
6063

6164
export function getAgentRunContext(runId: string) {
62-
return state.runContextById.get(runId);
65+
return getAgentEventState().runContextById.get(runId);
6366
}
6467

6568
export function clearAgentRunContext(runId: string) {
66-
state.runContextById.delete(runId);
69+
getAgentEventState().runContextById.delete(runId);
6770
}
6871

6972
export function resetAgentRunContextForTest() {
70-
state.runContextById.clear();
73+
getAgentEventState().runContextById.clear();
7174
}
7275

7376
export function emitAgentEvent(event: Omit<AgentEventPayload, "seq" | "ts">) {
77+
const state = getAgentEventState();
7478
const nextSeq = (state.seqByRun.get(event.runId) ?? 0) + 1;
7579
state.seqByRun.set(event.runId, nextSeq);
7680
const context = state.runContextById.get(event.runId);
@@ -88,10 +92,12 @@ export function emitAgentEvent(event: Omit<AgentEventPayload, "seq" | "ts">) {
8892
}
8993

9094
export function onAgentEvent(listener: (evt: AgentEventPayload) => void) {
95+
const state = getAgentEventState();
9196
return registerListener(state.listeners, listener);
9297
}
9398

9499
export function resetAgentEventsForTest() {
100+
const state = getAgentEventState();
95101
state.seqByRun.clear();
96102
state.listeners.clear();
97103
state.runContextById.clear();

src/process/command-queue.ts

Lines changed: 17 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -53,11 +53,13 @@ type LaneState = {
5353
*/
5454
const 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

6264
function normalizeLane(lane: string): string {
6365
return lane.trim() || CommandLane.Main;
@@ -68,6 +70,7 @@ function getLaneDepth(state: LaneState): number {
6870
}
6971

7072
function 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
*/
165168
export function markGatewayDraining(): void {
166-
queueState.gatewayDraining = true;
169+
getQueueState().gatewayDraining = true;
167170
}
168171

169172
export 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

214218
export 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

223227
export 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

231235
export 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
*/
259263
export 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
*/
280285
export 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

Comments
 (0)