File size: 4,823 Bytes
3144483 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 | // Shared command-queue runtime state, split out of command-queue.ts so the
// capacity-group policy can read lane state without importing the queue itself.
import { resolveGlobalSingleton } from "../shared/global-singleton.js";
import type { CommandQueueEnqueueOptions } from "./command-queue.types.js";
import { CommandLane } from "./lanes.js";
export type CommandLaneTaskMarker = Readonly<{
lane: string;
taskId: number;
generation: number;
}>;
export type QueuePriority = -1 | 0 | 1;
export type QueueEntry = {
queued?: true;
previous?: QueueEntry;
next?: QueueEntry;
task: (marker: CommandLaneTaskMarker) => Promise<unknown>;
resolve: (value: unknown) => void;
reject: (reason?: unknown) => void;
enqueuedAt: number;
sequence: number;
priority: QueuePriority;
warnAfterMs: number;
queuedAheadAtEnqueue: number;
activeAheadAtEnqueue: number;
taskTimeoutMs?: number;
taskTimeoutProgressAtMs?: () => number | undefined;
taskTimeoutSubscribe?: CommandQueueEnqueueOptions["taskTimeoutSubscribe"];
taskTimeoutAbortSignal?: AbortSignal;
taskTimeoutAbortGraceMs?: number;
taskTimeoutReleaseSignal?: AbortSignal;
onWait?: (waitMs: number, queuedAhead: number) => void;
releaseQueuedAbort?: () => void;
};
type QueueFifo = {
head: QueueEntry | undefined;
tail: QueueEntry | undefined;
length: number;
};
/** Three fixed FIFO lists, one for each supported priority. */
type LaneQueue = {
background: QueueFifo;
normal: QueueFifo;
foreground: QueueFifo;
length: number;
};
export type LaneState = {
lane: string;
queue: LaneQueue;
activeTaskIds: Set<number>;
maxConcurrent: number;
draining: boolean;
generation: number;
};
export type LaneGroupState = {
group: string;
budget: number;
members: Set<string>;
reservations: Map<string, number>;
};
function createQueueFifo(): QueueFifo {
return { head: undefined, tail: undefined, length: 0 };
}
export function createLaneQueue(): LaneQueue {
return {
background: createQueueFifo(),
normal: createQueueFifo(),
foreground: createQueueFifo(),
length: 0,
};
}
function getPriorityFifo(queue: LaneQueue, priority: QueuePriority): QueueFifo {
switch (priority) {
case 1:
return queue.foreground;
case -1:
return queue.background;
default:
return queue.normal;
}
}
/** Append to one of three fixed priority FIFOs and return the queued work ahead. */
export function enqueueLaneQueue(queue: LaneQueue, entry: QueueEntry): number {
const fifo = getPriorityFifo(queue, entry.priority);
const queuedAhead =
fifo.length +
(entry.priority <= 0 ? queue.foreground.length : 0) +
(entry.priority < 0 ? queue.normal.length : 0);
entry.queued = true;
entry.previous = fifo.tail;
entry.next = undefined;
if (fifo.tail) {
fifo.tail.next = entry;
} else {
fifo.head = entry;
}
fifo.tail = entry;
fifo.length += 1;
queue.length += 1;
return queuedAhead;
}
export function peekLaneQueue(queue: LaneQueue): QueueEntry | undefined {
return queue.foreground.head ?? queue.normal.head ?? queue.background.head;
}
export function dequeueLaneQueue(queue: LaneQueue): QueueEntry | undefined {
const entry = peekLaneQueue(queue);
if (entry) {
removeLaneQueueEntry(queue, entry);
}
return entry;
}
/** Unlink without scanning successors, including when one abort cancels an entire backlog. */
export function removeLaneQueueEntry(queue: LaneQueue, entry: QueueEntry): boolean {
if (!entry.queued) {
return false;
}
const fifo = getPriorityFifo(queue, entry.priority);
if (entry.previous) {
entry.previous.next = entry.next;
} else {
fifo.head = entry.next;
}
if (entry.next) {
entry.next.previous = entry.previous;
} else {
fifo.tail = entry.previous;
}
entry.queued = undefined;
fifo.length -= 1;
queue.length -= 1;
// A completed entry must not retain its neighbours or expose stale membership
// if listener cleanup reenters the queue.
const releaseQueuedAbort = entry.releaseQueuedAbort;
entry.previous = undefined;
entry.next = undefined;
entry.releaseQueuedAbort = undefined;
releaseQueuedAbort?.();
return true;
}
/**
* Keep queue runtime state on globalThis so every bundled entry/chunk shares
* the same lanes, counters, and draining flag in production builds.
*/
const COMMAND_QUEUE_STATE_KEY = Symbol.for("openclaw.commandQueueState");
export function getQueueState() {
return resolveGlobalSingleton(COMMAND_QUEUE_STATE_KEY, () => ({
lanes: new Map<string, LaneState>(),
nextTaskId: 1,
nextQueueSequence: 1,
laneGroups: new Map<string, LaneGroupState>(),
laneGroupByLane: new Map<string, string>(),
}));
}
export function normalizeLane(lane: string): string {
return lane.trim() || CommandLane.Main;
}
|