fix(zalouser): avoid Gateway stalls during credential persistence (#156003)

* fix(zalouser): avoid Gateway stalls during credential persistence

Move QR credential saves and logout revocation onto the shared SQLite worker,
await their durable completion, and keep restoration behind pending revocation
on its originally selected state directory. Preserve the credential format,
revocation markers, refresh CAS, and named older-host capability fallback.

Carry a live hosted-wizard assertion through setup and into the final write
grant so a detached QR login cannot persist after its owner is retired.

Validation: 172 focused tests across worker storage, registered setup, synthetic
SDK transport, and core wizard lifecycle; complete independent P2 review clean.
Repository gates: 34 passed, 3 unchanged core gates reused. The inherited unused
TSGO_CORE_TEST_MAX_ROOTS export is the sole remaining gate failure on this base;
canonical fix 1619aff478 supplies its repair.

Regression controls prove the former host-SQL logout path, a missing final
wizard-lifetime guard, and state-root drift across a pending logout. No live
provider account, operator Gateway restart, schema change, or protocol bump.

* test(gateway): await recap publication before projection assertion

Wait for the existing owner's next publication after releasing the model result,
assert its captured state without filtering successful outcomes, and check one
direct describe response. Preserve the durable watermark and model-call assertions.
No production logic or test deadline changes.

Validation: 22 cases on the exact tested CI merge with Node 24.21; real writer
queue probe preserves the old updating recap until write settlement, then exposes
the published current recap. Owning type graph, typed lint, formatting, and P2
review pass. The initial CI log redacts the returned summary, so the change does
not claim that its historical failure was proved timing-only.

* test(discord): scope policy reload fixture activation

Limit the synthetic policy fixture to its Discord plugin so real reloads do
not cold-load unrelated runtime plugins. Require the reload to report
applied while preserving the active-turn, held lookup and revocation flow.

Exact CI merge-source proof passed locally in 15.3s. Root fixture types,
typed lint, formatting and independent review passed. No production code
or test timeout changed.
This commit is contained in:
Peter Steinberger 2026-09-23 08:23:57 -07:00 • committed by GitHub
parent 470c69ced0
commit 4e9e68140a
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
24 changed files with 698 additions and 300 deletions

View file

@ -1357,7 +1357,7 @@ extensions/zalouser/src/text-styles-inline.ts 1
extensions/zalouser/src/text-styles-ranges.ts 1
extensions/zalouser/src/text-styles-source.ts 1
extensions/zalouser/src/tool.ts 1
extensions/zalouser/src/zalo-js.ts 17
extensions/zalouser/src/zalo-js.ts 14
extensions/zalouser/src/zca-client.ts 2
packages/acp-core/src/runtime/errors.ts 6
packages/acp-core/src/runtime/session-identity.ts 2

View file

@ -556,6 +556,15 @@ const setupWizard: ChannelSetupWizard = {
`ChannelSetupWizard` also supports `textInputs`, `dmPolicy`, `allowFrom`, `groupAccess`, `prepare`, `finalize`, and more. See the Discord plugin's `src/setup-core.ts` for a full bundled example.
Before a durable setup effect, await `options.beforePersistentEffect?.()` to run
host preparation. Hosted channel wizards also supply
`options.assertPersistentEffectCurrent`, a synchronous check of the live wizard
owner. Carry that check through asynchronous credential preparation and invoke it
at the storage owner's final write admission. Detached QR login callbacks must
stop after their wizard is disposed, replaced, or completed; a successful earlier
preparation check does not keep that wizard alive. This optional lifetime check
does not replace the existing asynchronous preparation callback.
<AccordionGroup>
<Accordion title="Shared allowFrom prompts">
For DM allowlist prompts that only need the standard `note -> prompt -> parse -> merge -> patch` flow, prefer the shared setup helpers from `openclaw/plugin-sdk/setup`: `createPromptParsedAllowFromForAccount(...)` and `createTopLevelChannelParsedAllowFromPrompt(...)`.

View file

@ -2,8 +2,8 @@
import fs from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import type { OpenAsyncKeyedStoreOptions } from "openclaw/plugin-sdk/plugin-state-runtime";
import {
createPluginStateSyncKeyedStoreForTests,
createPluginStateKeyedStoreForTests,
resetPluginStateStoreForTests,
} from "openclaw/plugin-sdk/plugin-state-test-runtime";
@ -135,13 +135,13 @@ describe("zalouser doctor state migration", () => {
}),
);
const runtime = createPluginRuntimeMock();
runtime.state.openSyncKeyedStore = <T>(options: OpenKeyedStoreOptions) =>
createPluginStateSyncKeyedStoreForTests<T>("zalouser", {
runtime.state.openKeyedStore = <T>(options: OpenAsyncKeyedStoreOptions) =>
createPluginStateKeyedStoreForTests<T>("zalouser", {
...options,
env: options.env ?? env,
});
setZalouserRuntime(runtime);
clearStoredZaloCredentials(profile, env);
await clearStoredZaloCredentials(profile, env);
const context = createDoctorContext(env);
const params = {
config: {},

View file

@ -0,0 +1,197 @@
import { normalizeOptionalString } from "openclaw/plugin-sdk/string-coerce-runtime";
import {
captureZalouserCredentialsEnv,
clearStoredZaloCredentials,
loadStoredZaloCredentials,
refreshStoredZaloCredentials,
saveStoredZaloCredentials,
type StoredZaloCredentials,
} from "./session-state.js";
import type { API, Credentials } from "./zca-client.js";
export type ZaloCredentialPayload = Omit<
StoredZaloCredentials,
"profile" | "createdAt" | "lastUsedAt"
>;
function credentialSignature(credentials: ZaloCredentialPayload): string {
return JSON.stringify({
imei: credentials.imei,
cookie: canonicalCredentialCookie(credentials.cookie),
userAgent: credentials.userAgent,
language: credentials.language,
});
}
function stableCanonicalValue(value: unknown): unknown {
if (Array.isArray(value)) {
return value.map(stableCanonicalValue);
}
if (!value || typeof value !== "object") {
return value;
}
return Object.fromEntries(
Object.entries(value)
.toSorted(([left], [right]) => left.localeCompare(right))
.map(([key, entry]) => [key, stableCanonicalValue(entry)]),
);
}
function stableSignatureValue(value: unknown): string {
return JSON.stringify(stableCanonicalValue(value)) ?? "undefined";
}
function canonicalCookieArray(value: unknown[]): unknown[] {
return value
.map(stableCanonicalValue)
.toSorted((left, right) =>
stableSignatureValue(left).localeCompare(stableSignatureValue(right)),
);
}
function canonicalCredentialCookie(cookie: Credentials["cookie"]): unknown {
if (Array.isArray(cookie)) {
return canonicalCookieArray(cookie);
}
if (!cookie || typeof cookie !== "object") {
return cookie;
}
return Object.fromEntries(
Object.entries(cookie)
.toSorted(([left], [right]) => left.localeCompare(right))
.map(([key, entry]) => [
key,
key === "cookies" && Array.isArray(entry)
? canonicalCookieArray(entry)
: stableCanonicalValue(entry),
]),
);
}
export function snapshotApiCredentials(
api: API,
fallback?: Partial<ZaloCredentialPayload>,
): ZaloCredentialPayload {
const ctx = api.getContext();
const cookieJson = api.getCookie().toJSON();
const refreshedCookies =
Array.isArray(cookieJson?.cookies) && cookieJson.cookies.length > 0
? cookieJson.cookies
: fallback?.cookie;
const imei = normalizeOptionalString(ctx.imei) ?? normalizeOptionalString(fallback?.imei);
const userAgent =
normalizeOptionalString(ctx.userAgent) ?? normalizeOptionalString(fallback?.userAgent);
if (!imei || !refreshedCookies || !userAgent) {
throw new Error("Zalo API session did not expose refreshed credentials");
}
const language =
normalizeOptionalString(ctx.language) ?? normalizeOptionalString(fallback?.language);
return {
imei,
cookie: refreshedCookies,
userAgent,
...(language ? { language } : {}),
};
}
export class ZaloCredentialPersistence {
private readonly credentialSignaturesByProfile = new Map<string, string>();
private readonly credentialRefreshesByProfile = new Map<string, Promise<void>>();
private readonly credentialRevocationsByProfile = new Map<string, Promise<boolean>>();
constructor(private readonly isCurrentApi: (profile: string, api: API) => boolean) {}
pendingRevocation(profile: string): Promise<boolean> | undefined {
return this.credentialRevocationsByProfile.get(profile);
}
rememberCredentials(profile: string, credentials: ZaloCredentialPayload): void {
this.credentialSignaturesByProfile.set(profile, credentialSignature(credentials));
}
private async writeCredentials(
profile: string,
credentials: ZaloCredentialPayload,
assertCurrent: () => void,
): Promise<void> {
const env = captureZalouserCredentialsEnv();
await this.credentialRevocationsByProfile.get(profile);
assertCurrent();
const existing = await loadStoredZaloCredentials(profile, env).catch(() => null);
assertCurrent();
const now = new Date().toISOString();
const next: StoredZaloCredentials = {
profile,
...credentials,
createdAt: existing?.createdAt ?? now,
lastUsedAt: now,
};
const { profile: _profile, ...stored } = next;
await saveStoredZaloCredentials(profile, stored, env, assertCurrent);
assertCurrent();
this.credentialSignaturesByProfile.set(profile, credentialSignature(next));
}
async writeApiCredentials(
profile: string,
api: API,
assertCurrent: () => void,
fallback?: Partial<ZaloCredentialPayload>,
): Promise<void> {
await this.writeCredentials(profile, snapshotApiCredentials(api, fallback), assertCurrent);
}
async persistApiCredentialsIfChanged(profile: string, api: API): Promise<void> {
const previous = this.credentialRefreshesByProfile.get(profile) ?? Promise.resolve();
const refresh = previous.then(async () => {
try {
const isCurrent = () => this.isCurrentApi(profile, api);
if (!isCurrent()) {
return;
}
const credentials = snapshotApiCredentials(api);
const signature = credentialSignature(credentials);
if (this.credentialSignaturesByProfile.get(profile) === signature) {
return;
}
if ((await refreshStoredZaloCredentials(profile, credentials, isCurrent)) && isCurrent()) {
this.credentialSignaturesByProfile.set(profile, signature);
}
} catch {
// Do not fail an already-successful Zalo operation only because the
// best-effort session refresh could not be persisted.
}
});
this.credentialRefreshesByProfile.set(profile, refresh);
try {
await refresh;
} finally {
if (this.credentialRefreshesByProfile.get(profile) === refresh) {
this.credentialRefreshesByProfile.delete(profile);
}
}
}
async clearCredentials(profile: string, assertCurrent?: () => void): Promise<boolean> {
const env = captureZalouserCredentialsEnv();
const previous = this.credentialRevocationsByProfile.get(profile);
const pending = (async () => {
await previous;
try {
const cleared = await clearStoredZaloCredentials(profile, env, assertCurrent);
this.credentialSignaturesByProfile.delete(profile);
return cleared;
} catch {
return false;
}
})();
this.credentialRevocationsByProfile.set(profile, pending);
try {
return await pending;
} finally {
if (this.credentialRevocationsByProfile.get(profile) === pending) {
this.credentialRevocationsByProfile.delete(profile);
}
}
}
}

View file

@ -15,7 +15,6 @@ import { setZalouserRuntime } from "./runtime.js";
import {
clearStoredZaloCredentials,
loadStoredZaloCredentials,
loadStoredZaloCredentialsAsync,
refreshStoredZaloCredentials,
type StoredZaloCredentials,
} from "./session-state.js";
@ -305,13 +304,12 @@ export function createSendHarness(options: { mediaFixture?: boolean } = {}): Sen
userAgent: "synthetic-user-agent",
createdAt: "2026-01-01T00:00:00.000Z",
};
vi.mocked(loadStoredZaloCredentials).mockReturnValue(stored);
vi.mocked(loadStoredZaloCredentialsAsync).mockResolvedValue(stored);
vi.mocked(loadStoredZaloCredentials).mockResolvedValue(stored);
vi.mocked(refreshStoredZaloCredentials).mockImplementation(async (_profile, credentials) => ({
...stored,
...credentials,
}));
vi.mocked(clearStoredZaloCredentials).mockReturnValue(true);
vi.mocked(clearStoredZaloCredentials).mockResolvedValue(true);
if (options.mediaFixture !== false) {
vi.mocked(loadOutboundMediaFromUrl).mockResolvedValue({
buffer: Buffer.from("fixture"),

View file

@ -14,9 +14,8 @@ vi.mock("./session-state.js", async (importOriginal) => ({
...(await importOriginal<typeof import("./session-state.js")>()),
clearStoredZaloCredentials: vi.fn(),
loadStoredZaloCredentials: vi.fn(),
loadStoredZaloCredentialsAsync: vi.fn(),
refreshStoredZaloCredentials: vi.fn(),
saveStoredZaloCredentials: vi.fn(),
saveStoredZaloCredentials: vi.fn(async () => {}),
}));
vi.mock("openclaw/plugin-sdk/outbound-media", () => ({

View file

@ -14,9 +14,8 @@ vi.mock("./session-state.js", async (importOriginal) => ({
...(await importOriginal<typeof import("./session-state.js")>()),
clearStoredZaloCredentials: vi.fn(),
loadStoredZaloCredentials: vi.fn(),
loadStoredZaloCredentialsAsync: vi.fn(),
refreshStoredZaloCredentials: vi.fn(),
saveStoredZaloCredentials: vi.fn(),
saveStoredZaloCredentials: vi.fn(async () => {}),
}));
const mediaUrl = "https://media.example.com/document.txt";

View file

@ -1,7 +1,6 @@
import { createHash } from "node:crypto";
import os from "node:os";
import path from "node:path";
import type { PluginStateSyncKeyedStore } from "openclaw/plugin-sdk/plugin-state-runtime";
import { resolveStateDir } from "openclaw/plugin-sdk/state-paths";
import { normalizeLowercaseStringOrEmpty } from "openclaw/plugin-sdk/string-coerce-runtime";
import { getZalouserRuntime } from "./runtime.js";
@ -101,42 +100,17 @@ export function isZaloCredentialRevocation(
);
}
function openZalouserCredentialsStore(
export function captureZalouserCredentialsEnv(
env: NodeJS.ProcessEnv = process.env,
): PluginStateSyncKeyedStore<ZaloCredentialStateRecord> {
return getZalouserRuntime().state.openSyncKeyedStore<ZaloCredentialStateRecord>({
namespace: ZALOUSER_CREDENTIALS_NAMESPACE,
maxEntries: ZALOUSER_CREDENTIALS_MAX_ENTRIES,
overflowPolicy: "reject-new",
env,
});
): NodeJS.ProcessEnv {
return {
...env,
OPENCLAW_STATE_DIR: resolveStateDir(env),
OPENCLAW_SUPERVISOR_MODE: env.OPENCLAW_SUPERVISOR_MODE,
};
}
export function loadStoredZaloCredentials(
profile: string,
env: NodeJS.ProcessEnv = process.env,
): StoredZaloCredentials | null {
const normalizedProfile = normalizeZalouserCredentialProfile(profile);
const stored = openZalouserCredentialsStore(env).lookup(
zalouserCredentialStoreKey(normalizedProfile),
);
const parsed = normalizeStoredZaloCredentials(stored, normalizedProfile);
return parsed?.profile === normalizedProfile ? parsed : null;
}
export function saveStoredZaloCredentials(
profile: string,
credentials: Omit<StoredZaloCredentials, "profile">,
env: NodeJS.ProcessEnv = process.env,
): void {
const normalizedProfile = normalizeZalouserCredentialProfile(profile);
openZalouserCredentialsStore(env).register(zalouserCredentialStoreKey(normalizedProfile), {
profile: normalizedProfile,
...credentials,
});
}
function openAsyncZalouserCredentialsStore(env: NodeJS.ProcessEnv) {
function openZalouserCredentialsStore(env: NodeJS.ProcessEnv) {
return getZalouserRuntime().state.openKeyedStore<ZaloCredentialStateRecord>({
namespace: ZALOUSER_CREDENTIALS_NAMESPACE,
maxEntries: ZALOUSER_CREDENTIALS_MAX_ENTRIES,
@ -145,18 +119,32 @@ function openAsyncZalouserCredentialsStore(env: NodeJS.ProcessEnv) {
});
}
export async function loadStoredZaloCredentialsAsync(
export async function loadStoredZaloCredentials(
profile: string,
env: NodeJS.ProcessEnv = process.env,
): Promise<StoredZaloCredentials | null> {
const normalizedProfile = normalizeZalouserCredentialProfile(profile);
const store = openAsyncZalouserCredentialsStore(env);
const store = openZalouserCredentialsStore(env);
return normalizeStoredZaloCredentials(
await store.lookup(zalouserCredentialStoreKey(normalizedProfile)),
normalizedProfile,
);
}
export async function saveStoredZaloCredentials(
profile: string,
credentials: Omit<StoredZaloCredentials, "profile">,
env: NodeJS.ProcessEnv = process.env,
assertCurrent?: () => void,
): Promise<void> {
const normalizedProfile = normalizeZalouserCredentialProfile(profile);
await openZalouserCredentialsStore(env).register(
zalouserCredentialStoreKey(normalizedProfile),
{ profile: normalizedProfile, ...credentials },
{ assertCurrent },
);
}
export async function refreshStoredZaloCredentials(
profile: string,
credentials: Omit<StoredZaloCredentials, "profile" | "createdAt" | "lastUsedAt">,
@ -164,7 +152,7 @@ export async function refreshStoredZaloCredentials(
env: NodeJS.ProcessEnv = process.env,
): Promise<StoredZaloCredentials | null> {
const normalizedProfile = normalizeZalouserCredentialProfile(profile);
const store = openAsyncZalouserCredentialsStore(env);
const store = openZalouserCredentialsStore(env);
const key = zalouserCredentialStoreKey(normalizedProfile);
const now = new Date().toISOString();
const prepare = (
@ -213,23 +201,47 @@ export async function refreshStoredZaloCredentials(
return null;
}
export function clearStoredZaloCredentials(
export async function clearStoredZaloCredentials(
profile: string,
env: NodeJS.ProcessEnv = process.env,
): boolean {
assertCurrent?: () => void,
): Promise<boolean> {
const normalizedProfile = normalizeZalouserCredentialProfile(profile);
const store = openZalouserCredentialsStore(env);
const hadCredentials =
normalizeStoredZaloCredentials(
store.lookup(zalouserCredentialStoreKey(normalizedProfile)),
normalizedProfile,
) !== null;
// Keep a durable revocation marker so doctor cannot resurrect explicitly
// cleared credentials from an older profile file.
store.register(zalouserCredentialStoreKey(normalizedProfile), {
const opened = openZalouserCredentialsStore(env);
const store =
assertCurrent && opened.withCurrent ? opened.withCurrent({ assertCurrent }) : opened;
const key = zalouserCredentialStoreKey(normalizedProfile);
const revoked: ZaloCredentialRevocationRecord = {
kind: "revoked",
profile: normalizedProfile,
revokedAt: new Date().toISOString(),
});
return hadCredentials;
};
if (!store.observe || !store.compareAndApply || (assertCurrent && !opened.withCurrent)) {
// Released hosts before worker comparisons retain their atomic update contract.
if (!opened.update) {
throw new Error("Zalo credential logout requires atomic plugin-state updates");
}
let hadCredentials = false;
await opened.update(key, (current) => {
assertCurrent?.();
hadCredentials = normalizeStoredZaloCredentials(current, normalizedProfile) !== null;
return revoked;
});
return hadCredentials;
}
let observed = await store.observe(key);
for (;;) {
const hadCredentials =
normalizeStoredZaloCredentials(observed.value, normalizedProfile) !== null;
// Revocation and the returned presence flag describe the same committed row.
const result = await store.compareAndApply(key, observed.comparison, {
operation: "update",
action: "set",
value: revoked,
});
if (result.status !== "conflict") {
return hadCredentials;
}
observed = result.current;
}
}

View file

@ -195,7 +195,7 @@ describe("zalouser setup wizard", () => {
);
expect(beforePersistentEffect).toHaveBeenCalledTimes(2);
expect(logoutZaloProfileMock).toHaveBeenCalledWith("default");
expect(logoutZaloProfileMock).toHaveBeenCalledWith("default", { assertCurrent: undefined });
expect(startZaloQrLoginMock).not.toHaveBeenCalled();
expect(beforePersistentEffect.mock.invocationCallOrder[0]).toBeLessThan(
logoutZaloProfileMock.mock.invocationCallOrder[0]!,

View file

@ -328,6 +328,9 @@ export const zalouserSetupWizard: ChannelSetupWizard = {
...(options?.beforePersistentEffect
? { beforeCredentialPersistence: options.beforePersistentEffect }
: {}),
...(options?.assertPersistentEffectCurrent
? { assertCredentialPersistenceCurrent: options.assertPersistentEffectCurrent }
: {}),
});
if (start.qrDataUrl) {
const qrPath = await writeQrDataUrlToTempFile(start.qrDataUrl, account.profile);
@ -366,7 +369,9 @@ export const zalouserSetupWizard: ChannelSetupWizard = {
});
if (!keepSession) {
await options?.beforePersistentEffect?.();
await logoutZaloProfile(account.profile);
await logoutZaloProfile(account.profile, {
assertCurrent: options?.assertPersistentEffectCurrent,
});
await options?.beforePersistentEffect?.();
const start = await startZaloQrLogin({
profile: account.profile,
@ -375,6 +380,9 @@ export const zalouserSetupWizard: ChannelSetupWizard = {
...(options?.beforePersistentEffect
? { beforeCredentialPersistence: options.beforePersistentEffect }
: {}),
...(options?.assertPersistentEffectCurrent
? { assertCredentialPersistenceCurrent: options.assertPersistentEffectCurrent }
: {}),
});
if (start.qrDataUrl) {
const qrPath = await writeQrDataUrlToTempFile(start.qrDataUrl, account.profile);

View file

@ -1,5 +1,6 @@
// Zalouser tests cover zalo js.credentials plugin behavior.
import { access, mkdtemp, rm } from "node:fs/promises";
import { createRequire } from "node:module";
import os from "node:os";
import path from "node:path";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
@ -14,7 +15,14 @@ import {
createPluginStateSyncKeyedStoreForTests,
resetPluginStateStoreForTests,
} from "openclaw/plugin-sdk/plugin-state-test-runtime";
import { createPluginRuntimeMock } from "openclaw/plugin-sdk/plugin-test-runtime";
import {
createPluginRuntimeMock,
createPluginSetupWizardConfigure,
createTestWizardPrompter,
runSetupWizardConfigure,
} from "openclaw/plugin-sdk/plugin-test-runtime";
import { openNodeSqliteDatabase } from "openclaw/plugin-sdk/sqlite-runtime";
import { closeOpenClawStateDatabaseAsync } from "openclaw/plugin-sdk/sqlite-runtime-testing";
import { useAutoCleanupTempDirTracker, withEnvAsync } from "openclaw/plugin-sdk/test-env";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import type { API, LoginQRCallbackEvent } from "./zca-client.js";
@ -27,10 +35,13 @@ const ISO_TIMESTAMP_RE = /^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\.\d{3}Z$/u;
vi.mock("./zca-client.js", () => ({
createZalo: createZaloMock,
}));
vi.mock("./qr-temp-file.js", () => ({ writeQrDataUrlToTempFile: async () => undefined }));
import { zalouserSetupPlugin } from "./channel.setup.js";
import { getZalouserRuntime, setZalouserRuntime } from "./runtime.js";
import {
clearStoredZaloCredentials,
isZaloCredentialRevocation,
loadStoredZaloCredentials,
refreshStoredZaloCredentials,
resolveLegacyZalouserCredentialsPath,
@ -52,19 +63,19 @@ async function readStoredCredentials(
stateDir: string,
profile: string,
): Promise<StoredZaloCredentials> {
const stored = loadStoredZaloCredentials(profile, { OPENCLAW_STATE_DIR: stateDir });
const stored = await loadStoredZaloCredentials(profile, { OPENCLAW_STATE_DIR: stateDir });
if (!stored) {
throw new Error("Expected stored Zalo credentials");
}
return stored;
}
function seedStoredCredentials(
async function seedStoredCredentials(
stateDir: string,
profile: string,
credentials: Omit<StoredZaloCredentials, "profile">,
): void {
saveStoredZaloCredentials(profile, credentials, { OPENCLAW_STATE_DIR: stateDir });
): Promise<void> {
await saveStoredZaloCredentials(profile, credentials, { OPENCLAW_STATE_DIR: stateDir });
}
// Credential reads and writes leave the shared state database open under the temporary
@ -112,18 +123,34 @@ function createMockApi(params: {
describe("zalouser credential persistence", () => {
let beforeCredentialApply: (() => Promise<void>) | undefined;
let beforeCredentialRevocation: (() => Promise<void>) | undefined;
let beforeCredentialRegister: (() => Promise<void>) | undefined;
let afterCredentialRegister: (() => void) | undefined;
beforeEach(() => {
resetPluginStateStoreForTests();
const runtime = createPluginRuntimeMock();
beforeCredentialApply = undefined;
beforeCredentialRevocation = undefined;
beforeCredentialRegister = undefined;
afterCredentialRegister = undefined;
runtime.state.openKeyedStore = <T>(
options: OpenAsyncKeyedStoreOptions,
): PluginStateKeyedStore<T> => {
const store = createPluginStateKeyedStoreForTests<T>("zalouser", options);
return {
...store,
register: async (...args) => {
await beforeCredentialRegister?.();
await store.register(...args);
afterCredentialRegister?.();
},
compareAndApply: async (...args) => {
await beforeCredentialApply?.();
const intent = args[2];
if (intent.action === "set" && isZaloCredentialRevocation(intent.value)) {
await beforeCredentialRevocation?.();
} else {
await beforeCredentialApply?.();
}
return await store.compareAndApply(...args);
},
};
@ -154,7 +181,7 @@ describe("zalouser credential persistence", () => {
createdAt: "2026-04-01T00:00:00.000Z",
};
try {
saveStoredZaloCredentials(profile, stored, env);
await saveStoredZaloCredentials(profile, stored, env);
const refreshedCookie = [{ key: "zpsid", value: "refreshed", domain: "chat.zalo.me" }];
await refreshStoredZaloCredentials(
profile,
@ -162,8 +189,8 @@ describe("zalouser credential persistence", () => {
() => true,
env,
);
expect(loadStoredZaloCredentials(profile, env)?.cookie).toEqual(refreshedCookie);
clearStoredZaloCredentials(profile, env);
expect((await loadStoredZaloCredentials(profile, env))?.cookie).toEqual(refreshedCookie);
await clearStoredZaloCredentials(profile, env);
expect(
await refreshStoredZaloCredentials(
@ -173,14 +200,151 @@ describe("zalouser credential persistence", () => {
env,
),
).toBeNull();
expect(loadStoredZaloCredentials(profile, env)).toBeNull();
expect(await loadStoredZaloCredentials(profile, env)).toBeNull();
} finally {
await removeCredentialStateDir(stateDir);
}
},
);
it("persists the final API cookie jar after QR login", async () => {
it("awaits durable revocation before logout completes without opening a synchronous store", async () => {
const stateDir = tempDirs.make("openclaw-zalouser-credentials-");
const profile = "logout-settlement";
await seedStoredCredentials(stateDir, profile, {
imei: "device",
cookie: [{ key: "zpsid", value: "stored", domain: "chat.zalo.me" }],
userAgent: "agent",
createdAt: "2026-04-01T00:00:00.000Z",
});
const entered = createDeferred<void>();
const release = createDeferred<void>();
beforeCredentialRevocation = async () => {
entered.resolve();
await release.promise;
};
await withEnvAsync({ OPENCLAW_STATE_DIR: stateDir }, async () => {
const syncOpen = vi.spyOn(getZalouserRuntime().state, "openSyncKeyedStore");
const result = logoutZaloProfile(profile);
try {
await expect(
Promise.race([entered.promise.then(() => "pending"), result.then(() => "completed")]),
).resolves.toBe("pending");
expect(syncOpen).not.toHaveBeenCalled();
} finally {
release.resolve();
await result;
syncOpen.mockRestore();
}
await expect(result).resolves.toMatchObject({ cleared: true, loggedOut: true });
expect(await loadStoredZaloCredentials(profile)).toBeNull();
});
});
it("keeps restoration on its original state directory while waiting for logout", async () => {
const stateDir = tempDirs.make("openclaw-zalouser-credentials-");
const otherStateDir = tempDirs.make("openclaw-zalouser-other-credentials-");
const profile = "logout-state-root";
const stored = {
imei: "device",
cookie: [{ key: "zpsid", value: "stored", domain: "chat.zalo.me" }],
userAgent: "agent",
createdAt: "2026-04-01T00:00:00.000Z",
};
await seedStoredCredentials(stateDir, profile, stored);
await seedStoredCredentials(otherStateDir, profile, stored);
createZaloMock.mockResolvedValue({
login: async () =>
createMockApi({ imei: stored.imei, userAgent: stored.userAgent, cookies: stored.cookie }),
});
const entered = createDeferred<void>();
const release = createDeferred<void>();
beforeCredentialRevocation = async () => {
entered.resolve();
await release.promise;
};
await withEnvAsync({ OPENCLAW_STATE_DIR: stateDir }, async () => {
const logout = logoutZaloProfile(profile);
await entered.promise;
const restored = getZaloUserInfo(profile).then(
() => "restored",
() => "missing",
);
try {
await withEnvAsync({ OPENCLAW_STATE_DIR: otherStateDir }, async () => {
release.resolve();
await logout;
await expect(restored).resolves.toBe("missing");
});
expect(createZaloMock).not.toHaveBeenCalled();
} finally {
release.resolve();
await Promise.all([logout, restored]);
}
});
});
it("does not persist a registered setup login after its host closes during the credential write", async () => {
const stateDir = tempDirs.make("openclaw-zalouser-credentials-");
const profile = "setup-closed";
const entered = createDeferred<void>();
const release = createDeferred<void>();
let closed = false;
const assertPersistentEffectCurrent = () => {
if (closed) {
throw new Error("setup host closed");
}
};
beforeCredentialRegister = async () => {
entered.resolve();
await release.promise;
};
const api = createMockApi({
imei: "device",
userAgent: "agent",
cookies: [{ key: "zpsid", value: "new", domain: "chat.zalo.me" }],
});
createZaloMock.mockResolvedValueOnce({
loginQR: async (_options: unknown, callback?: (event: LoginQRCallbackEvent) => unknown) => {
callback?.({
type: LoginQRCallbackEventType.QRCodeGenerated,
data: { code: "synthetic-qr", image: "data:image/png;base64,abc123" },
actions: { saveToFile: vi.fn(async () => undefined), retry: vi.fn(), abort: vi.fn() },
});
return api;
},
});
await withEnvAsync({ OPENCLAW_STATE_DIR: stateDir }, async () => {
const setup = runSetupWizardConfigure({
configure: createPluginSetupWizardConfigure(zalouserSetupPlugin),
cfg: { channels: { zalouser: { profile } } },
prompter: createTestWizardPrompter({
confirm: vi.fn(async ({ message }) =>
["Login via QR code now?", "Did you scan and approve the QR on your phone?"].includes(
message,
),
),
}),
options: {
beforePersistentEffect: async () => assertPersistentEffectCurrent(),
assertPersistentEffectCurrent,
},
});
try {
await expect(
Promise.race([entered.promise.then(() => "pending"), setup.then(() => "completed")]),
).resolves.toBe("pending");
closed = true;
} finally {
release.resolve();
await setup;
}
expect(await loadStoredZaloCredentials(profile)).toBeNull();
});
});
it("persists the final API cookie jar and logs out without parent SQLite work", async () => {
const registered = createDeferred<void>();
afterCredentialRegister = () => registered.resolve();
const stateDir = await mkdtemp(path.join(os.tmpdir(), "openclaw-zalouser-credentials-"));
const profile = "qr-refresh";
const callbackCookie = [{ key: "zpsid", value: "callback", domain: "chat.zalo.me" }];
@ -219,20 +383,58 @@ describe("zalouser credential persistence", () => {
},
});
const native = createRequire(import.meta.url)("node:sqlite") as typeof import("node:sqlite");
const Database = native.DatabaseSync;
const sql = [
vi.spyOn(native, "DatabaseSync"),
...(["close", "prepare", "exec"] as const).map((method) =>
vi.spyOn(Database.prototype, method),
),
...(["get", "all", "run", "iterate"] as const).map((method) =>
vi.spyOn(native.StatementSync.prototype, method),
),
];
try {
// Calibrate every observer and finish runtime capability checks before the cold flow.
const calibration = openNodeSqliteDatabase(":memory:");
calibration.exec("CREATE TABLE calibration (value INTEGER)");
calibration.prepare("INSERT INTO calibration VALUES (?)").run(1);
const query = calibration.prepare("SELECT value FROM calibration");
query.get();
query.all();
Array.from(query.iterate());
calibration.close();
for (const call of sql) {
expect(call).toHaveBeenCalled();
call.mockClear();
}
await expect(access(path.join(stateDir, "state", "openclaw.sqlite"))).rejects.toMatchObject({
code: "ENOENT",
});
await withEnvAsync({ OPENCLAW_STATE_DIR: stateDir }, async () => {
await startZaloQrLogin({ profile, timeoutMs: 1000 });
await registered.promise;
const loginResult = await waitForZaloQrLogin({ profile, timeoutMs: 1000 });
expect(loginResult.connected).toBe(true);
expect(loginResult.connected, loginResult.message).toBe(true);
const stored = await readStoredCredentials(stateDir, profile);
expect(stored.imei).toBe("api-imei");
expect(stored.userAgent).toBe("api-user-agent");
expect(stored.language).toBe("vi");
expect(stored.cookie).toEqual(refreshedCookie);
await expect(logoutZaloProfile(profile)).resolves.toMatchObject({
cleared: true,
loggedOut: true,
});
expect(await loadStoredZaloCredentials(profile)).toBeNull();
await closeOpenClawStateDatabaseAsync();
for (const call of sql) {
expect(call).not.toHaveBeenCalled();
}
});
} finally {
sql.forEach((call) => call.mockRestore());
await removeCredentialStateDir(stateDir);
}
});
@ -279,7 +481,7 @@ describe("zalouser credential persistence", () => {
expect(`${started.message} ${waited.message}`).toContain(guardError.message);
expect(beforeCredentialPersistence).toHaveBeenCalledTimes(1);
expect(loadStoredZaloCredentials(profile)).toBeNull();
expect(await loadStoredZaloCredentials(profile)).toBeNull();
});
} finally {
await removeCredentialStateDir(stateDir);
@ -321,7 +523,7 @@ describe("zalouser credential persistence", () => {
const profile = "restore-refresh";
const storedCookie = [{ key: "zpsid", value: "stored", domain: "chat.zalo.me" }];
const refreshedCookie = [{ key: "zpsid", value: "refreshed", domain: "chat.zalo.me" }];
seedStoredCredentials(stateDir, profile, {
await seedStoredCredentials(stateDir, profile, {
imei: "stored-imei",
cookie: storedCookie,
userAgent: "stored-user-agent",
@ -363,7 +565,7 @@ describe("zalouser credential persistence", () => {
const storedCookie = [{ key: "zpsid", value: "stored", domain: "chat.zalo.me" }];
const loginCookie = [{ key: "zpsid", value: "login", domain: "chat.zalo.me" }];
const refreshedCookie = [{ key: "zpsid", value: "refreshed", domain: "chat.zalo.me" }];
seedStoredCredentials(stateDir, profile, {
await seedStoredCredentials(stateDir, profile, {
imei: "stored-imei",
cookie: storedCookie,
userAgent: "stored-user-agent",
@ -405,7 +607,7 @@ describe("zalouser credential persistence", () => {
const refreshedCookie: unknown[] = [
{ key: "zpsid", value: "api-refreshed", domain: "chat.zalo.me" },
];
seedStoredCredentials(stateDir, profile, {
await seedStoredCredentials(stateDir, profile, {
imei: "stored-imei",
cookie: storedCookie,
userAgent: "stored-user-agent",
@ -460,7 +662,7 @@ describe("zalouser credential persistence", () => {
const profile = `settlement-${outcome}`;
const storedCookie = [{ key: "zpsid", value: "stored", domain: "chat.zalo.me" }];
const refreshedCookie = [{ key: "zpsid", value: "refreshed", domain: "chat.zalo.me" }];
seedStoredCredentials(stateDir, profile, {
await seedStoredCredentials(stateDir, profile, {
imei: "device",
cookie: storedCookie,
userAgent: "agent",
@ -519,14 +721,14 @@ describe("zalouser credential persistence", () => {
await Promise.all([result, secondResult]);
}
await expect(result).resolves.toEqual([]);
const stored = loadStoredZaloCredentials(profile);
const stored = await loadStoredZaloCredentials(profile);
expect(stored?.cookie ?? null).toEqual(
outcome === "logout" ? null : outcome === "failure" ? storedCookie : newestCookie,
);
if (outcome === "failure") {
beforeCredentialApply = undefined;
await listZaloFriends(profile);
expect(loadStoredZaloCredentials(profile)?.cookie).toEqual(refreshedCookie);
expect((await loadStoredZaloCredentials(profile))?.cookie).toEqual(refreshedCookie);
}
syncOpen.mockRestore();
});
@ -541,7 +743,7 @@ describe("zalouser credential persistence", () => {
const stateDir = tempDirs.make("openclaw-zalouser-credentials-");
const profile = "restore-logout";
const cookie = [{ key: "zpsid", value: "stored", domain: "chat.zalo.me" }];
seedStoredCredentials(stateDir, profile, {
await seedStoredCredentials(stateDir, profile, {
imei: "device",
cookie,
userAgent: "agent",
@ -578,13 +780,13 @@ describe("zalouser credential persistence", () => {
Promise.race([entered.promise.then(() => "pending"), restored.then(() => "completed")]),
).resolves.toBe("pending");
await logoutZaloProfile(profile);
seedStoredCredentials(stateDir, profile, {
beforeCredentialApply = undefined;
await seedStoredCredentials(stateDir, profile, {
imei: "replacement",
cookie: replacementCookie,
userAgent: "agent",
createdAt: "2026-04-02T00:00:00.000Z",
});
beforeCredentialApply = undefined;
await expect(getZaloUserInfo(profile.toUpperCase())).resolves.toMatchObject({
userId: "user-1",
});
@ -594,7 +796,7 @@ describe("zalouser credential persistence", () => {
}
await expect(restored).resolves.toBe(false);
await expect(listZaloFriends(profile)).resolves.toEqual([]);
expect(loadStoredZaloCredentials(profile)?.cookie).toEqual(replacementCookie);
expect((await loadStoredZaloCredentials(profile))?.cookie).toEqual(replacementCookie);
expect(originalFriends).not.toHaveBeenCalled();
expect(replacementFriends).toHaveBeenCalledOnce();
});
@ -612,7 +814,7 @@ describe("zalouser credential persistence", () => {
{ key: "zpw", value: "same-secondary", domain: "chat.zalo.me" },
];
const cookieB = [...cookieA].toReversed();
seedStoredCredentials(stateDir, profile, {
await seedStoredCredentials(stateDir, profile, {
imei: "stored-imei",
cookie: cookieA,
userAgent: "stored-user-agent",
@ -686,7 +888,7 @@ describe("zalouser credential persistence", () => {
it("writes plugin-state SQLite without recreating the retired credential blob", async () => {
const stateDir = await mkdtemp(path.join(os.tmpdir(), "openclaw-zalouser-credentials-"));
const profile = "sqlite-only";
seedStoredCredentials(stateDir, profile, {
await seedStoredCredentials(stateDir, profile, {
imei: "api-imei",
userAgent: "api-user-agent",
cookie: [{ key: "zpsid", value: "sqlite", domain: "chat.zalo.me" }],

View file

@ -6,7 +6,6 @@ import path from "node:path";
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
import {
createPluginStateKeyedStoreForTests,
createPluginStateSyncKeyedStoreForTests,
resetPluginStateStoreForTests,
} from "openclaw/plugin-sdk/plugin-state-test-runtime";
import { createPluginRuntimeMock } from "openclaw/plugin-sdk/plugin-test-runtime";
@ -16,7 +15,7 @@ import { WebSocketServer } from "ws";
import { withZalouserIngressTestQueue } from "./ingress.test-support.js";
import { monitorZalouserProvider } from "./monitor.js";
import { setZalouserRuntime } from "./runtime.js";
import { loadStoredZaloCredentialsAsync, saveStoredZaloCredentials } from "./session-state.js";
import { loadStoredZaloCredentials, saveStoredZaloCredentials } from "./session-state.js";
import { createDefaultResolvedZalouserAccount, createZalouserRuntimeEnv } from "./test-helpers.js";
import type { API } from "./zca-client.js";
@ -44,13 +43,13 @@ function sessionApi(
}
async function seedSession() {
saveStoredZaloCredentials("default", {
await saveStoredZaloCredentials("default", {
imei: "fixture",
userAgent: "openclaw-test",
cookie: [],
createdAt: new Date().toISOString(),
});
expect(await loadStoredZaloCredentialsAsync("default")).not.toBeNull();
expect(await loadStoredZaloCredentials("default")).not.toBeNull();
}
beforeEach(() => {
@ -58,8 +57,6 @@ beforeEach(() => {
const runtime = createPluginRuntimeMock();
runtime.state.openKeyedStore = (options) =>
createPluginStateKeyedStoreForTests("zalouser", options);
runtime.state.openSyncKeyedStore = (options) =>
createPluginStateSyncKeyedStoreForTests("zalouser", options);
setZalouserRuntime(runtime);
createZaloMock.mockReset();
});

View file

@ -19,18 +19,20 @@ import {
normalizeOptionalStringifiedId,
} from "openclaw/plugin-sdk/string-coerce-runtime";
import { sleep } from "openclaw/plugin-sdk/text-utility-runtime";
import {
ZaloCredentialPersistence,
snapshotApiCredentials,
type ZaloCredentialPayload,
} from "./credential-persistence.js";
import { normalizeZaloReactionIcon } from "./reaction.js";
import { sendZaloTextWithApi } from "./send-api.js";
import { withZaloSendContext } from "./send-context.js";
import { createZalouserSendReceipt } from "./send-receipt.js";
import {
clearStoredZaloCredentials,
captureZalouserCredentialsEnv,
loadStoredZaloCredentials,
loadStoredZaloCredentialsAsync,
normalizeZalouserCredentialProfile as normalizeProfile,
refreshStoredZaloCredentials,
saveStoredZaloCredentials,
type StoredZaloCredentials,
} from "./session-state.js";
import type {
ZaloAuthStatus,
@ -47,7 +49,6 @@ import type {
} from "./types.js";
import {
type API,
type Credentials,
type GroupInfo,
type LoginQRCallbackEvent,
type Message,
@ -71,18 +72,19 @@ const MAX_SAFE_ZALO_TIMESTAMP_SECONDS = Number.MAX_SAFE_INTEGER / 1000;
const apiByProfile = new Map<string, API>();
const apiInitByProfile = new Map<string, Promise<API>>();
const credentialSignaturesByProfile = new Map<string, string>();
const credentialRefreshesByProfile = new Map<string, Promise<void>>();
const credentials = new ZaloCredentialPersistence(
(profile, api) => apiByProfile.get(profile) === api,
);
type CredentialPersistenceMode = "persist" | "read-only";
type CredentialPersistenceOptions = { credentialPersistence?: CredentialPersistenceMode };
type ZaloCredentialPayload = Omit<StoredZaloCredentials, "profile" | "createdAt" | "lastUsedAt">;
type ActiveZaloQrLogin = {
id: string;
profile: string;
startedAt: number;
beforeCredentialPersistence?: () => Promise<void>;
assertCredentialPersistenceCurrent?: () => void;
qrDataUrl?: string;
connected: boolean;
error?: string;
@ -371,170 +373,17 @@ function mapGroup(groupId: string, group: GroupInfo & Record<string, unknown>):
};
}
function readCredentials(profile: string): StoredZaloCredentials | null {
try {
const credentials = loadStoredZaloCredentials(profile);
if (!credentials) {
return null;
}
credentialSignaturesByProfile.set(profile, credentialSignature(credentials));
return credentials;
} catch {
return null;
}
}
function credentialSignature(credentials: ZaloCredentialPayload): string {
return JSON.stringify({
imei: credentials.imei,
cookie: canonicalCredentialCookie(credentials.cookie),
userAgent: credentials.userAgent,
language: credentials.language,
});
}
function stableCanonicalValue(value: unknown): unknown {
if (Array.isArray(value)) {
return value.map(stableCanonicalValue);
}
if (!value || typeof value !== "object") {
return value;
}
return Object.fromEntries(
Object.entries(value as Record<string, unknown>)
.toSorted(([left], [right]) => left.localeCompare(right))
.map(([key, entry]) => [key, stableCanonicalValue(entry)]),
);
}
function stableSignatureValue(value: unknown): string {
return JSON.stringify(stableCanonicalValue(value)) ?? "undefined";
}
function canonicalCookieArray(value: unknown[]): unknown[] {
return value
.map(stableCanonicalValue)
.toSorted((left, right) =>
stableSignatureValue(left).localeCompare(stableSignatureValue(right)),
);
}
function canonicalCredentialCookie(cookie: Credentials["cookie"]): unknown {
if (Array.isArray(cookie)) {
return canonicalCookieArray(cookie);
}
if (!cookie || typeof cookie !== "object") {
return cookie;
}
return Object.fromEntries(
Object.entries(cookie as Record<string, unknown>)
.toSorted(([left], [right]) => left.localeCompare(right))
.map(([key, entry]) => [
key,
key === "cookies" && Array.isArray(entry)
? canonicalCookieArray(entry)
: stableCanonicalValue(entry),
]),
);
}
function writeCredentials(profile: string, credentials: ZaloCredentialPayload): void {
const existing = readCredentials(profile);
const now = new Date().toISOString();
const next: StoredZaloCredentials = {
profile,
...credentials,
createdAt: existing?.createdAt ?? now,
lastUsedAt: now,
};
const { profile: _profile, ...stored } = next;
saveStoredZaloCredentials(profile, stored);
credentialSignaturesByProfile.set(profile, credentialSignature(next));
}
function snapshotApiCredentials(
api: API,
fallback?: Partial<ZaloCredentialPayload>,
): ZaloCredentialPayload {
const ctx = api.getContext();
const cookieJson = api.getCookie().toJSON();
const refreshedCookies =
Array.isArray(cookieJson?.cookies) && cookieJson.cookies.length > 0
? cookieJson.cookies
: fallback?.cookie;
const imei = normalizeOptionalString(ctx.imei) ?? normalizeOptionalString(fallback?.imei);
const userAgent =
normalizeOptionalString(ctx.userAgent) ?? normalizeOptionalString(fallback?.userAgent);
if (!imei || !refreshedCookies || !userAgent) {
throw new Error("Zalo API session did not expose refreshed credentials");
}
const language =
normalizeOptionalString(ctx.language) ?? normalizeOptionalString(fallback?.language);
return {
imei,
cookie: refreshedCookies as Credentials["cookie"],
userAgent,
...(language ? { language } : {}),
};
}
function writeApiCredentials(
profile: string,
api: API,
fallback?: Partial<ZaloCredentialPayload>,
): void {
writeCredentials(profile, snapshotApiCredentials(api, fallback));
}
async function persistApiCredentialsIfChanged(profile: string, api: API): Promise<void> {
const previous = credentialRefreshesByProfile.get(profile) ?? Promise.resolve();
const refresh = previous.then(async () => {
try {
const isCurrent = () => apiByProfile.get(profile) === api;
if (!isCurrent()) {
return;
}
const credentials = snapshotApiCredentials(api);
const signature = credentialSignature(credentials);
if (credentialSignaturesByProfile.get(profile) === signature) {
return;
}
if ((await refreshStoredZaloCredentials(profile, credentials, isCurrent)) && isCurrent()) {
credentialSignaturesByProfile.set(profile, signature);
}
} catch {
// Do not fail an already-successful Zalo operation only because the
// best-effort session refresh could not be persisted.
}
});
credentialRefreshesByProfile.set(profile, refresh);
try {
await refresh;
} finally {
if (credentialRefreshesByProfile.get(profile) === refresh) {
credentialRefreshesByProfile.delete(profile);
}
}
}
function clearCredentials(profile: string): boolean {
try {
if (clearStoredZaloCredentials(profile)) {
credentialSignaturesByProfile.delete(profile);
return true;
}
} catch {
// ignore
}
return false;
}
async function ensureApi(
profileInput?: string | null,
timeoutMs = API_LOGIN_TIMEOUT_MS,
credentialPersistence: CredentialPersistenceMode = "persist",
): Promise<API> {
const profile = normalizeProfile(profileInput);
const env = captureZalouserCredentialsEnv();
const pendingRevocation = credentials.pendingRevocation(profile);
if (pendingRevocation) {
await pendingRevocation;
}
const cached = apiByProfile.get(profile);
if (cached) {
return cached;
@ -547,7 +396,7 @@ async function ensureApi(
const initPromise: Promise<API> = (async () => {
const isCurrent = () => apiInitByProfile.get(profile) === initPromise;
const stored = await loadStoredZaloCredentialsAsync(profile).catch(() => null);
const stored = await loadStoredZaloCredentials(profile, env).catch(() => null);
if (!stored || !isCurrent()) {
throw new Error(`No saved Zalo session for profile "${profile}"`);
}
@ -571,12 +420,13 @@ async function ensureApi(
profile,
snapshotApiCredentials(api, stored),
isCurrent,
env,
)
: stored;
if (!isCurrent() || !persisted) {
throw new Error(`Zalo session restore was superseded for profile "${profile}"`);
}
credentialSignaturesByProfile.set(profile, credentialSignature(persisted));
credentials.rememberCredentials(profile, persisted);
apiByProfile.set(profile, api);
return api;
})();
@ -616,7 +466,7 @@ async function withZaloApi<T>(
? await withZaloSendContext(options.handoff, () => operation(api))
: await operation(api);
if (credentialPersistence === "persist" && (options.shouldPersist?.(result) ?? true)) {
await persistApiCredentialsIfChanged(profile, api);
await credentials.persistApiCredentialsIfChanged(profile, api);
}
return result;
}
@ -821,7 +671,7 @@ export async function checkZaloAuthenticated(
options?: CredentialPersistenceOptions,
): Promise<boolean> {
const profile = normalizeProfile(profileInput);
if (!(await loadStoredZaloCredentialsAsync(profile).catch(() => null))) {
if (!(await loadStoredZaloCredentials(profile).catch(() => null))) {
return false;
}
try {
@ -1238,6 +1088,7 @@ export async function startZaloQrLogin(params: {
force?: boolean;
timeoutMs?: number;
beforeCredentialPersistence?: () => Promise<void>;
assertCredentialPersistenceCurrent?: () => void;
}): Promise<{ qrDataUrl?: string; message: string }> {
const profile = normalizeProfile(params.profile);
@ -1249,15 +1100,18 @@ export async function startZaloQrLogin(params: {
};
}
params.assertCredentialPersistenceCurrent?.();
if (params.force) {
await logoutZaloProfile(profile);
await logoutZaloProfile(profile, { assertCurrent: params.assertCredentialPersistenceCurrent });
}
let existing = activeQrLogins.get(profile);
if (
existing &&
params.beforeCredentialPersistence &&
existing.beforeCredentialPersistence !== params.beforeCredentialPersistence
((params.beforeCredentialPersistence &&
existing.beforeCredentialPersistence !== params.beforeCredentialPersistence) ||
(params.assertCredentialPersistenceCurrent &&
existing.assertCredentialPersistenceCurrent !== params.assertCredentialPersistenceCurrent))
) {
// A QR flow may outlive its setup turn. Never let a new setup owner adopt
// another owner's pending login and bypass its persistence revalidation.
@ -1283,6 +1137,9 @@ export async function startZaloQrLogin(params: {
...(params.beforeCredentialPersistence
? { beforeCredentialPersistence: params.beforeCredentialPersistence }
: {}),
...(params.assertCredentialPersistenceCurrent
? { assertCredentialPersistenceCurrent: params.assertCredentialPersistenceCurrent }
: {}),
connected: false,
waitPromise: Promise.resolve(),
};
@ -1357,12 +1214,22 @@ export async function startZaloQrLogin(params: {
};
}
const assertCurrent = () => {
login.assertCredentialPersistenceCurrent?.();
if (activeQrLogins.get(profile)?.id !== login.id) {
throw new Error("Zalo QR login was superseded before credential persistence");
}
};
assertCurrent();
await login.beforeCredentialPersistence?.();
const owned = activeQrLogins.get(profile);
if (!owned || owned.id !== login.id) {
return;
}
writeApiCredentials(profile, api, capturedCredentials ?? undefined);
assertCurrent();
await credentials.writeApiCredentials(
profile,
api,
assertCurrent,
capturedCredentials ?? undefined,
);
assertCurrent();
invalidateApi(profile);
apiByProfile.set(profile, api);
current.connected = true;
@ -1463,12 +1330,16 @@ export async function waitForZaloQrLogin(params: {
};
}
export async function logoutZaloProfile(profileInput?: string | null): Promise<{
export async function logoutZaloProfile(
profileInput?: string | null,
options?: { assertCurrent?: () => void },
): Promise<{
cleared: boolean;
loggedOut: boolean;
message: string;
}> {
const profile = normalizeProfile(profileInput);
options?.assertCurrent?.();
resetQrLogin(profile);
clearCachedGroupContext(profile);
@ -1483,7 +1354,7 @@ export async function logoutZaloProfile(profileInput?: string | null): Promise<{
}
invalidateApi(profile);
const cleared = clearCredentials(profile);
const cleared = await credentials.clearCredentials(profile, options?.assertCurrent);
return {
cleared,

View file

@ -2,13 +2,9 @@
import { mkdtemp, realpath, rm, writeFile } from "node:fs/promises";
import os from "node:os";
import path from "node:path";
import type {
OpenAsyncKeyedStoreOptions,
OpenKeyedStoreOptions,
} from "openclaw/plugin-sdk/plugin-state-runtime";
import type { OpenAsyncKeyedStoreOptions } from "openclaw/plugin-sdk/plugin-state-runtime";
import {
createPluginStateKeyedStoreForTests,
createPluginStateSyncKeyedStoreForTests,
resetPluginStateStoreForTests,
} from "openclaw/plugin-sdk/plugin-state-test-runtime";
import { createPluginRuntimeMock } from "openclaw/plugin-sdk/plugin-test-runtime";
@ -63,7 +59,7 @@ async function withStoredSession<T>(params: {
run: () => Promise<T>;
}): Promise<T> {
const stateDir = await mkdtemp(path.join(os.tmpdir(), "openclaw-zalouser-message-"));
saveStoredZaloCredentials(
await saveStoredZaloCredentials(
params.profile,
{
imei: "test-imei",
@ -112,8 +108,6 @@ beforeEach(() => {
const runtime = createPluginRuntimeMock();
runtime.state.openKeyedStore = <T>(options: OpenAsyncKeyedStoreOptions) =>
createPluginStateKeyedStoreForTests<T>("zalouser", options);
runtime.state.openSyncKeyedStore = <T>(options: OpenKeyedStoreOptions) =>
createPluginStateSyncKeyedStoreForTests<T>("zalouser", options);
setZalouserRuntime(runtime);
createZaloMock.mockReset();
});

View file

@ -309,6 +309,8 @@ export type SetupChannelsOptions = {
allowSignalInstall?: boolean;
/** Revalidate host authority immediately before an installer or other durable effect. */
beforePersistentEffect?: () => Promise<void>;
/** Pure live setup-owner check after asynchronous preparation and at the final write grant. */
assertPersistentEffectCurrent?: () => void;
onSelection?: (selection: ChannelId[]) => void;
onPostWriteHook?: (hook: ChannelOnboardingPostWriteHook) => void;
accountIds?: Partial<Record<ChannelId, string>>;

View file

@ -112,6 +112,7 @@ type ChannelsAddWizardFlowParams = {
prompter: WizardPrompter;
initialChannel?: ChannelChoice;
beforePersistentEffect?: () => Promise<void>;
assertPersistentEffectCurrent?: () => void;
/**
* The controlling client completes device linking itself after config is
* written (e.g. the Control UI renders the WhatsApp QR via web.login.*), so
@ -150,6 +151,9 @@ export async function runChannelsAddWizardFlow(params: ChannelsAddWizardFlowPara
...(params.beforePersistentEffect
? { beforePersistentEffect: params.beforePersistentEffect }
: {}),
...(params.assertPersistentEffectCurrent
? { assertPersistentEffectCurrent: params.assertPersistentEffectCurrent }
: {}),
...(params.deferDeviceLinkToClient ? { deferDeviceLinkToClient: true } : {}),
onPostWriteHook: (hook) => channelSetup.onPostWriteHook(hook),
promptAccountIds: true,
@ -321,6 +325,7 @@ export async function runChannelsSetupWizard(
onConfigured?: (accounts: Array<{ channel: string; accountId: string }>) => void;
/** Revalidate/lock cancellation immediately before durable effects. */
beforePersistentEffect?: () => Promise<void>;
assertPersistentEffectCurrent?: () => void;
},
runtime: RuntimeEnv,
prompter: WizardPrompter,
@ -348,5 +353,8 @@ export async function runChannelsSetupWizard(
deferDeviceLinkToClient: true,
...(opts.onConfigured ? { onConfigured: opts.onConfigured } : {}),
...(opts.beforePersistentEffect ? { beforePersistentEffect: opts.beforePersistentEffect } : {}),
...(opts.assertPersistentEffectCurrent
? { assertPersistentEffectCurrent: opts.assertPersistentEffectCurrent }
: {}),
});
}

View file

@ -43,6 +43,7 @@ export type ChannelSetupWizardRunner = (
channel?: string;
onConfigured?: (accounts: Array<{ channel: string; accountId: string }>) => void;
beforePersistentEffect?: () => Promise<void>;
assertPersistentEffectCurrent?: () => void;
},
runtime: RuntimeEnv,
prompter: WizardPrompter,
@ -125,6 +126,8 @@ export const wizardHandlers: GatewayRequestHandlers = {
// Durable effects (plugin installs, config commit) must finish
// even if the client cancels mid-write.
beforePersistentEffect: async () => wizardSession.lockCancellation(),
assertPersistentEffectCurrent: () =>
wizardSession.assertPersistentEffectCurrent(),
},
runtime,
prompter,

View file

@ -328,15 +328,19 @@ describe("Activity recap lifecycle with the canonical session store", () => {
totalMessages: 1,
});
const published = createDeferred<ReturnType<typeof view>>();
changed.mockImplementationOnce(() => published.resolve(view()));
completion.resolve(result("Completed the first turn."));
await vi.waitFor(async () => {
expect(await describeSession()).toMatchObject({
session: {
key: target.key,
sessionId: scope.sessionId,
activitySummary: { state: "current", text: "Completed the first turn." },
},
});
expect(await published.promise).toMatchObject({
state: "current",
text: "Completed the first turn.",
});
expect(await describeSession()).toMatchObject({
session: {
key: target.key,
sessionId: scope.sessionId,
activitySummary: { state: "current", text: "Completed the first turn." },
},
});
expect(read()?.activitySummary).toMatchObject({
...latestWatermark,

View file

@ -51,6 +51,7 @@ export type ChatWizardHostDependencies = {
channel: string,
prompter: WizardPrompter,
beforePersistentApply: (runtime: RuntimeEnv) => Promise<void>,
assertPersistentEffectCurrent?: () => void,
) => Promise<void | HostedSetupCompletion>;
runSkillsSetupWizard?: (
prompter: WizardPrompter,
@ -358,12 +359,24 @@ export class ChatWizardHost {
kind: "channel",
label: channel,
autoSelectChannel: channel,
run: async (prompter) =>
run
? await run(channel, prompter, this.options.beforePersistentApply)
run: async (prompter, assertPersistentEffectCurrent) => {
const beforePersistentApply = async (runtime: RuntimeEnv) => {
assertPersistentEffectCurrent();
await this.options.beforePersistentApply(runtime);
assertPersistentEffectCurrent();
};
return run
? await run(channel, prompter, beforePersistentApply, assertPersistentEffectCurrent)
: await (
await loadHostedRuntime()
).runHostedChannelSetup(channel, prompter, this.options.beforePersistentApply),
).runHostedChannelSetup(
channel,
prompter,
beforePersistentApply,
undefined,
assertPersistentEffectCurrent,
);
},
});
}
@ -442,7 +455,10 @@ export class ChatWizardHost {
label: string;
autoSelectChannel?: string;
memoryImportProviders?: MemoryImportProviderOutcome[];
run: (prompter: WizardPrompter) => Promise<HostedWizardRunResult>;
run: (
prompter: WizardPrompter,
assertPersistentEffectCurrent: () => void,
) => Promise<HostedWizardRunResult>;
}): Promise<ChatWizardResult> {
const completion: ActiveWizardBridge["completion"] = {
status: "applied",
@ -450,8 +466,16 @@ export class ChatWizardHost {
? { memoryImportProviders: params.memoryImportProviders }
: {}),
};
const session = new WizardSession(async (prompter) => {
const result = await params.run(prompter);
const session = new WizardSession(async (prompter, _signal, owner) => {
// Publish the bridge before a setup callback can use its retained authority.
await Promise.resolve();
const assertPersistentEffectCurrent = () => {
owner.assertPersistentEffectCurrent();
if (this.bridge?.session !== owner) {
throw new Error("Setup session is no longer active");
}
};
const result = await params.run(prompter, assertPersistentEffectCurrent);
if (typeof result === "string") {
completion.status = result;
} else if (result) {

View file

@ -3,6 +3,7 @@ import { afterEach, describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../test/helpers/promise.js";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import { applyAccountNameToChannelSection } from "../channels/plugins/setup-helpers.js";
import type { SetupChannelsOptions } from "../channels/plugins/setup-wizard-types.js";
import type { ChannelPlugin } from "../channels/plugins/types.public.js";
import { committedConfigFiles as hostedConfigFiles } from "../commands/committed-config.test-support.js";
import { withCommandPluginMetadata } from "../commands/config-validation.js";
@ -43,6 +44,61 @@ import {
import { ChatWizardHost } from "./chat-wizard-host.js";
describe("SystemAgentChatEngine runtime", () => {
it.each(["disposed", "completed"] as const)(
"retires retained channel write authority after hosted setup is %s",
async (ending) => {
let retainedOptions: SetupChannelsOptions | undefined;
const prepare = vi.fn(async () => {});
mocks.readSetupConfigFileSnapshot.mockResolvedValue({
exists: true,
valid: true,
hash: "channel-lifetime-hash",
config: {},
sourceConfig: {},
});
mocks.writeWizardConfigFile.mockImplementation(async (config: OpenClawConfig) =>
hostedConfigFiles.write(config),
);
mocks.setupChannels.mockImplementation(
async (
config: OpenClawConfig,
_runtime: unknown,
prompter: WizardPrompter,
options: SetupChannelsOptions,
) => {
retainedOptions = options;
await options.beforePersistentEffect?.();
await prompter.text({ message: "Waiting for channel login" });
return config;
},
);
const host = new ChatWizardHost({
surface: "gateway",
beforePersistentApply: prepare,
dependencies: { appendAuditEntry: vi.fn(async () => "state/openclaw.sqlite") },
});
try {
expect((await host.startChannel("zalouser")).text).toContain("Waiting for channel login");
const assertCurrent = retainedOptions?.assertPersistentEffectCurrent;
if (!assertCurrent) {
throw new Error("Hosted channel setup did not retain its live owner assertion");
}
expect(assertCurrent).not.toThrow();
if (ending === "disposed") {
host.dispose();
} else {
expect((await host.resolveReply("continue")).configWritten).toBe(true);
}
expect(assertCurrent).toThrow();
const preparationCount = prepare.mock.calls.length;
await expect(retainedOptions?.beforePersistentEffect?.()).rejects.toThrow();
expect(prepare).toHaveBeenCalledTimes(preparationCount);
} finally {
host.dispose();
}
},
);
it("hosts a channel setup wizard as chat turns", async () => {
useTempStateDir();
const wizardRuns: string[] = [];

View file

@ -92,6 +92,7 @@ export async function runHostedChannelSetup(
prompter: WizardPrompter,
beforePersistentApply: (runtime: RuntimeEnv) => Promise<void>,
runtime?: RuntimeEnv,
assertPersistentEffectCurrent?: () => void,
): Promise<HostedSetupCompletion> {
const { createChannelSetupHooks, setupChannels } =
await import("../commands/onboard-channels.js");
@ -116,6 +117,7 @@ export async function runHostedChannelSetup(
skipDmPolicyPrompt: true,
skipConfirm: true,
beforePersistentEffect: async () => await beforePersistentApply(setupRuntime),
...(assertPersistentEffectCurrent ? { assertPersistentEffectCurrent } : {}),
onPostWriteHook: (hook) => channelSetup.onPostWriteHook(hook),
}),
afterWrite: async (configPath) => {

View file

@ -456,9 +456,13 @@ describe("WizardSession", () => {
expect(session.cancel()).toBe(false);
expect(session.getStatus()).toBe("running");
expect(session.signal.aborted).toBe(false);
expect(() => session.assertPersistentEffectCurrent()).not.toThrow();
finish();
expect((await session.next()).status).toBe("done");
expect(() => session.assertPersistentEffectCurrent()).toThrow(
"Setup session is no longer active",
);
});
test("expires an abandoned interactive session", async () => {

View file

@ -444,6 +444,14 @@ export class WizardSession {
this.cancellationLocked = true;
}
/** A retained write callback cannot outlive the runner that owns setup. */
assertPersistentEffectCurrent(): void {
this.signal.throwIfAborted();
if (this.status !== "running" || this.settled) {
throw new Error("Setup session is no longer active");
}
}
/** Protect preparation until the next client checkpoint or final commit. */
lockCancellationForPreparation() {
this.signal.throwIfAborted();

View file

@ -70,6 +70,7 @@ describe("Discord admission through Gateway policy publication", () => {
artifactBasename: "api.js",
});
const cfg: OpenClawConfig = {
plugins: { allow: ["discord"] },
channels: { discord: { token: "synthetic-token", groupPolicy: "allowlist", guilds: {} } },
messages: { inbound: { debounceMs: 0 } },
};
@ -235,7 +236,7 @@ describe("Discord admission through Gateway policy publication", () => {
await waitForFast(() =>
expect(committed, "policy-only publication must not wait for the active turn").toBe(next),
);
await pendingReload;
expect(await pendingReload).toBe("applied");
};
const activeTurn = tryBeginGatewayRootWorkAdmission();
expect(activeTurn).not.toBeNull();