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
198 lines
6.4 KiB
TypeScript
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,
|
|
});
|
|
}
|
|
}
|