mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 09:39:25 +00:00
fix: resolve baseline-v3 Bun regressions (#163051)
* fix(auto-reply): keep cancellation out of transcript worker data * test(sessions): preserve database custody during fixture cleanup * test: synchronize admission fixtures with owned completion signals * test(gateway): isolate maintenance and scheduler clocks * test(gateway): keep background startup outside worker-free fixtures * test(gateway): join startup fixture cleanup after cancellation * test(update): publish valid receipts for systemctl probes * test(plugins): share the lease clock with subprocess workers * test(process): observe PTY termination on the signaled process * fix(cloudflare): include bundled plugin artwork
This commit is contained in:
parent
e3c763272c
commit
f9940f1dff
18 changed files with 254 additions and 126 deletions
32
extensions/cloudflare/assets/activity.svg
Normal file
32
extensions/cloudflare/assets/activity.svg
Normal file
|
|
@ -0,0 +1,32 @@
|
|||
<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 24 24">
|
||||
<!--
|
||||
Adapted from OpenClaw ui/public/provider-icons/ProviderIcon-cloudflare.svg.
|
||||
Source: @lobehub/icons-static-svg@1.94.0, cloudflare.svg
|
||||
https://www.npmjs.com/package/@lobehub/icons-static-svg/v/1.94.0
|
||||
https://github.com/lobehub/lobe-icons
|
||||
|
||||
MIT License
|
||||
Copyright (c) 2023 LobeHub
|
||||
|
||||
Permission is hereby granted, free of charge, to any person obtaining a copy
|
||||
of this software and associated documentation files (the "Software"), to deal
|
||||
in the Software without restriction, including without limitation the rights
|
||||
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
|
||||
copies of the Software, and to permit persons to whom the Software is
|
||||
furnished to do so, subject to the following conditions:
|
||||
|
||||
The above copyright notice and this permission notice shall be included in all
|
||||
copies or substantial portions of the Software.
|
||||
|
||||
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
|
||||
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
|
||||
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
|
||||
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
|
||||
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
|
||||
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
|
||||
SOFTWARE.
|
||||
|
||||
Names and marks remain the property of their respective owners.
|
||||
-->
|
||||
<path fill="currentColor" fill-rule="evenodd" stroke="none" d="M16.493 17.4c.135-.52.08-.983-.161-1.338-.215-.328-.592-.519-1.05-.519l-8.663-.109a.148.148 0 01-.135-.082c-.027-.054-.027-.109-.027-.163.027-.082.108-.164.189-.164l8.744-.11c1.05-.054 2.153-.9 2.556-1.937l.511-1.31c.027-.055.027-.11.027-.164C17.92 8.91 15.66 7 12.942 7c-2.503 0-4.628 1.638-5.381 3.903a2.432 2.432 0 00-1.803-.491c-1.21.109-2.153 1.092-2.287 2.32-.027.328 0 .628.054.9C1.56 13.688 0 15.326 0 17.319c0 .19.027.355.027.545 0 .082.08.137.161.137h15.983c.08 0 .188-.055.215-.164l.107-.437 M19.238 11.75h-.242c-.054 0-.108.054-.135.109l-.35 1.2c-.134.52-.08.983.162 1.338.215.328.592.518 1.05.518l1.855.11c.054 0 .108.027.135.082.027.054.027.109.027.163-.027.082-.108.164-.188.164l-1.91.11c-1.05.054-2.153.9-2.557 1.937l-.134.355c-.027.055.026.137.107.137h6.592c.081 0 .162-.055.162-.137.107-.41.188-.846.188-1.31-.027-2.62-2.153-4.777-4.762-4.777" />
|
||||
</svg>
|
||||
|
After Width: | Height: | Size: 2.3 KiB |
BIN
extensions/cloudflare/assets/icon.png
Normal file
BIN
extensions/cloudflare/assets/icon.png
Normal file
Binary file not shown.
|
After Width: | Height: | Size: 22 KiB |
|
|
@ -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<void>();
|
||||
const onMessage = vi.fn(async () => delivered.resolve());
|
||||
const authorizing = createDeferred<void>();
|
||||
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 () => {
|
||||
|
|
|
|||
|
|
@ -489,12 +489,11 @@ async function estimateProviderPromptTokens(
|
|||
: undefined;
|
||||
}
|
||||
|
||||
async function estimatePromptTokensFromSessionTranscript(params: {
|
||||
agentId?: string;
|
||||
async function estimatePromptTokensFromSessionTranscript({
|
||||
abortSignal,
|
||||
...params
|
||||
}: Parameters<typeof readPreflightTranscriptContextMessages>[0] & {
|
||||
abortSignal?: AbortSignal;
|
||||
sessionId?: string;
|
||||
sessionKey?: string;
|
||||
storePath?: string;
|
||||
contextWindowTokens: number;
|
||||
}): Promise<TranscriptTokenEstimate | undefined> {
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 });
|
||||
}
|
||||
});
|
||||
},
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
|
|
|
|||
|
|
@ -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))
|
||||
|
|
|
|||
|
|
@ -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({
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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<void> | 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;
|
||||
}
|
||||
},
|
||||
);
|
||||
|
|
|
|||
|
|
@ -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 () => {});
|
||||
|
|
|
|||
|
|
@ -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();
|
||||
|
|
|
|||
11
src/gateway/server-startup-background.test-support.ts
Normal file
11
src/gateway/server-startup-background.test-support.ts
Normal file
|
|
@ -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(),
|
||||
}));
|
||||
|
|
@ -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,
|
||||
}));
|
||||
|
|
|
|||
|
|
@ -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: {} } } };
|
||||
|
||||
|
|
|
|||
|
|
@ -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 => {
|
||||
|
|
|
|||
|
|
@ -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))) {
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
});
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue