File size: 5,409 Bytes
f778c12 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 | import { err, ok } from "@openclaw/normalization-core/result";
import {
prepareHostedGatewayStop,
type HostedGatewayStop,
type GatewayProcessOwner,
} from "../../daemon/hosted-stop.js";
import type { GatewayHostLifecycle } from "../../gateway/server-public.js";
import { formatErrorMessage } from "../../infra/errors.js";
import { disarmGatewaySuspendHandoff } from "../../infra/gateway-suspend-coordinator.js";
import { scheduleSafeGatewayRestart } from "../../infra/restart-coordinator.js";
/** The run loop retains this owner; kernels receive only its request capability. */
export function createGatewayHostLifecycle(params: {
isCurrent: () => boolean;
isServing: () => boolean;
acceptStop: () => void;
processOwner: GatewayProcessOwner;
}) {
const abort = new AbortController();
const processOwner = { ...params.processOwner };
const stopSignalReason = new Error("Gateway host stop signalled");
let state: "serving" | "preparing" | "accepted" | "finishing" | "retired" = "serving";
let stop: HostedGatewayStop | undefined;
let preparationFinished: Promise<void> | undefined;
let execution: ReturnType<HostedGatewayStop["execute"]> | undefined;
let retirement: Promise<void> | undefined;
const externalRestart = {
isCurrent: () => state === "serving" && params.isCurrent() && params.isServing(),
};
const assertCurrent = () => {
if (state === "retired" || !params.isCurrent()) {
throw new Error(
"Gateway host lifecycle is unavailable for this iteration. Reconnect and retry.",
);
}
};
const retire = () => {
if (retirement) {
return retirement;
}
state = "retired";
disarmGatewaySuspendHandoff(externalRestart);
abort.abort();
// Fence now; join the child, preparation, and execution before replacement.
// finishStop owns execution errors; retirement only waits for its unwind.
retirement = Promise.all([
stop?.dispose(),
preparationFinished,
execution?.catch(() => {}),
]).then(() => {});
stop = undefined;
return retirement;
};
const capability: GatewayHostLifecycle = {
...(processOwner.ownsProcessLifecycle ? { externalRestart } : {}),
async request(action, assertCaller) {
const assertRequest = () => {
assertCurrent();
if (!params.isServing()) {
throw new Error("Gateway host is not serving this iteration. Reconnect and retry.");
}
if (state === "accepted" || state === "finishing") {
throw new Error("Gateway stop is already scheduled.");
}
assertCaller();
};
let prepared: HostedGatewayStop | undefined;
let finishPreparation: (() => void) | undefined;
try {
assertRequest();
if (action === "start") {
return ok({ outcome: "already-running" });
}
if (!processOwner.ownsProcessLifecycle) {
throw new Error(
"This Gateway host does not own the process lifecycle. Use its owning host to stop or restart it.",
);
}
if (action === "restart") {
scheduleSafeGatewayRestart({ reason: "gateway.restart.safe", delayMs: 0 });
return ok({ outcome: "scheduled" });
}
if (state !== "serving") {
throw new Error(
"Gateway stop preparation is already in progress. Retry when it finishes.",
);
}
state = "preparing";
preparationFinished = new Promise((resolve) => {
finishPreparation = resolve;
});
prepared = await prepareHostedGatewayStop(processOwner, assertRequest, abort.signal);
assertRequest();
// Transfer exactly this intent before closing admission. Caller/kernel
// authority ends at close; only this private continuation crosses teardown.
stop = prepared;
state = "accepted";
params.acceptStop();
return ok({ outcome: "scheduled" });
} catch (error) {
await prepared?.dispose();
if (finishPreparation && state === "preparing") {
state = "serving";
}
return err(formatErrorMessage(error));
} finally {
finishPreparation?.();
}
},
};
return {
capability,
retire,
notifyStopSignal() {
if (state === "finishing" && params.isCurrent()) {
// Native stop can wait for our extinction. Cancel the client, not the
// already-drained stop, and join its close before allowing process exit.
abort.abort(stopSignalReason);
}
},
async finishStop() {
if (state !== "accepted" || !params.isCurrent() || !stop) {
return { outcome: "retired" as const };
}
state = "finishing";
const ownsStop = () => state === "finishing" && params.isCurrent();
try {
execution = stop.execute(assertCurrent);
const result = await execution;
if (!ownsStop()) {
return { outcome: "retired" as const };
}
return abort.signal.reason === stopSignalReason ? { outcome: "exit" as const } : result;
} catch (error) {
if (!ownsStop()) {
return { outcome: "retired" as const };
}
if (abort.signal.reason === stopSignalReason) {
return { outcome: "exit" as const };
}
throw error;
} finally {
await retire();
}
},
};
}
|