From f3adeb2ac677693f1fd3e5a2ff2f3c075f3c9cb5 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 15 Jul 2026 10:03:09 -0700 Subject: [PATCH] fix(discord): move command deploy cache to sqlite (#108381) --- docs/refactor/database-first.md | 4 + extensions/discord/src/channel.test.ts | 41 +- extensions/discord/src/channel.ts | 12 + .../discord/src/command-deploy-store.ts | 25 ++ .../discord/src/internal/client.test.ts | 18 +- extensions/discord/src/internal/client.ts | 5 +- .../src/internal/command-deploy.test.ts | 402 ++++++------------ .../discord/src/internal/command-deploy.ts | 129 +----- ...odel-picker-preferences-migrations.test.ts | 31 ++ .../model-picker-preferences-migrations.ts | 21 + .../src/monitor/provider.startup.test.ts | 8 + .../discord/src/monitor/provider.startup.ts | 10 +- extensions/discord/src/monitor/provider.ts | 3 + .../plugins/legacy-state-migration.types.ts | 2 +- src/commands/doctor-state-migrations.test.ts | 33 ++ src/infra/state-migrations.plugin-state.ts | 8 + 16 files changed, 346 insertions(+), 406 deletions(-) create mode 100644 extensions/discord/src/command-deploy-store.ts diff --git a/docs/refactor/database-first.md b/docs/refactor/database-first.md index d97ac77abcaf..5768d1138139 100644 --- a/docs/refactor/database-first.md +++ b/docs/refactor/database-first.md @@ -416,6 +416,10 @@ The branch already has a real shared SQLite base: excludes heartbeat sessions. `openclaw doctor --fix` strictly validates the old TUI JSON file, keeps newer SQLite rows, verifies the canonical result, and removes the unchanged legacy file instead of leaving an archive. +- Discord command deployment hashes now live in the shared plugin-state SQLite + store. Runtime reads and writes exact application-scoped keys only. Doctor + deletes the rebuildable legacy `discord/command-deploy-cache.json` file + without importing it, so the next startup performs one canonical reconcile. - Default TTS prefs now live in shared plugin-state SQLite rows keyed under the `speech-core` plugin. The old `settings/tts.json` file is doctor migration input only; runtime no longer reads or writes TTS prefs JSON files, and the diff --git a/extensions/discord/src/channel.test.ts b/extensions/discord/src/channel.test.ts index 2b81b92d30b4..f7dfab7382c5 100644 --- a/extensions/discord/src/channel.test.ts +++ b/extensions/discord/src/channel.test.ts @@ -115,7 +115,13 @@ async function expectDiscordStartupDelay( expect(sleepWithAbortMock).toHaveBeenCalledWith(expectedMs, ctx.abortSignal); } -function installDiscordRuntime(discord: Record) { +function installDiscordRuntime( + discord: Record, + openKeyedStore: (options: Record) => unknown = vi.fn(() => ({ + lookup: vi.fn(async () => undefined), + register: vi.fn(async () => undefined), + })), +) { setDiscordRuntime({ channel: { discord, @@ -123,6 +129,7 @@ function installDiscordRuntime(discord: Record) { logging: { shouldLogVerbose: () => false, }, + state: { openKeyedStore }, } as unknown as PluginRuntime); } @@ -757,6 +764,38 @@ describe("discordPlugin outbound", () => { ); }); + it("opens the SQLite command deployment cache and passes it to the provider", async () => { + prepareDiscordStartupMocks(); + const commandDeployHashStore = { + lookup: vi.fn(async () => undefined), + register: vi.fn(async () => undefined), + }; + const openKeyedStore = vi.fn(() => commandDeployHashStore); + installDiscordRuntime({}, openKeyedStore); + + await startDiscordAccount(createCfg()); + + expect(openKeyedStore).toHaveBeenCalledWith({ + namespace: "command-deploy-hashes", + maxEntries: 10_000, + overflowPolicy: "evict-oldest", + }); + expect(objectArgAt(monitorDiscordProviderMock, 0, 0).commandDeployHashStore).toBe( + commandDeployHashStore, + ); + }); + + it("continues Discord startup when the command deployment cache cannot open", async () => { + prepareDiscordStartupMocks(); + installDiscordRuntime({}, () => { + throw new Error("SQLite unavailable"); + }); + + await startDiscordAccount(createCfg()); + + expect(objectArgAt(monitorDiscordProviderMock, 0, 0).commandDeployHashStore).toBeUndefined(); + }); + it("clears stale Discord probe metadata when the async startup probe degrades", async () => { probeDiscordMock.mockResolvedValue({ ok: false, diff --git a/extensions/discord/src/channel.ts b/extensions/discord/src/channel.ts index 5094845f102a..7f8e76da0204 100644 --- a/extensions/discord/src/channel.ts +++ b/extensions/discord/src/channel.ts @@ -62,6 +62,7 @@ import { loadDiscordTargetResolverModule, loadDiscordThreadBindingsManagerModule, } from "./channel.loaders.js"; +import { openDiscordCommandDeployHashStore } from "./command-deploy-store.js"; import { shouldSuppressLocalDiscordExecApprovalPrompt } from "./exec-approvals.js"; import { resolveDiscordGroupRequireMention, @@ -724,6 +725,16 @@ export const discordPlugin: ChannelPlugin log: ctx.log, }); ctx.log?.info(`[${account.accountId}] starting provider`); + let commandDeployHashStore; + try { + commandDeployHashStore = openDiscordCommandDeployHashStore( + getDiscordRuntime().state.openKeyedStore, + ); + } catch (error) { + ctx.log?.warn?.( + `[${account.accountId}] Discord command deploy cache unavailable; continuing without persistence: ${formatErrorMessage(error)}`, + ); + } return (await loadDiscordProviderRuntime()).monitorDiscordProvider({ token, accountId: account.accountId, @@ -734,6 +745,7 @@ export const discordPlugin: ChannelPlugin mediaMaxMb: account.config.mediaMaxMb, historyLimit: account.config.historyLimit, setStatus: (patch) => ctx.setStatus({ accountId: account.accountId, ...patch }), + commandDeployHashStore, }); }, }, diff --git a/extensions/discord/src/command-deploy-store.ts b/extensions/discord/src/command-deploy-store.ts new file mode 100644 index 000000000000..5cf85de689fc --- /dev/null +++ b/extensions/discord/src/command-deploy-store.ts @@ -0,0 +1,25 @@ +import type { PluginStateKeyedStore } from "openclaw/plugin-sdk/plugin-state-runtime"; + +export const DISCORD_COMMAND_DEPLOY_HASH_NAMESPACE = "command-deploy-hashes"; +export const DISCORD_COMMAND_DEPLOY_HASH_MAX_ENTRIES = 10_000; + +export type DiscordCommandDeployHashStore = Pick< + PluginStateKeyedStore, + "lookup" | "register" +>; + +type OpenKeyedStore = (options: { + namespace: string; + maxEntries: number; + overflowPolicy: "evict-oldest"; +}) => PluginStateKeyedStore; + +export function openDiscordCommandDeployHashStore( + openKeyedStore: OpenKeyedStore, +): DiscordCommandDeployHashStore { + return openKeyedStore({ + namespace: DISCORD_COMMAND_DEPLOY_HASH_NAMESPACE, + maxEntries: DISCORD_COMMAND_DEPLOY_HASH_MAX_ENTRIES, + overflowPolicy: "evict-oldest", + }); +} diff --git a/extensions/discord/src/internal/client.test.ts b/extensions/discord/src/internal/client.test.ts index 875ba4364f12..f26ea3b5a1d7 100644 --- a/extensions/discord/src/internal/client.test.ts +++ b/extensions/discord/src/internal/client.test.ts @@ -1,7 +1,4 @@ // Discord tests cover client plugin behavior. -import fs from "node:fs/promises"; -import os from "node:os"; -import path from "node:path"; import { ApplicationCommandType, ComponentType, Routes } from "discord-api-types/v10"; import { MAX_TIMER_TIMEOUT_MS } from "openclaw/plugin-sdk/number-runtime"; import { afterEach, describe, expect, it, vi } from "vitest"; @@ -392,12 +389,15 @@ describe("Client.deployCommands", () => { }); it("skips unchanged command deploys across client restarts using the hash store", async () => { - const hashStorePath = path.join( - await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-discord-command-deploy-")), - "hashes.json", - ); + const hashes = new Map(); + const commandDeployHashStore = { + lookup: async (key: string) => hashes.get(key), + register: async (key: string, value: string) => { + hashes.set(key, value); + }, + }; const first = createInternalTestClient([createTestCommand({ name: "one" })], { - commandDeployHashStorePath: hashStorePath, + commandDeployHashStore, }); const firstGet = vi.fn(async () => []); const firstPost = vi.fn(async () => undefined); @@ -406,7 +406,7 @@ describe("Client.deployCommands", () => { await first.deployCommands({ mode: "reconcile" }); const second = createInternalTestClient([createTestCommand({ name: "one" })], { - commandDeployHashStorePath: hashStorePath, + commandDeployHashStore, }); const secondGet = vi.fn(async () => []); const secondPost = vi.fn(async () => undefined); diff --git a/extensions/discord/src/internal/client.ts b/extensions/discord/src/internal/client.ts index 358ee2c7bbb2..ab26b7b819ea 100644 --- a/extensions/discord/src/internal/client.ts +++ b/extensions/discord/src/internal/client.ts @@ -1,5 +1,6 @@ // Discord plugin module implements client behavior. import type { APIInteraction } from "discord-api-types/v10"; +import type { DiscordCommandDeployHashStore } from "../command-deploy-store.js"; import { DiscordCommandDeployer, type DeployCommandOptions } from "./command-deploy.js"; import type { BaseCommand } from "./commands.js"; import { ComponentRegistry } from "./component-registry.js"; @@ -46,7 +47,7 @@ interface ClientOptions { disableDeployRoute?: boolean; disableInteractionsRoute?: boolean; disableEventsRoute?: boolean; - commandDeployHashStorePath?: string; + commandDeployHashStore?: DiscordCommandDeployHashStore; devGuilds?: string[]; eventQueue?: DiscordEventQueueOptions; restCacheTtlMs?: number; @@ -99,7 +100,7 @@ export class Client { clientId: this.options.clientId, commands: this.commands, devGuilds: this.options.devGuilds, - hashStorePath: this.options.commandDeployHashStorePath, + hashStore: this.options.commandDeployHashStore, rest: () => this.rest, }); for (const component of handlers.components ?? []) { diff --git a/extensions/discord/src/internal/command-deploy.test.ts b/extensions/discord/src/internal/command-deploy.test.ts index e5d8e7464c15..fb531bc12df3 100644 --- a/extensions/discord/src/internal/command-deploy.test.ts +++ b/extensions/discord/src/internal/command-deploy.test.ts @@ -1,14 +1,12 @@ // Discord tests cover command deploy plugin behavior. /* oxlint-disable typescript/unbound-method -- vitest mocks of RequestClient methods (createRest) intentionally expose vi.fn refs via `restA.get`/`.post`; not unbound class methods. */ -import fs from "node:fs/promises"; -import os from "node:os"; -import path from "node:path"; import { ApplicationCommandType, type APIApplicationCommand, type APIApplicationCommandOption, } from "discord-api-types/v10"; import { describe, expect, test, vi } from "vitest"; +import type { DiscordCommandDeployHashStore } from "../command-deploy-store.js"; import { commandsEqual } from "./command-comparison.js"; import { DiscordCommandDeployer } from "./command-deploy.js"; import { BaseCommand } from "./commands.js"; @@ -209,24 +207,21 @@ describe("commandsEqual", () => { expect(commandsEqual(current, desired)).toBe(true); }); }); - /** - * Regression for #77359: when two Discord accounts share the same on-disk - * deploy-cache file (the default in multi-bot setups) the persisted hash key - * must be scoped by application/client id. Otherwise a later account whose - * command set hashes the same as the first account's reuses the first - * account's hash and skips reconciling its own Discord application — leaving - * "This application has no commands" in the secondary bot's Integrations panel. + * Regression for #77359: persisted hashes are scoped by Discord application id. + * Identical command sets for separate bots must never suppress each other's deploy. */ -describe("DiscordCommandDeployer cache scoping (multi-application)", () => { +describe("DiscordCommandDeployer SQLite cache", () => { class StaticCommand extends BaseCommand { name: string; override description = "ping the bot"; type = ApplicationCommandType.ChatInput; + constructor(name: string) { super(); this.name = name; } + serializeOptions() { return undefined; } @@ -242,299 +237,158 @@ describe("DiscordCommandDeployer cache scoping (multi-application)", () => { } as unknown as RequestClient; } + function createHashStore(initial: Record = {}): { + rows: Map; + store: DiscordCommandDeployHashStore; + } { + const rows = new Map(Object.entries(initial)); + return { + rows, + store: { + lookup: vi.fn(async (key: string) => rows.get(key)), + register: vi.fn(async (key: string, value: string) => { + rows.set(key, value); + }), + }, + }; + } + test("two applications with identical command sets each reconcile their own application", async () => { - const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-discord-multi-app-")); - const hashStorePath = path.join(dir, "command-deploy-cache.json"); + const { store, rows } = createHashStore(); const commands = [new StaticCommand("ping")]; - - const restA = createRest(); - const deployerA = new DiscordCommandDeployer({ - clientId: "app-default", - commands, - hashStorePath, - rest: () => restA, - }); - await deployerA.deploy({ mode: "reconcile" }); - - const restB = createRest(); - const deployerB = new DiscordCommandDeployer({ - clientId: "app-secondary", - commands, - hashStorePath, - rest: () => restB, - }); - await deployerB.deploy({ mode: "reconcile" }); - - // The first deploy issues a list + create against application "app-default". - expect(restA.get).toHaveBeenCalledTimes(1); - expect(restA.post).toHaveBeenCalledTimes(1); - // The second deploy MUST also list + create against "app-secondary"; before - // the fix it short-circuited on the shared `global:reconcile` hash and - // never touched its own Discord application. - expect(restB.get).toHaveBeenCalledTimes(1); - expect(restB.post).toHaveBeenCalledTimes(1); - }); - - test("re-deploying the same application still hits the persisted cache", async () => { - const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-discord-multi-app-")); - const hashStorePath = path.join(dir, "command-deploy-cache.json"); - const commands = [new StaticCommand("ping")]; - - const restFirst = createRest(); - await new DiscordCommandDeployer({ - clientId: "app-default", - commands, - hashStorePath, - rest: () => restFirst, - }).deploy({ mode: "reconcile" }); - - const restSecond = createRest(); - await new DiscordCommandDeployer({ - clientId: "app-default", - commands, - hashStorePath, - rest: () => restSecond, - }).deploy({ mode: "reconcile" }); - - expect(restFirst.get).toHaveBeenCalledTimes(1); - expect(restFirst.post).toHaveBeenCalledTimes(1); - // Same application, same command set, same hash file => skip reconcile. - expect(restSecond.get).not.toHaveBeenCalled(); - expect(restSecond.post).not.toHaveBeenCalled(); - }); - - test("persisted cache keys are namespaced by application id", async () => { - const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-discord-multi-app-")); - const hashStorePath = path.join(dir, "command-deploy-cache.json"); - const commands = [new StaticCommand("ping")]; - - await new DiscordCommandDeployer({ - clientId: "app-default", - commands, - hashStorePath, - rest: () => createRest(), - }).deploy({ mode: "reconcile" }); - - await new DiscordCommandDeployer({ - clientId: "app-secondary", - commands, - hashStorePath, - rest: () => createRest(), - }).deploy({ mode: "reconcile" }); - - const raw = await fs.readFile(hashStorePath, "utf8"); - const parsed = JSON.parse(raw) as { hashes: Record }; - const keys = Object.keys(parsed.hashes); - expect(keys).toContain("app:app-default:global:reconcile"); - expect(keys).toContain("app:app-secondary:global:reconcile"); - expect(keys).not.toContain("global:reconcile"); - }); - - test("successful deploy repairs a corrupt persisted cache file", async () => { - const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-discord-multi-app-")); - const hashStorePath = path.join(dir, "command-deploy-cache.json"); - await fs.writeFile(hashStorePath, "{not json", "utf8"); - - await new DiscordCommandDeployer({ - clientId: "app-default", - commands: [new StaticCommand("ping")], - hashStorePath, - rest: () => createRest(), - }).deploy({ mode: "reconcile" }); - - const raw = await fs.readFile(hashStorePath, "utf8"); - const parsed = JSON.parse(raw) as { hashes: Record }; - expect(parsed.hashes).toHaveProperty("app:app-default:global:reconcile"); - }); - - test("a deployer that loaded an empty cache before another deployer's write preserves the other deployer's entries on persist", async () => { - // Regression for the codex follow-up on PR #77367: `server-channels.ts` - // can start multiple Discord deployers concurrently. Before the fix, a - // deployer that loaded the (empty) cache file before another deployer's - // first write would later overwrite it on its own `persistHashes()`, - // serializing only its own in-memory `app::...` entry and dropping - // the other deployer's entry. The current implementation re-reads the - // on-disk hashes inside `persistHashes` and merges them with our - // in-memory entries before the rename. - const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-discord-multi-app-")); - const hashStorePath = path.join(dir, "command-deploy-cache.json"); - const commands = [new StaticCommand("ping")]; - - // Deployer B starts first, loads the empty cache. Then deployer A - // completes its full deploy + persist, writing `app:app-default:...` to - // disk. When deployer B finally persists, it must merge in deployer A's - // entry instead of overwriting it with just its own. - const deployerB = new DiscordCommandDeployer({ - clientId: "app-secondary", - commands, - hashStorePath, - rest: () => createRest(), - }); - // Trigger B's load of the (still missing) cache file by starting deploy - // and immediately awaiting just enough to clear the load. The deploy - // call awaits loadPersistedHashes inside putCommandSetIfChanged before - // calling deploy(); to keep the seam minimal here, we just race the load - // by running deployer A's full deploy in between. - const deployerA = new DiscordCommandDeployer({ - clientId: "app-default", - commands, - hashStorePath, - rest: () => createRest(), - }); - - // Step 1: A runs a full deploy (load -> reconcile -> persist) on the - // initially missing cache file; result: file now has app-default entry. - await deployerA.deploy({ mode: "reconcile" }); - - // Step 2: B runs its full deploy. Without the fix, B's persistHashes - // would write only `app:app-secondary:...` and drop A's entry. With the - // fix, B re-reads the on-disk file inside persistHashes, sees A's entry, - // and merges it into the write so both keys survive. - await deployerB.deploy({ mode: "reconcile" }); - - const raw = await fs.readFile(hashStorePath, "utf8"); - const parsed = JSON.parse(raw) as { hashes: Record }; - const keys = Object.keys(parsed.hashes); - expect(keys).toContain("app:app-default:global:reconcile"); - expect(keys).toContain("app:app-secondary:global:reconcile"); - - // And subsequent restarts must still hit the cache for both apps, - // proving the rate-limit protection survived the concurrent write. - const restA = createRest(); - await new DiscordCommandDeployer({ - clientId: "app-default", - commands, - hashStorePath, - rest: () => restA, - }).deploy({ mode: "reconcile" }); - const restB = createRest(); - await new DiscordCommandDeployer({ - clientId: "app-secondary", - commands, - hashStorePath, - rest: () => restB, - }).deploy({ mode: "reconcile" }); - expect(restA.get).not.toHaveBeenCalled(); - expect(restA.post).not.toHaveBeenCalled(); - expect(restB.get).not.toHaveBeenCalled(); - expect(restB.post).not.toHaveBeenCalled(); - }); - - test("truly parallel deployers serialize cache writes via the per-path mutex (codex follow-up on #77367)", async () => { - // Codex follow-up on PR #77367: re-read-before-write alone isn't enough - // when two deployers run `persistHashes` in real parallel — both can read - // the same snapshot before either writes. The in-process per-path mutex - // around the read-merge-write cycle makes the operation atomic. - const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-discord-multi-app-")); - const hashStorePath = path.join(dir, "command-deploy-cache.json"); - const commands = [new StaticCommand("ping")]; - - // Run BOTH deploys with Promise.all on the SAME process tick — pre-fix, - // both `persistHashes` calls would race on read-then-rename and one - // writer's `app::...` entry would be lost. const restA = createRest(); const restB = createRest(); - const restC = createRest(); + await Promise.all([ new DiscordCommandDeployer({ clientId: "app-default", commands, - hashStorePath, + hashStore: store, rest: () => restA, }).deploy({ mode: "reconcile" }), new DiscordCommandDeployer({ clientId: "app-secondary", commands, - hashStorePath, + hashStore: store, rest: () => restB, }).deploy({ mode: "reconcile" }), - new DiscordCommandDeployer({ - clientId: "app-tertiary", - commands, - hashStorePath, - rest: () => restC, - }).deploy({ mode: "reconcile" }), ]); - const raw = await fs.readFile(hashStorePath, "utf8"); - const parsed = JSON.parse(raw) as { hashes: Record }; - const keys = Object.keys(parsed.hashes); - // All three apps' entries must survive — pre-fix, one or two would be - // lost to the race. - expect(keys).toContain("app:app-default:global:reconcile"); - expect(keys).toContain("app:app-secondary:global:reconcile"); - expect(keys).toContain("app:app-tertiary:global:reconcile"); + expect(restA.get).toHaveBeenCalledTimes(1); + expect(restA.post).toHaveBeenCalledTimes(1); + expect(restB.get).toHaveBeenCalledTimes(1); + expect(restB.post).toHaveBeenCalledTimes(1); + expect([...rows.keys()].toSorted()).toEqual([ + "app:app-default:global:reconcile", + "app:app-secondary:global:reconcile", + ]); }); - test("parallel changed deploys preserve fresher sibling cache entries", async () => { - const dir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-discord-multi-app-")); - const hashStorePath = path.join(dir, "command-deploy-cache.json"); - const oldCommands = [new StaticCommand("ping")]; - const newCommands = [new StaticCommand("status")]; + test("skips unchanged command deploys across deployer restarts", async () => { + const { store } = createHashStore(); + const commands = [new StaticCommand("ping")]; + const firstRest = createRest(); await new DiscordCommandDeployer({ clientId: "app-default", - commands: oldCommands, - hashStorePath, - rest: () => createRest(), + commands, + hashStore: store, + rest: () => firstRest, }).deploy({ mode: "reconcile" }); + + const secondRest = createRest(); await new DiscordCommandDeployer({ - clientId: "app-secondary", - commands: oldCommands, - hashStorePath, - rest: () => createRest(), + clientId: "app-default", + commands, + hashStore: store, + rest: () => secondRest, }).deploy({ mode: "reconcile" }); - let postStarts = 0; - let releasePosts: () => void = () => {}; - const bothPostsStarted = new Promise((resolve) => { - releasePosts = resolve; - }); - function createWaitingRest(): RequestClient { - const rest = createRest(); - rest.post = vi.fn(async () => { - postStarts += 1; - if (postStarts === 2) { - releasePosts(); - } - await bothPostsStarted; - }) as RequestClient["post"]; - return rest; - } + expect(firstRest.get).toHaveBeenCalledTimes(1); + expect(firstRest.post).toHaveBeenCalledTimes(1); + expect(secondRest.get).not.toHaveBeenCalled(); + expect(secondRest.post).not.toHaveBeenCalled(); + expect(store.lookup).toHaveBeenLastCalledWith("app:app-default:global:reconcile"); + }); - await Promise.all([ + test("loads only the exact scoped key needed by a deployment", async () => { + const { store } = createHashStore({ + "app:other:global:reconcile": "unrelated", + }); + const rest = createRest(); + + await new DiscordCommandDeployer({ + clientId: "app-default", + commands: [new StaticCommand("ping")], + hashStore: store, + rest: () => rest, + }).deploy({ mode: "reconcile" }); + + expect(store.lookup).toHaveBeenCalledOnce(); + expect(store.lookup).toHaveBeenCalledWith("app:app-default:global:reconcile"); + }); + + test("treats a SQLite lookup failure as a cache miss and repairs the row", async () => { + const register = vi.fn(async () => undefined); + const store: DiscordCommandDeployHashStore = { + lookup: vi.fn(async () => { + throw new Error("database unavailable"); + }), + register, + }; + const rest = createRest(); + + await new DiscordCommandDeployer({ + clientId: "app-default", + commands: [new StaticCommand("ping")], + hashStore: store, + rest: () => rest, + }).deploy({ mode: "reconcile" }); + + expect(rest.get).toHaveBeenCalledTimes(1); + expect(rest.post).toHaveBeenCalledTimes(1); + expect(register).toHaveBeenCalledOnce(); + }); + + test("keeps successful deploys successful when SQLite persistence fails", async () => { + const store: DiscordCommandDeployHashStore = { + lookup: vi.fn(async () => undefined), + register: vi.fn(async () => { + throw new Error("database unavailable"); + }), + }; + const rest = createRest(); + const deployer = new DiscordCommandDeployer({ + clientId: "app-default", + commands: [new StaticCommand("ping")], + hashStore: store, + rest: () => rest, + }); + + await deployer.deploy({ mode: "reconcile" }); + await deployer.deploy({ mode: "reconcile" }); + + expect(rest.get).toHaveBeenCalledTimes(1); + expect(rest.post).toHaveBeenCalledTimes(1); + expect(store.register).toHaveBeenCalledOnce(); + }); + + test("does not persist a hash when Discord deployment fails", async () => { + const { store } = createHashStore(); + const rest = createRest(); + rest.post = vi.fn(async () => { + throw new Error("Discord rejected deploy"); + }) as RequestClient["post"]; + + await expect( new DiscordCommandDeployer({ clientId: "app-default", - commands: newCommands, - hashStorePath, - rest: () => createWaitingRest(), + commands: [new StaticCommand("ping")], + hashStore: store, + rest: () => rest, }).deploy({ mode: "reconcile" }), - new DiscordCommandDeployer({ - clientId: "app-secondary", - commands: newCommands, - hashStorePath, - rest: () => createWaitingRest(), - }).deploy({ mode: "reconcile" }), - ]); + ).rejects.toThrow("Discord rejected deploy"); - const restA = createRest(); - await new DiscordCommandDeployer({ - clientId: "app-default", - commands: newCommands, - hashStorePath, - rest: () => restA, - }).deploy({ mode: "reconcile" }); - const restB = createRest(); - await new DiscordCommandDeployer({ - clientId: "app-secondary", - commands: newCommands, - hashStorePath, - rest: () => restB, - }).deploy({ mode: "reconcile" }); - - expect(restA.get).not.toHaveBeenCalled(); - expect(restA.post).not.toHaveBeenCalled(); - expect(restB.get).not.toHaveBeenCalled(); - expect(restB.post).not.toHaveBeenCalled(); + expect(store.register).not.toHaveBeenCalled(); }); }); diff --git a/extensions/discord/src/internal/command-deploy.ts b/extensions/discord/src/internal/command-deploy.ts index 8d1a4b2d2804..65484a11eb2a 100644 --- a/extensions/discord/src/internal/command-deploy.ts +++ b/extensions/discord/src/internal/command-deploy.ts @@ -1,9 +1,7 @@ // Discord plugin module implements command deploy behavior. import { createHash } from "node:crypto"; -import path from "node:path"; import { ApplicationCommandType, type APIApplicationCommand } from "discord-api-types/v10"; -import { KeyedAsyncQueue } from "openclaw/plugin-sdk/keyed-async-queue"; -import { privateFileStore } from "openclaw/plugin-sdk/security-runtime"; +import type { DiscordCommandDeployHashStore } from "../command-deploy-store.js"; import { createApplicationCommand, deleteApplicationCommand, @@ -25,35 +23,16 @@ type SerializedCommand = ReturnType; const DISCORD_APPLICATION_COMMAND_LIMIT_REACHED = 30032; -/** - * Per-`command-deploy-cache.json` path async mutex. `server-channels.ts` can - * start several Discord deployers concurrently in the same Node.js process; - * each one shares the same on-disk cache file. Without this lock, two - * deployers can run `persistHashes` in parallel, both read the same on-disk - * snapshot before either writes, and the later `rename` then overwrites the - * earlier writer's entries — defeating the rate-limit cache. - * - * This is an in-process lock; cross-process serialization would need an OS - * file lock. Discord deployers only run inside the gateway process, so an - * in-process mutex is sufficient for the documented concurrency surface. - */ -const cachePersistLocks = new KeyedAsyncQueue(); - -async function withCachePersistLock(storePath: string, fn: () => Promise): Promise { - return await cachePersistLocks.enqueue(storePath, fn); -} - export class DiscordCommandDeployer { private readonly hashes = new Map(); - private readonly pendingHashes = new Map(); - private hashesLoaded = false; + private readonly loadedKeys = new Set(); constructor( private readonly params: { clientId: string; commands: BaseCommand[]; devGuilds?: string[]; - hashStorePath?: string; + hashStore?: DiscordCommandDeployHashStore; rest: () => RequestClient; }, ) {} @@ -124,7 +103,7 @@ export class DiscordCommandDeployer { /** * Scope cache keys by Discord application id so multi-bot setups that share a - * single deploy-cache file still reconcile each application separately. The + * single command-deploy store still reconcile each application separately. The * prior unscoped `global:reconcile` / `guild:` keys let a later account * with an identical command set reuse the first account's hash and skip its * own application's reconcile entirely (#77359). @@ -177,105 +156,31 @@ export class DiscordCommandDeployer { options: { force?: boolean }, ): Promise { const hash = stableCommandSetHash(commands); - await this.loadPersistedHashes(); + await this.loadPersistedHash(key); if (!options.force && this.hashes.get(key) === hash) { return; } await deploy(); this.hashes.set(key, hash); - this.pendingHashes.set(key, hash); - await this.persistHashes(); + try { + await this.params.hashStore?.register(key, hash); + } catch { + // Cache persistence must not turn a successful Discord deploy into a startup failure. + } } - private async loadPersistedHashes(): Promise { - if (this.hashesLoaded) { - return; - } - this.hashesLoaded = true; - const storePath = this.params.hashStorePath; - if (!storePath) { + private async loadPersistedHash(key: string): Promise { + if (this.loadedKeys.has(key)) { return; } + this.loadedKeys.add(key); try { - const parsed = await privateFileStore(path.dirname(storePath)).readJsonIfExists<{ - hashes?: unknown; - }>(path.basename(storePath)); - if (!parsed?.hashes || typeof parsed.hashes !== "object") { - return; - } - for (const [key, value] of Object.entries(parsed.hashes)) { - if (typeof value === "string" && key.trim() && value.trim()) { - this.hashes.set(key, value); - } + const hash = await this.params.hashStore?.lookup(key); + if (typeof hash === "string" && hash.trim()) { + this.hashes.set(key, hash); } } catch { - // Best-effort cache only. A corrupt or missing file should never block startup. - } - } - - private async persistHashes(): Promise { - const storePath = this.params.hashStorePath; - if (!storePath) { - return; - } - // Serialize concurrent persists for the same on-disk path. The earlier - // "re-read inside persistHashes" merge alone is not enough — two - // deployers running `persistHashes` in true parallel would both read the - // same snapshot before either writes, and the later `rename` would still - // overwrite the earlier one's `app::...` entries. The mutex makes the - // read-merge-write cycle atomic for in-process callers. - await withCachePersistLock(storePath, async () => { - await this.persistHashesLocked(storePath); - }); - } - - private async persistHashesLocked(storePath: string): Promise { - try { - // Re-read the on-disk hashes immediately before writing and merge only - // keys this deployer changed. Previously loaded hashes can be stale when - // sibling deployers update the same file, so on-disk wins for untouched - // keys while pending keys win because this deployer just produced them. - const storeFile = path.basename(storePath); - const fileStore = privateFileStore(path.dirname(storePath)); - const merged = new Map(); - let onDisk: { hashes?: unknown } | null = null; - try { - onDisk = await fileStore.readJsonIfExists<{ - hashes?: unknown; - }>(storeFile); - } catch { - // A corrupt cache should not become permanent. Treat the re-read as - // empty and replace it with the fresh pending hashes after deploy. - } - if (onDisk?.hashes && typeof onDisk.hashes === "object") { - for (const [key, value] of Object.entries(onDisk.hashes)) { - if (typeof value === "string" && key.trim() && value.trim()) { - merged.set(key, value); - } - } - } - for (const [key, value] of this.pendingHashes.entries()) { - merged.set(key, value); - } - await fileStore.writeJson( - storeFile, - { - version: 1, - updatedAt: new Date().toISOString(), - hashes: Object.fromEntries( - [...merged.entries()].toSorted(([left], [right]) => left.localeCompare(right)), - ), - }, - { trailingNewline: true }, - ); - // Refresh in-memory state so future writes from the same deployer also - // see entries that other deployers added concurrently. - for (const [key, value] of merged.entries()) { - this.hashes.set(key, value); - } - this.pendingHashes.clear(); - } catch { - // The cache is only an optimization to avoid redundant Discord writes. + // Cache lookup failure is a miss. Reconcile repairs the canonical row after success. } } diff --git a/extensions/discord/src/monitor/model-picker-preferences-migrations.test.ts b/extensions/discord/src/monitor/model-picker-preferences-migrations.test.ts index abff98398ea0..5974b281b848 100644 --- a/extensions/discord/src/monitor/model-picker-preferences-migrations.test.ts +++ b/extensions/discord/src/monitor/model-picker-preferences-migrations.test.ts @@ -18,6 +18,37 @@ afterEach(async () => { }); describe("Discord model picker preference migration", () => { + it("plans legacy command deployment cache deletion without importing hashes", async () => { + const stateDir = await makeStateDir(); + const sourcePath = path.join(stateDir, "discord", "command-deploy-cache.json"); + await fs.mkdir(path.dirname(sourcePath), { recursive: true }); + await fs.writeFile(sourcePath, "{malformed cache", "utf8"); + + const plans = await Promise.resolve( + detectDiscordLegacyStateMigrations({ + cfg: {}, + env: {}, + oauthDir: path.join(stateDir, "credentials"), + stateDir, + }), + ); + + expect(plans).toHaveLength(1); + const plan = plans?.[0]; + if (plan?.kind !== "plugin-state-import") { + throw new Error("expected plugin-state import plan"); + } + expect(plan).toMatchObject({ + label: "Discord command deployment cache", + pluginId: "discord", + namespace: "command-deploy-hashes", + maxEntries: 10_000, + cleanupSource: "remove", + cleanupWhenEmpty: true, + }); + expect(await plan.readEntries()).toEqual([]); + }); + it("plans legacy JSON import into plugin state", async () => { const stateDir = await makeStateDir(); const sourcePath = path.join(stateDir, "discord", "model-picker-preferences.json"); diff --git a/extensions/discord/src/monitor/model-picker-preferences-migrations.ts b/extensions/discord/src/monitor/model-picker-preferences-migrations.ts index a00f7f98d292..062faedbef58 100644 --- a/extensions/discord/src/monitor/model-picker-preferences-migrations.ts +++ b/extensions/discord/src/monitor/model-picker-preferences-migrations.ts @@ -6,6 +6,10 @@ import type { ChannelLegacyStateMigrationPlan } from "openclaw/plugin-sdk/channe import type { BundledChannelLegacyStateMigrationDetector } from "openclaw/plugin-sdk/channel-entry-contract"; import { MAX_DATE_TIMESTAMP_MS, timestampMsToIsoString } from "openclaw/plugin-sdk/number-runtime"; import { normalizeProviderId } from "openclaw/plugin-sdk/provider-model-shared"; +import { + DISCORD_COMMAND_DEPLOY_HASH_MAX_ENTRIES, + DISCORD_COMMAND_DEPLOY_HASH_NAMESPACE, +} from "../command-deploy-store.js"; import { normalizePersistedBinding, THREAD_BINDINGS_MAX_ENTRIES, @@ -167,6 +171,23 @@ export const detectDiscordLegacyStateMigrations: BundledChannelLegacyStateMigrat stateDir, }) => { const plans: ChannelLegacyStateMigrationPlan[] = []; + const commandDeployCacheSourcePath = path.join(stateDir, "discord", "command-deploy-cache.json"); + if (fileExists(commandDeployCacheSourcePath)) { + plans.push({ + kind: "plugin-state-import", + label: "Discord command deployment cache", + sourcePath: commandDeployCacheSourcePath, + targetPath: `plugin state:${DISCORD_COMMAND_DEPLOY_HASH_NAMESPACE}`, + pluginId: "discord", + namespace: DISCORD_COMMAND_DEPLOY_HASH_NAMESPACE, + maxEntries: DISCORD_COMMAND_DEPLOY_HASH_MAX_ENTRIES, + scopeKey: "", + cleanupSource: "remove", + cleanupWhenEmpty: true, + // Rebuildable cache: discard file-era hashes and reconcile once against Discord. + readEntries: () => [], + }); + } const modelPickerSourcePath = path.join(stateDir, "discord", "model-picker-preferences.json"); if (fileExists(modelPickerSourcePath)) { plans.push({ diff --git a/extensions/discord/src/monitor/provider.startup.test.ts b/extensions/discord/src/monitor/provider.startup.test.ts index 491e1b347569..ad9959aec428 100644 --- a/extensions/discord/src/monitor/provider.startup.test.ts +++ b/extensions/discord/src/monitor/provider.startup.test.ts @@ -239,6 +239,10 @@ describe("createDiscordMonitorClient", () => { it("configures internal Discord REST options explicitly", async () => { const createClient = vi.fn(createClientWithPlugins); + const commandDeployHashStore = { + lookup: vi.fn(async () => undefined), + register: vi.fn(async () => undefined), + }; await createDiscordMonitorClient({ accountId: "default", @@ -250,6 +254,7 @@ describe("createDiscordMonitorClient", () => { voiceEnabled: false, discordConfig: {}, runtime: createRuntime(), + commandDeployHashStore, createClient, createGatewayPlugin: () => ({ id: "gateway" }) as never, createGatewaySupervisor: () => ({ shutdown: vi.fn(), handleError: vi.fn() }) as never, @@ -264,6 +269,9 @@ describe("createDiscordMonitorClient", () => { runtimeProfile: "persistent", maxQueueSize: 1000, }); + expect((options as { commandDeployHashStore?: unknown }).commandDeployHashStore).toBe( + commandDeployHashStore, + ); if (!handlers) { throw new Error("expected Discord client handlers"); } diff --git a/extensions/discord/src/monitor/provider.startup.ts b/extensions/discord/src/monitor/provider.startup.ts index b57f2c5a57e7..f459abe6dada 100644 --- a/extensions/discord/src/monitor/provider.startup.ts +++ b/extensions/discord/src/monitor/provider.startup.ts @@ -1,11 +1,10 @@ // Discord provider module implements model/runtime integration. -import path from "node:path"; import type { OpenClawConfig } from "openclaw/plugin-sdk/config-contracts"; import { isDangerousNameMatchingEnabled } from "openclaw/plugin-sdk/dangerous-name-runtime"; import { danger } from "openclaw/plugin-sdk/runtime-env"; import type { RuntimeEnv } from "openclaw/plugin-sdk/runtime-env"; -import { resolveStateDir } from "openclaw/plugin-sdk/state-paths"; import { normalizeOptionalString } from "openclaw/plugin-sdk/string-coerce-runtime"; +import type { DiscordCommandDeployHashStore } from "../command-deploy-store.js"; import { Client, ReadyListener, @@ -104,6 +103,7 @@ export async function createDiscordMonitorClient(params: { >; }; runtime: RuntimeEnv; + commandDeployHashStore?: DiscordCommandDeployHashStore; createClient: CreateClientFn; createGatewayPlugin: typeof createDiscordGatewayPlugin; createGatewaySupervisor: typeof createDiscordGatewaySupervisor; @@ -142,11 +142,7 @@ export async function createDiscordMonitorClient(params: { publicKey: "a", token: params.token, autoDeploy: false, - commandDeployHashStorePath: path.join( - resolveStateDir(process.env), - "discord", - "command-deploy-cache.json", - ), + commandDeployHashStore: params.commandDeployHashStore, requestOptions: { timeout: DISCORD_REST_TIMEOUT_MS, runtimeProfile: "persistent", diff --git a/extensions/discord/src/monitor/provider.ts b/extensions/discord/src/monitor/provider.ts index b8e0bff0c12e..3e71c9f285f6 100644 --- a/extensions/discord/src/monitor/provider.ts +++ b/extensions/discord/src/monitor/provider.ts @@ -16,6 +16,7 @@ import { } from "openclaw/plugin-sdk/runtime-group-policy"; import { formatErrorMessage } from "openclaw/plugin-sdk/ssrf-runtime"; import { resolveDiscordAccountAllowFrom, resolveDiscordAccountDmPolicy } from "../accounts.js"; +import type { DiscordCommandDeployHashStore } from "../command-deploy-store.js"; import { GatewayCloseCodes } from "../internal/gateway.js"; import { parseApplicationIdFromToken } from "../probe.js"; import { normalizeDiscordToken } from "../token.js"; @@ -55,6 +56,7 @@ export type MonitorDiscordOpts = { historyLimit?: number; replyToMode?: ReplyToMode; setStatus?: DiscordMonitorStatusSink; + commandDeployHashStore?: DiscordCommandDeployHashStore; }; const DEFAULT_DISCORD_MEDIA_MAX_MB = 100; @@ -336,6 +338,7 @@ export async function monitorDiscordProvider(opts: MonitorDiscordOpts = {}) { voiceEnabled, discordConfig: discordCfg, runtime, + commandDeployHashStore: opts.commandDeployHashStore, createClient: discordProviderRuntime.createClient, createGatewayPlugin: createDiscordGatewayPlugin, createGatewaySupervisor: createDiscordGatewaySupervisor, diff --git a/src/channels/plugins/legacy-state-migration.types.ts b/src/channels/plugins/legacy-state-migration.types.ts index c439451e022e..e7b8604b47ca 100644 --- a/src/channels/plugins/legacy-state-migration.types.ts +++ b/src/channels/plugins/legacy-state-migration.types.ts @@ -16,7 +16,7 @@ export type ChannelLegacyStateMigrationPlan = defaultTtlMs?: number; scopeKey: string; stateDir?: string; - cleanupSource?: "rename"; + cleanupSource?: "rename" | "remove"; cleanupWhenEmpty?: boolean; /** Deletes a non-file legacy source (e.g. plugin-state rows) once all entries are covered. */ removeSource?: () => void | Promise; diff --git a/src/commands/doctor-state-migrations.test.ts b/src/commands/doctor-state-migrations.test.ts index cfb09644f936..e4b37d90d42d 100644 --- a/src/commands/doctor-state-migrations.test.ts +++ b/src/commands/doctor-state-migrations.test.ts @@ -1576,6 +1576,39 @@ describe("doctor legacy state migrations", () => { ); }); + it("deletes rebuildable legacy files after the SQLite target opens", async () => { + const root = await makeTempRoot(); + const sourcePath = path.join(root, "command-deploy-cache.json"); + fs.writeFileSync(sourcePath, "{malformed cache", "utf8"); + mockedChannelMigrationPlans.plans = [ + { + kind: "plugin-state-import", + label: "Test rebuildable cache", + sourcePath, + targetPath: "plugin state:test.rebuildable-cache", + pluginId: "discord", + namespace: "test.rebuildable-cache", + maxEntries: 4, + scopeKey: "", + cleanupSource: "remove", + cleanupWhenEmpty: true, + readEntries: () => [], + }, + ]; + + const detected = await detectLegacyStateMigrations({ + cfg: {}, + env: { OPENCLAW_STATE_DIR: root } as NodeJS.ProcessEnv, + }); + const result = await runLegacyStateMigrations({ detected }); + + expect(result.warnings).toStrictEqual([]); + expect(fs.existsSync(sourcePath)).toBe(false); + expect(result.changes).toContain( + `Removed Test rebuildable cache legacy source (${sourcePath})`, + ); + }); + it("replaces existing plugin-state entries when a channel import plan asks for it", async () => { const root = await makeTempRoot(); const sourcePath = path.join(root, "legacy-cache.json"); diff --git a/src/infra/state-migrations.plugin-state.ts b/src/infra/state-migrations.plugin-state.ts index b60edb767f3f..266398871d90 100644 --- a/src/infra/state-migrations.plugin-state.ts +++ b/src/infra/state-migrations.plugin-state.ts @@ -406,6 +406,14 @@ export async function runLegacyMigrationPlans( warnings, }); } + if (allEntriesCovered && plan.cleanupSource === "remove" && fileExists(plan.sourcePath)) { + try { + fs.unlinkSync(plan.sourcePath); + changes.push(`Removed ${plan.label} legacy source (${plan.sourcePath})`); + } catch (err) { + warnings.push(`Failed removing ${plan.label} legacy source: ${String(err)}`); + } + } if (allEntriesCovered && plan.removeSource) { try { await plan.removeSource();