refactor(agents): generalize provider local service leases

This commit is contained in:
Vincent Koc
2026-07-10 21:13:12 -07:00
committed by Vincent Koc
parent 20fe3f46db
commit 4228ae93c8
2 changed files with 237 additions and 87 deletions
+109 -74
View File
@@ -11,6 +11,8 @@ import { killPidIfAlive, readPidFile, waitForPidToExit } from "../test-utils/pro
import {
attachModelProviderLocalService,
ensureModelProviderLocalService,
ensureProviderLocalService,
getManagedProviderLocalServiceDiagnosticsForTest,
getModelProviderLocalService,
hasLocalServiceProcessExited,
stopManagedProviderLocalServicesForTest,
@@ -267,11 +269,23 @@ describe("provider local service", () => {
}
});
it("serializes concurrent cold starts for the same local service", async () => {
it("serializes concurrent chat and embedding starts with independent leases", async () => {
const port = await freePort();
const healthUrl = `http://127.0.0.1:${port}/v1/models`;
const tempDir = await fs.mkdtemp(path.join(os.tmpdir(), "openclaw-local-service-"));
const startsPath = path.join(tempDir, "starts.txt");
const service = {
command: process.execPath,
args: [
"-e",
`const fs=require("node:fs");const http=require("node:http");fs.appendFileSync(${JSON.stringify(
startsPath,
)},"start\\n");setTimeout(()=>{const server=http.createServer((req,res)=>{res.writeHead(200,{"content-type":"application/json"});res.end('{"ok":true}');});server.listen(${port},"127.0.0.1");process.on("SIGTERM",()=>server.close(()=>process.exit(0)));},100);`,
],
healthUrl,
readyTimeoutMs: 5_000,
idleStopMs: 1,
};
const model = attachModelProviderLocalService(
{
id: "demo",
@@ -279,31 +293,25 @@ describe("provider local service", () => {
api: "openai-completions",
baseUrl: `http://127.0.0.1:${port}/v1`,
} as unknown as Model<"openai-completions">,
{
command: process.execPath,
args: [
"-e",
`const fs=require("node:fs");const http=require("node:http");fs.appendFileSync(${JSON.stringify(
startsPath,
)},"start\\n");setTimeout(()=>{const server=http.createServer((req,res)=>{res.writeHead(200,{"content-type":"application/json"});res.end('{"ok":true}');});server.listen(${port},"127.0.0.1");process.on("SIGTERM",()=>server.close(()=>process.exit(0)));},100);`,
],
healthUrl,
readyTimeoutMs: 5_000,
idleStopMs: 1,
},
service,
);
try {
const leases = await Promise.all([
ensureModelProviderLocalService(model),
const [chatLease, embeddingLease] = await Promise.all([
ensureModelProviderLocalService(model),
ensureProviderLocalService({
providerId: "local-concurrent",
baseUrl: `http://127.0.0.1:${port}/v1`,
service,
}),
]);
expect(leases).toHaveLength(2);
expect(chatLease).toBeDefined();
expect(embeddingLease).toBeDefined();
expect((await fetch(healthUrl)).ok).toBe(true);
for (const lease of leases) {
lease?.release();
}
embeddingLease?.release();
expect((await fetch(healthUrl)).ok).toBe(true);
chatLease?.release();
await waitForProbeFailure(healthUrl);
const starts = (await fs.readFile(startsPath, "utf8")).trim().split("\n");
expect(starts).toHaveLength(1);
@@ -312,7 +320,7 @@ describe("provider local service", () => {
}
});
it("does not reuse a local service with different env and derived health endpoint", async () => {
it("keeps configured provider aliases on different local endpoints independent", async () => {
const firstPort = await freePort();
const secondPort = await freePort();
const firstHealthUrl = `http://127.0.0.1:${firstPort}/v1/models`;
@@ -323,41 +331,33 @@ describe("provider local service", () => {
"-e",
`const fs=require("node:fs");const http=require("node:http");fs.appendFileSync(process.env.STARTS_PATH,process.env.LOCAL_SERVICE_PORT+"\\n");const server=http.createServer((req,res)=>{res.writeHead(200,{"content-type":"application/json"});res.end('{"ok":true}');});server.listen(Number(process.env.LOCAL_SERVICE_PORT),"127.0.0.1");process.on("SIGTERM",()=>server.close(()=>process.exit(0)));`,
];
const firstModel = attachModelProviderLocalService(
{
id: "demo",
provider: "local-key",
api: "openai-completions",
baseUrl: `http://127.0.0.1:${firstPort}/v1`,
} as unknown as Model<"openai-completions">,
{
command: process.execPath,
args,
env: { LOCAL_SERVICE_PORT: String(firstPort), STARTS_PATH: startsPath },
readyTimeoutMs: 5_000,
idleStopMs: 1,
},
);
const secondModel = attachModelProviderLocalService(
{
id: "demo",
provider: "local-key",
api: "openai-completions",
baseUrl: `http://127.0.0.1:${secondPort}/v1`,
} as unknown as Model<"openai-completions">,
{
command: process.execPath,
args,
env: { LOCAL_SERVICE_PORT: String(secondPort), STARTS_PATH: startsPath },
readyTimeoutMs: 5_000,
idleStopMs: 1,
},
);
const firstService = {
command: process.execPath,
args,
env: { LOCAL_SERVICE_PORT: String(firstPort), STARTS_PATH: startsPath },
readyTimeoutMs: 5_000,
idleStopMs: 1,
};
const secondService = {
command: process.execPath,
args,
env: { LOCAL_SERVICE_PORT: String(secondPort), STARTS_PATH: startsPath },
readyTimeoutMs: 5_000,
idleStopMs: 1,
};
try {
const leases = await Promise.all([
ensureModelProviderLocalService(firstModel),
ensureModelProviderLocalService(secondModel),
ensureProviderLocalService({
providerId: "ollama-spark",
baseUrl: `http://127.0.0.1:${firstPort}/v1`,
service: firstService,
}),
ensureProviderLocalService({
providerId: "ollama-studio",
baseUrl: `http://127.0.0.1:${secondPort}/v1`,
service: secondService,
}),
]);
expect((await fetch(firstHealthUrl)).ok).toBe(true);
@@ -438,22 +438,18 @@ describe("provider local service", () => {
it("reports a local service startup exit without waiting for readiness timeout", async () => {
const port = await freePort();
const model = attachModelProviderLocalService(
{
id: "demo",
provider: "local-fast-exit",
api: "openai-completions",
baseUrl: `http://127.0.0.1:${port}/v1`,
} as unknown as Model<"openai-completions">,
{
const target = {
providerId: "local-fast-exit",
baseUrl: `http://127.0.0.1:${port}/v1`,
service: {
command: process.execPath,
args: ["-e", "process.exit(17)"],
readyTimeoutMs: 60_000,
},
);
};
const startedAt = Date.now();
await expect(ensureModelProviderLocalService(model)).rejects.toThrow(
await expect(ensureProviderLocalService(target)).rejects.toThrow(
"local-fast-exit local service exited before readiness with code 17",
);
expect(Date.now() - startedAt).toBeLessThan(5_000);
@@ -486,14 +482,10 @@ describe("provider local service", () => {
const port = await freePort();
const healthUrl = `http://127.0.0.1:${port}/v1/models`;
const controller = new AbortController();
const model = attachModelProviderLocalService(
{
id: "demo",
provider: "local-abort",
api: "openai-completions",
baseUrl: `http://127.0.0.1:${port}/v1`,
} as unknown as Model<"openai-completions">,
{
const target = {
providerId: "local-abort",
baseUrl: `http://127.0.0.1:${port}/v1`,
service: {
command: process.execPath,
args: [
"-e",
@@ -503,16 +495,59 @@ describe("provider local service", () => {
readyTimeoutMs: 60_000,
idleStopMs: 1,
},
);
};
const startedAt = Date.now();
const abortTimer = setTimeout(() => controller.abort(new Error("request aborted")), 100);
abortTimer.unref?.();
await expect(
ensureModelProviderLocalService(model, undefined, controller.signal),
).rejects.toThrow("request aborted");
await expect(ensureProviderLocalService(target, controller.signal)).rejects.toThrow(
"request aborted",
);
expect(Date.now() - startedAt).toBeLessThan(5_000);
await waitForProbeFailure(healthUrl);
});
it("retains only bounded redacted startup diagnostics", async () => {
const port = await freePort();
const healthUrl = `http://127.0.0.1:${port}/v1/models`;
const diagnosticSecret = "local-service-diagnostic-secret";
const lease = await ensureProviderLocalService({
providerId: "local-diagnostics",
baseUrl: `http://127.0.0.1:${port}/v1`,
service: {
command: process.execPath,
args: [
"-e",
`const http=require("node:http");const noise="x".repeat(9000);process.stdout.write(noise+" "+process.env.DIAGNOSTIC_SECRET);process.stderr.write(noise+" "+process.env.DIAGNOSTIC_SECRET);const server=http.createServer((req,res)=>{res.writeHead(200);res.end("ok");});server.listen(${port},"127.0.0.1");process.on("SIGTERM",()=>server.close(()=>process.exit(0)));`,
],
env: { DIAGNOSTIC_SECRET: diagnosticSecret },
healthUrl,
readyTimeoutMs: 5_000,
idleStopMs: 1,
},
});
try {
const [diagnostics] = getManagedProviderLocalServiceDiagnosticsForTest();
expect(diagnostics).toMatchObject({
providerId: "local-diagnostics",
healthUrl,
pid: expect.any(Number),
startedAt: expect.any(Number),
spawnedAt: expect.any(Number),
readyAt: expect.any(Number),
lastHealthyAt: expect.any(Number),
});
expect(Buffer.byteLength(diagnostics?.stdoutTail ?? "")).toBeLessThanOrEqual(8 * 1024);
expect(Buffer.byteLength(diagnostics?.stderrTail ?? "")).toBeLessThanOrEqual(8 * 1024);
expect(diagnostics?.stdoutTail).not.toContain(diagnosticSecret);
expect(diagnostics?.stderrTail).not.toContain(diagnosticSecret);
expect(diagnostics?.stdoutTail).toContain("[redacted]");
expect(diagnostics?.stderrTail).toContain("[redacted]");
} finally {
lease?.release();
await waitForProbeFailure(healthUrl);
}
});
});
+128 -13
View File
@@ -11,6 +11,7 @@ import {
import type { ModelProviderLocalServiceConfig } from "../config/types.models.js";
import { toErrorObject } from "../infra/errors.js";
import type { Model } from "../llm/types.js";
import { redactSensitiveText } from "../logging/redact.js";
import { createSubsystemLogger } from "../logging/subsystem.js";
import {
forceKillChildProcessTree,
@@ -23,6 +24,7 @@ const log = createSubsystemLogger("provider-local-service");
const DEFAULT_READY_TIMEOUT_MS = 120_000;
const DEFAULT_PROBE_TIMEOUT_MS = 2_000;
const PROBE_INTERVAL_MS = 250;
const LOCAL_SERVICE_OUTPUT_TAIL_MAX_BYTES = 8 * 1024;
const MODEL_PROVIDER_LOCAL_SERVICE_SYMBOL = Symbol.for("openclaw.modelProviderLocalService");
@@ -37,6 +39,7 @@ type ManagedLocalService = {
active: number;
idleTimer?: NodeJS.Timeout;
lastExit?: LocalServiceExit;
diagnostics?: LocalServiceDiagnostics;
};
const services = new Map<string, ManagedLocalService>();
@@ -47,6 +50,27 @@ type LocalServiceExit = {
signal: NodeJS.Signals | null;
};
type LocalServiceDiagnostics = {
providerId: string;
healthUrl: string;
pid?: number;
startedAt: number;
spawnedAt?: number;
readyAt?: number;
lastHealthyAt?: number;
stdoutTail: string;
stderrTail: string;
lastExit?: LocalServiceExit;
};
/** Exact provider endpoint whose optional local process should be leased. */
export type ProviderLocalServiceTarget = {
providerId: string;
baseUrl: string;
headers?: HeadersInit;
service?: ModelProviderLocalServiceConfig;
};
/** Lease returned for a started or already-running local provider service. */
export type ProviderLocalServiceLease = {
release: () => void;
@@ -79,15 +103,32 @@ export async function ensureModelProviderLocalService(
signal?: AbortSignal | null,
): Promise<ProviderLocalServiceLease | undefined> {
const service = getModelProviderLocalService(model);
return await ensureProviderLocalService(
{
providerId: model.provider,
baseUrl: model.baseUrl,
headers: buildHealthProbeHeaders((model as { headers?: HeadersInit }).headers, probeHeaders),
service,
},
signal,
);
}
/** Ensure a provider endpoint's local service is healthy and return a request lease. */
export async function ensureProviderLocalService(
target: ProviderLocalServiceTarget,
signal?: AbortSignal | null,
): Promise<ProviderLocalServiceLease | undefined> {
const service = target.service;
if (!service) {
return undefined;
}
throwIfAborted(signal);
validateLocalServiceConfig(service, model.provider);
const healthUrl = resolveHealthUrl(service, model.baseUrl);
const healthHeaders = buildHealthProbeHeaders(model, probeHeaders);
const key = localServiceKey(model.provider, service, healthUrl);
validateLocalServiceConfig(service, target.providerId);
const healthUrl = resolveHealthUrl(service, target.baseUrl);
const healthHeaders = filterHealthProbeHeaders(target.headers);
const key = localServiceKey(target.providerId, service, healthUrl);
installExitHandler();
const managed = services.get(key) ?? { active: 0 };
services.set(key, managed);
@@ -117,7 +158,7 @@ export async function ensureModelProviderLocalService(
const startupAbort = new AbortController();
managed.startupAbort = startupAbort;
managed.starting = startAndWaitForLocalService({
provider: model.provider,
provider: target.providerId,
service,
healthUrl,
healthHeaders,
@@ -159,6 +200,17 @@ export function stopManagedProviderLocalServicesForTest(): void {
services.clear();
}
/** Return bounded local-service state for focused lifecycle tests. */
export function getManagedProviderLocalServiceDiagnosticsForTest(): LocalServiceDiagnostics[] {
return [...services.values()]
.map((managed) => managed.diagnostics)
.filter((value): value is LocalServiceDiagnostics => value !== undefined)
.map((value) => ({
...value,
...(value.lastExit ? { lastExit: { ...value.lastExit } } : {}),
}));
}
function validateLocalServiceConfig(service: ModelProviderLocalServiceConfig, provider: string) {
if (!path.isAbsolute(service.command)) {
throw new Error(`models.providers.${provider}.localService.command must be an absolute path`);
@@ -198,7 +250,7 @@ function sortedStringRecord(record: Record<string, string> | undefined): Record<
}
function buildHealthProbeHeaders(
model: Model,
providerHeaders: HeadersInit | undefined,
requestHeaders: HeadersInit | undefined,
): Headers | undefined {
const headers = new Headers();
@@ -212,11 +264,15 @@ function buildHealthProbeHeaders(
}
}
};
appendHeaders((model as { headers?: HeadersInit }).headers);
appendHeaders(providerHeaders);
appendHeaders(requestHeaders);
return [...headers].length > 0 ? headers : undefined;
}
function filterHealthProbeHeaders(headers: HeadersInit | undefined): Headers | undefined {
return buildHealthProbeHeaders(headers, undefined);
}
async function probeHealth(
url: string,
headers: HeadersInit | undefined,
@@ -267,29 +323,58 @@ async function startAndWaitForLocalService(params: {
await stopManagedProcessForRestart(managed, signal);
}
const startedAt = Date.now();
const diagnostics: LocalServiceDiagnostics = {
providerId: provider,
healthUrl,
startedAt,
stdoutTail: "",
stderrTail: "",
};
managed.diagnostics = diagnostics;
log.info(`starting ${provider} local service: ${service.command}`);
managed.process = spawn(service.command, service.args ?? [], {
cwd: service.cwd,
env: service.env ? { ...process.env, ...service.env } : process.env,
stdio: "ignore",
stdio: ["ignore", "pipe", "pipe"],
detached: shouldDetachChildForProcessTree(),
});
const child = managed.process;
diagnostics.pid = child.pid;
managed.lastExit = undefined;
child.stdout?.on("data", (chunk: Buffer | string) => {
diagnostics.stdoutTail = appendLocalServiceOutputTail(
diagnostics.stdoutTail,
chunk,
service.env,
);
});
child.stderr?.on("data", (chunk: Buffer | string) => {
diagnostics.stderrTail = appendLocalServiceOutputTail(
diagnostics.stderrTail,
chunk,
service.env,
);
});
child.unref();
child.once("exit", (code, signalLocal) => {
const exit = { code, signal: signalLocal };
diagnostics.lastExit = exit;
log.info(
`${provider} local service exited: ${signalLocal ? `signal=${signalLocal}` : `code=${code ?? 0}`}`,
`${provider} local service exited: ${signalLocal ? `signal=${signalLocal}` : `code=${code ?? 0}`}${diagnostics.stderrTail ? ` stderr=${diagnostics.stderrTail}` : ""}`,
);
if (managed.process === child) {
managed.lastExit = { code, signal: signalLocal };
managed.lastExit = exit;
managed.process = undefined;
}
});
const spawnError = await waitForSpawnResult(child, signal);
if (spawnError) {
throw new Error(`${provider} local service failed to start: ${spawnError.message}`);
throw new Error(
`${provider} local service failed to start: ${spawnError.message}${formatLocalServiceDiagnosticTail(diagnostics)}`,
);
}
diagnostics.spawnedAt = Date.now();
const readyTimeoutMs = resolvePositiveTimerTimeoutMs(
service.readyTimeoutMs,
@@ -298,14 +383,18 @@ async function startAndWaitForLocalService(params: {
const deadline = Date.now() + readyTimeoutMs;
for (;;) {
if (await probeHealth(healthUrl, healthHeaders, signal)) {
log.info(`${provider} local service ready`);
diagnostics.readyAt = Date.now();
diagnostics.lastHealthyAt = diagnostics.readyAt;
log.info(
`${provider} local service ready: pid=${diagnostics.pid ?? "unknown"} spawnMs=${diagnostics.spawnedAt - startedAt} readyMs=${diagnostics.readyAt - startedAt}`,
);
return;
}
if (managed.lastExit) {
throw new Error(
`${provider} local service exited before readiness with ${formatLocalServiceExit(
managed.lastExit,
)}`,
)}${formatLocalServiceDiagnosticTail(diagnostics)}`,
);
}
if (Date.now() >= deadline) {
@@ -315,6 +404,32 @@ async function startAndWaitForLocalService(params: {
}
}
function appendLocalServiceOutputTail(
current: string,
chunk: Buffer | string,
serviceEnv: Record<string, string> | undefined,
): string {
let redacted = redactSensitiveText(`${current}${chunk.toString()}`, { mode: "tools" });
for (const value of Object.values(serviceEnv ?? {})) {
if (value) {
redacted = redacted.replaceAll(value, "[redacted]");
}
}
const bytes = Buffer.from(redacted);
if (bytes.byteLength <= LOCAL_SERVICE_OUTPUT_TAIL_MAX_BYTES) {
return redacted;
}
let start = bytes.byteLength - LOCAL_SERVICE_OUTPUT_TAIL_MAX_BYTES;
while (start < bytes.byteLength && (bytes[start] & 0xc0) === 0x80) {
start += 1;
}
return bytes.subarray(start).toString("utf8");
}
function formatLocalServiceDiagnosticTail(diagnostics: LocalServiceDiagnostics): string {
return diagnostics.stderrTail ? `; stderr: ${diagnostics.stderrTail}` : "";
}
function scheduleIdleStop(
key: string,
managed: ManagedLocalService,