refactor(telegram): normalize stored offsets through Doctor (#163418)

* fix(telegram): migrate stored update offsets through Doctor

* fix(doctor): own plugin state repairs in migration context

* fix(telegram): align Doctor tests with retired sidecars

* fix(telegram): use canonical identity fixtures after Doctor

* refactor(telegram): retire unsupported group mention alias

Remove groupMentionsOnly's plugin conversion, legacy rule, and dedicated helper. The shared retired-config detector preserves the original config and reports the existing 2026.9.5 upgrade path; that release's Telegram Doctor still converts the preference to groups.*.requireMention.

Source history first introduces this name as a March 2026 compatibility migration, and no supported producer was found. Both v2026.2.22 and v2026.2.23 omit it, so this change makes no last-writer release claim. Generic unknown-key preservation does not establish a supported config contract.

Extend the existing include/root/backup preservation case for an explicit false value. Offset migration and its prepared published-driver cell remain unchanged. Production +1/-50/net -49. Native P2 review and git diff --check passed; remote verification remains pending.
This commit is contained in:
Peter Steinberger 2026-10-03 00:34:30 -07:00 • committed by GitHub
parent 32bf9d4313
commit a45979ee0f
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
24 changed files with 672 additions and 490 deletions

View file

@ -971,7 +971,6 @@ extensions/telegram/src/telegram-ingress-supersede.ts 7
extensions/telegram/src/telegram-ingress-worker.runtime.ts 3
extensions/telegram/src/thread-bindings-store.ts 1
extensions/telegram/src/topic-name-cache.ts 1
extensions/telegram/src/update-offset-store.ts 1
extensions/telegram/src/webhook.ts 1
extensions/tencent/stream.ts 1
extensions/tlon/src/channel.runtime.ts 3

View file

@ -9,17 +9,20 @@ read_when:
`openclaw doctor --fix` owns the persistent file-to-SQLite migrations. This page
describes each migration source and what to do when one stays blocked.
Pre-June Telegram and iMessage caches, Active Memory session toggles, Nostr bus
Pre-June iMessage caches, Active Memory session toggles, Nostr bus
and profile state, and Microsoft Teams conversations, polls, SSO tokens, and
feedback learnings are no longer imported from JSON files. If those sources
remain, Doctor preserves them and directs you to [upgrade through `2026.9.5`](/install/updating#upgrading-very-old-versions)
and run its migrations first. Existing SQLite state remains authoritative.
Doctor archives a retired Telegram `thread-bindings-*.json` file as a completed
no-op only when it is a regular file containing exactly version `1` and an empty
`bindings` array. Nonempty, malformed, symlinked, or otherwise uncertain files
remain preserved for operator review. If archiving a verified empty file fails,
Doctor keeps the original bytes for a later retry and reports a recoverable
warning. This cleanup failure does not block an update.
Pre-July Telegram bot-info, sticker, thread-binding, update-offset, message,
sent-message, and topic-name JSON sidecars are no longer inspected or archived.
Doctor leaves their files untouched, including empty thread-binding files. If
you still need their state, restore a complete pre-update backup and run
`openclaw doctor --fix` on OpenClaw `2026.9.5` before updating again. The separate
Telegram JSON ingress-spool migration still imports pending updates, processing
claims, and failed tombstones with verified backups.
Retired `subagents/runs.json` files are also ignored and left untouched;
transient runs are never restored from them.

View file

@ -36,9 +36,9 @@ published later, including an extended-stable release, still counts. Retain a
transform whenever a release in that window can still write its input format.
A supported release that preserves a legacy
format when rewriting existing data also counts as a writer. A format last
written before the cutoff may be retired only with a clear refusal naming an
intermediate release to upgrade through before retrying. Retirement must never
silently discard persisted data.
written before the cutoff may be retired together with its Doctor checks.
When Doctor refuses a retired input, it names an intermediate release to upgrade
through before retrying. Retirement must leave persisted source data untouched.
Legacy normalization belongs to Doctor and migration owners, with the existing
backup and verification flow. Runtime readers consume canonical state.
@ -80,6 +80,20 @@ memory/auth/cache databases remain supported; see [agent schema
history](/reference/database-schemas/agent-schema-history) for the supported
layouts and recovery route.
Telegram's pre-July bot-info, sticker, thread-binding, update-offset, message,
sent-message, and topic-name JSON sidecars are no longer inspected or archived.
Their last file writers shipped in May 2026. To recover state held only in those
files, use a pre-update backup with OpenClaw `2026.9.5` Doctor before updating.
See [legacy state migration](/cli/doctor/state-migrations).
Telegram SQLite update-offset versions 1 and 2 remain supported because published
July-era Doctor imports can still write them. Doctor normalizes those rows to
version 3 after saving a verified SQLite backup. The cursor, row timestamps,
expiry, and unrelated fields are preserved. Missing bot identity and token
fingerprints remain null; account startup retains responsibility for token
rotation and any required ingress purge. Updates run this repair before account
startup. After a manual package replacement, run `openclaw doctor --fix` first.
Old `openclaw.extension.json` npm declaration stubs are ignored by discovery and
Doctor. They are not plugin manifests, and their files remain unchanged. Reinstall
the package with `openclaw plugins install npm:<package>` and update any explicit
@ -159,6 +173,7 @@ Doctor also refuses these retired config inputs:
`talk.model`, and `talk.voice`.
- `channels.telegram.requireMention`, `channels.feishu.accounts.<id>.botName`,
and the retired `channels.webchat` section.
- `channels.telegram.groupMentionsOnly`; use `channels.telegram.groups["*"].requireMention`.
- `session.threadBindings.ttlHours` and Discord/LINE/Matrix/Telegram `threadBindings.ttlHours`,
including per-account settings.
- Telegram `dm`, `direct.*.threadReplies`, native draft preview settings, and scalar

View file

@ -559,3 +559,20 @@ Telegram's normal inbound agent path after the handler succeeds. OpenClaw keeps
the callback button when inbound policy skips the text or processing fails, so
the user can retry after the blocking condition changes. This result field is
Telegram-specific; other channels keep their own interactive result contracts.
### Doctor plugin-state repairs
`PluginDoctorStateMigrationContext.repairPluginStateEntries(namespace, replacements)`
is available during the offline `after-session-repair` phase. Each replacement
contains an exact `PluginDoctorRawStateEntry` observation from
`readPluginStateEntriesInKeyRange` and a JSON-compatible `value`. An empty read
prefix scans the namespace in pages of at most 512 rows. The host binds plugin
identity and the state location; plugins never supply database paths or SQL.
The host freezes each batch, verifies a backup containing the original row bytes,
and compares the complete observations under current maintenance authority before
one transaction replaces their values. Keys, creation timestamps, and expiry
remain unchanged. Any changed row or database generation refuses the whole batch.
Plugins keep format interpretation in their Doctor contract and leave credential
binding and runtime lifecycle decisions with their existing owners. Older hosts
may omit this optional repair capability.

View file

@ -34,3 +34,12 @@ Verified cleanup-only failures warn without blocking an upgrade.
Plugin developers can follow the
[Doctor ingress migration contract](https://docs.openclaw.ai/plugins/sdk-migration/how-to-migrate#migrate-durable-ingress-files-through-doctor).
## Retired JSON sidecars
Bot-info, sticker, thread-binding, update-offset, message, sent-message, and
topic-name JSON sidecars from before July 2026 are no longer inspected or
archived by Doctor. The files remain untouched. If they contain state you
still need, restore a complete pre-update backup, run `openclaw doctor --fix`
with OpenClaw 2026.9.5, then update again. See
[older-version upgrades](https://docs.openclaw.ai/install/updating#upgrading-very-old-versions).

View file

@ -1,18 +1,89 @@
import type { PluginDoctorStateMigration } from "openclaw/plugin-sdk/runtime-doctor-migrations";
import {
asObjectRecord,
type PluginDoctorStateMigration,
type PluginDoctorStateMigrationContext,
} from "openclaw/plugin-sdk/runtime-doctor-migrations";
export * from "./config-doctor-api.js";
const offsetNamespace = "telegram.update-offsets";
function* legacyOffsets(context: PluginDoctorStateMigrationContext) {
const read = context.readPluginStateEntriesInKeyRange;
if (!read) {
throw new Error(
"Update OpenClaw before inspecting Telegram SQLite offsets, then run openclaw doctor --fix.",
);
}
let after: string | undefined;
while (true) {
const rows = read(offsetNamespace, { prefix: "", after, limit: 512 });
yield rows.flatMap((entry) => {
const value = asObjectRecord(entry.value);
if (!value || (value.version !== 1 && value.version !== 2)) {
return [];
}
const updateId = value.lastUpdateId;
if (
(updateId !== null &&
!(typeof updateId === "number" && Number.isSafeInteger(updateId) && updateId >= 0)) ||
(value.version === 2 && value.botId !== null && typeof value.botId !== "string")
) {
throw new Error(
`Telegram offset for account "${entry.key}" is malformed; restore its known-good state backup, then run openclaw doctor --fix.`,
);
}
return [
{
entry,
value: {
...value,
version: 3,
botId: value.version === 1 ? null : value.botId,
tokenFingerprint: null,
},
},
];
});
const last = rows.at(-1);
if (rows.length < 512 || !last) {
return;
}
after = last.key;
}
}
export const stateMigrations: PluginDoctorStateMigration[] = [
{
id: "telegram-legacy-state",
label: "Retired Telegram JSON state",
async detectLegacyState(params) {
const { telegramRetiredStateMigration } = await import("./src/state-migrations.js");
return telegramRetiredStateMigration.detectLegacyState(params);
id: "telegram-update-offsets",
label: "Telegram SQLite update offsets",
phase: "after-session-repair",
collectBackupResources: () => [],
detectLegacyState({ context }) {
for (const rows of legacyOffsets(context)) {
if (rows.length) {
return { preview: ["Normalize Telegram SQLite update offsets before account startup."] };
}
}
return null;
},
async migrateLegacyState(params) {
const { telegramRetiredStateMigration } = await import("./src/state-migrations.js");
return telegramRetiredStateMigration.migrateLegacyState(params);
async migrateLegacyState({ context }) {
const result: { changes: string[]; warnings: string[] } = { changes: [], warnings: [] };
const batches = [...legacyOffsets(context)];
for (const rows of batches) {
if (!rows.length) {
continue;
}
if (!context.repairPluginStateEntries) {
throw new Error("Update OpenClaw to repair Telegram SQLite offsets.");
}
const repaired = await context.repairPluginStateEntries(offsetNamespace, rows);
result.changes.push(
...repaired.changes,
`Normalized ${rows.length} Telegram SQLite update offsets.`,
);
result.warnings.push(...repaired.warnings);
}
return result;
},
},
{

View file

@ -5,7 +5,7 @@
"description": "OpenClaw Telegram channel plugin.",
"doctorContract": {
"configRepair": true,
"stateMigrations": [{ "id": "telegram-legacy-state" }, { "id": "telegram-json-ingress-spool" }]
"stateMigrations": [{ "id": "telegram-update-offsets", "phase": "after-session-repair" }, { "id": "telegram-json-ingress-spool" }]
},
"activation": {
"onStartup": false

View file

@ -115,31 +115,8 @@ function removeRetiredTelegramGroupHistoryContextConfig(params: {
return { entry: updated, changed: true };
}
function resolveCompatibleDefaultGroupEntry(section: Record<string, unknown>): {
groups: Record<string, unknown>;
entry: Record<string, unknown>;
} | null {
const existingGroups = section.groups;
if (existingGroups !== undefined && !asObjectRecord(existingGroups)) {
return null;
}
const groups = asObjectRecord(existingGroups) ?? {};
const defaultKey = "*";
const existingEntry = groups[defaultKey];
if (existingEntry !== undefined && !asObjectRecord(existingEntry)) {
return null;
}
const entry = asObjectRecord(existingEntry) ?? {};
return { groups, entry };
}
export const legacyConfigRules: ChannelDoctorLegacyConfigRule[] = [
...webhookListenerMigration.legacyConfigRules,
{
path: ["channels", "telegram", "groupMentionsOnly"],
message:
'channels.telegram.groupMentionsOnly was removed; use channels.telegram.groups."*".requireMention instead. Run "openclaw doctor --fix".',
},
{
path: ["channels", "telegram"],
message:
@ -205,32 +182,6 @@ export function normalizeCompatibilityConfig({
updated = retired.entry;
changed = changed || retired.changed;
if (updated.groupMentionsOnly !== undefined) {
const defaultGroupEntry = resolveCompatibleDefaultGroupEntry(updated);
if (!defaultGroupEntry) {
changes.push(
"Skipped channels.telegram.groupMentionsOnly migration because channels.telegram.groups already has an incompatible shape; fix remaining issues manually.",
);
} else {
const { groups, entry } = defaultGroupEntry;
if (entry.requireMention === undefined) {
entry.requireMention = updated.groupMentionsOnly;
groups["*"] = entry;
updated = { ...updated, groups };
changes.push(
'Moved channels.telegram.groupMentionsOnly → channels.telegram.groups."*".requireMention.',
);
} else {
changes.push(
'Removed channels.telegram.groupMentionsOnly (channels.telegram.groups."*" already set).',
);
}
const { groupMentionsOnly: _ignored, ...rest } = updated;
updated = rest;
changed = true;
}
}
const accounts = normalizeChannelAccounts({
entry: updated,
pathPrefix: "channels.telegram",

View file

@ -242,8 +242,13 @@ describe("monitorTelegramProvider", () => {
it.each([
{ name: "same-bot token rotation", version: 3, botId: "111111", tokenFingerprint: "old" },
{ name: "matching legacy identity", version: 2, botId: "111111", tokenFingerprint: null },
{ name: "unknown legacy identity", version: 1, botId: null, tokenFingerprint: null },
{
name: "matching identity without token fingerprint",
version: 3,
botId: "111111",
tokenFingerprint: null,
},
{ name: "unknown identity", version: 3, botId: null, tokenFingerprint: null },
])("keeps queue rows for $name", async (identity) => {
await withStateDirEnv("telegram-same-bot-", async ({ stateDir }) => {
const store = await vi.importActual<typeof OffsetStore>("./update-offset-store.js");

View file

@ -1,125 +0,0 @@
import fs from "node:fs/promises";
import path from "node:path";
import type { PluginDoctorStateMigration } from "openclaw/plugin-sdk/runtime-doctor-migrations";
import { resolveStorePath } from "openclaw/plugin-sdk/session-store-paths";
import { useAutoCleanupTempDirTracker } from "openclaw/plugin-sdk/test-env";
import { afterEach, beforeEach, describe, expect, it } from "vitest";
import { stateMigrations } from "../doctor-contract-api.js";
const migration = stateMigrations[0]!;
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
describe("retired Telegram state", () => {
let stateDir: string;
let input: Parameters<PluginDoctorStateMigration["detectLegacyState"]>[0];
beforeEach(() => {
stateDir = tempDirs.make("openclaw-telegram-retired-");
input = {
config: { agents: { ownership: "explicit", entries: { main: {}, ops: {} } } },
env: { OPENCLAW_STATE_DIR: stateDir },
stateDir,
oauthDir: path.join(stateDir, "credentials"),
context: {
openPluginStateKeyedStore() {
throw new Error("retired state inspection must not open canonical stores");
},
},
};
});
it("completes inspection when retired sources are absent or already archived", async () => {
await fs.mkdir(path.join(stateDir, "telegram"));
await fs.writeFile(path.join(stateDir, "telegram", "update-offset-old.json.migrated"), "old");
await fs.writeFile(path.join(stateDir, "telegram", "unrelated.json"), "keep");
expect(await migration.detectLegacyState(input)).toBeNull();
expect(await migration.migrateLegacyState(input)).toEqual({ changes: [], warnings: [] });
});
it("preserves each retired source and directs pending imports through the bridge release", async () => {
const storePath = resolveStorePath(undefined, { env: input.env, agentId: "ops" });
const paths = [
...["bot-info-old", "sticker-cache", "thread-bindings-old", "update-offset-old"].map((name) =>
path.join(stateDir, "telegram", `${name}.json`),
),
`${storePath}.telegram-messages.json`,
`${storePath}.telegram-sent-messages.json`,
`${storePath}.telegram-topic-names.json`,
path.join(stateDir, "sessions", "sessions.json.telegram-messages.json"),
];
for (const source of paths) {
await fs.mkdir(path.dirname(source), { recursive: true });
await fs.writeFile(source, "unparsed legacy bytes\n");
}
const detection = await migration.detectLegacyState(input);
expect(detection?.preview).toHaveLength(paths.length);
const result = await migration.migrateLegacyState(input);
expect(result.changes).toEqual([]);
expect(result.warnings).toHaveLength(paths.length);
expect(result.warningDisposition).toBeUndefined();
expect(
result.warnings.every((warning) => warning.includes("Run openclaw doctor --fix on 2026.9.5")),
).toBe(true);
for (const source of paths) {
expect(result.warnings).toContainEqual(
expect.stringContaining(`Preserved retired Telegram JSON state at ${source}.`),
);
expect(await fs.readFile(source, "utf8")).toBe("unparsed legacy bytes\n");
}
});
it("archives a verified empty version-1 thread bindings file", async () => {
const sourcePath = path.join(stateDir, "telegram", "thread-bindings-default.json");
const source = '{"version":1,"bindings":[]}\n';
await fs.mkdir(path.dirname(sourcePath), { recursive: true });
await fs.writeFile(sourcePath, source);
expect(await migration.detectLegacyState(input)).toEqual({
preview: [expect.stringContaining(sourcePath)],
});
const result = await migration.migrateLegacyState(input);
expect(result.warnings).toEqual([]);
expect(result.changes).toEqual([
`Archived empty Telegram thread bindings legacy source -> ${sourcePath}.migrated`,
]);
await expect(fs.readFile(`${sourcePath}.migrated`, "utf8")).resolves.toBe(source);
await expect(fs.stat(sourcePath)).rejects.toMatchObject({ code: "ENOENT" });
await expect(migration.detectLegacyState(input)).resolves.toBeNull();
});
it.each([
{
name: "nonempty bindings",
source: '{"version":1,"bindings":[{"chatId":"123"}]}\n',
},
{
name: "an unknown field",
source: '{"version":1,"bindings":[],"metadata":{}}\n',
},
{
name: "a duplicate escaped bindings key",
source: '{"version":1,"bindings":[{"chatId":"123"}],"\\u0062indings":[]}\n',
},
])("preserves a version-1 thread bindings file with $name", async ({ source }) => {
const sourcePath = path.join(stateDir, "telegram", "thread-bindings-default.json");
await fs.mkdir(path.dirname(sourcePath), { recursive: true });
await fs.writeFile(sourcePath, source);
const result = await migration.migrateLegacyState(input);
expect(result.changes).toEqual([]);
expect(result.warnings).toEqual([expect.stringContaining(sourcePath)]);
await expect(fs.readFile(sourcePath, "utf8")).resolves.toBe(source);
await expect(fs.stat(`${sourcePath}.migrated`)).rejects.toMatchObject({ code: "ENOENT" });
});
it("refuses an unreadable source directory instead of certifying inspection", async () => {
await fs.writeFile(path.join(stateDir, "telegram"), "not a directory");
await expect(migration.detectLegacyState(input)).rejects.toMatchObject({ code: "ENOTDIR" });
await expect(migration.migrateLegacyState(input)).rejects.toMatchObject({ code: "ENOTDIR" });
});
});

View file

@ -1,122 +0,0 @@
import fs from "node:fs/promises";
import path from "node:path";
import { listAgentIds } from "openclaw/plugin-sdk/agent-scope-runtime";
import { extractErrorCode } from "openclaw/plugin-sdk/error-runtime";
import {
archiveLegacyStateSource,
type PluginDoctorStateMigration,
} from "openclaw/plugin-sdk/runtime-doctor-migrations";
import { resolveStorePath } from "openclaw/plugin-sdk/session-store-paths";
import { listTelegramAccountIds } from "./account-selection.js";
type MigrationInput = Parameters<PluginDoctorStateMigration["detectLegacyState"]>[0];
// Exact key spellings keep duplicate or escaped property names in the uncertain path.
const EMPTY_THREAD_BINDINGS_PATTERNS = [
/^\s*\{\s*"version"\s*:\s*1\s*,\s*"bindings"\s*:\s*\[\s*\]\s*\}\s*$/,
/^\s*\{\s*"bindings"\s*:\s*\[\s*\]\s*,\s*"version"\s*:\s*1\s*\}\s*$/,
];
function retiredStateWarning(source: string): string {
const state = /^thread-bindings-.+\.json$/.test(path.basename(source))
? "Telegram thread bindings"
: "Telegram state";
return `${state} may contain unmigrated data. Run openclaw doctor --fix on 2026.9.5 with a pre-update backup. Preserved retired Telegram JSON state at ${source}. See https://docs.openclaw.ai/install/updating#upgrading-very-old-versions`;
}
async function isVerifiedEmptyThreadBindingsSource(source: string): Promise<boolean> {
try {
const entry = await fs.lstat(source);
if (!entry.isFile()) {
return false;
}
const raw = await fs.readFile(source, "utf8");
JSON.parse(raw);
return EMPTY_THREAD_BINDINGS_PATTERNS.some((pattern) => pattern.test(raw));
} catch {
return false;
}
}
async function collectRetiredStateSources(params: MigrationInput): Promise<string[]> {
const telegramDir = path.join(params.stateDir, "telegram");
const sources: string[] = [];
try {
for (const entry of await fs.readdir(telegramDir, { withFileTypes: true })) {
if (
(entry.isFile() || entry.isSymbolicLink()) &&
/^(?:bot-info-.+|sticker-cache|thread-bindings-.+|update-offset-.+)\.json$/.test(entry.name)
) {
sources.push(path.join(telegramDir, entry.name));
}
}
} catch (error) {
if (extractErrorCode(error) !== "ENOENT") {
throw error;
}
}
const agentIds = new Set([
...listAgentIds(params.config),
...listTelegramAccountIds(params.config),
"main",
]);
const storePaths = new Set([
path.join(params.stateDir, "sessions", "sessions.json"),
...[...agentIds].map((agentId) =>
resolveStorePath(params.config.session?.store, { env: params.env, agentId }),
),
]);
for (const storePath of storePaths) {
for (const suffix of ["telegram-messages", "telegram-sent-messages", "telegram-topic-names"]) {
const sourcePath = `${storePath}.${suffix}.json`;
try {
const entry = await fs.lstat(sourcePath);
if (entry.isFile() || entry.isSymbolicLink()) {
sources.push(sourcePath);
}
} catch (error) {
if (extractErrorCode(error) !== "ENOENT") {
throw error;
}
}
}
}
return sources;
}
// Keep the action identity so pending imports settle only after their old sources are gone.
export const telegramRetiredStateMigration: PluginDoctorStateMigration = {
id: "telegram-legacy-state",
label: "Retired Telegram JSON state",
async detectLegacyState(params) {
const preview = (await collectRetiredStateSources(params)).map(retiredStateWarning);
return preview.length > 0 ? { preview } : null;
},
async migrateLegacyState(params) {
const changes: string[] = [];
const warnings: string[] = [];
const archiveWarnings: string[] = [];
for (const source of await collectRetiredStateSources(params)) {
if (
/^thread-bindings-.+\.json$/.test(path.basename(source)) &&
(await isVerifiedEmptyThreadBindingsSource(source))
) {
await archiveLegacyStateSource({
filePath: source,
label: "empty Telegram thread bindings",
changes,
warnings: archiveWarnings,
});
} else {
warnings.push(retiredStateWarning(source));
}
}
return {
changes,
warnings: [...warnings, ...archiveWarnings],
...(warnings.length === 0 && archiveWarnings.length > 0
? { warningDisposition: "recoverable" as const }
: {}),
};
},
};

View file

@ -133,11 +133,13 @@ describe("deleteTelegramUpdateOffset", () => {
});
});
it("invokes onRotationDetected for imported legacy offsets without bot identity", async () => {
it("invokes onRotationDetected for canonical offsets without bot identity", async () => {
await withStateDirEnv("openclaw-tg-offset-", async () => {
await updateOffsetStore.register("default", {
version: 1,
version: 3,
lastUpdateId: 777,
botId: null,
tokenFingerprint: null,
});
const rotations: Array<Record<string, unknown>> = [];
@ -223,12 +225,13 @@ describe("deleteTelegramUpdateOffset", () => {
});
});
it("treats imported v2 bot-id-only offsets as stale when token identity cannot be verified", async () => {
it("treats canonical bot-id-only offsets as stale when token identity cannot be verified", async () => {
await withStateDirEnv("openclaw-tg-offset-", async () => {
await updateOffsetStore.register("default", {
version: 2,
version: 3,
lastUpdateId: 999,
botId: "111111",
tokenFingerprint: null,
});
const rotations: Array<Record<string, unknown>> = [];
@ -280,21 +283,43 @@ describe("deleteTelegramUpdateOffset", () => {
it("ignores invalid persisted update IDs from plugin-state", async () => {
await withStateDirEnv("openclaw-tg-offset-", async () => {
await updateOffsetStore.register("default", {
version: 2,
version: 3,
lastUpdateId: -1,
botId: "111111",
tokenFingerprint: null,
});
expect(await readTelegramUpdateOffset({ accountId: "default" })).toBeNull();
await updateOffsetStore.register("default", {
version: 2,
version: 3,
lastUpdateId: "not-a-number",
botId: "111111",
tokenFingerprint: null,
});
expect(await readTelegramUpdateOffset({ accountId: "default" })).toBeNull();
});
});
it.each([1, 2])("requires Doctor before reading or preparing a v%s offset", async (version) => {
await withStateDirEnv("openclaw-tg-offset-", async () => {
const original = {
version,
lastUpdateId: 777,
...(version === 2 ? { botId: "111111" } : {}),
};
await updateOffsetStore.register("default", original);
await expect(readTelegramUpdateOffset({})).rejects.toThrow("openclaw doctor --fix");
await expect(
prepareTelegramAccount({
accountId: "default",
botToken: "222222:fixture",
onRotationDetected() {},
}),
).rejects.toThrow("openclaw doctor --fix");
expect(await updateOffsetStore.lookup("default")).toEqual(original);
});
});
it("rejects writing invalid update IDs", async () => {
await withStateDirEnv("openclaw-tg-offset-", async () => {
await expect(

View file

@ -1,5 +1,6 @@
import { formatErrorMessage } from "openclaw/plugin-sdk/error-runtime";
import type { PluginStateKeyedStore } from "openclaw/plugin-sdk/plugin-state-runtime";
import { isRecord } from "openclaw/plugin-sdk/string-coerce-runtime";
import { getTelegramRuntime } from "./runtime.js";
import { normalizeTelegramStateAccountId } from "./state-account-id.js";
import {
@ -45,39 +46,27 @@ function fingerprintFromToken(token?: string): string | null {
return fingerprintTelegramBotToken(trimmed);
}
function safeParseState(parsed: unknown): TelegramUpdateOffsetState | null {
try {
const state = parsed as {
version?: number;
lastUpdateId?: number | null;
botId?: string | null;
tokenFingerprint?: string | null;
};
if (state?.version !== STORE_VERSION && state?.version !== 2 && state?.version !== 1) {
return null;
}
if (state.lastUpdateId !== null && !isValidUpdateId(state.lastUpdateId)) {
return null;
}
if (state.version >= 2 && state.botId !== null && typeof state.botId !== "string") {
return null;
}
if (
state.version === STORE_VERSION &&
state.tokenFingerprint !== null &&
typeof state.tokenFingerprint !== "string"
) {
return null;
}
return {
version: state.version,
lastUpdateId: state.lastUpdateId ?? null,
botId: state.version >= 2 ? (state.botId ?? null) : null,
tokenFingerprint: state.version === STORE_VERSION ? (state.tokenFingerprint ?? null) : null,
};
} catch {
function safeParseState(state: unknown): TelegramUpdateOffsetState | null {
if (!isRecord(state)) {
return null;
}
if (state.version === 1 || state.version === 2) {
throw new Error("Telegram update offsets require migration; run openclaw doctor --fix.");
}
if (
state.version !== STORE_VERSION ||
(state.lastUpdateId !== null && !isValidUpdateId(state.lastUpdateId)) ||
(state.botId !== null && typeof state.botId !== "string") ||
(state.tokenFingerprint !== null && typeof state.tokenFingerprint !== "string")
) {
return null;
}
return {
version: STORE_VERSION,
lastUpdateId: state.lastUpdateId,
botId: state.botId,
tokenFingerprint: state.tokenFingerprint,
};
}
export type TelegramOffsetRotationReason = "bot-id-changed" | "token-rotated" | "legacy-state";

View file

@ -814,7 +814,6 @@ export const PR_PROTECTED_RUNTIME_TEST_FILES: readonly string[] = [
"extensions/telegram/src/probe.response-body-timeout.integration.test.ts",
"extensions/telegram/src/send.telegram-http.test.ts",
"extensions/telegram/src/send.transport-close-proof.test.ts",
"extensions/telegram/src/state-migrations.test.ts",
"extensions/telegram/src/targets.test.ts",
"extensions/telegram/src/telegram-ingress-callback.integration.test.ts",
"extensions/telegram/src/update-offset-store.test.ts",

View file

@ -30,6 +30,7 @@ describe("doctor config persistence", () => {
name: "Telegram",
channels: {
telegram: {
groupMentionsOnly: false,
dm: {},
direct: { "42": { threadReplies: "always" } },
accounts: {
@ -51,6 +52,7 @@ describe("doctor config persistence", () => {
},
},
fields: [
"channels.telegram.groupMentionsOnly",
"channels.telegram.dm",
"channels.telegram.direct.42.threadReplies",
"channels.telegram.accounts.native.streaming.preview.nativeToolProgress",

View file

@ -46,9 +46,6 @@ function shouldUseCompatPreflight(path: ReadonlyArray<string>, value: unknown):
) {
return true;
}
if (joined === "channels.telegram.groupMentionsOnly") {
return true;
}
if (
last === "allow" &&
typeof value === "boolean" &&

View file

@ -10,6 +10,7 @@ import {
unlinkSync,
} from "node:fs";
import nodePath from "node:path";
import type { DatabaseSync } from "node:sqlite";
import { loadSqliteVecExtension } from "../../packages/memory-host-sdk/src/host/sqlite-vec.js";
import { requireDirectorySync, syncDirectorySync } from "../infra/directory-durability.js";
import { formatErrorMessage } from "../infra/errors.js";
@ -63,6 +64,7 @@ export async function backupDoctorSqliteDatabases(params: {
pendingDatabasePaths: readonly string[];
databasePaths: readonly string[];
authority: DoctorSqliteMaintenanceAuthority;
repair?: { key: string; validate: (database: DatabaseSync) => void };
verifiedSnapshots?: readonly BackupSqliteSnapshotFact[];
}): Promise<MigrationMessages> {
const pending = new Set(params.pendingDatabasePaths);
@ -120,6 +122,7 @@ export async function backupDoctorSqliteDatabases(params: {
const backupDigest = createHash("sha256")
.update(
JSON.stringify([
...(params.repair ? [params.repair.key] : []),
VERSION,
resolveRuntimeServiceBuildId(),
resolveRuntimeServiceCommit(),
@ -225,6 +228,7 @@ export async function backupDoctorSqliteDatabases(params: {
await loadSqliteVecExtension({ db: snapshot });
assertCapture();
assertSqliteIntegrity(snapshot, targetPath);
params.repair?.validate(snapshot);
} finally {
snapshot.close();
}
@ -253,6 +257,7 @@ export async function backupDoctorSqliteDatabases(params: {
preserveRowIds: true,
transform: sanitizeOpenClawStateLeaseRows,
beforePublish: assertCapture,
validate: params.repair?.validate,
});
assertCapture();
changes.push(`Saved pre-migration SQLite backup: ${backup.path}`);

View file

@ -88,7 +88,7 @@ export function findRetiredConfigUpgradeRequirement(
const channels = isRecord(config.channels) ? config.channels : {};
checkKeys(config.gateway, "gateway", ["webchat"]);
checkKeys(channels, "channels", ["webchat"]);
checkKeys(channels.telegram, "channels.telegram", ["requireMention"]);
checkKeys(channels.telegram, "channels.telegram", ["requireMention", "groupMentionsOnly"]);
const beforeDiscord = retired.length;
visitChannelEntries(config, "discord", (scope, configPath) => {
const voice = isRecord(scope.voice) ? scope.voice : {};

View file

@ -101,34 +101,6 @@ describe("legacy state migration caller plugin execution", () => {
});
});
it("completes Doctor after archiving verified empty Telegram thread bindings", async () => {
const fixture = await makeFixture();
const sourcePath = path.join(fixture.stateDir, "telegram", "thread-bindings-default.json");
const source = '{"version":1,"bindings":[]}\n';
fs.mkdirSync(path.dirname(sourcePath), { recursive: true });
fs.writeFileSync(sourcePath, source);
clearPluginDoctorContractRegistryCache();
const result = await autoMigrateLegacyState({
cfg: {},
doctorOnlyStateMigrations: true,
env: fixture.env,
homedir: () => fixture.homeDir,
legacySessionSurfaces: EMPTY_LEGACY_SESSION_SURFACES,
});
expect(
result.stepReceipts.find((receipt) => receipt.id === "plugin-doctor-state"),
).toMatchObject({
outcome: "completed",
changes: [`Archived empty Telegram thread bindings legacy source -> ${sourcePath}.migrated`],
warnings: [],
});
expect(() => throwIfDoctorStateMigrationRefused(result.stepReceipts)).not.toThrow();
expect(fs.existsSync(sourcePath)).toBe(false);
expect(fs.readFileSync(`${sourcePath}.migrated`, "utf8")).toBe(source);
});
it.each([
{
name: "reordered exports with only the second action pending",

View file

@ -1,13 +1,31 @@
import fs from "node:fs";
import path from "node:path";
import type { DatabaseSync } from "node:sqlite";
import { describe, expect, it, vi } from "vitest";
import { createChannelIngressQueue } from "../channels/message/ingress-queue.js";
import { withDoctorSqliteMaintenanceLock } from "../commands/doctor-sqlite-maintenance-lock.js";
import { replaceSessionEntry } from "../config/sessions/session-accessor.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import type { PluginDoctorStateMigration } from "../plugins/doctor-contract-module.js";
import { openOpenClawAgentDatabase } from "../state/openclaw-agent-db.js";
import {
openOpenClawStateDatabase,
closeOpenClawStateDatabaseAsync,
} from "../state/openclaw-state-db.js";
import { loadBundledPluginFacade } from "../test-utils/bundled-plugin-public-surface.js";
import { withOpenClawTestState } from "../test-utils/openclaw-test-state.js";
import {
readDeferredPluginMigrations,
recordDeferredPluginMigrations,
} from "./deferred-plugin-migrations.js";
import { openNodeSqliteDatabase } from "./node-sqlite.js";
import * as sqliteSnapshot from "./sqlite-snapshot.js";
import * as mutationAdmission from "./sqlite-worker-operation-admission.js";
import { createPluginDoctorStateMigrationContext } from "./state-migrations.plugin-doctor-context.js";
import {
runPluginDoctorStateMigrationPlans,
runPostSessionPluginDoctorStateRepairs,
} from "./state-migrations.plugin-doctor.js";
describe("plugin doctor ingress authority", () => {
it("rolls back a queued claim when the repair owner expires before native commit", async () => {
@ -340,3 +358,327 @@ describe("plugin doctor session identity evidence", () => {
);
});
});
describe("Telegram registered SQLite offset repair", () => {
const namespace = "telegram.update-offsets";
const config: OpenClawConfig = {
plugins: { allow: ["telegram"] },
channels: { telegram: { enabled: true } },
};
const pending = [
{
pluginId: "telegram",
requiresStateMigration: true as const,
reason: "offset repair pending",
command: "openclaw doctor --fix",
},
];
const firstOriginal = {
key: "first",
raw: '{ "version": 1, "lastUpdateId": 777, "extra": ["kept"] }',
createdAt: 11,
expiresAt: null,
};
const originals = [
firstOriginal,
{
key: "second",
raw: '{ "version": 2, "lastUpdateId": 999, "botId": "111111" }',
createdAt: 12,
expiresAt: 8_000_000_000_000,
},
];
function seed(db: DatabaseSync) {
for (const row of originals) {
db.prepare(
"INSERT INTO plugin_state_entries (plugin_id, namespace, entry_key, value_json, created_at, expires_at) VALUES (?, ?, ?, ?, ?, ?)",
).run("telegram", namespace, row.key, row.raw, row.createdAt, row.expiresAt);
}
}
function rows(db: DatabaseSync) {
return db
.prepare(
"SELECT entry_key, value_json, created_at, expires_at FROM plugin_state_entries WHERE plugin_id = ? AND namespace = ? ORDER BY entry_key",
)
.all("telegram", namespace);
}
async function registeredMigration() {
const contract = await loadBundledPluginFacade<{
stateMigrations: PluginDoctorStateMigration[];
}>({ pluginId: "telegram", artifactBasename: "doctor-contract-api.js" });
const migration = contract.stateMigrations.find(({ id }) => id === "telegram-update-offsets");
if (!migration) {
throw new Error("Missing registered Telegram offset migration");
}
return migration;
}
it("backs up exact originals and normalizes once without binding credentials or changing row age", async () => {
await withOpenClawTestState(
{ label: "telegram-offset-doctor", applyEnv: false },
async ({ env, stateDir }) => {
const { db } = openOpenClawStateDatabase({ env });
seed(db);
const before = rows(db);
const context = createPluginDoctorStateMigrationContext({
pluginId: "telegram",
env,
config: {},
repairAuthority: {
assertCurrent() {},
assertOwnedInTransaction(database) {
expect(database.isTransaction).toBe(true);
},
},
});
const migration = await registeredMigration();
const input = {
config: {},
env,
stateDir,
oauthDir: path.join(stateDir, "credentials"),
context,
};
expect(await migration.detectLegacyState(input)).not.toBeNull();
await recordDeferredPluginMigrations({ env, pending });
const earlier = await runPluginDoctorStateMigrationPlans({
config,
env,
detected: { stateDir, oauthDir: input.oauthDir, doctorOnlyStateMigrations: true },
});
expect(earlier.completedPluginIds ?? []).not.toContain("telegram");
expect(readDeferredPluginMigrations({ env })).toEqual(pending);
const result = await withDoctorSqliteMaintenanceLock({
env,
operation: "telegram-offset-test",
run: (maintenanceAuthority) =>
runPostSessionPluginDoctorStateRepairs({
env,
config,
maintenanceAuthority,
}),
});
expect(result.warnings).toEqual([]);
expect(readDeferredPluginMigrations({ env })).toEqual([]);
const backupPath = result.changes
.find((line) => line.startsWith("Saved pre-migration SQLite backup: "))
?.split(": ")[1];
if (!backupPath) {
throw new Error("Missing verified pre-repair backup");
}
const backup = openNodeSqliteDatabase(backupPath, { readOnly: true });
try {
expect(rows(backup)).toEqual(before);
} finally {
backup.close();
}
expect(rows(db)).toEqual([
{
entry_key: "first",
value_json: JSON.stringify({
version: 3,
lastUpdateId: 777,
extra: ["kept"],
botId: null,
tokenFingerprint: null,
}),
created_at: 11,
expires_at: null,
},
{
entry_key: "second",
value_json: JSON.stringify({
version: 3,
lastUpdateId: 999,
botId: "111111",
tokenFingerprint: null,
}),
created_at: 12,
expires_at: 8_000_000_000_000,
},
]);
expect(await migration.detectLegacyState(input)).toBeNull();
expect(await migration.migrateLegacyState(input)).toEqual({ changes: [], warnings: [] });
},
);
});
it.each([
{ raw: '{"version":1,"lastUpdateId":"777"}', padding: 511 },
{ raw: '{"version":1,"lastUpdateId":-1}', padding: 0 },
{ raw: '{"version":2,"lastUpdateId":777,"botId":{}}', padding: 0 },
])("leaves malformed durable offsets pending: $raw", async ({ raw, padding }) => {
await withOpenClawTestState(
{ label: "telegram-offset-malformed", applyEnv: false },
async ({ env }) => {
const { db } = openOpenClawStateDatabase({ env });
seed(db);
db.prepare("UPDATE plugin_state_entries SET value_json = ? WHERE entry_key = 'second'").run(
raw,
);
db.exec("BEGIN");
try {
const insert = db.prepare(
"INSERT INTO plugin_state_entries (plugin_id, namespace, entry_key, value_json, created_at, expires_at) VALUES (?, ?, ?, ?, ?, ?)",
);
for (let index = 0; index < padding; index++) {
insert.run("telegram", namespace, `middle-${index}`, firstOriginal.raw, 11, null);
}
db.exec("COMMIT");
} catch (error) {
db.exec("ROLLBACK");
throw error;
}
const before = rows(db);
await recordDeferredPluginMigrations({ env, pending });
const result = await withDoctorSqliteMaintenanceLock({
env,
operation: "telegram-offset-malformed-test",
run: (maintenanceAuthority) =>
runPostSessionPluginDoctorStateRepairs({ env, config, maintenanceAuthority }),
});
expect(result.warnings.join("\n")).toContain(
'account "second" is malformed; restore its known-good state backup',
);
expect(result.changes).toEqual([]);
expect(rows(db)).toEqual(before);
expect(readDeferredPluginMigrations({ env })).toEqual(pending);
},
);
});
it("refuses unsupported inspection instead of certifying empty Telegram state", async () => {
await withOpenClawTestState(
{ label: "telegram-offset-unsupported", applyEnv: false },
async ({ env, stateDir }) => {
const context = createPluginDoctorStateMigrationContext({
pluginId: "telegram",
env,
config,
});
delete context.readPluginStateEntriesInKeyRange;
const migration = await registeredMigration();
const input = {
config,
env,
stateDir,
oauthDir: path.join(stateDir, "credentials"),
context,
};
expect(() => migration.detectLegacyState(input)).toThrow(
"Update OpenClaw before inspecting Telegram SQLite offsets",
);
await expect(migration.migrateLegacyState(input)).rejects.toThrow(
"Update OpenClaw before inspecting Telegram SQLite offsets",
);
},
);
});
it("keeps Telegram pending when its registered after-session repair loses authority", async () => {
await withOpenClawTestState(
{ label: "telegram-offset-pending", applyEnv: false },
async ({ env }) => {
const { db } = openOpenClawStateDatabase({ env });
seed(db);
await recordDeferredPluginMigrations({ env, pending });
let active = true;
const createSnapshot = sqliteSnapshot.createVerifiedSqliteSnapshot;
const interception = vi
.spyOn(sqliteSnapshot, "createVerifiedSqliteSnapshot")
.mockImplementation(async (options) => {
const backup = await createSnapshot(options);
active = false;
return backup;
});
try {
const result = await withDoctorSqliteMaintenanceLock({
env,
operation: "telegram-offset-pending-test",
run: (authority) =>
runPostSessionPluginDoctorStateRepairs({
env,
config,
maintenanceAuthority: {
assertCurrent() {
authority.assertCurrent();
if (!active) {
throw new Error("repair owner expired");
}
},
},
}),
});
expect(result.warnings.join("\n")).toContain("repair owner expired");
expect(readDeferredPluginMigrations({ env })).toEqual(pending);
expect(rows(db)[0]).toMatchObject({ value_json: firstOriginal.raw });
} finally {
interception.mockRestore();
}
},
);
});
it.each(["row", "generation"] as const)(
"refuses a changed %s after backup without partially normalizing",
async (change) => {
await withOpenClawTestState(
{ label: `telegram-offset-${change}`, applyEnv: false },
async ({ env, stateDir }) => {
const database = openOpenClawStateDatabase({ env });
seed(database.db);
const assertCurrent = () => {};
const context = createPluginDoctorStateMigrationContext({
pluginId: "telegram",
env,
config: {},
repairAuthority: { assertCurrent, assertOwnedInTransaction: assertCurrent },
});
const migration = await registeredMigration();
const createSnapshot = sqliteSnapshot.createVerifiedSqliteSnapshot;
const interception = vi
.spyOn(sqliteSnapshot, "createVerifiedSqliteSnapshot")
.mockImplementation(async (options) => {
const backup = await createSnapshot(options);
if (change === "row") {
database.db
.prepare(
"UPDATE plugin_state_entries SET value_json = ? WHERE entry_key = 'second'",
)
.run('{"version":3,"lastUpdateId":123,"botId":null,"tokenFingerprint":null}');
}
if (change === "generation") {
await closeOpenClawStateDatabaseAsync();
fs.renameSync(database.path, `${database.path}.original`);
fs.copyFileSync(backup.path, database.path);
}
return backup;
});
try {
await expect(
migration.migrateLegacyState({
config: {},
env,
stateDir,
oauthDir: path.join(stateDir, "credentials"),
context,
}),
).rejects.toThrow(
change === "row" ? /Plugin state changed/ : /identity changed|source changed/,
);
const current = openOpenClawStateDatabase({ env }).db;
expect(rows(current)[0]).toMatchObject({
entry_key: "first",
value_json: firstOriginal.raw,
created_at: 11,
expires_at: null,
});
} finally {
interception.mockRestore();
}
},
);
},
);
});

View file

@ -1,5 +1,7 @@
import { createHash } from "node:crypto";
import fs from "node:fs";
import path from "node:path";
import type { DatabaseSync } from "node:sqlite";
import { isRecord } from "@openclaw/normalization-core/record-coerce";
import {
inspectAcpSessionClaimsForDoctor,
@ -20,6 +22,7 @@ import {
} from "../config/sessions/targets-read-availability.js";
import { dedupeSessionStoreTargetsBySqliteTarget } from "../config/sessions/targets.js";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { runWriteTransaction } from "../plugin-state/plugin-state-store.database.js";
import {
MAX_PLUGIN_STATE_BULK_DELETE_ENTRIES,
createPluginStateKeyedStore,
@ -29,16 +32,31 @@ import {
pluginStateDoctorEntriesInKeyRange,
type OpenKeyedStoreOptions,
} from "../plugin-state/plugin-state-store.js";
import { getPluginStateKysely } from "../plugin-state/plugin-state-store.kernel.js";
import {
observedPluginStateRow,
type PluginDoctorRawStateEntry,
} from "../plugin-state/plugin-state-store.sqlite.js";
import {
prepareRegisterParams,
validateNamespace,
} from "../plugin-state/plugin-state-store.validation.js";
import type {
PluginDoctorChannelIngressQueueAccess,
PluginDoctorChannelIngressQueueInspection,
PluginDoctorStateMigrationContext,
} from "../plugins/doctor-contract-module.js";
import { normalizeAgentId } from "../routing/session-key.js";
import { resolveOpenClawStateSqlitePath } from "../state/openclaw-state-db.paths.js";
import {
readDeferredPluginSessionImport,
resolveVerifiedSessionSource,
} from "./deferred-plugin-session-sources.js";
import { executeSqliteQuerySync } from "./kysely-sync.js";
import {
assertExistingDatabaseIdentity,
readDatabasePathIdentitySync,
} from "./sqlite-worker-identity.js";
import { readSessionStoreJson5 } from "./state-migrations.fs.js";
import type { PluginDoctorRepairAuthority } from "./state-migrations.types.js";
@ -284,6 +302,90 @@ export type PluginDoctorChannelIngressAccessOptions = {
mutation?: { assertCurrent(): void };
};
/** Backed-up same-schema repair; source observations and replacement values freeze before yielding. */
async function repairPluginStateEntriesForDoctor(params: {
pluginId: string;
namespace: string;
replacements: readonly { entry: PluginDoctorRawStateEntry; value: unknown }[];
authority: PluginDoctorRepairAuthority;
env: NodeJS.ProcessEnv;
}): Promise<{ changes: string[]; warnings: string[] }> {
const { pluginId, authority } = params;
const env = { ...params.env };
const namespace = validateNamespace(params.namespace);
const rows = structuredClone(params.replacements).map(({ entry, value }) => {
const { key, valueJson } = prepareRegisterParams(entry.key, value);
return { entry, key, valueJson };
});
if (
rows.length > MAX_PLUGIN_STATE_BULK_DELETE_ENTRIES ||
new Set(rows.map((row) => row.key)).size !== rows.length ||
rows.some(({ key, entry }) => key !== entry.key)
) {
throw new Error("Plugin Doctor repair requires a bounded batch of distinct rows.");
}
if (!rows.length) {
return { changes: [], warnings: [] };
}
authority.assertCurrent();
const databasePath = resolveOpenClawStateSqlitePath(env);
const identity = readDatabasePathIdentitySync(databasePath);
const assertCurrent = () => {
authority.assertCurrent();
assertExistingDatabaseIdentity(databasePath, identity.key, identity.birthtime);
};
const validate = (db: DatabaseSync) => {
for (const { entry } of rows) {
const query = getPluginStateKysely(db)
.selectFrom("plugin_state_entries")
.select("entry_key")
.where((eb) => observedPluginStateRow(eb, { pluginId, namespace }, entry));
if (!executeSqliteQuerySync(db, query).rows.length) {
throw new Error(
"Plugin state changed during Doctor repair; inspect again before retrying.",
);
}
}
};
const { backupDoctorSqliteDatabases } = await import("../commands/doctor-migration-backup.js");
assertCurrent();
const backup = await backupDoctorSqliteDatabases({
env,
pendingDatabasePaths: [databasePath],
databasePaths: [databasePath],
authority: { assertCurrent },
repair: {
key: createHash("sha256")
.update(JSON.stringify([pluginId, namespace, rows]))
.digest("hex"),
validate,
},
});
assertCurrent();
runWriteTransaction(
"register",
({ db }) => {
assertCurrent();
authority.assertOwnedInTransaction(db);
validate(db);
for (const { key, valueJson } of rows) {
executeSqliteQuerySync(
db,
getPluginStateKysely(db)
.updateTable("plugin_state_entries")
.set({ value_json: valueJson })
.where("plugin_id", "=", pluginId)
.where("namespace", "=", namespace)
.where("entry_key", "=", key),
);
}
authority.assertOwnedInTransaction(db);
},
{ env },
);
return backup;
}
export function createPluginDoctorStateMigrationContext(params: {
pluginId: string;
env: NodeJS.ProcessEnv;
@ -355,6 +457,8 @@ export function createPluginDoctorStateMigrationContext(params: {
}
if (params.repairAuthority) {
const authority = params.repairAuthority;
context.repairPluginStateEntries = (namespace, replacements) =>
repairPluginStateEntriesForDoctor({ pluginId, env, namespace, replacements, authority });
context.updateAcpSessionIdentity = (input) =>
updateAcpSessionIdentityForDoctor(params, authority, input);
context.deletePluginStateEntriesIfUnchanged = (namespace, entries) => {

View file

@ -1,90 +0,0 @@
import fs from "node:fs/promises";
import path from "node:path";
import { afterEach, beforeAll, describe, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import {
coercePluginDoctorContractModule,
type PluginDoctorContractModule,
type PluginDoctorStateMigration,
} from "../plugins/doctor-contract-module.js";
import {
createLegacyStateMigrationStepReceipt,
DoctorStateMigrationRefusalError,
throwIfDoctorStateMigrationRefused,
} from "./state-migrations.messages.js";
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
let migration: PluginDoctorStateMigration;
beforeAll(async () => {
const { stateMigrations } = coercePluginDoctorContractModule(
await vi.importActual<PluginDoctorContractModule>(
path.resolve("extensions/telegram/doctor-contract-api.ts"),
),
);
const selected = stateMigrations?.find((entry) => entry.id === "telegram-legacy-state");
expect(selected).toBeDefined();
migration = selected!;
});
describe("retired Telegram Doctor warning disposition", () => {
it.each([
{ empty: true, nonempty: false },
{ empty: false, nonempty: true },
{ empty: true, nonempty: true },
])("preserves state and diagnostics (empty=$empty, nonempty=$nonempty)", async (fixture) => {
const stateDir = tempDirs.make("openclaw-telegram-doctor-warnings-");
const telegramDir = path.join(stateDir, "telegram");
await fs.mkdir(telegramDir);
const emptyPath = path.join(telegramDir, "thread-bindings-default.json");
const nonemptyPath = path.join(telegramDir, "thread-bindings-private-account.json");
const emptyBytes = '{"version":1,"bindings":[]}';
const nonemptyBytes = '{"version":1,"bindings":[{"chatId":"123"}]}';
if (fixture.empty) {
await fs.writeFile(emptyPath, emptyBytes);
await fs.mkdir(`${emptyPath}.migrated`);
}
if (fixture.nonempty) {
await fs.writeFile(nonemptyPath, nonemptyBytes);
}
const result = await migration.migrateLegacyState({
config: {},
env: { OPENCLAW_STATE_DIR: stateDir },
stateDir,
oauthDir: path.join(stateDir, "credentials"),
context: {
openPluginStateKeyedStore() {
throw new Error("retired source inspection must not open canonical stores");
},
},
});
const receipt = createLegacyStateMigrationStepReceipt(
{
id: "plugin-doctor-state",
phase: "shared",
source: [],
target: [],
requiredness: "required",
reversibility: "checkpoint-required",
},
result,
);
expect(receipt.outcome).toBe(fixture.nonempty ? "refused" : "warning");
if (fixture.empty) {
expect(result.warnings).toContainEqual(expect.stringContaining("Failed archiving"));
expect(await fs.readFile(emptyPath, "utf8")).toBe(emptyBytes);
}
if (fixture.nonempty) {
const failure = new DoctorStateMigrationRefusalError([receipt]);
const message = failure.failureFacts[0]?.message;
expect(message).toContain("Telegram thread bindings may contain unmigrated data");
expect(message).toContain("openclaw doctor --fix on 2026.9.5");
expect(message).toContain("pre-update backup");
expect(message).not.toContain(stateDir);
expect(message).not.toContain("private-account");
expect(await fs.readFile(nonemptyPath, "utf8")).toBe(nonemptyBytes);
} else {
expect(() => throwIfDoctorStateMigrationRefused([receipt])).not.toThrow();
}
});
});

View file

@ -1,10 +1,12 @@
// Plugin state SQLite helpers persist plugin state in the OpenClaw state database.
import type { DatabaseSync } from "node:sqlite";
import { err, ok, type Result } from "@openclaw/normalization-core/result";
import type { ExpressionBuilder } from "kysely";
import { executeSqliteQuerySync } from "../infra/kysely-sync.js";
import { isSqliteCorruptionError } from "../infra/sqlite-error-diagnostics.js";
import { normalizeSqliteNumber } from "../infra/sqlite-number.js";
import { runSqliteImmediateTransactionSync } from "../infra/sqlite-transaction.js";
import type { DB } from "../state/openclaw-state-db.generated.js";
import {
closeOpenClawStateDatabase,
closeOpenClawStateDatabaseAsync,
@ -308,6 +310,22 @@ export function pluginStateDeleteIf(params: {
);
}
export function observedPluginStateRow(
eb: ExpressionBuilder<Pick<DB, "plugin_state_entries">, "plugin_state_entries">,
scope: { pluginId: string; namespace: string },
entry: PluginDoctorRawStateEntry,
) {
return eb
.and({
plugin_id: scope.pluginId,
namespace: scope.namespace,
entry_key: entry.key,
value_json: entry.valueJson,
created_at: entry.createdAt,
})
.and("expires_at", entry.expiresAt === null ? "is" : "=", entry.expiresAt);
}
/** Deletes one bounded set of exact observed rows in a single synchronous transaction. */
export function pluginStateDeleteEntriesIfUnchanged(params: {
pluginId: string;
@ -332,17 +350,9 @@ export function pluginStateDeleteEntriesIfUnchanged(params: {
params.assertOwnedInTransaction(db);
let deleted = 0;
for (const entry of observed) {
let query = getPluginStateKysely(db)
const query = getPluginStateKysely(db)
.deleteFrom("plugin_state_entries")
.where("plugin_id", "=", params.pluginId)
.where("namespace", "=", params.namespace)
.where("entry_key", "=", entry.key)
.where("value_json", "=", entry.valueJson)
.where("created_at", "=", entry.createdAt);
query =
entry.expiresAt === null
? query.where("expires_at", "is", null)
: query.where("expires_at", "=", entry.expiresAt);
.where((eb) => observedPluginStateRow(eb, params, entry));
deleted += Number(executeSqliteQuerySync(db, query).numAffectedRows ?? 0);
}
return { deleted, changed: observed.length - deleted };
@ -361,7 +371,6 @@ export function pluginStateDoctorEntriesInKeyRange(params: {
env?: NodeJS.ProcessEnv;
}): PluginDoctorRawStateEntry[] {
if (
!params.prefix ||
!Number.isSafeInteger(params.limit) ||
params.limit < 1 ||
params.limit > MAX_PLUGIN_STATE_BULK_DELETE_ENTRIES ||

View file

@ -85,6 +85,11 @@ export type PluginDoctorStateMigrationContext = {
namespace: string,
entries: readonly PluginDoctorRawStateEntry[],
) => { deleted: number; changed: number };
/** Offline repair only: verified backup, exact row comparison, and one atomic update. */
repairPluginStateEntries?: (
namespace: string,
replacements: readonly { entry: PluginDoctorRawStateEntry; value: unknown }[],
) => Promise<{ changes: string[]; warnings: string[] }>;
/** Owner-bound ingress queue access, one entry per manifest-declared channel;
* the host fixes the channel identity and doctor state directory. Older test
* hosts may omit it. */