mirror of
https://github.com/openclaw/openclaw.git
synced 2026-07-22 15:01:13 +00:00
refactor(telegram): trim internal exports (#107726)
* refactor(telegram): trim internal exports * chore(deadcode): refresh export baseline
This commit is contained in:
committed by
GitHub
parent
68f336e496
commit
65e2fddbee
@@ -1,13 +1,9 @@
|
||||
// Telegram tests cover account throttler plugin behavior.
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import {
|
||||
clearAccountThrottlersForTest,
|
||||
createTelegramAccountThrottler,
|
||||
getOrCreateAccountThrottler,
|
||||
} from "./account-throttler.js";
|
||||
import { clearAccountThrottlersForTest, getOrCreateAccountThrottler } from "./account-throttler.js";
|
||||
|
||||
type TelegramPreviousCall = Parameters<ReturnType<typeof createTelegramAccountThrottler>>[0];
|
||||
type TelegramTransform = ReturnType<typeof createTelegramAccountThrottler>;
|
||||
type TelegramTransform = ReturnType<typeof getOrCreateAccountThrottler>;
|
||||
type TelegramPreviousCall = Parameters<TelegramTransform>[0];
|
||||
|
||||
function callLooseSendMessage(
|
||||
throttler: TelegramTransform,
|
||||
@@ -48,7 +44,8 @@ describe("getOrCreateAccountThrottler", () => {
|
||||
it("round-robins group topic requests before entering the Telegram throttler", async () => {
|
||||
const firstGate = deferred<void>();
|
||||
const entered: string[] = [];
|
||||
const throttler = createTelegramAccountThrottler(
|
||||
const throttler = getOrCreateAccountThrottler(
|
||||
"round-robin",
|
||||
() => async (prev, method, payload, signal) => prev(method, payload, signal),
|
||||
);
|
||||
const prev = vi.fn(async (_method: string, payload: unknown) => {
|
||||
@@ -94,7 +91,8 @@ describe("getOrCreateAccountThrottler", () => {
|
||||
it("uses edited message ids as lanes when Telegram omits topic ids", async () => {
|
||||
const firstGate = deferred<void>();
|
||||
const entered: string[] = [];
|
||||
const throttler = createTelegramAccountThrottler(
|
||||
const throttler = getOrCreateAccountThrottler(
|
||||
"edited-message",
|
||||
() => async (prev, method, payload, signal) => prev(method, payload, signal),
|
||||
);
|
||||
const prev = vi.fn(async (_method: string, payload: unknown) => {
|
||||
@@ -138,7 +136,8 @@ describe("getOrCreateAccountThrottler", () => {
|
||||
it("does not group-throttle fractional chat ids", async () => {
|
||||
const firstGate = deferred<void>();
|
||||
const entered: string[] = [];
|
||||
const throttler = createTelegramAccountThrottler(
|
||||
const throttler = getOrCreateAccountThrottler(
|
||||
"direct-topic",
|
||||
() => async (prev, method, payload, signal) => prev(method, payload, signal),
|
||||
);
|
||||
const prev = vi.fn(async (_method: string, payload: unknown) => {
|
||||
@@ -173,7 +172,8 @@ describe("getOrCreateAccountThrottler", () => {
|
||||
it("uses strict decimal string ids for fair group lanes", async () => {
|
||||
const firstGate = deferred<void>();
|
||||
const entered: string[] = [];
|
||||
const throttler = createTelegramAccountThrottler(
|
||||
const throttler = getOrCreateAccountThrottler(
|
||||
"private-chat",
|
||||
() => async (prev, method, payload, signal) => prev(method, payload, signal),
|
||||
);
|
||||
const prev = vi.fn(async (_method: string, payload: unknown) => {
|
||||
|
||||
@@ -126,7 +126,7 @@ function resolveForumLaneKey(payload: TelegramApiPayload): string {
|
||||
return "main";
|
||||
}
|
||||
|
||||
export function createTelegramAccountThrottler(
|
||||
function createTelegramAccountThrottler(
|
||||
createThrottler: () => ApiThrottlerTransformer = apiThrottler,
|
||||
): ApiThrottlerTransformer {
|
||||
const baseThrottler = createThrottler();
|
||||
|
||||
@@ -4,11 +4,12 @@ import {
|
||||
deleteCachedTelegramBotInfo,
|
||||
readCachedTelegramBotInfo,
|
||||
setTelegramBotInfoCacheStoreForTest,
|
||||
TELEGRAM_BOT_INFO_CACHE_MAX_AGE_MS,
|
||||
writeCachedTelegramBotInfo,
|
||||
} from "./bot-info-cache.js";
|
||||
import type { TelegramBotInfo } from "./bot-info.js";
|
||||
|
||||
const BOT_INFO_CACHE_MAX_AGE_MS = 24 * 60 * 60 * 1000;
|
||||
|
||||
const botInfo: TelegramBotInfo = {
|
||||
id: 123456,
|
||||
is_bot: true,
|
||||
@@ -94,7 +95,7 @@ describe("Telegram bot info cache", () => {
|
||||
readCachedTelegramBotInfo({
|
||||
accountId: "ops",
|
||||
botToken: "123456:secret",
|
||||
now: new Date(Date.now() + TELEGRAM_BOT_INFO_CACHE_MAX_AGE_MS + 1),
|
||||
now: new Date(Date.now() + BOT_INFO_CACHE_MAX_AGE_MS + 1),
|
||||
}),
|
||||
).resolves.toBeNull();
|
||||
});
|
||||
|
||||
@@ -11,7 +11,7 @@ import { fingerprintTelegramBotToken } from "./token-fingerprint.js";
|
||||
const LEGACY_STORE_VERSION = 1;
|
||||
export const TELEGRAM_BOT_INFO_CACHE_NAMESPACE = "telegram.bot-info-cache";
|
||||
export const TELEGRAM_BOT_INFO_CACHE_MAX_ENTRIES = 128;
|
||||
export const TELEGRAM_BOT_INFO_CACHE_MAX_AGE_MS = 24 * 60 * 60 * 1000;
|
||||
const TELEGRAM_BOT_INFO_CACHE_MAX_AGE_MS = 24 * 60 * 60 * 1000;
|
||||
|
||||
type TelegramBotInfoCacheState = {
|
||||
tokenFingerprint: string;
|
||||
|
||||
@@ -219,7 +219,6 @@ const loadTelegramNativeCommandRuntime = createLazyRuntimeModule(
|
||||
export const testing = {
|
||||
loadNativeCommandRuntime: loadTelegramNativeCommandRuntime,
|
||||
};
|
||||
export { testing as __testing };
|
||||
|
||||
type TelegramNativeCommandRuntime = Awaited<ReturnType<typeof loadTelegramNativeCommandRuntime>>;
|
||||
|
||||
|
||||
@@ -156,13 +156,6 @@ export async function runWithTelegramSpooledReplayUpdate<T>(
|
||||
}
|
||||
}
|
||||
|
||||
export async function withTelegramSpooledReplayUpdate<T>(
|
||||
update: object,
|
||||
fn: () => Promise<T>,
|
||||
): Promise<T> {
|
||||
return (await runWithTelegramSpooledReplayUpdate(update, fn)).value;
|
||||
}
|
||||
|
||||
export function isTelegramSpooledReplayUpdate(update: unknown): boolean {
|
||||
return (
|
||||
telegramSpooledReplayFrames.getStore() !== undefined ||
|
||||
|
||||
@@ -67,11 +67,8 @@ const {
|
||||
} = harness;
|
||||
const { createTelegramBotCore: createTelegramBotBase, setTelegramBotRuntimeForTest } =
|
||||
await import("./bot-core.js");
|
||||
const {
|
||||
runWithTelegramSpooledReplayUpdate,
|
||||
runWithTelegramUpdateProcessingFrame,
|
||||
withTelegramSpooledReplayUpdate,
|
||||
} = await import("./bot-processing-outcome.js");
|
||||
const { runWithTelegramSpooledReplayUpdate, runWithTelegramUpdateProcessingFrame } =
|
||||
await import("./bot-processing-outcome.js");
|
||||
const { MediaFetchError } = await import("./telegram-media.runtime.js");
|
||||
|
||||
let createTelegramBot: (
|
||||
@@ -85,6 +82,13 @@ const TELEGRAM_TEST_TIMINGS = {
|
||||
textFragmentGapMs: 30,
|
||||
} as const;
|
||||
|
||||
async function withTelegramSpooledReplayUpdate<T>(
|
||||
update: object,
|
||||
fn: () => Promise<T>,
|
||||
): Promise<T> {
|
||||
return (await runWithTelegramSpooledReplayUpdate(update, fn)).value;
|
||||
}
|
||||
|
||||
function setOpenChannelPostConfig() {
|
||||
loadConfig.mockReturnValue({
|
||||
channels: {
|
||||
|
||||
@@ -68,7 +68,6 @@ const {
|
||||
recordTelegramMessageProcessingResult,
|
||||
runWithTelegramSpooledReplayUpdate,
|
||||
TelegramSpooledReplayProcessingError,
|
||||
withTelegramSpooledReplayUpdate,
|
||||
} = await import("./bot-processing-outcome.js");
|
||||
const { TELEGRAM_RICH_TEXT_LIMIT } = await import("./rich-message.js");
|
||||
const { resolveTelegramConversationRoute } = await import("./conversation-route.js");
|
||||
@@ -151,6 +150,7 @@ function installPerKeySequentializer(): void {
|
||||
key,
|
||||
current.catch(() => undefined),
|
||||
);
|
||||
|
||||
try {
|
||||
await current;
|
||||
} finally {
|
||||
@@ -162,6 +162,13 @@ function installPerKeySequentializer(): void {
|
||||
});
|
||||
}
|
||||
|
||||
async function withTelegramSpooledReplayUpdate<T>(
|
||||
update: object,
|
||||
fn: () => Promise<T>,
|
||||
): Promise<T> {
|
||||
return (await runWithTelegramSpooledReplayUpdate(update, fn)).value;
|
||||
}
|
||||
|
||||
function mockTelegramConfigWrites() {
|
||||
return vi.spyOn(configMutation, "mutateConfigFile").mockResolvedValue({} as never);
|
||||
}
|
||||
|
||||
@@ -57,11 +57,8 @@ const {
|
||||
wasSentByBot,
|
||||
} = await import("./bot.create-telegram-bot.test-harness.js");
|
||||
const { recordOutboundMessageForPromptContext } = await import("./outbound-message-context.js");
|
||||
const {
|
||||
runWithTelegramSpooledReplayUpdate,
|
||||
runWithTelegramUpdateProcessingFrame,
|
||||
withTelegramSpooledReplayUpdate,
|
||||
} = await import("./bot-processing-outcome.js");
|
||||
const { runWithTelegramSpooledReplayUpdate, runWithTelegramUpdateProcessingFrame } =
|
||||
await import("./bot-processing-outcome.js");
|
||||
|
||||
let createTelegramBotBase: typeof import("./bot-core.js").createTelegramBotCore;
|
||||
let setTelegramBotRuntimeForTest: typeof import("./bot-core.js").setTelegramBotRuntimeForTest;
|
||||
@@ -81,6 +78,13 @@ const PARTY_EMOJI = "\u{1F389}";
|
||||
const EYES_EMOJI = "\u{1F440}";
|
||||
const HEART_EMOJI = "\u{2764}\u{FE0F}";
|
||||
|
||||
async function withTelegramSpooledReplayUpdate<T>(
|
||||
update: object,
|
||||
fn: () => Promise<T>,
|
||||
): Promise<T> {
|
||||
return (await runWithTelegramSpooledReplayUpdate(update, fn)).value;
|
||||
}
|
||||
|
||||
function createSignal() {
|
||||
let resolve: (() => void) | undefined;
|
||||
const promise = new Promise<void>((res) => {
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
// Telegram tests cover group migration plugin behavior.
|
||||
import { expectDefined } from "@openclaw/normalization-core";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import { migrateTelegramGroupConfig, migrateTelegramGroupsInPlace } from "./group-migration.js";
|
||||
import { migrateTelegramGroupConfig } from "./group-migration.js";
|
||||
|
||||
function createTelegramGlobalGroupConfig(groups: Record<string, Record<string, unknown>>) {
|
||||
return {
|
||||
@@ -112,12 +112,17 @@ describe("migrateTelegramGroupConfig", () => {
|
||||
});
|
||||
|
||||
it("no-ops when old and new group ids are the same", () => {
|
||||
const groups = {
|
||||
const cfg = createTelegramGlobalGroupConfig({
|
||||
"-123": { requireMention: true },
|
||||
};
|
||||
const result = migrateTelegramGroupsInPlace(groups, "-123", "-123");
|
||||
expect(result).toEqual({ migrated: false, skippedExisting: false });
|
||||
expect(groups).toEqual({
|
||||
});
|
||||
const result = migrateTelegramGroupConfig({
|
||||
cfg,
|
||||
accountId: "default",
|
||||
oldChatId: "-123",
|
||||
newChatId: "-123",
|
||||
});
|
||||
expect(result).toEqual({ migrated: false, skippedExisting: false, scopes: [] });
|
||||
expect(cfg.channels.telegram.groups).toEqual({
|
||||
"-123": { requireMention: true },
|
||||
});
|
||||
});
|
||||
|
||||
@@ -37,7 +37,7 @@ function resolveAccountGroups(
|
||||
return { groups: matchKey ? accounts[matchKey]?.groups : undefined };
|
||||
}
|
||||
|
||||
export function migrateTelegramGroupsInPlace(
|
||||
function migrateTelegramGroupsInPlace(
|
||||
groups: TelegramGroups | undefined,
|
||||
oldChatId: string,
|
||||
newChatId: string,
|
||||
|
||||
@@ -9,9 +9,12 @@ import {
|
||||
resetTelegramMessageCacheBucketsForTest,
|
||||
resolveTelegramMessageCachePersistentScopeKey,
|
||||
TELEGRAM_MESSAGE_CACHE_PERSISTENT_MAX_MESSAGES,
|
||||
type TelegramMessageCachePersistentStore,
|
||||
} from "./message-cache.js";
|
||||
|
||||
type TelegramMessageCachePersistentStore = NonNullable<
|
||||
NonNullable<Parameters<typeof createTelegramMessageCache>[0]>["persistentStore"]
|
||||
>;
|
||||
|
||||
type PersistedCacheValue = {
|
||||
version: 1;
|
||||
sourceMessage: Message;
|
||||
|
||||
@@ -111,7 +111,7 @@ export type PersistedTelegramMessageCacheValue = {
|
||||
threadId?: string;
|
||||
};
|
||||
|
||||
export type TelegramMessageCachePersistentStore = {
|
||||
type TelegramMessageCachePersistentStore = {
|
||||
register(key: string, value: PersistedTelegramMessageCacheValue): Promise<void>;
|
||||
entries(): Promise<Array<{ key: string; value: unknown }>>;
|
||||
};
|
||||
|
||||
@@ -7,17 +7,17 @@ import { resetPluginStateStoreForTests } from "openclaw/plugin-sdk/plugin-state-
|
||||
import { afterEach, beforeEach, describe, expect, it } from "vitest";
|
||||
import {
|
||||
buildTelegramMessageDispatchAccountReplayKey,
|
||||
buildTelegramMessageDispatchReplayKey,
|
||||
claimTelegramMessageDispatchReplay,
|
||||
commitTelegramMessageDispatchReplay,
|
||||
createTelegramMessageDispatchReplayGuard,
|
||||
forgetTelegramMessageDispatchReplay,
|
||||
releaseTelegramMessageDispatchReplay,
|
||||
TELEGRAM_MESSAGE_DISPATCH_DEDUPE_NAMESPACE,
|
||||
TelegramMessageDispatchReplayForgetError,
|
||||
type TelegramMessageDispatchReplayGuard,
|
||||
} from "./message-dispatch-dedupe.js";
|
||||
|
||||
type TelegramMessageDispatchReplayGuard = Parameters<
|
||||
typeof claimTelegramMessageDispatchReplay
|
||||
>[0]["guard"];
|
||||
|
||||
const tempDirs: string[] = [];
|
||||
let previousStateDir: string | undefined;
|
||||
|
||||
@@ -36,10 +36,7 @@ function message(params?: { chatId?: number; messageId?: number }): Message {
|
||||
}
|
||||
|
||||
function storedReplayKey(accountId: string, msg: Message): string {
|
||||
const key = buildTelegramMessageDispatchReplayKey(msg);
|
||||
if (!key) {
|
||||
throw new Error("expected replay key");
|
||||
}
|
||||
const key = JSON.stringify(["message", String(msg.chat.id), msg.message_id]);
|
||||
return buildTelegramMessageDispatchAccountReplayKey({ accountId, key });
|
||||
}
|
||||
|
||||
@@ -89,13 +86,6 @@ afterEach(() => {
|
||||
});
|
||||
|
||||
describe("Telegram message dispatch replay guard", () => {
|
||||
it("keys messages by chat id and message id", () => {
|
||||
expect(buildTelegramMessageDispatchReplayKey(message())).toBe(
|
||||
JSON.stringify(["message", "1234", 42]),
|
||||
);
|
||||
expect(buildTelegramMessageDispatchReplayKey(message({ messageId: 0 }))).toBeNull();
|
||||
});
|
||||
|
||||
it("persists committed dispatches across guard recreation", async () => {
|
||||
const writer = createTelegramMessageDispatchReplayGuard();
|
||||
const first = await claimTelegramMessageDispatchReplay({
|
||||
@@ -351,27 +341,4 @@ describe("Telegram message dispatch replay guard", () => {
|
||||
key: first.key,
|
||||
});
|
||||
});
|
||||
|
||||
it("fails rollback when a committed dispatch key cannot be forgotten", async () => {
|
||||
const guard = {
|
||||
claim: async () => ({ kind: "claimed" }),
|
||||
commit: async () => true,
|
||||
forget: async (key: string) => key !== "failed-key",
|
||||
hasRecent: async () => false,
|
||||
warmup: async () => 0,
|
||||
clearMemory: () => {},
|
||||
memorySize: () => 0,
|
||||
release: () => {},
|
||||
} satisfies TelegramMessageDispatchReplayGuard;
|
||||
|
||||
await expect(
|
||||
forgetTelegramMessageDispatchReplay({
|
||||
guard,
|
||||
keys: ["ok-key", "failed-key", "failed-key"],
|
||||
}),
|
||||
).rejects.toMatchObject({
|
||||
name: TelegramMessageDispatchReplayForgetError.name,
|
||||
failures: [{ key: "failed-key" }],
|
||||
});
|
||||
});
|
||||
});
|
||||
|
||||
@@ -12,7 +12,7 @@ export const TELEGRAM_MESSAGE_DISPATCH_DEDUPE_STATE_PLUGIN_ID = "telegram-messag
|
||||
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 &
|
||||
type TelegramMessageDispatchReplayGuard = ClaimableDedupe &
|
||||
Required<Pick<ClaimableDedupe, "forget">>;
|
||||
|
||||
type TelegramMessageDispatchClaim =
|
||||
@@ -66,7 +66,7 @@ export function resolveTelegramMessageDispatchLegacyPath(params: {
|
||||
);
|
||||
}
|
||||
|
||||
export function buildTelegramMessageDispatchReplayKey(msg: Message): string | null {
|
||||
function buildTelegramMessageDispatchReplayKey(msg: Message): string | null {
|
||||
const chatId = msg.chat?.id;
|
||||
const messageId = msg.message_id;
|
||||
if (chatId == null || typeof messageId !== "number" || messageId <= 0) {
|
||||
@@ -223,30 +223,6 @@ export async function commitTelegramMessageDispatchReplay(params: {
|
||||
}
|
||||
}
|
||||
|
||||
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[];
|
||||
|
||||
@@ -1,4 +1,8 @@
|
||||
import type { PluginCommandContext } from "openclaw/plugin-sdk/plugin-entry";
|
||||
import { expectDefined } from "@openclaw/normalization-core";
|
||||
import type {
|
||||
OpenClawPluginCommandDefinition,
|
||||
PluginCommandContext,
|
||||
} from "openclaw/plugin-sdk/plugin-entry";
|
||||
import { createTestPluginApi } from "openclaw/plugin-sdk/plugin-test-api";
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
|
||||
@@ -9,7 +13,18 @@ vi.mock("./url.js", async (importOriginal) => ({
|
||||
resolveTelegramMiniAppUrls,
|
||||
}));
|
||||
|
||||
const { createTelegramMiniAppDashboardCommand } = await import("./command.js");
|
||||
const { registerTelegramMiniAppCommand } = await import("./command.js");
|
||||
|
||||
function registerDashboardCommand(
|
||||
api: Parameters<typeof registerTelegramMiniAppCommand>[0],
|
||||
): OpenClawPluginCommandDefinition {
|
||||
const commands: OpenClawPluginCommandDefinition[] = [];
|
||||
registerTelegramMiniAppCommand({
|
||||
...api,
|
||||
registerCommand: (command) => commands.push(command),
|
||||
});
|
||||
return expectDefined(commands[0], "registered Telegram dashboard command");
|
||||
}
|
||||
|
||||
function commandContext(overrides: Partial<PluginCommandContext>): PluginCommandContext {
|
||||
return {
|
||||
@@ -24,9 +39,9 @@ function commandContext(overrides: Partial<PluginCommandContext>): PluginCommand
|
||||
};
|
||||
}
|
||||
|
||||
describe("createTelegramMiniAppDashboardCommand", () => {
|
||||
describe("registerTelegramMiniAppCommand", () => {
|
||||
it("returns a DM-only message for group invocations", async () => {
|
||||
const command = createTelegramMiniAppDashboardCommand(
|
||||
const command = registerDashboardCommand(
|
||||
createTestPluginApi({
|
||||
config: {
|
||||
channels: {
|
||||
@@ -58,7 +73,7 @@ describe("createTelegramMiniAppDashboardCommand", () => {
|
||||
controlUiUrl: "https://host.tailnet.ts.net/openclaw",
|
||||
gatewayUrl: "wss://host.tailnet.ts.net",
|
||||
});
|
||||
const command = createTelegramMiniAppDashboardCommand(
|
||||
const command = registerDashboardCommand(
|
||||
createTestPluginApi({
|
||||
config: {
|
||||
channels: {
|
||||
|
||||
@@ -13,7 +13,7 @@ export function registerTelegramMiniAppCommand(api: OpenClawPluginApi): void {
|
||||
api.registerCommand(createTelegramMiniAppDashboardCommand(api));
|
||||
}
|
||||
|
||||
export function createTelegramMiniAppDashboardCommand(
|
||||
function createTelegramMiniAppDashboardCommand(
|
||||
api: OpenClawPluginApi,
|
||||
): OpenClawPluginCommandDefinition {
|
||||
return {
|
||||
|
||||
@@ -16,6 +16,7 @@ import {
|
||||
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { clearTelegramRuntime, setTelegramRuntime } from "./runtime.js";
|
||||
import type { TelegramRuntime } from "./runtime.types.js";
|
||||
import type { TelegramSpooledUpdate } from "./telegram-ingress-spool.types.js";
|
||||
import type { TelegramIngressWorkerMessage } from "./telegram-ingress-worker.js";
|
||||
|
||||
const runMock = vi.hoisted(() => vi.fn());
|
||||
@@ -70,7 +71,7 @@ let TelegramPollingSession: typeof import("./polling-session.js").TelegramPollin
|
||||
let pollingSessionTesting: typeof import("./polling-session.js").testing;
|
||||
// Mirrors the claim-owner lease default; update together with telegram-ingress-claim-owner.ts.
|
||||
const telegramSpooledUpdateClaimLeaseMs = 30 * 60 * 1000;
|
||||
let claimTelegramSpooledUpdate: typeof import("./telegram-ingress-spool.js").claimTelegramSpooledUpdate;
|
||||
let claimNextTelegramSpooledUpdate: typeof import("./telegram-ingress-spool.js").claimNextTelegramSpooledUpdate;
|
||||
let isTelegramSpooledUpdateClaimOwnedByOtherLiveProcess: typeof import("./telegram-ingress-claim-owner.js").isTelegramSpooledUpdateClaimOwnedByOtherLiveProcess;
|
||||
let listTelegramSpooledUpdateClaims: typeof import("./telegram-ingress-spool.js").listTelegramSpooledUpdateClaims;
|
||||
let listTelegramSpooledUpdates: typeof import("./telegram-ingress-spool.js").listTelegramSpooledUpdates;
|
||||
@@ -89,6 +90,13 @@ let buildTelegramReplyFenceLaneKey: typeof import("./telegram-reply-fence.js").b
|
||||
let endTelegramReplyFence: typeof import("./telegram-reply-fence.js").endTelegramReplyFence;
|
||||
let resetTelegramReplyFenceForTests: typeof import("./telegram-reply-fence.js").resetTelegramReplyFenceForTests;
|
||||
|
||||
async function claimSpooledUpdate(update: TelegramSpooledUpdate) {
|
||||
return await claimNextTelegramSpooledUpdate({
|
||||
spoolDir: path.dirname(update.path),
|
||||
candidateUpdateIds: [update.updateId],
|
||||
});
|
||||
}
|
||||
|
||||
type TelegramApiMiddleware = (
|
||||
prev: (...args: unknown[]) => Promise<unknown>,
|
||||
method: string,
|
||||
@@ -688,7 +696,7 @@ describe("TelegramPollingSession", () => {
|
||||
({ TelegramPollingSession, testing: pollingSessionTesting } =
|
||||
await import("./polling-session.js"));
|
||||
({
|
||||
claimTelegramSpooledUpdate,
|
||||
claimNextTelegramSpooledUpdate,
|
||||
listTelegramSpooledUpdateClaims,
|
||||
listTelegramSpooledUpdates,
|
||||
recoverStaleTelegramSpooledUpdateClaims,
|
||||
@@ -2775,7 +2783,7 @@ describe("TelegramPollingSession", () => {
|
||||
if (!interrupted) {
|
||||
throw new Error("Expected interrupted update");
|
||||
}
|
||||
await claimTelegramSpooledUpdate(interrupted);
|
||||
await claimSpooledUpdate(interrupted);
|
||||
|
||||
const { runPromise, stopWorker } = startIsolatedIngressSession({
|
||||
abort,
|
||||
@@ -2826,7 +2834,7 @@ describe("TelegramPollingSession", () => {
|
||||
if (!interrupted) {
|
||||
throw new Error("Expected interrupted update");
|
||||
}
|
||||
await claimTelegramSpooledUpdate(interrupted);
|
||||
await claimSpooledUpdate(interrupted);
|
||||
|
||||
await runPromise;
|
||||
expect(events).toEqual(["handled:40", "handled:42"]);
|
||||
@@ -2848,7 +2856,7 @@ describe("TelegramPollingSession", () => {
|
||||
if (!interrupted) {
|
||||
throw new Error("Expected interrupted update");
|
||||
}
|
||||
const claimed = await claimTelegramSpooledUpdate(interrupted);
|
||||
const claimed = await claimSpooledUpdate(interrupted);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected claimed update");
|
||||
}
|
||||
@@ -2895,7 +2903,7 @@ describe("TelegramPollingSession", () => {
|
||||
if (!interrupted) {
|
||||
throw new Error("Expected interrupted update");
|
||||
}
|
||||
const claimed = await claimTelegramSpooledUpdate(interrupted);
|
||||
const claimed = await claimSpooledUpdate(interrupted);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected claimed update");
|
||||
}
|
||||
@@ -2939,7 +2947,7 @@ describe("TelegramPollingSession", () => {
|
||||
if (!interrupted) {
|
||||
throw new Error("Expected interrupted update");
|
||||
}
|
||||
const claimed = await claimTelegramSpooledUpdate(interrupted);
|
||||
const claimed = await claimSpooledUpdate(interrupted);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected claimed update");
|
||||
}
|
||||
@@ -3002,7 +3010,7 @@ describe("TelegramPollingSession", () => {
|
||||
if (!adopted) {
|
||||
throw new Error("Expected adopted update");
|
||||
}
|
||||
const claimed = await claimTelegramSpooledUpdate(adopted);
|
||||
const claimed = await claimSpooledUpdate(adopted);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected claimed update");
|
||||
}
|
||||
@@ -3061,7 +3069,7 @@ describe("TelegramPollingSession", () => {
|
||||
if (!interrupted) {
|
||||
throw new Error("Expected interrupted update");
|
||||
}
|
||||
const claimed = await claimTelegramSpooledUpdate(interrupted);
|
||||
const claimed = await claimSpooledUpdate(interrupted);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected claimed update");
|
||||
}
|
||||
|
||||
@@ -12,8 +12,6 @@ import type { TelegramRuntime } from "./runtime.types.js";
|
||||
import { isTelegramSpooledUpdateClaimOwnedByOtherLiveProcess } from "./telegram-ingress-claim-owner.js";
|
||||
import {
|
||||
claimNextTelegramSpooledUpdate,
|
||||
claimTelegramSpooledUpdate,
|
||||
completeTelegramSpooledUpdate,
|
||||
completeTelegramSpooledUpdateWithRetry,
|
||||
failTelegramSpooledUpdateClaim,
|
||||
listTelegramSpooledUpdateClaims,
|
||||
@@ -21,9 +19,19 @@ import {
|
||||
recoverStaleTelegramSpooledUpdateClaims,
|
||||
refreshTelegramSpooledUpdateClaim,
|
||||
releaseTelegramSpooledUpdateClaim,
|
||||
TELEGRAM_SPOOLED_UPDATE_PROCESSING_STALE_MS,
|
||||
writeTelegramSpooledUpdate,
|
||||
} from "./telegram-ingress-spool.js";
|
||||
import type { TelegramSpooledUpdate } from "./telegram-ingress-spool.types.js";
|
||||
|
||||
// Mirrors the production stale-claim default; callers may override it explicitly.
|
||||
const telegramSpooledUpdateProcessingStaleMs = 6 * 60 * 60 * 1000;
|
||||
|
||||
async function claimSpooledUpdate(update: TelegramSpooledUpdate) {
|
||||
return await claimNextTelegramSpooledUpdate({
|
||||
spoolDir: path.dirname(update.path),
|
||||
candidateUpdateIds: [update.updateId],
|
||||
});
|
||||
}
|
||||
|
||||
function installTelegramIngressQueueRuntime(resolveStateDir: () => string): void {
|
||||
setTelegramRuntime({
|
||||
@@ -78,7 +86,11 @@ describe("Telegram ingress spool", () => {
|
||||
if (!updates[0]) {
|
||||
throw new Error("Expected a spooled update");
|
||||
}
|
||||
await completeTelegramSpooledUpdate(updates[0]);
|
||||
const claimed = await claimSpooledUpdate(updates[0]);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected a claimed update");
|
||||
}
|
||||
await completeTelegramSpooledUpdateWithRetry({ update: claimed });
|
||||
|
||||
expect(
|
||||
(await listTelegramSpooledUpdates({ spoolDir })).map((update) => update.updateId),
|
||||
@@ -106,7 +118,7 @@ describe("Telegram ingress spool", () => {
|
||||
throw new Error("Expected a spooled update");
|
||||
}
|
||||
|
||||
const claimed = await claimTelegramSpooledUpdate(update);
|
||||
const claimed = await claimSpooledUpdate(update);
|
||||
|
||||
expect(claimed?.updateId).toBe(20);
|
||||
expect(claimed?.path.endsWith(".json.processing")).toBe(true);
|
||||
@@ -124,7 +136,7 @@ describe("Telegram ingress spool", () => {
|
||||
if (!claimed) {
|
||||
throw new Error("Expected a claimed update");
|
||||
}
|
||||
await completeTelegramSpooledUpdate(claimed);
|
||||
await completeTelegramSpooledUpdateWithRetry({ update: claimed });
|
||||
expect(await listTelegramSpooledUpdateClaims({ spoolDir })).toEqual([]);
|
||||
|
||||
await writeTelegramSpooledUpdate({
|
||||
@@ -145,7 +157,7 @@ describe("Telegram ingress spool", () => {
|
||||
if (!pending) {
|
||||
throw new Error("Expected a spooled update");
|
||||
}
|
||||
const firstClaim = await claimTelegramSpooledUpdate(pending);
|
||||
const firstClaim = await claimSpooledUpdate(pending);
|
||||
if (!firstClaim) {
|
||||
throw new Error("Expected the first claim");
|
||||
}
|
||||
@@ -154,7 +166,7 @@ describe("Telegram ingress spool", () => {
|
||||
if (!retryPending) {
|
||||
throw new Error("Expected the released update");
|
||||
}
|
||||
const secondClaim = await claimTelegramSpooledUpdate(retryPending);
|
||||
const secondClaim = await claimSpooledUpdate(retryPending);
|
||||
if (!secondClaim) {
|
||||
throw new Error("Expected the replacement claim");
|
||||
}
|
||||
@@ -300,7 +312,7 @@ describe("Telegram ingress spool", () => {
|
||||
if (!update) {
|
||||
throw new Error("Expected a spooled update");
|
||||
}
|
||||
const claimed = await claimTelegramSpooledUpdate(update);
|
||||
const claimed = await claimSpooledUpdate(update);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected a claimed update");
|
||||
}
|
||||
@@ -323,7 +335,7 @@ describe("Telegram ingress spool", () => {
|
||||
if (!update) {
|
||||
throw new Error("Expected a spooled update");
|
||||
}
|
||||
const claimed = await claimTelegramSpooledUpdate(update);
|
||||
const claimed = await claimSpooledUpdate(update);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected a claimed update");
|
||||
}
|
||||
@@ -349,7 +361,7 @@ describe("Telegram ingress spool", () => {
|
||||
if (!update) {
|
||||
throw new Error("Expected a spooled update");
|
||||
}
|
||||
const claimed = await claimTelegramSpooledUpdate(update);
|
||||
const claimed = await claimSpooledUpdate(update);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected a claimed update");
|
||||
}
|
||||
@@ -389,9 +401,13 @@ describe("Telegram ingress spool", () => {
|
||||
if (!update) {
|
||||
throw new Error("Expected a spooled update");
|
||||
}
|
||||
await completeTelegramSpooledUpdate(update);
|
||||
const claimed = await claimSpooledUpdate(update);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected a claimed update");
|
||||
}
|
||||
await completeTelegramSpooledUpdateWithRetry({ update: claimed });
|
||||
|
||||
await expect(claimTelegramSpooledUpdate(update)).resolves.toBeNull();
|
||||
await expect(claimSpooledUpdate(update)).resolves.toBeNull();
|
||||
expect(await listTelegramSpooledUpdates({ spoolDir })).toEqual([]);
|
||||
});
|
||||
});
|
||||
@@ -407,7 +423,7 @@ describe("Telegram ingress spool", () => {
|
||||
if (!stale) {
|
||||
throw new Error("Expected spooled updates");
|
||||
}
|
||||
const claimedStale = await claimTelegramSpooledUpdate(stale);
|
||||
const claimedStale = await claimSpooledUpdate(stale);
|
||||
if (!claimedStale) {
|
||||
throw new Error("Expected claimed updates");
|
||||
}
|
||||
@@ -415,7 +431,7 @@ describe("Telegram ingress spool", () => {
|
||||
|
||||
const recovered = await recoverStaleTelegramSpooledUpdateClaims({
|
||||
spoolDir,
|
||||
now: now + TELEGRAM_SPOOLED_UPDATE_PROCESSING_STALE_MS + 1,
|
||||
now: now + telegramSpooledUpdateProcessingStaleMs + 1,
|
||||
});
|
||||
|
||||
expect(recovered).toBe(1);
|
||||
@@ -435,7 +451,7 @@ describe("Telegram ingress spool", () => {
|
||||
if (!update) {
|
||||
throw new Error("Expected a spooled update");
|
||||
}
|
||||
const claimed = await claimTelegramSpooledUpdate(update);
|
||||
const claimed = await claimSpooledUpdate(update);
|
||||
if (!claimed) {
|
||||
throw new Error("Expected a claimed update");
|
||||
}
|
||||
@@ -469,7 +485,7 @@ describe("Telegram ingress spool", () => {
|
||||
claim: {
|
||||
processId: `${process.pid}:1:other-process`,
|
||||
processPid: process.pid,
|
||||
claimedAt: now - TELEGRAM_SPOOLED_UPDATE_PROCESSING_STALE_MS - 1,
|
||||
claimedAt: now - telegramSpooledUpdateProcessingStaleMs - 1,
|
||||
},
|
||||
}),
|
||||
).toBe(false);
|
||||
|
||||
@@ -31,7 +31,7 @@ export type {
|
||||
|
||||
const SPOOL_VERSION = 1;
|
||||
const TELEGRAM_INGRESS_SPOOL_PREFIX = "ingress-spool-";
|
||||
export const TELEGRAM_SPOOLED_UPDATE_PROCESSING_STALE_MS = 6 * 60 * 60 * 1000;
|
||||
const TELEGRAM_SPOOLED_UPDATE_PROCESSING_STALE_MS = 6 * 60 * 60 * 1000;
|
||||
const TELEGRAM_SPOOLED_UPDATE_FAILED_TTL_MS = 30 * 24 * 60 * 60 * 1000;
|
||||
const TELEGRAM_SPOOLED_UPDATE_FAILED_MAX_ENTRIES = 1000;
|
||||
const TELEGRAM_SPOOLED_UPDATE_COMPLETED_TTL_MS = 30 * 24 * 60 * 60 * 1000;
|
||||
@@ -240,9 +240,7 @@ export async function listTelegramSpooledUpdates(params: {
|
||||
);
|
||||
}
|
||||
|
||||
export async function completeTelegramSpooledUpdate(
|
||||
update: TelegramSpooledUpdate,
|
||||
): Promise<boolean> {
|
||||
async function completeTelegramSpooledUpdate(update: TelegramSpooledUpdate): Promise<boolean> {
|
||||
const queue = createTelegramIngressQueue(path.dirname(update.path));
|
||||
// Successful rows stay as bounded tombstones: Telegram can refetch an update
|
||||
// after dispatch, and callbacks have side effects that plain delete would rerun.
|
||||
@@ -280,16 +278,6 @@ export async function completeTelegramSpooledUpdateWithRetry(params: {
|
||||
}
|
||||
}
|
||||
|
||||
export async function claimTelegramSpooledUpdate(
|
||||
update: TelegramSpooledUpdate,
|
||||
): Promise<ClaimedTelegramSpooledUpdate | null> {
|
||||
const spoolDir = path.dirname(update.path);
|
||||
const claimed = await createTelegramIngressQueue(spoolDir).claim(queueEventId(update.updateId), {
|
||||
ownerId: TELEGRAM_SPOOLED_UPDATE_PROCESS_ID,
|
||||
});
|
||||
return claimed ? parseQueueClaim(spoolDir, claimed) : null;
|
||||
}
|
||||
|
||||
export async function claimNextTelegramSpooledUpdate(params: {
|
||||
spoolDir: string;
|
||||
blockedLaneKeys?: Iterable<string>;
|
||||
|
||||
@@ -1030,5 +1030,4 @@ export const testing = {
|
||||
resolveBindingsPath,
|
||||
resolveStoredBindingKey,
|
||||
};
|
||||
export { testing as __testing };
|
||||
/* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */
|
||||
|
||||
@@ -10,14 +10,18 @@ import { fingerprintTelegramBotToken } from "./token-fingerprint.js";
|
||||
import {
|
||||
TELEGRAM_UPDATE_OFFSET_MAX_ENTRIES,
|
||||
TELEGRAM_UPDATE_OFFSET_NAMESPACE,
|
||||
type TelegramUpdateOffsetState,
|
||||
deleteTelegramUpdateOffset,
|
||||
listTelegramLegacyUpdateOffsetEntries,
|
||||
readTelegramUpdateOffset,
|
||||
setTelegramUpdateOffsetStoreForTest,
|
||||
shouldReplaceTelegramUpdateOffsetEntry,
|
||||
writeTelegramUpdateOffset,
|
||||
} from "./update-offset-store.js";
|
||||
|
||||
type TelegramUpdateOffsetState = Awaited<
|
||||
ReturnType<typeof listTelegramLegacyUpdateOffsetEntries>
|
||||
>[number]["value"];
|
||||
|
||||
describe("deleteTelegramUpdateOffset", () => {
|
||||
let updateOffsetStore: PluginStateKeyedStore<TelegramUpdateOffsetState>;
|
||||
|
||||
|
||||
@@ -10,7 +10,7 @@ const STORE_VERSION = 3;
|
||||
export const TELEGRAM_UPDATE_OFFSET_NAMESPACE = "telegram.update-offsets";
|
||||
export const TELEGRAM_UPDATE_OFFSET_MAX_ENTRIES = 1_000;
|
||||
|
||||
export type TelegramUpdateOffsetState = {
|
||||
type TelegramUpdateOffsetState = {
|
||||
version: number;
|
||||
lastUpdateId: number | null;
|
||||
botId: string | null;
|
||||
|
||||
@@ -137,24 +137,14 @@ export const KNIP_UNUSED_EXPORT_BASELINE = [
|
||||
"extensions/synology-chat/src/client.ts: fetchChatUsers (synologyClient)",
|
||||
"extensions/synology-chat/src/webhook-handler.ts: clearSynologyWebhookRateLimiterStateForTest",
|
||||
"extensions/telegram/src/account-throttler.ts: clearAccountThrottlersForTest",
|
||||
"extensions/telegram/src/account-throttler.ts: createTelegramAccountThrottler",
|
||||
"extensions/telegram/src/bot-info-cache.ts: setTelegramBotInfoCacheStoreForTest",
|
||||
"extensions/telegram/src/bot-info-cache.ts: TELEGRAM_BOT_INFO_CACHE_MAX_AGE_MS",
|
||||
"extensions/telegram/src/bot-message-dispatch.ts: resetTelegramReplyFenceForTests",
|
||||
"extensions/telegram/src/bot-native-commands.ts: __testing",
|
||||
"extensions/telegram/src/bot-native-commands.ts: testing",
|
||||
"extensions/telegram/src/bot-processing-outcome.ts: withTelegramSpooledReplayUpdate",
|
||||
"extensions/telegram/src/bot.ts: setTelegramBotRuntimeForTest",
|
||||
"extensions/telegram/src/channel-actions.ts: telegramMessageActionRuntime",
|
||||
"extensions/telegram/src/error-policy.ts: resetTelegramErrorPolicyStoreForTest",
|
||||
"extensions/telegram/src/group-migration.ts: migrateTelegramGroupsInPlace",
|
||||
"extensions/telegram/src/message-cache.ts: resetTelegramMessageCacheBucketsForTest",
|
||||
"extensions/telegram/src/message-cache.ts: TelegramMessageCachePersistentStore",
|
||||
"extensions/telegram/src/message-dispatch-dedupe.ts: buildTelegramMessageDispatchReplayKey",
|
||||
"extensions/telegram/src/message-dispatch-dedupe.ts: forgetTelegramMessageDispatchReplay",
|
||||
"extensions/telegram/src/message-dispatch-dedupe.ts: TelegramMessageDispatchReplayForgetError",
|
||||
"extensions/telegram/src/message-dispatch-dedupe.ts: TelegramMessageDispatchReplayGuard",
|
||||
"extensions/telegram/src/miniapp/command.ts: createTelegramMiniAppDashboardCommand",
|
||||
"extensions/telegram/src/network-config.ts: resetTelegramNetworkConfigStateForTests",
|
||||
"extensions/telegram/src/polling-lease.ts: resetTelegramPollingLeasesForTests",
|
||||
"extensions/telegram/src/polling-session.ts: testing",
|
||||
@@ -165,15 +155,10 @@ export const KNIP_UNUSED_EXPORT_BASELINE = [
|
||||
"extensions/telegram/src/startup-probe-limiter.ts: resetTelegramStartupProbeLimiterForTests",
|
||||
"extensions/telegram/src/sticker-cache-store.ts: clearTelegramStickerCacheForTest",
|
||||
"extensions/telegram/src/sticker-cache-store.ts: setTelegramStickerCacheStoreForTest",
|
||||
"extensions/telegram/src/telegram-ingress-spool.ts: claimTelegramSpooledUpdate",
|
||||
"extensions/telegram/src/telegram-ingress-spool.ts: completeTelegramSpooledUpdate",
|
||||
"extensions/telegram/src/telegram-ingress-spool.ts: TELEGRAM_SPOOLED_UPDATE_PROCESSING_STALE_MS",
|
||||
"extensions/telegram/src/thread-bindings.ts: __testing",
|
||||
"extensions/telegram/src/thread-bindings.ts: setTelegramThreadBindingStoreForTest",
|
||||
"extensions/telegram/src/topic-name-cache.ts: resetTopicNameCacheForTest",
|
||||
"extensions/telegram/src/topic-name-cache.ts: setTelegramTopicNameStoreFactoryForTest",
|
||||
"extensions/telegram/src/update-offset-store.ts: setTelegramUpdateOffsetStoreForTest",
|
||||
"extensions/telegram/src/update-offset-store.ts: TelegramUpdateOffsetState",
|
||||
"extensions/twitch/src/client-manager-registry.ts: clearRegistryForTest",
|
||||
"extensions/twitch/src/config.ts: ResolvedTwitchAccountContext",
|
||||
"extensions/twitch/src/monitor.ts: testing",
|
||||
|
||||
Reference in New Issue
Block a user