fix(matrix): keep bot discovery off the Gateway thread (#152122)

* fix(matrix): keep bot discovery off the Gateway thread

* fix(matrix): preserve bot discovery environment selectors

Prepare account authentication against native environment lookups before asynchronous reads. Capture the resolved credential state directory and supervisor setting once through existing host APIs, preserving Windows lookup behavior and the supported host floor without expanding the SDK.

* test(matrix): follow the runtime store options contract
This commit is contained in:
Peter Steinberger 2026-09-22 06:56:58 -07:00 • committed by GitHub
parent 976c5bfa46
commit 49d7f6c40d
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 400 additions and 58 deletions

View file

@ -527,9 +527,13 @@ completed imports and unrelated files. Device backfill remains nonblocking at st
monitor retirement cancels and joins it before releasing storage. Hosts without
data-only comparison support retain the existing native metadata and import decisions
under the declared plugin API floor. Worker failures never select that fallback.
Monitor bot-account discovery awaits selected credential reads after authentication,
then checks cancellation before acquiring the client or installing handlers. Each
identity uses one observed credential and the shared account-readiness rules; no
namespace scan or new bulk-read capability is required.
Approval actor and reaction approver lists resolve from account configuration without
reading credentials; native delivery eligibility still checks enabled and configured
account readiness. Synchronous credential readiness and package auth-presence probes
account readiness. Synchronous public account readiness and package auth-presence probes
retain their separate SDK contracts.
Node-host launch and turn journals execute on the same shared-state worker.

View file

@ -17,8 +17,11 @@ const loadMatrixCredentialsMock = vi.hoisted(() =>
);
vi.mock("./credentials-read.js", () => ({
captureMatrixCredentialsEnv: (env: NodeJS.ProcessEnv) => env,
loadMatrixCredentials: (env?: NodeJS.ProcessEnv, accountId?: string | null) =>
loadMatrixCredentialsMock(env, accountId),
loadMatrixCredentialsAsync: async (env?: NodeJS.ProcessEnv, accountId?: string | null) =>
loadMatrixCredentialsMock(env, accountId),
credentialsMatchConfig: () => false,
}));
@ -462,7 +465,7 @@ describe("resolveMatrixAccount", () => {
expect(resolveDefaultMatrixAccountId(cfg)).toBe("default");
});
it("collects other configured Matrix account user ids for bot detection", () => {
it("collects other configured Matrix account user ids for bot detection", async () => {
const cfg: CoreConfig = {
channels: {
matrix: {
@ -486,11 +489,11 @@ describe("resolveMatrixAccount", () => {
};
expect(
Array.from(resolveConfiguredMatrixBotUserIds({ cfg, accountId: "ops" })).toSorted(),
Array.from(await resolveConfiguredMatrixBotUserIds({ cfg, accountId: "ops" })).toSorted(),
).toEqual(["@alerts:example.org", "@main:example.org"]);
});
it("honors injected env when detecting configured bot accounts", () => {
it("honors injected env when detecting configured bot accounts", async () => {
const env = {
MATRIX_HOMESERVER: "https://matrix.example.org",
MATRIX_USER_ID: "@main:example.org",
@ -507,11 +510,13 @@ describe("resolveMatrixAccount", () => {
};
expect(
Array.from(resolveConfiguredMatrixBotUserIds({ cfg, accountId: "ops", env })).toSorted(),
Array.from(
await resolveConfiguredMatrixBotUserIds({ cfg, accountId: "ops", env }),
).toSorted(),
).toEqual(["@alerts:example.org", "@main:example.org"]);
});
it("falls back to stored credentials when an access-token-only account omits userId", () => {
it("falls back to stored credentials when an access-token-only account omits userId", async () => {
loadMatrixCredentialsMock.mockImplementation(
(env?: NodeJS.ProcessEnv, accountId?: string | null) =>
accountId === "ops"
@ -540,9 +545,9 @@ describe("resolveMatrixAccount", () => {
},
};
expect(Array.from(resolveConfiguredMatrixBotUserIds({ cfg, accountId: "default" }))).toEqual([
"@ops:example.org",
]);
expect(
Array.from(await resolveConfiguredMatrixBotUserIds({ cfg, accountId: "default" })),
).toEqual(["@ops:example.org"]);
});
it("preserves shared nested dm and actions config when an account overrides one field", () => {

View file

@ -13,7 +13,13 @@ import {
resolveMatrixBaseConfig,
} from "./account-config.js";
import { resolveGlobalMatrixEnvConfig, resolveScopedMatrixEnvConfig } from "./client/env-auth.js";
import { credentialsMatchConfig, loadMatrixCredentials } from "./credentials-read.js";
import {
captureMatrixCredentialsEnv,
credentialsMatchConfig,
loadMatrixCredentials,
loadMatrixCredentialsAsync,
} from "./credentials-read.js";
import type { MatrixStoredCredentials } from "./credentials-state.js";
export type ResolvedMatrixAccount = {
accountId: string;
@ -71,23 +77,14 @@ function resolveMatrixAccountAuthView(params: {
};
}
function resolveMatrixAccountUserId(params: {
cfg: CoreConfig;
accountId: string;
env?: NodeJS.ProcessEnv;
}): string | null {
const env = params.env ?? process.env;
const authView = resolveMatrixAccountAuthView({
cfg: params.cfg,
accountId: params.accountId,
env,
});
function resolveMatrixAccountUserId(
authView: ReturnType<typeof resolveMatrixAccountAuthView>,
stored: MatrixStoredCredentials | null,
): string | null {
const configuredUserId = authView.userId.trim();
if (configuredUserId) {
return configuredUserId;
}
const stored = loadMatrixCredentials(env, params.accountId);
if (!stored) {
return null;
}
@ -109,44 +106,48 @@ export function resolveDefaultMatrixAccountId(cfg: CoreConfig): string {
return normalizeAccountId(resolveMatrixDefaultOrOnlyAccountId(cfg));
}
export function resolveConfiguredMatrixBotUserIds(params: {
export async function resolveConfiguredMatrixBotUserIds(params: {
cfg: CoreConfig;
accountId?: string | null;
env?: NodeJS.ProcessEnv;
}): Set<string> {
abortSignal?: AbortSignal;
}): Promise<Set<string>> {
const env = params.env ?? process.env;
const currentAccountId = normalizeAccountId(params.accountId);
const accountIds = new Set(resolveConfiguredMatrixAccountIds(params.cfg, env));
if (resolveMatrixAccount({ cfg: params.cfg, accountId: DEFAULT_ACCOUNT_ID, env }).configured) {
accountIds.add(DEFAULT_ACCOUNT_ID);
}
const accountIds = new Set([
...resolveConfiguredMatrixAccountIds(params.cfg, env),
DEFAULT_ACCOUNT_ID,
]);
// Capture config/env facts before storage yields; each identity uses one credential observation.
const accounts = [...accountIds]
.filter((accountId) => normalizeAccountId(accountId) !== currentAccountId)
.map((accountId) => prepareMatrixAccount({ cfg: params.cfg, accountId, env }));
const ids = new Set<string>();
for (const accountId of accountIds) {
if (normalizeAccountId(accountId) === currentAccountId) {
if (accounts.length === 0 || params.abortSignal?.aborted) {
return ids;
}
const credentialsEnv = captureMatrixCredentialsEnv(env);
for (const prepared of accounts) {
if (params.abortSignal?.aborted) {
break;
}
const stored = await loadMatrixCredentialsAsync(credentialsEnv, prepared.account.accountId);
if (!isMatrixAccountConfigured(prepared, stored)) {
continue;
}
if (!resolveMatrixAccount({ cfg: params.cfg, accountId, env }).configured) {
continue;
}
const userId = resolveMatrixAccountUserId({
cfg: params.cfg,
accountId,
env,
});
const userId = resolveMatrixAccountUserId(prepared.authView, stored);
if (userId) {
ids.add(userId);
}
}
return ids;
}
export function resolveMatrixAccount(params: {
function prepareMatrixAccount(params: {
cfg: CoreConfig;
accountId?: string | null;
env?: NodeJS.ProcessEnv;
}): ResolvedMatrixAccount {
}) {
const env = params.env ?? process.env;
const accountId = normalizeAccountId(
params.accountId ?? resolveDefaultMatrixAccountId(params.cfg),
@ -171,7 +172,26 @@ export function resolveMatrixAccount(params: {
const hasPassword = Boolean(authView.password);
const hasPasswordAuth =
hasUserId && (hasPassword || hasConfiguredSecretInput(explicitAuthConfig.password));
const stored = loadMatrixCredentials(env, accountId);
return {
authView,
hasHomeserver,
hasConfiguredAuth: hasAccessToken || hasPasswordAuth,
account: {
accountId,
enabled,
name: normalizeOptionalString(base.name),
homeserver: authView.homeserver || undefined,
userId: authView.userId || undefined,
config: base,
},
};
}
function isMatrixAccountConfigured(
prepared: ReturnType<typeof prepareMatrixAccount>,
stored: MatrixStoredCredentials | null,
): boolean {
const { authView } = prepared;
const hasStored =
stored && authView.homeserver
? credentialsMatchConfig(stored, {
@ -179,16 +199,17 @@ export function resolveMatrixAccount(params: {
userId: authView.userId || "",
})
: false;
const configured = hasHomeserver && (hasAccessToken || hasPasswordAuth || hasStored);
return {
accountId,
enabled,
name: normalizeOptionalString(base.name),
configured,
homeserver: authView.homeserver || undefined,
userId: authView.userId || undefined,
config: base,
};
return prepared.hasHomeserver && (prepared.hasConfiguredAuth || hasStored);
}
export function resolveMatrixAccount(params: {
cfg: CoreConfig;
accountId?: string | null;
env?: NodeJS.ProcessEnv;
}): ResolvedMatrixAccount {
const prepared = prepareMatrixAccount(params);
const stored = loadMatrixCredentials(params.env ?? process.env, prepared.account.accountId);
return { ...prepared.account, configured: isMatrixAccountConfigured(prepared, stored) };
}
export { resolveMatrixAccountConfig } from "./account-config.js";

View file

@ -42,6 +42,15 @@ export function openMatrixCredentialsAsyncStore(env: NodeJS.ProcessEnv = process
);
}
export function captureMatrixCredentialsEnv(env: NodeJS.ProcessEnv): NodeJS.ProcessEnv {
// Resolve selectors before awaiting reads; spreading Windows process.env loses its lookup semantics.
return {
...env,
OPENCLAW_STATE_DIR: getMatrixRuntime().state.resolveStateDir(env),
OPENCLAW_SUPERVISOR_MODE: env.OPENCLAW_SUPERVISOR_MODE,
};
}
export async function loadMatrixCredentialsAsync(
env: NodeJS.ProcessEnv = process.env,
accountId?: string | null,

View file

@ -6,10 +6,12 @@ import type { OpenAsyncKeyedStoreOptions } from "openclaw/plugin-sdk/plugin-stat
import { createPluginStateSyncKeyedStore } from "openclaw/plugin-sdk/plugin-state-store-runtime";
import { resetPluginStateStoreForTests } from "openclaw/plugin-sdk/plugin-state-test-runtime";
import { closeOpenClawStateDatabaseAsync } from "openclaw/plugin-sdk/sqlite-runtime-testing";
import { resolveStateDir } from "openclaw/plugin-sdk/state-paths";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { hasAnyMatrixAuth } from "../../auth-presence.js";
import { getMatrixRuntime } from "../runtime.js";
import { installMatrixTestRuntime } from "../test-runtime.js";
import { resolveConfiguredMatrixBotUserIds } from "./accounts.js";
import { loadMatrixCredentialsAsync, openMatrixCredentialsStore } from "./credentials-read.js";
import {
clearMatrixCredentials,
@ -75,6 +77,91 @@ describe("matrix credentials storage", () => {
expect(fs.existsSync(path.join(stateDir, "credentials", "matrix"))).toBe(false);
});
it.each([
{ platform: "win32", mixedCase: true, expected: ["@alerts:example.org", "@main:example.org"] },
{ platform: "linux", mixedCase: true, expected: [] },
{ platform: "linux", mixedCase: false, expected: ["@alerts:example.org", "@main:example.org"] },
] as const)(
"keeps $platform mixedCase=$mixedCase bot discovery on its captured credential source",
async ({ platform, mixedCase, expected }) => {
for (const [accountId, userId] of [
["default", "@main:example.org"],
["alerts", "@alerts:example.org"],
] as const) {
await saveMatrixCredentials(
{
homeserver: "https://matrix.example.org",
userId,
accessToken: "synthetic-token",
},
{ OPENCLAW_STATE_DIR: stateDir },
accountId,
);
}
const values = {
Matrix_Homeserver: "https://matrix.example.org",
Matrix_Access_Token: "synthetic-token",
Matrix_Alerts_Homeserver: "https://matrix.example.org",
Matrix_Alerts_Access_Token: "synthetic-token",
OpenClaw_State_Dir: stateDir,
OpenClaw_Supervisor_Mode: "internal",
};
const env: NodeJS.ProcessEnv = Object.fromEntries(
Object.entries(values).map(([key, value]) => [mixedCase ? key : key.toUpperCase(), value]),
);
if (platform === "win32" && mixedCase) {
// Windows lookups ignore casing without changing the enumerated key spelling.
for (const key of Object.keys(env)) {
Object.defineProperty(env, key.toUpperCase(), {
get: () => env[key],
set: (value: string | undefined) => {
env[key] = value;
},
});
}
}
const runtime = getMatrixRuntime();
vi.spyOn(runtime.state, "resolveStateDir").mockImplementation((input) =>
input?.OPENCLAW_STATE_DIR ? resolveStateDir(input) : stateDir,
);
const openStore = runtime.state.openKeyedStore.bind(runtime.state);
const observedSources: Array<{ root: string | undefined; supervisor: string | undefined }> =
[];
vi.spyOn(runtime.state, "openKeyedStore").mockImplementation(
<T>(options: Parameters<typeof runtime.state.openKeyedStore>[0]) => {
observedSources.push({
root: options.env?.OPENCLAW_STATE_DIR,
supervisor: options.env?.OPENCLAW_SUPERVISOR_MODE,
});
const store = openStore<T>(options);
const lookup = store.lookup.bind(store);
return {
...store,
lookup: async (key) => {
const value = await lookup(key);
env.OPENCLAW_STATE_DIR = path.join(stateDir, "replacement");
env.OPENCLAW_SUPERVISOR_MODE = "external";
env.MATRIX_ALERTS_ACCESS_TOKEN = "replacement-token";
return value;
},
};
},
);
const ids = await resolveConfiguredMatrixBotUserIds({
cfg: { channels: { matrix: { accounts: { alerts: {} } } } },
accountId: "ops",
env,
});
expect([...ids].toSorted()).toEqual(expected);
expect(observedSources).toEqual([
{ root: stateDir, supervisor: mixedCase && platform === "linux" ? undefined : "internal" },
{ root: stateDir, supervisor: mixedCase && platform === "linux" ? undefined : "internal" },
]);
expect(fs.existsSync(path.join(stateDir, "replacement"))).toBe(false);
},
);
it("touch updates lastUsedAt while preserving createdAt", async () => {
vi.useFakeTimers();
vi.setSystemTime(new Date("2026-03-01T10:00:00.000Z"));

View file

@ -0,0 +1,208 @@
import { DatabaseSync, StatementSync } from "node:sqlite";
import { setImmediate } from "node:timers/promises";
import {
implicitMentionKindWhen,
resolveInboundMentionDecision,
} from "openclaw/plugin-sdk/channel-mention-gating";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
import { closeOpenClawStateDatabaseAsync } from "openclaw/plugin-sdk/sqlite-runtime-testing";
import { useAutoCleanupTempDirTracker } from "openclaw/plugin-sdk/test-env";
import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest";
import type { CoreConfig } from "../../types.js";
import { getMatrixMonitorIndexTestHarness } from "./index.test-helpers.js";
const harness = getMatrixMonitorIndexTestHarness();
let monitorMatrixProvider: typeof import("./index.js").monitorMatrixProvider;
let getMatrixRuntime: typeof import("../../runtime.js").getMatrixRuntime;
let credentials: typeof import("../credentials-read.js");
let installMatrixTestRuntime: typeof import("../../test-runtime.js").installMatrixTestRuntime;
beforeAll(async () => {
vi.doUnmock("../accounts.js");
vi.doUnmock("../../runtime.js");
vi.doUnmock("./handler.js");
({ installMatrixTestRuntime } = await import("../../test-runtime.js"));
({ getMatrixRuntime } = await import("../../runtime.js"));
credentials = await import("../credentials-read.js");
({ monitorMatrixProvider } = await import("./index.js"));
});
describe("Matrix monitor credential discovery", () => {
let cfg: CoreConfig;
const tempDirs = useAutoCleanupTempDirTracker((cleanup) =>
afterEach(async () => {
vi.restoreAllMocks();
await closeOpenClawStateDatabaseAsync();
vi.unstubAllEnvs();
cleanup();
}),
);
beforeEach(() => {
vi.clearAllMocks();
harness.callOrder.length = 0;
harness.state.leaseAbortController = new AbortController();
harness.state.monitorRetirement = null;
harness.state.monitorRetirementPromise = null;
harness.registeredOnRoomMessage = null;
harness.client.removeAllListeners();
Object.assign(harness.client, { getUserId: async () => "@bot:example.org" });
const stateDir = tempDirs.make("matrix-monitor-credentials-");
vi.stubEnv("OPENCLAW_STATE_DIR", stateDir);
cfg = {
channels: {
matrix: {
homeserver: "https://matrix.example.org",
userId: "@bot:example.org",
accessToken: "synthetic-main-token",
allowBots: false,
accounts: {
ops: { homeserver: "https://matrix.example.org", accessToken: "synthetic-ops-token" },
},
},
},
};
installMatrixTestRuntime({
stateDir,
cfg,
logging: { getChildLogger: () => harness.logger, shouldLogVerbose: () => true },
channel: {
mentions: {
buildMentionRegexes: () => [],
matchesMentionPatterns: () => false,
matchesMentionWithExplicit: () => false,
implicitMentionKindWhen,
resolveInboundMentionDecision,
},
},
});
Object.assign(getMatrixRuntime(), { system: { formatNativeDependencyHint: () => "" } });
credentials.openMatrixCredentialsStore().register("account:ops", {
accountId: "ops",
homeserver: "https://matrix.example.org",
userId: "@ops:example.org",
accessToken: "synthetic-ops-token",
createdAt: "2026-09-01T00:00:00.000Z",
});
});
it.each(["matching", "changed-token", "changed-homeserver", "revoked"])(
"uses %s credentials for real bot ingress without host SQLite",
async (kind) => {
if (kind === "revoked") {
credentials.openMatrixCredentialsStore().register("account:ops", {
accountId: "ops",
kind: "revoked",
revokedAt: "2026-09-02T00:00:00.000Z",
});
} else if (kind !== "matching") {
const stored = credentials.loadMatrixCredentials(undefined, "ops");
if (!stored) {
throw new Error("missing seeded credential");
}
credentials.openMatrixCredentialsStore().register("account:ops", {
...stored,
accountId: "ops",
...(kind === "changed-token"
? { accessToken: "synthetic-replacement-token" }
: { homeserver: "https://other.example.org" }),
});
}
await closeOpenClawStateDatabaseAsync();
const started = createDeferred<void>();
harness.resolveSharedMatrixClient.mockImplementationOnce(async () => {
const client = await harness.resolveSharedMatrixClientImpl();
started.resolve();
return client;
});
const sql = [
vi.spyOn(DatabaseSync.prototype, "prepare"),
vi.spyOn(DatabaseSync.prototype, "exec"),
...(["get", "all", "run", "iterate"] as const).map((method) =>
vi.spyOn(StatementSync.prototype, method),
),
];
const controller = new AbortController();
const monitoring = monitorMatrixProvider({ abortSignal: controller.signal });
try {
await Promise.race([
started.promise,
monitoring.then(() => {
throw new Error("monitor stopped before startup");
}),
]);
expect(harness.registeredOnRoomMessage).not.toBeNull();
await harness.registeredOnRoomMessage?.("!room:example.org", {
type: "m.room.message",
event_id: "$synthetic-bot",
sender: "@ops:example.org",
content: { msgtype: "m.text", body: "synthetic bot message" },
});
const botDrop = expect.stringContaining("drop configured bot sender=@ops:example.org");
if (kind === "matching") {
expect(harness.logger.debug).toHaveBeenCalledWith(botDrop);
} else {
expect(harness.logger.debug).not.toHaveBeenCalledWith(botDrop);
expect(harness.logger.debug).toHaveBeenCalledWith(
expect.stringContaining("no allowlist"),
);
}
expect(harness.inboundReplayClaim.commit).toHaveBeenCalledOnce();
expect(harness.logger.error).not.toHaveBeenCalled();
console.log(
"matrix-monitor bot discovery host SQL",
kind,
sql.map((spy) => spy.mock.calls.length),
);
for (const spy of sql) {
expect(spy).not.toHaveBeenCalled();
}
} finally {
controller.abort();
await monitoring;
}
},
);
it("joins admitted credential discovery after abort without starting a client", async () => {
const accounts = cfg.channels?.matrix?.accounts;
if (!accounts) {
throw new Error("missing synthetic Matrix accounts");
}
accounts.secondary = {
homeserver: "https://matrix.example.org",
accessToken: "synthetic-secondary-token",
};
const observed = createDeferred<void>();
const release = createDeferred<void>();
const read = credentials.loadMatrixCredentialsAsync;
const lookup = vi
.spyOn(credentials, "loadMatrixCredentialsAsync")
.mockImplementation(async (...args) => {
const stored = await read(...args);
observed.resolve();
await release.promise;
return stored;
});
const controller = new AbortController();
const monitoring = monitorMatrixProvider({ abortSignal: controller.signal });
let settled = false;
void monitoring.then(() => {
settled = true;
});
try {
await observed.promise;
controller.abort();
await setImmediate();
expect(settled).toBe(false);
expect(harness.acquireSharedMatrixClient).not.toHaveBeenCalled();
expect(harness.registeredOnRoomMessage).toBeNull();
} finally {
release.resolve();
controller.abort();
await monitoring;
}
expect(lookup.mock.calls.map(([, accountId]) => accountId)).toEqual(["ops"]);
expect(harness.acquireSharedMatrixClient).not.toHaveBeenCalled();
expect(harness.registerMatrixMonitorEvents).not.toHaveBeenCalled();
expect(harness.resolveSharedMatrixClient).not.toHaveBeenCalled();
});
});

View file

@ -149,10 +149,6 @@ export async function monitorMatrixProvider(opts: MonitorMatrixOpts = {}): Promi
let needsRoomAliasesForConfig = false;
const initialAllowFrom = (accountConfig.dm?.allowFrom ?? []).map(String);
const initialGroupAllowFrom = (accountConfig.groupAllowFrom ?? []).map(String);
const configuredBotUserIds = resolveConfiguredMatrixBotUserIds({
cfg,
accountId: effectiveAccountId,
});
const {
allowFrom,
@ -190,6 +186,17 @@ export async function monitorMatrixProvider(opts: MonitorMatrixOpts = {}): Promi
};
const auth = await resolveMatrixAuth({ cfg, accountId: effectiveAccountId });
if (opts.abortSignal?.aborted) {
return;
}
const configuredBotUserIds = await resolveConfiguredMatrixBotUserIds({
cfg,
accountId: effectiveAccountId,
abortSignal: opts.abortSignal,
});
if (opts.abortSignal?.aborted) {
return;
}
const resolvedInitialSyncLimit =
resolveOptionalIntegerOption(opts.initialSyncLimit, { min: 0 }) ?? auth.initialSyncLimit;
const authWithLimit =

View file

@ -195,6 +195,7 @@ export const databaseWorkerExtensionTestFiles = [
"extensions/matrix/doctor-contract-api.test.ts",
"extensions/matrix/src/approval-config.storage.test.ts",
"extensions/matrix/src/matrix/client/create-client.storage.test.ts",
"extensions/matrix/src/matrix/monitor/index.credentials.test.ts",
"extensions/matrix/src/matrix/client/file-sync-store.test.ts",
"extensions/matrix/src/matrix/client/file-sync-store.sdk.test.ts",
"extensions/matrix/src/matrix/client/storage.test.ts",