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();
      }
    },
  };
}