mirror of
https://github.com/openclaw/openclaw.git
synced 2026-09-28 14:12:28 +08:00
perf(infra): reuse SQLite readers in detached Gateway callbacks (#155779)
This commit is contained in:
@@ -116,6 +116,9 @@ and personal-account selection still run on each request. Isolated agent scopes
|
||||
and private database snapshots do not share this cache. Gateway cache misses reuse
|
||||
a read-only child whose lifetime ends at shutdown; each read reacquires its source
|
||||
admission and closes its SQLite handles before returning.
|
||||
Detached connection, cron, heartbeat, and hook callbacks retain that Gateway's
|
||||
read-only worker scope without inheriting startup or request authority. Shutdown
|
||||
refuses late callbacks before they can create another reader.
|
||||
Usage bookkeeping invalidates later cache reuse while admitted reads can finish
|
||||
their snapshots. Credential, selection, ownership, and lifecycle changes still
|
||||
invalidate in-flight preparation.
|
||||
|
||||
@@ -7,6 +7,7 @@
|
||||
* their caller and must keep it.
|
||||
*/
|
||||
import { withoutGatewayToolCallerIdentity } from "../agents/tools/gateway-caller-context.js";
|
||||
import { captureSqliteReadOnlyWorkerScope } from "../infra/sqlite-readonly-worker-context.js";
|
||||
import {
|
||||
bindGatewayContextResolver,
|
||||
withPluginRuntimeGatewayContextResolver,
|
||||
@@ -52,8 +53,13 @@ export function createScheduledGatewayRunner(
|
||||
resolveGatewayContext?: ScheduledGatewayContextResolver,
|
||||
) {
|
||||
const spawnBroker = getSpawnBroker();
|
||||
const runWithReadOnlyWorkers = captureSqliteReadOnlyWorkerScope();
|
||||
return <T>(run: () => Promise<T>): Promise<T> =>
|
||||
runWithScheduledGatewayContext({ resolveGatewayContext, spawnBroker, run });
|
||||
runWithScheduledGatewayContext({
|
||||
resolveGatewayContext,
|
||||
spawnBroker,
|
||||
run: () => runWithReadOnlyWorkers(run),
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import { captureSqliteReadOnlyWorkerScope } from "../infra/sqlite-readonly-worker-context.js";
|
||||
import { getSpawnBroker, runWithSpawnBroker } from "../process/spawn-broker/context.js";
|
||||
import { AsyncWorkScope } from "../shared/async-work-scope.js";
|
||||
import { createDeferredCore } from "../shared/deferred.js";
|
||||
@@ -5,12 +6,15 @@ import { createDeferredCore } from "../shared/deferred.js";
|
||||
/** Owns received work and connection cleanup until this Gateway generation settles. */
|
||||
export class GatewayConnectionWork extends AsyncWorkScope {
|
||||
private readonly spawnBroker = getSpawnBroker();
|
||||
private readonly runWithReadOnlyWorkers = captureSqliteReadOnlyWorkerScope();
|
||||
private readonly connections = new Set<() => void>();
|
||||
private failure: { error: unknown } | undefined;
|
||||
|
||||
override track<T>(run: () => T | Promise<T>): Promise<T> {
|
||||
// Socket callbacks retain this Gateway's transport without borrowing startup admission.
|
||||
return runWithSpawnBroker(this.spawnBroker, () => super.track(run));
|
||||
// Socket callbacks retain this Gateway's workers without borrowing startup admission.
|
||||
return runWithSpawnBroker(this.spawnBroker, () =>
|
||||
super.track(() => this.runWithReadOnlyWorkers(run)),
|
||||
);
|
||||
}
|
||||
|
||||
trackCleanup(run: () => Promise<void>): Promise<void> {
|
||||
|
||||
@@ -8,6 +8,10 @@ import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import { createMockCronStateForJobs } from "../cron/service.test-harness.js";
|
||||
import { listPage } from "../cron/service/ops-read.js";
|
||||
import type { CronJob } from "../cron/types.js";
|
||||
import {
|
||||
createSqliteReadOnlyWorkerScope,
|
||||
isSqliteInspectionDeadlineOwnedByCaller,
|
||||
} from "../infra/sqlite-readonly-worker.js";
|
||||
import { getSpawnBroker, runWithSpawnBroker } from "../process/spawn-broker/context.js";
|
||||
import { useSpawnBrokerTestFixture } from "../process/spawn-broker/host.test-support.js";
|
||||
import { runInDetachedAsyncContext } from "../shared/async-work-scope.js";
|
||||
@@ -65,22 +69,35 @@ describe("createLazyGatewayCronState", () => {
|
||||
const state = createCronState(cron);
|
||||
hoisted.setState(state);
|
||||
let observedBroker: unknown = "not-built";
|
||||
let observedReadOnlyScope = false;
|
||||
hoisted.buildGatewayCronService.mockImplementationOnce(() => {
|
||||
observedBroker = getSpawnBroker();
|
||||
observedReadOnlyScope = isSqliteInspectionDeadlineOwnedByCaller();
|
||||
return state;
|
||||
});
|
||||
|
||||
const lazy = runWithSpawnBroker(broker, () => createLazyGatewayCronState(createParams()));
|
||||
const readers = createSqliteReadOnlyWorkerScope({
|
||||
signal: new AbortController().signal,
|
||||
deadlineOwnedByCaller: true,
|
||||
});
|
||||
const lazy = readers.run(() =>
|
||||
runWithSpawnBroker(broker, () => createLazyGatewayCronState(createParams())),
|
||||
);
|
||||
|
||||
expect(hoisted.buildGatewayCronService).not.toHaveBeenCalled();
|
||||
expect(lazy.cron.getJob("demo")).toBeUndefined();
|
||||
expect(lazy.cron.getDefaultAgentId()).toBeUndefined();
|
||||
|
||||
await runInDetachedAsyncContext(() => lazy.cron.status());
|
||||
try {
|
||||
await runInDetachedAsyncContext(() => lazy.cron.status());
|
||||
} finally {
|
||||
await readers.close();
|
||||
}
|
||||
|
||||
expect(hoisted.buildGatewayCronService).toHaveBeenCalledTimes(1);
|
||||
expect(cron["status"]).toHaveBeenCalledTimes(1);
|
||||
expect(observedBroker === broker).toBe(true);
|
||||
expect(observedReadOnlyScope).toBe(true);
|
||||
});
|
||||
|
||||
it("loads the cron service for direct job reads", async () => {
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
import type { CliDeps } from "../cli/deps.types.js";
|
||||
import type { OpenClawConfig } from "../config/types.openclaw.js";
|
||||
import { resolveCronJobsStorePathFromConfig } from "../cron/store.js";
|
||||
import { captureSqliteReadOnlyWorkerScope } from "../infra/sqlite-readonly-worker-context.js";
|
||||
import { getSpawnBroker, runWithSpawnBroker } from "../process/spawn-broker/context.js";
|
||||
import { createLazyPromiseLoader } from "../shared/lazy-runtime.js";
|
||||
import type { GatewayCronServiceContract } from "./server-cron-contract.js";
|
||||
@@ -35,6 +36,7 @@ type LoadedGatewayCronState = {
|
||||
/** Creates a cron state proxy that imports the real cron service on first use. */
|
||||
export function createLazyGatewayCronState(params: LazyGatewayCronParams): GatewayCronState {
|
||||
const spawnBroker = getSpawnBroker();
|
||||
const runWithReadOnlyWorkers = captureSqliteReadOnlyWorkerScope();
|
||||
const env = params.env ?? process.env;
|
||||
const storePath = resolveCronJobsStorePathFromConfig(params.cfg, env);
|
||||
const cronEnabled = env.OPENCLAW_SKIP_CRON !== "1" && params.cfg.cron?.enabled !== false;
|
||||
@@ -64,7 +66,9 @@ export function createLazyGatewayCronState(params: LazyGatewayCronParams): Gatew
|
||||
() =>
|
||||
import("./server-cron.js").then(({ buildGatewayCronService }) => {
|
||||
loaded = {
|
||||
state: runWithSpawnBroker(spawnBroker, () => buildGatewayCronService(params)),
|
||||
state: runWithSpawnBroker(spawnBroker, () =>
|
||||
runWithReadOnlyWorkers(() => buildGatewayCronService(params)),
|
||||
),
|
||||
phase: "idle",
|
||||
startPromise: null,
|
||||
startGeneration: null,
|
||||
|
||||
@@ -8,6 +8,7 @@ import { resolveSandboxHostPort } from "../agents/sandbox-host.js";
|
||||
import { isCoreCanvasHostEnabled } from "../canvas/config.js";
|
||||
import { resolveCanvasNodeCapability } from "../canvas/constants.js";
|
||||
import type { CliDeps } from "../cli/deps.types.js";
|
||||
import { captureSqliteReadOnlyWorkerScope } from "../infra/sqlite-readonly-worker-context.js";
|
||||
import type { GatewayTlsRuntime } from "../infra/tls/gateway.js";
|
||||
import type { createSubsystemLogger } from "../logging/subsystem.js";
|
||||
import type { PluginRegistry } from "../plugins/registry.js";
|
||||
@@ -154,6 +155,7 @@ export async function createGatewayHttpTransport(params: {
|
||||
) => ReturnType<PluginRuntimeCore["hooks"]["dispatchHookAgentTurn"]>;
|
||||
}> {
|
||||
const spawnBroker = getSpawnBroker();
|
||||
const runWithReadOnlyWorkers = captureSqliteReadOnlyWorkerScope();
|
||||
if (params.testListener) {
|
||||
const address = params.testListener.address();
|
||||
if (
|
||||
@@ -178,13 +180,15 @@ export async function createGatewayHttpTransport(params: {
|
||||
const getHookDispatcher = async () => {
|
||||
const { createGatewayHookDispatcher } = await import("./server/hooks.js");
|
||||
return (loadedHookDispatcher ??= runWithSpawnBroker(spawnBroker, () =>
|
||||
createGatewayHookDispatcher({
|
||||
deps: params.deps,
|
||||
logHooks: params.logHooks,
|
||||
...(params.getGatewayRequestContext
|
||||
? { resolveGatewayContext: params.getGatewayRequestContext }
|
||||
: {}),
|
||||
}),
|
||||
runWithReadOnlyWorkers(() =>
|
||||
createGatewayHookDispatcher({
|
||||
deps: params.deps,
|
||||
logHooks: params.logHooks,
|
||||
...(params.getGatewayRequestContext
|
||||
? { resolveGatewayContext: params.getGatewayRequestContext }
|
||||
: {}),
|
||||
}),
|
||||
),
|
||||
));
|
||||
};
|
||||
const handleHooksRequest: HooksRequestHandler = async (req, res) => {
|
||||
|
||||
@@ -0,0 +1,40 @@
|
||||
import { AsyncLocalStorage } from "node:async_hooks";
|
||||
import type {
|
||||
SqliteAuthProfileReadOptions,
|
||||
SqliteReadOnlyWorkerValue,
|
||||
} from "./sqlite-readonly-worker-protocol.js";
|
||||
import type {
|
||||
createSqliteReadOnlyWorkerSession,
|
||||
SqliteReadOnlyWorkerLaunch,
|
||||
} from "./sqlite-readonly-worker-session.js";
|
||||
|
||||
export type SqliteReadOnlyWorkerScope = {
|
||||
active: boolean;
|
||||
busy: boolean;
|
||||
controller: AbortController;
|
||||
pending: Set<Promise<SqliteReadOnlyWorkerValue>>;
|
||||
deadlineOwnedByCaller: boolean;
|
||||
worker?: ReturnType<typeof createSqliteReadOnlyWorkerSession>;
|
||||
authWorker?: {
|
||||
source: SqliteAuthProfileReadOptions["source"];
|
||||
launch: SqliteReadOnlyWorkerLaunch;
|
||||
session: ReturnType<typeof createSqliteReadOnlyWorkerSession>;
|
||||
};
|
||||
authTail: Promise<void>;
|
||||
};
|
||||
export const readOnlyWorkerScope = new AsyncLocalStorage<SqliteReadOnlyWorkerScope>();
|
||||
|
||||
/** Carry the owning readers into callbacks without retaining startup or request authority. */
|
||||
export function captureSqliteReadOnlyWorkerScope(): <T>(operation: () => T) => T {
|
||||
const scope = readOnlyWorkerScope.getStore();
|
||||
return (operation) => {
|
||||
if (!scope) {
|
||||
return readOnlyWorkerScope.exit(operation);
|
||||
}
|
||||
if (!scope.active) {
|
||||
throw new Error("SQLite read-only worker scope closed");
|
||||
}
|
||||
scope.controller.signal.throwIfAborted();
|
||||
return readOnlyWorkerScope.run(scope, operation);
|
||||
};
|
||||
}
|
||||
@@ -1,10 +1,16 @@
|
||||
import { AsyncLocalStorage } from "node:async_hooks";
|
||||
import { EventEmitter } from "node:events";
|
||||
import { setImmediate as nextTurn } from "node:timers/promises";
|
||||
import { isRecord } from "@openclaw/normalization-core/record-coerce";
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||
import { emitChildProcessSpawnSample } from "../process/spawn-diagnostics.js";
|
||||
import { onDiagnosticEvent, setDiagnosticsEnabledForProcess } from "./diagnostic-events.js";
|
||||
import { encodeSqliteAuthTransferFrame } from "./sqlite-readonly-auth-transfer.js";
|
||||
import { captureSqliteReadOnlyWorkerScope } from "./sqlite-readonly-worker-context.js";
|
||||
import { createSqliteReadOnlyWorkerSession } from "./sqlite-readonly-worker-session.js";
|
||||
import {
|
||||
createSqliteReadOnlyWorkerScope,
|
||||
resolveSqliteInspectionSignal,
|
||||
createScopedSqliteReadOnlyWorker,
|
||||
withSqliteReadOnlyWorkerScope,
|
||||
} from "./sqlite-readonly-worker.js";
|
||||
@@ -253,3 +259,69 @@ it("keeps a detached staging command budget inside a caller-owned deadline scope
|
||||
timer.mockRestore();
|
||||
}
|
||||
});
|
||||
|
||||
it("carries only its captured read scope and refuses callbacks after owner retirement", async () => {
|
||||
const caller = new AsyncLocalStorage<string>();
|
||||
const withoutOwner = captureSqliteReadOnlyWorkerScope();
|
||||
const controller = new AbortController();
|
||||
const owner = createSqliteReadOnlyWorkerScope({
|
||||
signal: controller.signal,
|
||||
deadlineOwnedByCaller: false,
|
||||
});
|
||||
const other = createSqliteReadOnlyWorkerScope();
|
||||
const run = owner.run(() => caller.run("startup", captureSqliteReadOnlyWorkerScope));
|
||||
const signal = owner.run(() => resolveSqliteInspectionSignal());
|
||||
try {
|
||||
await other.run(() =>
|
||||
caller.run("request", () =>
|
||||
run(async () => {
|
||||
await Promise.resolve();
|
||||
expect(resolveSqliteInspectionSignal()).toBe(signal);
|
||||
expect(caller.getStore()).toBe("request");
|
||||
expect(withoutOwner(() => resolveSqliteInspectionSignal())).toBeUndefined();
|
||||
expect(resolveSqliteInspectionSignal()).toBe(signal);
|
||||
}),
|
||||
),
|
||||
);
|
||||
const late = vi.fn();
|
||||
const aborted = new Error("captured owner cancelled");
|
||||
controller.abort(aborted);
|
||||
expect(() => run(late)).toThrow(aborted);
|
||||
await owner.close();
|
||||
expect(() => other.run(() => run(late))).toThrow("scope closed");
|
||||
expect(late).not.toHaveBeenCalled();
|
||||
} finally {
|
||||
await Promise.all([owner.close(), other.close()]);
|
||||
}
|
||||
});
|
||||
|
||||
it("counts admitted read-only session children in node spawn diagnostics", () => {
|
||||
let now = 0;
|
||||
const clock = vi.spyOn(performance, "now").mockImplementation(() => now);
|
||||
const events: unknown[] = [];
|
||||
const stop = onDiagnosticEvent((event) => {
|
||||
if (event.type === "diagnostic.child_process.spawn") {
|
||||
events.push(event);
|
||||
}
|
||||
});
|
||||
try {
|
||||
setDiagnosticsEnabledForProcess(false);
|
||||
emitChildProcessSpawnSample();
|
||||
setDiagnosticsEnabledForProcess(true);
|
||||
const { child } = createSession();
|
||||
now = 60_000;
|
||||
emitChildProcessSpawnSample();
|
||||
expect(events).toEqual([]);
|
||||
child.emit("spawn");
|
||||
now = 120_000;
|
||||
emitChildProcessSpawnSample();
|
||||
expect(events).toEqual([
|
||||
expect.objectContaining({ family: process.versions.bun ? "other" : "node", count: 1 }),
|
||||
]);
|
||||
} finally {
|
||||
stop();
|
||||
setDiagnosticsEnabledForProcess(false);
|
||||
emitChildProcessSpawnSample();
|
||||
clock.mockRestore();
|
||||
}
|
||||
});
|
||||
|
||||
@@ -2,6 +2,7 @@ import { spawn, type ChildProcess, type SpawnOptions } from "node:child_process"
|
||||
import { sliceUtf16Safe } from "@openclaw/normalization-core/utf16-slice";
|
||||
import { BrokerChild } from "../process/spawn-broker/child.js";
|
||||
import type { SpawnBrokerHost } from "../process/spawn-broker/host.js";
|
||||
import { recordChildProcessSpawn } from "../process/spawn-diagnostics.js";
|
||||
import { createSqliteAuthTransferReceiver } from "./sqlite-readonly-auth-transfer.js";
|
||||
import { retainSnapshotWork } from "./sqlite-readonly-location-cleanup.js";
|
||||
import {
|
||||
@@ -79,6 +80,7 @@ export function createSqliteReadOnlyWorkerSession(
|
||||
transport.kind === "broker"
|
||||
? transport.owner.spawn(process.execPath, argv, spawnOptions)
|
||||
: spawn(process.execPath, argv, spawnOptions);
|
||||
recordChildProcessSpawn(process.execPath, child);
|
||||
let retired = false;
|
||||
let sequence = 0;
|
||||
let stderr = "";
|
||||
|
||||
@@ -6,10 +6,13 @@ import { DatabaseSync, StatementSync } from "node:sqlite";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
|
||||
import { prepareAgentAuthProfileRowsRead } from "../agents/auth-profiles/sqlite-read.js";
|
||||
import { createScheduledGatewayRunner } from "../gateway/scheduled-run-gateway-context.js";
|
||||
import { GatewayConnectionWork } from "../gateway/server-connection-work.js";
|
||||
import { BrokerChild } from "../process/spawn-broker/child.js";
|
||||
import { runWithSpawnBroker } from "../process/spawn-broker/context.js";
|
||||
import { createSpawnBrokerHost } from "../process/spawn-broker/host.js";
|
||||
import { SpawnBrokerError } from "../process/spawn-broker/protocol.js";
|
||||
import { runInDetachedAsyncContext } from "../shared/async-work-scope.js";
|
||||
import { createDeferredCore } from "../shared/deferred.js";
|
||||
import { requireNodeSqlite } from "./node-sqlite.js";
|
||||
import { SQLITE_READONLY_CHILD_ARG } from "./runtime-process-entrypoints.js";
|
||||
@@ -166,6 +169,25 @@ describe.each([
|
||||
status: "readable",
|
||||
raw: { lastGood: {} },
|
||||
});
|
||||
const connection = new GatewayConnectionWork();
|
||||
const scheduled = createScheduledGatewayRunner();
|
||||
try {
|
||||
for (const enter of [
|
||||
<T>(run: () => Promise<T>) => connection.track(run),
|
||||
scheduled,
|
||||
]) {
|
||||
const detachedRows = await runInDetachedAsyncContext(() =>
|
||||
enter(() => read(source, transport.sourceKind)),
|
||||
);
|
||||
expect(detachedRows).toEqual({
|
||||
store: { status: "readable", raw: store },
|
||||
state: { status: "readable", raw: { lastGood: {} } },
|
||||
cacheable: true,
|
||||
});
|
||||
}
|
||||
} finally {
|
||||
await connection.drain();
|
||||
}
|
||||
const child =
|
||||
transport.label === "broker"
|
||||
? brokerSpawn?.mock.results[0]?.value
|
||||
|
||||
@@ -1,4 +1,3 @@
|
||||
import { AsyncLocalStorage } from "node:async_hooks";
|
||||
import { execFile, spawnSync } from "node:child_process";
|
||||
import fs from "node:fs";
|
||||
import path from "node:path";
|
||||
@@ -14,6 +13,10 @@ import {
|
||||
} from "./runtime-process-entrypoints.js";
|
||||
import { resolveRuntimeWorkerArgv, resolveRuntimeWorkerUrl } from "./runtime-worker-url.js";
|
||||
import { retainSnapshotWork } from "./sqlite-readonly-location-cleanup.js";
|
||||
import {
|
||||
readOnlyWorkerScope,
|
||||
type SqliteReadOnlyWorkerScope,
|
||||
} from "./sqlite-readonly-worker-context.js";
|
||||
import {
|
||||
SQLITE_READONLY_WORKER_MAX_BUFFER,
|
||||
readSqliteReadOnlyWorkerValue,
|
||||
@@ -118,22 +121,6 @@ export function sqliteInspectionTimeoutError(
|
||||
);
|
||||
}
|
||||
|
||||
type SqliteReadOnlyWorkerScope = {
|
||||
active: boolean;
|
||||
busy: boolean;
|
||||
controller: AbortController;
|
||||
pending: Set<Promise<SqliteReadOnlyWorkerValue>>;
|
||||
deadlineOwnedByCaller: boolean;
|
||||
worker?: ReturnType<typeof createScopedSqliteReadOnlyWorker>;
|
||||
authWorker?: {
|
||||
source: SqliteAuthProfileReadOptions["source"];
|
||||
launch: SqliteReadOnlyWorkerLaunch;
|
||||
session: ReturnType<typeof createScopedSqliteReadOnlyWorker>;
|
||||
};
|
||||
authTail: Promise<void>;
|
||||
};
|
||||
const readOnlyWorkerScope = new AsyncLocalStorage<SqliteReadOnlyWorkerScope>();
|
||||
|
||||
/** Reuse child imports until the lifecycle owner closes; reads reacquire source admission. */
|
||||
export function createSqliteReadOnlyWorkerScope(options?: {
|
||||
signal: AbortSignal;
|
||||
|
||||
@@ -11,7 +11,7 @@ import {
|
||||
type DiagnosticPhaseSnapshot,
|
||||
type DiagnosticLivenessWarningReason,
|
||||
} from "../infra/diagnostic-events.js";
|
||||
import { emitChildProcessSpawnSample } from "../process/spawn-utils.js";
|
||||
import { emitChildProcessSpawnSample } from "../process/spawn-diagnostics.js";
|
||||
import { createLazyRuntimeModule } from "../shared/lazy-runtime.js";
|
||||
import { reconcileDiagnosticGcObserver, stopDiagnosticGcObserver } from "./diagnostic-gc.js";
|
||||
import {
|
||||
|
||||
@@ -21,7 +21,7 @@ import {
|
||||
type CommandSubprocess,
|
||||
} from "./spawn-broker/execa-client.js";
|
||||
import type { CommandSpawnOptions } from "./spawn-broker/execa-types.js";
|
||||
import { recordChildProcessSpawn } from "./spawn-utils.js";
|
||||
import { recordChildProcessSpawn } from "./spawn-diagnostics.js";
|
||||
import { resolveSafeChildProcessInvocation } from "./windows-command.js";
|
||||
|
||||
export const COMMAND_PROCESS_TREE_KILL_GRACE_MS = 300;
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
import type { ChildProcess } from "node:child_process";
|
||||
import path from "node:path";
|
||||
import {
|
||||
areDiagnosticsEnabledForProcess,
|
||||
emitInternalDiagnosticEvent,
|
||||
} from "../infra/diagnostic-events.js";
|
||||
import { createSubsystemLogger } from "../logging/subsystem.js";
|
||||
import { resolveGlobalSingleton } from "../shared/global-singleton.js";
|
||||
|
||||
const spawnCounts = resolveGlobalSingleton(Symbol.for("openclaw.childProcessSpawnCounts"), () => ({
|
||||
counts: new Map<string, number>(),
|
||||
sampledAt: performance.now(),
|
||||
}));
|
||||
let spawnLog: ReturnType<typeof createSubsystemLogger> | undefined;
|
||||
const COMMAND_FAMILIES =
|
||||
/^(node|git|ps|pgrep|lsof|sh|bash|zsh|cmd|powershell|pwsh|npm|pnpm|python|python3|uv|ssh|openclaw)$/;
|
||||
|
||||
/** Count admitted local and broker launches, without recording paths or arguments. */
|
||||
export function recordChildProcessSpawn(command: string, child: ChildProcess): void {
|
||||
if (!areDiagnosticsEnabledForProcess()) {
|
||||
return;
|
||||
}
|
||||
const name = path.win32
|
||||
.basename(command)
|
||||
.toLowerCase()
|
||||
.replace(/\.(exe|cmd)$/, "");
|
||||
const family = COMMAND_FAMILIES.test(name) ? name : "other";
|
||||
child.once("spawn", () => {
|
||||
if (areDiagnosticsEnabledForProcess()) {
|
||||
spawnCounts.counts.set(family, (spawnCounts.counts.get(family) ?? 0) + 1);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/** The existing diagnostics heartbeat owns sampling; rates use actual elapsed time. */
|
||||
export function emitChildProcessSpawnSample(): void {
|
||||
const now = performance.now();
|
||||
if (!areDiagnosticsEnabledForProcess()) {
|
||||
spawnCounts.counts.clear();
|
||||
spawnCounts.sampledAt = now;
|
||||
return;
|
||||
}
|
||||
const intervalMs = now - spawnCounts.sampledAt;
|
||||
if (intervalMs < 60_000) {
|
||||
return;
|
||||
}
|
||||
for (const [family, count] of spawnCounts.counts) {
|
||||
emitInternalDiagnosticEvent({
|
||||
type: "diagnostic.child_process.spawn",
|
||||
family,
|
||||
count,
|
||||
intervalMs,
|
||||
});
|
||||
(spawnLog ??= createSubsystemLogger("gateway/diagnostics/process")).debug(
|
||||
`child process spawns: family=${family} count=${count} ratePerMinute=${((count * 60_000) / intervalMs).toFixed(2)}`,
|
||||
);
|
||||
}
|
||||
spawnCounts.counts.clear();
|
||||
spawnCounts.sampledAt = now;
|
||||
}
|
||||
@@ -13,12 +13,8 @@ import { withTempDir } from "../test-utils/temp-dir.js";
|
||||
import { spawnCommand } from "./exec-spawn.js";
|
||||
import { runWithSpawnBroker } from "./spawn-broker/context.js";
|
||||
import { createSpawnBrokerHost } from "./spawn-broker/host.js";
|
||||
import {
|
||||
emitChildProcessSpawnSample,
|
||||
recordChildProcessSpawn,
|
||||
spawnProcess,
|
||||
spawnWithFallback,
|
||||
} from "./spawn-utils.js";
|
||||
import { emitChildProcessSpawnSample, recordChildProcessSpawn } from "./spawn-diagnostics.js";
|
||||
import { spawnProcess, spawnWithFallback } from "./spawn-utils.js";
|
||||
|
||||
type SpawnImplementation = NonNullable<Parameters<typeof spawnWithFallback>[0]["spawnImpl"]>;
|
||||
|
||||
|
||||
@@ -1,69 +1,11 @@
|
||||
import type { ChildProcess, SpawnOptions } from "node:child_process";
|
||||
import { spawn } from "node:child_process";
|
||||
import { once } from "node:events";
|
||||
import path from "node:path";
|
||||
import { expectDefined } from "@openclaw/normalization-core";
|
||||
import {
|
||||
areDiagnosticsEnabledForProcess,
|
||||
emitInternalDiagnosticEvent,
|
||||
} from "../infra/diagnostic-events.js";
|
||||
import { toErrorObject } from "../infra/errors.js";
|
||||
import { createSubsystemLogger } from "../logging/subsystem.js";
|
||||
import { resolveGlobalSingleton } from "../shared/global-singleton.js";
|
||||
import { getSpawnBroker } from "./spawn-broker/context.js";
|
||||
import { brokerSpawnOptions } from "./spawn-broker/host.js";
|
||||
|
||||
const spawnCounts = resolveGlobalSingleton(Symbol.for("openclaw.childProcessSpawnCounts"), () => ({
|
||||
counts: new Map<string, number>(),
|
||||
sampledAt: performance.now(),
|
||||
}));
|
||||
let spawnLog: ReturnType<typeof createSubsystemLogger> | undefined;
|
||||
const COMMAND_FAMILIES =
|
||||
/^(node|git|ps|pgrep|lsof|sh|bash|zsh|cmd|powershell|pwsh|npm|pnpm|python|python3|uv|ssh|openclaw)$/;
|
||||
|
||||
/** Count admitted local and broker launches, without recording paths or arguments. */
|
||||
export function recordChildProcessSpawn(command: string, child: ChildProcess): void {
|
||||
if (!areDiagnosticsEnabledForProcess()) {
|
||||
return;
|
||||
}
|
||||
const name = path.win32
|
||||
.basename(command)
|
||||
.toLowerCase()
|
||||
.replace(/\.(exe|cmd)$/, "");
|
||||
const family = COMMAND_FAMILIES.test(name) ? name : "other";
|
||||
child.once("spawn", () => {
|
||||
if (areDiagnosticsEnabledForProcess()) {
|
||||
spawnCounts.counts.set(family, (spawnCounts.counts.get(family) ?? 0) + 1);
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/** The existing diagnostics heartbeat owns sampling; rates use actual elapsed time. */
|
||||
export function emitChildProcessSpawnSample(): void {
|
||||
const now = performance.now();
|
||||
if (!areDiagnosticsEnabledForProcess()) {
|
||||
spawnCounts.counts.clear();
|
||||
spawnCounts.sampledAt = now;
|
||||
return;
|
||||
}
|
||||
const intervalMs = now - spawnCounts.sampledAt;
|
||||
if (intervalMs < 60_000) {
|
||||
return;
|
||||
}
|
||||
for (const [family, count] of spawnCounts.counts) {
|
||||
emitInternalDiagnosticEvent({
|
||||
type: "diagnostic.child_process.spawn",
|
||||
family,
|
||||
count,
|
||||
intervalMs,
|
||||
});
|
||||
(spawnLog ??= createSubsystemLogger("gateway/diagnostics/process")).debug(
|
||||
`child process spawns: family=${family} count=${count} ratePerMinute=${((count * 60_000) / intervalMs).toFixed(2)}`,
|
||||
);
|
||||
}
|
||||
spawnCounts.counts.clear();
|
||||
spawnCounts.sampledAt = now;
|
||||
}
|
||||
import { recordChildProcessSpawn } from "./spawn-diagnostics.js";
|
||||
|
||||
/** Select the process-scoped native spawn transport without changing launch options. */
|
||||
export function spawnProcess(command: string, args: string[], options: SpawnOptions): ChildProcess {
|
||||
|
||||
Reference in New Issue
Block a user