Files
adolf/extensions/telegram/src/message-dispatch-dedupe.ts
alvis bedb527145
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
Vendor OpenClaw source as Adolf fork baseline
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
2026-07-05 09:36:54 +00:00

198 lines
6.4 KiB
TypeScript

// Telegram plugin module implements message dispatch dedupe behavior.
import path from "node:path";
import type { Message } from "grammy/types";
import { createClaimableDedupe, type ClaimableDedupe } from "openclaw/plugin-sdk/persistent-dedupe";
import { normalizeStringEntries, uniqueStrings } from "openclaw/plugin-sdk/string-coerce-runtime";
export const TELEGRAM_MESSAGE_DISPATCH_DEDUPE_TTL_MS = 7 * 24 * 60 * 60 * 1000;
export const TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE = "global";
export const TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE_PREFIX = "telegram.message-dispatch-dedupe";
export const TELEGRAM_MESSAGE_DISPATCH_DEDUPE_STATE_PLUGIN_ID = "telegram-message-dispatch-dedupe";
export const TELEGRAM_MESSAGE_DISPATCH_DEDUPE_MEMORY_MAX_ENTRIES = 50_000;
export const TELEGRAM_MESSAGE_DISPATCH_DEDUPE_STATE_MAX_ENTRIES = 50_000;
export type TelegramMessageDispatchReplayGuard = ClaimableDedupe &
Required<Pick<ClaimableDedupe, "forget">>;
export type TelegramMessageDispatchClaim =
| { kind: "claimed"; key: string }
| { kind: "duplicate" }
| { kind: "invalid" };
export type TelegramMessageDispatchReplayForgetFailure = {
key: string;
error?: unknown;
};
export class TelegramMessageDispatchReplayForgetError extends Error {
readonly failures: TelegramMessageDispatchReplayForgetFailure[];
override readonly cause: unknown;
constructor(failures: readonly TelegramMessageDispatchReplayForgetFailure[]) {
const count = failures.length;
super(`telegram message dispatch dedupe rollback failed for ${count} key(s)`, {
cause: failures.find((failure) => failure.error !== undefined)?.error,
});
this.name = "TelegramMessageDispatchReplayForgetError";
this.failures = [...failures];
this.cause = failures.find((failure) => failure.error !== undefined)?.error;
}
}
export function isTelegramMessageDispatchReplayForgetError(
error: unknown,
): error is TelegramMessageDispatchReplayForgetError {
return error instanceof TelegramMessageDispatchReplayForgetError;
}
function sanitizeFileSegment(value: string): string {
const trimmed = value.trim();
if (!trimmed) {
return "default";
}
return trimmed.replace(/[^a-zA-Z0-9_-]/g, "_");
}
export function resolveTelegramMessageDispatchLegacyPath(params: {
storePath: string;
namespace: string;
}): string {
return path.join(
path.dirname(params.storePath),
`${path.basename(params.storePath)}.telegram-message-dispatch-${sanitizeFileSegment(
params.namespace,
)}.json`,
);
}
export function buildTelegramMessageDispatchReplayKey(msg: Message): string | null {
const chatId = msg.chat?.id;
const messageId = msg.message_id;
if (chatId == null || typeof messageId !== "number" || messageId <= 0) {
return null;
}
return JSON.stringify(["message", String(chatId), messageId]);
}
export function buildTelegramMessageDispatchAccountReplayKey(params: {
accountId: string;
key: string;
}): string {
return JSON.stringify(["account", params.accountId, params.key]);
}
function buildTelegramMessageDispatchStoredReplayKey(params: {
accountId: string;
msg: Message;
}): string | null {
const key = buildTelegramMessageDispatchReplayKey(params.msg);
return key
? buildTelegramMessageDispatchAccountReplayKey({ accountId: params.accountId, key })
: null;
}
export function createTelegramMessageDispatchReplayGuard(
params: {
onDiskError?: (error: unknown) => void;
} = {},
): TelegramMessageDispatchReplayGuard {
return createClaimableDedupe({
ttlMs: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_TTL_MS,
memoryMaxSize: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_MEMORY_MAX_ENTRIES,
pluginId: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_STATE_PLUGIN_ID,
namespacePrefix: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE_PREFIX,
stateMaxEntries: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_STATE_MAX_ENTRIES,
...(params.onDiskError ? { onDiskError: params.onDiskError } : {}),
});
}
export async function claimTelegramMessageDispatchReplay(params: {
guard: TelegramMessageDispatchReplayGuard;
accountId: string;
msg: Message;
}): Promise<TelegramMessageDispatchClaim> {
const key = buildTelegramMessageDispatchStoredReplayKey({
accountId: params.accountId,
msg: params.msg,
});
if (!key) {
return { kind: "invalid" };
}
let releaseRetries = 0;
while (true) {
const claim = await params.guard.claim(key, {
namespace: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE,
});
if (claim.kind === "claimed") {
return { kind: "claimed", key };
}
if (claim.kind === "duplicate") {
return { kind: "duplicate" };
}
try {
await claim.pending;
return { kind: "duplicate" };
} catch {
releaseRetries += 1;
if (releaseRetries > 1) {
return { kind: "duplicate" };
}
}
}
}
function normalizeReplayKeys(keys?: readonly string[]): string[] {
return uniqueStrings(normalizeStringEntries(keys ?? []));
}
export async function commitTelegramMessageDispatchReplay(params: {
guard: TelegramMessageDispatchReplayGuard;
keys?: readonly string[];
}): Promise<void> {
const keys = normalizeReplayKeys(params.keys);
await Promise.all(
keys.map((key) =>
params.guard.commit(key, { namespace: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE }),
),
);
}
export async function forgetTelegramMessageDispatchReplay(params: {
guard: TelegramMessageDispatchReplayGuard;
keys?: readonly string[];
}): Promise<void> {
const keys = normalizeReplayKeys(params.keys);
const failures = (
await Promise.all(
keys.map(async (key): Promise<TelegramMessageDispatchReplayForgetFailure | null> => {
try {
const forgotten = await params.guard.forget(key, {
namespace: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE,
});
return forgotten ? null : { key };
} catch (error) {
return { key, error };
}
}),
)
).filter((failure): failure is TelegramMessageDispatchReplayForgetFailure => Boolean(failure));
if (failures.length > 0) {
throw new TelegramMessageDispatchReplayForgetError(failures);
}
}
export function releaseTelegramMessageDispatchReplay(params: {
guard: TelegramMessageDispatchReplayGuard;
keys?: readonly string[];
error?: unknown;
}): void {
const keys = normalizeReplayKeys(params.keys);
for (const key of keys) {
params.guard.release(key, {
namespace: TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE,
error: params.error,
});
}
}