Skip to content

Commit 11e9042

Browse files
fix(server): finish a plugin's host work before disable returns
When a plugin's process had already exited, disable returned at once, even while the exit was still ending that process's host calls. Disable now waits for that, so their cleanup cannot overlap a re-enable or what runs after the disable. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
1 parent dd543f8 commit 11e9042

2 files changed

Lines changed: 65 additions & 3 deletions

File tree

‎apps/server/src/plugins/PluginSettings.test.ts‎

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -726,6 +726,52 @@ it.layer(NodeServices.layer)("PluginSettings", (it) => {
726726
}),
727727
);
728728

729+
it.effect("waits for host work still ending when disabling a process that exited", () =>
730+
Effect.gen(function* () {
731+
const supervisor = yield* makeSupervisor(yield* Scope.Scope, BIN_PATH);
732+
const started = yield* Deferred.make<void>();
733+
const ending = yield* Deferred.make<void>();
734+
const release = yield* Deferred.make<void>();
735+
const order: Array<string> = [];
736+
yield* supervisor.serveHostMethod("settings.get", () => Effect.succeed({ value: null }));
737+
yield* supervisor.serveHostMethod("storage.get", () =>
738+
Deferred.succeed(started, undefined).pipe(
739+
Effect.andThen(Effect.never),
740+
// Stands for cleanup that outlasts the process, such as releasing a lock.
741+
Effect.onInterrupt(() =>
742+
Deferred.succeed(ending, undefined).pipe(
743+
Effect.andThen(Deferred.await(release)),
744+
Effect.andThen(Effect.sync(() => order.push("host work ended"))),
745+
),
746+
),
747+
),
748+
);
749+
const registration = yield* loadPluginDirectory(yield* preparePlugin());
750+
const pluginId = registration.manifest.id;
751+
yield* supervisor.enable(registration);
752+
yield* supervisor
753+
.invoke(pluginId, "load", { key: "a" })
754+
.pipe(Effect.ignore, Effect.forkChild({ startImmediately: true }));
755+
yield* Deferred.await(started);
756+
yield* supervisor
757+
.invoke(pluginId, "exit", null)
758+
.pipe(Effect.ignore, Effect.forkChild({ startImmediately: true }));
759+
// The process has exited and its host work is ending.
760+
yield* Deferred.await(ending);
761+
762+
const disabling = yield* supervisor
763+
.disable(pluginId)
764+
.pipe(
765+
Effect.andThen(Effect.sync(() => order.push("disabled"))),
766+
Effect.forkChild({ startImmediately: true }),
767+
);
768+
yield* Effect.yieldNow;
769+
yield* Deferred.succeed(release, undefined);
770+
yield* Fiber.join(disabling);
771+
expect(order).toEqual(["host work ended", "disabled"]);
772+
}),
773+
);
774+
729775
it.effect("drops a write that waited for the settings lock past its generation", () =>
730776
withDatabase(
731777
Effect.gen(function* () {

‎apps/server/src/plugins/PluginSupervisor.ts‎

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -278,6 +278,8 @@ interface Child {
278278
hostCalls: number;
279279
/** Owns the fibers serving this child's host calls; closed on revocation and exit. */
280280
readonly hostWork: Scope.Closeable;
281+
/** Set once `hostWork` starts closing; done once it has closed. */
282+
hostWorkEnded: Deferred.Deferred<void> | undefined;
281283
/** One host-call answer at a time waits for the child to read earlier messages. */
282284
readonly replies: Semaphore.Semaphore;
283285
/** Set while answers wait for the child to read; logged once per episode. */
@@ -639,14 +641,25 @@ export const make = Effect.fn("PluginSupervisor.make")(function* (
639641
}
640642
});
641643

644+
/** Closes a child's host work; a caller that comes while it is closing waits for the end. */
645+
const endHostWork = (child: Child) =>
646+
Effect.suspend(() => {
647+
if (child.hostWorkEnded) return Deferred.await(child.hostWorkEnded);
648+
const ended = Deferred.makeUnsafe<void>();
649+
child.hostWorkEnded = ended;
650+
return Scope.close(child.hostWork, Exit.void).pipe(
651+
Effect.ensuring(Deferred.succeed(ended, undefined)),
652+
);
653+
});
654+
642655
const handleExit = Effect.fnUntraced(function* (
643656
entry: Entry,
644657
child: Child,
645658
code: number | null,
646659
signal: string | null,
647660
) {
648661
// Host work of a dead process ends before anything else can run for this plugin.
649-
yield* Scope.close(child.hostWork, Exit.void);
662+
yield* endHostWork(child);
650663
// Anything still unread from the dead process is discarded.
651664
child.channel.destroy();
652665
const reason = child.killReason ?? describeExit(child, code, signal, options.heapLimitMb);
@@ -701,6 +714,7 @@ export const make = Effect.fn("PluginSupervisor.make")(function* (
701714
outOfMemory: false,
702715
hostCalls: 0,
703716
hostWork: Scope.forkUnsafe(fibers, "parallel"),
717+
hostWorkEnded: undefined,
704718
replies: Semaphore.makeUnsafe(1),
705719
backedUp: false,
706720
};
@@ -950,10 +964,12 @@ export const make = Effect.fn("PluginSupervisor.make")(function* (
950964
// Only reachable between claiming a start and spawning; the start then fails fast.
951965
if (!entry.child && entry.starting) yield* Deferred.await(entry.starting).pipe(Effect.ignore);
952966
const child = entry.child;
953-
if (!child || Deferred.isDoneUnsafe(child.exited)) return;
967+
if (!child) return;
968+
// A process that already exited may still be ending its host work.
969+
if (Deferred.isDoneUnsafe(child.exited)) return yield* endHostWork(child);
954970
revokeChild(entry, child);
955971
// The revoked generation's host work ends before its plugin is asked to deactivate.
956-
yield* Scope.close(child.hostWork, Exit.void);
972+
yield* endHostWork(child);
957973
write(child, { _tag: "Deactivate" });
958974
const exited = yield* Deferred.await(child.exited).pipe(Effect.timeoutOption(stopGrace));
959975
if (Option.isNone(exited)) {

0 commit comments

Comments
 (0)