import {
diagnosticLogger as diag,
logLaneDequeue,
logLaneEnqueue,
} from "../logging/diagnostic-runtime.js"; import { resolveGlobalSingleton } from "../shared/global-singleton.js"; import { CommandLane } from "./lanes.js"; /** *Dedicatederrortypethrownwhenaqueuedcommandisrejectedbecause *itslanewascleared.Callersthatfire-and-forgetenqueuedtaskscan *catch(orignore)thisspecifictypetoavoidunhandled-rejectionnoise.
*/
export class CommandLaneClearedError extends Error {
constructor(lane?: string) { super(lane ? `Command lane "${lane}" cleared` : "Command lane cleared"); this.name = "CommandLaneClearedError";
}
}
/** *Dedicatederrortypethrownwhenanewcommandisrejectedbecausethe *gatewayiscurrentlydrainingforrestart.
*/
export class GatewayDrainingError extends Error {
constructor() { super("Gateway is draining for restart; new tasks are not accepted"); this.name = "GatewayDrainingError";
}
}
// Minimal in-process queue to serialize command executions. // Default lane ("main") preserves the existing behavior. Additional lanes allow // low-risk parallelism (e.g. cron jobs) without interleaving stdin / logs for // the main auto-reply workflow.
function getQueueState() { const state = resolveGlobalSingleton(COMMAND_QUEUE_STATE_KEY, () => ({
gatewayDraining: false,
lanes: new Map<string, LaneState>(),
activeTaskWaiters: new Set<ActiveTaskWaiter>(),
nextTaskId: 1,
})); // Schema migration: the singleton may have been created by an older code // version (e.g. v2026.4.2) that did not include `activeTaskWaiters`. After // a SIGUSR1 in-process restart the new code inherits the stale object via // `resolveGlobalSingleton` because the Symbol key already exists on // globalThis. Patch the missing field so all downstream consumers see a // valid Set instead of `undefined`. if (!state.activeTaskWaiters) {
state.activeTaskWaiters = new Set<ActiveTaskWaiter>();
} return state;
}
function normalizeLane(lane: string): string { return lane.trim() || CommandLane.Main;
}
function getLaneDepth(state: LaneState): number { return state.queue.length + state.activeTaskIds.size;
}
function completeTask(state: LaneState, taskId: number, taskGeneration: number): boolean { if (taskGeneration !== state.generation) { returnfalse;
}
state.activeTaskIds.delete(taskId); returntrue;
}
function hasPendingActiveTasks(taskIds: Set<number>): boolean { const queueState = getQueueState(); for (const state of queueState.lanes.values()) { for (const taskId of state.activeTaskIds) { if (taskIds.has(taskId)) { returntrue;
}
}
} returnfalse;
}
function resolveActiveTaskWaiter(waiter: ActiveTaskWaiter, result: { drained: boolean }): void { const queueState = getQueueState(); if (!queueState.activeTaskWaiters.delete(waiter)) { return;
} if (waiter.timeout) {
clearTimeout(waiter.timeout);
}
waiter.resolve(result);
}
function notifyActiveTaskWaiters(): void { const queueState = getQueueState(); for (const waiter of Array.from(queueState.activeTaskWaiters)) { if (waiter.activeTaskIds.size === 0 || !hasPendingActiveTasks(waiter.activeTaskIds)) {
resolveActiveTaskWaiter(waiter, { drained: true });
}
}
}
function drainLane(lane: string) { const state = getLaneState(lane); if (state.draining) { if (state.activeTaskIds.size === 0 && state.queue.length > 0) {
diag.warn(
`drainLane blocked: lane=${lane} draining=true active=0 queue=${state.queue.length}`,
);
} return;
}
state.draining = true;
/** *Resetalllaneruntimestatetoidle.UsedafterSIGUSR1in-process *restartswhereinterruptedtasks'finallyblocksmaynotrun,leaving *staleactivetaskIDsthatpermanentlyblocknewworkfromdraining. * *Bumpslanegenerationandclearsexecutioncounterssostalecompletions *fromoldin-flighttasksareignored.Queuedentriesareintentionally *preserved—theyrepresentpendinguserworkthatshouldstillexecute *afterrestart. * *Afterresetting,drainsanylanesthatstillhavequeuedentriesso *preservedworkispumpedimmediatelyratherthanwaitingforafuture *`enqueueCommandInLane()`call(whichmaynevercome).
*/
export function resetAllLanes(): void { const queueState = getQueueState();
queueState.gatewayDraining = false; const lanesToDrain: string[] = []; for (const state of queueState.lanes.values()) {
state.generation += 1;
state.activeTaskIds.clear();
state.draining = false; if (state.queue.length > 0) {
lanesToDrain.push(state.lane);
}
} // Drain after the full reset pass so all lanes are in a clean state first. for (const lane of lanesToDrain) {
drainLane(lane);
}
notifyActiveTaskWaiters();
}
/** *Returnsthetotalnumberofactivelyexecutingtasksacrossalllanes *(excludesqueued-but-not-startedentries).
*/
export function getActiveTaskCount(): number { const queueState = getQueueState();
let total = 0; for (const s of queueState.lanes.values()) {
total += s.activeTaskIds.size;
} return total;
}
/** *Waitforallcurrentlyactivetasksacrossalllanestofinish. *Pollsatashortinterval;resolveswhennotasksareactiveor *when`timeoutMs`elapses(whichevercomesfirst). * *Newtasksenqueuedafterthiscallareignored—onlytasksthatare *alreadyexecutingarewaitedon.
*/
export function waitForActiveTasks(timeoutMs: number): Promise<{ drained: boolean }> { const queueState = getQueueState(); const activeAtStart = new Set<number>(); for (const state of queueState.lanes.values()) { for (const taskId of state.activeTaskIds) {
activeAtStart.add(taskId);
}
}
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.