From 0903eaddfbcec63bf7b3867860b7eff44a7ef3cf Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Tue, 25 Aug 2026 02:07:47 -0700 Subject: [PATCH] fix(telegram): deduplicate state migration targets (#129096) --- .../telegram/src/state-migrations.test.ts | 73 ++++-- extensions/telegram/src/state-migrations.ts | 235 ++++++++---------- 2 files changed, 152 insertions(+), 156 deletions(-) diff --git a/extensions/telegram/src/state-migrations.test.ts b/extensions/telegram/src/state-migrations.test.ts index c02b145a1664..5b57ff6f9b8b 100644 --- a/extensions/telegram/src/state-migrations.test.ts +++ b/extensions/telegram/src/state-migrations.test.ts @@ -166,30 +166,67 @@ describe("telegram state migrations", () => { } }); - it("fails closed when global legacy state spans multiple Telegram route owners", async () => { + it("keeps agent sidecars local while ambiguous global state fails closed", async () => { const dir = await mkdtemp(path.join(os.tmpdir(), "openclaw-telegram-state-migration-")); const env = { ...process.env, OPENCLAW_STATE_DIR: dir }; const legacyStorePath = path.join(dir, "sessions", "sessions.json"); const sentMessagePath = `${legacyStorePath}.telegram-sent-messages.json`; - try { - await mkdir(path.dirname(legacyStorePath), { recursive: true }); - await writeFile(sentMessagePath, JSON.stringify({ 7: { 54: Date.now() } })); - - const cfg = { - agents: { ownership: "explicit", entries: { main: {}, ops: {} } }, - channels: { - telegram: { - accounts: { - primary: { botToken: "123456:primary" }, - alerts: { botToken: "123456:alerts" }, - }, + const stores = ["main", "ops"].map((agentId) => { + const storePath = resolveStorePath(undefined, { env, agentId }); + return { + storePath, + messagePath: resolveTelegramMessageCachePath(storePath), + sentPath: `${storePath}.telegram-sent-messages.json`, + }; + }); + const cfg = { + agents: { ownership: "explicit", entries: { main: {}, ops: {} } }, + channels: { + telegram: { + accounts: { + primary: { botToken: "123456:primary" }, + alerts: { botToken: "123456:alerts" }, }, }, - bindings: [ - { agentId: "main", match: { channel: "telegram", accountId: "primary" } }, - { agentId: "ops", match: { channel: "telegram", accountId: "alerts" } }, - ], - } as OpenClawConfig; + }, + bindings: [ + { agentId: "main", match: { channel: "telegram", accountId: "primary" } }, + { agentId: "ops", match: { channel: "telegram", accountId: "alerts" } }, + ], + } as OpenClawConfig; + try { + for (const [index, store] of stores.entries()) { + await mkdir(path.dirname(store.storePath), { recursive: true }); + await writeFile(store.messagePath, JSON.stringify([persistedCacheEntry(60 + index, "")])); + await writeFile(store.sentPath, JSON.stringify({ 7: { [60 + index]: Date.now() } })); + } + + const plans = (await detectTelegramLegacyStateMigrations({ cfg, env })).filter((plan) => + plan.label.includes("message cache"), + ); + expect(plans).toHaveLength(4); + for (const store of stores) { + expect(plans.find((plan) => plan.sourcePath === store.messagePath)).toMatchObject({ + kind: "plugin-state-import", + scopeKey: resolveTelegramMessageCachePersistentScopeKey(store.messagePath), + }); + const sentPlan = plans.find((plan) => plan.sourcePath === store.sentPath); + if (!sentPlan || sentPlan.kind !== "plugin-state-import") { + throw new Error("expected agent-local sent-message import plan"); + } + expect((await sentPlan.readEntries())[0]?.value).toMatchObject({ + scopeKey: createHash("sha256").update(store.storePath, "utf8").digest("hex").slice(0, 24), + }); + } + + const fixedPlans = await detectTelegramLegacyStateMigrations({ + cfg: { ...cfg, session: { store: stores[0]!.storePath } }, + env, + }); + expect(fixedPlans.filter((plan) => plan.label.includes("message cache"))).toHaveLength(2); + + await mkdir(path.dirname(legacyStorePath), { recursive: true }); + await writeFile(sentMessagePath, JSON.stringify({ 7: { 54: Date.now() } })); await expect(detectTelegramLegacyStateMigrations({ cfg, env })).rejects.toThrow( /^Legacy Telegram state has multiple routed owners \((?:main, ops|ops, main)\)/, ); diff --git a/extensions/telegram/src/state-migrations.ts b/extensions/telegram/src/state-migrations.ts index 12a522db1e9f..7d436402733e 100644 --- a/extensions/telegram/src/state-migrations.ts +++ b/extensions/telegram/src/state-migrations.ts @@ -74,18 +74,45 @@ function resolveAgentSessionStorePath(params: { }); } -function listLegacyAgentSessionStorePaths(params: { +function listLegacyAgentSessionStoreSources(params: { cfg: OpenClawConfig; env: NodeJS.ProcessEnv; stateDir?: string; -}): string[] { - return uniqueStrings([ - ...listAgentIds(params.cfg).map((agentId) => + resolveSourcePath: (storePath: string) => string; + preserveLegacyAccountScope?: boolean; +}): Array<{ sourcePath: string; targetStorePath: string }> { + const asSource = (targetStorePath: string) => ({ + sourcePath: params.resolveSourcePath(targetStorePath), + targetStorePath, + }); + const sources = uniqueStrings( + [...listAgentIds(params.cfg), "main"].map((agentId) => resolveAgentSessionStorePath({ ...params, agentId }), ), - resolveAgentSessionStorePath({ ...params, agentId: "main" }), - resolveLegacySessionStorePath(params), - ]); + ) + .map(asSource) + .filter(({ sourcePath }) => fileExists(sourcePath)); + const legacySourcePath = params.resolveSourcePath(resolveLegacySessionStorePath(params)); + if ( + !fileExists(legacySourcePath) || + sources.some(({ sourcePath }) => sourcePath === legacySourcePath) + ) { + return sources; + } + // Only the global sidecar needs a selected owner; agent-local sidecars retain theirs. + const targetAgentId = + params.preserveLegacyAccountScope && + params.cfg.agents?.entries === undefined && + params.cfg.agents?.list === undefined + ? resolveDefaultTelegramAccountId(params.cfg) + : resolveTelegramLegacyStateOwnerAgentId(params.cfg); + return [ + ...sources, + { + sourcePath: legacySourcePath, + targetStorePath: resolveAgentSessionStorePath({ ...params, agentId: targetAgentId }), + }, + ]; } function resolveTelegramLegacyStateOwnerAgentId(cfg: OpenClawConfig): string { @@ -116,6 +143,28 @@ function resolveMigrationStateDir(params: { env: NodeJS.ProcessEnv; stateDir?: s ); } +type TelegramStateImportPlan = Extract< + ChannelLegacyStateMigrationPlan, + { kind: "plugin-state-import" } +>; + +function telegramStateImport( + params: Omit< + TelegramStateImportPlan, + "kind" | "targetPath" | "pluginId" | "scopeKey" | "cleanupSource" | "preview" + > & { scopeKey?: string }, +): TelegramStateImportPlan { + return { + ...params, + kind: "plugin-state-import", + targetPath: `plugin state:${params.namespace}`, + pluginId: "telegram", + scopeKey: params.scopeKey ?? "", + cleanupSource: "rename", + preview: `- ${params.label}: ${params.sourcePath} → plugin state (${params.namespace})`, + }; +} + function parseLegacyMessageCacheJson(text: string): unknown[] | undefined { try { const value: unknown = JSON.parse(text); @@ -225,34 +274,22 @@ function detectTelegramMessageCacheLegacyStateMigration(params: { env: NodeJS.ProcessEnv; stateDir?: string; }): ChannelLegacyStateMigrationPlan[] { - const persistedPaths = listLegacyAgentSessionStorePaths(params) - .map(resolveTelegramMessageCachePath) - .filter(fileExists); - if (persistedPaths.length === 0) { - return []; - } - const ownerStorePath = resolveAgentSessionStorePath({ + const sources = listLegacyAgentSessionStoreSources({ ...params, - agentId: resolveTelegramLegacyStateOwnerAgentId(params.cfg), + resolveSourcePath: resolveTelegramMessageCachePath, }); - const scopeKey = resolveTelegramMessageCachePersistentScopeKey( - resolveTelegramMessageCachePath(ownerStorePath), - ); - return persistedPaths.map((persistedPath) => { - return { - kind: "plugin-state-import", + return sources.map(({ sourcePath, targetStorePath }) => + telegramStateImport({ label: "Telegram prompt-context message cache", - sourcePath: persistedPath, - targetPath: `plugin state:${TELEGRAM_MESSAGE_CACHE_PERSISTENT_NAMESPACE}`, - pluginId: "telegram", + sourcePath, namespace: TELEGRAM_MESSAGE_CACHE_PERSISTENT_NAMESPACE, maxEntries: TELEGRAM_MESSAGE_CACHE_PERSISTENT_MAX_MESSAGES, - scopeKey, - cleanupSource: "rename", - preview: `- Telegram prompt-context message cache: ${persistedPath} → plugin state (${TELEGRAM_MESSAGE_CACHE_PERSISTENT_NAMESPACE})`, - readEntries: () => listTelegramLegacyMessageCacheEntries(persistedPath), - }; - }); + scopeKey: resolveTelegramMessageCachePersistentScopeKey( + resolveTelegramMessageCachePath(targetStorePath), + ), + readEntries: () => listTelegramLegacyMessageCacheEntries(sourcePath), + }), + ); } function detectTelegramBotInfoCacheLegacyStateMigration(params: { @@ -264,24 +301,18 @@ function detectTelegramBotInfoCacheLegacyStateMigration(params: { if (!fileExists(persistedPath)) { return []; } - return { - kind: "plugin-state-import", + return telegramStateImport({ label: "Telegram startup bot info cache", sourcePath: persistedPath, - targetPath: `plugin state:${TELEGRAM_BOT_INFO_CACHE_NAMESPACE}`, - pluginId: "telegram", namespace: TELEGRAM_BOT_INFO_CACHE_NAMESPACE, maxEntries: TELEGRAM_BOT_INFO_CACHE_MAX_ENTRIES, - scopeKey: "", - cleanupSource: "rename", - preview: `- Telegram startup bot info cache: ${persistedPath} → plugin state (${TELEGRAM_BOT_INFO_CACHE_NAMESPACE})`, readEntries: () => { return listTelegramLegacyBotInfoCacheEntries({ accountId, persistedPath, }); }, - }; + }); }); } @@ -315,17 +346,11 @@ async function detectTelegramUpdateOffsetLegacyStateMigration(params: { } catch { botToken = undefined; } - return { - kind: "plugin-state-import", + return telegramStateImport({ label: "Telegram update offset", sourcePath: persistedPath, - targetPath: `plugin state:${TELEGRAM_UPDATE_OFFSET_NAMESPACE}`, - pluginId: "telegram", namespace: TELEGRAM_UPDATE_OFFSET_NAMESPACE, maxEntries: TELEGRAM_UPDATE_OFFSET_MAX_ENTRIES, - scopeKey: "", - cleanupSource: "rename", - preview: `- Telegram update offset: ${persistedPath} → plugin state (${TELEGRAM_UPDATE_OFFSET_NAMESPACE})`, readEntries: () => listTelegramLegacyUpdateOffsetEntries({ accountId, persistedPath }), shouldReplaceExistingEntry: ({ existingValue, incomingValue }) => shouldReplaceTelegramUpdateOffsetEntry({ @@ -333,7 +358,7 @@ async function detectTelegramUpdateOffsetLegacyStateMigration(params: { incomingValue, botToken, }), - }; + }); }); } @@ -347,19 +372,13 @@ function detectTelegramStickerCacheLegacyStateMigration(params: { return []; } return [ - { - kind: "plugin-state-import", + telegramStateImport({ label: "Telegram sticker cache", sourcePath: persistedPath, - targetPath: `plugin state:${TELEGRAM_STICKER_CACHE_NAMESPACE}`, - pluginId: "telegram", namespace: TELEGRAM_STICKER_CACHE_NAMESPACE, maxEntries: TELEGRAM_STICKER_CACHE_MAX_ENTRIES, - scopeKey: "", - cleanupSource: "rename", - preview: `- Telegram sticker cache: ${persistedPath} → plugin state (${TELEGRAM_STICKER_CACHE_NAMESPACE})`, readEntries: () => listTelegramLegacyStickerCacheEntries({ persistedPath }), - }, + }), ]; } @@ -368,39 +387,24 @@ function detectTelegramSentMessageCacheLegacyStateMigration(params: { env: NodeJS.ProcessEnv; stateDir?: string; }): ChannelLegacyStateMigrationPlan[] { - const sourcePaths = listLegacyAgentSessionStorePaths(params) - .map((storePath) => `${storePath}.telegram-sent-messages.json`) - .filter(fileExists); - if (sourcePaths.length === 0) { - return []; - } - const ownerAgentId = resolveTelegramLegacyStateOwnerAgentId(params.cfg); - const targetStorePath = resolveAgentSessionStorePath({ + const sources = listLegacyAgentSessionStoreSources({ ...params, - agentId: ownerAgentId, + resolveSourcePath: (storePath) => `${storePath}.telegram-sent-messages.json`, }); - return sourcePaths.map((sourcePath) => { - return { - kind: "plugin-state-import", + return sources.map(({ sourcePath, targetStorePath }) => + telegramStateImport({ label: "Telegram sent-message cache", sourcePath, - targetPath: `plugin state:${TELEGRAM_SENT_MESSAGE_CACHE_NAMESPACE}`, - pluginId: "telegram", namespace: TELEGRAM_SENT_MESSAGE_CACHE_NAMESPACE, maxEntries: TELEGRAM_SENT_MESSAGE_CACHE_MAX_ENTRIES, - scopeKey: "", - cleanupSource: "rename", cleanupWhenEmpty: true, - preview: `- Telegram sent-message cache: ${sourcePath} → plugin state (${TELEGRAM_SENT_MESSAGE_CACHE_NAMESPACE})`, readEntries: () => listTelegramLegacySentMessageCacheEntries({ - cfg: params.cfg, - agentId: ownerAgentId, persistedPath: sourcePath, targetStorePath, }), - }; - }); + }), + ); } function detectTelegramThreadBindingLegacyStateMigration(params: { @@ -419,31 +423,20 @@ function detectTelegramThreadBindingLegacyStateMigration(params: { if (!fileExists(persistedPath)) { return []; } - return { - kind: "plugin-state-import", + return telegramStateImport({ label: "Telegram thread bindings", sourcePath: persistedPath, - targetPath: `plugin state:${TELEGRAM_THREAD_BINDINGS_NAMESPACE}`, - pluginId: "telegram", namespace: TELEGRAM_THREAD_BINDINGS_NAMESPACE, maxEntries: TELEGRAM_THREAD_BINDINGS_MAX_ENTRIES, - scopeKey: "", - cleanupSource: "rename", - preview: `- Telegram thread bindings: ${persistedPath} → plugin state (${TELEGRAM_THREAD_BINDINGS_NAMESPACE})`, readEntries: () => listTelegramLegacyThreadBindingEntries({ accountId, persistedPath }), - }; + }); }); } -function topicNameCacheImportSource(params: { - sourceStorePath: string; - targetStorePath?: string; -}): { sourcePath: string; namespace: string } { - const targetStorePath = params.targetStorePath ?? params.sourceStorePath; - const scope = resolveTopicNameCacheScope(targetStorePath); +function topicNameCacheImportSource(sourcePath: string, targetStorePath: string) { return { - sourcePath: resolveTopicNameCachePath(params.sourceStorePath), - namespace: resolveTopicNameCacheNamespace(scope), + sourcePath, + namespace: resolveTopicNameCacheNamespace(resolveTopicNameCacheScope(targetStorePath)), }; } @@ -457,69 +450,35 @@ function detectTelegramTopicNameCacheLegacyStateMigration(params: { env: params.env, agentId: accountId, }); - return topicNameCacheImportSource({ sourceStorePath: storePath }); + return topicNameCacheImportSource(resolveTopicNameCachePath(storePath), storePath); }); - const agentSources = listAgentIds(params.cfg).map((agentId) => - topicNameCacheImportSource({ - sourceStorePath: resolveAgentSessionStorePath({ ...params, agentId }), - }), + const sessionSources = listLegacyAgentSessionStoreSources({ + ...params, + resolveSourcePath: resolveTopicNameCachePath, + // Pre-roster topic caches used the account id as their session-store scope. + preserveLegacyAccountScope: true, + }).map(({ sourcePath, targetStorePath }) => + topicNameCacheImportSource(sourcePath, targetStorePath), ); - const legacyMainStorePath = resolveAgentSessionStorePath({ ...params, agentId: "main" }); - const legacyStorePath = resolveLegacySessionStorePath(params); - const legacySourcePath = resolveTopicNameCachePath(legacyStorePath); - const fixedSources = [ - ...accountSources, - ...agentSources, - topicNameCacheImportSource({ sourceStorePath: legacyMainStorePath }), - ].filter((source) => fileExists(source.sourcePath)); - if (fixedSources.length === 0 && !fileExists(legacySourcePath)) { - return []; - } - let legacySource: ReturnType | undefined; - if (fileExists(legacySourcePath)) { - const ownerStorePath = resolveAgentSessionStorePath({ - ...params, - agentId: resolveTelegramLegacyStateOwnerAgentId(params.cfg), - }); - // Pre-roster Telegram scoped this legacy cache by account id. Once an agent roster exists, - // routing owns the migration target just as it owns new Telegram conversations. - const legacyTargetStorePath = - params.cfg.agents?.entries !== undefined || params.cfg.agents?.list !== undefined - ? ownerStorePath - : resolveStorePath(params.cfg.session?.store, { - env: params.env, - agentId: resolveDefaultTelegramAccountId(params.cfg), - }); - legacySource = topicNameCacheImportSource({ - sourceStorePath: legacyStorePath, - targetStorePath: legacyTargetStorePath, - }); - } const sourcesByKey = new Map( - [...fixedSources, ...(legacySource ? [legacySource] : [])].map( + [...accountSources.filter((source) => fileExists(source.sourcePath)), ...sessionSources].map( (source) => [`${source.sourcePath}\0${source.namespace}`, source] as const, ), ); - return [...sourcesByKey.values()].map((source) => { - return { - kind: "plugin-state-import", + return [...sourcesByKey.values()].map((source) => + telegramStateImport({ label: "Telegram forum topic-name cache", sourcePath: source.sourcePath, - targetPath: `plugin state:${source.namespace}`, - pluginId: "telegram", namespace: source.namespace, maxEntries: TELEGRAM_TOPIC_NAME_CACHE_MAX_ENTRIES, - scopeKey: "", - cleanupSource: "rename", - preview: `- Telegram forum topic-name cache: ${source.sourcePath} → plugin state (${source.namespace})`, readEntries: () => { return listTelegramLegacyTopicNameCacheEntries({ persistedPath: source.sourcePath, maxEntries: TELEGRAM_TOPIC_NAME_CACHE_MAX_ENTRIES, }); }, - }; - }); + }), + ); } export async function detectTelegramLegacyStateMigrations(params: {