Some checks failed
ClawSweeper Dispatch / dispatch (push) Has been cancelled
CodeQL / Security High (actions) (push) Has been cancelled
CodeQL / Security High (channel-runtime-boundary) (push) Has been cancelled
CodeQL / Security High (core-auth-secrets) (push) Has been cancelled
CodeQL / Security High (mcp-process-tool-boundary) (push) Has been cancelled
CodeQL / Security High (network-ssrf-boundary) (push) Has been cancelled
CodeQL / Security High (plugin-trust-boundary) (push) Has been cancelled
CodeQL / Security High (process-exec-boundary) (push) Has been cancelled
Docs Sync Publish Repo / sync-publish-repo (push) Has been cancelled
Docs / docs (push) Has been cancelled
OpenClaw Stable Main Closeout / Resolve stable release closeout inputs (push) Has been cancelled
OpenClaw Stable Main Closeout / Verify stable main closeout (push) Has been cancelled
Workflow Sanity / no-tabs (push) Has been cancelled
Workflow Sanity / actionlint (push) Has been cancelled
Workflow Sanity / generated-doc-baselines (push) Has been cancelled
CI / runner-admission (push) Has been cancelled
CI / preflight (push) Has been cancelled
CI / security-fast (push) Has been cancelled
CI / pnpm-store-warmup (push) Has been cancelled
CI / build-artifacts (push) Has been cancelled
CI / native-i18n (push) Has been cancelled
CI / ${{ matrix.check_name }} (push) Has been cancelled
CI / ${{ matrix.checkName }} (push) Has been cancelled
CI / checks-node-compat-node22 (push) Has been cancelled
CI / check-bundled-channel-config-metadata (push) Has been cancelled
CI / check-dependencies (push) Has been cancelled
CI / check-guards (push) Has been cancelled
CI / check-lint (push) Has been cancelled
CI / check-prod-types (push) Has been cancelled
CI / check-shrinkwrap (push) Has been cancelled
CI / check-test-types (push) Has been cancelled
CI / check-additional-boundaries-a (push) Has been cancelled
CI / check-additional-boundaries-bcd (push) Has been cancelled
CI / check-additional-extension-bundled (push) Has been cancelled
CI / check-additional-extension-channels (push) Has been cancelled
CI / check-additional-extension-package-boundary (push) Has been cancelled
CI / check-additional-runtime-topology-architecture (push) Has been cancelled
CI / check-session-accessor-boundary (push) Has been cancelled
CI / check-session-transcript-reader-boundary (push) Has been cancelled
CI / check-docs (push) Has been cancelled
CI / skills-python (push) Has been cancelled
CI / macos-swift (push) Has been cancelled
CI / ios-build (push) Has been cancelled
CI / ci-timings-summary (push) Has been cancelled
Native App Locale Refresh / Refresh native fa (push) Has been cancelled
Native App Locale Refresh / Refresh native fr (push) Has been cancelled
Native App Locale Refresh / Refresh native hi (push) Has been cancelled
Native App Locale Refresh / Refresh native id (push) Has been cancelled
Native App Locale Refresh / Refresh native it (push) Has been cancelled
Native App Locale Refresh / Refresh native ja-JP (push) Has been cancelled
Control UI Locale Refresh / plan (push) Has been cancelled
Control UI Locale Refresh / Refresh ${{ matrix.locale }} (push) Has been cancelled
Control UI Locale Refresh / Commit control UI locale refresh (push) Has been cancelled
Live Media Runner Image / Build live media runner image (push) Has been cancelled
Native App Locale Refresh / Refresh native ar (push) Has been cancelled
Native App Locale Refresh / Refresh native de (push) Has been cancelled
Native App Locale Refresh / Refresh native es (push) Has been cancelled
Native App Locale Refresh / Refresh native ko (push) Has been cancelled
Native App Locale Refresh / Refresh native nl (push) Has been cancelled
Native App Locale Refresh / Refresh native pl (push) Has been cancelled
Native App Locale Refresh / Refresh native pt-BR (push) Has been cancelled
Native App Locale Refresh / Refresh native ru (push) Has been cancelled
Native App Locale Refresh / Refresh native sv (push) Has been cancelled
Native App Locale Refresh / Refresh native th (push) Has been cancelled
Native App Locale Refresh / Refresh native tr (push) Has been cancelled
Native App Locale Refresh / Refresh native uk (push) Has been cancelled
Native App Locale Refresh / Refresh native vi (push) Has been cancelled
Native App Locale Refresh / Refresh native zh-CN (push) Has been cancelled
Native App Locale Refresh / Refresh native zh-TW (push) Has been cancelled
Native App Locale Refresh / Commit native locale refresh (push) Has been cancelled
Plugin Init Scaffold Validation / Validate provider scaffold (push) Has been cancelled
Plugin NPM Release / preview_plugins_npm (push) Has been cancelled
Plugin NPM Release / Validate release publish approval (push) Has been cancelled
Plugin NPM Release / preview_plugin_pack (push) Has been cancelled
Plugin NPM Release / publish_plugins_npm (push) Has been cancelled
Sandbox Common Smoke / sandbox-common-smoke (push) Has been cancelled
Website Installer Sync / static (push) Has been cancelled
Website Installer Sync / linux-docker (push) Has been cancelled
Website Installer Sync / macos-installer (push) Has been cancelled
Website Installer Sync / windows-installer (push) Has been cancelled
Website Installer Sync / sync-website (push) Has been cancelled
Adolf is a fork/vendored clone of github.com/openclaw/openclaw (v2026.6.11), free to diverge. Tree copied sans upstream .git; upstream remote added for future syncs. Node pinned to 24 (.nvmrc); engines already require >=22.19. Preserves docs/ARCHITECTURE.md. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01LeqyaxJF2nbRXJtae2kNB2
164 lines
4.7 KiB
TypeScript
164 lines
4.7 KiB
TypeScript
// Telegram plugin module implements account throttler behavior.
|
|
import { parseStrictInteger } from "openclaw/plugin-sdk/number-runtime";
|
|
import { apiThrottler } from "./bot.runtime.js";
|
|
|
|
type ApiThrottlerTransformer = ReturnType<typeof apiThrottler>;
|
|
type TelegramApiPayload = {
|
|
chat_id?: unknown;
|
|
direct_messages_topic_id?: unknown;
|
|
message_id?: unknown;
|
|
message_thread_id?: unknown;
|
|
};
|
|
type QueuedApiRequest<T> = {
|
|
run: () => Promise<T>;
|
|
resolve: (value: T) => void;
|
|
reject: (err: unknown) => void;
|
|
};
|
|
|
|
class GroupFairQueue {
|
|
private readonly lanes = new Map<string, Array<QueuedApiRequest<unknown>>>();
|
|
private laneOrder: string[] = [];
|
|
private nextLaneIndex = 0;
|
|
private running = false;
|
|
|
|
enqueue<T>(laneKey: string, run: () => Promise<T>): Promise<T> {
|
|
return new Promise<T>((resolve, reject) => {
|
|
const request: QueuedApiRequest<unknown> = {
|
|
run,
|
|
resolve: resolve as (value: unknown) => void,
|
|
reject,
|
|
};
|
|
const existing = this.lanes.get(laneKey);
|
|
if (existing) {
|
|
existing.push(request);
|
|
} else {
|
|
this.lanes.set(laneKey, [request]);
|
|
this.laneOrder.push(laneKey);
|
|
}
|
|
this.start();
|
|
});
|
|
}
|
|
|
|
private start(): void {
|
|
if (this.running) {
|
|
return;
|
|
}
|
|
this.running = true;
|
|
void this.drain();
|
|
}
|
|
|
|
private async drain(): Promise<void> {
|
|
try {
|
|
while (true) {
|
|
const request = this.takeNext();
|
|
if (!request) {
|
|
return;
|
|
}
|
|
try {
|
|
request.resolve(await request.run());
|
|
} catch (err) {
|
|
request.reject(err);
|
|
}
|
|
}
|
|
} finally {
|
|
this.running = false;
|
|
if (this.laneOrder.length > 0) {
|
|
this.start();
|
|
}
|
|
}
|
|
}
|
|
|
|
private takeNext(): QueuedApiRequest<unknown> | undefined {
|
|
for (let remaining = this.laneOrder.length; remaining > 0; remaining -= 1) {
|
|
this.nextLaneIndex %= this.laneOrder.length;
|
|
const laneKey = this.laneOrder[this.nextLaneIndex];
|
|
const queue = this.lanes.get(laneKey);
|
|
if (!queue || queue.length === 0) {
|
|
this.lanes.delete(laneKey);
|
|
this.laneOrder.splice(this.nextLaneIndex, 1);
|
|
if (this.laneOrder.length === 0) {
|
|
this.nextLaneIndex = 0;
|
|
return undefined;
|
|
}
|
|
continue;
|
|
}
|
|
|
|
const request = queue.shift();
|
|
this.nextLaneIndex += 1;
|
|
return request;
|
|
}
|
|
return undefined;
|
|
}
|
|
}
|
|
|
|
const throttlerByToken = new Map<string, ApiThrottlerTransformer>();
|
|
|
|
function readNumericId(value: unknown): number | undefined {
|
|
return parseStrictInteger(value);
|
|
}
|
|
|
|
function readPayload(payload: unknown): TelegramApiPayload | undefined {
|
|
return payload && typeof payload === "object" ? (payload as TelegramApiPayload) : undefined;
|
|
}
|
|
|
|
function resolveGroupChatKey(payload: TelegramApiPayload): string | undefined {
|
|
const chatId = readNumericId(payload.chat_id);
|
|
return chatId !== undefined && chatId < 0 ? String(chatId) : undefined;
|
|
}
|
|
|
|
function resolveForumLaneKey(payload: TelegramApiPayload): string {
|
|
const threadId = readNumericId(payload.message_thread_id);
|
|
if (threadId !== undefined) {
|
|
return `topic:${threadId}`;
|
|
}
|
|
const directTopicId = readNumericId(payload.direct_messages_topic_id);
|
|
if (directTopicId !== undefined) {
|
|
return `direct-topic:${directTopicId}`;
|
|
}
|
|
const messageId = readNumericId(payload.message_id);
|
|
if (messageId !== undefined) {
|
|
return `message:${messageId}`;
|
|
}
|
|
return "main";
|
|
}
|
|
|
|
export function createTelegramAccountThrottler(
|
|
createThrottler: () => ApiThrottlerTransformer = apiThrottler,
|
|
): ApiThrottlerTransformer {
|
|
const baseThrottler = createThrottler();
|
|
const fairQueuesByChat = new Map<string, GroupFairQueue>();
|
|
|
|
return (prev, method, payload, signal) => {
|
|
const apiPayload = readPayload(payload);
|
|
const groupChatKey = apiPayload ? resolveGroupChatKey(apiPayload) : undefined;
|
|
if (!apiPayload || !groupChatKey) {
|
|
return baseThrottler(prev, method, payload, signal);
|
|
}
|
|
|
|
let fairQueue = fairQueuesByChat.get(groupChatKey);
|
|
if (!fairQueue) {
|
|
fairQueue = new GroupFairQueue();
|
|
fairQueuesByChat.set(groupChatKey, fairQueue);
|
|
}
|
|
|
|
const laneKey = resolveForumLaneKey(apiPayload);
|
|
return fairQueue.enqueue(laneKey, () => baseThrottler(prev, method, payload, signal));
|
|
};
|
|
}
|
|
|
|
export function getOrCreateAccountThrottler(
|
|
token: string,
|
|
createThrottler: () => ApiThrottlerTransformer = apiThrottler,
|
|
): ApiThrottlerTransformer {
|
|
let throttler = throttlerByToken.get(token);
|
|
if (!throttler) {
|
|
throttler = createTelegramAccountThrottler(createThrottler);
|
|
throttlerByToken.set(token, throttler);
|
|
}
|
|
return throttler;
|
|
}
|
|
|
|
export function clearAccountThrottlersForTest(): void {
|
|
throttlerByToken.clear();
|
|
}
|