mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-25 20:05:46 -06:00
fix(telegram): deduplicate state migration targets (#129096)
This commit is contained in:
committed by
GitHub
parent
53900c72e4
commit
0903eaddfb
@@ -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)\)/,
|
||||
);
|
||||
|
||||
@@ -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<typeof topicNameCacheImportSource> | 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: {
|
||||
|
||||
Reference in New Issue
Block a user