fix: preserve postprocessing run ownership across stops
Bind main, deferred, and hybrid postprocessing controllers to the run context that owns their package generation. Abort only the stopped run and unowned work so a later run cannot cancel deferred completion for an earlier run. Treat stop and shutdown aborts as lifecycle cancellation instead of cleanup failure and cover the real deferred handoff through package, run, history, and cleanup results.
This commit is contained in:
@@ -1918,14 +1918,20 @@ export class DownloadManager extends EventEmitter {
|
|||||||
|
|
||||||
private packagePostProcessTasks = new Map<string, Promise<void>>();
|
private packagePostProcessTasks = new Map<string, Promise<void>>();
|
||||||
|
|
||||||
private packagePostProcessAbortControllers = new Map<string, AbortController>();
|
private packagePostProcessAbortControllers = new Map<string, AbortController>();
|
||||||
|
|
||||||
|
private packagePostProcessRunOwnerByController = new WeakMap<AbortController, string | null>();
|
||||||
|
|
||||||
private packageDeferredPostProcessAbortControllers = new Map<string, AbortController>();
|
private packageDeferredPostProcessAbortControllers = new Map<string, AbortController>();
|
||||||
|
|
||||||
|
private packageDeferredRunOwnerByController = new WeakMap<AbortController, string | null>();
|
||||||
|
|
||||||
private packageDeferredPostProcessTasks = new Map<string, Set<Promise<void>>>();
|
private packageDeferredPostProcessTasks = new Map<string, Set<Promise<void>>>();
|
||||||
|
|
||||||
private packageHybridPostProcessControllers = new Map<string, Set<AbortController>>();
|
private packageHybridPostProcessControllers = new Map<string, Set<AbortController>>();
|
||||||
|
|
||||||
|
private packageHybridRunOwnerByController = new WeakMap<AbortController, string | null>();
|
||||||
|
|
||||||
private packageHybridPostProcessTasks = new Map<string, Set<Promise<void>>>();
|
private packageHybridPostProcessTasks = new Map<string, Set<Promise<void>>>();
|
||||||
|
|
||||||
private packagePostProcessVersions = new Map<string, number>();
|
private packagePostProcessVersions = new Map<string, number>();
|
||||||
@@ -6805,7 +6811,7 @@ export class DownloadManager extends EventEmitter {
|
|||||||
this.speedBytesLastWindow = 0;
|
this.speedBytesLastWindow = 0;
|
||||||
this.speedBytesPerPackage.clear();
|
this.speedBytesPerPackage.clear();
|
||||||
this.speedEventsHead = 0;
|
this.speedEventsHead = 0;
|
||||||
this.abortPostProcessing("stop");
|
this.abortPostProcessing("stop", stoppedRunContext?.id);
|
||||||
for (const waiter of this.packagePostProcessWaiters) { waiter.resolve(); }
|
for (const waiter of this.packagePostProcessWaiters) { waiter.resolve(); }
|
||||||
this.packagePostProcessWaiters = [];
|
this.packagePostProcessWaiters = [];
|
||||||
this.packagePostProcessActive = 0;
|
this.packagePostProcessActive = 0;
|
||||||
@@ -8381,9 +8387,13 @@ export class DownloadManager extends EventEmitter {
|
|||||||
return claimed;
|
return claimed;
|
||||||
}
|
}
|
||||||
|
|
||||||
private abortPostProcessing(reason: string): void {
|
private abortPostProcessing(reason: string, runContextId?: string): void {
|
||||||
for (const [packageId, controller] of this.packagePostProcessAbortControllers.entries()) {
|
for (const [packageId, controller] of this.packagePostProcessAbortControllers.entries()) {
|
||||||
if (!controller.signal.aborted) {
|
const owner = this.packagePostProcessRunOwnerByController.get(controller);
|
||||||
|
if (runContextId !== undefined && owner !== undefined && owner !== null && owner !== runContextId) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if (!controller.signal.aborted) {
|
||||||
controller.abort(reason);
|
controller.abort(reason);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -8410,16 +8420,24 @@ export class DownloadManager extends EventEmitter {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
for (const controller of this.packageDeferredPostProcessAbortControllers.values()) {
|
for (const controller of this.packageDeferredPostProcessAbortControllers.values()) {
|
||||||
if (!controller.signal.aborted) {
|
const owner = this.packageDeferredRunOwnerByController.get(controller);
|
||||||
|
if (runContextId !== undefined && owner !== undefined && owner !== null && owner !== runContextId) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if (!controller.signal.aborted) {
|
||||||
controller.abort(reason);
|
controller.abort(reason);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
for (const hybridSet of this.packageHybridPostProcessControllers.values()) {
|
for (const hybridSet of this.packageHybridPostProcessControllers.values()) {
|
||||||
for (const controller of hybridSet) {
|
for (const controller of hybridSet) {
|
||||||
if (!controller.signal.aborted) {
|
const owner = this.packageHybridRunOwnerByController.get(controller);
|
||||||
|
if (runContextId !== undefined && owner !== undefined && owner !== null && owner !== runContextId) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
if (!controller.signal.aborted) {
|
||||||
controller.abort(reason);
|
controller.abort(reason);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -8473,6 +8491,7 @@ export class DownloadManager extends EventEmitter {
|
|||||||
|
|
||||||
const abortController = new AbortController();
|
const abortController = new AbortController();
|
||||||
this.packagePostProcessAbortControllers.set(packageId, abortController);
|
this.packagePostProcessAbortControllers.set(packageId, abortController);
|
||||||
|
this.packagePostProcessRunOwnerByController.set(abortController, this.getPackageResultRunOwner(packageId));
|
||||||
const queuedPackage = this.session.packages[packageId];
|
const queuedPackage = this.session.packages[packageId];
|
||||||
if (queuedPackage) {
|
if (queuedPackage) {
|
||||||
queuedPackage.postProcessQueuedAt = queuedPackage.postProcessQueuedAt || nowMs();
|
queuedPackage.postProcessQueuedAt = queuedPackage.postProcessQueuedAt || nowMs();
|
||||||
@@ -12367,6 +12386,16 @@ export class DownloadManager extends EventEmitter {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private getPackageResultRunOwner(packageId: string): string | null {
|
||||||
|
const generation = this.getPackageResultGeneration(packageId);
|
||||||
|
for (const context of this.runContexts.values()) {
|
||||||
|
if (context.packageGenerations.get(packageId) === generation) {
|
||||||
|
return context.id;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
private queueNotificationEvent(notification: NotificationEvent): void {
|
private queueNotificationEvent(notification: NotificationEvent): void {
|
||||||
if (!this.enqueueNotificationCallback || !String(this.settings.notifyUrl || "").trim()) {
|
if (!this.enqueueNotificationCallback || !String(this.settings.notifyUrl || "").trim()) {
|
||||||
return;
|
return;
|
||||||
@@ -13389,6 +13418,7 @@ export class DownloadManager extends EventEmitter {
|
|||||||
if (result.extracted > 0) {
|
if (result.extracted > 0) {
|
||||||
this.trackPackagePostProcessResult(packageId);
|
this.trackPackagePostProcessResult(packageId);
|
||||||
const hybridController = new AbortController();
|
const hybridController = new AbortController();
|
||||||
|
this.packageHybridRunOwnerByController.set(hybridController, this.getPackageResultRunOwner(packageId));
|
||||||
let hybridSet = this.packageHybridPostProcessControllers.get(packageId);
|
let hybridSet = this.packageHybridPostProcessControllers.get(packageId);
|
||||||
if (!hybridSet) {
|
if (!hybridSet) {
|
||||||
hybridSet = new Set<AbortController>();
|
hybridSet = new Set<AbortController>();
|
||||||
@@ -14093,8 +14123,9 @@ export class DownloadManager extends EventEmitter {
|
|||||||
if (replacedController && !replacedController.signal.aborted) {
|
if (replacedController && !replacedController.signal.aborted) {
|
||||||
replacedController.abort("deferred_replaced");
|
replacedController.abort("deferred_replaced");
|
||||||
}
|
}
|
||||||
const deferredController = new AbortController();
|
const deferredController = new AbortController();
|
||||||
this.packageDeferredPostProcessAbortControllers.set(packageId, deferredController);
|
this.packageDeferredPostProcessAbortControllers.set(packageId, deferredController);
|
||||||
|
this.packageDeferredRunOwnerByController.set(deferredController, this.getPackageResultRunOwner(packageId));
|
||||||
const deferredVersion = this.getPackagePostProcessVersion(packageId);
|
const deferredVersion = this.getPackagePostProcessVersion(packageId);
|
||||||
const shouldAbort = (): boolean => !this.isDeferredPostProcessStillCurrent(packageId, pkg, deferredVersion, deferredController.signal);
|
const shouldAbort = (): boolean => !this.isDeferredPostProcessStillCurrent(packageId, pkg, deferredVersion, deferredController.signal);
|
||||||
const throwIfAborted = (): void => this.throwIfDeferredPostProcessAborted(packageId, pkg, deferredVersion, deferredController.signal);
|
const throwIfAborted = (): void => this.throwIfDeferredPostProcessAborted(packageId, pkg, deferredVersion, deferredController.signal);
|
||||||
@@ -14252,9 +14283,13 @@ export class DownloadManager extends EventEmitter {
|
|||||||
|| reason.includes("package_removed")
|
|| reason.includes("package_removed")
|
||||||
|| reason === "reset"
|
|| reason === "reset"
|
||||||
|| reason === "cancel"
|
|| reason === "cancel"
|
||||||
|| reason === "overwrite"
|
|| reason === "overwrite"
|
||||||
|| reason === "skip"
|
|| reason === "skip"
|
||||||
|| reason === "package_toggle") {
|
|| reason === "package_toggle"
|
||||||
|
|| reason === "stop"
|
||||||
|
|| reason === "shutdown"
|
||||||
|
|| reason === "Error: stop"
|
||||||
|
|| reason === "Error: shutdown") {
|
||||||
logger.info(`Deferred Post-Extraction abgebrochen: pkg=${pkg.name}, reason=${reason}`);
|
logger.info(`Deferred Post-Extraction abgebrochen: pkg=${pkg.name}, reason=${reason}`);
|
||||||
} else {
|
} else {
|
||||||
pkg.cleanupErrorCategory = reason.slice(0, 256) || "cleanup";
|
pkg.cleanupErrorCategory = reason.slice(0, 256) || "cleanup";
|
||||||
|
|||||||
@@ -755,6 +755,7 @@ describe("authoritative run completion", () => {
|
|||||||
await Promise.allSettled(deferredTasks);
|
await Promise.allSettled(deferredTasks);
|
||||||
await flushNotifications();
|
await flushNotifications();
|
||||||
|
|
||||||
|
expect(packageA.cleanupErrorCategory || "").toBe("");
|
||||||
expect(events.filter((event) => event.type === "package_completed")).toHaveLength(1);
|
expect(events.filter((event) => event.type === "package_completed")).toHaveLength(1);
|
||||||
expect(events.filter((event) => event.type === "run_completed")).toHaveLength(1);
|
expect(events.filter((event) => event.type === "run_completed")).toHaveLength(1);
|
||||||
expect(history.map((entry) => entry.name)).toEqual([packageA.name]);
|
expect(history.map((entry) => entry.name)).toEqual([packageA.name]);
|
||||||
|
|||||||
Reference in New Issue
Block a user