diff --git a/extensions/cloudflare/assets/activity.svg b/extensions/cloudflare/assets/activity.svg new file mode 100644 index 000000000000..6c2837636e32 --- /dev/null +++ b/extensions/cloudflare/assets/activity.svg @@ -0,0 +1,32 @@ + + + + diff --git a/extensions/cloudflare/assets/icon.png b/extensions/cloudflare/assets/icon.png new file mode 100644 index 000000000000..2cae8a9f9b76 Binary files /dev/null and b/extensions/cloudflare/assets/icon.png differ diff --git a/extensions/nostr/src/nostr-bus.inbound.test.ts b/extensions/nostr/src/nostr-bus.inbound.test.ts index 2be88d6d0a3c..185a98e24408 100644 --- a/extensions/nostr/src/nostr-bus.inbound.test.ts +++ b/extensions/nostr/src/nostr-bus.inbound.test.ts @@ -6,6 +6,8 @@ import { closeOpenClawStateDatabaseForTest, createChannelIngressQueueForTests, } from "openclaw/plugin-sdk/channel-ingress-test-runtime"; +import { createDeferred } from "openclaw/plugin-sdk/extension-shared"; +import { withinTest } from "openclaw/plugin-sdk/test-fixtures"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import type { PluginRuntime } from "../runtime-api.js"; import { startNostrBus } from "./nostr-bus.js"; @@ -777,15 +779,19 @@ describe("startNostrBus inbound guards", () => { await bus.close(); }); - it("does not rate limit an allowed sender while another authorization is still pending", async () => { - const onMessage = vi.fn(async () => {}); - let resolveBlocked: ((value: "block") => void) | undefined; - const blockedPromise = new Promise<"block">((resolve) => { - resolveBlocked = resolve; - }); + it("does not rate limit an allowed sender while another authorization is still pending", async ({ + signal, + }) => { + const delivered = createDeferred(); + const onMessage = vi.fn(async () => delivered.resolve()); + const authorizing = createDeferred(); + const blocked = createDeferred<"block">(); const authorizeSender = vi .fn<(params: { senderPubkey: string }) => Promise<"allow" | "block" | "pairing">>() - .mockImplementationOnce(async () => await blockedPromise) + .mockImplementationOnce(async () => { + authorizing.resolve(); + return await blocked.promise; + }) .mockResolvedValueOnce("allow"); const bus = await startTestNostrBus({ privateKey: TEST_HEX_PRIVATE_KEY, @@ -802,30 +808,33 @@ describe("startNostrBus inbound guards", () => { }, }); - const handlers = mockState.handlers[0]; - if (!handlers) { - throw new Error("missing subscription handlers"); + try { + const handlers = mockState.handlers[0]; + if (!handlers) { + throw new Error("missing subscription handlers"); + } + void handlers.onevent( + createEvent({ id: "blocked-pending", pubkey: `blocked${"a".repeat(57)}` }), + ); + await withinTest(authorizing.promise, signal); + void handlers.onevent( + createEvent({ + id: "allowed-during-pending-auth", + pubkey: `allowed${"b".repeat(57)}`, + }), + ); + await withinTest(delivered.promise, signal); + blocked.resolve("block"); + await Promise.all(ingressTasks.splice(0)); + + expect(authorizeSender).toHaveBeenCalledTimes(2); + expect(mockState.decrypt).toHaveBeenCalledTimes(1); + expect(onMessage).toHaveBeenCalledTimes(1); + expect(bus.getMetrics().eventsRejected.rateLimited).toBe(0); + } finally { + blocked.resolve("block"); + await bus.close(); } - void handlers.onevent( - createEvent({ id: "blocked-pending", pubkey: `blocked${"a".repeat(57)}` }), - ); - await vi.waitFor(() => expect(authorizeSender).toHaveBeenCalledTimes(1)); - void handlers.onevent( - createEvent({ - id: "allowed-during-pending-auth", - pubkey: `allowed${"b".repeat(57)}`, - }), - ); - await vi.waitFor(() => expect(onMessage).toHaveBeenCalledTimes(1)); - resolveBlocked?.("block"); - await Promise.all(ingressTasks.splice(0)); - - expect(authorizeSender).toHaveBeenCalledTimes(2); - expect(mockState.decrypt).toHaveBeenCalledTimes(1); - expect(onMessage).toHaveBeenCalledTimes(1); - expect(bus.getMetrics().eventsRejected.rateLimited).toBe(0); - - await bus.close(); }); it("rate limits repeated invalid signatures before authorization work fans out", async () => { diff --git a/src/auto-reply/reply/agent-runner-memory.ts b/src/auto-reply/reply/agent-runner-memory.ts index 83391d0acc21..32ca38072463 100644 --- a/src/auto-reply/reply/agent-runner-memory.ts +++ b/src/auto-reply/reply/agent-runner-memory.ts @@ -489,12 +489,11 @@ async function estimateProviderPromptTokens( : undefined; } -async function estimatePromptTokensFromSessionTranscript(params: { - agentId?: string; +async function estimatePromptTokensFromSessionTranscript({ + abortSignal, + ...params +}: Parameters[0] & { abortSignal?: AbortSignal; - sessionId?: string; - sessionKey?: string; - storePath?: string; contextWindowTokens: number; }): Promise { const sessionId = normalizeOptionalString(params.sessionId); @@ -550,10 +549,9 @@ async function estimatePromptTokensFromSessionTranscript(params: { const messages = await readPreflightTranscriptContextMessages( { ...params, - agentId: params.agentId ?? resolveAgentIdFromSessionKey(params.sessionKey), sessionId, }, - params.abortSignal, + abortSignal, ); const estimatedTokens = await estimateProviderPromptTokens( messages, @@ -572,7 +570,7 @@ async function estimatePromptTokensFromSessionTranscript(params: { transcriptByteSize: snapshot.byteSize, }; } catch (error) { - params.abortSignal?.throwIfAborted(); + abortSignal?.throwIfAborted(); return error instanceof SessionTranscriptReadFenceError ? Promise.reject(error) : undefined; } } diff --git a/src/auto-reply/reply/directive-handling.model.test.ts b/src/auto-reply/reply/directive-handling.model.test.ts index 822418d2e5cd..bb6c94031579 100644 --- a/src/auto-reply/reply/directive-handling.model.test.ts +++ b/src/auto-reply/reply/directive-handling.model.test.ts @@ -119,9 +119,8 @@ import { type InternalHookEvent, } from "../../hooks/internal-hooks.js"; import { enqueueSystemEvent } from "../../infra/system-events.js"; -import { closeOpenClawAgentDatabasesForTest } from "../../state/openclaw-agent-db.js"; -import { closeOpenClawStateDatabaseForTest } from "../../state/openclaw-state-db.js"; import { withEnvAsync } from "../../test-utils/env.js"; +import { cleanupSessionStateForTest } from "../../test-utils/session-state-cleanup.js"; import type { ElevatedLevel } from "../thinking.js"; import { registerModelRuntimeDirectiveTests } from "./directive-handling.model-runtime.test-support.js"; import { registerModelStatusDirectiveTests } from "./directive-handling.model-status.test-support.js"; @@ -687,7 +686,8 @@ describe("/model chat UX", () => { "%s reads terminal fallback from the transcript scope, not the runtime-policy key", async (command) => { const tempRoot = tempDirs.make("openclaw-model-terminal-display-"); - await withEnvAsync({ OPENCLAW_STATE_DIR: path.join(tempRoot, "state") }, async () => { + const stateDir = path.join(tempRoot, "state"); + await withEnvAsync({ OPENCLAW_STATE_DIR: stateDir }, async () => { const sessionKey = "agent:main:main"; const storePath = path.join(tempRoot, "custom-store", "openclaw-agent.sqlite"); const scope = { agentId: "main", sessionKey, sessionId: "terminal-display", storePath }; @@ -740,8 +740,7 @@ describe("/model chat UX", () => { expect(sessionEntry).toEqual(before); expect(loadSessionEntry(scope)).toEqual(before); } finally { - closeOpenClawAgentDatabasesForTest(tempRoot); - closeOpenClawStateDatabaseForTest(); + await cleanupSessionStateForTest({ stateDir, rootPath: tempRoot }); } }); }, diff --git a/src/auto-reply/reply/reply-turn-admission.database-claim.test.ts b/src/auto-reply/reply/reply-turn-admission.database-claim.test.ts index a77343f6efab..e2c3ebee0df1 100644 --- a/src/auto-reply/reply/reply-turn-admission.database-claim.test.ts +++ b/src/auto-reply/reply/reply-turn-admission.database-claim.test.ts @@ -1,7 +1,11 @@ import fs from "node:fs"; import path from "node:path"; import { afterEach, expect, it, vi } from "vitest"; -import { createDeferred } from "../../../test/helpers/promise.js"; +import { + awaitGateBeforeSettlement, + createDeferred, + withinTest, +} from "../../../test/helpers/promise.js"; import { useAutoCleanupTempDirTracker } from "../../../test/helpers/temp-dir.js"; import * as sessionEntries from "../../config/sessions/session-accessor.sqlite-entry.js"; import { runExclusiveSessionStoreWrite } from "../../config/sessions/store-writer.js"; @@ -170,7 +174,7 @@ it.each(["cancelled", "request-changed", "later-rebound-store"] as const)( }, ); -it.each( +it.for( (["writer", "active", "delivery"] as const).flatMap((wait) => (["unchanged", "same-inode", "other-inode"] as const).map((replacement) => ({ wait, @@ -179,7 +183,7 @@ it.each( ), )( "keeps the exact database owner across $wait wait, replacement=$replacement", - async ({ wait, replacement }) => { + async ({ wait, replacement }, { signal }) => { const root = tempDirs.make("reply-admission-claim-"); const originalPath = path.join(root, "original.sqlite"); const replacementPath = path.join(root, "replacement.sqlite"); @@ -210,13 +214,29 @@ it.each( owner.completeWithAfterClearBarrier(release.promise); } } - const loaded = vi.spyOn(sessionEntries, "loadSessionEntryForAdmission"); - const waiting = - wait === "active" - ? vi.spyOn(registry.replyRunRegistry, "waitForIdle") - : wait === "delivery" - ? vi.spyOn(registry, "waitForReplyRunFollowupAdmission") - : loaded; + const enteredWait = createDeferred(); + const load = sessionEntries.loadSessionEntryForAdmission; + const loaded = vi + .spyOn(sessionEntries, "loadSessionEntryForAdmission") + .mockImplementation((...args) => { + if (wait === "writer") { + enteredWait.resolve(); + } + return load(...args); + }); + if (wait === "active") { + const waitForIdle = registry.replyRunRegistry.waitForIdle.bind(registry.replyRunRegistry); + vi.spyOn(registry.replyRunRegistry, "waitForIdle").mockImplementation((...args) => { + enteredWait.resolve(); + return waitForIdle(...args); + }); + } else if (wait === "delivery") { + const waitForAdmission = registry.waitForReplyRunFollowupAdmission; + vi.spyOn(registry, "waitForReplyRunFollowupAdmission").mockImplementation((...args) => { + enteredWait.resolve(); + return waitForAdmission(...args); + }); + } const controller = new AbortController(); const pending = admit(storePath, { expectedSessionId: sessionId, @@ -225,7 +245,14 @@ it.each( }); void pending.catch(() => {}); try { - await vi.waitFor(() => expect(waiting).toHaveBeenCalled()); + await withinTest( + awaitGateBeforeSettlement( + enteredWait.promise, + pending, + "Admission completed before reaching its owned wait", + ), + signal, + ); const observed = loaded.mock.results.at(-1); if (observed?.type !== "return") { throw new Error("fixture requires a completed authoritative row read"); diff --git a/src/config/sessions/session-cold-storage.test.ts b/src/config/sessions/session-cold-storage.test.ts index 35a3d9ca9e06..d36d7343a2cc 100644 --- a/src/config/sessions/session-cold-storage.test.ts +++ b/src/config/sessions/session-cold-storage.test.ts @@ -337,6 +337,7 @@ describe("cold transcript storage workers", () => { clearOpenClawAgentIntegrityVerification(fixture.options.path); await restoreSessionColdTranscript(fixture.secondScope); expect(fixture.snapshot()).toEqual(fixture.original); + await closeOpenClawAgentDatabaseByPathAsync(fixture.options.path); await flushLogger(); const summaries = (await fs.readFile(file, "utf8")) .split("\n") @@ -803,6 +804,7 @@ describe("cold transcript storage workers", () => { db.prepare("UPDATE schema_meta SET schema_version = ? WHERE meta_key = 'primary'").run( version, ); + await closeOpenClawAgentDatabaseByPathAsync(fixture.options.path); closeOpenClawAgentDatabasesForTest(); const before = createHash("sha256") .update(await fs.readFile(fixture.scope.storePath)) diff --git a/src/config/sessions/session-history-budget-owner.test.ts b/src/config/sessions/session-history-budget-owner.test.ts index 95cf4391a2a8..d6f90a94a6b3 100644 --- a/src/config/sessions/session-history-budget-owner.test.ts +++ b/src/config/sessions/session-history-budget-owner.test.ts @@ -182,8 +182,21 @@ it.each([ protectedKey, ), ).toEqual({ current_session_id: currentId }); - // Forget both the handle and process validation, exposing registration as well as lease drift. - closeOpenClawAgentDatabasesForTest(state.root); + // Forget cached reads without revoking the active sweep's database workers. + const databaseOptions = { + agentId: target.agentId ?? "main", + path: databasePath, + env: state.env, + }; + const forgetCachedDatabase = () => { + const cached = getOpenClawAgentDatabaseIfOpen(databaseOptions); + if (cached) { + closeCachedOpenClawAgentDatabase(cached, { eviction: true }); + } + clearOpenClawAgentDatabaseValidationCache(state.root); + expect(getOpenClawAgentDatabaseIfOpen(databaseOptions)).toBeUndefined(); + }; + forgetCachedDatabase(); let capEntryCalls = 0; const deleteEntry = entryEviction.deleteDiskBudgetArchivedSessionEntry; @@ -198,18 +211,7 @@ it.each([ sessionKey, ), ).toEqual({ current_session_id: originalId }); - // Evict the host handle before the lazy loader without revoking this active sweep's workers. - const databaseOptions = { - agentId: target.agentId ?? "main", - path: databasePath, - env: state.env, - }; - const cached = getOpenClawAgentDatabaseIfOpen(databaseOptions); - if (cached) { - closeCachedOpenClawAgentDatabase(cached, { eviction: true }); - } - clearOpenClawAgentDatabaseValidationCache(state.root); - expect(getOpenClawAgentDatabaseIfOpen(databaseOptions)).toBeUndefined(); + forgetCachedDatabase(); } return await deleteEntry(...args); }, @@ -382,7 +384,7 @@ it.each([ moveRelativeCwd?.(); // Patch commit reopened A. Remove its handle and validation before allowing // the REAL first measurement to return to enforcement/preview. - closeOpenClawAgentDatabasesForTest(state.root); + forgetCachedDatabase(); release.resolve(); if (trigger === "inspect") { await expect(sweep).resolves.toMatchObject({ diff --git a/src/gateway/health/collector.session-row-projection.test.ts b/src/gateway/health/collector.session-row-projection.test.ts index 4c579254a1ce..4aa15380b371 100644 --- a/src/gateway/health/collector.session-row-projection.test.ts +++ b/src/gateway/health/collector.session-row-projection.test.ts @@ -1,7 +1,7 @@ import assert from "node:assert/strict"; import path from "node:path"; import { DatabaseSync } from "node:sqlite"; -import { afterEach, describe, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; import { replaceSessionEntrySync, upsertSessionEntryCore, @@ -18,7 +18,12 @@ import { } from "../session-row-projection.js"; import { buildHealthAgentSummaries, resolveHealthAgentOrder } from "./collector.js"; -afterEach(() => vi.restoreAllMocks()); +// Periodic WAL maintenance is independent of the request SQL budget. +beforeEach(() => vi.useFakeTimers({ toFake: ["setInterval", "clearInterval"] })); +afterEach(() => { + vi.useRealTimers(); + vi.restoreAllMocks(); +}); async function settleProjection(projection: SessionRowProjection) { do { diff --git a/src/gateway/server-agent-database-startup.test.ts b/src/gateway/server-agent-database-startup.test.ts index ba9fa0983d8b..ea526bab8c19 100644 --- a/src/gateway/server-agent-database-startup.test.ts +++ b/src/gateway/server-agent-database-startup.test.ts @@ -2,6 +2,7 @@ import fs from "node:fs"; import path from "node:path"; import { DatabaseSync } from "node:sqlite"; import { afterEach, expect, it, vi } from "vitest"; +import { withinTest } from "../../test/helpers/promise.js"; import { observeHostDataSql } from "../../test/helpers/sqlite-statement-execution-counter.js"; import { saveAuthProfileStore } from "../agents/auth-profiles.js"; import { listConfiguredOwnerInputs } from "../agents/prepared-model-runtime.configured.js"; @@ -43,9 +44,15 @@ import { testState } from "./test-helpers.runtime-state.js"; import { installGatewayTestHooks, startTestGatewayServer } from "./test-helpers.server.js"; installGatewayTestHooks(); -afterEach(() => { - vi.restoreAllMocks(); - vi.unstubAllEnvs(); +let pendingFixtureCleanup: Promise | undefined; +afterEach(async () => { + try { + await pendingFixtureCleanup; + } finally { + pendingFixtureCleanup = undefined; + vi.restoreAllMocks(); + vi.unstubAllEnvs(); + } }); function pauseIntegrityInspections(params: { @@ -117,7 +124,7 @@ DatabaseSync.prototype.prepare = function(sql) { }; } -it.each([ +it.for([ { outcome: "recover", agentId: "worker" }, { outcome: "corrupt", agentId: "worker" }, { outcome: "physical-corrupt", agentId: "main" }, @@ -130,7 +137,7 @@ it.each([ { outcome: "shutdown-preparation", agentId: "worker" }, ] as const)( "applies startup admission while $agentId follows its $outcome lifecycle", - async ({ outcome, agentId }) => { + async ({ outcome, agentId }, { signal }) => { const nativeBroker = process.platform === "linux" && !process.versions.bun; const brokerExpected = nativeBroker || @@ -400,10 +407,13 @@ it.each([ expect(() => process.kill(pid, 0)).toThrow(); } else if (outcome === "recover" || outcome === "superseded") { fs.writeFileSync(releasePath, "resume"); - await Promise.race([ - preparationEntered.promise, - hostJournalRead.promise.then(() => expect(hostJournalReads).toBe(0)), - ]); + await withinTest( + Promise.race([ + preparationEntered.promise, + hostJournalRead.promise.then(() => expect(hostJournalReads).toBe(0)), + ]), + signal, + ); expect(sessionPrepared).toBe(true); if (outcome === "recover") { expect(preparationParent).toBe(brokerExpected ? brokerPid : process.pid); @@ -468,20 +478,23 @@ it.each([ expect(readAgentDatabaseAdmissionRefusal(agentId, { env })).toBeUndefined(); } } finally { - preparationRelease.resolve(); - fs.writeFileSync(releasePath, "resume"); - if (pause) { - fs.writeFileSync(pause.preparationReleasePath, "resume"); - } - try { - await server?.close(); - } finally { - try { - await suppliedBroker?.close(); - } finally { - await unadoptedPortClaim?.release(); + pendingFixtureCleanup = (async () => { + preparationRelease.resolve(); + fs.writeFileSync(releasePath, "resume"); + if (pause) { + fs.writeFileSync(pause.preparationReleasePath, "resume"); } - } + try { + await server?.close(); + } finally { + try { + await suppliedBroker?.close(); + } finally { + await unadoptedPortClaim?.release(); + } + } + })(); + await pendingFixtureCleanup; } }, ); diff --git a/src/gateway/server-channels.test.ts b/src/gateway/server-channels.test.ts index 8213095ef3d3..337fb379e70f 100644 --- a/src/gateway/server-channels.test.ts +++ b/src/gateway/server-channels.test.ts @@ -241,7 +241,7 @@ describe("server-channels auto restart", () => { resetGatewayWorkAdmission(); previousRegistry = getActivePluginRegistry(); vi.useRealTimers(); - vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout", "Date"] }); + vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout", "Date", "performance"] }); hoisted.sleepWithAbort.mockClear(); hoisted.startChannelApprovalHandlerBootstrap.mockReset(); hoisted.startChannelApprovalHandlerBootstrap.mockResolvedValue(async () => {}); diff --git a/src/gateway/server-plugin-subagent-runtime.test.ts b/src/gateway/server-plugin-subagent-runtime.test.ts index 495086bd9233..67375e972392 100644 --- a/src/gateway/server-plugin-subagent-runtime.test.ts +++ b/src/gateway/server-plugin-subagent-runtime.test.ts @@ -1,5 +1,5 @@ import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; -import { createDeferred } from "../../test/helpers/promise.js"; +import { awaitGateBeforeSettlement, createDeferred } from "../../test/helpers/promise.js"; import { FailoverError } from "../agents/failover-error.js"; import type { runIsolatedCompletion } from "../agents/isolated-completion.js"; import { withGatewayToolCallerIdentity } from "../agents/tools/gateway-caller-context.js"; @@ -23,7 +23,9 @@ import { createBackgroundWorkOwner, getBackgroundWorkSnapshot, } from "../process/background-work.js"; +import * as commandQueue from "../process/command-queue.js"; import { resetCommandQueueStateForTest } from "../process/command-queue.test-support.js"; +import { CommandLane } from "../process/lanes.js"; import { ensureProfileForEmail } from "../state/user-profiles.js"; import { withOpenClawTestState } from "../test-utils/openclaw-test-state.js"; import { captureGatewayOperatorRunAuthority } from "./operator-run-authority.js"; @@ -385,6 +387,19 @@ describe("plugin background completions", () => { const profile = ensureProfileForEmail("completion-mutation@example.com"); config.agents!.entries!.main!.model = "test-provider/main-model@main-profile"; const blockers = blockBackgroundSlots(3); + const queued = createDeferred(); + const enqueue = commandQueue.enqueueCommandInLane; + vi.spyOn(commandQueue, "enqueueCommandInLane").mockImplementation((lane, task, options) => + enqueue(lane, task, { + ...options, + onQueued: () => { + options?.onQueued?.(); + if (lane === `${CommandLane.Background}:plugin:${PLUGIN_ID}`) { + queued.resolve(); + } + }, + }), + ); const runtime = createRuntime(); const request = { agentId: "main", message: "Review these notes" }; const result = withPluginRuntimeGatewayRequestScope( @@ -401,8 +416,12 @@ describe("plugin background completions", () => { (value) => ({ value }), (error: unknown) => ({ error }), ); - await vi.dynamicImportSettled(); - await vi.waitFor(() => expect(getBackgroundWorkSnapshot().queuedCount).toBe(1)); + await awaitGateBeforeSettlement( + queued.promise, + result, + "Completion settled before its queue admission", + ); + expect(getBackgroundWorkSnapshot().queuedCount).toBe(1); request.agentId = "research"; blockers.release(); await blockers.settled(); diff --git a/src/gateway/server-startup-background.test-support.ts b/src/gateway/server-startup-background.test-support.ts new file mode 100644 index 000000000000..24e26e6e8a69 --- /dev/null +++ b/src/gateway/server-startup-background.test-support.ts @@ -0,0 +1,11 @@ +import { vi } from "vitest"; + +// Post-attach orchestration keeps unrelated background discovery worker-free. +vi.mock("../agents/session-dirs.js", () => ({ + resolveAgentSessionDirs: vi.fn(async () => []), +})); + +vi.mock("./update-run-watcher.js", () => ({ + startUpdateRunWatcher: vi.fn(() => ({ stop: vi.fn(async () => {}) })), + wakeUpdateRunWatcher: vi.fn(), +})); diff --git a/src/gateway/server-startup-post-attach.test.ts b/src/gateway/server-startup-post-attach.test.ts index 35b32cc27566..58f02b51034d 100644 --- a/src/gateway/server-startup-post-attach.test.ts +++ b/src/gateway/server-startup-post-attach.test.ts @@ -2,6 +2,7 @@ * Gateway post-attach startup task tests. */ import "./server-worker-free.test-support.js"; +import "./server-startup-background.test-support.js"; import fs from "node:fs"; import { performance } from "node:perf_hooks"; import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; @@ -150,10 +151,6 @@ const hoisted = vi.hoisted(() => { }; }); -vi.mock("../agents/session-dirs.js", () => ({ - resolveAgentSessionDirs: vi.fn(async () => []), -})); - vi.mock("../agents/subagents/registry/subagent-registry.js", () => ({ activateSubagentRegistry: hoisted.activateSubagentRegistry, })); diff --git a/src/gateway/session-row-projection.category.test.ts b/src/gateway/session-row-projection.category.test.ts index b6d6e78f63f2..6d85ece68fb7 100644 --- a/src/gateway/session-row-projection.category.test.ts +++ b/src/gateway/session-row-projection.category.test.ts @@ -1,4 +1,4 @@ -import { afterEach, expect, it, vi } from "vitest"; +import { afterEach, beforeEach, expect, it, vi } from "vitest"; import { observeHostDataSql } from "../../test/helpers/sqlite-statement-execution-counter.js"; import { replaceSessionEntrySync } from "../config/sessions/session-accessor.js"; import * as history from "../config/sessions/session-transcript-worker-runtime.js"; @@ -12,7 +12,12 @@ import { createSessionRowProjection } from "./session-row-projection.js"; import { listProjectedSessions } from "./session-utils-list.js"; import { createWorkerSessionPlacementStore } from "./worker-environments/placement-store.js"; -afterEach(() => vi.restoreAllMocks()); +// Periodic WAL maintenance is independent of the request SQL budget. +beforeEach(() => vi.useFakeTimers({ toFake: ["setInterval", "clearInterval"] })); +afterEach(() => { + vi.useRealTimers(); + vi.restoreAllMocks(); +}); const cfg = { agents: { entries: { main: {} } } }; diff --git a/src/infra/update-managed-service-triage.test-support.ts b/src/infra/update-managed-service-triage.test-support.ts index fc85fb89f6ac..2d92235b7148 100644 --- a/src/infra/update-managed-service-triage.test-support.ts +++ b/src/infra/update-managed-service-triage.test-support.ts @@ -183,7 +183,7 @@ const action = args.find(x => ['show','start','stop','restart','reset-failed'].i const name = args[args.indexOf(action)+1]; const scope = JSON.parse(fs.readFileSync(scopeFile,'utf8')); let primary = JSON.parse(fs.readFileSync(primaryFile,'utf8')); -event(action, {name}); +event(action ?? 'probe', {name, args}); if (action === 'show') { // A stopped scope retains its cgroup until its registered processes have exited. const populated = !scope.active && name.endsWith('.scope') && fs.readdirSync(root+'/members').some(member => { diff --git a/src/plugins/plugin-lifecycle-lease.test.ts b/src/plugins/plugin-lifecycle-lease.test.ts index 8e8fdc410b61..51dd1967e70a 100644 --- a/src/plugins/plugin-lifecycle-lease.test.ts +++ b/src/plugins/plugin-lifecycle-lease.test.ts @@ -535,9 +535,21 @@ describe("plugin lifecycle lease", () => { const betaGoMarker = state.path("beta-go"); const releaseAlphaMarker = state.path("release-alpha"); // Both processes and their SQLite workers share the lease clock for this cache handoff. + // Bun's explicit worker env skips inherited preloads in these non-Vitest children. const clockPreload = await state.writeText( "lease-clock.cjs", - `Date.now = () => ${Date.now()};\n`, + `Date.now = () => ${Date.now()}; +if (process.versions.bun) { + const threads = require("node:worker_threads"); + const Worker = threads.Worker; + threads.Worker = class extends Worker { + constructor(url, options) { + super(url, { ...options, execArgv: [...(options?.execArgv ?? process.execArgv), "--preload", __filename] }); + } + }; + require("node:module").syncBuiltinESMExports(); +} +`, ); const childEnv = { ...process.env }; for (const [key, value] of Object.entries(sqliteWorkerPreloadEnv(clockPreload))) { diff --git a/src/process/terminal-pty-bun.test.ts b/src/process/terminal-pty-bun.test.ts index 0a128a89f465..11ff1de05055 100644 --- a/src/process/terminal-pty-bun.test.ts +++ b/src/process/terminal-pty-bun.test.ts @@ -261,33 +261,30 @@ describe.runIf(Boolean(process.versions.bun) && process.platform !== "win32" && expect(observed.output).toBe("READY\r\ntail 🦞\r\n"); }); - it("reports exit after kill while the consumer keeps re-pausing output", async () => { + it("reports exit after kill while the consumer keeps re-pausing output", async ({ + signal, + }) => { const cwd = tempDirs.make("openclaw-bun-pty-kill-"); - fs.writeFileSync(path.join(cwd, "payload"), "x".repeat(4 * 1024 * 1024)); - const { handle, observed } = await start( - [ - "-c", - 'stty -echo; printf "READY\\n"; read input; printf started > progress; cat payload', - ], - { cwd }, - ); - await vi.waitFor(() => expect(observed.output).toBe("READY\r\n"), deadline); - handle.pause(); - handle.write("go\r"); - await vi.waitFor( - () => expect(fs.readFileSync(path.join(cwd, "progress"), "utf8")).toBe("started"), - deadline, - ); + const payload = path.join(cwd, "payload"); + fs.writeFileSync(payload, "x".repeat(4 * 1024 * 1024)); + // A shell can exit with its killed child's status before the tree kill reaches it. + const { handle, observed, exited, waitForOutput } = await start([payload], { + file: "/bin/cat", + cwd, + }); // A viewer whose backlog stays full pauses again on every chunk it receives. handle.onData(() => handle.pause()); + await waitForOutput("x", signal); + expect(observed.exit).toBeUndefined(); + const beforeKill = observed.output.length; handle.kill(); - await vi.waitFor( - () => expect(observed.exit).toEqual({ exitCode: 0, signal: constants.signals.SIGKILL }), - deadline, - ); + expect(await withinTest(exited, signal)).toEqual({ + exitCode: 0, + signal: constants.signals.SIGKILL, + }); // Teardown delivered the dying tree's output before exit; nothing trails it. const atExit = observed.output.length; - expect(atExit).toBeGreaterThan("READY\r\n".length); + expect(atExit).toBeGreaterThan(beforeKill); handle.resume(); expect(observed.output.length).toBe(atExit); });