Skip to content
File

Blob: src/worker/dispatch/shared/run-lease.ts

typescript185 lines
1import type { ProjectRunStatus } from "@/worker/contracts";
2 
3import {
4 CANCEL_GRACE_MS,
5 HEARTBEAT_INTERVAL_MS,
6 logger,
7 now,
8 sleep,
9 type RunExecutionContext,
10} from "@/worker/dispatch/shared/run-execution-context";
11 
12type RunLeaseContext = Pick<RunExecutionContext, "control" | "logs" | "projectControl" | "runtime" | "scope" | "state">;
13 
14export interface RunLeaseControl {
15 stop(): Promise<void>;
16 throwIfOwnershipLost(): void;
17 isCancellationRequested(): boolean;
18 refreshControl(): Promise<void>;
19 applyCancellationIfNeeded(): Promise<void>;
20}
21 
22const isOwnedProjectRunStatus = (status: ProjectRunStatus): boolean =>
23 status === "active" || status === "cancel_requested";
24 
25export class RunOwnershipLostError extends Error {
26 constructor(readonly observedStatus: ProjectRunStatus | null) {
27 super(`Run ownership lost${observedStatus ? ` with ProjectDO status ${observedStatus}` : ""}.`);
28 this.name = "RunOwnershipLostError";
29 }
30}
31 
32export class RunLease implements RunLeaseControl {
33 private heartbeatPromise: Promise<void> | null = null;
34 private stopRequested = false;
35 
36 constructor(private readonly context: RunLeaseContext) {}
37 
38 start(): void {
39 if (this.heartbeatPromise) {
40 return;
41 }
42 
43 this.stopRequested = false;
44 this.heartbeatPromise = this.runHeartbeatLoop();
45 }
46 
47 async stop(): Promise<void> {
48 this.stopRequested = true;
49 if (!this.heartbeatPromise) {
50 return;
51 }
52 
53 await this.heartbeatPromise;
54 }
55 
56 throwIfOwnershipLost(): void {
57 if (this.context.state.ownershipLost) {
58 throw new RunOwnershipLostError(this.context.state.ownershipLossStatus);
59 }
60 }
61 
62 isCancellationRequested(): boolean {
63 return this.context.state.cancelRequestedAt !== null;
64 }
65 
66 async refreshControl(): Promise<void> {
67 const control = await this.context.projectControl.recordHeartbeat();
68 
69 if (control === null || !isOwnedProjectRunStatus(control.status)) {
70 this.context.control.markOwnershipLost(control?.status ?? null);
71 this.stopRequested = true;
72 return;
73 }
74 
75 if (control.status === "cancel_requested") {
76 this.context.state.cancelRequestedAt = control.cancelRequestedAt ?? this.context.state.cancelRequestedAt ?? now();
77 }
78 }
79 
80 async applyCancellationIfNeeded(): Promise<void> {
81 if (this.context.state.ownershipLost) {
82 await this.stopForLostOwnershipIfNeeded();
83 return;
84 }
85 
86 if (this.context.state.cancelRequestedAt === null) {
87 return;
88 }
89 
90 this.context.state.phase = "canceling";
91 
92 try {
93 await this.context.control.ensureRunCancelRequested();
94 } catch (error) {
95 logger.warn("run_cancel_requested_state_failed", {
96 ...this.context.scope.logContext,
97 error: error instanceof Error ? error.message : String(error),
98 });
99 }
100 
101 if (this.context.state.session === null) {
102 return;
103 }
104 
105 const cancellableProcess = this.context.runtime.getLiveCurrentProcess();
106 
107 if (!this.context.state.softCancelIssued) {
108 this.context.state.softCancelIssued = true;
109 try {
110 await this.context.control.ensureRunCanceling();
111 } catch (error) {
112 logger.warn("run_canceling_state_failed", {
113 ...this.context.scope.logContext,
114 error: error instanceof Error ? error.message : String(error),
115 });
116 }
117 
118 await this.context.logs.appendSystemLog(
119 cancellableProcess
120 ? "Cancellation requested. Sending SIGTERM."
121 : "Cancellation requested. Stopping active process tree.",
122 );
123 await this.context.runtime.softCancelProcessTree(this.context.state.session, cancellableProcess);
124 return;
125 }
126 
127 if (
128 !this.context.state.hardCancelIssued &&
129 Date.now() - this.context.state.cancelRequestedAt >= CANCEL_GRACE_MS &&
130 (await this.context.runtime.isProcessTreeAlive(this.context.state.session))
131 ) {
132 this.context.state.hardCancelIssued = true;
133 logger.warn("run_cancel_hard_kill", { ...this.context.scope.logContext });
134 await this.context.logs.appendSystemLog(
135 cancellableProcess
136 ? "Cancellation grace window expired. Sending SIGKILL."
137 : "Cancellation grace window expired. Force-stopping active process tree.",
138 );
139 await this.context.runtime.hardCancelProcessTree(this.context.state.session, cancellableProcess);
140 }
141 }
142 
143 private async runHeartbeatLoop(): Promise<void> {
144 while (!this.stopRequested) {
145 try {
146 await this.refreshControl();
147 await this.applyCancellationIfNeeded();
148 } catch (error) {
149 logger.warn("run_heartbeat_failed", {
150 ...this.context.scope.logContext,
151 error: error instanceof Error ? error.message : String(error),
152 });
153 }
154 
155 if (this.stopRequested) {
156 break;
157 }
158 
159 await sleep(HEARTBEAT_INTERVAL_MS);
160 }
161 }
162 
163 private async stopForLostOwnershipIfNeeded(): Promise<void> {
164 if (this.context.state.session === null) {
165 return;
166 }
167 
168 const stoppableProcess = this.context.runtime.getLiveCurrentProcess();
169 
170 if (!this.context.state.softCancelIssued) {
171 this.context.state.softCancelIssued = true;
172 await this.context.runtime.softCancelProcessTree(this.context.state.session, stoppableProcess);
173 return;
174 }
175 
176 if (
177 !this.context.state.hardCancelIssued &&
178 (await this.context.runtime.isProcessTreeAlive(this.context.state.session))
179 ) {
180 this.context.state.hardCancelIssued = true;
181 await this.context.runtime.hardCancelProcessTree(this.context.state.session, stoppableProcess);
182 }
183 }
184}