Files
openclaw/src/utils/queue-helpers.ts
T
Peter Steinberger 0b8aabe864 docs: document auth profile failure policy contract (#89613)
* docs: document markdown marker renderer

* docs: document rendered markdown chunking

* docs: document markdown text chunking

* docs: document shared text chunking

* docs: document plugin text chunking exports

* docs: document avatar policy constants

* docs: document node match candidates

* docs: document scoped expiring id cache

* docs: document runtime import normalization

* docs: document string sample summaries

* docs: document session usage timeseries types

* docs: document session usage response types

* docs: document manifest frontmatter shapes

* docs: document channel route input metadata

* docs: document pair loop guard settings

* docs: document migration config patch helpers

* docs: document api provider registry

* docs: document tool call repair payloads

* docs: document plugin tool payload helpers

* docs: document lazy promise loader

* docs: document store writer queue state

* docs: document thread binding lifecycle

* docs: document concurrency helper contract

* docs: document gateway client info contract

* docs: document delivery context contracts

* docs: document secret ref defaults contract

* docs: document command gating contract

* docs: document avatar policy contract

* docs: document node match policy

* docs: document message channel normalization

* docs: document boolean parsing contract

* docs: document zod parse helpers

* docs: document direct dm guard policy

* docs: document fixed window limiter contract

* docs: document node presence event contract

* docs: document secret normalization contract

* docs: document progress draft line removal

* docs: document usage formatting contracts

* docs: document agent run status contract

* docs: document runtime import helpers

* docs: document provider utility ownership

* docs: document invalid config helpers

* docs: document json compat parser

* docs: document channel config metadata ownership

* docs: document channel logging helpers

* docs: document sender identity validation ownership

* docs: document string sampling helper

* docs: document global singleton helpers

* docs: document transcript tool helpers

* docs: document exec safe-bin normalization

* docs: document reaction level resolver

* docs: document account snapshot redaction boundary

* docs: document messaging target helpers

* docs: document thread binding messages

* docs: document conversation binding context

* docs: document conversation resolution helper

* docs: document owner display secret retention

* docs: document provider request config types

* docs: document skills config types

* docs: document memory config types

* docs: document imessage config types

* docs: document crestodian config types

* docs: document tools config policies

* docs: document shared config base types

* docs: document channel config contracts

* docs: document openclaw config state types

* docs: document model config contracts

* docs: document shared agent config types

* docs: document agent defaults config types

* docs: document secret input contracts

* docs: document auth config contracts

* docs: document gateway config contracts

* docs: document tool call stream repair contracts

* docs: document memory host facades

* docs: document llm core contracts

* docs: document markdown core contracts

* docs: document gateway connect error contracts

* docs: document gateway protocol primitives

* docs: document gateway frame schemas

* docs: document gateway device schemas

* docs: document gateway environment schemas

* docs: document gateway push schemas

* docs: document gateway plugin schemas

* docs: document gateway artifact schemas

* docs: document gateway command schemas

* docs: document gateway task schemas

* docs: document gateway exec approval schemas

* docs: document gateway secret schemas

* docs: document gateway config schemas

* docs: document gateway snapshot schemas

* docs: document gateway chat schemas

* docs: document gateway wizard schemas

* docs: document gateway node schemas

* docs: document gateway plugin approval schemas

* docs: document gateway talk schemas

* docs: document gateway agent schemas

* docs: document gateway session schemas

* docs: document gateway cron schemas

* docs: document gateway agent model skill schemas

* docs: document gateway skill proposal tool schemas

* docs: document gateway protocol registry

* docs: document gateway channel status schemas

* docs: document gateway schema regression tests

* docs: document gateway schema barrel

* docs: document gateway validator tests

* docs: document gateway primitive push tests

* docs: document gateway contract tests

* docs: document native protocol guard

* docs: document channel schema tests

* docs: document gateway protocol smoke tests

* docs: document gateway protocol entrypoint

* docs: document gateway protocol type exports

* docs: document gateway error codes

* docs: document protocol schema registry

* docs: document talk audio codec

* docs: document talk activation names

* docs: document talk consult questions

* docs: document talk consult tool

* docs: document talk run control contracts

* docs: document talk run control adapter

* docs: document talkback consult queue

* docs: document talk consult transcript guard

* docs: document talk fast context runtime

* docs: document forced talk consult coordinator

* docs: document talk output activity tracker

* docs: document talk event metrics

* docs: document talk diagnostics

* docs: document talk observability hook

* docs: document talk provider resolver

* docs: document talk provider registry

* docs: document talk runtime primitives

* docs: document talk consult controller logs

* docs: document channel identity helpers

* docs: document channel account allowlist helpers

* docs: document channel metadata draft controls

* docs: document channel ingress policy

* docs: document channel sender access gates

* docs: document channel catalog message contracts

* docs: document channel account plugin helpers

* docs: document configured binding helpers

* docs: document channel acp approval config helpers

* docs: document channel bundled config write helpers

* docs: document channel plugin utility contracts

* docs: document channel config access helpers

* docs: document channel message action helpers

* docs: document channel outbound runtime helpers

* docs: document channel pairing promotion helpers

* docs: document channel registry helpers

* docs: document channel setup wizard helpers

* docs: document channel lifecycle status helpers

* docs: document channel target thread helpers

* docs: document channel session binding helpers

* docs: document channel package module probes

* docs: document channel setup wizard contracts

* docs: document channel plugin API barrels

* docs: document channel contract test helpers

* docs: document channel core helpers

* docs: document small core facades

* docs: document provider runtime helpers

* docs: document persistence and realtime helpers

* docs: document mcp and state helpers

* docs: document tool planner contracts

* docs: document music generation runtime

* docs: document crestodian command flow

* docs: document utility helpers

* docs: document node host helpers

* docs: document transcript contracts

* docs: document trajectory export contracts

* docs: document image generation contracts

* docs: document routing helper contracts

* docs: document session helper contracts

* docs: document video generation contracts

* docs: document model catalog contracts

* docs: document proxy capture contracts

* docs: document status rendering contracts

* docs: document test helper contracts

* docs: document wizard setup contracts

* docs: document process contracts

* docs: document memory host sdk contracts

* docs: document tts contracts

* docs: document secrets runtime contracts

* docs: document shared helper contracts

* docs: document hook runtime contracts

* docs: document security audit contracts

* docs: document flow contracts

* docs: document media understanding contracts

* docs: document tui contracts

* docs: document logging contracts

* docs: document llm contracts

* docs: document cron contracts

* docs: document daemon contracts

* docs: document task contracts

* docs: document acp contracts

* docs: document test utility contracts

* docs: document skill contracts

* docs: document config contracts

* docs: document outbound infra contracts

* docs: document command analysis contracts

* docs: document provider usage infra contracts

* docs: document file safety infra contracts

* docs: document exec approval infra contracts

* docs: document gateway runtime infra contracts

* docs: document infra utility contracts

* docs: document infra queue storage contracts

* docs: document heartbeat infra contracts

* docs: document remaining infra contracts

* docs: document gateway auth contracts

* docs: document gateway display helpers

* docs: document gateway http helpers

* docs: document gateway node helpers

* docs: document gateway mcp helpers

* docs: document gateway support helpers

* docs: document gateway server runtime helpers

* docs: document gateway runtime bootstrap helpers

* docs: document gateway session events

* docs: document gateway utility helpers

* docs: document gateway talk helpers

* docs: document gateway helper contracts

* docs: document gateway server method helpers

* docs: document gateway server auth helpers

* docs: document gateway server tests

* docs: document gateway test helpers

* docs: document gateway node tests

* docs: document gateway channel tests

* docs: document gateway session tests

* docs: document gateway server startup tests

* docs: document gateway tool test helpers

* docs: document gateway server test helpers

* docs: document gateway server method tests

* docs: document remaining gateway tests

* docs: document plugin sdk public subpaths

* docs: document plugin sdk runtime helpers

* docs: document plugin sdk memory provider helpers

* docs: document plugin sdk runtime facades

* docs: document plugin sdk command approval helpers

* docs: document plugin sdk runtime types

* docs: document plugin sdk browser account helpers

* docs: document plugin sdk media memory helpers

* docs: document plugin sdk core tests

* docs: document plugin sdk contract helpers

* docs: document plugin sdk test helpers

* docs: document remaining plugin sdk tests

* docs: document cli utility helpers

* docs: document cli runtime helpers

* docs: document cli command registration helpers

* docs: document node cli helpers

* docs: document cli program registration

* docs: document message cli registration

* docs: document daemon cli helpers

* docs: document cli route parsers
2026-06-03 15:20:39 -07:00

281 lines
8.3 KiB
TypeScript

/**
* Shared queue overflow, debounce, and collection helpers.
*
* Queue owners use these helpers to cap pending work, summarize dropped items,
* debounce drains, and force individual collection when cross-channel ordering matters.
*/
/** Mutable summary state for a capped queue. */
export type QueueSummaryState = {
dropPolicy: "summarize" | "old" | "new";
droppedCount: number;
summaryLines: string[];
};
/** Queue overflow strategy. */
export type QueueDropPolicy = QueueSummaryState["dropPolicy"];
/** Generic capped queue state with shared overflow summary fields. */
export type QueueState<T> = QueueSummaryState & {
items: T[];
cap: number;
};
/** Clear accumulated overflow summary state after it has been emitted. */
export function clearQueueSummaryState(state: QueueSummaryState): void {
state.droppedCount = 0;
state.summaryLines = [];
}
/** Build a summary prompt preview without mutating the source queue state. */
export function previewQueueSummaryPrompt(params: {
state: QueueSummaryState;
noun: string;
title?: string;
}): string | undefined {
return buildQueueSummaryPrompt({
state: {
dropPolicy: params.state.dropPolicy,
droppedCount: params.state.droppedCount,
summaryLines: [...params.state.summaryLines],
},
noun: params.noun,
title: params.title,
});
}
/** Apply runtime queue settings while preserving previous values for omitted fields. */
export function applyQueueRuntimeSettings<TMode extends string>(params: {
target: {
mode: TMode;
debounceMs: number;
cap: number;
dropPolicy: QueueDropPolicy;
};
settings: {
mode: TMode;
debounceMs?: number;
cap?: number;
dropPolicy?: QueueDropPolicy;
};
}): void {
params.target.mode = params.settings.mode;
params.target.debounceMs =
typeof params.settings.debounceMs === "number"
? Math.max(0, params.settings.debounceMs)
: params.target.debounceMs;
params.target.cap =
typeof params.settings.cap === "number" && params.settings.cap > 0
? Math.floor(params.settings.cap)
: params.target.cap;
params.target.dropPolicy = params.settings.dropPolicy ?? params.target.dropPolicy;
}
/** Trim queue summary text to a bounded single-line preview. */
export function elideQueueText(text: string, limit = 140): string {
if (text.length <= limit) {
return text;
}
return `${text.slice(0, Math.max(0, limit - 1)).trimEnd()}…`;
}
/** Normalize whitespace and elide one dropped item for queue summaries. */
export function buildQueueSummaryLine(text: string, limit = 160): string {
const cleaned = text.replace(/\s+/g, " ").trim();
return elideQueueText(cleaned, limit);
}
/** Run optional duplicate detection before an item enters a queue. */
export function shouldSkipQueueItem<T>(params: {
item: T;
items: T[];
dedupe?: (item: T, items: T[]) => boolean;
}): boolean {
if (!params.dedupe) {
return false;
}
return params.dedupe(params.item, params.items);
}
/** Apply overflow policy before enqueueing another item. */
export function applyQueueDropPolicy<T>(params: {
queue: QueueState<T>;
summarize: (item: T) => string;
summaryLimit?: number;
onDrop?: (items: T[]) => void;
}): boolean {
const cap = params.queue.cap;
if (cap <= 0 || params.queue.items.length < cap) {
return true;
}
if (params.queue.dropPolicy === "new") {
return false;
}
const dropCount = params.queue.items.length - cap + 1;
const dropped = params.queue.items.splice(0, dropCount);
params.onDrop?.(dropped);
if (params.queue.dropPolicy === "summarize") {
for (const item of dropped) {
params.queue.droppedCount += 1;
params.queue.summaryLines.push(buildQueueSummaryLine(params.summarize(item)));
}
// Summary memory is bounded independently from the item cap to avoid prompt blowups.
const limit = Math.max(0, params.summaryLimit ?? cap);
while (params.queue.summaryLines.length > limit) {
params.queue.summaryLines.shift();
}
}
return true;
}
/** Wait until the queue has been quiet for its debounce window. */
export function waitForQueueDebounce(queue: {
debounceMs: number;
lastEnqueuedAt: number;
}): Promise<void> {
if (process.env.OPENCLAW_TEST_FAST === "1") {
// Tests use this escape hatch so debounce logic does not slow deterministic queue specs.
return Promise.resolve();
}
const debounceMs = Math.max(0, queue.debounceMs);
if (debounceMs <= 0) {
return Promise.resolve();
}
return new Promise<void>((resolve) => {
const check = () => {
const since = Date.now() - queue.lastEnqueuedAt;
if (since >= debounceMs) {
resolve();
return;
}
setTimeout(check, debounceMs - since);
};
check();
});
}
/** Mark one queue as draining unless another drain is already active. */
export function beginQueueDrain<T extends { draining: boolean }>(
map: Map<string, T>,
key: string,
): T | undefined {
const queue = map.get(key);
if (!queue || queue.draining) {
return undefined;
}
queue.draining = true;
return queue;
}
/** Run and remove the next queued item, returning false when empty. */
export async function drainNextQueueItem<T>(
items: T[],
run: (item: T) => Promise<void>,
): Promise<boolean> {
const next = items[0];
if (!next) {
return false;
}
await run(next);
items.shift();
return true;
}
/** Drain one item when collect mode requires individual processing. */
export async function drainCollectItemIfNeeded<T>(params: {
forceIndividualCollect: boolean;
isCrossChannel: boolean;
setForceIndividualCollect?: (next: boolean) => void;
items: T[];
run: (item: T) => Promise<void>;
}): Promise<"skipped" | "drained" | "empty"> {
if (!params.forceIndividualCollect && !params.isCrossChannel) {
return "skipped";
}
if (params.isCrossChannel) {
// Once cross-channel items appear, future collection stays individual to preserve ordering.
params.setForceIndividualCollect?.(true);
}
const drained = await drainNextQueueItem(params.items, params.run);
return drained ? "drained" : "empty";
}
/** Drain one collect step using mutable queue collection state. */
export async function drainCollectQueueStep<T>(params: {
collectState: { forceIndividualCollect: boolean };
isCrossChannel: boolean;
items: T[];
run: (item: T) => Promise<void>;
}): Promise<"skipped" | "drained" | "empty"> {
return await drainCollectItemIfNeeded({
forceIndividualCollect: params.collectState.forceIndividualCollect,
isCrossChannel: params.isCrossChannel,
setForceIndividualCollect: (next) => {
params.collectState.forceIndividualCollect = next;
},
items: params.items,
run: params.run,
});
}
/** Build and consume the queue overflow summary prompt. */
export function buildQueueSummaryPrompt(params: {
state: QueueSummaryState;
noun: string;
title?: string;
}): string | undefined {
if (params.state.dropPolicy !== "summarize" || params.state.droppedCount <= 0) {
return undefined;
}
const noun = params.noun;
const title =
params.title ??
`[Queue overflow] Dropped ${params.state.droppedCount} ${noun}${params.state.droppedCount === 1 ? "" : "s"} due to cap.`;
const lines = [title];
if (params.state.summaryLines.length > 0) {
lines.push("Summary:");
for (const line of params.state.summaryLines) {
lines.push(`- ${line}`);
}
}
clearQueueSummaryState(params.state);
return lines.join("\n");
}
/** Render a collect prompt from queued items and optional overflow summary. */
export function buildCollectPrompt<T>(params: {
title: string;
items: T[];
summary?: string;
renderItem: (item: T, index: number) => string;
}): string {
const blocks: string[] = [params.title];
if (params.summary) {
blocks.push(params.summary);
}
params.items.forEach((item, idx) => {
blocks.push(params.renderItem(item, idx));
});
return blocks.join("\n\n");
}
/** Return true when queued items span keys or explicitly mark cross-channel state. */
export function hasCrossChannelItems<T>(
items: T[],
resolveKey: (item: T) => { key?: string; cross?: boolean },
): boolean {
const keys = new Set<string>();
for (const item of items) {
const resolved = resolveKey(item);
if (resolved.cross) {
return true;
}
if (!resolved.key) {
continue;
}
keys.add(resolved.key);
}
return keys.size > 1;
}