diff --git a/extensions/browser/src/browser/session-tab-cleanup.test.ts b/extensions/browser/src/browser/session-tab-cleanup.test.ts index d91b4e08995..ffc19560d60 100644 --- a/extensions/browser/src/browser/session-tab-cleanup.test.ts +++ b/extensions/browser/src/browser/session-tab-cleanup.test.ts @@ -43,13 +43,13 @@ describe("session tab cleanup timer", () => { const onWarn = vi.fn(); const stop = startTrackedBrowserTabCleanupTimer({ onWarn }); - await vi.advanceTimersByTimeAsync(60_000); + await vi.advanceTimersByTimeAsync(300_000); await vi.waitFor(() => expect(onWarn).toHaveBeenCalledWith( "failed to sweep tracked browser tabs: Error: sqlite read failed", ), ); - await vi.advanceTimersByTimeAsync(60_000); + await vi.advanceTimersByTimeAsync(300_000); await vi.waitFor(() => expect(registryMocks.sweepTrackedBrowserTabs).toHaveBeenCalledTimes(2)); await stop(); diff --git a/extensions/discord/src/monitor/message-handler.process.ack.test.ts b/extensions/discord/src/monitor/message-handler.process.ack.test.ts index fc00633b921..c8c9896ce11 100644 --- a/extensions/discord/src/monitor/message-handler.process.ack.test.ts +++ b/extensions/discord/src/monitor/message-handler.process.ack.test.ts @@ -519,17 +519,19 @@ describe("processDiscordMessage ack reactions", () => { outcome: "done", timingKey: "doneHoldMs", configuredHoldMs: 2_000, + builtInHoldMs: DEFAULT_TIMING.doneHoldMs, terminalEmoji: DEFAULT_EMOJIS.done, }, { outcome: "error", timingKey: "errorHoldMs", configuredHoldMs: 4_000, + builtInHoldMs: DEFAULT_TIMING.errorHoldMs, terminalEmoji: DEFAULT_EMOJIS.error, }, ] as const)( - "uses configured statusReactions.timing.$timingKey for $outcome cleanup", - async ({ outcome, timingKey, configuredHoldMs, terminalEmoji }) => { + "uses built-in statusReactions.timing.$timingKey for $outcome cleanup", + async ({ outcome, timingKey, configuredHoldMs, builtInHoldMs, terminalEmoji }) => { vi.useFakeTimers(); dispatchInboundMessage.mockImplementationOnce(async (params?: DispatchInboundParams) => { if (outcome === "done") { @@ -559,7 +561,7 @@ describe("processDiscordMessage ack reactions", () => { await runProcessDiscordMessage(ctx); expect(getReactionEmojis()).toContain(terminalEmoji); - await vi.advanceTimersByTimeAsync(configuredHoldMs - 1); + await vi.advanceTimersByTimeAsync(builtInHoldMs - 1); expect(sendMocks.removeReactionDiscord).not.toHaveBeenCalledWith( expect.anything(), expect.anything(), diff --git a/extensions/discord/src/monitor/provider.test.ts b/extensions/discord/src/monitor/provider.test.ts index f1979a1ee2a..5911da046c4 100644 --- a/extensions/discord/src/monitor/provider.test.ts +++ b/extensions/discord/src/monitor/provider.test.ts @@ -806,7 +806,7 @@ describe("monitorDiscordProvider", () => { }); }); - it("forwards custom eventQueue config from discord config to the Discord client", async () => { + it("ignores removed custom eventQueue config", async () => { resolveDiscordAccountMock.mockReturnValue({ accountId: "default", token: "MTIz.abc.def", @@ -825,7 +825,7 @@ describe("monitorDiscordProvider", () => { }); const eventQueue = getConstructedEventQueue(); - expect(eventQueue?.listenerTimeout).toBe(300_000); + expect(eventQueue?.listenerTimeout).toBe(120_000); }); it("does not pass eventQueue.listenerTimeout into the message run queue", async () => { diff --git a/extensions/memory-core/src/dreaming-narrative.test.ts b/extensions/memory-core/src/dreaming-narrative.test.ts index 58257da16c1..fa29619e19d 100644 --- a/extensions/memory-core/src/dreaming-narrative.test.ts +++ b/extensions/memory-core/src/dreaming-narrative.test.ts @@ -104,6 +104,7 @@ function readSessionStoreEntries(storePath: string): Record { const updatedAt = Date.now(); await seedSessionStore(storePath, { "agent:main:dreaming-narrative-light-1": { - sessionId: "missing", + sessionId: "orphan", + sessionFile: orphanPath, updatedAt, }, "agent:main:kept-session": { @@ -886,12 +888,14 @@ describe("generateAndAppendDreamNarrative", () => { }); await seedDreamingTranscriptEvent({ sessionId: "orphan", + sessionKey: "agent:main:dreaming-narrative-light-1", storePath, timestampMs: Date.now() - 600_000, runId: "dreaming-narrative-light-123", }); await seedDreamingTranscriptEvent({ sessionId: "still-live", + sessionKey: "agent:main:kept-session", storePath, timestampMs: Date.now(), runId: "dreaming-narrative-light-keep", @@ -969,12 +973,14 @@ describe("generateAndAppendDreamNarrative", () => { }); await seedDreamingTranscriptEvent({ sessionId: "orphan-dreaming", + sessionKey: "agent:main:dreaming-narrative-deep-orphan", storePath, timestampMs: Date.now() - 600_000, runId: "dreaming-narrative-deep-orphan", }); await seedDreamingTranscriptEvent({ sessionId: "live-dreaming", + sessionKey: "agent:main:dreaming-narrative-deep-live", storePath, timestampMs: Date.now(), runId: "dreaming-narrative-deep-live", diff --git a/extensions/memory-core/src/memory/manager.watcher-config.test.ts b/extensions/memory-core/src/memory/manager.watcher-config.test.ts index e487635a59f..95c2ec57ebf 100644 --- a/extensions/memory-core/src/memory/manager.watcher-config.test.ts +++ b/extensions/memory-core/src/memory/manager.watcher-config.test.ts @@ -11,6 +11,8 @@ import { afterAll, afterEach, beforeEach, describe, expect, it, vi } from "vites type WatchIgnoredFn = (watchPath: string, stats?: { isDirectory?: () => boolean }) => boolean; +const BUILT_IN_WATCH_DEBOUNCE_MS = 1_500; + const { createdChokidarWatchers, createdNativeWatchers, @@ -380,7 +382,7 @@ describe("memory watcher config", () => { .mockResolvedValue(undefined); createdChokidarWatchers[0]?.emit(event); - await vi.advanceTimersByTimeAsync(25); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(syncSpy).toHaveBeenCalledWith({ reason: "watch" }); }, @@ -407,7 +409,7 @@ describe("memory watcher config", () => { (w) => w.dir === path.join(workspaceDir, "memory"), ); memoryWatcher?.emit(eventType, "notes.md"); - await vi.advanceTimersByTimeAsync(25); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(syncSpy).toHaveBeenCalledWith({ reason: "watch" }); }, @@ -434,7 +436,7 @@ describe("memory watcher config", () => { // Node docs warn that filename may be null on some platforms; conservative // dirty must still be scheduled. memoryWatcher?.emit("rename", null as unknown as string); - await vi.advanceTimersByTimeAsync(50); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(syncSpy).toHaveBeenCalledWith({ reason: "watch" }); }); @@ -495,7 +497,7 @@ describe("memory watcher config", () => { ); memoryWatcher?.emitError(new Error("watcher error: ENOSPC")); - await vi.advanceTimersByTimeAsync(50); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(memoryLoggerWarn).toHaveBeenCalledWith( expect.stringContaining("memory native watcher error"), @@ -511,7 +513,7 @@ describe("memory watcher config", () => { // continues to schedule sync. syncSpy.mockClear(); existingChokidar?.emit("change"); - await vi.advanceTimersByTimeAsync(25); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(syncSpy).toHaveBeenCalledWith({ reason: "watch" }); }); @@ -611,7 +613,7 @@ describe("memory watcher config", () => { await fs.mkdir(newDir); const memoryWatcher = createdNativeWatchers.find((w) => w.dir === memoryDir); memoryWatcher?.emit("rename", "new-topic"); - await vi.advanceTimersByTimeAsync(25); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(syncSpy).toHaveBeenCalledWith({ reason: "watch" }); expect( @@ -1005,10 +1007,10 @@ describe("memory watcher config", () => { extraWatcher?.emit("change", "notes.md"); await fs.writeFile(notesPath, "hello updated"); - await vi.advanceTimersByTimeAsync(25); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(syncSpy).not.toHaveBeenCalled(); - await vi.advanceTimersByTimeAsync(25); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(syncSpy).toHaveBeenCalledWith({ reason: "watch" }); // Recorded path should match the resolved absolute path under extraDir. const recordedStats = (initialStats as unknown as { isDirectory: () => boolean }).isDirectory(); diff --git a/extensions/memory-core/src/memory/qmd-manager.test.ts b/extensions/memory-core/src/memory/qmd-manager.test.ts index 1ecaf547772..91c87cd7b2d 100644 --- a/extensions/memory-core/src/memory/qmd-manager.test.ts +++ b/extensions/memory-core/src/memory/qmd-manager.test.ts @@ -49,6 +49,7 @@ const MEMORY_EMBEDDING_PROVIDERS_KEY = Symbol.for("openclaw.memoryEmbeddingProvi const MCPORTER_STATE_KEY = Symbol.for("openclaw.mcporterState"); const QMD_EMBED_QUEUE_KEY = Symbol.for("openclaw.qmdEmbedQueueTail"); const QMD_UPDATE_QUEUE_KEY = Symbol.for("openclaw.qmdUpdateQueueState"); +const BUILT_IN_WATCH_DEBOUNCE_MS = 1_500; type WatchOptions = { ignored?: (watchPath: string) => boolean; @@ -1193,7 +1194,7 @@ describe("QmdMemoryManager", () => { }); expect(manager.status().dirty).toBe(true); - await vi.advanceTimersByTimeAsync(25); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); const updateCalls = spawnMock.mock.calls.filter((call) => call[1]?.[0] === "update"); expect(updateCalls).toHaveLength(1); @@ -1328,10 +1329,10 @@ describe("QmdMemoryManager", () => { }); await fs.writeFile(notesPath, "hello updated"); - await vi.advanceTimersByTimeAsync(25); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(spawnMock.mock.calls.filter((call) => call[1]?.[0] === "update")).toHaveLength(0); - await vi.advanceTimersByTimeAsync(25); + await vi.advanceTimersByTimeAsync(BUILT_IN_WATCH_DEBOUNCE_MS); expect(spawnMock.mock.calls.filter((call) => call[1]?.[0] === "update")).toHaveLength(1); await manager.close(); diff --git a/extensions/telegram/src/bot-message-dispatch.pipeline-init.test.ts b/extensions/telegram/src/bot-message-dispatch.pipeline-init.test.ts index fa4262828be..a2f0bbb078e 100644 --- a/extensions/telegram/src/bot-message-dispatch.pipeline-init.test.ts +++ b/extensions/telegram/src/bot-message-dispatch.pipeline-init.test.ts @@ -1,3 +1,4 @@ +import { DEFAULT_TIMING } from "openclaw/plugin-sdk/channel-feedback"; import { expect, it, vi } from "vitest"; import { createChannelMessageReplyPipeline, @@ -12,6 +13,7 @@ import { notifyTelegramInboundEventOutboundSuccess } from "./inbound-event-deliv describeTelegramDispatch("dispatchTelegramMessage pipeline-init", () => { it("cleans delivery correlation when reply-pipeline initialization fails", async () => { + vi.useFakeTimers(); const sessionKey = "agent:main:telegram:direct:pipeline-init-failure"; const statusReactionController = createStatusReactionController(); const reactionApi = vi.fn(async () => undefined); @@ -41,7 +43,8 @@ describeTelegramDispatch("dispatchTelegramMessage pipeline-init", () => { suppressFailureFallback: true, }); - await vi.waitFor(() => expect(statusReactionController.restoreInitial).toHaveBeenCalled()); + await vi.advanceTimersByTimeAsync(DEFAULT_TIMING.errorHoldMs); + expect(statusReactionController.restoreInitial).toHaveBeenCalled(); expect(reactionApi).not.toHaveBeenCalled(); }); }); diff --git a/extensions/telegram/src/bot.create-telegram-bot.test.ts b/extensions/telegram/src/bot.create-telegram-bot.test.ts index b27ea136caa..6250fb53b94 100644 --- a/extensions/telegram/src/bot.create-telegram-bot.test.ts +++ b/extensions/telegram/src/bot.create-telegram-bot.test.ts @@ -325,14 +325,14 @@ describe("createTelegramBot", () => { globalThis.fetch = originalFetch; } }); - it("applies global and per-account timeoutSeconds", () => { + it("ignores removed global and per-account timeoutSeconds", () => { loadConfig.mockReturnValue({ channels: { telegram: { dmPolicy: "open", allowFrom: ["*"], timeoutSeconds: 60 }, }, }); createTelegramBot({ token: "tok" }); - expectBotClientFields({ timeoutSeconds: 60 }); + expectBotClientFields({ timeoutSeconds: undefined }); botCtorSpy.mockClear(); loadConfig.mockReturnValue({ @@ -348,7 +348,7 @@ describe("createTelegramBot", () => { }, }); createTelegramBot({ token: "tok", accountId: "foo" }); - expectBotClientFields({ timeoutSeconds: 61 }); + expectBotClientFields({ timeoutSeconds: undefined }); }); it("keeps low timeoutSeconds above the outbound request guard", () => { @@ -358,7 +358,7 @@ describe("createTelegramBot", () => { }, }); createTelegramBot({ token: "tok" }); - expectBotClientFields({ timeoutSeconds: 60 }); + expectBotClientFields({ timeoutSeconds: undefined }); }); it("keeps polling client timeout above the outbound request guard", () => { @@ -368,7 +368,7 @@ describe("createTelegramBot", () => { }, }); createTelegramBot({ token: "tok", minimumClientTimeoutSeconds: 45 }); - expectBotClientFields({ timeoutSeconds: 60 }); + expectBotClientFields({ timeoutSeconds: undefined }); }); it("passes startup probe botInfo to grammY", () => { diff --git a/extensions/telegram/src/channel.gateway.test.ts b/extensions/telegram/src/channel.gateway.test.ts index 76349792625..8c2799e4f96 100644 --- a/extensions/telegram/src/channel.gateway.test.ts +++ b/extensions/telegram/src/channel.gateway.test.ts @@ -552,7 +552,7 @@ describe("telegramPlugin gateway startup", () => { ).resolves.toBeNull(); }); - it("honors higher per-account timeoutSeconds for startup probe", async () => { + it("uses the built-in startup probe timeout", async () => { installTelegramRuntime(); probeTelegram.mockResolvedValue({ ok: true, @@ -565,7 +565,7 @@ describe("telegramPlugin gateway startup", () => { const { ctx, task } = startTelegramAccount("ops", { timeoutSeconds: 60 }); await expect(task).resolves.toBeUndefined(); - expect(probeTelegram).toHaveBeenCalledWith("123456:bad-token", 60_000, { + expect(probeTelegram).toHaveBeenCalledWith("123456:bad-token", 15_000, { abortSignal: ctx.abortSignal, accountId: "ops", proxyUrl: undefined, diff --git a/extensions/telegram/src/send.test.ts b/extensions/telegram/src/send.test.ts index d8e8dc5839b..a23e76f9390 100644 --- a/extensions/telegram/src/send.test.ts +++ b/extensions/telegram/src/send.test.ts @@ -853,13 +853,13 @@ describe("sendMessageTelegram", () => { expect(botCtorSpy).not.toHaveBeenCalled(); }); - it("applies timeoutSeconds config precedence", async () => { + it("ignores removed timeoutSeconds config", async () => { const cases = [ { name: "global telegram timeout", cfg: { channels: { telegram: { timeoutSeconds: 60 } } }, opts: { cfg: TELEGRAM_TEST_CFG, token: "tok" }, - expectedTimeout: 60, + expectedTimeout: undefined, }, { name: "per-account timeout override", @@ -872,7 +872,7 @@ describe("sendMessageTelegram", () => { }, }, opts: { cfg: TELEGRAM_TEST_CFG, token: "tok", accountId: "foo" }, - expectedTimeout: 61, + expectedTimeout: undefined, }, ] as const; for (const testCase of cases) { diff --git a/scripts/e2e/lib/gateway-network/client.d.mts b/scripts/e2e/lib/gateway-network/client.d.mts index 5edbff0cfd9..448f5545ccd 100644 --- a/scripts/e2e/lib/gateway-network/client.d.mts +++ b/scripts/e2e/lib/gateway-network/client.d.mts @@ -20,6 +20,14 @@ export function assertSuspendedProbes( health: Record, readiness: Record, ): void; +export function prepareReadySuspension( + options: { + deadline: number; + requestId: string; + rpc: (method: string, params: Record) => Promise>; + }, + deps?: { delayImpl?: (ms: number) => Promise; now?: () => number }, +): Promise>; export function runGatewaySuspensionPreRestartClient( options: { statePath: string; token: string; url: string; timeoutMs?: number }, deps?: { fetchImpl?: typeof fetch }, diff --git a/scripts/e2e/lib/gateway-network/client.mjs b/scripts/e2e/lib/gateway-network/client.mjs index 8d8071026e4..92d25e586df 100644 --- a/scripts/e2e/lib/gateway-network/client.mjs +++ b/scripts/e2e/lib/gateway-network/client.mjs @@ -148,6 +148,30 @@ export function assertReadySuspensionResponse(response, now = Date.now()) { return payload; } +export async function prepareReadySuspension( + { deadline, requestId, rpc }, + { delayImpl = delay, now = Date.now } = {}, +) { + while (true) { + if (now() >= deadline) { + throw new DOMException("gateway suspension preparation timeout", "TimeoutError"); + } + const response = await rpc("gateway.suspend.prepare", { requestId }); + if (response?.status !== 200 || response.body?.ok !== true) { + return assertReadySuspensionResponse(response, now()); + } + const payload = response.body.payload; + if (payload?.status !== "busy") { + return assertReadySuspensionResponse(response, now()); + } + const retryAfterMs = + typeof payload.retryAfterMs === "number" && Number.isFinite(payload.retryAfterMs) + ? Math.max(1, Math.floor(payload.retryAfterMs)) + : 100; + await delayImpl(Math.min(retryAfterMs, Math.max(1, deadline - now()))); + } +} + export function assertGatewaySuspendingError(response) { assert.equal(response?.ok, false, "normal RPC must fail during suspension"); assert.equal(response.error?.code, "UNAVAILABLE", "normal RPC must be unavailable"); @@ -206,9 +230,11 @@ export async function runGatewaySuspensionPreRestartClient( url, }; const rpc = (method, params) => adminRpc(requestContext, method, params); - const firstLease = assertReadySuspensionResponse( - await rpc("gateway.suspend.prepare", { requestId: "gateway-network-live-contract" }), - ); + const firstLease = await prepareReadySuspension({ + deadline: requestContext.deadline, + requestId: "gateway-network-live-contract", + rpc, + }); assertSuspendedProbes( await readProbe(requestContext, "/healthz"), @@ -263,9 +289,11 @@ export async function runGatewaySuspensionPreRestartClient( ); const requestId = "gateway-network-restart-contract"; - const secondLease = assertReadySuspensionResponse( - await rpc("gateway.suspend.prepare", { requestId }), - ); + const secondLease = await prepareReadySuspension({ + deadline: requestContext.deadline, + requestId, + rpc, + }); await writeFile( statePath, JSON.stringify({ @@ -309,9 +337,11 @@ export async function runGatewaySuspensionPostRestartClient( ); assertAdminSuccess(await rpc("health"), "Admin health after restart"); - const replacement = assertReadySuspensionResponse( - await rpc("gateway.suspend.prepare", { requestId: state.requestId }), - ); + const replacement = await prepareReadySuspension({ + deadline: requestContext.deadline, + requestId: state.requestId, + rpc, + }); assert( replacement.suspensionId !== state.suspensionId, "reused request id must create a fresh suspension lease after restart", diff --git a/test/scripts/gateway-network-client.test.ts b/test/scripts/gateway-network-client.test.ts index e0dc993df1b..df04876915c 100644 --- a/test/scripts/gateway-network-client.test.ts +++ b/test/scripts/gateway-network-client.test.ts @@ -8,6 +8,7 @@ import { assertGatewaySuspendingError, assertReadySuspensionResponse, assertSuspendedProbes, + prepareReadySuspension, runGatewayNetworkClient, runGatewaySuspensionPostRestartClient, runGatewaySuspensionPreRestartClient, @@ -80,6 +81,68 @@ describe("gateway network client", () => { ).toBe(3000); }); + it("retries busy suspension preparation using the server delay", async () => { + const expiresAtMs = Date.now() + 10_000; + const rpc = vi + .fn() + .mockResolvedValueOnce({ + status: 200, + body: { + ok: true, + payload: { status: "busy", retryAfterMs: 250, activeCount: 1, blockers: ["agent"] }, + }, + }) + .mockResolvedValueOnce({ + status: 200, + body: { + ok: true, + payload: { + status: "ready", + suspensionId: "lease-1", + expiresAtMs, + activeCount: 0, + blockers: [], + }, + }, + }); + const delayImpl = vi.fn(async () => undefined); + + await expect( + prepareReadySuspension( + { deadline: Date.now() + 5_000, requestId: "request-1", rpc }, + { delayImpl }, + ), + ).resolves.toMatchObject({ status: "ready", suspensionId: "lease-1" }); + expect(rpc).toHaveBeenCalledTimes(2); + expect(rpc).toHaveBeenNthCalledWith(1, "gateway.suspend.prepare", { + requestId: "request-1", + }); + expect(delayImpl).toHaveBeenCalledWith(250); + }); + + it("stops busy suspension retries at the client deadline", async () => { + let nowMs = 1_000; + const rpc = vi.fn(async () => ({ + status: 200, + body: { + ok: true, + payload: { status: "busy", retryAfterMs: 250, activeCount: 1, blockers: ["agent"] }, + }, + })); + const delayImpl = vi.fn(async (ms: number) => { + nowMs += ms; + }); + + await expect( + prepareReadySuspension( + { deadline: 1_250, requestId: "request-busy", rpc }, + { delayImpl, now: () => nowMs }, + ), + ).rejects.toMatchObject({ name: "TimeoutError" }); + expect(rpc).toHaveBeenCalledOnce(); + expect(delayImpl).toHaveBeenCalledWith(250); + }); + it("bounds a stalled suspension admin request by the client deadline", async () => { let requestSignal: AbortSignal | null | undefined; const fetchImpl = vi.fn((_input, init) => {