From 4e9e68140aa3d80d7b1cd45e47faa1c8a6809375 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 23 Sep 2026 08:23:57 -0700 Subject: [PATCH] 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 1619aff4781603cb998646f5ac8068603dbd3e20 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. --- config/assertion-safety-baseline.txt | 2 +- docs/plugins/sdk-setup.md | 9 + .../zalouser/doctor-contract-api.test.ts | 8 +- .../zalouser/src/credential-persistence.ts | 197 ++++++++++++++ .../zalouser/src/send.handoff.test-support.ts | 6 +- extensions/zalouser/src/send.handoff.test.ts | 3 +- .../zalouser/src/send.media-handoff.test.ts | 3 +- extensions/zalouser/src/session-state.ts | 112 ++++---- extensions/zalouser/src/setup-surface.test.ts | 2 +- extensions/zalouser/src/setup-surface.ts | 10 +- .../zalouser/src/zalo-js.credentials.test.ts | 252 ++++++++++++++++-- .../zalouser/src/zalo-js.listener.test.ts | 9 +- extensions/zalouser/src/zalo-js.ts | 231 ++++------------ .../zalouser/src/zalo-quote-metadata.test.ts | 10 +- src/channels/plugins/setup-wizard-types.ts | 2 + src/commands/channels/add-wizard.ts | 8 + src/gateway/server-methods/wizard.ts | 3 + .../session-activity-summaries.test.ts | 20 +- src/system-agent/chat-wizard-host.ts | 38 ++- src/system-agent/hosted-setup.runtime.test.ts | 56 ++++ src/system-agent/hosted-setup.runtime.ts | 2 + src/wizard/session.test.ts | 4 + src/wizard/session.ts | 8 + test/discord-live-policy.integration.test.ts | 3 +- 24 files changed, 698 insertions(+), 300 deletions(-) create mode 100644 extensions/zalouser/src/credential-persistence.ts diff --git a/config/assertion-safety-baseline.txt b/config/assertion-safety-baseline.txt index fd5733ee94dc..a84033adf30d 100644 --- a/config/assertion-safety-baseline.txt +++ b/config/assertion-safety-baseline.txt @@ -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 diff --git a/docs/plugins/sdk-setup.md b/docs/plugins/sdk-setup.md index b526124e22ff..4e096ad2c08e 100644 --- a/docs/plugins/sdk-setup.md +++ b/docs/plugins/sdk-setup.md @@ -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. + 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(...)`. diff --git a/extensions/zalouser/doctor-contract-api.test.ts b/extensions/zalouser/doctor-contract-api.test.ts index b26316c9a8c4..b3cf2f2a12e9 100644 --- a/extensions/zalouser/doctor-contract-api.test.ts +++ b/extensions/zalouser/doctor-contract-api.test.ts @@ -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 = (options: OpenKeyedStoreOptions) => - createPluginStateSyncKeyedStoreForTests("zalouser", { + runtime.state.openKeyedStore = (options: OpenAsyncKeyedStoreOptions) => + createPluginStateKeyedStoreForTests("zalouser", { ...options, env: options.env ?? env, }); setZalouserRuntime(runtime); - clearStoredZaloCredentials(profile, env); + await clearStoredZaloCredentials(profile, env); const context = createDoctorContext(env); const params = { config: {}, diff --git a/extensions/zalouser/src/credential-persistence.ts b/extensions/zalouser/src/credential-persistence.ts new file mode 100644 index 000000000000..c3ab19f0b15b --- /dev/null +++ b/extensions/zalouser/src/credential-persistence.ts @@ -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 { + 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(); + private readonly credentialRefreshesByProfile = new Map>(); + private readonly credentialRevocationsByProfile = new Map>(); + + constructor(private readonly isCurrentApi: (profile: string, api: API) => boolean) {} + + pendingRevocation(profile: string): Promise | 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 { + 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, + ): Promise { + await this.writeCredentials(profile, snapshotApiCredentials(api, fallback), assertCurrent); + } + + async persistApiCredentialsIfChanged(profile: string, api: API): Promise { + 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 { + 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); + } + } + } +} diff --git a/extensions/zalouser/src/send.handoff.test-support.ts b/extensions/zalouser/src/send.handoff.test-support.ts index 2b4eba75d845..9ed70685d62b 100644 --- a/extensions/zalouser/src/send.handoff.test-support.ts +++ b/extensions/zalouser/src/send.handoff.test-support.ts @@ -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"), diff --git a/extensions/zalouser/src/send.handoff.test.ts b/extensions/zalouser/src/send.handoff.test.ts index 9fd4710a9f18..0023f4ce48be 100644 --- a/extensions/zalouser/src/send.handoff.test.ts +++ b/extensions/zalouser/src/send.handoff.test.ts @@ -14,9 +14,8 @@ vi.mock("./session-state.js", async (importOriginal) => ({ ...(await importOriginal()), 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", () => ({ diff --git a/extensions/zalouser/src/send.media-handoff.test.ts b/extensions/zalouser/src/send.media-handoff.test.ts index 5aa40d6b928e..c3ee72529774 100644 --- a/extensions/zalouser/src/send.media-handoff.test.ts +++ b/extensions/zalouser/src/send.media-handoff.test.ts @@ -14,9 +14,8 @@ vi.mock("./session-state.js", async (importOriginal) => ({ ...(await importOriginal()), 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"; diff --git a/extensions/zalouser/src/session-state.ts b/extensions/zalouser/src/session-state.ts index a46c04965083..6fbbe945f361 100644 --- a/extensions/zalouser/src/session-state.ts +++ b/extensions/zalouser/src/session-state.ts @@ -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 { - return getZalouserRuntime().state.openSyncKeyedStore({ - 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, - 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({ 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 { 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, + env: NodeJS.ProcessEnv = process.env, + assertCurrent?: () => void, +): Promise { + const normalizedProfile = normalizeZalouserCredentialProfile(profile); + await openZalouserCredentialsStore(env).register( + zalouserCredentialStoreKey(normalizedProfile), + { profile: normalizedProfile, ...credentials }, + { assertCurrent }, + ); +} + export async function refreshStoredZaloCredentials( profile: string, credentials: Omit, @@ -164,7 +152,7 @@ export async function refreshStoredZaloCredentials( env: NodeJS.ProcessEnv = process.env, ): Promise { 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 { 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; + } } diff --git a/extensions/zalouser/src/setup-surface.test.ts b/extensions/zalouser/src/setup-surface.test.ts index ddf819826301..2eaca6934374 100644 --- a/extensions/zalouser/src/setup-surface.test.ts +++ b/extensions/zalouser/src/setup-surface.test.ts @@ -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]!, diff --git a/extensions/zalouser/src/setup-surface.ts b/extensions/zalouser/src/setup-surface.ts index 8b5971d3bbce..cd6118564883 100644 --- a/extensions/zalouser/src/setup-surface.ts +++ b/extensions/zalouser/src/setup-surface.ts @@ -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); diff --git a/extensions/zalouser/src/zalo-js.credentials.test.ts b/extensions/zalouser/src/zalo-js.credentials.test.ts index 9280c935ede0..20952dace29e 100644 --- a/extensions/zalouser/src/zalo-js.credentials.test.ts +++ b/extensions/zalouser/src/zalo-js.credentials.test.ts @@ -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 { - 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, -): void { - saveStoredZaloCredentials(profile, credentials, { OPENCLAW_STATE_DIR: stateDir }); +): Promise { + 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) | undefined; + let beforeCredentialRevocation: (() => Promise) | undefined; + let beforeCredentialRegister: (() => Promise) | undefined; + let afterCredentialRegister: (() => void) | undefined; beforeEach(() => { resetPluginStateStoreForTests(); const runtime = createPluginRuntimeMock(); beforeCredentialApply = undefined; + beforeCredentialRevocation = undefined; + beforeCredentialRegister = undefined; + afterCredentialRegister = undefined; runtime.state.openKeyedStore = ( options: OpenAsyncKeyedStoreOptions, ): PluginStateKeyedStore => { const store = createPluginStateKeyedStoreForTests("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(); + const release = createDeferred(); + 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(); + const release = createDeferred(); + 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(); + const release = createDeferred(); + 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(); + 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" }], diff --git a/extensions/zalouser/src/zalo-js.listener.test.ts b/extensions/zalouser/src/zalo-js.listener.test.ts index f5660b8f1953..20f05bc3a6e7 100644 --- a/extensions/zalouser/src/zalo-js.listener.test.ts +++ b/extensions/zalouser/src/zalo-js.listener.test.ts @@ -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(); }); diff --git a/extensions/zalouser/src/zalo-js.ts b/extensions/zalouser/src/zalo-js.ts index 5126e6c7ce49..970829566cf0 100644 --- a/extensions/zalouser/src/zalo-js.ts +++ b/extensions/zalouser/src/zalo-js.ts @@ -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(); const apiInitByProfile = new Map>(); -const credentialSignaturesByProfile = new Map(); -const credentialRefreshesByProfile = new Map>(); +const credentials = new ZaloCredentialPersistence( + (profile, api) => apiByProfile.get(profile) === api, +); type CredentialPersistenceMode = "persist" | "read-only"; type CredentialPersistenceOptions = { credentialPersistence?: CredentialPersistenceMode }; -type ZaloCredentialPayload = Omit; type ActiveZaloQrLogin = { id: string; profile: string; startedAt: number; beforeCredentialPersistence?: () => Promise; + assertCredentialPersistenceCurrent?: () => void; qrDataUrl?: string; connected: boolean; error?: string; @@ -371,170 +373,17 @@ function mapGroup(groupId: string, group: GroupInfo & Record): }; } -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) - .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) - .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 { - 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, -): void { - writeCredentials(profile, snapshotApiCredentials(api, fallback)); -} - -async function persistApiCredentialsIfChanged(profile: string, api: API): Promise { - 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 { 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 = (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( ? 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 { 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; + 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, diff --git a/extensions/zalouser/src/zalo-quote-metadata.test.ts b/extensions/zalouser/src/zalo-quote-metadata.test.ts index a733e6b45d71..ed30b81e6b7d 100644 --- a/extensions/zalouser/src/zalo-quote-metadata.test.ts +++ b/extensions/zalouser/src/zalo-quote-metadata.test.ts @@ -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(params: { run: () => Promise; }): Promise { 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 = (options: OpenAsyncKeyedStoreOptions) => createPluginStateKeyedStoreForTests("zalouser", options); - runtime.state.openSyncKeyedStore = (options: OpenKeyedStoreOptions) => - createPluginStateSyncKeyedStoreForTests("zalouser", options); setZalouserRuntime(runtime); createZaloMock.mockReset(); }); diff --git a/src/channels/plugins/setup-wizard-types.ts b/src/channels/plugins/setup-wizard-types.ts index 971a33aa15f1..13c29f8424d3 100644 --- a/src/channels/plugins/setup-wizard-types.ts +++ b/src/channels/plugins/setup-wizard-types.ts @@ -309,6 +309,8 @@ export type SetupChannelsOptions = { allowSignalInstall?: boolean; /** Revalidate host authority immediately before an installer or other durable effect. */ beforePersistentEffect?: () => Promise; + /** 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>; diff --git a/src/commands/channels/add-wizard.ts b/src/commands/channels/add-wizard.ts index d04e5438b06e..15c4db84b987 100644 --- a/src/commands/channels/add-wizard.ts +++ b/src/commands/channels/add-wizard.ts @@ -112,6 +112,7 @@ type ChannelsAddWizardFlowParams = { prompter: WizardPrompter; initialChannel?: ChannelChoice; beforePersistentEffect?: () => Promise; + 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; + 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 } + : {}), }); } diff --git a/src/gateway/server-methods/wizard.ts b/src/gateway/server-methods/wizard.ts index 9acf0c2c475d..4cf97491152e 100644 --- a/src/gateway/server-methods/wizard.ts +++ b/src/gateway/server-methods/wizard.ts @@ -43,6 +43,7 @@ export type ChannelSetupWizardRunner = ( channel?: string; onConfigured?: (accounts: Array<{ channel: string; accountId: string }>) => void; beforePersistentEffect?: () => Promise; + 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, diff --git a/src/gateway/session-activity-summaries.test.ts b/src/gateway/session-activity-summaries.test.ts index b96aca840d3b..031a2dd8e8d4 100644 --- a/src/gateway/session-activity-summaries.test.ts +++ b/src/gateway/session-activity-summaries.test.ts @@ -328,15 +328,19 @@ describe("Activity recap lifecycle with the canonical session store", () => { totalMessages: 1, }); + const published = createDeferred>(); + 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, diff --git a/src/system-agent/chat-wizard-host.ts b/src/system-agent/chat-wizard-host.ts index 555f1f42bef8..ba18ab4bc21a 100644 --- a/src/system-agent/chat-wizard-host.ts +++ b/src/system-agent/chat-wizard-host.ts @@ -51,6 +51,7 @@ export type ChatWizardHostDependencies = { channel: string, prompter: WizardPrompter, beforePersistentApply: (runtime: RuntimeEnv) => Promise, + assertPersistentEffectCurrent?: () => void, ) => Promise; 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; + run: ( + prompter: WizardPrompter, + assertPersistentEffectCurrent: () => void, + ) => Promise; }): Promise { 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) { diff --git a/src/system-agent/hosted-setup.runtime.test.ts b/src/system-agent/hosted-setup.runtime.test.ts index 7a9da32af8c7..9e42128eb0eb 100644 --- a/src/system-agent/hosted-setup.runtime.test.ts +++ b/src/system-agent/hosted-setup.runtime.test.ts @@ -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[] = []; diff --git a/src/system-agent/hosted-setup.runtime.ts b/src/system-agent/hosted-setup.runtime.ts index 0a1e4b3916af..eb033ef655df 100644 --- a/src/system-agent/hosted-setup.runtime.ts +++ b/src/system-agent/hosted-setup.runtime.ts @@ -92,6 +92,7 @@ export async function runHostedChannelSetup( prompter: WizardPrompter, beforePersistentApply: (runtime: RuntimeEnv) => Promise, runtime?: RuntimeEnv, + assertPersistentEffectCurrent?: () => void, ): Promise { 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) => { diff --git a/src/wizard/session.test.ts b/src/wizard/session.test.ts index 677e4bcd8f26..ff0c8e374244 100644 --- a/src/wizard/session.test.ts +++ b/src/wizard/session.test.ts @@ -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 () => { diff --git a/src/wizard/session.ts b/src/wizard/session.ts index 8d54e0919199..c7f2e5f60573 100644 --- a/src/wizard/session.ts +++ b/src/wizard/session.ts @@ -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(); diff --git a/test/discord-live-policy.integration.test.ts b/test/discord-live-policy.integration.test.ts index cce5e088aa2d..0e08eca56fc9 100644 --- a/test/discord-live-policy.integration.test.ts +++ b/test/discord-live-policy.integration.test.ts @@ -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();