diff --git a/src/main/download-manager.ts b/src/main/download-manager.ts index 67e801b..ea570d3 100644 --- a/src/main/download-manager.ts +++ b/src/main/download-manager.ts @@ -64,7 +64,7 @@ function releaseTlsSkip(): void { import { cleanupCancelledPackageArtifactsAsync, removeDownloadLinkArtifactsFromScope, removeSampleArtifactsFromScope } from "./cleanup"; import { planDownloadCompletion, reconcileFinalizedSize, validateDownloadedFileCompletion } from "./download-completion"; import { AllDebridWebUnrestrictor, BestDebridWebUnrestrictor, DebridService, MegaWebUnrestrictor, RealDebridWebUnrestrictor, checkDdownloadOnline, checkOneFichierLinks, checkRapidgatorOnline, fetchAllDebridHostInfo, filenameFromDdownloadUrlPath, getAvailableDebridLinkApiKeys, getAvailableMegaDebridAccounts, getAvailableRealDebridAccounts, getMegaDebridAccountCooldownState, getMegaDebridInFlightCountForMode, getRealDebridAccountAttemptTimeoutMs, isDdownloadLink, isOneFichierLink, isProviderDisabledForSelection, pruneExpiredDebridLinkRuntimeState, pruneExpiredMegaDebridRuntimeState, pruneExpiredRealDebridRuntimeState, releaseRealDebridAccountCooldown, type DdownloadCheckResult, type OneFichierCheckResult } from "./debrid"; -import { cleanupArchives, clearExtractResumeState, collectArchiveCleanupTargets, detectArchiveSignature, extractPackageArchives, findArchiveCandidates, hasAnyFilesRecursive, removeEmptyDirectoryTree, resetExtractorCachesForPasswordChange, type ExtractArchiveFailureInfo, type ExtractProgressUpdate } from "./extractor"; +import { clearExtractResumeState, collectArchiveCleanupTargets, detectArchiveSignature, extractPackageArchives, findArchiveCandidates, resetExtractorCachesForPasswordChange, type ExtractArchiveFailureInfo, type ExtractProgressUpdate } from "./extractor"; import { validateFileAgainstManifest } from "./integrity"; import { classifyDiskError } from "./fs-error"; import { processVideoFile, resolveVideoTooling, stripDualLangMarker, hasDualLangMarker, isRemuxableVideoFile, type GermanAudioMode, type VideoProcessResult } from "./video-processor"; @@ -5883,45 +5883,53 @@ export class DownloadManager extends EventEmitter { return removed; } - private async cleanupRemainingArchiveArtifacts(packageDir: string, shouldAbort?: () => boolean): Promise { - if (this.settings.cleanupMode === "none") { - return 0; - } - if (shouldAbort?.()) { - return 0; - } - const candidates = await findArchiveCandidates(packageDir); - if (candidates.length === 0) { - return 0; - } - - let removed = 0; - const dirFilesCache = new Map(); - const targets = new Set(); - for (const sourceFile of candidates) { - if (shouldAbort?.()) { - return removed; - } - const dir = path.dirname(sourceFile); - let filesInDir = dirFilesCache.get(dir); - if (!filesInDir) { - try { - filesInDir = (await fs.promises.readdir(dir, { withFileTypes: true })) - .filter((entry) => entry.isFile()) - .map((entry) => entry.name); - } catch { - filesInDir = []; - } - dirFilesCache.set(dir, filesInDir); - } - for (const target of collectArchiveCleanupTargets(sourceFile, filesInDir)) { - targets.add(target); - } - } - - for (const targetPath of targets) { - if (shouldAbort?.()) { - return removed; + private async cleanupRemainingArchiveArtifacts(pkg: PackageEntry, shouldAbort?: () => boolean): Promise { + if (this.settings.cleanupMode === "none") { + return 0; + } + if (shouldAbort?.()) { + return 0; + } + const ownedPaths = new Map(); + for (const itemId of pkg.itemIds) { + const item = this.session.items[itemId]; + const rawPath = String(item?.targetPath || (item?.fileName ? path.join(pkg.outputDir, item.fileName) : "")).trim(); + if (!rawPath || !isArchiveLikePath(rawPath) || !isPathInsideDir(rawPath, pkg.outputDir)) { + continue; + } + const resolved = path.resolve(rawPath); + ownedPaths.set(pathKey(resolved), resolved); + } + if (ownedPaths.size === 0) { + return 0; + } + + let removed = 0; + const dirFiles = new Map(); + for (const ownedPath of ownedPaths.values()) { + const directory = path.dirname(ownedPath); + const directoryKey = pathKey(directory); + const files = dirFiles.get(directoryKey) || []; + files.push(path.basename(ownedPath)); + dirFiles.set(directoryKey, files); + } + const targets = new Set(); + for (const sourceFile of ownedPaths.values()) { + if (shouldAbort?.()) { + return removed; + } + const dir = path.dirname(sourceFile); + for (const target of collectArchiveCleanupTargets(sourceFile, dirFiles.get(pathKey(dir)) || [])) { + const resolved = path.resolve(target); + if (ownedPaths.has(pathKey(resolved))) { + targets.add(resolved); + } + } + } + + for (const targetPath of targets) { + if (shouldAbort?.()) { + return removed; } try { if (!await this.existsAsync(targetPath)) { @@ -5941,20 +5949,20 @@ export class DownloadManager extends EventEmitter { await this.renamePathWithExdevFallback(targetPath, candidate, { label: "mkv-move (Konflikt-Aufloesung)" }); moved = true; break; - } - if (moved) { - removed += 1; - } - continue; - } - await fs.promises.rm(toWindowsLongPathIfNeeded(targetPath), { force: true }); - removed += 1; - } catch { - } - } - - return removed; - } + } + if (moved) { + removed += 1; + } + continue; + } + await fs.promises.rm(toWindowsLongPathIfNeeded(targetPath), { force: true }); + removed += 1; + } catch { + } + } + + return removed; + } private hasDeferredPostProcessPending(packageId: string): boolean { if ((this.packageDeferredPostProcessTasks.get(packageId)?.size || 0) > 0) { @@ -14597,26 +14605,15 @@ export class DownloadManager extends EventEmitter { if (hasBlockingExtractError) { logger.info(`Deferred Archive-Cleanup uebersprungen: pkg=${pkg.name}, reason=extract_error`); } else { - const sourceAndTargetEqual = path.resolve(pkg.outputDir).toLowerCase() === path.resolve(pkg.extractDir).toLowerCase(); - if (!sourceAndTargetEqual) { - const candidates = await findArchiveCandidates(pkg.outputDir); - if (candidates.length > 0) { - const removed = await cleanupArchives(candidates, this.settings.cleanupMode, { shouldAbort }); - if (removed > 0) { - logger.info(`Deferred Archive-Cleanup: pkg=${pkg.name}, entfernt=${removed}`); - } - } - } - } - } - - if (this.settings.autoExtract && alreadyMarkedExtracted && failed === 0 && success > 0 && this.settings.cleanupMode !== "none" && !hasBlockingExtractError) { - throwIfAborted(); - const removedArchives = await this.cleanupRemainingArchiveArtifacts(pkg.outputDir, shouldAbort); - if (removedArchives > 0) { - logger.info(`Hybrid-Post-Cleanup entfernte Archive: pkg=${pkg.name}, entfernt=${removedArchives}`); - } - } + const sourceAndTargetEqual = path.resolve(pkg.outputDir).toLowerCase() === path.resolve(pkg.extractDir).toLowerCase(); + if (!sourceAndTargetEqual) { + const removed = await this.cleanupRemainingArchiveArtifacts(pkg, shouldAbort); + if (removed > 0) { + logger.info(`Deferred Archive-Cleanup: pkg=${pkg.name}, entfernt=${removed}`); + } + } + } + } if (extractedCount > 0 || alreadyMarkedExtracted) { throwIfAborted(); @@ -14636,23 +14633,22 @@ export class DownloadManager extends EventEmitter { } } - if ((extractedCount > 0 || alreadyMarkedExtracted) && failed === 0) { - throwIfAborted(); - await clearExtractResumeState(pkg.outputDir, packageId); - await clearExtractResumeState(pkg.outputDir); - } + if ((extractedCount > 0 || alreadyMarkedExtracted) && failed === 0) { + throwIfAborted(); + await clearExtractResumeState(pkg.outputDir, packageId); + await clearExtractResumeState(pkg.outputDir); + const archiveParents = new Set(); + for (const itemId of pkg.itemIds) { + const item = this.session.items[itemId]; + const itemPath = String(item?.targetPath || (item?.fileName ? path.join(pkg.outputDir, item.fileName) : "")).trim(); + if (itemPath && isArchiveLikePath(itemPath) && isPathInsideDir(itemPath, pkg.outputDir)) { + archiveParents.add(path.dirname(path.resolve(itemPath))); + } + } + await this.removeEmptyScopedParentChains(pkg.outputDir, archiveParents); + } - if ((extractedCount > 0 || alreadyMarkedExtracted) && failed === 0 && this.settings.cleanupMode === "delete") { - throwIfAborted(); - if (!(await hasAnyFilesRecursive(pkg.outputDir))) { - const removedDirs = await removeEmptyDirectoryTree(pkg.outputDir); - if (removedDirs > 0) { - logger.info(`Deferred leere Download-Ordner entfernt: pkg=${pkg.name}, dirs=${removedDirs}`); - } - } - } - - if (success > 0 && (pkg.status === "completed" || pkg.status === "failed")) { + if (success > 0 && (pkg.status === "completed" || pkg.status === "failed")) { throwIfAborted(); pkg.postProcessLabel = "Verschiebe Videos..."; this.emitState(); diff --git a/src/main/extractor.ts b/src/main/extractor.ts index fc84599..5c123f0 100644 --- a/src/main/extractor.ts +++ b/src/main/extractor.ts @@ -1863,6 +1863,9 @@ function handleDaemonLine(line: string): void { if (daemonCurrentRequest) { const req = daemonCurrentRequest; + if (req.terminationStarted) { + return; + } parseJvmLine(trimmed, req.onArchiveProgress, req.parseState, req.onOutput); failDaemonOutputCallback(req); } @@ -1912,12 +1915,15 @@ function startDaemon(layout: JvmExtractorLayout): boolean { child.stderr!.on("data", (chunk) => { const raw = String(chunk || ""); daemonOutput = appendLimited(daemonOutput, raw); - daemonStderrBuffer += raw; - const lines = daemonStderrBuffer.split(/\r?\n/); - daemonStderrBuffer = lines.pop() || ""; - for (const line of lines) { + daemonStderrBuffer += raw; + const lines = daemonStderrBuffer.split(/\r?\n/); + daemonStderrBuffer = lines.pop() || ""; + for (const line of lines) { if (daemonCurrentRequest) { const req = daemonCurrentRequest; + if (req.terminationStarted) { + continue; + } parseJvmLine(line, req.onArchiveProgress, req.parseState, req.onOutput); failDaemonOutputCallback(req); } diff --git a/tests/download-manager.test.ts b/tests/download-manager.test.ts index acc0f19..2714555 100644 --- a/tests/download-manager.test.ts +++ b/tests/download-manager.test.ts @@ -43,6 +43,19 @@ function writePackageOutputOwnerMarker(pkg: PackageEntry): void { })); } +function setExtractOutputRecords(pkg: PackageEntry, outputPaths: string[]): void { + pkg.outputProvenanceVersion = 1; + pkg.outputRecords = outputPaths.map((outputPath) => ({ + version: 1, + archivePath: path.join(pkg.outputDir, "source.zip"), + entryPath: path.relative(pkg.extractDir, outputPath).replace(/\\/g, "/"), + outputPath, + state: "complete", + disposition: "written" + })); + pkg.outputCount = outputPaths.length; +} + describe("runWithLimitedConcurrency", () => { it("processes the full batch without exceeding the configured worker count", async () => { let active = 0; @@ -12015,7 +12028,7 @@ describe("download manager", () => { expect(remainingPackage?.outputDir).toBe(outputDir); }); - it("removes link and sample artifacts from extracted output even when deferred post-processing has failures", async () => { + it("removes link and sample artifacts from extracted output even when deferred post-processing has failures", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "rd-dm-")); tempDirs.push(root); @@ -12060,12 +12073,16 @@ describe("download manager", () => { resumable: true, attempts: 1, lastError: "", - fullStatus: "Entpackt - Done (1s)", - createdAt, - updatedAt: createdAt - }; - - const manager = new DownloadManager( + fullStatus: "Entpackt - Done (1s)", + createdAt, + updatedAt: createdAt + }; + setExtractOutputRecords(session.packages[packageId], [ + path.join(extractDir, "episode.links.txt"), + path.join(extractDir, "sample", "sample.mkv") + ]); + + const manager = new DownloadManager( { ...defaultSettings(), token: "rd-token", @@ -12093,10 +12110,87 @@ describe("download manager", () => { ); expect(fs.existsSync(path.join(extractDir, "episode.links.txt"))).toBe(false); - expect(fs.existsSync(path.join(extractDir, "sample", "sample.mkv"))).toBe(false); - }); - - it("does not delete startup archives when any completed item has an extract error", async () => { + expect(fs.existsSync(path.join(extractDir, "sample", "sample.mkv"))).toBe(false); + }); + + it("deletes only package-owned archive members from a shared download root", async () => { + const root = fs.mkdtempSync(path.join(os.tmpdir(), "rd-shared-cleanup-")); + tempDirs.push(root); + const outputDir = path.join(root, "downloads"); + const extractDir = path.join(root, "extract", "owned-package"); + fs.mkdirSync(outputDir, { recursive: true }); + fs.mkdirSync(extractDir, { recursive: true }); + const ownedPaths = [path.join(outputDir, "owned.part1.rar"), path.join(outputDir, "owned.part2.rar")]; + const foreignPaths = [path.join(outputDir, "foreign.part1.rar"), path.join(outputDir, "foreign.part2.rar")]; + for (const archivePath of [...ownedPaths, ...foreignPaths]) { + fs.writeFileSync(archivePath, archivePath); + } + const extractedPath = path.join(extractDir, "episode.mkv"); + fs.writeFileSync(extractedPath, "video"); + const session = emptySession(); + const packageId = "owned-package"; + const createdAt = Date.now() - 20_000; + session.packageOrder = [packageId]; + session.packages[packageId] = { + id: packageId, + name: packageId, + outputDir, + extractDir, + status: "completed", + itemIds: ["owned-part-1", "owned-part-2"], + cancelled: false, + enabled: true, + createdAt, + updatedAt: createdAt + }; + for (let index = 0; index < ownedPaths.length; index += 1) { + const itemId = `owned-part-${index + 1}`; + session.items[itemId] = { + id: itemId, + packageId, + url: `https://example.com/${path.basename(ownedPaths[index])}`, + provider: "realdebrid", + status: "completed", + retries: 0, + speedBps: 0, + downloadedBytes: fs.statSync(ownedPaths[index]).size, + totalBytes: fs.statSync(ownedPaths[index]).size, + progressPercent: 100, + fileName: path.basename(ownedPaths[index]), + targetPath: ownedPaths[index], + resumable: true, + attempts: 1, + lastError: "", + fullStatus: "Entpackt - Done (1s)", + createdAt, + updatedAt: createdAt + }; + } + setExtractOutputRecords(session.packages[packageId], [extractedPath]); + const manager = new DownloadManager( + { + ...defaultSettings(), + outputDir, + extractDir: path.join(root, "extract"), + autoExtract: true, + autoRename4sf4sj: false, + collectMkvToLibrary: false, + removeLinkFilesAfterExtract: false, + removeSamplesAfterExtract: false, + enableIntegrityCheck: false, + cleanupMode: "delete" + }, + session, + createStoragePaths(path.join(root, "state")) + ); + + await (manager as any).runDeferredPostExtraction(packageId, session.packages[packageId], 1, 0, true, 1); + + expect(ownedPaths.map((archivePath) => fs.existsSync(archivePath))).toEqual([false, false]); + expect(foreignPaths.map((archivePath) => fs.existsSync(archivePath))).toEqual([true, true]); + }); + + it("does not delete startup archives when any completed item has an extract error", async () => { const root = fs.mkdtempSync(path.join(os.tmpdir(), "rd-dm-")); tempDirs.push(root); @@ -13087,14 +13181,16 @@ describe("download manager", () => { }, session, createStoragePaths(path.join(root, "state")) - ); - (manager as any).fileStabilizeMinAgeMs = 30_000; - - const expectedBase = "Test.Show.S02E05.Title.GERMAN.WS.720p.HDTV.x264-aWake"; + ); + (manager as any).fileStabilizeMinAgeMs = 30_000; + setExtractOutputRecords(session.packages[packageId], [scenePath]); + + const expectedBase = "Test.Show.S02E05.Title.GERMAN.WS.720p.HDTV.x264-aWake"; const renamedLibPath = path.join(mkvLibraryDir, `${expectedBase}.mkv`); const sceneLibPath = path.join(mkvLibraryDir, sceneName); - await (manager as any).autoRenameExtractedVideoFiles(extractDir, session.packages[packageId], undefined, true); + const outputScope = (manager as any).getPackageOutputScope(session.packages[packageId]); + await (manager as any).autoRenameExtractedVideoFiles(extractDir, outputScope, session.packages[packageId], undefined, true); await (manager as any).collectMkvFilesToLibrary(packageId, session.packages[packageId], undefined, false); expect(fs.existsSync(renamedLibPath)).toBe(true); @@ -13149,11 +13245,12 @@ describe("download manager", () => { }, session, createStoragePaths(path.join(root, "state")) - ); - (manager as any).fileStabilizeMinAgeMs = 30_000; - - const expectedBase = "Test.Show.S02E05.Title.GERMAN.WS.720p.HDTV.x264-aWake"; - await (manager as any).runDeferredPostExtraction(packageId, session.packages[packageId], 1, 0, true, 1); + ); + (manager as any).fileStabilizeMinAgeMs = 30_000; + + const expectedBase = "Test.Show.S02E05.Title.GERMAN.WS.720p.HDTV.x264-aWake"; + setExtractOutputRecords(session.packages[packageId], [path.join(epFolder, sceneName)]); + await (manager as any).runDeferredPostExtraction(packageId, session.packages[packageId], 1, 0, true, 1); expect(fs.existsSync(path.join(mkvLibraryDir, `${expectedBase}.mkv`))).toBe(true); expect(fs.existsSync(path.join(mkvLibraryDir, sceneName))).toBe(false); diff --git a/tests/notify-hooks.test.ts b/tests/notify-hooks.test.ts index 7f131f5..6336c3f 100644 --- a/tests/notify-hooks.test.ts +++ b/tests/notify-hooks.test.ts @@ -789,7 +789,21 @@ describe("authoritative run completion", () => { packageAItem.fullStatus = "Fertig"; packageA.status = "completed"; state.runOutcomes.set(packageAItem.id, "completed"); - state.packagePostProcessActive = 1; + const handlePackagePostProcessing = state.handlePackagePostProcessing.bind(state); + let releaseMainPostProcess = (): void => {}; + let markMainPostProcessEntered = (): void => {}; + const mainPostProcessGate = new Promise((resolve) => { + releaseMainPostProcess = resolve; + }); + const mainPostProcessEntered = new Promise((resolve) => { + markMainPostProcessEntered = resolve; + }); + vi.spyOn(state, "handlePackagePostProcessing").mockImplementation(async (...args: unknown[]) => { + const [packageId, signal] = args as [string, AbortSignal?]; + markMainPostProcessEntered(); + await mainPostProcessGate; + await handlePackagePostProcessing(packageId, signal); + }); let releaseCollection = (): void => {}; const collectionGate = new Promise((resolve) => { @@ -797,14 +811,14 @@ describe("authoritative run completion", () => { }); const collect = vi.spyOn(state, "collectMkvFilesToLibrary").mockImplementation(async () => collectionGate); const packageAMainPostProcess = state.runPackagePostProcessing(packageA.id); - await vi.waitFor(() => expect(state.packagePostProcessWaiters).toHaveLength(1)); + await mainPostProcessEntered; state.finishRun(); const packageB = addPackage(session, ["queued"], "active-run-package"); await manager.start(); expect(state.runPackageIds).toEqual(new Set([packageB.id])); - state.releasePostProcessSlot(); + releaseMainPostProcess(); await packageAMainPostProcess; await vi.waitFor(() => expect(collect).toHaveBeenCalled()); const deferredTasks = [...(state.packageDeferredPostProcessTasks.get(packageA.id) || [])];