Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
119 changes: 86 additions & 33 deletions utils/background-task-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -182,32 +182,65 @@ export class BackgroundTaskQueue<M extends TaskModule> {
const msg: CallRequest = { id, type: "call", input };
this.worker.postMessage(msg);

this.worker.stdout.pipe(split2()).on("data", log.onstdout);
this.worker.stderr.pipe(split2()).on("data", log.onstderr);
const stdoutSplit = this.worker.stdout.pipe(split2());
const stderrSplit = this.worker.stderr.pipe(split2());
stdoutSplit.on("data", log.onstdout);
stderrSplit.on("data", log.onstderr);

const cleanup = () => {
stdoutSplit.off("data", log.onstdout);
stderrSplit.off("data", log.onstderr);
this.worker.stdout.unpipe(stdoutSplit);
this.worker.stderr.unpipe(stderrSplit);
log.markTime("finish");
};

try {
return await new Promise<RetType>((resolve, reject) => {
const listener = (response: Response) => {
if (response.id === id) {
this.lock.release();
this.worker.removeListener("message", listener);

const { type, result } = response;

log.setResult(type, result);
this.worker.stdout.off("data", log.onstdout);
this.worker.stderr.off("data", log.onstderr);
log.markTime("finish");

if (type === "success") {
resolve(result as RetType);
} else {
reject(deserializeError(result));
}
const { promise, resolve, reject } = Promise.withResolvers<RetType>();

const onWorkerError = (err: Error) => {
this.worker.removeListener("message", listener);
this.worker.removeListener("exit", onWorkerExit);
this.lock.release();
cleanup();
reject(err);
};

const onWorkerExit = (code: number | null) => {
this.worker.removeListener("message", listener);
this.worker.removeListener("error", onWorkerError);
this.lock.release();
cleanup();
reject(
new Error(
`Worker exited unexpectedly with code ${code} (job: ${id})`,
),
);
};

const listener = (response: Response) => {
if (response.id === id) {
this.worker.removeListener("message", listener);
Comment thread
marcoscaceres marked this conversation as resolved.
this.worker.removeListener("error", onWorkerError);
this.worker.removeListener("exit", onWorkerExit);
this.lock.release();

const { type, result } = response;
log.setResult(type, result);
cleanup();

if (type === "success") {
resolve(result as RetType);
} else {
reject(deserializeError(result));
}
};
this.worker.addListener("message", listener);
});
}
};

this.worker.once("error", onWorkerError);
this.worker.once("exit", onWorkerExit);
this.worker.addListener("message", listener);
return await promise;
} finally {
await log.write();
}
Expand All @@ -224,19 +257,39 @@ export class BackgroundTaskQueue<M extends TaskModule> {
const msg: OperationRequest = { id, type: "init", modulePath };
this.worker.postMessage(msg);

await new Promise<void>((resolve, reject) => {
this.worker.once("message", (response: Response) => {
if (response.id !== id) {
reject(new Error(`Failed to register worker module: ${modulePath}`));
const { promise, resolve, reject } = Promise.withResolvers<void>();

const onError = (err: Error) => {
this.worker.removeListener("exit", onExit);
reject(err);
};

const onExit = (code: number | null) => {
this.worker.removeListener("error", onError);
reject(
new Error(
`Worker exited with code ${code} during module initialization: ${modulePath}`,
),
);
};

this.worker.once("error", onError);
this.worker.once("exit", onExit);

this.worker.once("message", (response: Response) => {
this.worker.removeListener("error", onError);
this.worker.removeListener("exit", onExit);
if (response.id !== id) {
reject(new Error(`Failed to register worker module: ${modulePath}`));
} else {
if (response.type === "success") {
resolve();
} else {
if (response.type === "success") {
resolve();
} else {
reject(deserializeError(response.result));
}
reject(deserializeError(response.result));
}
});
}
});
await promise;
Comment thread
marcoscaceres marked this conversation as resolved.
}

private generateId() {
Expand Down
64 changes: 37 additions & 27 deletions utils/sh.ts
Original file line number Diff line number Diff line change
Expand Up @@ -31,34 +31,44 @@ export default async function sh(
}

try {
return await new Promise<string>((resolve, reject) => {
let stdout: string[] = [];
let stderr: string[] = [];
const child = exec(command, {
...execOptions,
env: { ...process.env, ...execOptions.env },
encoding: "utf-8",
});
child.stdout!.pipe(split()).on("data", (line: string) => {
if (shouldStream) log.out(line);
stdout.push(line);
});
child.stderr!.pipe(split()).on("data", (line: string) => {
if (shouldStream) log.err(line);
stderr.push(line);
});
child.on("exit", code => {
if (output === "buffer") {
if (stdout.length) log.out(stdout.join("\n"));
if (stderr.length) log.err(stderr.join("\n"));
}
if (code === 0) {
resolve(stdout.join("\n"));
} else {
reject({ command, stdout, stderr, code });
}
});
const { promise, resolve, reject } = Promise.withResolvers<string>();
let stdout: string[] = [];
let stderr: string[] = [];
const child = exec(command, {
...execOptions,
env: { ...process.env, ...execOptions.env },
encoding: "utf-8",
});
child.stdout!.pipe(split()).on("data", (line: string) => {
if (shouldStream) log.out(line);
stdout.push(line);
});
child.stderr!.pipe(split()).on("data", (line: string) => {
if (shouldStream) log.err(line);
stderr.push(line);
});
child.on("error", err => {
reject(Object.assign(err, { command, stdout, stderr, code: null }));
});
child.on("close", code => {
if (output === "buffer") {
if (stdout.length) log.out(stdout.join("\n"));
if (stderr.length) log.err(stderr.join("\n"));
}
if (code === 0) {
resolve(stdout.join("\n"));
} else {
reject(
Object.assign(new Error(`Command failed: ${command}`), {
command,
stdout,
stderr,
code,
}),
);
}
});
return await promise;
} finally {
if (output !== "silent") {
console.groupEnd();
Expand Down
Loading