fix: serialize provenance capture for shared extract roots
Protect each shared extraction target with a path-scoped asynchronous handoff so concurrent packages cannot claim each other's files between provenance snapshots. Independent extract directories remain parallel, and the handoff is released on both success and failure. Add a concurrent two-package regression proving delayed entry and one attributed output per package.
This commit is contained in:
@@ -1910,6 +1910,8 @@ export class DownloadManager extends EventEmitter {
|
||||
|
||||
private cleanupQueue: Promise<void> = Promise.resolve();
|
||||
|
||||
private packageOutputProvenanceTails = new Map<string, Promise<void>>();
|
||||
|
||||
private packagePostProcessQueue: Promise<void> = Promise.resolve();
|
||||
|
||||
private packagePostProcessActive = 0;
|
||||
@@ -4410,6 +4412,30 @@ export class DownloadManager extends EventEmitter {
|
||||
pkg.outputCount = Math.max(pkg.outputCount || 0, provenance.size);
|
||||
}
|
||||
|
||||
private async runWithPackageOutputProvenance<T>(pkg: PackageEntry, operation: () => Promise<T>): Promise<T> {
|
||||
const key = pathKey(pkg.extractDir);
|
||||
const previous = this.packageOutputProvenanceTails.get(key) || Promise.resolve();
|
||||
let release!: () => void;
|
||||
const current = new Promise<void>((resolve) => {
|
||||
release = resolve;
|
||||
});
|
||||
this.packageOutputProvenanceTails.set(key, current);
|
||||
await previous;
|
||||
const before = await this.snapshotPackageOutputFiles(pkg.extractDir);
|
||||
try {
|
||||
return await operation();
|
||||
} finally {
|
||||
try {
|
||||
await this.recordPackageOutputFiles(pkg, before);
|
||||
} finally {
|
||||
release();
|
||||
if (this.packageOutputProvenanceTails.get(key) === current) {
|
||||
this.packageOutputProvenanceTails.delete(key);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private async removeEmptyDirectoryTree(rootDir: string): Promise<number> {
|
||||
if (!rootDir) {
|
||||
return 0;
|
||||
@@ -13265,8 +13291,7 @@ export class DownloadManager extends EventEmitter {
|
||||
return 0;
|
||||
}
|
||||
|
||||
const packageOutputBefore = await this.snapshotPackageOutputFiles(pkg.extractDir);
|
||||
const result = await extractPackageArchives({
|
||||
const result = await this.runWithPackageOutputProvenance(pkg, () => extractPackageArchives({
|
||||
packageDir: pkg.outputDir,
|
||||
targetDir: pkg.extractDir,
|
||||
cleanupMode: this.settings.cleanupMode,
|
||||
@@ -13424,8 +13449,7 @@ export class DownloadManager extends EventEmitter {
|
||||
this.emitState();
|
||||
}
|
||||
}
|
||||
});
|
||||
await this.recordPackageOutputFiles(pkg, packageOutputBefore);
|
||||
}));
|
||||
|
||||
logger.info(`Hybrid-Extract Ende: pkg=${pkg.name}, extracted=${result.extracted}, failed=${result.failed}`);
|
||||
this.logPackageForPackage(pkg, "INFO", "Hybrid-Extract abgeschlossen", {
|
||||
@@ -13842,8 +13866,7 @@ export class DownloadManager extends EventEmitter {
|
||||
entry.updatedAt = pendingAt;
|
||||
}
|
||||
this.emitState();
|
||||
const packageOutputBefore = await this.snapshotPackageOutputFiles(pkg.extractDir);
|
||||
const result = await extractPackageArchives({
|
||||
const result = await this.runWithPackageOutputProvenance(pkg, () => extractPackageArchives({
|
||||
packageDir: pkg.outputDir,
|
||||
targetDir: pkg.extractDir,
|
||||
cleanupMode: this.settings.cleanupMode,
|
||||
@@ -13984,8 +14007,7 @@ export class DownloadManager extends EventEmitter {
|
||||
}
|
||||
emitExtractStatus(overallLabel);
|
||||
}
|
||||
});
|
||||
await this.recordPackageOutputFiles(pkg, packageOutputBefore);
|
||||
}));
|
||||
logger.info(`Post-Processing Entpacken Ende: pkg=${pkg.name}, extracted=${result.extracted}, failed=${result.failed}, lastError=${result.lastError || ""}`);
|
||||
this.logPackageForPackage(pkg, "INFO", "Post-Processing Entpacken Ende", {
|
||||
extracted: result.extracted,
|
||||
@@ -14192,8 +14214,7 @@ export class DownloadManager extends EventEmitter {
|
||||
});
|
||||
const nestedFailureCategories = new Map<string, string>();
|
||||
const nestedItems = pkg.itemIds.map((itemId) => this.session.items[itemId]).filter(Boolean) as DownloadItem[];
|
||||
const packageOutputBefore = await this.snapshotPackageOutputFiles(pkg.extractDir);
|
||||
const nestedResult = await extractPackageArchives({
|
||||
const nestedResult = await this.runWithPackageOutputProvenance(pkg, () => extractPackageArchives({
|
||||
packageDir: pkg.extractDir,
|
||||
targetDir: pkg.extractDir,
|
||||
cleanupMode: this.settings.cleanupMode,
|
||||
@@ -14218,9 +14239,8 @@ export class DownloadManager extends EventEmitter {
|
||||
nestedFailureCategories.get(progress.archiveName.toLowerCase()) || ""
|
||||
);
|
||||
}
|
||||
});
|
||||
}));
|
||||
throwIfAborted();
|
||||
await this.recordPackageOutputFiles(pkg, packageOutputBefore);
|
||||
extractedCount += nestedResult.extracted;
|
||||
logger.info(`Deferred Nested-Extraction Ende: extracted=${nestedResult.extracted}, failed=${nestedResult.failed}`);
|
||||
this.logPackageForPackage(pkg, "INFO", "Deferred Nested-Extraction Ende", {
|
||||
|
||||
@@ -15114,6 +15114,51 @@ describe("package priority ordering", () => {
|
||||
});
|
||||
|
||||
describe("package lifecycle telemetry boundaries", () => {
|
||||
it("serializes provenance capture for packages sharing one extract directory", async () => {
|
||||
const root = fs.mkdtempSync(path.join(os.tmpdir(), "rd-output-provenance-lock-"));
|
||||
tempDirs.push(root);
|
||||
const extractDir = path.join(root, "extract");
|
||||
fs.mkdirSync(extractDir, { recursive: true });
|
||||
const manager = new DownloadManager(defaultSettings(), emptySession(), createStoragePaths(path.join(root, "state")));
|
||||
const createPackage = (id: string): PackageEntry => ({
|
||||
id,
|
||||
name: id,
|
||||
outputDir: path.join(root, "downloads", id),
|
||||
extractDir,
|
||||
status: "completed",
|
||||
itemIds: [],
|
||||
cancelled: false,
|
||||
enabled: true,
|
||||
createdAt: 1_000,
|
||||
updatedAt: 1_000
|
||||
});
|
||||
const packageA = createPackage("package-a");
|
||||
const packageB = createPackage("package-b");
|
||||
let releaseA!: () => void;
|
||||
const gateA = new Promise<void>((resolve) => {
|
||||
releaseA = resolve;
|
||||
});
|
||||
let enteredB = false;
|
||||
const state = manager as any;
|
||||
|
||||
const first = state.runWithPackageOutputProvenance(packageA, async () => {
|
||||
fs.writeFileSync(path.join(extractDir, "package-a.mkv"), "a");
|
||||
await gateA;
|
||||
});
|
||||
await vi.waitFor(() => expect(fs.existsSync(path.join(extractDir, "package-a.mkv"))).toBe(true));
|
||||
const second = state.runWithPackageOutputProvenance(packageB, async () => {
|
||||
enteredB = true;
|
||||
fs.writeFileSync(path.join(extractDir, "package-b.mkv"), "b");
|
||||
});
|
||||
await Promise.resolve();
|
||||
|
||||
expect(enteredB).toBe(false);
|
||||
releaseA();
|
||||
await Promise.all([first, second]);
|
||||
expect(packageA.outputCount).toBe(1);
|
||||
expect(packageB.outputCount).toBe(1);
|
||||
});
|
||||
|
||||
it("uses item-path provenance for archive identity and leaves unknown part counts at zero", () => {
|
||||
const root = fs.mkdtempSync(path.join(os.tmpdir(), "rd-archive-identity-"));
|
||||
tempDirs.push(root);
|
||||
|
||||
Reference in New Issue
Block a user