fix worker lifecycle edge cases
This commit is contained in:
parent
dc0627c285
commit
0531fc1dbf
2 changed files with 24 additions and 16 deletions
|
|
@ -577,6 +577,12 @@ async function handleGetData(
|
||||||
if (getDataMs >= 1500) {
|
if (getDataMs >= 1500) {
|
||||||
logger.warn(`[teams:getData] slow team=${tn} ms=${getDataMs}`);
|
logger.warn(`[teams:getData] slow team=${tn} ms=${getDataMs}`);
|
||||||
}
|
}
|
||||||
|
const teamDataService = getTeamDataService();
|
||||||
|
if (data.processes.some((process) => !process.stoppedAt)) {
|
||||||
|
teamDataService.trackProcessHealthForTeam?.(tn);
|
||||||
|
} else {
|
||||||
|
teamDataService.untrackProcessHealthForTeam?.(tn);
|
||||||
|
}
|
||||||
const provisioning = getTeamProvisioningService();
|
const provisioning = getTeamProvisioningService();
|
||||||
const isAlive = provisioning.isTeamAlive(tn);
|
const isAlive = provisioning.isTeamAlive(tn);
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -60,6 +60,18 @@ export class TeamDataWorkerClient {
|
||||||
private warnedUnavailable = false;
|
private warnedUnavailable = false;
|
||||||
private pending = new Map<string, PendingEntry>();
|
private pending = new Map<string, PendingEntry>();
|
||||||
|
|
||||||
|
private failWorker(worker: Worker, error: Error): void {
|
||||||
|
if (this.worker !== worker) return;
|
||||||
|
|
||||||
|
this.worker = null;
|
||||||
|
const pendingEntries = Array.from(this.pending.values());
|
||||||
|
this.pending.clear();
|
||||||
|
|
||||||
|
for (const entry of pendingEntries) {
|
||||||
|
entry.reject(error);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
isAvailable(): boolean {
|
isAvailable(): boolean {
|
||||||
if (!this.workerPath && !this.warnedUnavailable) {
|
if (!this.workerPath && !this.warnedUnavailable) {
|
||||||
this.warnedUnavailable = true;
|
this.warnedUnavailable = true;
|
||||||
|
|
@ -90,23 +102,13 @@ export class TeamDataWorkerClient {
|
||||||
// Without this guard, a stale worker's exit event can reject
|
// Without this guard, a stale worker's exit event can reject
|
||||||
// pending requests that belong to a newer replacement worker.
|
// pending requests that belong to a newer replacement worker.
|
||||||
w.on('error', (err) => {
|
w.on('error', (err) => {
|
||||||
if (this.worker !== w) return;
|
|
||||||
logger.error('Worker error', err);
|
logger.error('Worker error', err);
|
||||||
for (const [, entry] of this.pending) {
|
this.failWorker(w, err instanceof Error ? err : new Error(String(err)));
|
||||||
entry.reject(err instanceof Error ? err : new Error(String(err)));
|
|
||||||
}
|
|
||||||
this.pending.clear();
|
|
||||||
this.worker = null;
|
|
||||||
});
|
});
|
||||||
|
|
||||||
w.on('exit', (code) => {
|
w.on('exit', (code) => {
|
||||||
if (this.worker !== w) return;
|
|
||||||
if (code !== 0) logger.warn(`Worker exited with code ${code}`);
|
if (code !== 0) logger.warn(`Worker exited with code ${code}`);
|
||||||
for (const [, entry] of this.pending) {
|
this.failWorker(w, new Error(`Worker exited with code ${code}`));
|
||||||
entry.reject(new Error(`Worker exited with code ${code}`));
|
|
||||||
}
|
|
||||||
this.pending.clear();
|
|
||||||
this.worker = null;
|
|
||||||
});
|
});
|
||||||
|
|
||||||
return w;
|
return w;
|
||||||
|
|
@ -121,10 +123,10 @@ export class TeamDataWorkerClient {
|
||||||
|
|
||||||
return new Promise((resolve, reject) => {
|
return new Promise((resolve, reject) => {
|
||||||
const timeout = setTimeout(() => {
|
const timeout = setTimeout(() => {
|
||||||
this.pending.delete(id);
|
const timeoutError = new Error(`Worker call timeout after ${WORKER_CALL_TIMEOUT_MS}ms`);
|
||||||
this.worker?.terminate().catch(() => undefined);
|
this.failWorker(worker, timeoutError);
|
||||||
this.worker = null;
|
worker.terminate().catch(() => undefined);
|
||||||
reject(new Error(`Worker call timeout after ${WORKER_CALL_TIMEOUT_MS}ms`));
|
reject(timeoutError);
|
||||||
}, WORKER_CALL_TIMEOUT_MS);
|
}, WORKER_CALL_TIMEOUT_MS);
|
||||||
|
|
||||||
this.pending.set(id, {
|
this.pending.set(id, {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue