fix(gateway): keep model metadata available during plugin drains (#162301)

* fix(gateway): keep model metadata available during plugin drains

Keep the active model publication readable while admitted plugin work settles,
then retire it before resource replacement. Preserve execution fencing, auth
revocation, decision cancellation, rollback, and cleanup ownership.

Retain an automatic drain failure only while its original plugin configuration
delta remains unresolved. Allow explicit wait recovery and reversion alongside
changes to a different plugin without retrying unrelated failed work.

Validate on Blacksmith Testbox with real Gateway reader latency proof, 378
focused tests, 145 follow-up tests, negative controls, types, lint, and guards.

* fix(plugins): fence admission while preserving owned cleanup

* test(gateway): preserve plugin record helpers in reload fixture

* test: repair reload mocks and supervised process joins

* test(models): preserve queue receiver in drain observer
This commit is contained in:
Peter Steinberger 2026-10-01 01:38:43 -07:00 • committed by GitHub
parent 40d9ed0924
commit a3e4005ebc
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
48 changed files with 1150 additions and 233 deletions

View file

@ -28,6 +28,13 @@ No additional config edit is needed. The previous runtime stays active until the
change applies, and shutdown cancels pending retries. Other reload failures remain
visible in the Gateway log.
If an automatic plugin reload cannot drain active work, the Gateway records the
failure and keeps the last-good runtime. Later config edits do not repeat that
drain while its plugin generation is still active. Run `openclaw plugins reload
<id> --wait` to finish the replacement, or revert the pending plugin settings.
Edits that still include the unapplied plugin settings remain pending until
recovery; they cannot publish those settings through an unrelated hot update.
Direct file edits are treated as untrusted until they validate. The source's file adapter waits
for editor temp-write/rename churn to settle, reads the final file, and rejects
invalid external edits without rewriting `openclaw.json`. OpenClaw-owned config
@ -105,6 +112,9 @@ Model runtime selection keeps your authored settings separate from catalog defau
Hot reload and secrets reload preserve that distinction: catalog compatibility
metadata does not become a custom request override that switches a native runtime
back to OpenClaw.
Changes to `agents.defaults.models`, agent model selection and fallbacks, and
`models.providers` hot-apply without draining the Codex plugin. Changing Codex's
own plugin settings still follows its plugin reload policy.
Channel transport edits, such as `channels.slack.streaming.mode`, retain prepared
session rows and model catalogs. Agent rosters, session policy, store topology,

View file

@ -170,7 +170,9 @@ Durable final channel replies can use the admitting Gateway's current registry a
Replacement reserves the affected instance even when agent turns or unfinished cleanup retain it. The prepared-model replacement gate holds new runs while already admitted runs finish using their original callbacks. New top-level retained work cannot acquire the old instance; already admitted consumers can still derive work needed to finish their runs. Detailed readiness and logs report the retained-work count and drain deadline; the final RPC receipt reports application and any drain notices. Reload from the instance's own active callback still fails during preparation to avoid waiting on itself. Idle prepared publications do not block replacement.
Retained work and in-flight calls share a 60-second pre-stop budget. Retained runs finish before their callbacks close; sidecars release their capability consumers before the remaining finite consumers drain. Replacement then pauses ordinary calls before stopping services and channels. Idle service and channel custody stops later with its owners. If work does not finish, the reload fails once, resumes admission, and clears the reload status; the previous plugin generation keeps serving. The deadline does not cancel agent runs or permit disposal of unfinished writes. Retry `openclaw plugins reload <id>` after that work finishes, or explicitly select `--wait` to wait until admitted work settles. The explicit wait belongs to its requesting connection: Ctrl+C or disconnect cancels the pre-publication wait and runs the same rollback. Cancellation never disposes admitted work, bypasses cleanup, or reverses a committed publication. Service shutdown and recovery retain their bounded deadlines. Successful publication logs the applied replacement and emits `plugins.changed`.
While admitted work drains, ordinary model catalog, auth-status, and chat metadata reads continue using the active publication. Catalog refresh work, downloaded catalog adoption, and new execution still wait. Published facts retire when resource replacement begins; an admitted-work timeout keeps those facts available without rebuilding them. An unfinished startup publication still follows its existing cancellation and recovery path.
Retained work and in-flight calls share a 60-second pre-stop budget. Retained runs finish before ordinary call admission closes. Existing calls then settle while published metadata remains readable. Sidecars release their capability consumers before the remaining finite consumers drain; memory teardown retains cleanup authority for its exact retiring provider instances until the raw cleanup settles. Replacement stops services and channels only after that handoff. Idle service and channel custody stops later with its owners. If work does not finish, the reload fails once, resumes admission, and clears the reload status; the previous plugin generation keeps serving. The deadline does not cancel agent runs or permit disposal of unfinished writes. Retry `openclaw plugins reload <id>` after that work finishes, or explicitly select `--wait` to wait until admitted work settles. The explicit wait belongs to its requesting connection: Ctrl+C or disconnect cancels the pre-publication wait and runs the same rollback. Cancellation never disposes admitted work, bypasses cleanup, or reverses a committed publication. Service shutdown and recovery retain their bounded deadlines. Successful publication logs the applied replacement and emits `plugins.changed`.
A provider or harness plugin load failure remains recorded in its runtime generation. It makes that plugin unavailable without superseding the generation or blocking models that use healthy plugins. Inspect the failing owner with `openclaw plugins inspect <id> --runtime --json`. Use `openclaw doctor --fix` for supported installation repairs, or fix the reported problem in plugin code, then request `plugins.reload` through the admin Gateway API to load the repaired plugin.

View file

@ -1115,6 +1115,7 @@ src/plugins/plugin-instance-invocation.ts
src/plugins/plugin-instance-module-loader.ts
src/plugins/plugin-instance-owned-values.ts
src/plugins/plugin-instance-scope.ts
src/plugins/plugin-instance-settlement.ts
src/plugins/plugin-instance-value-views.ts
src/plugins/plugin-instance.ts
src/plugins/plugin-lifecycle-trace.ts

View file

@ -545,7 +545,16 @@ export async function loadPreparedModelCatalogOwnerSnapshot(
export async function loadPublishedPreparedModelCatalogOwnerSnapshot(
params: LoadPreparedModelCatalogParams = {},
): Promise<PreparedModelRuntimeSnapshot> {
return await withPreparedModelCatalogOwnerPolicy(params, "published", (snapshot) => snapshot);
return await withPreparedModelCatalogOwnerPolicy(
params,
"published",
(snapshot) => snapshot,
async (input) => ({
snapshot: await prepareModelRuntimeSnapshot(input, {
readPublished: params.readOnly !== false && params.refreshFullCatalog !== true,
}),
}),
);
}
/** Resolves a complete published owner for long-lived runtime consumers. */

View file

@ -66,6 +66,10 @@ export class PreparedModelRuntimeAuthPublicationOwner {
#transaction: PreparedModelRuntimeAuthTransaction | undefined;
#drainTail: Promise<void> = Promise.resolve();
get hasPendingPublication(): boolean {
return this.#transaction !== undefined;
}
enqueue(
invalidatedOwners: readonly PreparedModelRuntimeOwner[],
profileSetChanged = false,

View file

@ -5,6 +5,7 @@ import { resolvePluginMetadataSnapshot } from "../plugins/plugin-metadata-snapsh
import { isReservedSystemAgentId } from "../system-agent/agent-id.js";
import { getPreparedModelRuntimeBorrowedSnapshot } from "./prepared-model-runtime-generation-scope.js";
import { capturePreparedModelRuntimeCatalog } from "./prepared-model-runtime.capture.js";
import { isPreparedModelRuntimePluginLifecycleFailure } from "./prepared-model-runtime.errors.js";
import {
PreparedModelRuntimeOwnerNotPublishedError,
PreparedModelRuntimePublicationSupersededError,
@ -365,6 +366,9 @@ export async function acquirePreparedModelRuntimeLeaseFromOwners(
supersededPublication = error;
continue;
}
if (context.getPendingReplacement() && isPreparedModelRuntimePluginLifecycleFailure(error)) {
continue;
}
throw error;
}
const published = context.owners.get(key);

View file

@ -187,6 +187,7 @@ function applyRemoteModelCatalogUpdateNow(
signal: host.getCancellationSignal(),
isPublicationCurrent: () =>
isCurrent() &&
!host.getPendingReplacement() &&
!attempt.signal.aborted &&
host.getEpoch() === epoch &&
preparedModelRuntimeConfigsMatch(attemptConfig, getConfig()),

View file

@ -201,10 +201,14 @@ describe("Gateway plugin reload run admission", () => {
}
});
const handler = createPluginReloadHandler(async ({ prepareConfigEffects, commitRuntime }) => {
prepareConfigEffects({ pluginIds: new Set(["synthetic"]), channels: new Set() });
const effects = prepareConfigEffects({
pluginIds: new Set(["synthetic"]),
channels: new Set(),
});
events.push("plugin-drain");
drainageStarted.resolve();
await finishDrainage.promise;
effects.retire();
await commitRuntime({ publish: () => setRuntimeConfigSnapshot(committed, committed) });
pluginCommitted.resolve();
return {
@ -305,16 +309,17 @@ describe("Gateway plugin reload run admission", () => {
},
);
const handler = createPluginReloadHandler(async ({ prepareConfigEffects, commitRuntime }) => {
const restorePreparedRuntime = prepareConfigEffects({
const effects = prepareConfigEffects({
pluginIds: new Set(["synthetic"]),
channels: new Set(),
});
drainageStarted.resolve();
await finishDrainage.promise;
if (outcome === "rollback") {
await restorePreparedRuntime();
await effects.rollback();
throw pluginFailure;
}
effects.retire();
await commitRuntime({ publish: () => setRuntimeConfigSnapshot(committed, committed) });
if (outcome === "activation failure") {
throw pluginFailure;
@ -403,21 +408,18 @@ describe("Gateway plugin reload run admission", () => {
await nextTurn();
expect(settled).toBe(false);
expect(requestSettled).toBe(false);
expect(catalogSettled).toBe(false);
expect(catalogSettled).toBe(true);
await expect(catalogRequest).resolves.toMatchObject({ config: retained });
finishDrainage.resolve();
if (outcome !== "commit") {
await expect(reload).rejects.toBe(pluginFailure);
} else {
await expect(reload).resolves.toMatchObject({ status: "applied" });
}
// Both rollback and committed failure must replace the drained model owner
// before readers resume, while retaining the original lifecycle error.
// Execution resumes only after rollback or replacement publication settles.
await expect(request).resolves.toMatchObject({
config: outcome === "rollback" ? retained : committed,
});
await expect(catalogRequest).resolves.toMatchObject({
config: outcome === "rollback" ? retained : committed,
});
const lease = await admission!;
expect(lease.snapshot.config).toEqual(outcome === "rollback" ? retained : committed);
expect(lease.snapshot.workspaceDir).toBe(input.workspaceDir);
@ -475,7 +477,10 @@ it.each([
});
}
const handler = createPluginReloadHandler(async ({ prepareConfigEffects, commitRuntime }) => {
prepareConfigEffects({ pluginIds: new Set(pluginLifecycle.pluginIds), channels: new Set() });
prepareConfigEffects({
pluginIds: new Set(pluginLifecycle.pluginIds),
channels: new Set(),
}).retire();
await commitRuntime({ publish: () => setRuntimeConfigSnapshot(committed, committed) });
if (event === "refresh failure") {
mocks.configuredAgentIdsError = refreshFailure;
@ -776,7 +781,10 @@ it.each(["success", "activation failure"] as const)(
committed: true,
});
const handler = createPluginReloadHandler(async ({ prepareConfigEffects, commitRuntime }) => {
prepareConfigEffects({ pluginIds: new Set(pluginLifecycle.pluginIds), channels: new Set() });
prepareConfigEffects({
pluginIds: new Set(pluginLifecycle.pluginIds),
channels: new Set(),
}).retire();
await commitRuntime({ publish: () => setRuntimeConfigSnapshot(committed, committed) });
if (outcome === "activation failure") {
throw activationFailure;

View file

@ -0,0 +1,44 @@
import { expect, it, vi } from "vitest";
import { createDeferred } from "../../test/helpers/promise.js";
import { createPreparedModelRuntimePluginDrain } from "./prepared-model-runtime.lifecycle.js";
import { PreparedModelRuntimePublicationQueue } from "./prepared-model-runtime.publication-queue.js";
it("rechecks a successor plugin reservation after a publication reaches the queue", async () => {
const signal = new AbortController().signal;
const drain = createPreparedModelRuntimePluginDrain(
() => signal,
() => false,
);
const first = drain.begin();
const queue = new PreparedModelRuntimePublicationQueue();
const releaseQueue = createDeferred();
const queued = createDeferred();
const blocker = queue.enqueue(() => releaseQueue.promise);
const enqueue = queue.enqueue.bind(queue);
vi.spyOn(queue, "enqueue").mockImplementationOnce((task) => {
const publication = enqueue(task);
queued.resolve();
return publication;
});
const writes: string[] = [];
const publication = drain.runAfter(queue, async () => {
writes.push("published");
});
let successor: ReturnType<typeof drain.begin> | undefined;
try {
first.release();
await queued.promise;
successor = drain.begin();
releaseQueue.resolve();
await queue.settle();
expect(writes).toEqual([]);
successor.release();
await publication;
expect(writes).toEqual(["published"]);
} finally {
first.release();
successor?.release();
releaseQueue.resolve();
await Promise.allSettled([blocker, publication]);
}
});

View file

@ -96,3 +96,69 @@ export function createPreparedModelRuntimeReplacement(): PreparedModelRuntimeRep
void promise.catch(() => undefined);
return { gateId: Symbol("prepared-model-runtime-replacement"), promise, resolve, reject };
}
/** Execution waits for plugin replacement while the active publication remains readable. */
export function createPreparedModelRuntimePluginDrain(
getCancellationSignal: () => AbortSignal,
hasPendingPublication: () => boolean,
) {
let pending: PreparedModelRuntimeReplacement | undefined;
return {
get pending() {
return pending;
},
begin: () => {
capturePreparedModelRuntimeLifetime();
getCancellationSignal().throwIfAborted();
if (pending) {
throw new Error("Prepared model runtime plugin drain is already pending");
}
const drain = createPreparedModelRuntimeReplacement();
pending = drain;
const unregister = registerPreparedModelRuntimeClose(async (error) => {
unregister();
if (pending === drain) {
pending = undefined;
}
drain.reject(error);
});
return {
pendingPublication: hasPendingPublication(),
release: () => {
unregister();
if (pending === drain) {
pending = undefined;
drain.resolve();
}
},
};
},
async runAfter(
queue: { enqueue: (task: () => Promise<void>) => Promise<void> },
run: () => Promise<void>,
): Promise<void> {
const assertCurrent = capturePreparedModelRuntimeLifetime();
for (;;) {
const reservation = pending;
if (reservation) {
await reservation.promise;
continue;
}
assertCurrent();
let started = false;
await queue.enqueue(async () => {
assertCurrent();
if (pending) {
return;
}
started = true;
await run();
});
if (started) {
return;
}
// A newer reservation must release the queue before waiting for its drain.
}
},
};
}

View file

@ -72,6 +72,24 @@ type PublishedModelRuntimeContext = {
owners: Map<string, PreparedModelRuntimeOwner>;
};
/** Bind passive reads and retained acquisitions to the same publication owner. */
export function createPublishedModelRuntimeAccess(
context: PublishedModelRuntimeContext,
getPublishedReplacement: PublishedModelRuntimeContext["getPendingReplacement"],
) {
const readContext = { ...context, getPendingReplacement: getPublishedReplacement };
return {
acquire: (input: PreparedModelRuntimeInput) =>
projectPublishedModelRuntimeOwner(input, context, retainPublishedModelRuntimeOwner),
prepare: (input: PreparedModelRuntimeInput, options: { readPublished?: boolean } = {}) =>
projectPublishedModelRuntimeOwner(
input,
options.readPublished ? readContext : context,
(_owner, snapshot) => snapshot,
),
};
}
/** Project or retain the exact published owner before its snapshot crosses an await. */
export async function projectPublishedModelRuntimeOwner<T>(
rawInput: PreparedModelRuntimeInput,

View file

@ -2,12 +2,13 @@
// oxfmt-ignore
import { usePreparedModelRuntimeHarness } from "./prepared-model-runtime.test-harness.js";
import { isDeepStrictEqual } from "node:util";
import { describe, expect, it } from "vitest";
import { describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../test/helpers/promise.js";
import { createEmptyPluginRegistry } from "../plugins/registry-empty.js";
import { getPluginLoaderCacheState } from "../plugins/registry-lifecycle.js";
import { createPluginRecord } from "../plugins/status.test-helpers.js";
import {
beginPreparedModelRuntimePluginDrain,
loadPublishedGatewayReplyDispatchRuntime,
prepareModelRuntimeSnapshot,
refreshPreparedModelRuntimeSnapshots,
@ -18,6 +19,60 @@ const fixture = usePreparedModelRuntimeHarness({ label: "prepared-model-runtime"
const { mocks } = fixture;
describe("prepared model runtime reload auth adoption", () => {
it.each(["rollback", "replacement"] as const)(
"revokes changed auth immediately and publishes it after plugin drain %s",
async (outcome) => {
mocks.configuredAgentIds = ["default"];
const initialConfig = {};
const options = { gatewayLifecycle: true, catalogMode: "static" as const };
await refreshPreparedModelRuntimeSnapshots(initialConfig, options);
const input = fixture.agentInput("default", initialConfig);
const original = await prepareModelRuntimeSnapshot(input);
const drain = beginPreparedModelRuntimePluginDrain();
let draining = true;
mocks.prepareStaticCatalog.mockClear();
mocks.prepareStaticCatalog.mockImplementation(async () => {
expect(draining).toBe(false);
return { entries: [] };
});
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout"] });
let read: ReturnType<typeof prepareModelRuntimeSnapshot> | undefined;
let publication: Promise<void> | undefined;
try {
mocks.mutationListener?.({ agentDir: input.agentDir, affectsInheritedStores: false });
expect(original.isCurrent()).toBe(false);
let readSettled = false;
read = prepareModelRuntimeSnapshot(input, { readPublished: true });
void read.then(
() => {
readSettled = true;
},
() => {
readSettled = true;
},
);
await vi.advanceTimersByTimeAsync(0);
expect(mocks.prepareStaticCatalog).not.toHaveBeenCalled();
expect(readSettled).toBe(false);
draining = false;
drain.release();
if (outcome === "replacement") {
publication = refreshPreparedModelRuntimeSnapshots({ plugins: {} }, options);
await publication;
}
const refreshed = await read;
expect(refreshed).not.toBe(original);
expect(refreshed.isCurrent()).toBe(true);
expect(mocks.prepareStaticCatalog).toHaveBeenCalled();
} finally {
draining = false;
drain.release();
vi.useRealTimers();
await Promise.allSettled([read, publication]);
}
},
);
it("releases a rejected replacement's cached registry after the old catalog finishes", async () => {
mocks.configuredAgentIds = ["default"];
const cache = getPluginLoaderCacheState();

View file

@ -34,11 +34,13 @@ import { PreparedModelRuntimePublicationSupersededError } from "./prepared-model
import {
acquireReadOnlyPreparedModelRuntime,
applyRemoteModelCatalogUpdate,
beginPreparedModelRuntimePluginDrain,
prepareModelRuntimeSnapshot,
refreshPreparedModelRuntimeSnapshots,
} from "./prepared-model-runtime.js";
import { closePreparedModelRuntimeSnapshots } from "./prepared-model-runtime.lifecycle.js";
import { registerPreparedModelRuntimePublicationListener } from "./prepared-model-runtime.publication-events.js";
import { PreparedModelRuntimePublicationQueue } from "./prepared-model-runtime.publication-queue.js";
const fixture = usePreparedModelRuntimeHarness({
label: "remote-publication",
@ -137,6 +139,60 @@ async function refresh() {
}
afterEach(() => setRemoteModelCatalogOverlaySourcesForTest());
it("keeps downloaded catalogs pending while plugin work drains", async ({ signal }) => {
await setup();
const preparing = createDeferred();
const releasePricing = createDeferred();
const preparePricing = pricing.prepareModelPricingContext;
const pricingSpy = vi
.spyOn(pricing, "prepareModelPricingContext")
.mockImplementationOnce(async (...args) => {
preparing.resolve();
await releasePricing.promise;
return await preparePricing(...args);
});
const attempted = createDeferred();
const queueSpy = vi.spyOn(PreparedModelRuntimePublicationQueue.prototype, "enqueue");
queueSpy.mockImplementationOnce(function (this: PreparedModelRuntimePublicationQueue, ...args) {
queueSpy.mockRestore();
const publication = this.enqueue(...args);
void publication.then(
() => attempted.resolve(),
() => attempted.resolve(),
);
return publication;
});
const adoption = applyRemoteModelCatalogUpdate(() => config);
let drain: ReturnType<typeof beginPreparedModelRuntimePluginDrain> | undefined;
try {
await withinTest(preparing.promise, signal);
drain = beginPreparedModelRuntimePluginDrain();
releasePricing.resolve();
await withinTest(attempted.promise, signal);
const active = await withinTest(
loadPreparedGatewayModelCatalogSnapshot({ agentId: "default", getConfig: () => config }),
signal,
);
expect(active.entries.map((entry) => entry.id)).toContain("remote-200");
expect(captureRemoteModelCatalogStartupSnapshot()?.generatedAt).toBe(200);
// Lifecycle publication must not queue behind adoption's pending drain wait.
await withinTest(
refreshPreparedModelRuntimeSnapshots(config, { catalogMode: "static" }),
signal,
);
expect(captureRemoteModelCatalogStartupSnapshot()?.generatedAt).toBe(200);
drain.release();
expect(await withinTest(adoption, signal)).toBe("published");
expect(captureRemoteModelCatalogStartupSnapshot()?.generatedAt).toBe(300);
} finally {
drain?.release();
releasePricing.resolve();
await adoption;
queueSpy.mockRestore();
pricingSpy.mockRestore();
}
});
it("does not reuse a dynamic build captured before a remote publication", async () => {
await setup();
const preparing = createDeferred();

View file

@ -20,6 +20,7 @@ import {
closePreparedModelRuntimeSnapshots,
registerPreparedModelRuntimeClose,
createPreparedModelRuntimeReplacement,
createPreparedModelRuntimePluginDrain,
retirePreparedModelRuntimeGeneration,
} from "./prepared-model-runtime.lifecycle.js";
import {
@ -50,6 +51,7 @@ import {
} from "./prepared-model-runtime.publication-events.js";
import { PreparedModelRuntimePublicationQueue } from "./prepared-model-runtime.publication-queue.js";
import {
createPublishedModelRuntimeAccess,
projectPublishedModelRuntimeOwner,
refreshPublishedModelRuntimeCatalog,
retainPublishedModelRuntimeOwner,
@ -104,9 +106,17 @@ const publicationQueue = new PreparedModelRuntimePublicationQueue();
let refreshRequestEpoch = 0;
let refreshCancellation = new AbortController();
let pendingModelRuntimeReplacement: PreparedModelRuntimeReplacement | undefined;
const modelRuntimeDrain = createPreparedModelRuntimePluginDrain(
() => {
captureModelRuntimeLifetime();
return refreshCancellation.signal;
},
() => pendingModelRuntimeReplacement !== undefined || authPublication.hasPendingPublication,
);
const authPublication = new PreparedModelRuntimeAuthPublicationOwner();
const getBlockingReplacement = () =>
pendingModelRuntimeReplacement?.degraded ? undefined : pendingModelRuntimeReplacement;
const getAdmissionReplacement = () => modelRuntimeDrain.pending ?? getBlockingReplacement();
const replyDispatchPublication = new PreparedReplyDispatchPublicationOwner({
isGatewayLifecycleActive: () => gatewayLifecycleActive,
@ -116,7 +126,7 @@ const replyDispatchPublication = new PreparedReplyDispatchPublicationOwner({
agentDir: ".",
config: {},
}).pending,
getPendingReplacement: () => getBlockingReplacement()?.promise,
getPendingReplacement: () => getAdmissionReplacement()?.promise,
});
export const loadPublishedGatewayReplyDispatchRuntime = replyDispatchPublication.load;
@ -189,17 +199,6 @@ export async function acquirePublishedPreparedModelRuntime(
return await loadPreparedModelRuntimeOwner(rawInput, retainPublishedModelRuntimeOwner);
}
/** Retains the selected publication without activating an unpublished owner. */
export async function acquirePreparedModelRuntimeSnapshot(
rawInput: PreparedModelRuntimeInput,
): Promise<PreparedModelRuntimeLease> {
return await projectPublishedModelRuntimeOwner(
rawInput,
preparedModelRuntimeLeaseContext,
retainPublishedModelRuntimeOwner,
);
}
/** Retains existing execution owners, including switched-away models, without loading plugins. */
export async function acquireAgentRuntimeCleanupRegistries(agentDir: string) {
return await acquireRetainedAgentRuntimeCleanupRegistries(
@ -220,10 +219,10 @@ async function loadPreparedModelRuntimeOwner<T>(
});
for (;;) {
assertLifetime();
const replacement = getBlockingReplacement();
const replacement = getAdmissionReplacement();
if (replacement) {
await replacement.promise;
if (getBlockingReplacement()) {
if (getAdmissionReplacement()) {
continue;
}
input = rebindInputToCommittedConfiguredOwner(owners, input);
@ -240,12 +239,12 @@ async function loadPreparedModelRuntimeOwner<T>(
throw error;
}
}
if (getBlockingReplacement()) {
if (getAdmissionReplacement()) {
continue;
}
assertLifetime();
const activated = await activateStandalonePreparedModelRuntime(input);
if (getBlockingReplacement()) {
if (getAdmissionReplacement()) {
continue;
}
try {
@ -276,6 +275,9 @@ export function getPendingPreparedModelRuntimeReplacement(): Promise<void> | und
return getBlockingReplacement()?.promise;
}
/** Fence new execution while plugin work drains, without withdrawing the active catalog. */
export const beginPreparedModelRuntimePluginDrain = modelRuntimeDrain.begin;
/** Publishes one owner from an explicit startup/activation lifecycle boundary. */
export async function publishPreparedModelRuntimeSnapshot(
rawInput: PreparedModelRuntimeInput,
@ -390,8 +392,16 @@ const preparedModelRuntimeLeaseContext = {
retainedGatewayRunOwners,
getBuildTimeoutMs: () => modelRuntimeBuildTimeoutMs,
getGatewayLifecycleActive: () => gatewayLifecycleActive,
getPendingReplacement: getBlockingReplacement,
getPendingReplacement: getAdmissionReplacement,
};
const publishedModelRuntime = createPublishedModelRuntimeAccess(
preparedModelRuntimeLeaseContext,
getBlockingReplacement,
);
/** Retains the selected publication without activating an unpublished owner. */
export const acquirePreparedModelRuntimeSnapshot = publishedModelRuntime.acquire;
/** Returns the published snapshot; only passive catalog reads may bypass execution drainage. */
export const prepareModelRuntimeSnapshot = publishedModelRuntime.prepare;
/** Acquires a run generation from configured facts; full catalog discovery is explicit. */
export async function acquireAgentRunPreparedModelRuntime(
@ -419,17 +429,6 @@ export async function acquireReadOnlyPreparedModelRuntime(
);
}
/** Returns the snapshot published by the lifecycle owner. Request config cannot replace it. */
export async function prepareModelRuntimeSnapshot(
rawInput: PreparedModelRuntimeInput,
): Promise<PreparedModelRuntimeSnapshot> {
return await projectPublishedModelRuntimeOwner(
rawInput,
preparedModelRuntimeLeaseContext,
(_owner, snapshot) => snapshot,
);
}
/** Initializes or refreshes inventory on catalog demand; turn admission remains static. */
export async function refreshPreparedModelRuntimeCatalog(
snapshot: PreparedModelRuntimeSnapshot,
@ -518,7 +517,9 @@ const remoteCatalogPublication = configuredRefresh.createRemoteCatalogPublicatio
replyDispatchPublication,
getEpoch: () => refreshRequestEpoch,
getCancellationSignal: () => refreshCancellation.signal,
getPendingReplacement: () => pendingModelRuntimeReplacement?.promise,
// Catalog adoption also waits for a degraded startup's final publication.
getPendingReplacement: () =>
(modelRuntimeDrain.pending ?? pendingModelRuntimeReplacement)?.promise,
onPluginGenerationRetired: recoverRetiredConfiguredPluginGeneration,
});
export const { applyRemoteModelCatalogUpdate, advancePreparedModelRuntimeConfig } =
@ -711,7 +712,7 @@ function invalidateForAuthMutation(event: PreparedModelRuntimeAuthMutation): voi
notifyPreparedModelRuntimePublication({ phase: "invalidated" });
return;
}
const publication = publicationQueue.enqueue(async () => {
const publish = async () => {
// A pending replacement gate means a queued config publication owns the next generation:
// it drains queued auth mutations against the new config and rebuilds/announces the
// dispatch publication. Rebuilding here would revive stale owners with the old config or
@ -734,7 +735,10 @@ function invalidateForAuthMutation(event: PreparedModelRuntimeAuthMutation): voi
refreshCommittedProviderCatalogs(owners.values());
}
});
});
};
// Auth revocation fences its affected owners immediately; publication must wait for the
// plugin reservation to settle without occupying the queue needed by replacement recovery.
const publication = modelRuntimeDrain.runAfter(publicationQueue, publish);
notifyPreparedModelRuntimePublication({ phase: "invalidated" });
void publication.catch((error: unknown) => {
if (!authPublication.isCurrent(transaction)) {

View file

@ -74,6 +74,12 @@ export function diffGatewayReloadPaths(
reloadPrefixes: Iterable<string>,
): string[] {
const refinementPrefixes = new Set(reloadPrefixes);
// Preserve individual plugin owners when the entries dictionary is added or removed.
for (const config of [prevConfig, nextConfig]) {
for (const pluginId of Object.keys(config.plugins?.entries ?? {})) {
refinementPrefixes.add(`plugins.entries.${pluginId}`);
}
}
// Decision selectors refine to authored leaves; other wildcard owners retain parent lifecycle rules.
if (refinementPrefixes.delete("agents.entries.*.decisionModel")) {
for (const config of [prevConfig, nextConfig]) {

View file

@ -27,6 +27,8 @@ export type GatewayReloadPlan = {
restartHeartbeat: boolean;
reconcileSystemJobs?: boolean;
reloadPlugins: boolean;
/** Canonical config/install deltas that require a plugin replacement. */
reloadPluginPaths?: string[];
/** Plugin owners whose undeclared channel settings require fresh registration. */
reloadPluginIds?: Set<string>;
pluginLifecycle?: {
@ -574,6 +576,9 @@ export function buildGatewayReloadPlan(
plan.hotReasons.push(path);
for (const action of rule?.actions ?? []) {
plan[action] = true;
if (action === "reloadPlugins") {
(plan.reloadPluginPaths ??= []).push(path);
}
}
if (rule?.replaceChannelPlugins) {
// Manifest channel IDs survive even when registration has no active channel.

View file

@ -0,0 +1,74 @@
import {
getPluginRuntimeGeneration,
PluginRuntimeApplicationError,
type PluginRuntimeApplication,
} from "../plugins/lifecycle.js";
import type { GatewayReloadPlan } from "./config-reload-plan.js";
import { GatewayConfigReloadSupersededError } from "./server-reload-contracts.js";
export function isConfigReloadSuperseded(error: unknown): boolean {
// Only completed rollback preserves the direct cause. Cleanup failures and
// published replacements must settle instead of transferring the write.
const cause =
error instanceof PluginRuntimeApplicationError && !error.details.committed
? error.cause
: error;
return cause instanceof GatewayConfigReloadSupersededError;
}
/** Retain a failed automatic drain until its generation or pending settings settle. */
export function createConfigPluginDrainTracker() {
let failure:
| { error: PluginRuntimeApplicationError; paths: readonly string[]; reported: boolean }
| undefined;
const hasPendingFailure = (plan?: GatewayReloadPlan) =>
failure?.paths.some((failedPath) =>
plan?.changedPaths.some(
(path) =>
path === failedPath ||
path.startsWith(`${failedPath}.`) ||
failedPath.startsWith(`${path}.`),
),
) ?? false;
return {
assertCanApply(plan: GatewayReloadPlan) {
if (failure && failure.error.details.generation !== getPluginRuntimeGeneration()) {
failure = undefined;
}
if (
plan.reloadPlugins &&
failure &&
!plan.pluginLifecycle?.waitForDrain &&
hasPendingFailure(plan)
) {
// Publishing this candidate would expose plugin settings whose replacement never committed.
throw failure.error;
}
},
recordFailure(plan: GatewayReloadPlan, error: unknown) {
if (
!plan.pluginLifecycle &&
error instanceof PluginRuntimeApplicationError &&
error.details.phase === "drain" &&
!error.details.committed &&
!isConfigReloadSuperseded(error)
) {
failure = { error, paths: plan.reloadPluginPaths ?? [], reported: false };
}
},
applied(plan?: GatewayReloadPlan, runtime?: PluginRuntimeApplication) {
// Clear reverted settings only after the candidate applies, including unrelated plugin edits.
if (!plan?.reloadPlugins || runtime || !hasPendingFailure(plan)) {
failure = undefined;
}
},
shouldReport(error: unknown) {
if (!failure || error !== failure.error) {
return true;
}
const report = !failure.reported;
failure.reported = true;
return report;
},
};
}

View file

@ -0,0 +1,217 @@
import { afterEach, beforeEach, expect, it, vi } from "vitest";
import type { OpenClawConfig } from "../config/types.openclaw.js";
import { getPluginRuntimeGeneration, PluginRuntimeApplicationError } from "../plugins/lifecycle.js";
import {
closeTestConfigReloaders,
createReloaderHarness,
flushWatcherChange,
makeSnapshot,
prepareConfigReloadTest,
} from "./config-reload.test-support.js";
vi.mock("../config/io.audit.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../config/io.audit.js")>()),
appendConfigAuditRecord: vi.fn(),
}));
vi.mock("../config/config-journal-snapshot.js", async (importOriginal) => ({
...(await importOriginal<typeof import("../config/config-journal-snapshot.js")>()),
readLatestConfigSnapshotAuditRecordAsync: vi.fn(async () => null),
upsertConfigSnapshotAuditRecordAsync: vi.fn(),
}));
beforeEach((context) => {
prepareConfigReloadTest(context);
vi.useFakeTimers();
});
afterEach(async () => {
await closeTestConfigReloaders();
vi.useRealTimers();
vi.restoreAllMocks();
});
it("hot-applies model settings without replacing the Codex generation", async () => {
const initialConfig: OpenClawConfig = {
plugins: { entries: { codex: { enabled: true } } },
};
let config = initialConfig;
const harness = createReloaderHarness(async () => makeSnapshot({ config }), {
initialConfig,
});
await harness.reloader.ready;
for (const update of [
{ agents: { defaults: { models: { "openai/gpt-5.6-sol": { alias: "primary" } } } } },
{ agents: { entries: { worker: { model: { primary: "openai/gpt-5.6-sol" } } } } },
{
models: {
providers: { openai: { baseUrl: "https://api.openai.com/v1", models: [] } },
},
},
] satisfies OpenClawConfig[]) {
config = { ...initialConfig, ...update };
await flushWatcherChange(harness);
expect(harness.onHotReload.mock.lastCall?.[0].reloadPlugins).toBe(false);
expect(harness.onConfigApplied.mock.lastCall?.[1]).toEqual(config);
}
expect(harness.onHotReload).toHaveBeenCalledTimes(3);
expect(harness.onRestart).not.toHaveBeenCalled();
});
it.each(["explicit wait", "revert", "revert with model edit"] as const)(
"retains a failed automatic plugin drain until recovery: %s",
async (recovery) => {
const initialConfig: OpenClawConfig = {
plugins: { entries: { codex: { enabled: true, config: { sandbox: "read-only" } } } },
};
let config: OpenClawConfig = {
plugins: { entries: { codex: { enabled: true, config: { sandbox: "workspace-write" } } } },
};
const failure = new PluginRuntimeApplicationError("admitted work did not settle", {
operationId: "failed-automatic-drain",
generation: getPluginRuntimeGeneration(),
pluginIds: ["codex"],
phase: "drain",
committed: false,
});
const runtime = {
operationId: "explicit-wait-recovery",
generation: getPluginRuntimeGeneration(),
pluginIds: ["codex"],
};
const harness = createReloaderHarness(async () => makeSnapshot({ config }), {
initialConfig,
onHotReload: async (plan) => {
if (plan.reloadPlugins && !plan.pluginLifecycle?.waitForDrain) {
throw failure;
}
return plan.reloadPlugins ? { status: "applied", runtime } : "applied";
},
});
await harness.reloader.ready;
await flushWatcherChange(harness);
expect(harness.onHotReload).toHaveBeenCalledOnce();
await flushWatcherChange(harness);
config = {
...config,
agents: { defaults: { models: { "openai/gpt-5.6-sol": { alias: "primary" } } } },
};
await flushWatcherChange(harness);
expect(harness.onHotReload).toHaveBeenCalledOnce();
expect(harness.log.error).toHaveBeenCalledOnce();
expect(harness.onConfigApplied).not.toHaveBeenCalled();
if (recovery !== "explicit wait") {
config = {
...initialConfig,
...(recovery === "revert with model edit" ? { agents: config.agents } : {}),
};
await flushWatcherChange(harness);
expect(harness.onConfigAccepted.mock.lastCall?.[0]).toEqual(config);
const attempts = harness.onHotReload.mock.calls.length;
config = {
...config,
plugins: { entries: { codex: { enabled: true, config: { sandbox: "new-settings" } } } },
};
await flushWatcherChange(harness);
expect(harness.onHotReload).toHaveBeenCalledTimes(attempts + 1);
expect(harness.log.error).toHaveBeenCalledTimes(2);
await flushWatcherChange(harness);
expect(harness.onHotReload).toHaveBeenCalledTimes(attempts + 1);
return;
}
await expect(
harness.reloader.applyPluginLifecycleChange({
config,
pluginIds: ["codex"],
reason: "reload",
}),
).rejects.toBe(failure);
expect(harness.onHotReload).toHaveBeenCalledOnce();
await expect(
harness.reloader.applyPluginLifecycleChange({
config,
pluginIds: ["codex"],
reason: "reload",
waitForDrain: true,
}),
).resolves.toBe(runtime);
expect(harness.onHotReload).toHaveBeenCalledTimes(2);
expect(harness.onConfigApplied.mock.lastCall?.[1]).toEqual(config);
config = { ...config, agents: { defaults: { model: "openai/gpt-5.6-sol" } } };
await flushWatcherChange(harness);
expect(harness.onHotReload).toHaveBeenCalledTimes(3);
expect(harness.onHotReload.mock.lastCall?.[0].reloadPlugins).toBe(false);
expect(harness.onConfigApplied.mock.lastCall?.[1]).toEqual(config);
},
);
it.each([
{ existingEntries: true, partialRevertFirst: false },
{ existingEntries: true, partialRevertFirst: true },
{ existingEntries: false, partialRevertFirst: false },
{ existingEntries: false, partialRevertFirst: true },
])(
"applies a different plugin after reverting the failed delta (existing entries: $existingEntries, partial revert first: $partialRevertFirst)",
async ({ existingEntries, partialRevertFirst }) => {
const codex = {
enabled: true,
config: { sandbox: "read-only", appServer: { args: ["--original"] } },
};
const other = { enabled: true, config: { value: "original" } };
const initialConfig: OpenClawConfig = existingEntries
? { plugins: { entries: { codex, other } } }
: {};
const pendingCodex = {
...codex,
config: { sandbox: "workspace-write", appServer: { args: ["--replacement"] } },
};
let config: OpenClawConfig = {
plugins: { entries: { codex: pendingCodex, ...(existingEntries ? { other } : {}) } },
};
const harness = createReloaderHarness(async () => makeSnapshot({ config }), {
initialConfig,
});
await harness.reloader.ready;
harness.onHotReload.mockRejectedValueOnce(
new PluginRuntimeApplicationError("admitted work did not settle", {
operationId: "failed-codex-drain",
generation: getPluginRuntimeGeneration(),
pluginIds: ["codex"],
phase: "drain",
committed: false,
}),
);
await flushWatcherChange(harness);
expect(harness.onHotReload).toHaveBeenCalledOnce();
const nextOther = { ...other, config: { value: "replacement" } };
if (partialRevertFirst) {
config = {
plugins: {
entries: {
codex: {
...pendingCodex,
config: existingEntries
? { ...pendingCodex.config, sandbox: codex.config.sandbox }
: { appServer: pendingCodex.config.appServer },
},
other: nextOther,
},
},
};
await flushWatcherChange(harness);
expect(harness.onHotReload).toHaveBeenCalledOnce();
expect(harness.onConfigApplied).not.toHaveBeenCalled();
}
config = { plugins: { entries: { ...(existingEntries ? { codex } : {}), other: nextOther } } };
await flushWatcherChange(harness);
expect(harness.onHotReload).toHaveBeenCalledTimes(2);
expect(harness.onHotReload.mock.lastCall?.[0].reloadPlugins).toBe(true);
expect(harness.onConfigApplied.mock.lastCall?.[1]).toEqual(config);
expect(harness.log.error).toHaveBeenCalledOnce();
},
);

View file

@ -34,7 +34,10 @@ import type { ConfigFileSnapshot, OpenClawConfig } from "../config/types.opencla
import type { GatewayScheduledJob } from "../infra/gateway-scheduler.js";
import { getProcessGatewayPluginMetadataSnapshot } from "../plugins/current-plugin-metadata-state.js";
import { hashStableJson } from "../plugins/installed-plugin-index-hash.js";
import { loadInstalledPluginIndexInstallRecords } from "../plugins/installed-plugin-index-records.js";
import {
loadInstalledPluginIndexInstallRecords,
withPluginInstallRecords,
} from "../plugins/installed-plugin-index-records.js";
import {
getPluginRuntimeGeneration,
PluginRuntimeApplicationError,
@ -59,6 +62,10 @@ import {
resolvePluginInstallReloadMetadata,
type GatewayReloadPlan,
} from "./config-reload-plan.js";
import {
createConfigPluginDrainTracker,
isConfigReloadSuperseded,
} from "./config-reload-plugin-drain.js";
import { resolveGatewayReloadSettings } from "./config-reload-settings.js";
import type {
GatewayConfigReloader,
@ -83,24 +90,6 @@ const MISSING_CONFIG_MAX_RETRIES = 2;
const LEASE_RETRY_INITIAL_DELAY_MS = 250;
const LEASE_RETRY_MAX_DELAY_MS = 5_000;
function asPluginInstallConfig(records: PluginInstallRecords): OpenClawConfig {
return {
plugins: {
installs: records,
},
};
}
function isConfigReloadSuperseded(error: unknown): boolean {
// Only completed rollback preserves the direct cause. Cleanup failures and
// published replacements must settle instead of transferring the write.
const cause =
error instanceof PluginRuntimeApplicationError && !error.details.committed
? error.cause
: error;
return cause instanceof GatewayConfigReloadSupersededError;
}
export function startGatewayConfigReloader(
opts: GatewayConfigReloaderOptions,
): GatewayConfigReloader {
@ -246,6 +235,7 @@ export function startGatewayConfigReloader(
installRecords: PluginInstallRecords;
}
| undefined;
const pluginDrain = createConfigPluginDrainTracker();
const readPluginInstallRecords = opts.readPluginInstallRecords ?? readCurrentInstallRecords;
const appliedRevision = createConfigAppliedRevisionTracker({
onConfigApplied: opts.onConfigApplied,
@ -506,8 +496,8 @@ export function startGatewayConfigReloader(
);
await checkpoint();
assertCurrent();
const previousPluginInstallConfig = asPluginInstallConfig(currentPluginInstallRecords);
const nextPluginInstallConfig = asPluginInstallConfig(nextPluginInstallRecords);
const previousPluginInstallConfig = withPluginInstallRecords({}, currentPluginInstallRecords);
const nextPluginInstallConfig = withPluginInstallRecords({}, nextPluginInstallRecords);
const pluginInstallRecordChangedPaths = diffConfigPaths(
previousPluginInstallConfig,
nextPluginInstallConfig,
@ -694,6 +684,7 @@ export function startGatewayConfigReloader(
};
if (changedPaths.length === 0 && !pluginLifecycle) {
await commitReloadBaseline();
pluginDrain.applied();
publishedSource?.commit?.();
opts.onConfigRevisionApplied?.(nextConfigRevisionHash);
settleRuntimeApplication();
@ -757,6 +748,7 @@ export function startGatewayConfigReloader(
return completeApplication();
}
pluginDrain.assertCanApply(plan);
// No-op plans also publish the runtime snapshot before its applied receipt.
const applyRuntime = isNoopGatewayReloadPlan(plan) ? opts.onNoopConfigCommit : opts.onHotReload;
await opts.onConfigChange?.(plan, nextConfig);
@ -765,6 +757,7 @@ export function startGatewayConfigReloader(
applicationStatus = await applyRuntime(plan, nextConfig, ownership, nextSourceConfig);
} catch (error) {
ownership.rollbackRuntimeEnv();
pluginDrain.recordFailure(plan, error);
throw error;
}
await checkpoint();
@ -776,6 +769,7 @@ export function startGatewayConfigReloader(
typeof applicationStatus === "object" && applicationStatus.status === "applied"
? applicationStatus.runtime
: undefined;
pluginDrain.applied(plan, runtime);
if (runtime) {
completedPluginApplication = {
runtime,
@ -1061,8 +1055,10 @@ export function startGatewayConfigReloader(
}
if (superseded) {
opts.log.info(`config reload superseded: ${String(err)}`);
} else {
} else if (pluginDrain.shouldReport(err)) {
opts.log.error(`config reload failed: ${String(err)}`);
} else {
opts.log.info("config reload deferred: retry the failed plugin reload with --wait");
}
} finally {
running = false;

View file

@ -168,7 +168,7 @@ it.each(["success", "failure"] as const)(
nextConfig: first.cfgAtStart,
sourceConfig: first.cfgAtStart,
changedPaths: [],
prepareConfigEffects: () => async () => {},
prepareConfigEffects: () => ({ retire: () => {}, rollback: async () => {} }),
pluginLifecycle: {
reason: "reload",
operationId: "concurrent-bootstrap",

View file

@ -163,7 +163,7 @@ it.each(["active", "closing-memory"] as const)(
nextConfig: b.cfgAtStart,
sourceConfig: b.cfgAtStart,
changedPaths: [],
prepareConfigEffects: () => async () => {},
prepareConfigEffects: () => ({ retire: () => {}, rollback: async () => {} }),
pluginLifecycle: {
reason: "reload",
operationId: "queued-model-b",

View file

@ -28,7 +28,7 @@ class PluginAdmittedWorkTimeoutError extends Error {
constructor(pluginIds: ReadonlySet<string>, cause: PluginHostCleanupTimeoutError) {
const ids = [...pluginIds].join(", ");
super(
`plugin ${ids} admitted work did not settle within 60s; the previous plugin generation stays active. Use \`openclaw plugins reload ${[...pluginIds].join(" ")} --wait\` to wait until it finishes, or retry after it finishes.`,
`plugin ${ids} admitted work did not settle within 60s; the previous plugin generation stays active. Use \`openclaw plugins reload ${[...pluginIds].join(" ")} --wait\` to wait until it finishes.`,
{ cause },
);
}
@ -279,25 +279,31 @@ export function createPluginReloadCleanup({
pluginIds: ReadonlySet<string>,
signal: AbortSignal,
reportStatus: (status: GatewayPluginReloadStatus) => void,
includeConsumers = false,
{ includeConsumers = false, includeCalls = false } = {},
) => {
const instances = previousRegistry.plugins.flatMap((record) => {
const instance = pluginIds.has(record.id) && getPluginInstance(record);
return instance ? [instance] : [];
});
const count = instances.reduce((total, instance) => total + instance.retainedWorkCount, 0);
const count = instances.reduce(
(total, instance) =>
total + instance.retainedWorkCount + (includeCalls ? instance.ordinaryCallCount : 0),
0,
);
if (!count) {
return;
}
retainedWorkQueued = true;
const deadlineAtMs = admittedWorkDeadline();
const reason = `Plugin replacement queued behind ${count} retained work item(s); ${waitForDrain ? "waiting until they finish or the request is cancelled" : "applies when they finish within the 60s drain budget"}.`;
const reason = `Plugin replacement queued behind ${count} ${includeCalls ? "admitted" : "retained"} work item(s) for plugin ${[...pluginIds].join(", ")}; ${waitForDrain ? "waiting until they finish or the request is cancelled" : "applies when they finish within the 60s drain budget"}.`;
reportStatus({ phase: "reloading", pluginIds: [...changedPluginIds], deadlineAtMs, reason });
log.info(reason);
try {
await observeDrain("retained plugin work", signal, deadlineAtMs, (current) =>
Promise.all(
instances.map((instance) => instance.waitForRetainedWork(current, includeConsumers)),
instances.map((instance) =>
instance.waitForRetainedWork(current, { includeConsumers, includeCalls }),
),
),
);
} catch (error) {
@ -381,6 +387,30 @@ export function createPluginReloadCleanup({
return release;
},
drainInstances,
drainMemory: async (drain: () => Promise<{ errors: readonly unknown[] }>) => {
try {
const result = await drain();
for (const error of result.errors) {
const warning = `Memory cleanup failed: ${formatErrorMessage(error)}`;
log.warn(warning);
recordWarning(warning);
}
} catch (error) {
if (!(error instanceof PluginHostCleanupTimeoutError)) {
throw error;
}
log.warn(error.message);
recordWarning(error.message);
}
},
quiesceInstances: () => {
for (const record of previousRegistry.plugins) {
const instance = changedPluginIds.has(record.id) && getPluginInstance(record);
if (instance && instance.quiesce()) {
quiescedInstances.push(instance);
}
}
},
drainBeforeReplacement: async (
pluginIds: ReadonlySet<string>,
signal: AbortSignal,
@ -388,17 +418,14 @@ export function createPluginReloadCleanup({
assertCurrent: () => void,
) => {
// Sidecars release their capability consumers before finite work and callbacks drain.
await drainRetainedWork(pluginIds, signal, reportStatus, true);
await drainRetainedWork(pluginIds, signal, reportStatus, {
includeConsumers: true,
includeCalls: true,
});
assertCurrent();
if (retainedWorkQueued) {
recordWarning("Plugin replacement waited for retained work to finish.");
}
for (const record of previousRegistry.plugins) {
const instance = changedPluginIds.has(record.id) && getPluginInstance(record);
if (instance && instance.quiesce()) {
quiescedInstances.push(instance);
}
}
try {
if (pluginIds.size) {
await drainWithDeadline(pluginIds, signal, reportStatus, "reloading");

View file

@ -161,11 +161,13 @@ export async function verifyActiveCallDrainLease(
const instance = getPluginInstance(record);
assert(instance);
const drainStarted = createDeferredCore();
const drain = instance.drain.bind(instance);
const drainObservation = vi.spyOn(instance, "drain").mockImplementation((options) => {
drainStarted.resolve();
return drain(options);
});
const drain = instance.waitForRetainedWork.bind(instance);
const drainObservation = vi
.spyOn(instance, "waitForRetainedWork")
.mockImplementation((...options) => {
drainStarted.resolve();
return drain(...options);
});
const handler = fixture.previousRegistry.gatewayHandlers["first.call"];
assert(handler);
const invoke = (hold: boolean, respond: GatewayRequestHandlerOptions["respond"]) =>
@ -224,7 +226,7 @@ export async function verifyActiveCallDrainLease(
expect(fixture.owner.getReloadStatus()).toMatchObject({
phase: "reloading",
deadlineAtMs: expect.any(Number),
reason: expect.stringMatching(/admitted work.*first/),
reason: expect.stringContaining("plugin first"),
});
const deadlineAtMs = fixture.owner.getReloadStatus()?.deadlineAtMs;
assert(deadlineAtMs);
@ -242,7 +244,7 @@ export async function verifyActiveCallDrainLease(
pluginReload: {
phase: "reloading",
deadlineAtMs,
reason: expect.stringMatching(/admitted work.*first/),
reason: expect.stringContaining("plugin first"),
},
});
expect(fixture.candidates).toHaveLength(0);

View file

@ -26,10 +26,13 @@ export async function verifyCancelledDrainRollbackLease(
assert(reloadLease);
reloadLease.assertOwned();
},
prepareConfigEffects: () => async () => {
rollbackStarted.resolve();
await finishRollback.promise;
},
prepareConfigEffects: () => ({
retire: () => {},
rollback: async () => {
rollbackStarted.resolve();
await finishRollback.promise;
},
}),
});
const instance = getPluginInstance(fixture.previousRegistry.plugins[0]!);
assert(instance);

View file

@ -413,7 +413,7 @@ module.exports = { id: ${JSON.stringify(id)}, register(api) {
nextConfig,
sourceConfig,
changedPaths: [],
prepareConfigEffects: () => async () => {},
prepareConfigEffects: () => ({ retire: () => {}, rollback: async () => {} }),
assertInvokerOwned,
pluginLifecycle: {
reason,

View file

@ -338,10 +338,14 @@ it.each(["plugins.reload", "auth refresh"] as const)(
reason: "reload",
operationId: "borrow-reload",
},
prepareConfigEffects: () => {
markPreparedModelRuntimeSnapshotsStale("plugin reload", { waitForReplacement: true });
return async () => {};
},
prepareConfigEffects: () => ({
retire: () => {
markPreparedModelRuntimeSnapshotsStale("plugin reload", {
waitForReplacement: true,
});
},
rollback: async () => {},
}),
env,
commitRuntime: async (publication) => {
publication?.publish();

View file

@ -325,7 +325,9 @@ export async function createPluginReloadRecoveryFixture(
sourceConfig: nextConfig,
changedPaths,
checkpoint: options.checkpoint,
prepareConfigEffects: options.prepareConfigEffects ?? (() => rollbackConfigEffects),
prepareConfigEffects:
options.prepareConfigEffects ??
(() => ({ retire: () => {}, rollback: rollbackConfigEffects })),
pluginLifecycle: {
reason: "reload",
waitForDrain: options.waitForDrain,

View file

@ -254,7 +254,7 @@ async function verifySelfConsumerReload(
| "final checkpoint"
| "later replacement target",
) {
const prepareConfigEffects = vi.fn(() => async () => {});
const prepareConfigEffects = vi.fn(() => ({ retire: () => {}, rollback: async () => {} }));
let checkpoints = 0;
let consumer: PluginInstanceConsumer | undefined;
const fixture = await createFixture({
@ -334,7 +334,7 @@ async function verifyOverlappingRetainedWork(createRecoveryFixture: RecoveryFixt
abortOnCandidateStart: false,
prepareConfigEffects: () => {
reserved.resolve();
return async () => {};
return { retire: () => {}, rollback: async () => {} };
},
register(api, owner) {
if (owner !== "first") {
@ -422,20 +422,12 @@ async function verifyExplicitDrainWait(
: () => released.resolve();
const call = kind === "active call" ? instance.run(() => released.promise) : undefined;
const drainEntered = createDeferredCore();
const drain = instance.drain.bind(instance);
const waitForWork = instance.waitForRetainedWork.bind(instance);
const observation =
kind === "active call"
? vi.spyOn(instance, "drain").mockImplementation((...args) => {
const pending = drain(...args);
drainEntered.resolve();
return pending;
})
: vi.spyOn(instance, "waitForRetainedWork").mockImplementation((...args) => {
const pending = waitForWork(...args);
drainEntered.resolve();
return pending;
});
const observation = vi.spyOn(instance, "waitForRetainedWork").mockImplementation((...args) => {
const pending = waitForWork(...args);
drainEntered.resolve();
return pending;
});
let settled = false;
vi.useFakeTimers();
const reloading = fixture

View file

@ -257,13 +257,15 @@ it.each(["cold start", "hot enable"] as const)(
nextConfig,
sourceConfig: nextConfig,
changedPaths,
prepareConfigEffects: () => {
markPreparedModelRuntimeSnapshotsStale(
"prepared model runtime owner is stale before plugin drain",
{ waitForReplacement: true },
);
return async () => {};
},
prepareConfigEffects: () => ({
retire: () => {
markPreparedModelRuntimeSnapshotsStale(
"prepared model runtime owner is stale before plugin replacement",
{ waitForReplacement: true },
);
},
rollback: async () => {},
}),
env,
commitRuntime: async (publication) => {
publication?.publish();

View file

@ -329,7 +329,7 @@ export function registerPluginServiceRecoveryTests(createRecoveryFixture: Recove
expect(() => getPluginInstance(record)?.retainWork()).toThrow(
"replacement is in progress",
);
return rollback;
return { retire: () => {}, rollback };
},
});
const failure = new Error("fixture channel pause failed");
@ -358,7 +358,7 @@ export function registerPluginServiceRecoveryTests(createRecoveryFixture: Recove
}
});
const fixture = await createRecoveryFixture({
prepareConfigEffects: () => rollback,
prepareConfigEffects: () => ({ retire: () => {}, rollback }),
recoveryStart: async () => {
recoveryStarted.resolve();
await releaseRecovery.promise;

View file

@ -9,10 +9,7 @@ import { isTruthyEnvValue } from "../infra/env.js";
import { formatErrorMessage } from "../infra/errors.js";
import { prepareGatewayPluginMetadataSnapshotPublication } from "../plugins/current-plugin-metadata-snapshot.js";
import type { PluginHookGatewayCronService } from "../plugins/hook-gateway.types.js";
import {
PluginHostCleanupTimeoutError,
withPluginHostCleanupTimeout,
} from "../plugins/host-hook-cleanup-timeout.js";
import { withPluginHostCleanupTimeout } from "../plugins/host-hook-cleanup-timeout.js";
import {
createPluginRuntimeApplication,
getPluginRuntimeGeneration,
@ -147,8 +144,10 @@ export async function reloadGatewayPlugins(
reserveResourceHandoff,
selectResourceHandoff,
drainInstances,
drainMemory,
drainRetainedWork,
drainBeforeReplacement,
quiesceInstances,
resumeInstances,
drainForRecovery,
disposeInstances,
@ -250,16 +249,23 @@ export async function reloadGatewayPlugins(
assertCurrent();
// Reserve and gate new model runs atomically; admitted runs keep their callbacks until settled.
releaseResourceHandoff = reserveResourceHandoff(resourceHandoffIds);
rollbackConfigEffects = params.prepareConfigEffects({
const configEffects = params.prepareConfigEffects({
pluginIds: changedPluginIds,
channels: channelTargets,
});
rollbackConfigEffects = configEffects.rollback;
phase = "drain";
replacement.setReloadStatus({ phase: "reloading", pluginIds: [...changedPluginIds] });
await drainRetainedWork(resourceHandoffIds, drainSignal, replacement.setReloadStatus);
assertCurrent();
channels.pause();
decisionReplacement = prepareDecisionProviderReload(previousRegistry, changedPluginIds);
channels.pause();
quiesceInstances();
await drainRetainedWork(resourceHandoffIds, drainSignal, replacement.setReloadStatus, {
includeCalls: true,
});
assertCurrent();
configEffects.retire();
for (const sidecar of runtimeState.gatewayLifetimeSidecars.snapshot()) {
const prepared = sidecar.preparePluginReload?.({
previousRegistry,
@ -273,24 +279,11 @@ export async function reloadGatewayPlugins(
}
}
memoryReplacement = prepareMemoryRuntimeReload(previousRegistry, nextRegistry);
// Consumers release their handles while the producing instance is callable.
// Retained consumers and cleanup calls remain usable after ordinary admission closes.
for (const sidecar of sidecarReplacements) {
await sidecar.drain();
}
try {
const result = await memoryReplacement.drain();
for (const error of result.errors) {
const warning = `Memory cleanup failed: ${formatErrorMessage(error)}`;
log.warn(warning);
recordWarning(warning);
}
} catch (error) {
if (!(error instanceof PluginHostCleanupTimeoutError)) {
throw error;
}
log.warn(error.message);
recordWarning(error.message);
}
await drainMemory(memoryReplacement.drain);
await runtimeState.discovery?.update({
gatewayDiscoveryServices: previousRegistry.gatewayDiscoveryServices.filter(
(entry) => !changedPluginIds.has(entry.pluginId),

View file

@ -0,0 +1,143 @@
import assert from "node:assert/strict";
import fs from "node:fs/promises";
import { afterEach, expect, it, vi } from "vitest";
import { useAutoCleanupTempDirTracker } from "../../test/helpers/temp-dir.js";
import { getPluginInstance } from "../plugins/plugin-instance-scope.js";
import { getActivePluginRegistry } from "../plugins/runtime.js";
import { createDeferredCore } from "../shared/deferred.js";
import { acquireTestPortBlock } from "../test-utils/port-claims.js";
import { clearInstanceBindingProbeCoordinators } from "./server-plugins.lifecycle.test-fixtures.js";
import {
installInstanceBindingConfigIo,
prepareInstanceBindingFixture,
} from "./server-plugins.lifecycle.test-support.js";
import {
connectWebchatClient,
installGatewayTestHooks,
rpcReq,
startTestGatewayServer,
} from "./test-helpers.server.js";
vi.doUnmock("../plugins/loader.js");
installGatewayTestHooks({ scope: "suite" });
const tempDirs = useAutoCleanupTempDirTracker(afterEach);
installInstanceBindingConfigIo();
it("serves active model and chat metadata throughout an admitted plugin call drain", async () => {
const fixture = await prepareInstanceBindingFixture(tempDirs.make("openclaw-drain-readers-"));
const entered = createDeferredCore();
const release = createDeferredCore();
fixture.coordinator.heldCall = { entered: entered.resolve, completion: release.promise };
const config = JSON.parse(await fs.readFile(fixture.configPath, "utf8"));
config.agents = {
defaults: {
model: { primary: "openai/gpt-reader-fixture" },
models: { "openai/gpt-reader-fixture": {} },
},
entries: { main: { default: true } },
};
config.models = {
providers: {
openai: {
baseUrl: "https://openai.example.com/v1",
models: [{ id: "gpt-reader-fixture", name: "Reader fixture" }],
},
},
};
await fs.writeFile(fixture.configPath, JSON.stringify(config));
const claim = await acquireTestPortBlock({ offsets: [0, 1, 2, 3, 4] });
const server = await startTestGatewayServer(claim, {
auth: { mode: "none" },
controlUiEnabled: false,
sidecarStartup: "start",
});
let socket: Awaited<ReturnType<typeof connectWebchatClient>> | undefined;
let held: ReturnType<typeof rpcReq> | undefined;
let reloading: ReturnType<typeof rpcReq> | undefined;
const drainObservations: Array<{ mockRestore: () => void }> = [];
try {
await server.startupSettled;
socket = await connectWebchatClient({ port: claim.port, scopes: ["operator.admin"] });
const connected = socket;
const reads = async () =>
await Promise.all(
["models.list", "chat.metadata"].map(async (method) => {
const start = performance.now();
const response = await rpcReq(connected, method, { agentId: "main" }, 120_000);
expect(response.ok, `${method}: ${JSON.stringify(response)}`).toBe(true);
return { method, durationMs: performance.now() - start, payload: response.payload };
}),
);
const before = await reads();
const registry = getActivePluginRegistry();
const record = registry?.plugins.find((plugin) => plugin.id === "instance-binding-probe");
assert(record);
const instance = getPluginInstance(record);
assert(instance);
const draining = createDeferredCore();
const wait = instance.waitForRetainedWork.bind(instance);
drainObservations.push(
vi.spyOn(instance, "waitForRetainedWork").mockImplementation((...args) => {
const pending = wait(...args);
draining.resolve();
return pending;
}),
);
const drain = instance.drain.bind(instance);
drainObservations.push(
vi.spyOn(instance, "drain").mockImplementation((...args) => {
const pending = drain(...args);
draining.resolve();
return pending;
}),
);
held = rpcReq(connected, "instanceBinding.hold", {}, 120_000);
await entered.promise;
vi.useFakeTimers({ toFake: ["setTimeout", "clearTimeout", "Date"] });
let reloadSettled = false;
reloading = rpcReq(
connected,
"plugins.reload",
{ plugins: [{ pluginId: "instance-binding-probe" }] },
120_000,
).then((result) => {
reloadSettled = true;
return result;
});
await draining.promise;
expect(instance.acceptingCalls).toBe(false);
const during = await reads();
expect(reloadSettled).toBe(false);
expect(getActivePluginRegistry()).toBe(registry);
expect(during.map((entry) => entry.payload)).toEqual(before.map((entry) => entry.payload));
await vi.advanceTimersByTimeAsync(60_000);
vi.useRealTimers();
expect(await reloading).toMatchObject({
ok: false,
error: { details: { runtime: { committed: false, phase: "drain" } } },
});
expect(getActivePluginRegistry()).toBe(registry);
const after = await reads();
console.info(
"PLUGIN_DRAIN_READER_LATENCY " +
JSON.stringify({
before: before.map(({ method, durationMs }) => ({ method, durationMs })),
during: during.map(({ method, durationMs }) => ({ method, durationMs })),
after: after.map(({ method, durationMs }) => ({ method, durationMs })),
}),
);
release.resolve();
expect((await held).ok).toBe(true);
} finally {
vi.useRealTimers();
release.resolve();
await Promise.allSettled([held, reloading]);
for (const observation of drainObservations) {
observation.mockRestore();
}
socket?.close();
await server.close({ reason: "plugin drain reader fixture complete" });
clearInstanceBindingProbeCoordinators();
delete process.env.OPENCLAW_TEST_TRUST_BUNDLED_PLUGINS_DIR;
}
}, 60_000);

View file

@ -59,6 +59,7 @@ export type InstanceBindingProbeCoordinator = {
channelIds?: readonly string[];
channelStops?: Array<Pick<ChannelBindingMonitor, "channelId" | "runtimeId" | "abortSignal">>;
channelCleanup?: Map<ChannelBindingMonitor, { release: () => void; finished: Promise<void> }>;
heldCall?: { entered: () => void; completion: Promise<void> };
};
export async function withPluginServiceStopDeadline<T>(
@ -172,6 +173,13 @@ export async function writeInstanceBindingProbePlugin(
const coordinator = request.coordinator;
const reportReloadSettlement = Boolean(coordinator.reportReloadSettlement || coordinator.channelProof || coordinator.channel);
const registryId = coordinator.nextRegistryId++;
if (coordinator.heldCall) {
api.registerGatewayMethod("instanceBinding.hold", async ({ respond }) => {
coordinator.heldCall.entered();
await coordinator.heldCall.completion;
respond(true, { registryId });
}, { scope: "operator.read" });
}
coordinator.runtimes.push(api.runtime);
coordinator.registrationModes.push(api.registrationMode);
if (coordinator.contextEngineId) {

View file

@ -172,11 +172,11 @@ export type GatewayReloadHandlerParams = {
changedPaths: readonly string[];
reloadPluginIds?: ReadonlySet<string>;
pluginLifecycle?: GatewayReloadPlan["pluginLifecycle"];
/** Fence config consumers before drain; return their publication after successful rollback. */
/** Fence execution before drain; retire active facts only when replacement can begin. */
prepareConfigEffects: (replacement: {
pluginIds: ReadonlySet<string>;
channels: ReadonlySet<ChannelKind>;
}) => () => Promise<void>;
}) => { retire: () => void; rollback: () => Promise<void> };
commitRuntime: (publication?: GatewayRuntimePublication) => Promise<void>;
env: NodeJS.ProcessEnv;
isAborted?: () => boolean;

View file

@ -11,7 +11,10 @@ import {
} from "../state/openclaw-state-db.js";
import type { GatewayReloadPlan } from "./config-reload-plan.js";
import type { GatewayCronState } from "./server-cron.js";
import type { ManagedGatewayConfigReloaderParams } from "./server-reload-contracts.js";
import type {
GatewayPluginReloadResult,
ManagedGatewayConfigReloaderParams,
} from "./server-reload-contracts.js";
export function createMonitorPublicationFailure() {
const database = openOpenClawStateDatabase();
@ -46,6 +49,16 @@ export function createMonitorPublicationFailure() {
type ConfigWriteListener = (event: ConfigWriteNotification) => void;
type ConfigWriteListenerRef = { current: ConfigWriteListener | null };
export function makePluginReloadResult(
overrides: Partial<GatewayPluginReloadResult> = {},
): GatewayPluginReloadResult {
return {
runtime: { operationId: "test-reload", generation: 1, pluginIds: [] },
activeChannels: new Set(),
...overrides,
};
}
export function enableChannelReloadsForTest() {
const previousSkipChannels = process.env.OPENCLAW_SKIP_CHANNELS;
const previousSkipProviders = process.env.OPENCLAW_SKIP_PROVIDERS;

View file

@ -132,6 +132,7 @@ import {
createTestCronState,
createValidConfigSnapshot,
enableChannelReloadsForTest,
makePluginReloadResult,
publishConfigWrite,
} from "./server-reload-handlers.config.test-support.js";
import { createGatewayReloadHandlers as createGatewayReloadHandlersImpl } from "./server-reload-hot.js";
@ -215,7 +216,7 @@ function createGatewayReloadHandlers(
pruneInactiveChannelAccountState: vi.fn(),
stopPostReadySidecars: vi.fn(),
reloadPlugins: vi.fn<ReloadHandlerParams["reloadPlugins"]>(async ({ prepareConfigEffects }) => {
prepareConfigEffects({ pluginIds: new Set(), channels: new Set() });
prepareConfigEffects({ pluginIds: new Set(), channels: new Set() }).retire();
return makePluginReloadResult();
}),
logHooks: { info: vi.fn(), warn: vi.fn(), error: vi.fn() },
@ -252,7 +253,7 @@ function startManagedGatewayConfigReloader(params: ManagedReloaderTestParams) {
startChannel: vi.fn(async () => new Map()),
stopChannel: vi.fn(async () => {}),
reloadPlugins: vi.fn<ReloadHandlerParams["reloadPlugins"]>(async ({ prepareConfigEffects }) => {
prepareConfigEffects({ pluginIds: new Set(), channels: new Set() });
prepareConfigEffects({ pluginIds: new Set(), channels: new Set() }).retire();
return makePluginReloadResult();
}),
logHooks: { info: vi.fn(), warn: vi.fn(), error: vi.fn() },
@ -403,6 +404,7 @@ vi.mock("../agents/model-catalog.js", () => ({
vi.mock("../agents/prepared-model-runtime.js", () => ({
advancePreparedModelRuntimeConfig: hoisted.advancePreparedModelRuntimeConfig,
beginPreparedModelRuntimePluginDrain: () => ({ pendingPublication: false, release: () => {} }),
markPreparedModelRuntimeSnapshotsStale: (
reason?: string,
options?: { waitForReplacement?: boolean; preserveReplacementWait?: boolean },
@ -433,7 +435,7 @@ vi.mock("../agents/agent-bundle-mcp-tools.js", () => ({
reloadSessionMcpRuntimes: hoisted.reloadSessionMcpRuntimes,
}));
vi.mock("../plugins/installed-plugin-index-records.js", () => ({
vi.mock("../plugins/installed-plugin-index-record-reader.js", () => ({
clearLoadInstalledPluginIndexInstallRecordsCache: vi.fn(),
loadInstalledPluginIndexInstallRecords: vi.fn(async () => ({})),
loadInstalledPluginIndexInstallRecordsSync: vi.fn(() => ({})),
@ -536,16 +538,6 @@ async function withReloadChannelManager(
}
}
function makePluginReloadResult(
overrides: Partial<GatewayPluginReloadResult> = {},
): GatewayPluginReloadResult {
return {
runtime: { operationId: "test-reload", generation: 1, pluginIds: [] },
activeChannels: new Set(),
...overrides,
};
}
function createTestCronReconciliation() {
const complete = vi.fn<() => Promise<void>>(async () => {});
return {
@ -1192,7 +1184,7 @@ async function withManagedChannelSecretFixture(
startChannel: manager.startChannel,
stopChannel: manager.stopChannel,
reloadPlugins: async ({ commitRuntime, prepareConfigEffects }) => {
prepareConfigEffects({ pluginIds: new Set(["notes"]), channels: new Set() });
prepareConfigEffects({ pluginIds: new Set(["notes"]), channels: new Set() }).retire();
await commitRuntime();
return makePluginReloadResult({
activeChannels: new Set(["mattermost"]),
@ -2570,7 +2562,10 @@ describe("gateway targeted service reload", () => {
getPluginRegistry: () => registry,
reloadPluginServices,
reloadPlugins: async ({ commitRuntime, prepareConfigEffects }) => {
prepareConfigEffects({ pluginIds: new Set(runtime.pluginIds), channels: new Set() });
prepareConfigEffects({
pluginIds: new Set(runtime.pluginIds),
channels: new Set(),
}).retire();
await commitRuntime();
events.push("replaced-plugin");
if (outcome === "plugin failure") {
@ -5775,10 +5770,12 @@ describe("gateway plugin hot reload handlers", () => {
const channels = { start: vi.fn(async () => new Map()), stop: vi.fn(async () => {}) };
const pruneInactiveChannelAccountState = vi.fn();
const reloadPlugins = vi.fn<ReloadHandlerParams["reloadPlugins"]>(async (params) => {
params.prepareConfigEffects({
pluginIds: new Set(runtime.pluginIds),
channels: new Set(["discord"]),
});
params
.prepareConfigEffects({
pluginIds: new Set(runtime.pluginIds),
channels: new Set(["discord"]),
})
.retire();
await params.commitRuntime();
current = false;
return makePluginReloadResult({

View file

@ -3,6 +3,7 @@ import { tryResolveConfiguredAgentWorkspaceDir } from "../agents/agent-scope-con
import { refreshContextWindowCache } from "../agents/context.js";
import {
advancePreparedModelRuntimeConfig,
beginPreparedModelRuntimePluginDrain,
markPreparedModelRuntimeSnapshotsStale,
rejectPendingPreparedModelRuntimeReplacement,
type PreparedModelRuntimeReplacementGateId,
@ -202,23 +203,40 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams)
};
assertIrreversibleReloadPlanHasRecoveryOwner(remainingPlan, restartRecoveryAvailable);
const previousConfig = getRuntimeConfig();
// Drain revokes plugin calls before commit; new and unfinished model preparation must wait.
preparedModelRuntimeReplacementGateId = markPreparedModelRuntimeSnapshotsStale(
"prepared model runtime owner is stale before plugin drain",
{ waitForReplacement: true, ...modelRuntimeRefreshScope },
);
return async () => {
await withPluginRuntimeRegistryScope(params.getPluginRegistry(), () =>
mrReload.refreshModelRuntimeAfterHotReload({
config: previousConfig,
agentIds: modelRuntimeAgentIds,
pluginMetadataSnapshot: params.getPluginMetadataSnapshot?.(),
isPublicationCurrent: () =>
isCurrentGatewayReloadGeneration(myGeneration) &&
!isLifecycleReloadAborted() &&
!isRestartRetryStopped(),
}),
const drain = beginPreparedModelRuntimePluginDrain();
releasePreparedModelRuntimeDrain = drain.release;
const retire = () => {
if (preparedModelRuntimeReplacementGateId) {
return;
}
preparedModelRuntimeReplacementGateId = markPreparedModelRuntimeSnapshotsStale(
"prepared model runtime owner is stale before plugin replacement",
{ waitForReplacement: true, ...modelRuntimeRefreshScope },
);
drain.release();
};
// Unfinished publication must cancel its acquisition before the plugin owner drains it.
if (drain.pendingPublication) {
retire();
}
return {
retire,
rollback: async () => {
if (!preparedModelRuntimeReplacementGateId) {
return;
}
await withPluginRuntimeRegistryScope(params.getPluginRegistry(), () =>
mrReload.refreshModelRuntimeAfterHotReload({
config: previousConfig,
agentIds: modelRuntimeAgentIds,
pluginMetadataSnapshot: params.getPluginMetadataSnapshot?.(),
isPublicationCurrent: () =>
isCurrentGatewayReloadGeneration(myGeneration) &&
!isLifecycleReloadAborted() &&
!isRestartRetryStopped(),
}),
);
},
};
};
let activePluginChannelsAfterReload: ReadonlySet<ChannelKind> | null = null;
@ -241,6 +259,7 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams)
isRestartRetryStopped() ||
isLifecycleReloadAborted();
let preparedModelRuntimeReplacementGateId: PreparedModelRuntimeReplacementGateId | undefined;
let releasePreparedModelRuntimeDrain: (() => void) | undefined;
let recoveryRestartScheduled = false;
const laneConcurrency = resolveGatewayLaneConcurrency(nextConfig);
// Use one candidate env snapshot before publication and through later channel starts.
@ -594,6 +613,8 @@ export function createGatewayReloadHandlers(params: GatewayReloadHandlerParams)
rejectPendingPreparedModelRuntimeReplacement(preparedModelRuntimeReplacementGateId, error);
}
throw error;
} finally {
releasePreparedModelRuntimeDrain?.();
}
try {
await commitRuntime();

View file

@ -50,6 +50,7 @@ import { createTestRuntimeSecretsActivator } from "./server-startup-config.test-
// real reload transaction and async owners, not a cold model/plugin runtime.
vi.mock("../agents/prepared-model-runtime.js", () => ({
advancePreparedModelRuntimeConfig: vi.fn(),
beginPreparedModelRuntimePluginDrain: () => ({ pendingPublication: false, release: () => {} }),
markPreparedModelRuntimeSnapshotsStale: vi.fn(),
rejectPendingPreparedModelRuntimeReplacement: vi.fn(),
refreshPreparedModelRuntimeSnapshots: vi.fn(async () => {}),
@ -151,7 +152,7 @@ function startManagedGatewayConfigReloader(
startChannel: vi.fn(async () => new Map()),
stopChannel: vi.fn(async () => {}),
reloadPlugins: vi.fn(async ({ prepareConfigEffects }) => {
prepareConfigEffects({ pluginIds: new Set(), channels: new Set() });
prepareConfigEffects({ pluginIds: new Set(), channels: new Set() }).retire();
return {
runtime: { operationId: "test-reload", generation: 1, pluginIds: [] },
activeChannels: new Set(),

View file

@ -56,11 +56,12 @@ vi.mock("../agents/context.js", () => ({
}));
vi.mock("../agents/prepared-model-runtime.js", () => ({
advancePreparedModelRuntimeConfig: vi.fn(),
beginPreparedModelRuntimePluginDrain: () => ({ pendingPublication: false, release: () => {} }),
markPreparedModelRuntimeSnapshotsStale: vi.fn(() => Symbol("model-replacement")),
rejectPendingPreparedModelRuntimeReplacement: vi.fn(),
refreshPreparedModelRuntimeSnapshots: vi.fn(async () => {}),
}));
vi.mock("../plugins/installed-plugin-index-records.js", () => ({
vi.mock("../plugins/installed-plugin-index-record-reader.js", () => ({
clearLoadInstalledPluginIndexInstallRecordsCache: vi.fn(),
loadInstalledPluginIndexInstallRecords: vi.fn(async () => ({})),
loadInstalledPluginIndexInstallRecordsSync: vi.fn(() => ({})),
@ -209,7 +210,12 @@ describe("channel reload with a retained wizard waiting for the lifecycle lease"
const releasePlugin = createDeferred();
const reloadPlugins = vi.fn<ManagedGatewayConfigReloaderParams["reloadPlugins"]>(
async (params) => {
params.prepareConfigEffects({ pluginIds: new Set(["fixture"]), channels: new Set() });
params
.prepareConfigEffects({
pluginIds: new Set(["fixture"]),
channels: new Set(),
})
.retire();
await params.commitRuntime();
pluginCommitted.resolve();
await releasePlugin.promise;

View file

@ -397,7 +397,7 @@ it.each([
await expect(reload.close()).resolves.toMatchObject({
errors: [
expect.objectContaining({
message: expect.stringContaining(`Plugin ${targetId} was reloaded or disabled`),
message: `Plugin ${targetId} is retiring`,
}),
],
});

View file

@ -16,6 +16,7 @@ import {
setStandaloneMemoryManagerActive,
} from "./memory-state.js";
import { getPluginValueInstance, runPluginCleanup } from "./plugin-instance-scope.js";
import { runPluginCleanupScope } from "./plugin-invocation-scope.js";
import type {
MemoryPluginRuntime,
RegisteredMemorySearchManager,
@ -356,19 +357,23 @@ export function prepareMemoryRuntimeReload(
// the Gateway owner. Final shutdown must join close() before disposing shared state.
const close = () => {
if (!cleanup) {
cleanup = Promise.allSettled(
prepared.map(({ runtime, handle }) =>
Promise.resolve().then(() =>
runPluginCleanup(runtime, async () => {
// Admission stays outside this catch; only admitted teardown reports faults.
try {
return await handle.drain();
} catch (error) {
return { errors: [error] };
}
}),
cleanup = runPluginCleanupScope(
[...prepared.map(({ runtime }) => runtime), ...retiringEmbeddingProviders],
() =>
Promise.allSettled(
prepared.map(({ runtime, handle }) =>
Promise.resolve().then(() =>
runPluginCleanup(runtime, async () => {
// Admission stays outside this catch; only admitted teardown reports faults.
try {
return await handle.drain();
} catch (error) {
return { errors: [error] };
}
}),
),
),
),
),
).then((results) => {
const failures = results.flatMap((result) =>
result.status === "rejected" ? [result.reason] : [],

View file

@ -23,7 +23,11 @@ export interface PluginInstanceHandle extends PluginInvocationInstance, PluginIn
adopt<T>(value: T): T;
retainWork(): () => void;
readonly retainedWorkCount: number;
waitForRetainedWork(signal: AbortSignal, includeConsumers?: boolean): Promise<void>;
readonly ordinaryCallCount: number;
waitForRetainedWork(
signal: AbortSignal,
options?: { includeConsumers?: boolean; includeCalls?: boolean },
): Promise<void>;
reserveReplacement(): () => void;
retainConsumer(
invoke?: <T>(run: () => T) => T,
@ -45,6 +49,7 @@ export type PluginInvocationBinding = {
};
export type PluginInvocationContext = {
assertCurrent?: (instance: PluginInstanceHandle) => void;
lookup: (instance: PluginInstanceHandle) => PluginInvocationBinding | undefined;
};

View file

@ -0,0 +1,32 @@
import { toErrorObject } from "../infra/errors.js";
/** Observers never revoke admitted work; cancellation removes only their waiter. */
export async function waitForPluginInstanceSettlement(
pluginId: string,
waiters: Set<() => void>,
settled: () => boolean,
signal: AbortSignal,
): Promise<void> {
signal.throwIfAborted();
if (settled()) {
return;
}
await new Promise<void>((resolve, reject) => {
const cleanup = () => {
waiters.delete(wake);
signal.removeEventListener("abort", abort);
};
const wake = () => {
if (settled()) {
cleanup();
resolve();
}
};
const abort = () => {
cleanup();
reject(toErrorObject(signal.reason, `Plugin ${pluginId} work drain aborted`));
};
waiters.add(wake);
signal.addEventListener("abort", abort, { once: true });
});
}

View file

@ -1,4 +1,4 @@
import { formatErrorMessage, toErrorObject } from "../infra/errors.js";
import { formatErrorMessage } from "../infra/errors.js";
import { createSubsystemLogger } from "../logging/subsystem.js";
import { AsyncWorkScope, trackAsyncWork } from "../shared/async-work-scope.js";
import { createDeferredCore } from "../shared/deferred.js";
@ -16,6 +16,7 @@ import {
resolvePluginInstanceOwner,
type PluginInstanceOwner,
} from "./plugin-instance-scope.js";
import { waitForPluginInstanceSettlement } from "./plugin-instance-settlement.js";
import { createPluginValueView } from "./plugin-instance-value-views.js";
import type {
PluginInstanceCallLease,
@ -214,37 +215,23 @@ export class PluginInstance {
);
}
/** Observe host settlement without closing the callbacks that retained runs still need. */
async waitForRetainedWork(signal: AbortSignal, includeConsumers = true): Promise<void> {
await this.waitForSettlement(
() => (includeConsumers ? this.retainedWorkCount : this.retainedWork.size) === 0,
signal,
);
get ordinaryCallCount(): number {
return [...this.calls.values()].filter((call) => !call.cleanup).length;
}
private async waitForSettlement(settled: () => boolean, signal: AbortSignal): Promise<void> {
signal.throwIfAborted();
if (settled()) {
return;
}
await new Promise<void>((resolve, reject) => {
const cleanup = () => {
this.waiters.delete(wake);
signal.removeEventListener("abort", abort);
};
const wake = () => {
if (settled()) {
cleanup();
resolve();
}
};
const abort = () => {
cleanup();
reject(toErrorObject(signal.reason, `Plugin ${this.pluginId} work drain aborted`));
};
this.waiters.add(wake);
signal.addEventListener("abort", abort, { once: true });
});
/** Observe admitted work without closing callbacks or idle publication custody. */
async waitForRetainedWork(
signal: AbortSignal,
{ includeConsumers = true, includeCalls = false } = {},
): Promise<void> {
await waitForPluginInstanceSettlement(
this.pluginId,
this.waiters,
() =>
(!includeCalls || this.ordinaryCallCount === 0) &&
(includeConsumers ? this.retainedWorkCount : this.retainedWork.size) === 0,
signal,
);
}
/** Reserve replacement atomically before host owners invalidate or stop this instance. */
@ -396,6 +383,7 @@ export class PluginInstance {
}
private enter<T>(token: object, run: () => T): T {
pluginInvocationContext.getStore()?.assertCurrent?.(this);
const current = invocation.getStore();
const call =
current?.instance === this && current.token === token ? current : { instance: this, token };
@ -542,7 +530,9 @@ export class PluginInstance {
try {
await this.waitForCalls(ownToken, options?.signal);
if (options?.includeConsumers) {
await this.waitForSettlement(
await waitForPluginInstanceSettlement(
this.pluginId,
this.waiters,
() => this.consumers.size === 0,
options.signal ?? new AbortController().signal,
);
@ -558,7 +548,7 @@ export class PluginInstance {
const settled = () => [...this.calls.keys()].every((token) => token === ownToken);
if (signal) {
// Reload owns its observation budget; disposal keeps its independent deadline.
return this.waitForSettlement(settled, signal);
return waitForPluginInstanceSettlement(this.pluginId, this.waiters, settled, signal);
}
const deadline = new AbortController();
const timer = setTimeout(() => {
@ -567,7 +557,7 @@ export class PluginInstance {
);
}, SHUTDOWN_TIMEOUT_MS);
try {
await this.waitForSettlement(settled, deadline.signal);
await waitForPluginInstanceSettlement(this.pluginId, this.waiters, settled, deadline.signal);
} finally {
clearTimeout(timer);
}

View file

@ -0,0 +1,39 @@
import { expect, it, vi } from "vitest";
import { createDeferredCore } from "../shared/deferred.js";
import { PluginInstanceUnavailableError } from "./plugin-instance-error.js";
import { PluginInstance } from "./plugin-instance.js";
import { runPluginCleanupScope } from "./plugin-invocation-scope.js";
it.each(["idle", "active"])(
"bounds %s cross-plugin teardown authority to its host cleanup",
async (mode) => {
const first = new PluginInstance("first-cleanup");
const second = new PluginInstance("second-cleanup");
const unrelated = new PluginInstance("unrelated-cleanup");
const called = vi.fn();
const callback = second.wrap(called);
const forbidden = unrelated.wrap(called);
const service = first.wrap({ stop: async () => callback() });
for (const instance of [first, second, unrelated]) {
instance.quiesce();
}
const delayed = createDeferredCore();
let late: Promise<void> | undefined;
await runPluginCleanupScope([service, callback], async () => {
await service.stop();
expect(() => forbidden()).toThrow(PluginInstanceUnavailableError);
const resume = async () => {
await delayed.promise;
callback();
};
late = mode === "active" ? second.wrap(resume)() : resume();
});
const refused = expect(late).rejects.toThrow("Plugin cleanup scope is closed");
delayed.resolve();
await refused;
expect(called).toHaveBeenCalledOnce();
for (const instance of [first, second, unrelated]) {
await instance.dispose();
}
},
);

View file

@ -10,6 +10,58 @@ import {
import type { PluginInstanceConsumer } from "./plugin-instance.types.js";
import type { PluginRegistry } from "./registry-types.js";
/** Host teardown may cross exact plugin owners after ordinary admission closes. */
export async function runPluginCleanupScope<T>(values: readonly object[], run: () => Promise<T>) {
const instances = new Set(
values.map(getPluginValueInstance).filter((instance) => instance !== undefined),
);
const parent = pluginInvocationContext.getStore();
let closed = false;
const assertOpen = () => {
if (closed) {
throw new Error("Plugin cleanup scope is closed");
}
};
const bindings = new Map(
[...instances].map((instance) => [
instance,
{
run: <R>(operation: () => R) => {
assertOpen();
return instance.runCleanup(operation);
},
wrap: <R>(value: R) => {
assertOpen();
return instance.wrap(value);
},
},
]),
);
try {
return await pluginInvocationContext.run(
{
assertCurrent: (instance) => {
if (bindings.has(instance)) {
assertOpen();
} else {
parent?.assertCurrent?.(instance);
}
},
lookup: (instance) => {
const binding = bindings.get(instance);
if (binding) {
assertOpen();
}
return binding ?? parent?.lookup(instance);
},
},
run,
);
} finally {
closed = true;
}
}
/** Finite execution custody for one host-selected registry and its exact instances. */
export class PluginInvocationScope {
private readonly bindings = new Map<PluginInstanceHandle, PluginInvocationBinding>();

View file

@ -909,8 +909,7 @@ try:
reaped.append({"pid": reaped_pid, "status": reaped_status})
if group_present():
raise RuntimeError("renderer group remains after exact leaf reap")
wait(lambda: controller.poll() is not None)
report["controllerCode"] = controller.returncode
report["controllerCode"] = controller.wait(timeout=max(0, min(5, deadline - time.monotonic())))
except BaseException as error:
report["fixtureError"] = type(error).__name__ + ": " + str(error).replace(str(root), "<fixture>")
finally:

View file

@ -17,6 +17,7 @@ export const gatewayDatabaseWorkerTestFiles = [
"src/gateway/chat-display-projection.cron.test.ts",
"src/gateway/config-reload.activation.integration.test.ts",
"src/gateway/config-reload.lease-retry.test.ts",
"src/gateway/config-reload.plugin-drain.test.ts",
"src/gateway/config-reload.plugin-observation.test.ts",
"src/gateway/config-reload.test.ts",
"src/gateway/config-reload.transcripts.test.ts",