Files
openclaw/src/gateway/node-pending-work.ts
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

232 lines
7.4 KiB
TypeScript

import { randomUUID } from "node:crypto";
import {
asDateTimestampMs,
isFutureDateTimestampMs,
resolveDateTimestampMs,
resolveExpiresAtMsFromDurationMs,
} from "@openclaw/normalization-core/number-coercion";
// Pending node work is an in-memory per-node queue for gateway prompts such as
// status/location requests. Nodes drain it opportunistically and acknowledge
// item ids after handling them.
const NODE_PENDING_WORK_TYPES = ["status.request", "location.request"] as const;
/** Work item types that connected nodes understand today. */
export type NodePendingWorkType = (typeof NODE_PENDING_WORK_TYPES)[number];
const NODE_PENDING_WORK_PRIORITIES = ["default", "normal", "high"] as const;
/** Priority labels used for pending work drain ordering. */
export type NodePendingWorkPriority = (typeof NODE_PENDING_WORK_PRIORITIES)[number];
type NodePendingWorkItem = {
id: string;
type: NodePendingWorkType;
priority: NodePendingWorkPriority;
createdAtMs: number;
expiresAtMs: number | null;
payload?: Record<string, unknown>;
};
type NodePendingWorkState = {
revision: number;
itemsById: Map<string, NodePendingWorkItem>;
};
type DrainOptions = {
maxItems?: number;
includeDefaultStatus?: boolean;
nowMs?: number;
};
type DrainResult = {
revision: number;
items: NodePendingWorkItem[];
hasMore: boolean;
};
const DEFAULT_STATUS_ITEM_ID = "baseline-status";
const DEFAULT_STATUS_PRIORITY: NodePendingWorkPriority = "default";
const DEFAULT_PRIORITY: NodePendingWorkPriority = "normal";
const DEFAULT_MAX_ITEMS = 4;
const MAX_ITEMS = 10;
const PRIORITY_RANK: Record<NodePendingWorkPriority, number> = {
high: 3,
normal: 2,
default: 1,
};
const stateByNodeId = new Map<string, NodePendingWorkState>();
function getOrCreateState(nodeId: string): NodePendingWorkState {
let state = stateByNodeId.get(nodeId);
if (!state) {
state = {
revision: 0,
itemsById: new Map(),
};
stateByNodeId.set(nodeId, state);
}
return state;
}
function pruneExpired(state: NodePendingWorkState, nowMs: number): boolean {
const validNowMs = asDateTimestampMs(nowMs);
if (validNowMs === undefined) {
return false;
}
let changed = false;
for (const [id, item] of state.itemsById) {
if (
item.expiresAtMs !== null &&
!isFutureDateTimestampMs(item.expiresAtMs, { nowMs: validNowMs })
) {
state.itemsById.delete(id);
changed = true;
}
}
if (changed) {
state.revision += 1;
}
return changed;
}
function pruneStateIfEmpty(nodeId: string, state: NodePendingWorkState) {
if (state.itemsById.size === 0) {
stateByNodeId.delete(nodeId);
}
}
function sortedItems(state: NodePendingWorkState): NodePendingWorkItem[] {
// Higher priority wins, then older work, then id for deterministic paging.
return [...state.itemsById.values()].toSorted((a, b) => {
const priorityDelta = PRIORITY_RANK[b.priority] - PRIORITY_RANK[a.priority];
if (priorityDelta !== 0) {
return priorityDelta;
}
if (a.createdAtMs !== b.createdAtMs) {
return a.createdAtMs - b.createdAtMs;
}
return a.id.localeCompare(b.id);
});
}
function makeBaselineStatusItem(nowMs: number): NodePendingWorkItem {
return {
id: DEFAULT_STATUS_ITEM_ID,
type: "status.request",
priority: DEFAULT_STATUS_PRIORITY,
createdAtMs: resolveDateTimestampMs(nowMs),
expiresAtMs: null,
};
}
function resolvePendingWorkExpiresAtMs(expiresInMs: unknown, nowMs: number): number | null {
if (typeof expiresInMs !== "number" || !Number.isFinite(expiresInMs)) {
return null;
}
return resolveExpiresAtMsFromDurationMs(Math.max(1_000, Math.trunc(expiresInMs)), { nowMs }) ?? 0;
}
export function enqueueNodePendingWork(params: {
nodeId: string;
type: NodePendingWorkType;
priority?: NodePendingWorkPriority;
expiresInMs?: number;
payload?: Record<string, unknown>;
}): { revision: number; item: NodePendingWorkItem; deduped: boolean } {
const nodeId = params.nodeId.trim();
if (!nodeId) {
throw new Error("nodeId required");
}
const rawNowMs = Date.now();
const nowMs = resolveDateTimestampMs(rawNowMs);
const state = getOrCreateState(nodeId);
pruneExpired(state, nowMs);
// Keep one outstanding item per type so repeated status/location requests
// collapse until the node has a chance to drain and acknowledge them.
const existing = [...state.itemsById.values()].find((item) => item.type === params.type);
if (existing) {
return { revision: state.revision, item: existing, deduped: true };
}
const item: NodePendingWorkItem = {
id: randomUUID(),
type: params.type,
priority: params.priority ?? DEFAULT_PRIORITY,
createdAtMs: nowMs,
expiresAtMs: resolvePendingWorkExpiresAtMs(params.expiresInMs, rawNowMs),
...(params.payload ? { payload: params.payload } : {}),
};
state.itemsById.set(item.id, item);
state.revision += 1;
return { revision: state.revision, item, deduped: false };
}
/** Drains pending work for a node, including a baseline status request unless disabled. */
export function drainNodePendingWork(nodeId: string, opts: DrainOptions = {}): DrainResult {
const normalizedNodeId = nodeId.trim();
if (!normalizedNodeId) {
return { revision: 0, items: [], hasMore: false };
}
const nowMs = resolveDateTimestampMs(opts.nowMs ?? Date.now());
const state = stateByNodeId.get(normalizedNodeId);
if (state) {
pruneExpired(state, nowMs);
pruneStateIfEmpty(normalizedNodeId, state);
}
const revision = state?.revision ?? 0;
const maxItems = Math.min(MAX_ITEMS, Math.max(1, Math.trunc(opts.maxItems ?? DEFAULT_MAX_ITEMS)));
const explicitItems = state ? sortedItems(state) : [];
const items = explicitItems.slice(0, maxItems);
const hasExplicitStatus = explicitItems.some((item) => item.type === "status.request");
const includeBaseline = opts.includeDefaultStatus !== false && !hasExplicitStatus;
if (includeBaseline && items.length < maxItems) {
items.push(makeBaselineStatusItem(nowMs));
}
const explicitReturnedCount = items.filter((item) => item.id !== DEFAULT_STATUS_ITEM_ID).length;
const baselineIncluded = items.some((item) => item.id === DEFAULT_STATUS_ITEM_ID);
return {
revision,
items,
hasMore: explicitItems.length > explicitReturnedCount || (includeBaseline && !baselineIncluded),
};
}
/** Acknowledges completed pending-work ids and advances the node revision. */
export function acknowledgeNodePendingWork(params: { nodeId: string; itemIds: string[] }): {
revision: number;
removedItemIds: string[];
} {
const nodeId = params.nodeId.trim();
if (!nodeId) {
return { revision: 0, removedItemIds: [] };
}
const state = stateByNodeId.get(nodeId);
if (!state) {
return { revision: 0, removedItemIds: [] };
}
const removedItemIds: string[] = [];
for (const itemId of params.itemIds) {
const trimmedId = itemId.trim();
if (!trimmedId || trimmedId === DEFAULT_STATUS_ITEM_ID) {
continue;
}
if (state.itemsById.delete(trimmedId)) {
removedItemIds.push(trimmedId);
}
}
if (removedItemIds.length > 0) {
state.revision += 1;
}
pruneStateIfEmpty(nodeId, state);
return { revision: state.revision, removedItemIds };
}
/** Clears all pending work state for tests. */
export function resetNodePendingWorkForTests() {
stateByNodeId.clear();
}
/** Returns the number of node queues retained in memory for tests. */
export function getNodePendingWorkStateCountForTests(): number {
return stateByNodeId.size;
}