Skip to content

Commit 9a1b077

Browse files
marcoscaceresCopilotCopilot
committed
refactor: use Promise.withResolvers() (Node 22+) (#503)
* refactor: use Promise.withResolvers() (Node 22+) Replace new Promise((resolve, reject) => { ... }) with the cleaner Promise.withResolvers() pattern in background-task-queue.ts (2 sites) and sh.ts (1 site). * Apply suggestion from @Copilot Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> * fix: address all PR review feedback on sh.ts and background-task-queue.ts - sh.ts: Remove broken duplicate Promise.withResolvers + new Promise code; use Promise.withResolvers() properly; switch exit->close for complete stdio flush; add error event handler for spawn failures; make rejection shapes consistent (both reject with Error instances with command/stdout/stderr/code fields) - background-task-queue: Store split2 stream refs so off() targets the right stream and unpipe on cleanup; add worker error/exit handling in run() to prevent deadlocks on unexpected worker failures; add worker error/exit handling in activate() to prevent indefinite hangs during initialization Agent-Logs-Url: https://github.com/speced/respec-web-services/sessions/4640ebb0-d6ac-4269-b7bf-e563038f39e1 * fix: add job ID to worker unexpected-exit error message for easier debugging Agent-Logs-Url: https://github.com/speced/respec-web-services/sessions/4640ebb0-d6ac-4269-b7bf-e563038f39e1 * fix: cross-remove error/exit listeners in run() to prevent double lock release Agent-Logs-Url: https://github.com/speced/respec-web-services/sessions/5a7c6559-a65f-44e5-865b-b8a2f0b8c955 --------- Co-authored-by: Copilot <175728472+Copilot@users.noreply.github.com> Co-authored-by: copilot-swe-agent[bot] <198982749+Copilot@users.noreply.github.com>
1 parent de9d6e6 commit 9a1b077

2 files changed

Lines changed: 123 additions & 60 deletions

File tree

utils/background-task-queue.ts

Lines changed: 86 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -182,32 +182,65 @@ export class BackgroundTaskQueue<M extends TaskModule> {
182182
const msg: CallRequest = { id, type: "call", input };
183183
this.worker.postMessage(msg);
184184

185-
this.worker.stdout.pipe(split2()).on("data", log.onstdout);
186-
this.worker.stderr.pipe(split2()).on("data", log.onstderr);
185+
const stdoutSplit = this.worker.stdout.pipe(split2());
186+
const stderrSplit = this.worker.stderr.pipe(split2());
187+
stdoutSplit.on("data", log.onstdout);
188+
stderrSplit.on("data", log.onstderr);
189+
190+
const cleanup = () => {
191+
stdoutSplit.off("data", log.onstdout);
192+
stderrSplit.off("data", log.onstderr);
193+
this.worker.stdout.unpipe(stdoutSplit);
194+
this.worker.stderr.unpipe(stderrSplit);
195+
log.markTime("finish");
196+
};
187197

188198
try {
189-
return await new Promise<RetType>((resolve, reject) => {
190-
const listener = (response: Response) => {
191-
if (response.id === id) {
192-
this.lock.release();
193-
this.worker.removeListener("message", listener);
194-
195-
const { type, result } = response;
196-
197-
log.setResult(type, result);
198-
this.worker.stdout.off("data", log.onstdout);
199-
this.worker.stderr.off("data", log.onstderr);
200-
log.markTime("finish");
201-
202-
if (type === "success") {
203-
resolve(result as RetType);
204-
} else {
205-
reject(deserializeError(result));
206-
}
199+
const { promise, resolve, reject } = Promise.withResolvers<RetType>();
200+
201+
const onWorkerError = (err: Error) => {
202+
this.worker.removeListener("message", listener);
203+
this.worker.removeListener("exit", onWorkerExit);
204+
this.lock.release();
205+
cleanup();
206+
reject(err);
207+
};
208+
209+
const onWorkerExit = (code: number | null) => {
210+
this.worker.removeListener("message", listener);
211+
this.worker.removeListener("error", onWorkerError);
212+
this.lock.release();
213+
cleanup();
214+
reject(
215+
new Error(
216+
`Worker exited unexpectedly with code ${code} (job: ${id})`,
217+
),
218+
);
219+
};
220+
221+
const listener = (response: Response) => {
222+
if (response.id === id) {
223+
this.worker.removeListener("message", listener);
224+
this.worker.removeListener("error", onWorkerError);
225+
this.worker.removeListener("exit", onWorkerExit);
226+
this.lock.release();
227+
228+
const { type, result } = response;
229+
log.setResult(type, result);
230+
cleanup();
231+
232+
if (type === "success") {
233+
resolve(result as RetType);
234+
} else {
235+
reject(deserializeError(result));
207236
}
208-
};
209-
this.worker.addListener("message", listener);
210-
});
237+
}
238+
};
239+
240+
this.worker.once("error", onWorkerError);
241+
this.worker.once("exit", onWorkerExit);
242+
this.worker.addListener("message", listener);
243+
return await promise;
211244
} finally {
212245
await log.write();
213246
}
@@ -224,19 +257,39 @@ export class BackgroundTaskQueue<M extends TaskModule> {
224257
const msg: OperationRequest = { id, type: "init", modulePath };
225258
this.worker.postMessage(msg);
226259

227-
await new Promise<void>((resolve, reject) => {
228-
this.worker.once("message", (response: Response) => {
229-
if (response.id !== id) {
230-
reject(new Error(`Failed to register worker module: ${modulePath}`));
260+
const { promise, resolve, reject } = Promise.withResolvers<void>();
261+
262+
const onError = (err: Error) => {
263+
this.worker.removeListener("exit", onExit);
264+
reject(err);
265+
};
266+
267+
const onExit = (code: number | null) => {
268+
this.worker.removeListener("error", onError);
269+
reject(
270+
new Error(
271+
`Worker exited with code ${code} during module initialization: ${modulePath}`,
272+
),
273+
);
274+
};
275+
276+
this.worker.once("error", onError);
277+
this.worker.once("exit", onExit);
278+
279+
this.worker.once("message", (response: Response) => {
280+
this.worker.removeListener("error", onError);
281+
this.worker.removeListener("exit", onExit);
282+
if (response.id !== id) {
283+
reject(new Error(`Failed to register worker module: ${modulePath}`));
284+
} else {
285+
if (response.type === "success") {
286+
resolve();
231287
} else {
232-
if (response.type === "success") {
233-
resolve();
234-
} else {
235-
reject(deserializeError(response.result));
236-
}
288+
reject(deserializeError(response.result));
237289
}
238-
});
290+
}
239291
});
292+
await promise;
240293
}
241294

242295
private generateId() {

utils/sh.ts

Lines changed: 37 additions & 27 deletions
Original file line numberDiff line numberDiff line change
@@ -31,34 +31,44 @@ export default async function sh(
3131
}
3232

3333
try {
34-
return await new Promise<string>((resolve, reject) => {
35-
let stdout: string[] = [];
36-
let stderr: string[] = [];
37-
const child = exec(command, {
38-
...execOptions,
39-
env: { ...process.env, ...execOptions.env },
40-
encoding: "utf-8",
41-
});
42-
child.stdout!.pipe(split()).on("data", (line: string) => {
43-
if (shouldStream) log.out(line);
44-
stdout.push(line);
45-
});
46-
child.stderr!.pipe(split()).on("data", (line: string) => {
47-
if (shouldStream) log.err(line);
48-
stderr.push(line);
49-
});
50-
child.on("exit", code => {
51-
if (output === "buffer") {
52-
if (stdout.length) log.out(stdout.join("\n"));
53-
if (stderr.length) log.err(stderr.join("\n"));
54-
}
55-
if (code === 0) {
56-
resolve(stdout.join("\n"));
57-
} else {
58-
reject({ command, stdout, stderr, code });
59-
}
60-
});
34+
const { promise, resolve, reject } = Promise.withResolvers<string>();
35+
let stdout: string[] = [];
36+
let stderr: string[] = [];
37+
const child = exec(command, {
38+
...execOptions,
39+
env: { ...process.env, ...execOptions.env },
40+
encoding: "utf-8",
6141
});
42+
child.stdout!.pipe(split()).on("data", (line: string) => {
43+
if (shouldStream) log.out(line);
44+
stdout.push(line);
45+
});
46+
child.stderr!.pipe(split()).on("data", (line: string) => {
47+
if (shouldStream) log.err(line);
48+
stderr.push(line);
49+
});
50+
child.on("error", err => {
51+
reject(Object.assign(err, { command, stdout, stderr, code: null }));
52+
});
53+
child.on("close", code => {
54+
if (output === "buffer") {
55+
if (stdout.length) log.out(stdout.join("\n"));
56+
if (stderr.length) log.err(stderr.join("\n"));
57+
}
58+
if (code === 0) {
59+
resolve(stdout.join("\n"));
60+
} else {
61+
reject(
62+
Object.assign(new Error(`Command failed: ${command}`), {
63+
command,
64+
stdout,
65+
stderr,
66+
code,
67+
}),
68+
);
69+
}
70+
});
71+
return await promise;
6272
} finally {
6373
if (output !== "silent") {
6474
console.groupEnd();

0 commit comments

Comments
 (0)